diff --git a/.env.example b/.env.example index 9945ff73..f2210431 100644 --- a/.env.example +++ b/.env.example @@ -7,7 +7,7 @@ PRODUCTION=false # ===================================================================== # Agent & AI Provider Settings # ===================================================================== -# Agent Provider Options: openai, julep +# Agent Provider Options: openai AGENT_PROVIDER=openai # Microsoft Azure Foundry (OpenAI-compatible) @@ -20,10 +20,6 @@ OPENAI_API_KEY=your-openai-api-key OPENAI_MODEL=gpt-4o-mini OPENAI_BASE_URL= -# Julep Settings (when AGENT_PROVIDER=julep or RAG_BACKEND=julep) -# JULEP_API_KEY=your-julep-api-key -# JULEP_MODEL=claude-3.5-sonnet - # ===================================================================== # RAG & Embedding Settings # ===================================================================== diff --git a/Dockerfile.mcp b/Dockerfile.mcp new file mode 100644 index 00000000..67369b64 --- /dev/null +++ b/Dockerfile.mcp @@ -0,0 +1,14 @@ +FROM python:3.13-slim + +WORKDIR /app + +COPY mcp_server/requirements.txt mcp_server/requirements.txt +RUN pip install --no-cache-dir -r mcp_server/requirements.txt + +COPY config.py . +COPY service/ service/ +COPY mcp_server/ mcp_server/ + +EXPOSE 8002 + +CMD ["python", "-m", "mcp_server"] diff --git a/config.py b/config.py index e79abf40..1b150104 100644 --- a/config.py +++ b/config.py @@ -72,22 +72,6 @@ def OPENAI_EMBEDDING_MODEL(self) -> str: or os.getenv("OPENAI_EMBEDDING_MODEL", "text-embedding-3-small") ) - # Julep Provider - @property - def JULEP_API_KEY(self) -> str | None: - return os.getenv("JULEP_API_KEY") - - @property - def JULEP_MODEL(self) -> str: - return ( - os.getenv("AGENT_MODEL") - or os.getenv("JULEP_MODEL", "claude-3.5-sonnet") - ) - - @property - def JULEP_ENVIRONMENT(self) -> str: - return os.getenv("JULEP_ENVIRONMENT", "production") - # ------------------------------------------------------------------ # Local Sentence Transformers Settings # ------------------------------------------------------------------ diff --git a/data/raw/2026-07-29T15-39-50Z/default/www.moneycontrol.com/direct.html b/data/raw/2026-07-29T15-39-50Z/default/www.moneycontrol.com/direct.html new file mode 100644 index 00000000..654703e1 --- /dev/null +++ b/data/raw/2026-07-29T15-39-50Z/default/www.moneycontrol.com/direct.html @@ -0,0 +1,5854 @@ + Double whammy: Bengaluru airport cab crunch deepens as BluSmart, Refex pull out + +
+ + + + +
Personal Loan Offers
Personal Loan Offers
+ +
+ + +

Double whammy: Bengaluru airport cab crunch deepens as BluSmart, Refex pull out

Passengers travelling to and from Kempegowda International Airport face a double whammy as two electric cab operators, BluSmart and Refex eVeelz, halt services.
April 18, 2025 / 19:05 IST
+
Bengaluru airport
+ + + + + + + + + + +

Passengers using Kempegowda International Airport (KIA) in Bengaluru, already grappling with a shortage of Ola and Uber cabs, are now facing a double-whammy as two electric cab operators—BluSmart and Refex eVeelz—have ceased operations.

BluSmart, a ride-hailing startup linked to Gensol Engineering, has abruptly stopped accepting new ride bookings across major parts of Delhi-NCR, Mumbai, and Bengaluru. This move follows a probe by the Securities and Exchange Board of India (SEBI) into Gensol Engineering, raising serious questions about BluSmart’s viability.

Also readGensol crisis: BluSmart suspends cab bookings via app in parts of NCR, Mumbai and Bengaluru

Meanwhile, Refex Green Mobility Ltd (RGML), a subsidiary of Refex Industries, has confirmed that it is streamlining its operations in Bengaluru by phasing out its airport-based EV taxi services. These services were run via Refex EV Fleet Services. “The move aims to focus on scaling electric mobility solutions for enterprise and institutional clients through long-term partnerships and stable demand,” the company said.

Interestingly, both Refex and BluSmart have a linkage. Earlier this year, Refex Green Mobility had announced a Rs 315 crore deal to acquire 2,997 electric vehicles from Gensol Engineering, with plans to lease them to BluSmart. However, Refex later withdrew the proposal, citing “evolving commitments” on both sides.

Also readRefex Green withdraws plan to takeover Gensol's 3,000 EVs, cites 'challenges' to conclude deal

Travellers hit

Kempegowda International Airport handled 41.88 million travellers in the last financial year, up from 37.53 million the previous year—an 11.6 percent increase. Despite this growth, the availability of airport cabs continues to be a challenge, with frequent complaints of long waits and inadequate fleets at the designated Ola-Uber pick-up zones.

Bangalore International Airport Ltd (BIAL), which operates KIA, was unavailable to comment.

Also readBluSmart takes on Uber, Ola, launches cab pickup zone at Bengaluru airport

Sources said the exit of both Refex and BluSmart is likely to compound the shortage. “BIAL tied up with Refex and commenced services in June 2024. They had also introduced pink taxis driven by women, exclusively for female passengers. BIAL had ended partnerships with earlier operators like Meru and Mega in favour of Refex. Now, with both Refex and BluSmart halting operations, it’s causing inconvenience to passengers,” a source said.

Another source added that airport taxis operated through BIAL’s official 'BLR Pulse' mobile application were mostly linked to Refex, which had a 100 percent electric fleet. "Many passengers who couldn't find cabs via Ola and Uber relied on BluSmart and Refex, and are now the most affected. These services were also environmentally friendly," the source added.

In March 2023, BluSmart Mobility launched a dedicated cab pickup zone at Bengaluru airport, following Ola and Uber. Other aggregators like Rapido, Namma Yatri, QuickRide and Shoffr do not have a dedicated pickup zone at the airport. However, dedicated airport taxi services such as Refex and KSTDC are allowed to pick up passengers directly in front of the terminals.

However, a BIAL official said that passengers can continue to book rides through Ola, Uber, Karnataka State Tourism Development Corporation (KSTDC), WTi rental cabs, as well as BMTC’s Vayu Vajra and KSRTC’s inter-city Flybus services.

In a statement, Refex said: “Refex continues to offer high-quality mobility services through its integrated B2B and B2B2C EV platform, featuring company-owned EVs, trained driver partners, and a tech-driven operations framework. Since its launch in March 2023 with 24 vehicles, Refex’s Green Mobility vertical has grown to nearly 1,300 EVs. We serve Bengaluru, Chennai, Hyderabad, and Mumbai across use cases such as employee transport, corporate rentals, airport transfers, and partnerships with leading ride-hailing platforms.”

Also readBengaluru airport launches electric taxis for passengers + +

+
+
Christin Mathew Philip is an Assistant editor at moneycontrol.com. Based in Bengaluru, he writes on mobility, infrastructure and start-ups. He is a Ramnath Goenka excellence in journalism awardee. You can find him on Twitter here: twitter.com/ChristinMP_
+ +

Discover the latest Business News, Sensex, and Nifty updates. Obtain Personal Finance insights, tax queries, and expert opinions on Moneycontrol or download the Moneycontrol App to stay updated! +

+
+
+
+
+
+
+
+ +

Subscribe to Tech Newsletters

  • On Saturdays

    Find the best of Al News in one place, specially curated for you every + weekend. +

  • Daily-Weekdays

    Stay on top of the latest tech trends and biggest startup news. +

+ + + + + +
+ +
+
+ +
+ + +

Advisory Alert:

It has come to our attention that certain individuals are representing themselves as affiliates of Moneycontrol and soliciting funds on the false promise of assured returns on their investments. We wish to reiterate that Moneycontrol does not solicit funds from investors and neither does it promise any assured returns. In case you are approached by anyone making such claims, please write to us at grievanceofficer@nw18.com or call on 02268882347
+ + \ No newline at end of file diff --git a/dev.sh b/dev.sh index a158fce8..25527fed 100755 --- a/dev.sh +++ b/dev.sh @@ -57,6 +57,7 @@ if [ "$DO_INSTALL" = true ]; then $VENV_PYTHON -m pip install -r server/requirements.txt \ -r embedding_server/requirements.txt \ -r pipeline/requirements.txt \ + -r mcp_server/requirements.txt \ pytest pytest-asyncio email-validator azure-storage-blob if [ -d "frontend" ] && [ -f "frontend/package.json" ]; then echo -e "${CYAN}📦 Installing frontend npm dependencies...${NC}" @@ -81,7 +82,7 @@ PIDS=() cleanup() { trap - INT TERM EXIT - echo -e "\n${CYAN}🛑 Stopping local servers...${NC}" + echo -e "\n${CYAN}D Stopping local servers...${NC}" for pid in "${PIDS[@]}"; do if kill -0 "$pid" 2>/dev/null; then kill -TERM "$pid" 2>/dev/null || true @@ -104,6 +105,13 @@ echo -e "${CYAN}🧠 Starting Embedding Server on http://localhost:8001...${NC}" $VENV_PYTHON -m uvicorn embedding_server.app:app --host 0.0.0.0 --port 8001 & PIDS+=($!) +# Start MCP Server (SSE) in background +if $VENV_PYTHON -c "import mcp" 2>/dev/null; then + echo -e "${CYAN}🔌 Starting MCP Server (SSE) on http://localhost:8002...${NC}" + $VENV_PYTHON -c "from mcp_server.app import mcp; mcp.run(transport='sse', port=8002)" & + PIDS+=($!) +fi + # Start Web Backend Server with Hot-Reloading in background echo -e "${CYAN}🌐 Starting Web Backend Server (hot-reloading) on http://localhost:8000...${NC}" $VENV_PYTHON -m uvicorn server.app:app --reload --host 0.0.0.0 --port 8000 & diff --git a/docker-compose.yml b/docker-compose.yml index 4d46a69b..dd9af9b6 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -43,6 +43,27 @@ services: networks: - default + mcp-server: + build: + context: . + dockerfile: Dockerfile.mcp + ports: + - "8002:8002" + environment: + EMBEDDING_SERVICE_URL: http://embedding-server:8001 + DB_URL: ${DB_URL:-mongodb://mongo:27017/evolution} + ARTICLE_STORE_BACKEND: file + RAG_BACKEND: memory + env_file: + - .env + volumes: + - ./data:/app/data + depends_on: + - embedding-server + restart: unless-stopped + networks: + - default + backend: build: context: . diff --git a/mcp_server/__init__.py b/mcp_server/__init__.py new file mode 100644 index 00000000..fde45f13 --- /dev/null +++ b/mcp_server/__init__.py @@ -0,0 +1 @@ +"""DistillNews MCP Server package.""" diff --git a/mcp_server/__main__.py b/mcp_server/__main__.py new file mode 100644 index 00000000..3ceda7d0 --- /dev/null +++ b/mcp_server/__main__.py @@ -0,0 +1,2 @@ +from mcp_server.app import main +main() diff --git a/mcp_server/app.py b/mcp_server/app.py new file mode 100644 index 00000000..b13a65df --- /dev/null +++ b/mcp_server/app.py @@ -0,0 +1,32 @@ +from mcp.server.fastmcp import FastMCP + +mcp = FastMCP("DistillNews Engine") + +@mcp.tool() +def news_search(query: str, limit: int = 5, category: str | None = None) -> list[dict]: + """Search the DistillNews corpus for relevant news articles using vector similarity and keyword matching.""" + from mcp_server.tools.search import search_news + return search_news(query=query, limit=limit, category=category) + +@mcp.tool() +def get_article(article_id: str) -> dict: + """Retrieve the full content and metadata of a specific article by its ID.""" + from mcp_server.tools.articles import fetch_article + return fetch_article(article_id=article_id) + +@mcp.tool() +def list_categories() -> list[str]: + """List all available news categories in the corpus.""" + return ["World", "Business", "Technology", "Entertainment", "Sports", "Science", "Health"] + +@mcp.tool() +def get_article_count() -> dict: + """Get the total number of articles in the corpus.""" + from mcp_server.tools.articles import count_articles + return count_articles() + +def main(): + mcp.run(transport="stdio") + +if __name__ == "__main__": + main() diff --git a/mcp_server/requirements.txt b/mcp_server/requirements.txt new file mode 100644 index 00000000..25e26b21 --- /dev/null +++ b/mcp_server/requirements.txt @@ -0,0 +1,5 @@ +mcp>=1.0.0 +fastmcp +pymongo +requests +python-dotenv diff --git a/mcp_server/tools/__init__.py b/mcp_server/tools/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/mcp_server/tools/articles.py b/mcp_server/tools/articles.py new file mode 100644 index 00000000..7be9dc22 --- /dev/null +++ b/mcp_server/tools/articles.py @@ -0,0 +1,30 @@ +from service.db import create_article_store + +_article_store = None + +def _get_store(): + global _article_store + if _article_store is None: + _article_store = create_article_store() + return _article_store + +def fetch_article(article_id: str) -> dict: + store = _get_store() + article = store.load_article(article_id) + if not article: + return {"error": f"Article with ID {article_id} not found."} + + return { + "title": article.get("title", ""), + "content": article.get("content", article.get("markdown_content", "")), + "category": article.get("category", ""), + "tags": article.get("tags", []), + "summary": article.get("summary", ""), + "publication_date": article.get("publication_date", "") + } + +def count_articles() -> dict: + store = _get_store() + # list_articles returns lightweight metadata for stored articles + articles = store.list_articles() + return {"total": len(articles)} diff --git a/mcp_server/tools/search.py b/mcp_server/tools/search.py new file mode 100644 index 00000000..a436a3e0 --- /dev/null +++ b/mcp_server/tools/search.py @@ -0,0 +1,59 @@ +from service.rag.providers.remote_embedding import RemoteEmbeddingProvider +from service.rag.backends.memory import InMemoryVectorStore +from service.rag.base import Document +from service.db import create_article_store + +_store = None +_loaded = False + +def _ensure_store(): + global _store, _loaded + if _loaded: + return + + embedder = RemoteEmbeddingProvider() + _store = InMemoryVectorStore(embedder=embedder) + + article_store = create_article_store() + articles = article_store.load_all_articles() + + documents = [] + for art in articles: + content = art.get("content", art.get("markdown_content", "")) + metadata = { + "id": art.get("id", ""), + "category": art.get("category", ""), + "tags": art.get("tags", []), + "summary": art.get("summary", ""), + "publication_date": art.get("publication_date", "") + } + documents.append(Document( + title=art.get("title", ""), + content=content, + metadata=metadata + )) + + _store.upload(documents) + _loaded = True + +def search_news(query: str, limit: int = 5, category: str | None = None) -> list[dict]: + _ensure_store() + results = _store.search(query=query, limit=limit if not category else limit * 5) + + formatted_results = [] + for res in results: + if category and res.metadata.get("category", "").lower() != category.lower(): + continue + + formatted_results.append({ + "id": res.metadata.get("id", ""), + "title": res.title, + "snippet": res.snippet, + "category": res.metadata.get("category", ""), + "score": res.score + }) + + if len(formatted_results) >= limit: + break + + return formatted_results diff --git a/pipeline/extraction_schemas.py b/pipeline/extraction_schemas.py new file mode 100644 index 00000000..ab2404be --- /dev/null +++ b/pipeline/extraction_schemas.py @@ -0,0 +1,65 @@ +"""JSON Schema definitions for structured extraction via LLM tool calls. + +These schemas replace raw text/JSON parsing with type-safe function call schemas, +ensuring 100% valid structured output from the extraction pipeline. +""" + +from service.agents.base import ToolDefinition + + +ARTICLE_EXTRACTION_TOOL = ToolDefinition( + name="submit_extracted_article", + description="Submit the extracted and structured news article metadata. Call this exactly once with all extracted fields.", + parameters={ + "type": "object", + "properties": { + "title": {"type": "string", "description": "The article headline."}, + "publication_date": {"type": "string", "description": "Publication date in ISO 8601 format."}, + "summary": {"type": "string", "description": "A concise summary of the article, maximum 100 words."}, + "content": { + "type": "string", + "description": "The full article body text. Preserve paragraph breaks. Write in third-person voice for community sources.", + }, + "category": { + "type": "string", + "enum": ["World", "Business", "Technology", "Entertainment", "Sports", "Science", "Health"], + "description": "The primary news category.", + }, + "tags": { + "type": "array", + "items": {"type": "string"}, + "description": "Relevant keyword tags (e.g., politics, economy, AI).", + }, + "location": { + "type": "string", + "description": "Geographic location mentioned or inferred from the article. Use 'unknown' if not inferable.", + }, + }, + "required": ["title", "publication_date", "summary", "content", "category", "tags", "location"], + }, +) + +NEWS_CLASSIFICATION_TOOL = ToolDefinition( + name="submit_classification", + description="Submit whether this content is a newsworthy article or not.", + parameters={ + "type": "object", + "properties": { + "is_news": {"type": "boolean", "description": "True if the content is newsworthy, false otherwise."}, + "reason": {"type": "string", "description": "Brief explanation of the classification decision."}, + }, + "required": ["is_news"], + }, +) + +MARKDOWN_FORMAT_TOOL = ToolDefinition( + name="submit_formatted_content", + description="Submit the markdown-formatted version of the article content.", + parameters={ + "type": "object", + "properties": { + "markdown": {"type": "string", "description": "The article content formatted with markdown headings, bullet points, and block quotes."}, + }, + "required": ["markdown"], + }, +) diff --git a/service/agents/__init__.py b/service/agents/__init__.py index ec2df8b9..e6916de9 100644 --- a/service/agents/__init__.py +++ b/service/agents/__init__.py @@ -10,6 +10,14 @@ """ from .factory import create_agent -from .base import AgentProvider, CompletionResult +from .base import AgentProvider, CompletionResult, ToolCallingProvider, ToolDefinition, ToolCall, AgentMessage -__all__ = ["create_agent", "AgentProvider", "CompletionResult"] +__all__ = [ + "create_agent", + "AgentProvider", + "CompletionResult", + "ToolCallingProvider", + "ToolDefinition", + "ToolCall", + "AgentMessage", +] diff --git a/service/agents/base.py b/service/agents/base.py index 34c8c477..2a836405 100644 --- a/service/agents/base.py +++ b/service/agents/base.py @@ -17,6 +17,41 @@ class CompletionResult: raw: dict | None = field(default=None, repr=False) # Provider-specific raw response +@dataclass +class ToolDefinition: + """Schema definition for a tool the LLM can invoke.""" + name: str + description: str + parameters: dict # JSON Schema + + def to_openai_schema(self) -> dict: + return { + "type": "function", + "function": { + "name": self.name, + "description": self.description, + "parameters": self.parameters, + }, + } + + +@dataclass +class ToolCall: + """A tool invocation returned by the LLM.""" + id: str + name: str + arguments: dict + + +@dataclass +class AgentMessage: + """A single message in a multi-turn tool-calling conversation.""" + role: str # system | user | assistant | tool + content: str | None = None + tool_calls: list[ToolCall] | None = None + tool_call_id: str | None = None + + class AgentProvider(ABC): """Abstract base for LLM completion backends. @@ -100,3 +135,16 @@ def _replacer(match: re.Match) -> str: return re.sub(r"\{steps\[0\]\.input\.(\w+)\}", _replacer, text) return _substitute(system_content), _substitute(user_content) + + +class ToolCallingProvider(AgentProvider): + """Extended provider supporting native tool / function calling.""" + + @abstractmethod + def chat_with_tools( + self, + messages: list[AgentMessage], + tools: list[ToolDefinition] | None = None, + tool_choice: str | dict = "auto", + ) -> AgentMessage: + """Execute a chat completion with optional tool definitions.""" diff --git a/service/agents/factory.py b/service/agents/factory.py index 2a3d4b40..e8fd60d3 100644 --- a/service/agents/factory.py +++ b/service/agents/factory.py @@ -25,12 +25,8 @@ def create_agent(provider: str | None = None, **kwargs) -> AgentProvider: from .providers.openai import OpenAIAgent return OpenAIAgent(**kwargs) - elif provider == "julep": - from .providers.julep import JulepAgent - - return JulepAgent(**kwargs) else: raise ValueError( f"Unknown agent provider: {provider!r}. " - f"Available: openai, julep" + f"Available: openai" ) diff --git a/service/agents/orchestrator.py b/service/agents/orchestrator.py new file mode 100644 index 00000000..f2412ac5 --- /dev/null +++ b/service/agents/orchestrator.py @@ -0,0 +1,53 @@ +import json +from typing import Callable, Any +from service.agents.base import ToolCallingProvider, ToolDefinition, ToolCall, AgentMessage + +class AgentOrchestrator: + """Multi-turn agent loop that auto-executes tool calls returned by the LLM.""" + + def __init__( + self, + agent: ToolCallingProvider, + tools: dict[str, tuple[ToolDefinition, Callable[..., Any]]], + max_turns: int = 5, + ): + self._agent = agent + self._tools = tools # name -> (schema, callable) + self._max_turns = max_turns + + @property + def tool_definitions(self) -> list[ToolDefinition]: + return [schema for schema, _ in self._tools.values()] + + def run(self, user_prompt: str, system_prompt: str) -> str: + messages = [ + AgentMessage(role="system", content=system_prompt), + AgentMessage(role="user", content=user_prompt), + ] + + for _ in range(self._max_turns): + response = self._agent.chat_with_tools( + messages, tools=self.tool_definitions + ) + messages.append(response) + + if not response.tool_calls: + return response.content or "" + + for call in response.tool_calls: + _, fn = self._tools[call.name] + try: + result = fn(**call.arguments) + except Exception as e: + result = {"error": str(e)} + messages.append(AgentMessage( + role="tool", + tool_call_id=call.id, + content=json.dumps(result) if not isinstance(result, str) else result, + )) + + # Exhausted turns, return last assistant content + for msg in reversed(messages): + if msg.role == "assistant" and msg.content: + return msg.content + return "" diff --git a/service/agents/providers/julep.py b/service/agents/providers/julep.py deleted file mode 100644 index d3e11cbb..00000000 --- a/service/agents/providers/julep.py +++ /dev/null @@ -1,129 +0,0 @@ -""" -Julep agent provider. - -Wraps the Julep SDK behind the abstract AgentProvider interface. -Supports both simple completions and native Julep task execution -(agents → tasks → executions with polling). -""" - -import os -import time -import yaml -from pathlib import Path -from functools import lru_cache - -from config import config -from service.agents.base import AgentProvider, CompletionResult - - -class JulepAgent(AgentProvider): - """AgentProvider backed by the Julep AI platform. - - Constructor args can override env-var defaults: - - ``api_key`` → ``JULEP_API_KEY`` - - ``model`` → ``AGENT_MODEL`` / ``JULEP_MODEL`` - - ``environment`` → ``JULEP_ENVIRONMENT`` - """ - - def __init__( - self, - api_key: str | None = None, - model: str | None = None, - environment: str | None = None, - ): - from julep import Julep # defer import so the SDK is optional - - self._api_key = api_key or config.JULEP_API_KEY - self._model = model or config.JULEP_MODEL - self._environment = environment or config.JULEP_ENVIRONMENT - - self._client = Julep( - api_key=self._api_key, - environment=self._environment, - ) - - # ------------------------------------------------------------------ - # AgentProvider interface - # ------------------------------------------------------------------ - - def complete(self, system_prompt: str, user_prompt: str) -> CompletionResult: - """Create a one-shot Julep agent + inline task, execute and poll.""" - agent = self._client.agents.create( - name="OneShot", - model=self._model, - about="Temporary agent for a single completion.", - ) - - task_def = { - "name": "inline_completion", - "description": "Single-turn completion task.", - "main": [ - { - "prompt": [ - {"role": "system", "content": system_prompt}, - {"role": "user", "content": user_prompt}, - ] - } - ], - } - task = self._client.tasks.create(agent_id=agent.id, **task_def) - execution = self._client.executions.create(task_id=task.id, input={}) - - result = self._poll_execution(execution.id) - content = result["choices"][0]["message"]["content"] - return CompletionResult(content=content, raw=result) - - def complete_from_template( - self, template_path: str | Path, input_data: dict - ) -> CompletionResult: - """Use Julep's native task system — the YAML is passed through as-is. - - This preserves the ``{steps[0].input.field}`` syntax that Julep - understands natively, so no variable substitution is needed on our side. - """ - task = self._get_or_create_task(str(template_path)) - execution = self._client.executions.create( - task_id=task.id, input=input_data - ) - - result = self._poll_execution(execution.id) - content = result["choices"][0]["message"]["content"] - return CompletionResult(content=content, raw=result) - - # ------------------------------------------------------------------ - # Internal helpers - # ------------------------------------------------------------------ - - def _poll_execution(self, execution_id: str, poll_interval: float = 1.0) -> dict: - """Poll a Julep execution until it succeeds or fails.""" - while True: - res = self._client.executions.get(execution_id) - if res.status == "succeeded": - return res.output - if res.status == "failed": - raise RuntimeError( - f"Julep execution {execution_id} failed: {res}" - ) - time.sleep(poll_interval) - - @lru_cache(maxsize=128) - def _get_or_create_task(self, yaml_path: str) -> object: - """Load a YAML task definition and register it with a Julep agent. - - Results are cached so repeated calls with the same path reuse - the same agent + task. - """ - with open(yaml_path, "r", encoding="utf-8") as f: - task_definition = yaml.safe_load(f) - - # Derive agent name / description from the task YAML - agent_name = task_definition.get("name", "TaskAgent") - agent_about = task_definition.get("description", "Agent for a YAML-defined task.") - - agent = self._client.agents.create( - name=agent_name, - model=self._model, - about=agent_about, - ) - task = self._client.tasks.create(agent_id=agent.id, **task_definition) - return task diff --git a/service/agents/providers/openai.py b/service/agents/providers/openai.py index 2c07160f..dbb98944 100644 --- a/service/agents/providers/openai.py +++ b/service/agents/providers/openai.py @@ -6,11 +6,12 @@ """ import os +import json from config import config -from service.agents.base import AgentProvider, CompletionResult +from service.agents.base import ToolCallingProvider, CompletionResult, ToolDefinition, ToolCall, AgentMessage -class OpenAIAgent(AgentProvider): +class OpenAIAgent(ToolCallingProvider): """AgentProvider backed by an OpenAI-compatible chat completions API. Constructor args can override env-var defaults via config. @@ -65,6 +66,60 @@ def complete(self, system_prompt: str, user_prompt: str) -> CompletionResult: content=content, raw=response.model_dump(), ) + + def chat_with_tools( + self, + messages: list[AgentMessage], + tools: list[ToolDefinition] | None = None, + tool_choice: str | dict = "auto", + ) -> AgentMessage: + openai_messages = [] + for msg in messages: + m = {"role": msg.role} + if msg.content is not None: + m["content"] = msg.content + if msg.tool_call_id is not None: + m["tool_call_id"] = msg.tool_call_id + if msg.tool_calls: + m["tool_calls"] = [ + { + "id": tc.id, + "type": "function", + "function": { + "name": tc.name, + "arguments": json.dumps(tc.arguments) + } + } for tc in msg.tool_calls + ] + openai_messages.append(m) + + kwargs = { + "model": self._model, + "messages": openai_messages, + } + + if tools: + kwargs["tools"] = [t.to_openai_schema() for t in tools] + kwargs["tool_choice"] = tool_choice + + response = self._client.chat.completions.create(**kwargs) + resp_msg = response.choices[0].message + + parsed_tool_calls = None + if resp_msg.tool_calls: + parsed_tool_calls = [] + for tc in resp_msg.tool_calls: + parsed_tool_calls.append(ToolCall( + id=tc.id, + name=tc.function.name, + arguments=json.loads(tc.function.arguments) + )) + + return AgentMessage( + role=resp_msg.role, + content=resp_msg.content, + tool_calls=parsed_tool_calls, + ) # complete_from_template() uses the default base-class implementation, # which parses the YAML, substitutes {steps[0].input.field}, and diff --git a/service/chatbot/service.py b/service/chatbot/service.py index e3138a54..c021f9af 100644 --- a/service/chatbot/service.py +++ b/service/chatbot/service.py @@ -1,82 +1,157 @@ -"""Chatbot orchestration over abstract chat and retrieval interfaces.""" +"""Chatbot orchestration with dual-mode context: autonomous global search vs. article-focused grounding.""" from collections import defaultdict, deque from pathlib import Path from typing import Any -from service.agents.base import AgentProvider +from service.agents.base import ToolCallingProvider, ToolDefinition, AgentMessage +from service.agents.orchestrator import AgentOrchestrator from service.rag.base import DocumentStore +# Tool schema for news_search +NEWS_SEARCH_TOOL = ToolDefinition( + name="news_search", + description=( + "Search the DistillNews article corpus for relevant news articles. " + "Use this to find factual, grounded information when answering user questions about news, events, or current affairs. " + "You decide the search keywords and how many results to retrieve." + ), + parameters={ + "type": "object", + "properties": { + "query": {"type": "string", "description": "Search keywords or phrase to find relevant articles."}, + "limit": {"type": "integer", "description": "Maximum number of articles to return (1-10).", "default": 5}, + }, + "required": ["query"], + }, +) -class ChatbotService: - """Generate grounded responses without knowing pipeline or embedding details. +GET_ARTICLE_TOOL = ToolDefinition( + name="get_article", + description=( + "Retrieve the full content of a specific article by its ID. " + "Use this when a search result snippet is insufficient and you need the complete article text." + ), + parameters={ + "type": "object", + "properties": { + "article_id": {"type": "string", "description": "The unique article identifier."}, + }, + "required": ["article_id"], + }, +) + +GLOBAL_SYSTEM_PROMPT = """You are DistillNews AI, a knowledgeable and friendly news assistant. + +Your role is to help users stay informed by answering questions about news and current events. +You have access to a curated corpus of news articles through the `news_search` tool. + +Guidelines: +- Use `news_search` to find relevant articles when answering factual questions about news, events, or current affairs. +- You decide the best search keywords and how many results to retrieve. +- If a search snippet is insufficient, use `get_article` to read the full article. +- Synthesize information from multiple sources when appropriate. +- Be concise (under 150 words) unless the user asks for detail. +- If no relevant articles are found, say so honestly rather than fabricating information. +- Use a conversational, engaging tone — like a knowledgeable friend sharing updates. +- End with a brief closing that invites further questions. +""" + +ARTICLE_SYSTEM_PROMPT = """You are DistillNews AI, a knowledgeable and friendly news assistant. - The service only asks a ``DocumentStore`` for results. Whether that store - uses local Ollama vectors, an OpenAI-compatible endpoint, or no embeddings - at all is determined outside the chatbot. - """ +The user is currently reading the article below. Answer their questions about it directly. +You have the `news_search` tool available if the user asks about related or comparative news topics. + +Guidelines: +- Answer questions about the active article immediately using the context provided — no tool call needed. +- Use `news_search` only if the user asks about external or related news topics. +- Be concise (under 150 words) unless the user asks for detail. +- Use a conversational, engaging tone. + +ACTIVE ARTICLE: +{article_context} +""" + +class ChatbotService: + """Dual-mode chatbot: autonomous tool-calling for global queries, article-grounded for reader sessions.""" def __init__( self, - agent: AgentProvider, + agent: ToolCallingProvider, document_store: DocumentStore, - prompts_dir: Path, + prompts_dir: Path | None = None, logger: Any | None = None, ): self._agent = agent self._document_store = document_store - self._prompts_dir = prompts_dir self._logger = logger - self._user_memory = defaultdict(lambda: deque(maxlen=6)) + self._user_memory: dict[str, deque] = defaultdict(lambda: deque(maxlen=6)) + + # Build tool executors + self._tools = { + "news_search": (NEWS_SEARCH_TOOL, self._execute_news_search), + "get_article": (GET_ARTICLE_TOOL, self._execute_get_article), + } + self._orchestrator = AgentOrchestrator( + agent=agent, tools=self._tools, max_turns=5 + ) def get_response( self, query: str, user_id: str = "debug", reading: str | None = None, - prompt: str = "chatbot.yaml", + prompt: str = "chatbot.yaml", # kept for API compat, ignored in new path ) -> str | None: - """Generate a response grounded in the configured document store.""" self._log("chat_query", user_id, query) memory = self._user_memory[user_id] - filtered = self._filter_prompt(query) - search_results = self._document_store.search(filtered, limit=5) - self._log("rag_search", filtered, len(search_results)) - - if not search_results: - context = f"Article Being Read:\n{reading}" if reading else "No relevant articles found." - self._log("warn", "No RAG results for query") + # Build system prompt based on mode + if reading: + system_prompt = ARTICLE_SYSTEM_PROMPT.format(article_context=reading) else: - rag_context = "\n\n".join( - result.snippet or result.content for result in search_results - ) - context = f"Article Being Read:\n{reading}\n\nRelated News Context:\n{rag_context}" if reading else rag_context + system_prompt = GLOBAL_SYSTEM_PROMPT + + # Append conversation memory + if memory: + system_prompt += "\n\nRecent conversation history:\n" + "\n".join(memory) self._log("ai_call", "chatbot_response", query) - result = self._agent.complete_from_template( - self._prompts_dir / prompt, - { - "query": query, - "reading": reading or "", - "content": context, - "memory": "\n".join(memory), - }, + response = self._orchestrator.run( + user_prompt=query, system_prompt=system_prompt ) - response = result.content self._log("chat_response", response) memory.append(f"User: {query}") memory.append(f"Assistant: {response}") return response - def _filter_prompt(self, query: str) -> str: - self._log("ai_call", "keyword_extraction", query) - result = self._agent.complete_from_template( - self._prompts_dir / "filter_prompt.yaml", {"query": query} - ) - self._log("ai_result", "keyword_extraction", result.content) - return result.content + def _execute_news_search(self, query: str, limit: int = 5) -> list[dict]: + self._log("rag_search", query, limit) + results = self._document_store.search(query, limit=limit) + return [ + { + "id": r.metadata.get("id", ""), + "title": r.title, + "snippet": r.snippet or r.content[:300], + "category": r.metadata.get("category", ""), + "score": r.score, + } + for r in results + ] + + def _execute_get_article(self, article_id: str) -> dict: + from service.db import create_article_store + store = create_article_store() + article = store.load_article(article_id) + if article is None: + return {"error": f"Article '{article_id}' not found"} + return { + "title": article.get("title", ""), + "content": article.get("content") or article.get("markdown_content") or article.get("summary", ""), + "category": article.get("category", ""), + "tags": article.get("tags", []), + } def _log(self, method: str, *args: object) -> None: if self._logger: diff --git a/service/chatbot/wiring.py b/service/chatbot/wiring.py index 675d9a4d..7eb791a5 100644 --- a/service/chatbot/wiring.py +++ b/service/chatbot/wiring.py @@ -16,7 +16,7 @@ agent = create_agent() doc_store = create_doc_store() article_store = create_article_store() -chatbot = ChatbotService(agent, doc_store, prompts_dir, logger=log) +chatbot = ChatbotService(agent, doc_store, logger=log) def _load_and_upload_articles(): diff --git a/service/logger.py b/service/logger.py index a8a93703..cd7097c0 100644 --- a/service/logger.py +++ b/service/logger.py @@ -58,7 +58,7 @@ def truncate(text: str, max_len: int = MAX_INLINE) -> str: def _badge(label: str, bg: str, fg: str = _C.WHITE) -> str: padded = label.center(6) - return f"{bg}{fg}{_C.BOLD}[{padded}]{_C.RESET}" + return f"{bg}{fg}{_C.BOLD} {padded} {_C.RESET}" # ── Logger class ──────────────────────────────────────────────────────────── @@ -106,7 +106,7 @@ def remove_listener(cls, fn: Callable): @staticmethod def _print(badge_key: str, message: str, detail: str | None = None, truncate_detail: bool = True): - badge = Logger._BADGES.get(badge_key, f"[{badge_key.upper()}]") + badge = Logger._BADGES.get(badge_key, f" {badge_key.upper()} ") line = f"{badge} {message}" if detail is not None: detail_str = truncate(detail) if truncate_detail else str(detail).replace("\n", " ").replace("\r", "") diff --git a/service/rag/backends/julep.py b/service/rag/backends/julep.py deleted file mode 100644 index 5f97b519..00000000 --- a/service/rag/backends/julep.py +++ /dev/null @@ -1,81 +0,0 @@ -""" -Julep document store backend. - -Uses the Julep SDK's built-in document storage and search API. -""" - -import os -from config import config -from service.rag.base import DocumentStore, Document, SearchResult - - -class JulepDocStore(DocumentStore): - """DocumentStore backed by Julep's agent document API. - - Creates a dedicated Julep agent for document storage. - """ - - def __init__( - self, - api_key: str | None = None, - model: str | None = None, - environment: str | None = None, - agent_name: str = "RAG Doc Store", - ): - from julep import Julep, ConflictError # defer import - - self._ConflictError = ConflictError - - self._client = Julep( - api_key=api_key or config.JULEP_API_KEY, - environment=environment or config.JULEP_ENVIRONMENT, - ) - - _model = model or config.JULEP_MODEL - - # Create a dedicated agent for doc storage - self._agent = self._client.agents.create( - name=agent_name, - model=_model, - about="Agent used for document storage and retrieval.", - ) - - # ------------------------------------------------------------------ - # DocumentStore interface - # ------------------------------------------------------------------ - - def upload(self, documents: list[Document]) -> None: - """Upload documents to Julep's agent doc store.""" - for doc in documents: - try: - self._client.agents.docs.create( - agent_id=self._agent.id, - title=doc.title, - content=doc.content, - metadata=doc.metadata or {}, - ) - except self._ConflictError: - # Document already exists — skip - continue - - def search(self, query: str, limit: int = 5) -> list[SearchResult]: - """Search Julep's agent doc store.""" - results = self._client.agents.docs.search( - agent_id=self._agent.id, - text=query, - limit=limit, - ) - - results_dict = results.model_dump() - docs = results_dict.get("docs", []) - - return [ - SearchResult( - title=doc.get("title", "Unknown"), - content=doc.get("content", ""), - snippet=doc.get("snippet", {}).get("content", ""), - score=doc.get("score"), - metadata=doc.get("metadata") or {}, - ) - for doc in docs - ] diff --git a/service/rag/factory.py b/service/rag/factory.py index a218ac8e..c0b37c3f 100644 --- a/service/rag/factory.py +++ b/service/rag/factory.py @@ -16,11 +16,7 @@ def create_doc_store(backend: str | None = None, **kwargs) -> DocumentStore: """ backend = (backend or config.RAG_BACKEND).lower() - if backend == "julep": - from .backends.julep import JulepDocStore - - return JulepDocStore(**kwargs) - elif backend == "memory": + if backend == "memory": from .backends.memory import InMemoryVectorStore embedder = kwargs.pop("embedder", None) @@ -42,5 +38,5 @@ def create_doc_store(backend: str | None = None, **kwargs) -> DocumentStore: else: raise ValueError( f"Unknown RAG backend: {backend!r}. " - "Available: memory, bm25, julep, none" + "Available: memory, bm25, none" ) diff --git a/tests/test_agents_and_factory.py b/tests/test_agents_and_factory.py index b5e67295..d500a79b 100644 --- a/tests/test_agents_and_factory.py +++ b/tests/test_agents_and_factory.py @@ -3,7 +3,14 @@ import pytest from pathlib import Path from service.agents.factory import create_agent -from service.agents.base import AgentProvider, CompletionResult +from service.agents.base import ( + AgentProvider, + CompletionResult, + ToolCallingProvider, + ToolDefinition, + ToolCall, + AgentMessage +) class DummyAgent(AgentProvider): @@ -11,8 +18,22 @@ def complete(self, system_prompt: str, user_prompt: str) -> CompletionResult: return CompletionResult(content=f"System: {system_prompt} | User: {user_prompt}") +class FakeToolCallingAgent(ToolCallingProvider): + def complete(self, system_prompt: str, user_prompt: str) -> CompletionResult: + return CompletionResult(content="dummy") + + def chat_with_tools(self, messages: list[AgentMessage], tools: list[ToolDefinition] | None = None, tool_choice: str | dict = "auto") -> AgentMessage: + has_tool_call = any(m.role == "assistant" and m.tool_calls for m in messages) + if not has_tool_call and tools: + return AgentMessage( + role="assistant", + tool_calls=[ToolCall(id="call_123", name=tools[0].name, arguments={"arg": "val"})] + ) + return AgentMessage(role="assistant", content="final result") + + def test_agent_factory_unknown_provider_raises_error(): - with pytest.raises(ValueError, match="Unknown agent provider: 'invalid_provider'"): + with pytest.raises(ValueError, match="Available: openai"): create_agent("invalid_provider") @@ -50,3 +71,37 @@ def test_dummy_agent_complete_from_template(tmp_path: Path): agent = DummyAgent() res = agent.complete_from_template(template_file, {"query": "Hello world"}) assert res.content == "System: System prompt | User: Input: Hello world" + + +def test_tool_calling_provider_interface(): + tool = ToolDefinition(name="test", description="desc", parameters={"type": "object"}) + assert tool.name == "test" + + call = ToolCall(id="1", name="test", arguments={}) + assert call.id == "1" + + msg = AgentMessage(role="user", content="hello") + assert msg.role == "user" + + agent = FakeToolCallingAgent() + assert isinstance(agent, ToolCallingProvider) + + +def test_orchestrator_basic_flow(): + agent = FakeToolCallingAgent() + tool = ToolDefinition(name="test_tool", description="test", parameters={}) + + messages = [AgentMessage(role="user", content="do something")] + response1 = agent.chat_with_tools(messages, tools=[tool]) + + assert response1.role == "assistant" + assert response1.tool_calls is not None + assert len(response1.tool_calls) == 1 + assert response1.tool_calls[0].name == "test_tool" + + messages.append(response1) + messages.append(AgentMessage(role="tool", content="success", tool_call_id=response1.tool_calls[0].id)) + + response2 = agent.chat_with_tools(messages, tools=[tool]) + assert response2.role == "assistant" + assert response2.content == "final result" diff --git a/tests/test_chatbot_and_embeddings.py b/tests/test_chatbot_and_embeddings.py index ce4cdaa6..078943b1 100644 --- a/tests/test_chatbot_and_embeddings.py +++ b/tests/test_chatbot_and_embeddings.py @@ -13,6 +13,7 @@ ) from service.rag.base import Document, DocumentStore, SearchResult from service.rag.factory import create_doc_store +from service.agents.base import ToolCallingProvider, ToolDefinition, AgentMessage class KeywordEmbedder(EmbeddingProvider): @@ -37,15 +38,21 @@ def search(self, query: str, limit: int = 5) -> list[SearchResult]: ] -class FakeAgent: +class FakeToolCallingAgent(ToolCallingProvider): def __init__(self): self.calls = [] - def complete_from_template(self, template_path, input_data): - self.calls.append((Path(template_path).name, input_data)) - if Path(template_path).name == "filter_prompt.yaml": - return CompletionResult(content="climate") - return CompletionResult(content=f"Grounded answer: {input_data['content']}") + def complete(self, system_prompt: str, user_prompt: str) -> CompletionResult: + return CompletionResult(content="dummy") + + def chat_with_tools( + self, + messages: list[AgentMessage], + tools: list[ToolDefinition] | None = None, + tool_choice: str | dict = "auto", + ) -> AgentMessage: + self.calls.append(("chat_with_tools", messages)) + return AgentMessage(role="assistant", content="Grounded answer: Climate article excerpt") def test_memory_store_uses_injected_embedder_and_preserves_metadata(): @@ -100,7 +107,7 @@ def test_noop_embedding_provider_makes_no_vectors(): def test_embedding_backends_are_not_rag_backends(): for backend in ("openai", "foundry", "sentence_transformers", "in_memory"): - with pytest.raises(ValueError, match="Available: memory, bm25, julep, none"): + with pytest.raises(ValueError, match="Available: memory, bm25, none"): create_doc_store(backend) @@ -151,16 +158,20 @@ def encode(self, texts, **kwargs): def test_chatbot_only_depends_on_agent_and_document_store(): - agent = FakeAgent() + agent = FakeToolCallingAgent() chatbot = ChatbotService( agent=agent, document_store=StaticDocumentStore(), - prompts_dir=Path("chatbot/prompts"), ) response = chatbot.get_response("What happened?", user_id="reader-1") assert response == "Grounded answer: Climate article excerpt" - assert agent.calls[0][0] == "filter_prompt.yaml" - assert agent.calls[1][1]["content"] == "Climate article excerpt" - assert agent.calls[1][1]["memory"] == "" + assert len(agent.calls) > 0 + assert agent.calls[0][0] == "chat_with_tools" + # messages[0] = system prompt, messages[1] = user message + messages = agent.calls[0][1] + user_messages = [m for m in messages if m.role == "user"] + assert len(user_messages) == 1 + assert user_messages[0].content == "What happened?" + diff --git a/tests/test_logger.py b/tests/test_logger.py index cb46b5ca..fb33bfa1 100644 --- a/tests/test_logger.py +++ b/tests/test_logger.py @@ -10,25 +10,25 @@ def test_truncate_helper(): def test_logger_badges_formatting(capsys): log.info("Test Info Message") captured = capsys.readouterr().out - assert "[ INFO ]" in captured + assert " INFO " in captured assert "Test Info Message" in captured log.warn("Test Warning Message") captured = capsys.readouterr().out - assert "[ WARN ]" in captured + assert " WARN " in captured assert "Test Warning Message" in captured log.error("Test Error Message") captured = capsys.readouterr().out - assert "[ FAIL ]" in captured + assert " FAIL " in captured assert "Test Error Message" in captured log.success("Test Success Message") captured = capsys.readouterr().out - assert "[ OK ]" in captured + assert " OK " in captured assert "Test Success Message" in captured log.db("MongoDB Action", "connected") captured = capsys.readouterr().out - assert "[ DB ]" in captured + assert " DB " in captured assert "MongoDB Action" in captured diff --git a/tests/test_server_routes.py b/tests/test_server_routes.py index 2b982f03..923f0351 100644 --- a/tests/test_server_routes.py +++ b/tests/test_server_routes.py @@ -53,7 +53,7 @@ async def test_send_otp_failure_logging_output(capsys): captured = capsys.readouterr() out = captured.out - assert "[ INFO ]" in out - assert "[ FAIL ]" in out - assert "[ WARN ]" in out + assert " INFO " in out + assert " FAIL " in out + assert " WARN " in out assert "Email delivery failed" in out diff --git a/utils/logger.py b/utils/logger.py index a8a93703..cd7097c0 100644 --- a/utils/logger.py +++ b/utils/logger.py @@ -58,7 +58,7 @@ def truncate(text: str, max_len: int = MAX_INLINE) -> str: def _badge(label: str, bg: str, fg: str = _C.WHITE) -> str: padded = label.center(6) - return f"{bg}{fg}{_C.BOLD}[{padded}]{_C.RESET}" + return f"{bg}{fg}{_C.BOLD} {padded} {_C.RESET}" # ── Logger class ──────────────────────────────────────────────────────────── @@ -106,7 +106,7 @@ def remove_listener(cls, fn: Callable): @staticmethod def _print(badge_key: str, message: str, detail: str | None = None, truncate_detail: bool = True): - badge = Logger._BADGES.get(badge_key, f"[{badge_key.upper()}]") + badge = Logger._BADGES.get(badge_key, f" {badge_key.upper()} ") line = f"{badge} {message}" if detail is not None: detail_str = truncate(detail) if truncate_detail else str(detail).replace("\n", " ").replace("\r", "")