Skip to content

orchestrator

orchestrator

Executes user queries through the engine or through an agent.

Classes

QueryOrchestrator

QueryOrchestrator(system: OrchestratorDeps)
Source code in src/diapason/system/orchestrator.py
def __init__(self, system: OrchestratorDeps) -> None:
    self._system = system
Methods:
ask
ask(
    query: str,
    *,
    context: bool = True,
    temperature: Optional[float] = None,
    max_tokens: Optional[int] = None,
    agent: Optional[str] = None,
    tools: Optional[List[str]] = None,
    system_prompt: Optional[str] = None,
    operator_id: Optional[str] = None,
    prior_messages: Optional[List[Message]] = None,
    action_mode: str = "off",
) -> Dict[str, Any]

Execute a query through the system and return a result dict.

Source code in src/diapason/system/orchestrator.py
def ask(
    self,
    query: str,
    *,
    context: bool = True,
    temperature: Optional[float] = None,
    max_tokens: Optional[int] = None,
    agent: Optional[str] = None,
    tools: Optional[List[str]] = None,
    system_prompt: Optional[str] = None,
    operator_id: Optional[str] = None,
    prior_messages: Optional[List[Message]] = None,
    action_mode: str = "off",
) -> Dict[str, Any]:
    """Execute a query through the system and return a result dict."""
    s = self._system
    # CLI/SDK counterpart of the desktop chat fast path.  A caller that
    # explicitly chooses an agent, tools, a custom prompt or prior turns
    # retains the exact orchestration semantics they requested.
    if (
        action_mode == "auto"
        and not agent
        and not tools
        and system_prompt is None
        and not prior_messages
    ):
        try:
            from diapason.actions import LightningActionService

            action = LightningActionService(s.config).handle(query)
            if action.handled:
                return {
                    "content": action.message,
                    "usage": {},
                    "model": s.model,
                    "engine": "lightning",
                    "lightning": action.public_metadata(),
                }
        except Exception:
            logger.exception("Lightning action routing failed")
    if temperature is None:
        temperature = s.config.intelligence.temperature
    if max_tokens is None:
        max_tokens = s.config.intelligence.max_tokens

    messages = [Message(role=Role.USER, content=query)]

    if context and s.memory_backend and s.config.agent.context_from_memory:
        try:
            from diapason.tools.storage.context import (
                ContextConfig,
                inject_context,
            )

            ctx_cfg = ContextConfig(
                top_k=s.config.memory.context_top_k,
                min_score=s.config.memory.context_min_score,
                max_context_tokens=s.config.memory.context_max_tokens,
            )
            messages = inject_context(
                query,
                messages,
                s.memory_backend,
                config=ctx_cfg,
            )
        except Exception as exc:
            logger.warning("Failed to inject memory context: %s", exc)

    use_agent = agent or s.agent_name
    if not agent and use_agent != "none":
        detected = self._detect_agent_intent(query)
        if detected:
            use_agent = detected
    if use_agent and use_agent != "none":
        return self._run_agent(
            query,
            messages,
            use_agent,
            tools,
            temperature,
            max_tokens,
            system_prompt=system_prompt,
            operator_id=operator_id,
            prior_messages=prior_messages,
        )

    result = s.engine.generate(
        messages,
        model=s.model,
        temperature=temperature,
        max_tokens=max_tokens,
    )
    return {
        "content": result.get("content", ""),
        "usage": result.get("usage", {}),
        "model": s.model,
        "engine": s.engine_key,
    }