diff --git a/backend/app/agents/executor.py b/backend/app/agents/executor.py index 8e03557..7b6de6f 100644 --- a/backend/app/agents/executor.py +++ b/backend/app/agents/executor.py @@ -171,8 +171,9 @@ def _build_agent_executor(self) -> AgentExecutor: name="research_token", func=self._execute_research_token_wrapper, description=( - "Research a token for news and any market sentiments, " - "Input: token symbol and recency (e.g., 'ETH 7d' for last 7 days)" + "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)" ), ), Tool( @@ -249,7 +250,7 @@ def _build_agent_executor(self) -> AgentExecutor: - Always provide reasoning in your summary - This is a simulation - trades are not executed on-chain - Before making significant trades, research tokens using the research_token tool to check recent news and sentiment -- If recent research already exists, you may reuse it instead of researching again +- The research_token tool uses a shared cache automatically. Just call it when you need research. TOOLS: ------ @@ -347,15 +348,14 @@ def _execute_trade_wrapper(self, trade_input: str) -> str: def _execute_research_token_wrapper(self, input_str: str) -> str: """ - Wrapper for executing research token from LangChain tool. + Cache-first research wrapper. - Args: - input_str: "ETH" or "ETH 7d" - - Returns: - Research summary string + 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 """ - from .tools.research_tool import research_result_to_dict + from .tools.research_tool import research_result_to_dict, derive_opinion try: parts = input_str.split() @@ -363,24 +363,42 @@ def _execute_research_token_wrapper(self, input_str: str) -> str: return "Error: Invalid format. Expected 'TOKEN' or 'TOKEN RECENCY'" token = parts[0].upper() - recency = parts[1] if len(parts) > 1 else "7d" + _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"] - # Call sync research method + # ── Cache MISS — call Perplexity ───────────────────────── research_result = self.research_tool.research_token( token_symbol=token, recency=recency, ) - # Persist result to DB (convert dataclass to dict) - if self.database_tool and self.agent_uuid: + # 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: result_dict = research_result_to_dict(research_result) self._run_async( - self.database_tool.save_research_result( - agent_uuid=self.agent_uuid, + self.database_tool.upsert_research_result( + crypto_token=token, query=f"Research {token}", - result=result_dict, + summary_markdown=result_dict["summary_markdown"], + citations=result_dict["citations"], 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 331a030..b78d187 100644 --- a/backend/app/agents/tools/__init__.py +++ b/backend/app/agents/tools/__init__.py @@ -17,6 +17,7 @@ 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 @@ -29,6 +30,7 @@ "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 cbba8ac..382e039 100644 --- a/backend/app/agents/tools/database_tool.py +++ b/backend/app/agents/tools/database_tool.py @@ -8,18 +8,31 @@ import logging from typing import Optional, Dict, Any, List from uuid import UUID -from datetime import datetime, timezone, timedelta +from datetime import datetime, timezone from decimal import Decimal from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy import select, update from sqlalchemy.dialects.postgresql import insert -from ...db.models import Agent, AgentState, Trade, Tournament, ActionEnum, PlanActionEnum, PlanStatusEnum, PlanItem +from ...db.models import Agent, AgentState, Trade, Tournament, ActionEnum, PlanActionEnum, PlanStatusEnum, PlanItem, AgentResearchArtifact from ..data_classes import Trade as TradeData, Portfolio 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: """ @@ -402,72 +415,139 @@ async def create_agent_if_not_exists( logger.error(f"Failed to create agent: {e}") raise - 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) + # ── 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. + """ try: - 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), + stmt = select(AgentResearchArtifact).where( + AgentResearchArtifact.crypto_token == crypto_token ) - await self.session.execute(stmt) - await self.session.commit() - logger.info(f"Research result saved for agent={agent_uuid}, query='{query}'") + 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)" + ) + 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(), + } + except Exception as e: - 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 + 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) try: - cutoff_date = datetime.now(timezone.utc) - timedelta(days=days) - stmt = select(AgentResearchArtifact).where( - AgentResearchArtifact.agent_id == agent_uuid, - AgentResearchArtifact.created_at >= cutoff_date + 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, + }, + ) ) - 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 - + + 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})" + ) + except Exception as e: - logger.error(f"Failed to retrieve research results: {e}") + await self.session.rollback() + logger.error(f"Failed to upsert research result: {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, @@ -526,7 +606,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.utcnow(), + "updated_at": datetime.now(timezone.utc), } if reason: @@ -553,7 +633,7 @@ async def reschedule_plan_item(self, plan_item_id: UUID, new_execute_at: datetim ) .values( execute_at=new_execute_at, - updated_at=datetime.utcnow(), + updated_at=datetime.now(timezone.utc), ) .returning(PlanItem) ) diff --git a/backend/app/agents/tools/research_tool.py b/backend/app/agents/tools/research_tool.py index 3a4f22a..2383a73 100644 --- a/backend/app/agents/tools/research_tool.py +++ b/backend/app/agents/tools/research_tool.py @@ -34,18 +34,59 @@ 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. - - get info ab tokens, protocols, - market narratives, and crypto news with citations. + """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. """ 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}", @@ -53,6 +94,20 @@ 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 = { @@ -67,24 +122,27 @@ 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 + sources: Optional[List[Literal["web", "news", "docs"]]] = None, + max_results: Optional[int] = None, + system_prompt: Optional[str] = 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": ( - "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." - ) + "content": system_prompt or default_system, }, { "role": "user", @@ -168,16 +226,14 @@ def research( return result def research_token(self, token_symbol: str, recency: str = "7d") -> ResearchResult: - """ - method for researching a specific token. - - """ + """Research a specific token with a structured system prompt.""" 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." ) - return self.research(query, recency=recency) + system_prompt = self._build_system_prompt(token_symbol) + return self.research(query, recency=recency, system_prompt=system_prompt) # 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 0baf2a4..2395a11 100644 --- a/backend/app/db/models.py +++ b/backend/app/db/models.py @@ -157,20 +157,25 @@ class Bet(Base): agent: Mapped["Agent"] = relationship(back_populates="bets") class AgentResearchArtifact(Base): - """Stores research artifacts generated by agents via Perplexity API.""" - + """Shared research cache keyed by crypto_token. One row per token, UPSERT on refresh.""" + __tablename__ = "agent_research_artifact" __table_args__ = ( - Index("ix_research_agent_created", "agent_id", "created_at"), + Index("ix_research_token_updated", "crypto_token", "updated_at"), ) - crypto_token: Mapped[str] = mapped_column(primary_key = True, default="BTC") + crypto_token: Mapped[str] = mapped_column(primary_key=True) id: Mapped[int] = mapped_column(autoincrement=True, unique=True) - agent_id: Mapped[UUID] = mapped_column(ForeignKey("agent.id"), index=True) + last_researched_by: Mapped[Optional[UUID]] = mapped_column(ForeignKey("agent.id"), index=True, default=None) 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") @@ -178,7 +183,9 @@ 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[OpinionEnum] = mapped_column(SQLEnum(OpinionEnum, native_enum = False)) + agent_opinion: Mapped[Optional[str]] = mapped_column( + SQLEnum(OpinionEnum, native_enum=False), default=None, nullable=True + ) 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 b7ae094..c7c0acb 100644 --- a/backend/app/scripts/research_table_db.py +++ b/backend/app/scripts/research_table_db.py @@ -22,34 +22,49 @@ 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 ( - id UUID PRIMARY KEY DEFAULT gen_random_uuid(), - agent_id UUID NOT NULL REFERENCES agent(id), + crypto_token VARCHAR NOT NULL PRIMARY KEY, + id SERIAL UNIQUE, + last_researched_by UUID 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 + related_tokens JSONB, + agent_opinion VARCHAR(20) ) """ create_index_sql = """ - CREATE INDEX IF NOT EXISTS ix_research_agent_created - ON agent_research_artifact(agent_id, created_at DESC) + 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) """ async with engine.begin() as conn: try: - print("🏗️ Creating agent_research_artifact table...") + print("Dropping old agent_research_artifact table...") + await conn.execute(text(drop_old_sql)) + print("Creating agent_research_artifact table (shared cache schema)...") await conn.execute(text(create_table_sql)) await conn.execute(text(create_index_sql)) - print(" Migration completed successfully!") + await conn.execute(text(create_fk_index_sql)) + 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 cd3a593..8713169 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 +from app.agents.tools.research_tool import ResearchTool, research_result_to_dict, derive_opinion from app.agents.tools.database_tool import DatabaseTool async def test_research_flow(): @@ -19,28 +19,41 @@ async def test_research_flow(): result_dict = research_result_to_dict(result) print(f"2. Converted to dict: {list(result_dict.keys())}") - # Testing the DB save - print("3. Testing DB save...") + # 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") + async with AsyncSessionLocal() as session: db_tool = DatabaseTool(session) - agent_uuid = UUID("22693db3-f7ec-4d5d-a6f2-e9c48a2ac3dc") - - await db_tool.save_research_result( - agent_uuid=agent_uuid, + await db_tool.upsert_research_result( + crypto_token="ETH", query="Research ETH", - result=result_dict, + summary_markdown=result_dict["summary_markdown"], + citations=result_dict["citations"], 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(" Saved to DB!") + print(" Upserted to DB!") - # Testing the DB fetch - print("4. Fetching from DB...") + # Testing the cache fetch + print("4. Fetching from cache...") async with AsyncSessionLocal() as session: db_tool = DatabaseTool(session) - results = await db_tool.get_recent_research_results(agent_uuid) - print(f" Found {len(results)} research results") + 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)") if __name__ == "__main__": asyncio.run(test_research_flow())