From d132ca23a4430934740fc5dff92f0a83d5142474 Mon Sep 17 00:00:00 2001 From: SumanthPal <55633923+SumanthPal@users.noreply.github.com> Date: Wed, 11 Feb 2026 11:44:17 -0800 Subject: [PATCH] Revert "fixing research tool" --- backend/app/agents/executor.py | 52 ++---- backend/app/agents/tools/__init__.py | 2 - backend/app/agents/tools/database_tool.py | 204 +++++++--------------- backend/app/agents/tools/research_tool.py | 94 ++-------- backend/app/db/models.py | 19 +- backend/app/scripts/research_table_db.py | 31 +--- backend/tests/test_research_tool.py | 39 ++--- 7 files changed, 125 insertions(+), 316 deletions(-) diff --git a/backend/app/agents/executor.py b/backend/app/agents/executor.py index c24b852..dfdbb79 100644 --- a/backend/app/agents/executor.py +++ b/backend/app/agents/executor.py @@ -186,9 +186,8 @@ def _build_agent_executor(self) -> AgentExecutor: name="research_token", func=self._execute_research_token_wrapper, description=( - "Research a token for news and market sentiment. Uses a shared cache — " - "calling it is cheap, so use it whenever you need research. " - "Input: token symbol and optional recency (e.g., 'ETH 7d' for last 7 days)" + "Research a token for news and any market sentiments, " + "Input: token symbol and recency (e.g., 'ETH 7d' for last 7 days)" ), ), Tool( @@ -358,7 +357,7 @@ def _build_agent_executor(self) -> AgentExecutor: - Do not create duplicate plans - cancel outdated ones first - DOING NOTHING IS VALID: If no plans are due and market conditions don't warrant action, it's perfectly fine to take no action this cycle. Don't trade just to trade. - Before making significant trades, research tokens using the research_token tool to check recent news and sentiment -- The research_token tool uses a shared cache automatically. Just call it when you need research. +- If recent research already exists, you may reuse it instead of researching again TOOLS: ------ @@ -580,14 +579,15 @@ def _get_technical_indicator_wrapper(self, input_str: str) -> str: return error_msg def _execute_research_token_wrapper(self, input_str: str) -> str: """ - Cache-first research wrapper. + Wrapper for executing research token from LangChain tool. - 1. Parse input -> token + recency - 2. Check shared cache via database_tool.get_fresh_research() - 3. Cache HIT -> return cached summary (zero API cost) - 4. Cache MISS -> call Perplexity, derive opinion, upsert to cache + Args: + input_str: "ETH" or "ETH 7d" + + Returns: + Research summary string """ - from .tools.research_tool import research_result_to_dict, derive_opinion + from .tools.research_tool import research_result_to_dict try: parts = input_str.split() @@ -595,42 +595,24 @@ def _execute_research_token_wrapper(self, input_str: str) -> str: return "Error: Invalid format. Expected 'TOKEN' or 'TOKEN RECENCY'" token = parts[0].upper() - _valid_recencies = {"1d", "7d", "30d", "365d"} - raw_recency = parts[1] if len(parts) > 1 else "7d" - recency = raw_recency if raw_recency in _valid_recencies else "7d" - - # ── Cache check ────────────────────────────────────────── - if self.database_tool: - cached = self._run_async( - self.database_tool.get_fresh_research(token, recency) - ) - if cached is not None: - logger.info(f"Research cache HIT for {token} — skipping API call") - return cached["summary_markdown"] + recency = parts[1] if len(parts) > 1 else "7d" - # ── Cache MISS — call Perplexity ───────────────────────── + # Call sync research method research_result = self.research_tool.research_token( token_symbol=token, recency=recency, ) - # Derive opinion from summary - opinion = derive_opinion(research_result.summary_markdown) - - # Persist to shared cache (skip raw_results to keep rows lean) - if self.database_tool: + # Persist result to DB (convert dataclass to dict) + if self.database_tool and self.agent_uuid: result_dict = research_result_to_dict(research_result) self._run_async( - self.database_tool.upsert_research_result( - crypto_token=token, + self.database_tool.save_research_result( + agent_uuid=self.agent_uuid, query=f"Research {token}", - summary_markdown=result_dict["summary_markdown"], - citations=result_dict["citations"], + result=result_dict, recency=recency, - provider=result_dict.get("provider", "perplexity"), related_tokens=[token], - agent_opinion=opinion, - last_researched_by=self.agent_uuid, ) ) diff --git a/backend/app/agents/tools/__init__.py b/backend/app/agents/tools/__init__.py index b78d187..331a030 100644 --- a/backend/app/agents/tools/__init__.py +++ b/backend/app/agents/tools/__init__.py @@ -17,7 +17,6 @@ from .tweet_post_tool import TweetPostTool from .database_tool import DatabaseTool from .plan_tool import PlanTool -from .research_tool import ResearchTool # Deprecated: These tools are kept for backwards compatibility # Use MakeTradeTool instead for simulated trading @@ -30,7 +29,6 @@ "TweetPostTool", "DatabaseTool", "PlanTool", - "ResearchTool", # Deprecated "TradeTool", "PortfolioTool", diff --git a/backend/app/agents/tools/database_tool.py b/backend/app/agents/tools/database_tool.py index 4246e98..7cf45c3 100644 --- a/backend/app/agents/tools/database_tool.py +++ b/backend/app/agents/tools/database_tool.py @@ -8,7 +8,7 @@ import logging from typing import Optional, Dict, Any, List from uuid import UUID -from datetime import datetime, timezone +from datetime import datetime, timezone, timedelta from decimal import Decimal from sqlalchemy.ext.asyncio import AsyncSession @@ -20,19 +20,6 @@ logger = logging.getLogger(__name__) -# Freshness thresholds per recency tier (seconds) -FRESHNESS_THRESHOLDS = { - "1d": 2 * 3600, # 2 hours - "7d": 6 * 3600, # 6 hours - "30d": 24 * 3600, # 24 hours - "365d": 72 * 3600, # 72 hours -} -DEFAULT_FRESHNESS = 6 * 3600 # 6 hours - -# Tightness ordering — lower = tighter (more recent Perplexity search window). -# A cache row fetched with "7d" should NOT satisfy a "1d" request. -_RECENCY_TIGHTNESS = {"1d": 0, "7d": 1, "30d": 2, "365d": 3} - class DatabaseTool: """ @@ -415,139 +402,72 @@ async def create_agent_if_not_exists( logger.error(f"Failed to create agent: {e}") raise - # ── Shared research cache ────────────────────────────────────────── - - async def get_fresh_research( - self, crypto_token: str, recency: str = "7d" - ) -> Optional[Dict]: - """ - Return cached research for *crypto_token* if it is still fresh. - - Freshness is determined by comparing the row's `updated_at` against - the threshold for the requested *recency* tier. - - Returns: - Dict with cached research data if fresh, None if stale/missing. - """ + async def save_research_result( + self, + agent_uuid: UUID, + query: str, + result: Dict, + recency: Optional[str] = None, + related_tokens: Optional[List[str]] = None + ) -> None: + #Saving the research to Postgres (agent_research_artifact table) try: - stmt = select(AgentResearchArtifact).where( - AgentResearchArtifact.crypto_token == crypto_token - ) - result = await self.session.execute(stmt) - artifact = result.scalar_one_or_none() - - if artifact is None: - logger.info(f"Research cache MISS for {crypto_token} (no row)") - return None - - # Tightness check: cached "7d" data can't satisfy a "1d" request - cached_tight = _RECENCY_TIGHTNESS.get(artifact.recency, 99) - requested_tight = _RECENCY_TIGHTNESS.get(recency, 1) - if cached_tight > requested_tight: - logger.info( - f"Research cache MISS for {crypto_token} " - f"(cached recency={artifact.recency} too loose for requested={recency})" - ) - return None - - age_seconds = ( - datetime.now(timezone.utc) - artifact.updated_at - ).total_seconds() - threshold = FRESHNESS_THRESHOLDS.get(recency, DEFAULT_FRESHNESS) - - if age_seconds > threshold: - logger.info( - f"Research cache STALE for {crypto_token} " - f"(age={age_seconds:.0f}s, threshold={threshold}s)" - ) - return None - - logger.info( - f"Research cache HIT for {crypto_token} " - f"(age={age_seconds:.0f}s, threshold={threshold}s)" + stmt = insert(AgentResearchArtifact).values( + agent_id=agent_uuid, + query=query, + recency=recency, + provider=result.get("provider", "perplexity"), + summary_markdown=result["summary_markdown"], + citations=result.get("citations", []), + raw_results=result.get("raw_results", {}), + related_tokens=related_tokens or [], + created_at=datetime.now(timezone.utc), ) - return { - "crypto_token": artifact.crypto_token, - "query": artifact.query, - "recency": artifact.recency, - "provider": artifact.provider, - "summary_markdown": artifact.summary_markdown, - "citations": artifact.citations, - "related_tokens": artifact.related_tokens, - "agent_opinion": artifact.agent_opinion, - "updated_at": artifact.updated_at.isoformat(), - } - + await self.session.execute(stmt) + await self.session.commit() + logger.info(f"Research result saved for agent={agent_uuid}, query='{query}'") except Exception as e: - logger.error(f"Failed to check research cache: {e}") - raise - - async def upsert_research_result( - self, - crypto_token: str, - query: str, - summary_markdown: str, - citations: List[Dict], - recency: Optional[str] = None, - provider: str = "perplexity", - raw_results: Optional[Dict] = None, - related_tokens: Optional[List[str]] = None, - agent_opinion: Optional[str] = None, - last_researched_by: Optional[UUID] = None, - ) -> None: - """ - Insert or update the shared research row for *crypto_token*. - - Uses INSERT ... ON CONFLICT (crypto_token) DO UPDATE so every token - has exactly one row that gets refreshed in place. - """ - now = datetime.now(timezone.utc) + await self.session.rollback() + logger.error(f"Failed to save research result: {e}") + raise + + async def get_recent_research_results( + self, + agent_uuid: UUID, + token: Optional[str] = None, + days: int = 7 + ) -> List[Dict]: + #Fetch the recent research data for an agent from Postgres try: - stmt = ( - insert(AgentResearchArtifact) - .values( - crypto_token=crypto_token, - last_researched_by=last_researched_by, - query=query, - recency=recency, - provider=provider, - summary_markdown=summary_markdown, - citations=citations, - raw_results=raw_results, - related_tokens=related_tokens, - agent_opinion=agent_opinion, - created_at=now, - updated_at=now, - ) - .on_conflict_do_update( - index_elements=["crypto_token"], - set_={ - "last_researched_by": last_researched_by, - "query": query, - "recency": recency, - "provider": provider, - "summary_markdown": summary_markdown, - "citations": citations, - "raw_results": raw_results, - "related_tokens": related_tokens, - "agent_opinion": agent_opinion, - "updated_at": now, - }, - ) - ) - - await self.session.execute(stmt) - await self.session.commit() - logger.info( - f"Research upserted for {crypto_token} " - f"(by={last_researched_by}, opinion={agent_opinion})" + cutoff_date = datetime.now(timezone.utc) - timedelta(days=days) + stmt = select(AgentResearchArtifact).where( + AgentResearchArtifact.agent_id == agent_uuid, + AgentResearchArtifact.created_at >= cutoff_date ) - + if token: + stmt = stmt.where(AgentResearchArtifact.related_tokens.contains([token])) + + result = await self.session.execute(stmt) + artifacts = result.scalars().all() + + research_results = [] + for artifact in artifacts: + research_results.append({ + "query": artifact.query, + "recency": artifact.recency, + "provider": artifact.provider, + "summary_markdown": artifact.summary_markdown, + "citations": artifact.citations, + "raw_results": artifact.raw_results, + "created_at": artifact.created_at.isoformat(), + }) + + logger.info(f"Retrieved {len(research_results)} research results for agent={agent_uuid}") + return research_results + except Exception as e: - await self.session.rollback() - logger.error(f"Failed to upsert research result: {e}") + logger.error(f"Failed to retrieve research results: {e}") raise - async def create_plan_item(self, agent_uuid: UUID, tournament_uuid: UUID, action_type: PlanActionEnum, execute_at: datetime, payload, idempotency_key: str, max_attempts: int = 3) -> PlanItem: stmt = (insert(PlanItem).values(agent_id = agent_uuid, @@ -606,7 +526,7 @@ async def list_plan_items(self, agent_uuid: UUID, tournament_uuid: UUID, async def cancel_plan_item(self, plan_item_id: UUID, reason: str | None = None) -> None: values = { "status": PlanStatusEnum.cancelled, - "updated_at": datetime.now(timezone.utc), + "updated_at": datetime.utcnow(), } if reason: @@ -633,7 +553,7 @@ async def reschedule_plan_item(self, plan_item_id: UUID, new_execute_at: datetim ) .values( execute_at=new_execute_at, - updated_at=datetime.now(timezone.utc), + updated_at=datetime.utcnow(), ) .returning(PlanItem) ) diff --git a/backend/app/agents/tools/research_tool.py b/backend/app/agents/tools/research_tool.py index 2383a73..3a4f22a 100644 --- a/backend/app/agents/tools/research_tool.py +++ b/backend/app/agents/tools/research_tool.py @@ -34,59 +34,18 @@ class ResearchResult: created_at: str -_BULLISH_KEYWORDS = [ - "bullish", "breakout", "growth", "upgrade", "momentum", "rally", "surge", - "uptrend", "accumulation", "outperform", -] -_BEARISH_KEYWORDS = [ - "bearish", "decline", "risk", "hack", "crash", "sell-off", "downgrade", - "downtrend", "liquidation", "underperform", -] - - -def derive_opinion(summary_markdown: str) -> str: - """Keyword-based sentiment derivation from a research summary. - - Returns one of 'positive', 'negative', or 'neutral' (matches OpinionEnum values). - """ - text_lower = summary_markdown.lower() - bull_count = sum(text_lower.count(kw) for kw in _BULLISH_KEYWORDS) - bear_count = sum(text_lower.count(kw) for kw in _BEARISH_KEYWORDS) - - if bull_count > bear_count: - return "positive" - elif bear_count > bull_count: - return "negative" - return "neutral" - - -def select_recency(context: str) -> str: - """Map a context hint to the optimal Perplexity recency filter. - - Args: - context: One of 'breaking', 'trading', 'analysis', 'background'. - """ - mapping = { - "breaking": "1d", - "trading": "1d", - "analysis": "7d", - "background": "30d", - } - return mapping.get(context, "7d") - - class ResearchTool: - """Fetching and summarizing external research via Perplexity API. - - Provides structured prompts, opinion derivation, and recency selection - so every agent gets consistent, cache-friendly research output. + """ fetching and summarizing external research via Perplexity API. + + get info ab tokens, protocols, + market narratives, and crypto news with citations. """ def __init__(self, api_key: str = None): self.api_key = api_key or PERPLEXITY_API_KEY if not self.api_key: raise ResearchError("PERPLEXITY_API_KEY not found in environment") - + self.api_base_url = "https://api.perplexity.ai" self.headers = { "Authorization": f"Bearer {self.api_key}", @@ -94,20 +53,6 @@ def __init__(self, api_key: str = None): } logger.info("ResearchTool initialized") - @staticmethod - def _build_system_prompt(token_symbol: str) -> str: - """Return a structured system prompt that forces consistent sections.""" - return ( - f"You are a crypto research assistant. Organize your response about " - f"{token_symbol} using exactly these markdown sections:\n\n" - f"## {token_symbol} Market Overview\n" - f"## Recent Developments\n" - f"## Sentiment & Narrative\n" - f"## Risks & Catalysts\n" - f"## Summary\n\n" - f"Be concise and factual. Always cite your sources." - ) - def _map_recency_filter(self, recency: str) -> Optional[str]: """Map our recency format to Perplexity's search_recency_filter.""" mapping = { @@ -122,27 +67,24 @@ def research( self, query: str, recency: Optional[Literal["1d", "7d", "30d", "365d"]] = None, - sources: Optional[List[Literal["web", "news", "docs"]]] = None, - max_results: Optional[int] = None, - system_prompt: Optional[str] = None, + sources: Optional[List[Literal["web", "news", "docs"]]] = None, + max_results: Optional[int] = None ) -> ResearchResult: logger.info(f"Research query: {query}") - - default_system = ( - "You are a crypto research assistant. Provide concise, factual summaries " - "about tokens, protocols, market conditions, and crypto news. " - "Focus on key facts, recent developments, risks, and catalysts. " - "Always cite your sources." - ) - + # request payload payload = { "model": "sonar", # mayb switch to "sonar-pro" for more better results "messages": [ { "role": "system", - "content": system_prompt or default_system, + "content": ( + "You are a crypto research assistant. Provide concise, factual summaries " + "about tokens, protocols, market conditions, and crypto news. " + "Focus on key facts, recent developments, risks, and catalysts. " + "Always cite your sources." + ) }, { "role": "user", @@ -226,14 +168,16 @@ def research( return result def research_token(self, token_symbol: str, recency: str = "7d") -> ResearchResult: - """Research a specific token with a structured system prompt.""" + """ + method for researching a specific token. + + """ query = ( f"What are the latest news, developments, and market sentiment for {token_symbol} " f"cryptocurrency? Include any recent catalysts, risks, protocol updates, " f"and narrative shifts." ) - system_prompt = self._build_system_prompt(token_symbol) - return self.research(query, recency=recency, system_prompt=system_prompt) + return self.research(query, recency=recency) # Helper function to convert ResearchResult to dict (for JSON serialization) diff --git a/backend/app/db/models.py b/backend/app/db/models.py index 124571c..f327ac4 100644 --- a/backend/app/db/models.py +++ b/backend/app/db/models.py @@ -157,25 +157,20 @@ class Bet(Base): agent: Mapped["Agent"] = relationship(back_populates="bets") class AgentResearchArtifact(Base): - """Shared research cache keyed by crypto_token. One row per token, UPSERT on refresh.""" - + """Stores research artifacts generated by agents via Perplexity API.""" + __tablename__ = "agent_research_artifact" __table_args__ = ( - Index("ix_research_token_updated", "crypto_token", "updated_at"), + Index("ix_research_agent_created", "agent_id", "created_at"), ) - crypto_token: Mapped[str] = mapped_column(primary_key=True) + crypto_token: Mapped[str] = mapped_column(primary_key = True, default="BTC") id: Mapped[int] = mapped_column(autoincrement=True, unique=True) - last_researched_by: Mapped[Optional[UUID]] = mapped_column(ForeignKey("agent.id"), index=True, default=None) + agent_id: Mapped[UUID] = mapped_column(ForeignKey("agent.id"), index=True) created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), default=lambda: datetime.now(timezone.utc) ) - updated_at: Mapped[datetime] = mapped_column( - DateTime(timezone=True), - default=lambda: datetime.now(timezone.utc), - onupdate=lambda: datetime.now(timezone.utc), - ) query: Mapped[str] = mapped_column(String) recency: Mapped[Optional[str]] = mapped_column(String, default=None) # "1d", "7d", "30d", "365d" provider: Mapped[str] = mapped_column(String, default="perplexity") @@ -183,9 +178,7 @@ class AgentResearchArtifact(Base): citations: Mapped[list[dict[str, Any]]] = mapped_column(JSON, default=list) # [{title, url, date}] raw_results: Mapped[Optional[dict[str, Any]]] = mapped_column(JSON, default=None) related_tokens: Mapped[Optional[list[str]]] = mapped_column(JSON, default=None) - agent_opinion: Mapped[Optional[str]] = mapped_column( - SQLEnum(OpinionEnum, native_enum=False), default=None, nullable=True - ) + agent_opinion: Mapped[OpinionEnum] = mapped_column(SQLEnum(OpinionEnum, native_enum = False)) class PlanItem(Base): __tablename__ = "plan_item" diff --git a/backend/app/scripts/research_table_db.py b/backend/app/scripts/research_table_db.py index c7c0acb..b7ae094 100644 --- a/backend/app/scripts/research_table_db.py +++ b/backend/app/scripts/research_table_db.py @@ -22,49 +22,34 @@ async def migrate(): connect_args={"ssl": "require"}, ) - drop_old_sql = """ - DROP TABLE IF EXISTS agent_research_artifact CASCADE - """ - create_table_sql = """ CREATE TABLE IF NOT EXISTS agent_research_artifact ( - crypto_token VARCHAR NOT NULL PRIMARY KEY, - id SERIAL UNIQUE, - last_researched_by UUID REFERENCES agent(id), + id UUID PRIMARY KEY DEFAULT gen_random_uuid(), + agent_id UUID NOT NULL REFERENCES agent(id), created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), - updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), query TEXT NOT NULL, recency VARCHAR(10), provider VARCHAR(50) DEFAULT 'perplexity', summary_markdown TEXT NOT NULL, citations JSONB DEFAULT '[]'::jsonb, raw_results JSONB, - related_tokens JSONB, - agent_opinion VARCHAR(20) + related_tokens JSONB ) """ create_index_sql = """ - CREATE INDEX IF NOT EXISTS ix_research_token_updated - ON agent_research_artifact(crypto_token, updated_at DESC) - """ - - create_fk_index_sql = """ - CREATE INDEX IF NOT EXISTS ix_research_last_researched_by - ON agent_research_artifact(last_researched_by) + CREATE INDEX IF NOT EXISTS ix_research_agent_created + ON agent_research_artifact(agent_id, created_at DESC) """ async with engine.begin() as conn: try: - print("Dropping old agent_research_artifact table...") - await conn.execute(text(drop_old_sql)) - print("Creating agent_research_artifact table (shared cache schema)...") + print("🏗️ Creating agent_research_artifact table...") await conn.execute(text(create_table_sql)) await conn.execute(text(create_index_sql)) - await conn.execute(text(create_fk_index_sql)) - print("Migration completed successfully!") + print(" Migration completed successfully!") except Exception as e: - print(f"Migration failed: {e}") + print(f" Migration failed: {e}") await engine.dispose() diff --git a/backend/tests/test_research_tool.py b/backend/tests/test_research_tool.py index 8713169..cd3a593 100644 --- a/backend/tests/test_research_tool.py +++ b/backend/tests/test_research_tool.py @@ -4,7 +4,7 @@ load_dotenv() from app.db.database import AsyncSessionLocal -from app.agents.tools.research_tool import ResearchTool, research_result_to_dict, derive_opinion +from app.agents.tools.research_tool import ResearchTool, research_result_to_dict from app.agents.tools.database_tool import DatabaseTool async def test_research_flow(): @@ -19,41 +19,28 @@ async def test_research_flow(): result_dict = research_result_to_dict(result) print(f"2. Converted to dict: {list(result_dict.keys())}") - # Derive opinion - opinion = derive_opinion(result.summary_markdown) - print(f" Derived opinion: {opinion}") - - # Testing the DB upsert - print("3. Testing DB upsert...") - agent_uuid = UUID("22693db3-f7ec-4d5d-a6f2-e9c48a2ac3dc") - + # Testing the DB save + print("3. Testing DB save...") async with AsyncSessionLocal() as session: db_tool = DatabaseTool(session) - await db_tool.upsert_research_result( - crypto_token="ETH", + agent_uuid = UUID("22693db3-f7ec-4d5d-a6f2-e9c48a2ac3dc") + + await db_tool.save_research_result( + agent_uuid=agent_uuid, query="Research ETH", - summary_markdown=result_dict["summary_markdown"], - citations=result_dict["citations"], + result=result_dict, recency="7d", - provider=result_dict.get("provider", "perplexity"), - raw_results=result_dict.get("raw_results"), related_tokens=["ETH"], - agent_opinion=opinion, - last_researched_by=agent_uuid, ) - print(" Upserted to DB!") + print(" Saved to DB!") - # Testing the cache fetch - print("4. Fetching from cache...") + # Testing the DB fetch + print("4. Fetching from DB...") async with AsyncSessionLocal() as session: db_tool = DatabaseTool(session) - cached = await db_tool.get_fresh_research("ETH", "7d") - if cached: - print(f" Cache HIT: {cached['summary_markdown'][:80]}...") - print(f" Opinion: {cached['agent_opinion']}") - else: - print(" Cache MISS (unexpected)") + results = await db_tool.get_recent_research_results(agent_uuid) + print(f" Found {len(results)} research results") if __name__ == "__main__": asyncio.run(test_research_flow())