Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,11 @@ GEMINI_API_KEY=
ANTHROPIC_API_KEY=
GROQ_API_KEY=
OPENROUTER_API_KEY=
# Optional: needed only when the Paperclip literature tool is enabled.
PAPERCLIP_API_KEY=
# Optional: 1/true forces Paperclip's MCP-only transport. Leave blank for the
# default (REST enabled); any non-empty value other than a false-y one enables it.
PAPERCLIP_DISABLE_REST=

NEO4J_USER=
NEO4J_PASSWORD=
Expand Down
8 changes: 8 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,8 @@ GEMINI_API_KEY=
ANTHROPIC_API_KEY=
GROQ_API_KEY=
OPENROUTER_API_KEY=
PAPERCLIP_API_KEY=
PAPERCLIP_DISABLE_REST=

NEO4J_USER=
NEO4J_PASSWORD=
Expand All @@ -68,6 +70,12 @@ BROWSER_COOKIE_SECRET=
RATE_LIMIT_IP_HASH_SECRET=
```

Paperclip and PubTator3 are optional literature-evidence tools and are disabled
by default. Each can be enabled independently from the chat settings. When both
are enabled, they run concurrently after the biological-relevance check;
Paperclip requires `PAPERCLIP_API_KEY`. Set `PAPERCLIP_DISABLE_REST=1` to force
its MCP-only transport.

## Run the Backend API

Start the FastAPI backend from the repository root:
Expand Down
52 changes: 52 additions & 0 deletions crossbar_llm/agent_tools/callback_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -319,6 +319,58 @@ def reset(self) -> None:
)


def merge_usage_summaries(*summaries: dict[str, Any]) -> dict[str, Any]:
"""Combine several `get_summary()` results into one request-level total.

Used where one request is served by more than one callback — the core agent
runs strict (a provider that hides usage metadata should fail loudly rather
than bill silently), while optional side agents run lenient because their
JSON-fallback path legitimately produces responses without it. Keeping two
handlers preserves both behaviours; this puts the numbers back together.

Node names are expected to be unique across summaries (side agents
namespace theirs). On a collision the later summary wins for `model`, and
token counts are summed.
"""
merged_nodes: dict[str, dict[str, Union[int, str]]] = {}
totals = UsageCounter()
models_by_node: dict[str, list[str]] = {}
session_id: str | None = None

for summary in summaries:
if not summary:
continue
session_id = session_id or summary.get("session_id")

for node, record in (summary.get("per_node_usage") or {}).items():
existing = merged_nodes.get(node)
if existing is None:
merged_nodes[node] = dict(record)
continue
for key in ("input_tokens", "output_tokens", "total_tokens",
"cache_read", "cache_write", "reasoning", "call_count"):
existing[key] = (existing.get(key, 0) or 0) + (record.get(key, 0) or 0)
if record.get("model"):
existing["model"] = record["model"]

aggregated = summary.get("aggregated_usage") or {}
totals.add_usage(aggregated.get("totals") or {})
for node, models in (aggregated.get("models_by_node") or {}).items():
bucket = models_by_node.setdefault(node, [])
for model in models:
if model not in bucket:
bucket.append(model)

return {
"session_id": session_id,
"per_node_usage": merged_nodes,
"aggregated_usage": {
"totals": totals.to_dict(),
"models_by_node": models_by_node,
},
}


# ---------------------------------------------------------------
# CHECK:
# WHAT HAPPEN WITH DIFFERENT LLM PROVIDERS OTHER THAN OPENAI?
Expand Down
96 changes: 68 additions & 28 deletions crossbar_llm/agent_tools/cypher_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -194,49 +194,89 @@ def trim_messages(self, state: CypherAgentState, new_messages: list[BaseMessage]
# when the list is still below the limit.
return [RemoveMessage(id=REMOVE_ALL_MESSAGES), *combined_messages[-AgentConfig().keep_last_n_messages:]]

@log_execution_time(logger, component="CypherAgent.initialize_state")
def biological_relevance_validation_node(self, state: CypherAgentState):

logger.info(
"Starting biological relevance validation",
event_type="biological_relevance_validation_started",
component="CypherAgent.biological_relevance_validation_node",
question=state["question"]
)

# The sync node and the async preflight below share these three helpers so
# the prompt, the metadata tag and the result shape exist in exactly one
# place. Only the invoke/ainvoke call itself differs, which is irreducible
# while the graph still supports `graph.invoke()`.
def _relevance_request(self, question: str):
prompt = ChatPromptTemplate.from_messages([
SystemMessagePromptTemplate.from_template(
BIOLOGICAL_RELEVANCE_VALIDATION_TEMPLATE
),
HumanMessagePromptTemplate.from_template("{question}")
])

verdict_llm = self.llm_factory.create_biological_relevance_validator_llm()

messages = prompt.format_messages(question=state["question"])
logger.info(
"Starting biological relevance validation",
event_type="biological_relevance_validation_started",
component="CypherAgent._relevance_request",
question=question
)

verdict = verdict_llm.invoke(
messages,
config={
"metadata": {
"node_name": "biological_relevance_validation",
}
}
return (
self.llm_factory.create_biological_relevance_validator_llm(),
prompt.format_messages(question=question),
{"metadata": {"node_name": "biological_relevance_validation"}},
)

@staticmethod
def _relevance_result(verdict, question: str) -> dict[str, object]:
parsed = verdict["parsed"]

logger.info(
"Completed biological relevance validation",
event_type="biological_relevance_validation_completed",
component="CypherAgent.biological_relevance_validation_node",
question=state["question"],
verdict=verdict["parsed"].relevant,
reason=verdict["parsed"].reason
component="CypherAgent._relevance_result",
question=question,
verdict=parsed.relevant,
reason=parsed.reason
)

return {
"biological_relevance": verdict["parsed"].relevant,
"final_answer": f"Your question is outside the biological/biomedical domain.\nReason: {verdict['parsed'].reason}" if verdict["parsed"].relevant is False else None,
}

"biological_relevance": parsed.relevant,
"final_answer": (
"Your question is outside the biological/biomedical domain.\n"
f"Reason: {parsed.reason}"
if parsed.relevant is False
else None
),
}

@log_execution_time(logger, component="CypherAgent.initialize_state")
def biological_relevance_validation_node(self, state: CypherAgentState):
# The API orchestrator runs `avalidate_biological_relevance` ahead of the
# graph so it can decide whether to launch the literature tools. When it
# has, the verdict is already in state and paying for a second identical
# LLM call would be pure waste — the router downstream reads the same key
# either way, so returning no update keeps the routing identical.
if state.get("biological_relevance") is not None:
logger.info(
"Reusing pre-validated biological relevance verdict",
event_type="biological_relevance_validation_reused",
component="CypherAgent.biological_relevance_validation_node",
question=state["question"],
verdict=state["biological_relevance"],
)
return {}

verdict_llm, messages, config = self._relevance_request(state["question"])
verdict = verdict_llm.invoke(messages, config=config)
return self._relevance_result(verdict, state["question"])

async def avalidate_biological_relevance(
self,
question: str,
) -> dict[str, object]:
"""Async relevance preflight used by the API request orchestrator.

Seed the returned keys into the graph's initial state and
`biological_relevance_validation_node` will reuse the verdict instead of
re-running it.
"""
verdict_llm, messages, config = self._relevance_request(question)
verdict = await verdict_llm.ainvoke(messages, config=config)
return self._relevance_result(verdict, question)

def route_after_biological_relevance_validation(self, state: CypherAgentState) -> Literal["entity_resolution", "generate_cypher", "end"]:
if state["biological_relevance"] is False:
logger.info(
Expand Down
43 changes: 42 additions & 1 deletion crossbar_llm/api/core/deps.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,49 @@
from functools import lru_cache

from fastapi import HTTPException, status
from pydantic import ValidationError

from crossbar_llm.agent_tools.logging_config import get_logger
from crossbar_llm.api.core.settings import Settings
from crossbar_llm.api.services.agent_service import AgentService

logger = get_logger(__name__)


@lru_cache()
def get_runtime_service() -> AgentService:
return AgentService()
try:
return AgentService()
except ValidationError as exc:
missing_fields = sorted({
str(error["loc"][-1])
for error in exc.errors()
if error.get("type") == "missing" and error.get("loc")
})

# Always logged in full — the operator fixing this needs the list.
logger.error(
"Runtime service could not be configured",
event_type="runtime_service_misconfigured",
component="deps.get_runtime_service",
missing_fields=missing_fields,
exc_info=exc,
)

# Only echoed to the client in development. The names alone leak no
# secrets, but on a public deployment they confirm the backend stack to
# anyone who happens to hit the API while it is misconfigured, and the
# sentence below already tells a user everything they can act on.
settings = Settings()
missing_hint = (
f" Missing environment variables: {', '.join(missing_fields)}."
if missing_fields and settings.is_dev
else ""
)
raise HTTPException(
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
detail=(
"The knowledge-graph backend is not configured, so chat and "
f"vector queries cannot run.{missing_hint}"
),
) from exc
Loading