{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"codemirror_mode":{"name":"ipython","version":3},"file_extension":".py","mimetype":"text/x-python","name":"python","nbconvert_exporter":"python","pygments_lexer":"ipython3","version":"3.10.0"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":113204,"databundleVersionId":13751849,"sourceType":"competition"}],"dockerImageVersionId":31153,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Chunk Ranking Pipeline for ACM ICAIF '25 Challenge\n\nThis notebook implements a two-stage agentic retrieval system, retrieval-rerank pipeline:\n1. **Hybrid Retrieval**: Combines sparse (BM25) and dense (embeddings) retrieval to narrow down candidates\n2. **Setwise Reranking**: Uses LLM-based setwise comparisons to produce final ranking\n\n## Overview\n\n**Task**: Given a financial question and text chunks from SEC filings, rank the top 5 most relevant chunks.\n\n**Pipeline**:\n```\nQuestion + chunks\n    ↓\n[1] Hybrid Retrieval → Top 60 chunks\n    ↓\n[2] Setwise Reranking → Top 5 chunks\n    ↓\nSubmission CSV\n```","metadata":{}},{"cell_type":"markdown","source":"## Setup and Imports","metadata":{}},{"cell_type":"code","source":"import asyncio\nimport csv\nimport json\nimport os\nimport re\nimport traceback\nfrom typing import Dict, List, Tuple\n\nimport boto3\nimport numpy as np\nfrom dotenv import load_dotenv\nfrom openai import AsyncOpenAI\nfrom tqdm.asyncio import tqdm\nfrom rank_bm25 import BM25Okapi\nfrom collections import defaultdict\nimport httpx\n\nload_dotenv()\n\nprint(\"✓ Imports loaded\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Initialize Clients\n\nWe use:\n- **AsyncOpenAI**: For LLM calls to AWS Bedrock (via OpenAI-compatible API)\n- **Boto3 Bedrock**: For structured output extraction using Amazon Nova","metadata":{}},{"cell_type":"code","source":"# AWS credentials check\nsession = boto3.Session()\ncredentials = session.get_credentials()\nprint(\"AWS credentials:\", \"Found\" if credentials else \"Not found\")\n\n# Initialize clients\nBEDROCK_TOKEN = os.environ.get('BEDROCK_TOKEN')\nif not BEDROCK_TOKEN:\n    raise ValueError(\"BEDROCK_TOKEN not found in environment variables\")\n\nclient = AsyncOpenAI(\n    api_key=BEDROCK_TOKEN,\n    base_url=\"https://bedrock-runtime.us-west-2.amazonaws.com/openai/v1\"\n)\n\nprint(\"✓ Clients initialized\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Embedding Model Setup\n\nFor dense retrieval, we use a vLLM-served embedding model with async batch processing.\nThe embedding model used is [Linq-Embed-Mistral](https://huggingface.co/Linq-AI-Research/Linq-Embed-Mistral)","metadata":{}},{"cell_type":"markdown","source":"### Setting Up vLLM Embedding Server (Optional but Recommended)\n\nIf you want to run the embedding model locally, follow these instructions based on your GPU setup.\n\n#### Prerequisites\n```bash\npip install vllm\n```\n\n#### Check Available GPUs\n```python\nimport torch\nprint(f\"CUDA available: {torch.cuda.is_available()}\")\nprint(f\"GPU count: {torch.cuda.device_count()}\")\nfor i in range(torch.cuda.device_count()):\n    print(f\"  GPU {i}: {torch.cuda.get_device_name(i)}\")\n```\n\n#### Option 1: Single GPU Setup (Simplest)\n\nFor 1 GPU, run this command in a terminal:\n\n```bash\nvllm serve Linq-AI-Research/Linq-Embed-Mistral \\\n    --host 0.0.0.0 \\\n    --port 8080 \\\n    --task embed \\\n    --max-model-len 8192 \\\n    --dtype float16 \\\n    --gpu-memory-utilization 0.9 \\\n    --max-num-batched-tokens 16384 \\\n    --max-num-seqs 512 \\\n    --disable-log-requests\n```\n\nThen set in the code below:\n```python\nEMBEDDING_MODEL_URL = 'http://localhost:8080/v1/embeddings'\nembedding_semaphore = asyncio.Semaphore(1)  # Only 1 server\n```\n\n#### Option 2: Production Setup with Systemd (Linux servers)\n\nFor persistent services that auto-restart, use systemd. Example for 4 GPUs:\n\n```bash\n# Create service for GPU 0\nsudo tee /etc/systemd/system/vllm-gpu0.service > /dev/null << 'EOF'\n[Unit]\nDescription=vLLM Embedding Service on GPU 0\nAfter=network.target\n\n[Service]\nType=simple\nUser=YOUR_USERNAME\nWorkingDirectory=/home/YOUR_USERNAME\nEnvironment=\"CUDA_VISIBLE_DEVICES=0\"\nExecStart=/usr/local/bin/vllm serve Linq-AI-Research/Linq-Embed-Mistral \\\n    --host 0.0.0.0 --port 8000 --task embed \\\n    --max-model-len 8192 --dtype float16 \\\n    --gpu-memory-utilization 0.9 --max-num-batched-tokens 16384 \\\n    --max-num-seqs 512 --disable-log-requests\nRestart=always\nRestartSec=10\n\n[Install]\nWantedBy=multi-user.target\nEOF\n\n# Repeat for other GPUs (ports 8001, 8002, 8003)\n# Then enable and start:\nsudo systemctl daemon-reload\nsudo systemctl enable vllm-gpu0.service\nsudo systemctl start vllm-gpu0.service\n\n# Check status\nsudo systemctl status vllm-gpu0.service\n```\n\nFull setup script available at: `embedding/setup_vllm.sh`\n\n#### Option 3: Use External Embedding API\n\nIf you don't have GPUs, you can use external embedding APIs:\n- OpenAI embeddings: `text-embedding-3-large`\n- Cohere embeddings\n- Voyage AI embeddings\n\nModify the `VLLMEmbeddingModel` class to use the appropriate API client.\n\n---\n\n**Note**: The configuration below assumes you have a vLLM server running. If not, the dense retrieval will be skipped and only BM25 (sparse) will be used.","metadata":{}},{"cell_type":"code","source":"class VLLMEmbeddingModel:\n    \"\"\"Async wrapper for vLLM embedding API with batch processing\"\"\"\n    \n    def __init__(self, api_url: str, model_name: str):\n        self.api_url = api_url\n        self.model_name = model_name\n    \n    async def encode(self, texts: List[str], client: httpx.AsyncClient) -> np.ndarray:\n        \"\"\"Get embeddings for a batch of texts\"\"\"\n        payload = {\n            \"model\": self.model_name,\n            \"input\": texts,\n            \"encoding_format\": \"float\"\n        }\n        \n        response = await client.post(self.api_url, json=payload, timeout=30.0)\n        response.raise_for_status()\n        result = response.json()\n        \n        embeddings = [item['embedding'] for item in result['data']]\n        return np.array(embeddings, dtype=np.float32)\n\n# Configure embedding model\nEMBEDDING_MODEL_NAME = \"Linq-AI-Research/Linq-Embed-Mistral\"\nEMBEDDING_MODEL_URL = 'http://10.50.1.6:8080/v1/embeddings'\nembedding_model = VLLMEmbeddingModel(EMBEDDING_MODEL_URL, EMBEDDING_MODEL_NAME)\n\nprint(f\"✓ Embedding model: {EMBEDDING_MODEL_NAME}\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Utility Functions\n\nHelper functions for:\n- Extracting rankings from LLM responses\n- Loading/saving data\n- Parsing question and chunks from messages","metadata":{}},{"cell_type":"code","source":"def load_data(filepath: str) -> List[Dict]:\n    \"\"\"Load JSONL data\"\"\"\n    print(f\"Loading data from: {filepath}\")\n    try:\n        with open(filepath, 'r', encoding='utf-8') as f:\n            data = [json.loads(line.strip()) for line in f]\n        print(f\"Loaded {len(data)} items\")\n        return data\n    except Exception as e:\n        print(f\"Error loading file: {e}\")\n        return []\n\n\ndef extract_question_and_chunks(content: str) -> Tuple[str, List[str], List[int]]:\n    \"\"\"Extract question and chunks from message content\"\"\"\n    question_start = content.find('Question:')\n    text_chunks_start = content.find('Text chunks:')\n    \n    if question_start == -1 or text_chunks_start == -1:\n        return None, [], []\n    \n    question = content[question_start + len('Question:'):text_chunks_start].strip()\n    \n    chunks = []\n    chunk_indices = []\n    chunk_pattern = r'\\[Chunk Index (\\d+)\\]\\s*(.*?)(?=\\[Chunk Index|\\nTask:|$)'\n    matches = re.findall(chunk_pattern, content, re.DOTALL)\n    \n    for match in matches:\n        orig_idx = int(match[0])\n        chunk_content = match[1].strip()\n        if '\\nTask:' in chunk_content:\n            chunk_content = chunk_content.split('\\nTask:')[0].strip()\n        if chunk_content:\n            chunks.append(chunk_content)\n            chunk_indices.append(orig_idx)\n    \n    return question, chunks, chunk_indices\n\n\ndef append_to_csv(entry: Dict, filename: str):\n    \"\"\"Append a single entry to CSV file\"\"\"\n    file_exists = os.path.isfile(filename)\n    \n    with open(filename, 'a', newline='', encoding='utf-8') as f:\n        writer = csv.writer(f)\n        if not file_exists:\n            writer.writerow(['sample_id', 'target_index'])\n        writer.writerow([entry['sample_id'], entry['target_index']])\n\n\ndef get_processed_ids(filename: str) -> set:\n    \"\"\"Get set of sample_ids that have already been processed\"\"\"\n    if not os.path.isfile(filename):\n        return set()\n    \n    processed = set()\n    with open(filename, 'r', encoding='utf-8') as f:\n        reader = csv.DictReader(f)\n        for row in reader:\n            processed.add(row['sample_id'])\n    return processed\n\nprint(\"✓ Utility functions loaded\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Stage 1: Hybrid Retrieval\n\nThe `HybridRetriever` combines three retrieval strategies:\n\n### 1. Query Enrichment\n- Expands query with domain-specific synonyms (e.g., \"FCF\" → \"free cash flow\", \"operating cash flow\")\n- Uses 90+ financial term mappings for better keyword recall\n\n### 2. Sparse Retrieval (BM25)\n- Traditional keyword-based ranking using BM25Okapi\n- Works well for exact terminology matches\n- Returns top 100 candidates\n\n### 3. Dense Retrieval (Embeddings)\n- Semantic similarity using embedding model\n- Captures conceptual relevance beyond keywords\n- Async batch processing for efficiency\n- Returns top 100 candidates\n\n### 4. Reciprocal Rank Fusion (RRF)\n- Merges sparse and dense results using the formula: `score = Σ(1 / (k + rank))`\n- `k=60` balances contribution from both methods\n- Returns top 60 fused candidates for reranking","metadata":{}},{"cell_type":"code","source":"class HybridRetriever:\n    \"\"\"Combines sparse (BM25) and dense (embeddings) retrieval with RRF fusion\"\"\"\n    \n    def __init__(\n        self,\n        embedding_model,\n        embedding_semaphore: asyncio.Semaphore,\n        rrf_k: int = 60,\n        vllm_batch_size: int = 32\n    ):\n        self.embedding_model = embedding_model\n        self.embedding_semaphore = embedding_semaphore\n        self.rrf_k = rrf_k\n        self.vllm_batch_size = vllm_batch_size\n        \n        # Domain-specific financial synonyms for query expansion\n        self.domain_synonyms = {\n            # Compensation & Equity\n            'equity award burn rate': ['dilution rate', 'share pool availability'],\n            'restricted stock units': ['rsu', 'equity grants', 'stock grants'],\n            'executive compensation': ['executive pay', 'management compensation'],\n            'share repurchase': ['stock buyback', 'share buyback'],\n            'dividend policy': ['dividend strategy', 'payout policy'],\n            \n            # Financial Metrics\n            'profitability': ['operating margin', 'ebitda margin'],\n            'revenue': ['sales', 'top-line', 'turnover'],\n            'recurring revenue': ['subscription revenue', 'contracted revenue'],\n            'guidance': ['outlook', 'forecast', 'projection'],\n            'margin': ['operating margin', 'profit margin'],\n            'fcf': ['free cash flow', 'operating cash flow'],\n            'capex': ['capital expenditures', 'capital spending'],\n            'inventory': ['working capital', 'stock levels'],\n            'combined ratio': ['underwriting ratio', 'loss ratio'],\n            \n            # Strategy & Operations\n            'innovation cycles': ['technology advancement', 'product development cycles'],\n            'market competitiveness': ['competitive positioning', 'market position'],\n            'geographic expansion': ['international expansion', 'market expansion'],\n            'customer engagement': ['user engagement', 'customer retention'],\n            'supply chain': ['logistics', 'distribution', 'procurement'],\n            \n            # Risk & Governance\n            'dependency risks': ['concentration risk', 'customer concentration'],\n            'foreign exchange': ['fx', 'currency', 'forex'],\n            'scenario analyses': ['stress testing', 'risk assessment'],\n            'climate risks': ['environmental risks', 'esg risks'],\n            'geopolitical': ['political risk', 'country risk'],\n            'governance': ['corporate governance', 'board oversight'],\n            \n            # ESG\n            'esg': ['environmental social governance', 'sustainability'],\n            \n            # Customer/Financing\n            'customer financing': ['vendor financing', 'lease program'],\n            'penetration rate': ['adoption rate', 'market penetration'],\n            'retention': ['retention rate', 'churn'],\n            'booking': ['reservations', 'orders'],\n            'occupancy': ['occupancy rate', 'utilization'],\n        }\n    \n    def _enrich_query(self, question: str) -> str:\n        \"\"\"Expand query with domain synonyms for better keyword recall\"\"\"\n        query_lower = question.lower()\n        enriched_terms = [question]\n        \n        # Match multi-word phrases (longer matches first)\n        matched_terms = set()\n        sorted_terms = sorted(self.domain_synonyms.keys(), key=len, reverse=True)\n        \n        for term in sorted_terms:\n            if term in query_lower and term not in matched_terms:\n                # Add top 2 most relevant synonyms\n                synonyms = self.domain_synonyms[term][:2]\n                enriched_terms.extend(synonyms)\n                matched_terms.add(term)\n                \n                # Mark overlapping terms to avoid duplication\n                for other_term in sorted_terms:\n                    if other_term != term and (other_term in term or term in other_term):\n                        matched_terms.add(other_term)\n        \n        result = ' '.join(enriched_terms)\n        \n        # Safety: limit expansion to 2x original length\n        if len(result) > len(question) * 2:\n            enriched_terms = enriched_terms[:5]\n            result = ' '.join(enriched_terms)\n        \n        return result\n    \n    def _sparse_retrieval(\n        self,\n        query: str,\n        chunks: List[str],\n        chunk_indices: List[int],\n        top_k: int = 100\n    ) -> List[Tuple[int, float]]:\n        \"\"\"BM25 keyword-based retrieval\"\"\"\n        if not chunks:\n            return []\n        \n        top_k = min(top_k, len(chunks))\n        \n        # Tokenize and build BM25 index\n        tokenized_chunks = [chunk.lower().split() for chunk in chunks]\n        bm25 = BM25Okapi(tokenized_chunks)\n        \n        # Score query\n        query_tokens = query.lower().split()\n        scores = bm25.get_scores(query_tokens)\n        \n        # Return top-k with scores\n        ranked = sorted(\n            zip(chunk_indices, scores),\n            key=lambda x: x[1],\n            reverse=True\n        )\n        \n        return ranked[:top_k]\n    \n    async def _get_embeddings_concurrently(self, texts: List[str]) -> np.ndarray:\n        \"\"\"Batch texts and get embeddings concurrently\"\"\"\n        if not texts:\n            return np.array([], dtype=np.float32)\n\n        all_embeddings = []\n        \n        async with httpx.AsyncClient() as http_client:\n            tasks = []\n            batch_size = self.vllm_batch_size\n            \n            for i in range(0, len(texts), batch_size):\n                batch = texts[i : i + batch_size]\n                \n                async def get_batch_embeddings(b):\n                    async with self.embedding_semaphore:\n                        return await self.embedding_model.encode(b, http_client)\n                \n                tasks.append(get_batch_embeddings(batch))\n\n            results = await asyncio.gather(*tasks)\n\n            for batch_embeddings in results:\n                all_embeddings.append(batch_embeddings)\n                    \n        if all_embeddings:\n            return np.vstack(all_embeddings)\n        return np.array([], dtype=np.float32)\n    \n    async def _dense_retrieval(\n        self,\n        query: str,\n        chunks: List[str],\n        chunk_indices: List[int],\n        top_k: int = 100\n    ) -> List[Tuple[int, float]]:\n        \"\"\"Embedding-based semantic retrieval\"\"\"\n        if not chunks:\n            return []\n        \n        top_k = min(top_k, len(chunks))\n        \n        try:\n            # Get embeddings for query + all chunks\n            embeddings = await self._get_embeddings_concurrently([query] + chunks)\n        except Exception as e:\n            print(f\"Embedding API failed: {e}\")\n            return []\n\n        if embeddings.shape[0] < 1:\n            return []\n            \n        # Split: first is query, rest are chunks\n        query_emb = embeddings[0]\n        chunk_embs = embeddings[1:]\n        \n        # Compute cosine similarities\n        similarities = np.dot(chunk_embs, query_emb)\n        \n        # Get top-k\n        ranked_idx = np.argsort(similarities)[::-1][:top_k]\n        ranked = [(chunk_indices[i], float(similarities[i])) for i in ranked_idx]\n        \n        return ranked\n    \n    def _reciprocal_rank_fusion(\n        self,\n        sparse_results: List[Tuple[int, float]],\n        dense_results: List[Tuple[int, float]],\n        k: int = 60\n    ) -> List[int]:\n        \"\"\"Fuse sparse and dense results using RRF\"\"\"\n        rrf_scores = defaultdict(float)\n        \n        # Add sparse scores\n        for rank, (idx, score) in enumerate(sparse_results, start=1):\n            rrf_scores[idx] += 1.0 / (k + rank)\n        \n        # Add dense scores\n        for rank, (idx, score) in enumerate(dense_results, start=1):\n            rrf_scores[idx] += 1.0 / (k + rank)\n        \n        # Sort by RRF score\n        ranked = sorted(rrf_scores.items(), key=lambda x: x[1], reverse=True)\n        return [idx for idx, score in ranked]\n    \n    async def retrieve(\n        self,\n        question: str,\n        chunks: List[str],\n        chunk_indices: List[int],\n        k: int = 30\n    ) -> Tuple[List[str], List[int]]:\n        \"\"\"Full hybrid retrieval pipeline\"\"\"\n        if len(chunks) <= k:\n            print(f\"Only {len(chunks)} chunks available, returning all\")\n            return chunks, chunk_indices\n        \n        print(f\"Hybrid retrieval: {len(chunks)} chunks → {k}\")\n        \n        # Step 1: Enrich query with domain synonyms\n        enriched_query = self._enrich_query(question)\n        print(f\"Enriched query: {enriched_query[:100]}...\")\n        \n        # Step 2: Sparse retrieval (BM25)\n        sparse_top_k = min(100, len(chunks))\n        sparse_results = self._sparse_retrieval(\n            enriched_query, chunks, chunk_indices, top_k=sparse_top_k\n        )\n        print(f\"Sparse (BM25): {len(sparse_results)} results\")\n        \n        # Step 3: Dense retrieval (embeddings)\n        dense_top_k = min(100, len(chunks))\n        dense_results = await self._dense_retrieval(\n            question, chunks, chunk_indices, top_k=dense_top_k\n        )\n        print(f\"Dense (embeddings): {len(dense_results)} results\")\n        \n        # Step 4: Fusion with RRF\n        fused_indices = self._reciprocal_rank_fusion(\n            sparse_results, dense_results, k=self.rrf_k\n        )\n        print(f\"Fused: {len(fused_indices)} unique candidates\")\n        \n        # Take top-k from fused results\n        final_indices = fused_indices[:k]\n        \n        # Get corresponding chunks\n        idx_to_chunk = {idx: chunk for idx, chunk in zip(chunk_indices, chunks)}\n        final_chunks = [idx_to_chunk[idx] for idx in final_indices if idx in idx_to_chunk]\n        \n        return final_chunks, final_indices\n\nprint(\"✓ HybridRetriever loaded\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Stage 2: Setwise Reranking\n\nThe `SetwiseReranker` implements the approach from \"[A Setwise Approach for Effective and Highly Efficient Zero-shot Ranking with Large Language Models](https://arxiv.org/abs/2310.09497)\" (SIGIR '24).\n\n### Key Concepts\n\n**Setwise Comparison**: Instead of comparing documents pairwise (A vs B, B vs C, ...), we compare multiple documents at once:\n- Given chunks A, B, C, D, E, ask the LLM: \"Which is most relevant?\"\n- This reduces the number of LLM calls significantly\n\n**Heapsort Algorithm**: \n- Efficiently finds top-k items using setwise comparisons\n- Complexity: O(k × log_c(N)) where c is `compare_size`\n\n**How it works**:\n1. Process candidates in batches of 5\n2. For each batch, find the most relevant chunk\n3. Compare batch winners to find overall best\n4. Remove the best and repeat k times\n\nThis is much more efficient than listwise reranking (1 call for all 30 chunks, but may miss nuances) or pairwise comparisons (435 calls for 30 chunks).","metadata":{}},{"cell_type":"code","source":"class SetwiseReranker:\n    \"\"\"Setwise reranking using heapsort with c-way comparisons\"\"\"\n    \n    def __init__(\n        self,\n        client: AsyncOpenAI,\n        semaphore: asyncio.Semaphore,\n        compare_size: int = 5,\n        model: str = \"deepseek.v3-v1:0\"\n    ):\n        self.client = client\n        self.semaphore = semaphore\n        self.compare_size = compare_size\n        self.model = model\n        \n        if compare_size < 2:\n            raise ValueError(\"compare_size must be at least 2\")\n    \n    async def _setwise_compare(\n        self,\n        question: str,\n        chunks: List[str],\n        chunk_indices: List[int]\n    ) -> int:\n        \"\"\"\n        Compare multiple chunks at once and return the index of the most relevant one.\n        Uses LLM to determine which chunk is most relevant to the question.\n        \"\"\"\n        if len(chunks) != len(chunk_indices):\n            raise ValueError(\"chunks and chunk_indices must have same length\")\n        \n        # Create prompt with labeled chunks (A, B, C, D, E)\n        prompt = f\"\"\"Given a query, which of the following chunk is more relevant to the query?\n\nQuery: {question}\n\n\"\"\"\n        for i, (chunk, idx) in enumerate(zip(chunks, chunk_indices)):\n            label = chr(65 + i)  # A, B, C, D, E...\n            prompt += f\"Chunk {label}: {chunk}\\n\\n\"\n        \n        prompt += \"Output only the chunk label of the most relevant chunk, only a single letter (A, B, C, etc.):\"\n        \n        system_message = {\n            \"role\": \"system\",\n            \"content\": \"\"\"You are a helpful financial analyst.\nGiven a question about a company and multiple chunks from the financial documents about the company, your task is to find the most relevant chunk to the question.\n\nUse this relevance hierarchy when ranking chunks:\n\nTier 1 - Direct Answers (Highest Priority):\n- Chunks that explicitly address the question with specific data, policies, or direct statements\n- Exact keyword matches for the topic (e.g., \"dividend policy,\" \"ESG,\" \"free cash flow\")\n- Direct quotes from management discussing the specific topic\n- Specific financial figures or metrics requested in the question\n\nTier 2 - Strong Contextual Relevance:\n- Related financial metrics that inform or support the answer\n- Management discussion of broader themes that encompass the question topic\n- Risk disclosures or business factors that directly relate to the question\n- Forward-looking statements and guidance relevant to the topic\n\nTier 3 - Supporting Information:\n- Background information that provides important context for understanding the answer\n- Historical trends or comparative data that illuminate the topic\n- Industry or market factors that affect the specific area being questioned\n- Explanatory text that clarifies related financial data\n\nTier 4 - Weak/Tangential (Lowest Priority):\n- General company information with minimal connection to the question\n- Unrelated financial data that doesn't inform the specific topic\n- Boilerplate text without substantive relevance\n\nPrioritize chunk from higher tiers, but consider that a comprehensive answer may require information from multiple tiers. Focus on chunk that would be most valuable for providing a thorough, accurate response to the specific question asked.\"\"\"\n        }\n        \n        try:\n            async with self.semaphore:\n                response = await self.client.chat.completions.create(\n                    messages=[system_message, {\"role\": \"user\", \"content\": prompt}],\n                    model=self.model,\n                    temperature=0.0,\n                    reasoning_effort='high'\n                )\n                \n                result = response.choices[0].message.content.strip().upper()\n                \n                # Parse the letter response (A, B, C, etc.)\n                for i, label in enumerate([chr(65 + j) for j in range(len(chunks))]):\n                    if result and result[-1] == label:\n                        return chunk_indices[i]\n                \n                # Fallback: return first if parsing fails\n                return chunk_indices[0]\n                \n        except Exception as e:\n            print(f\"Error in setwise comparison: {e}\")\n            return chunk_indices[0]\n    \n    async def _setwise_heapsort(\n        self,\n        question: str,\n        chunks: List[str],\n        chunk_indices: List[int],\n        k: int\n    ) -> List[int]:\n        \"\"\"\n        Heap sort with setwise comparisons and parallel batch processing.\n        Finds top-k items efficiently using c-way comparisons.\n        \"\"\"\n        print(f\"Setwise heap sort: top-{k}, compare_size={self.compare_size}\")\n        \n        current_ranking = list(zip(chunks, chunk_indices))\n        result = []\n        \n        for iteration in range(min(k, len(current_ranking))):\n            remaining = current_ranking[iteration:]\n            \n            if len(remaining) <= 1:\n                result.extend([idx for _, idx in remaining])\n                break\n            \n            # Process in batches of compare_size in parallel\n            batch_tasks = []\n            batches = []\n            for i in range(0, len(remaining), self.compare_size):\n                batch = remaining[i:i + self.compare_size]\n                batch_chunks = [item[0] for item in batch]\n                batch_indices = [item[1] for item in batch]\n                \n                batches.append(batch)\n                batch_tasks.append(\n                    self._setwise_compare(question, batch_chunks, batch_indices)\n                )\n            \n            # Execute all batch comparisons in parallel\n            best_indices = await asyncio.gather(*batch_tasks)\n            \n            # Build candidates from batch winners\n            candidates = []\n            for batch, best_idx in zip(batches, best_indices):\n                best_item = next(item for item in batch if item[1] == best_idx)\n                candidates.append(best_item)\n            \n            # If multiple batch winners, compare them to find overall best\n            if len(candidates) > 1:\n                cand_chunks = [item[0] for item in candidates]\n                cand_indices = [item[1] for item in candidates]\n                best_idx = await self._setwise_compare(\n                    question, cand_chunks, cand_indices\n                )\n                best_item = next(item for item in candidates if item[1] == best_idx)\n            else:\n                best_item = candidates[0]\n            \n            result.append(best_item[1])\n            \n            # Remove best_item from remaining candidates\n            current_ranking = [item for item in current_ranking if item[1] != best_item[1]]\n            \n            print(f\"Iteration {iteration + 1}: found rank {iteration + 1}\")\n        \n        return result\n    \n    async def rerank(\n        self,\n        question: str,\n        chunks: List[str],\n        chunk_indices: List[int],\n        k: int = 5\n    ) -> List[int]:\n        \"\"\"Rerank using setwise heapsort approach\"\"\"\n        return await self._setwise_heapsort(question, chunks, chunk_indices, k)\n\nprint(\"✓ SetwiseReranker loaded\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Main Processing Pipeline\n\nThe `process_item` function orchestrates the full pipeline:\n1. Extract question and chunks from the message\n2. Run hybrid retrieval → 30 candidates\n3. Run setwise reranking → 5 final chunks\n4. Return ranked indices for submission","metadata":{}},{"cell_type":"code","source":"async def process_item(\n    messages: List[Dict],\n    query_id: str,\n    retriever: HybridRetriever,\n    reranker: SetwiseReranker\n) -> List[int]:\n    \"\"\"Process single item through the full pipeline\"\"\"\n    try:\n        content = messages[0].get('content', '')\n        question, chunks, chunk_indices = extract_question_and_chunks(content)\n        \n        if not question or not chunks:\n            print(f\"Could not extract data for {query_id}\")\n            return list(range(5))\n        \n        print(f\"\\n{'='*60}\")\n        print(f\"Processing {query_id}: {len(chunks)} chunks\")\n        print(f\"Question: {question[:100]}...\")\n        print(f\"{'='*60}\")\n \n        # Stage 1: Hybrid Retrieval (60 candidates)\n        retrieved_chunks, retrieved_indices = await retriever.retrieve(\n            question, chunks, chunk_indices, k=60\n        )\n        \n        # Stage 2: Setwise Reranking (top 5)\n        final_ranking = await reranker.rerank(\n            question,\n            retrieved_chunks, \n            retrieved_indices, \n            k=5\n        )\n        \n        print(f\"Final ranking: {final_ranking}\")\n        \n        return final_ranking\n        \n    except Exception as e:\n        traceback.print_exc()\n        print(f\"Error processing {query_id}: {e}\")\n        return list(range(5))\n\nprint(\"✓ Processing pipeline loaded\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Evaluation Function\n\nProcesses all samples in the eval set with:\n- **Incremental saving**: Results saved after each sample\n- **Resume capability**: Skips already-processed samples if run is interrupted\n- **CSV output**: Direct Kaggle submission format","metadata":{}},{"cell_type":"code","source":"async def evaluate(\n    data_path: str,\n    retriever: HybridRetriever,\n    reranker: SetwiseReranker,\n    output_file: str\n):\n    \"\"\"Run evaluation with incremental saving and resume capability\"\"\"\n    data = load_data(data_path)\n    if not data:\n        return\n    \n\n    # Create new file with header\n    with open(output_file, 'w', newline='', encoding='utf-8') as f:\n        writer = csv.writer(f)\n        writer.writerow(['sample_id', 'target_index'])\n    \n    # Process each item sequentially\n    for item in tqdm(data, desc=\"Processing samples\"):\n        result = await process_item(\n            item['messages'],\n            item['_id'],\n            retriever,\n            reranker\n        )\n        \n        if result:\n            # Save top 5 results\n            for rank, doc_idx in enumerate(result[:5]):\n                entry = {\n                    'sample_id': item['_id'], \n                    'target_index': doc_idx\n                }\n                append_to_csv(entry, output_file)\n    \n    print(f\"\\n✓ Evaluation complete! Results saved to: {output_file}\")\n\nprint(\"✓ Evaluation function loaded\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## Initialize Components and Run Evaluation\n\nConfigure and run the full pipeline on the eval dataset.","metadata":{}},{"cell_type":"code","source":"# Configure concurrency\nllm_semaphore = asyncio.Semaphore(20)  # Max 20 concurrent LLM calls\nembedding_semaphore = asyncio.Semaphore(4)  # Max 4 concurrent embedding calls\n\n# Initialize retriever\nretriever = HybridRetriever(\n    embedding_model=embedding_model,\n    embedding_semaphore=embedding_semaphore,\n    rrf_k=60,\n    vllm_batch_size=32\n)\n\n# Initialize reranker\nreranker = SetwiseReranker(\n    client=client,\n    semaphore=llm_semaphore,\n    compare_size=5,\n    model=\"deepseek.v3-v1:0\"\n)\n\nprint(f\"\\n{'='*60}\")\nprint(f\"Configuration:\")\nprint(f\"  Retriever: HybridRetriever (BM25 + Embeddings + RRF)\")\nprint(f\"  Reranker: SetwiseReranker (heapsort, compare_size=5)\")\nprint(f\"  Model: {reranker.model}\")\nprint(f\"  Embedding: {EMBEDDING_MODEL_NAME}\")\nprint(f\"{'='*60}\\n\")","metadata":{},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Configure paths\nDATA_PATH = \"./output/chunk_ranking_kaggle_eval.jsonl\"\nOUTPUT_FILE = f\"submission_chunk_ranking.csv\"\n\n# Run evaluation\nawait evaluate(DATA_PATH, retriever, reranker, OUTPUT_FILE)\n\nprint(f\"\\n{'='*60}\")\nprint(f\"✓ Complete! Output: {OUTPUT_FILE}\")\nprint(f\"{'='*60}\")","metadata":{},"outputs":[],"execution_count":null}]}