Skip to content

Core

Runtime

The one loop, and the policies that make it the lead or a role.

helioai.runtime.runner

The agent loop, written once.

Runner.run is the turn-taking skeleton both stream_chat and stream_subagent used to carry as separate copies: call the model with a compacted history, start the turn's tool calls together, dispatch each in the model's order, review the figures, emit the events, append the results, retry once when the answer quotes ids that exist in no catalogue, stop at the cap. What differs between the lead and a role is a Policy and two small hooks — how a tool call may be intercepted before it reaches the registry (the lead's task and internal tools, a role's whitelist), and what to do with each model call's usage.

The run yields the same {"event", "data"} dicts the loops yielded, in the same order, and ends with a RunEnd the wrapper turns into its own closing events — the lead's reply … done, a role's sub_agent_end.

RunEnd dataclass

How a run ended, handed to the wrapper as the last item of Runner.run.

Attributes:

Name Type Description
final_text str | None

The model's closing text; None when the run was capped and there is no answer to check.

turns int

Model calls made.

artifacts list[dict]

Every artifact the run produced, without a role's sub_agent_ctx — what check_answer confronts the answer with.

usage dict

Token counts summed over the run's model calls, as the providers reported them.

capped bool

The turn budget ran out before the model answered.

empty bool

The model returned neither text nor a tool call and the policy stops there (the lead's output-budget failure).

claims list[dict]

The numbers the model says its answer states, each with name, value, units and source, when it closed with final_answer; empty for prose.

Source code in helioai/runtime/runner.py
@dataclass
class RunEnd:
    """How a run ended, handed to the wrapper as the last item of `Runner.run`.

    Attributes:
        final_text: The model's closing text; `None` when the run was capped and there
            is no answer to check.
        turns: Model calls made.
        artifacts: Every artifact the run produced, without a role's `sub_agent_ctx` —
            what `check_answer` confronts the answer with.
        usage: Token counts summed over the run's model calls, as the providers
            reported them.
        capped: The turn budget ran out before the model answered.
        empty: The model returned neither text nor a tool call and the policy stops
            there (the lead's output-budget failure).
        claims: The numbers the model says its answer states, each with `name`, `value`,
            `units` and `source`, when it closed with `final_answer`; empty for prose.
    """

    final_text: str | None
    turns: int
    artifacts: list[dict]
    usage: dict
    capped: bool = False
    empty: bool = False
    claims: list[dict] = field(default_factory=list)

Runner dataclass

One run of the loop under one policy.

Attributes:

Name Type Description
policy Policy

What makes this run the lead or a role.

llm LLMClient

The provider client the run calls.

intercept Intercept | None

Called with each tool call and the turn before the registry is consulted. Returns None to let the call reach the registry, or an async iterator that yields events to forward, then exactly one ToolResult, then any events to emit after the result's own (tool_result, artifacts) — how the lead re-emits a sub_agent_end and a plan.

on_llm_call OnLLMCall | None

Called with the turn and the model's reply after each call; the lead records the usage row there. The run also sums usage into usage.

registry ToolRegistry

Where tool calls are dispatched. The process-wide registry unless a caller hands in another — the wrappers pass their own module's reference, which is what the tests that stub a registry patch.

ctx RunContext | None

Who runs, in which session, writing where. Bound to the workspace contextvars for the whole run — the one place they are set — and the source of the directories the writing tools receive as trusted arguments. None leaves the ambient bindings alone, for a caller that manages its own.

artifacts, (usage, turns)

Progress so far, readable while the run is in flight and after it raised — a role that blew up still reports what it measured.

Source code in helioai/runtime/runner.py
@dataclass
class Runner:
    """One run of the loop under one policy.

    Attributes:
        policy: What makes this run the lead or a role.
        llm: The provider client the run calls.
        intercept: Called with each tool call and the turn before the registry is
            consulted. Returns `None` to let the call reach the registry, or an async
            iterator that yields events to forward, then exactly one `ToolResult`, then
            any events to emit *after* the result's own (`tool_result`, artifacts) — how
            the lead re-emits a `sub_agent_end` and a `plan`.
        on_llm_call: Called with the turn and the model's reply after each call; the
            lead records the usage row there. The run also sums usage into `usage`.
        registry: Where tool calls are dispatched. The process-wide registry unless a
            caller hands in another — the wrappers pass their own module's reference,
            which is what the tests that stub a registry patch.
        ctx: Who runs, in which session, writing where. Bound to the workspace
            contextvars for the whole run — the one place they are set — and the source
            of the directories the writing tools receive as trusted arguments. None
            leaves the ambient bindings alone, for a caller that manages its own.
        artifacts, usage, turns: Progress so far, readable while the run is in flight
            and after it raised — a role that blew up still reports what it measured.
    """

    policy: Policy
    llm: LLMClient
    intercept: Intercept | None = None
    on_llm_call: OnLLMCall | None = None
    registry: ToolRegistry = field(default_factory=lambda: _registry_module.registry)
    ctx: RunContext | None = None
    artifacts: list[dict] = field(default_factory=list, init=False)
    usage: dict = field(default_factory=_zero_usage, init=False)
    turns: int = field(default=0, init=False)
    revealed: set[str] = field(default_factory=set, init=False)

    async def run(self, history: list[Message]) -> AsyncIterator[dict | RunEnd]:
        """Drive the model over `history` until it answers, stops or hits the cap.

        Args:
            history: The conversation so far, appended to in place — the assistant's
                replies and every tool result land here, so the caller persists the
                same list it passed.

        Yields:
            Event dicts, then one `RunEnd`.
        """
        with self.ctx.bound() if self.ctx is not None else nullcontext():
            started: dict = {}
            try:
                async for item in self._turns(history, started):
                    yield item
            finally:
                cancel_pending(started)

    async def _turns(self, history: list[Message], started: dict) -> AsyncIterator[dict | RunEnd]:
        policy = self.policy
        extra = policy.event_extra
        retried_bogus_ids = False
        searches, touched_data, search_corrected = 0, False, False
        for i in range(policy.max_turns):
            turn = self.turns = i + 1
            history[:] = strip_orphan_tool_calls(history)
            response: Message | None = None
            async for item in self._model_turn(history, turn, first=(i == 0)):
                if isinstance(item, Message):
                    response = item
                else:
                    yield make("reply_delta", text=item, **extra)
            assert response is not None
            history.append(response)

            claims: list[dict] = []
            if policy.final_answer and self._is_final_answer(response):
                # The answer arrived as a tool call. The history keeps it as the plain
                # assistant message it is — the export, the replay and the next turn's
                # context read assistant text, and the wire needs no tool reply for a
                # call that no longer exists.
                args = response.tool_calls[0].arguments or {}
                final_text = str(args.get("answer") or "")
                claims = _claims_from(args.get("claims"))
                history[-1] = replace(response, content=final_text, tool_calls=None)
                response = history[-1]

            if not response.tool_calls:
                final_text = response.content or ""
                if policy.stop_on_empty_reply and not final_text.strip():
                    yield self._end("", empty=True)
                    return
                if policy.bogus_retry and not retried_bogus_ids:
                    from helioai.tools.rag import extract_ids, unknown_ids

                    bogus = unknown_ids(extract_ids(final_text))
                    if bogus:
                        retried_bogus_ids = True
                        log.warning("invented_ids_retry", agent=policy.name, ids=bogus, turn=turn)
                        note = unknown_id_correction(bogus)
                        history.append(Message(role="user", content=note, origin="correction"))
                        yield make("correction", ids=bogus, text=note, **extra)
                        continue
                yield self._end(final_text, claims=claims)
                return

            if policy.comment_replies and response.content and response.content.strip():
                yield make("reply", text=response.content)

            names = [tc.name for tc in response.tool_calls]
            searches += sum(n in policy.search_tools_names for n in names)
            touched_data = touched_data or any(n in policy.data_tools_names for n in names)
            # The dict is shared with `run`, whose `finally` cancels whatever is left in
            # it: cleared and refilled rather than rebound, so it stays the same object.
            started.clear()
            started.update(
                start_tool_calls(
                    response.tool_calls,
                    allowed=policy.allowed,
                    registry=self.registry,
                    ctx=self.ctx,
                )
            )
            for tc in response.tool_calls:
                log.info("tool_call_issued", agent=policy.name, turn=turn, tool=tc.name)
                yield make(
                    "tool_call",
                    turn=turn,
                    name=tc.name,
                    arguments=tc.arguments,
                    display=describe_tool_call(tc.name, tc.arguments),
                    **extra,
                )
                result, trailing = None, []
                try:
                    if tc.name == SEARCH_TOOLS_NAME and policy.deferred:
                        result = self._search_tools(tc)
                    elif tc.name == FINAL_ANSWER_NAME and policy.final_answer:
                        result = ToolResult.failure(
                            tc.name,
                            "final_answer must be called alone, once the other tool calls of "
                            "the turn have returned; call it again by itself",
                        )
                    elif tc.name in policy.deferred and tc.name not in self.revealed:
                        # The model named a deferred tool without asking for it: the name
                        # was right, so it already knows the tool, and refusing would only
                        # cost a turn. Revealed and run.
                        self.revealed.add(tc.name)
                        log.info("deferred_tool_called_directly", agent=policy.name, tool=tc.name)
                    handler = (
                        self.intercept(tc, turn) if self.intercept and result is None else None
                    )
                    if handler is not None:
                        async for item in handler:
                            if isinstance(item, ToolResult):
                                result = item
                            elif result is None:
                                yield item
                            else:
                                trailing.append(item)
                    if result is None:
                        result = await self._dispatch(tc, started)
                except Exception as e:
                    log.exception("tool_call_failed", agent=policy.name, turn=turn, tool=tc.name)
                    result = ToolResult.failure(tc.name, str(e) or type(e).__name__)

                result, figure_verdict = await maybe_review(tc.name, result)
                if figure_verdict:
                    yield make("figure_review", turn=turn, text=figure_verdict, **extra)
                result = recipe_available(tc.name, tc.arguments, result, history)
                for ev in emit_post_tool_events(
                    tc.name,
                    result,
                    tool_result_extra={"turn": turn, **extra},
                    common_extra=extra,
                ):
                    if ev["event"] == "artifact":
                        self.artifacts.append(
                            {k: v for k, v in ev["data"].items() if k != "sub_agent_ctx"}
                        )
                    yield ev
                for ev in trailing:
                    yield ev
                history.append(
                    Message(
                        role="tool",
                        tool_call_id=tc.id,
                        name=tc.name,
                        content=_history_tool_result(tc.name, result.for_llm()),
                    )
                )
            if (
                policy.search_budget
                and not touched_data
                and not search_corrected
                and searches > policy.search_budget
            ):
                search_corrected = True
                ids = _ids_in_results(history, policy.search_tools_names)
                note = search_loop_correction(searches, ids)
                log.warning(
                    "search_loop_corrected", agent=policy.name, searches=searches, turn=turn
                )
                history.append(Message(role="user", content=note, origin="correction"))
                yield make("correction", ids=ids, text=note, **extra)
        yield self._end(None, capped=True)

    async def _model_turn(
        self, history: list[Message], turn: int, *, first: bool
    ) -> AsyncIterator[str | Message]:
        """One model call: text deltas when the policy streams, then the reply."""
        policy = self.policy
        log.info("llm_call_start", agent=policy.name, turn=turn, n_messages=len(history))
        t0 = time.monotonic()
        kwargs: dict[str, Any] = {"system_prompt": policy.system_prompt}
        # A policy with no opinion on tool_choice leaves the client's default alone,
        # exactly as the lead always did; a role asks for a tool on its first turn.
        if policy.tool_choice_first != "auto":
            kwargs["tool_choice"] = policy.tool_choice_first if first else "auto"
        payload, tools = compact_history(history), self.visible_tools()
        # A duck-typed client without `stream_chat` is a client that does not stream.
        stream = getattr(self.llm, "stream_chat", None) if policy.stream_replies else None
        if stream is not None:
            response = None
            async for item in stream(payload, tools, **kwargs):
                if isinstance(item, Message):
                    response = item
                elif item:
                    yield item
            assert response is not None, "stream_chat must end with the reply"
        else:
            response = await self.llm.chat(payload, tools, **kwargs)
        log.info(
            "llm_call_end",
            agent=policy.name,
            turn=turn,
            duration_ms=int((time.monotonic() - t0) * 1000),
            has_tool_calls=bool(response.tool_calls),
        )
        self.usage["prompt_tokens"] += response.prompt_tokens
        self.usage["completion_tokens"] += response.completion_tokens
        self.usage["cached_tokens"] += response.cached_tokens
        self.usage["n_calls"] += 1
        if self.on_llm_call is not None:
            self.on_llm_call(turn, response)
        yield response

    def visible_tools(self) -> list[ToolDef]:
        """The definitions the model is shown this call: the policy's tools minus the
        deferred ones it has not asked for, plus `search_tools` while any are withheld."""
        policy = self.policy
        withheld = policy.deferred - self.revealed
        shown = [t for t in policy.tools if t.name not in withheld]
        if withheld:
            shown.append(SEARCH_TOOLS_DEF)
        if policy.final_answer:
            shown.append(FINAL_ANSWER_DEF)
        return shown

    @staticmethod
    def _is_final_answer(response: Message) -> bool:
        calls = response.tool_calls or []
        return len(calls) == 1 and calls[0].name == FINAL_ANSWER_NAME

    def _search_tools(self, tc: ToolCall) -> ToolResult:
        query = str((tc.arguments or {}).get("query") or "").lower()
        words = [w for w in re.findall(r"[a-z0-9_]+", query) if len(w) > 2]
        deferred = [t for t in self.policy.tools if t.name in self.policy.deferred]
        matches = [
            t
            for t in deferred
            if any(w in t.name.lower() or w in t.description.lower() for w in words)
        ]
        # A query that matches nothing still means "I need one of these": all of them
        # are revealed rather than sending the model on a second guess.
        chosen = matches or deferred
        self.revealed.update(t.name for t in chosen)
        return ToolResult.from_raw(
            SEARCH_TOOLS_NAME,
            {
                "enabled": [t.name for t in chosen],
                "tools": [{"name": t.name, "description": t.description} for t in chosen],
                "note": "these tools are callable from your next turn",
            },
        )

    async def _dispatch(self, tc: ToolCall, started: dict) -> ToolResult:
        if tc.id in started:
            return await started[tc.id]
        trusted = trusted_args(tc.name, self.ctx, no_network=self.policy.sandbox_no_network)
        return await self.registry.call_tool(tc.name, tc.arguments, trusted=trusted)

    def _end(
        self,
        final_text: str | None,
        *,
        capped: bool = False,
        empty: bool = False,
        claims: list[dict] | None = None,
    ) -> RunEnd:
        return RunEnd(
            final_text=final_text,
            turns=self.turns,
            artifacts=self.artifacts,
            usage=self.usage,
            capped=capped,
            empty=empty,
            claims=list(claims or []),
        )

run async

run(history: list[Message]) -> AsyncIterator[dict | RunEnd]

Drive the model over history until it answers, stops or hits the cap.

Parameters:

Name Type Description Default
history list[Message]

The conversation so far, appended to in place — the assistant's replies and every tool result land here, so the caller persists the same list it passed.

required

Yields:

Type Description
AsyncIterator[dict | RunEnd]

Event dicts, then one RunEnd.

Source code in helioai/runtime/runner.py
async def run(self, history: list[Message]) -> AsyncIterator[dict | RunEnd]:
    """Drive the model over `history` until it answers, stops or hits the cap.

    Args:
        history: The conversation so far, appended to in place — the assistant's
            replies and every tool result land here, so the caller persists the
            same list it passed.

    Yields:
        Event dicts, then one `RunEnd`.
    """
    with self.ctx.bound() if self.ctx is not None else nullcontext():
        started: dict = {}
        try:
            async for item in self._turns(history, started):
                yield item
        finally:
            cancel_pending(started)

visible_tools

visible_tools() -> list[ToolDef]

The definitions the model is shown this call: the policy's tools minus the deferred ones it has not asked for, plus search_tools while any are withheld.

Source code in helioai/runtime/runner.py
def visible_tools(self) -> list[ToolDef]:
    """The definitions the model is shown this call: the policy's tools minus the
    deferred ones it has not asked for, plus `search_tools` while any are withheld."""
    policy = self.policy
    withheld = policy.deferred - self.revealed
    shown = [t for t in policy.tools if t.name not in withheld]
    if withheld:
        shown.append(SEARCH_TOOLS_DEF)
    if policy.final_answer:
        shown.append(FINAL_ANSWER_DEF)
    return shown

search_loop_correction

search_loop_correction(searches: int, ids: list[str]) -> str

The note handed to a role that keeps searching instead of downloading.

Parameters:

Name Type Description Default
searches int

How many lookups it has made.

required
ids list[str]

The product ids its searches already returned, best first.

required

Returns:

Type Description
str

The correction text; the ids are listed so the next turn has nothing to look up.

Source code in helioai/runtime/runner.py
def search_loop_correction(searches: int, ids: list[str]) -> str:
    """The note handed to a role that keeps searching instead of downloading.

    Args:
        searches: How many lookups it has made.
        ids: The product ids its searches already returned, best first.

    Returns:
        The correction text; the ids are listed so the next turn has nothing to look up.
    """
    listed = "\n".join(f"  - {i}" for i in ids[:12]) or "  (none parsed — read your last results)"
    return (
        f"⚠️ AUTOMATED CORRECTION — {searches} searches and no data yet. The products you need "
        f"are in the results you already have; a variable that is not listed under a dataset "
        f"does not exist under that name (SWE has no `Proton_Temp`; its temperature is "
        f"`Proton_W_*`). Do not search again: pick from these ids and call the data tool now:\n"
        f"{listed}"
    )

helioai.runtime.policies

What makes the one loop the lead agent or a delegated role.

A Policy is data: the prompt, the tools the model sees and the ones it may actually call, the turn budget, the first turn's tool_choice, and the few behavioural switches the two loops disagreed on. The lead's policy is built from the constants of agent_loop; a role's from sub_agents.AGENT_ROLES, which stays the one place roles are declared.

Policy dataclass

How one run of the loop behaves.

Attributes:

Name Type Description
name str

"lead" or the role name; used in logs and as the agent of usage rows.

system_prompt str

The instructions placed before the history on every call.

tools tuple[ToolDef, ...]

The definitions the model is shown.

max_turns int

How many model calls a run may make before it is capped.

tool_choice_first str

tool_choice of the first call — required for a role, so it cannot answer from memory without looking anything up; auto for the lead.

allowed frozenset[str] | None

The tools a run may call. None for the lead, who may call anything the registry holds; a role's whitelist is enforced at dispatch, not advised.

sandbox_no_network bool

Whether run_python is denied a network namespace.

sub_agent_ctx Mapping[str, str] | None

{role, task_id} added to every event a role yields, so an interface can nest its trace under the lead's; None for the lead.

comment_replies bool

Whether text the model writes alongside tool calls is shown as a reply. The lead thinks aloud for the reader; a role's asides are noise the lead never sees.

stream_replies bool

Whether the model's text is shown as it is generated (reply_delta events, then the usual reply). The lead's answer is read by a person waiting for it; a role's goes to the lead, whole.

stop_on_empty_reply bool

Whether a response with neither text nor tool calls ends the run as a failure. The lead's does — the usual cause is the output budget, and the error names the setting to raise; a role's simply ends with an empty summary.

bogus_retry bool

Whether an answer quoting ids absent from the catalogue buys the model one more turn, with the correction appended, before it is accepted.

provider str | None

The provider this run's client talks to — for the usage rows, which bill per provider. None means the lead's configured one.

model str | None

The model, when the run was given one of its own (HELIOAI_ROLE_MODELS).

deferred frozenset[str]

Tools whose definitions are withheld from the model until it asks for them with search_tools, or calls one by name. Twenty-one definitions rode on every call of a lead turn; the ten formulary and catalogue ones are used in a minority of sessions and cost a third of that payload each time.

final_answer bool

Whether the model may close the run with final_answer(answer, claims) — the answer plus the list of numbers it states, each with its source — instead of a plain message. Prose is still accepted; the claims are what lets a verdict compare numbers by name instead of by regex.

search_budget int

How many lookups (search_tools_names) a run may make before it has touched any data. A data_analyst spent all twelve of its turns on search_parameters, re-asking for ids it had been handed on the first call because they were not spelled the way it imagined; the run ended with nothing downloaded. Past the budget, one correction is appended telling it to use the ids it has. 0 disables — a parameter_hunter's whole job is to search.

search_tools_names frozenset[str]

The tools that count as lookups.

data_tools_names frozenset[str]

The tools whose first call ends the lookup phase.

Source code in helioai/runtime/policies.py
@dataclass(frozen=True)
class Policy:
    """How one run of the loop behaves.

    Attributes:
        name: `"lead"` or the role name; used in logs and as the `agent` of usage rows.
        system_prompt: The instructions placed before the history on every call.
        tools: The definitions the model is shown.
        max_turns: How many model calls a run may make before it is capped.
        tool_choice_first: `tool_choice` of the first call — `required` for a role, so
            it cannot answer from memory without looking anything up; `auto` for the lead.
        allowed: The tools a run may call. `None` for the lead, who may call anything
            the registry holds; a role's whitelist is enforced at dispatch, not advised.
        sandbox_no_network: Whether `run_python` is denied a network namespace.
        sub_agent_ctx: `{role, task_id}` added to every event a role yields, so an
            interface can nest its trace under the lead's; `None` for the lead.
        comment_replies: Whether text the model writes alongside tool calls is shown as
            a `reply`. The lead thinks aloud for the reader; a role's asides are noise
            the lead never sees.
        stream_replies: Whether the model's text is shown as it is generated
            (`reply_delta` events, then the usual `reply`). The lead's answer is read by
            a person waiting for it; a role's goes to the lead, whole.
        stop_on_empty_reply: Whether a response with neither text nor tool calls ends
            the run as a failure. The lead's does — the usual cause is the output budget,
            and the error names the setting to raise; a role's simply ends with an empty
            summary.
        bogus_retry: Whether an answer quoting ids absent from the catalogue buys the
            model one more turn, with the correction appended, before it is accepted.
        provider: The provider this run's client talks to — for the usage rows, which
            bill per provider. None means the lead's configured one.
        model: The model, when the run was given one of its own (`HELIOAI_ROLE_MODELS`).
        deferred: Tools whose definitions are withheld from the model until it asks for
            them with `search_tools`, or calls one by name. Twenty-one definitions rode
            on every call of a lead turn; the ten formulary and catalogue ones are used
            in a minority of sessions and cost a third of that payload each time.
        final_answer: Whether the model may close the run with `final_answer(answer,
            claims)` — the answer plus the list of numbers it states, each with its
            source — instead of a plain message. Prose is still accepted; the claims are
            what lets a verdict compare numbers by name instead of by regex.
        search_budget: How many lookups (`search_tools_names`) a run may make before it
            has touched any data. A data_analyst spent all twelve of its turns on
            `search_parameters`, re-asking for ids it had been handed on the first call
            because they were not spelled the way it imagined; the run ended with nothing
            downloaded. Past the budget, one correction is appended telling it to use the
            ids it has. 0 disables — a parameter_hunter's whole job is to search.
        search_tools_names: The tools that count as lookups.
        data_tools_names: The tools whose first call ends the lookup phase.
    """

    name: str
    system_prompt: str
    tools: tuple[ToolDef, ...]
    max_turns: int
    tool_choice_first: str = "auto"
    allowed: frozenset[str] | None = None
    sandbox_no_network: bool = False
    sub_agent_ctx: Mapping[str, str] | None = None
    comment_replies: bool = False
    stream_replies: bool = False
    stop_on_empty_reply: bool = False
    bogus_retry: bool = True
    provider: str | None = None
    model: str | None = None
    deferred: frozenset[str] = frozenset()
    final_answer: bool = False
    search_budget: int = 0
    search_tools_names: frozenset[str] = frozenset({"search_parameters", "list_missions"})
    data_tools_names: frozenset[str] = frozenset(
        {"get_timeseries", "get_events_timeseries", "run_python", "run_recipe", "get_catalog"}
    )

    @property
    def event_extra(self) -> dict:
        """The keys added to every event of the run: a role's context, or nothing."""
        return {"sub_agent_ctx": dict(self.sub_agent_ctx)} if self.sub_agent_ctx else {}

event_extra property

event_extra: dict

The keys added to every event of the run: a role's context, or nothing.

helioai.runtime.context

Who is running, in which session, writing where — as an object, not an ambience.

Three contextvars (workspace._current_user, _current_session, _current_label) used to be set by each loop, by the MCP server and by nobody in a test, and read back from inside the tools, the datastore and the sandbox: an argument nobody passed and everybody depended on. A tool called from a test wrote under the default user; a sub-agent worked only because the lead had bound the label first; a save landed in a temporary directory when the caller forgot the session.

RunContext carries the same facts explicitly and is handed to the runner, which binds them for the duration of a run — bound() is the one place the three contextvars are set and reset. The tools that write receive their directories as trusted arguments (tool_exec.trusted_args); the contextvars remain the hot path for everything else, RunContext.current() reads them back, and the fallback logs when it is used.

RunContext dataclass

The facts a run writes under.

Attributes:

Name Type Description
user_id str

Owner of every path derived during the run.

session_id str

The conversation; a sub-agent runs under its parent's.

session_dir Path

The workspace directory — figures, scripts, npz, manifest, ledger. A sub-agent shares its lead's, which is how load_data() reaches what the lead downloaded.

label str | None

The human-readable directory name, when the session has one.

agent str

"lead" or the sub-agent role.

task_id str | None

The delegation's correlation id, for a sub-agent.

no_network bool

Whether run_python is denied a network namespace.

query str | None

The user's question for this turn, verbatim. It used to stop at the lead's loop: a sub-agent saw only the lead's brief, and nothing downstream could ask whether a result answered what the person typed. Carried on the context, it reaches every sub-agent through child() for free.

Source code in helioai/runtime/context.py
@dataclass(frozen=True)
class RunContext:
    """The facts a run writes under.

    Attributes:
        user_id: Owner of every path derived during the run.
        session_id: The conversation; a sub-agent runs under its parent's.
        session_dir: The workspace directory — figures, scripts, npz, manifest, ledger.
            A sub-agent shares its lead's, which is how `load_data()` reaches what the
            lead downloaded.
        label: The human-readable directory name, when the session has one.
        agent: `"lead"` or the sub-agent role.
        task_id: The delegation's correlation id, for a sub-agent.
        no_network: Whether `run_python` is denied a network namespace.
        query: The user's question for this turn, verbatim. It used to stop at the lead's
            loop: a sub-agent saw only the lead's brief, and nothing downstream could ask
            whether a result answered what the person typed. Carried on the context, it
            reaches every sub-agent through `child()` for free.
    """

    user_id: str
    session_id: str
    session_dir: Path
    label: str | None = None
    agent: str = "lead"
    task_id: str | None = None
    no_network: bool = False
    query: str | None = None

    @classmethod
    def for_session(
        cls, user_id: str, session_id: str, *, label: str | None = None, **extra: Any
    ) -> RunContext:
        """Build a context from ids, deriving the session directory the way
        `workspace.get_session_dir` does: by label when there is one, else by id.

        Args:
            user_id: Owner of the session.
            session_id: The conversation.
            label: Its workspace directory name, when already minted.
            **extra: `agent`, `task_id`, `no_network`, `query`.
        """
        return cls(
            user_id=user_id,
            session_id=session_id,
            session_dir=_ws.session_dir_for(user_id, session_id, label),
            label=label,
            **extra,
        )

    @classmethod
    def current(cls, **extra) -> RunContext | None:
        """The context the ambient contextvars describe, or None when no session is bound.

        The transition path: a caller that has not been handed a context yet can still
        obtain the one its caller bound.
        """
        session_id = _ws._current_session.get()
        if not session_id:
            return None
        return cls.for_session(
            _ws.current_user(), session_id, label=_ws._current_label.get(), **extra
        )

    def child(self, *, agent: str, task_id: str | None, no_network: bool = False) -> RunContext:
        """A sub-agent's context: the same user, session and directory, its own identity."""
        return replace(self, agent=agent, task_id=task_id, no_network=no_network)

    @property
    def data_dir(self) -> Path:
        """Where the session's downloads (npz + manifest) live."""
        from helioai.datastore import DATA_SUBDIR

        return self.session_dir / DATA_SUBDIR

    @property
    def catalogs_dir(self) -> Path:
        """Where the user's saved catalogues live: beside the workspaces, not in one."""
        return _ws.user_home(self.user_id) / "catalogs"

    @contextmanager
    def bound(self) -> Iterator[RunContext]:
        """Bind the workspace contextvars to this context for the duration of a block.

        Tokens are reset in reverse order, so a sub-agent bound inside its lead's run
        hands the lead's bindings back untouched.
        """
        tokens = [_ws.set_user(self.user_id), _ws.set_session(self.session_id)]
        resets = [_ws.reset_user, _ws.reset_session]
        if self.label:
            tokens.append(_ws.set_label(self.label))
            resets.append(_ws.reset_label)
        try:
            yield self
        finally:
            for reset, token in zip(reversed(resets), reversed(tokens), strict=True):
                reset(token)

data_dir property

data_dir: Path

Where the session's downloads (npz + manifest) live.

catalogs_dir property

catalogs_dir: Path

Where the user's saved catalogues live: beside the workspaces, not in one.

for_session classmethod

for_session(user_id: str, session_id: str, *, label: str | None = None, **extra: Any) -> RunContext

Build a context from ids, deriving the session directory the way workspace.get_session_dir does: by label when there is one, else by id.

Parameters:

Name Type Description Default
user_id str

Owner of the session.

required
session_id str

The conversation.

required
label str | None

Its workspace directory name, when already minted.

None
**extra Any

agent, task_id, no_network, query.

{}
Source code in helioai/runtime/context.py
@classmethod
def for_session(
    cls, user_id: str, session_id: str, *, label: str | None = None, **extra: Any
) -> RunContext:
    """Build a context from ids, deriving the session directory the way
    `workspace.get_session_dir` does: by label when there is one, else by id.

    Args:
        user_id: Owner of the session.
        session_id: The conversation.
        label: Its workspace directory name, when already minted.
        **extra: `agent`, `task_id`, `no_network`, `query`.
    """
    return cls(
        user_id=user_id,
        session_id=session_id,
        session_dir=_ws.session_dir_for(user_id, session_id, label),
        label=label,
        **extra,
    )

current classmethod

current(**extra) -> RunContext | None

The context the ambient contextvars describe, or None when no session is bound.

The transition path: a caller that has not been handed a context yet can still obtain the one its caller bound.

Source code in helioai/runtime/context.py
@classmethod
def current(cls, **extra) -> RunContext | None:
    """The context the ambient contextvars describe, or None when no session is bound.

    The transition path: a caller that has not been handed a context yet can still
    obtain the one its caller bound.
    """
    session_id = _ws._current_session.get()
    if not session_id:
        return None
    return cls.for_session(
        _ws.current_user(), session_id, label=_ws._current_label.get(), **extra
    )

child

child(*, agent: str, task_id: str | None, no_network: bool = False) -> RunContext

A sub-agent's context: the same user, session and directory, its own identity.

Source code in helioai/runtime/context.py
def child(self, *, agent: str, task_id: str | None, no_network: bool = False) -> RunContext:
    """A sub-agent's context: the same user, session and directory, its own identity."""
    return replace(self, agent=agent, task_id=task_id, no_network=no_network)

bound

bound() -> Iterator[RunContext]

Bind the workspace contextvars to this context for the duration of a block.

Tokens are reset in reverse order, so a sub-agent bound inside its lead's run hands the lead's bindings back untouched.

Source code in helioai/runtime/context.py
@contextmanager
def bound(self) -> Iterator[RunContext]:
    """Bind the workspace contextvars to this context for the duration of a block.

    Tokens are reset in reverse order, so a sub-agent bound inside its lead's run
    hands the lead's bindings back untouched.
    """
    tokens = [_ws.set_user(self.user_id), _ws.set_session(self.session_id)]
    resets = [_ws.reset_user, _ws.reset_session]
    if self.label:
        tokens.append(_ws.set_label(self.label))
        resets.append(_ws.reset_label)
    try:
        yield self
    finally:
        for reset, token in zip(reversed(resets), reversed(tokens), strict=True):
            reset(token)

helioai.runtime.validator

One verdict on a finished answer: its ids, its recipes, its figures, its numbers.

Four checks judged the lead's answer from four places — _flag_unknown_ids and _flag_recipe_bypass in tool_exec, provenance_check.check_reply for the numbers in the prose, vision.maybe_review for the figures — each with its own event and none aware of the others. validate runs them in one place and returns a Verdict the wrapper turns into the events the interfaces already render, plus one verdict event that carries the whole judgement.

What is new is the judgement of claims: when the model closed with final_answer, it named each number and where it came from. Those are compared to the provenance ledger by name, with a unit-aware tolerance — 2.5 states a recorded 2.519, 57.2 deg states 57.16 deg, 0.0025 pT states 2.5 nT — instead of being found again in the prose and attributed by the words around them. The regex check stays as the net under the prose: a number the model stated without claiming it is still judged.

Verdict dataclass

How an answer holds up, in every respect the runtime can check.

Attributes:

Name Type Description
matched list[dict]

Claims a ledger entry of the same name states (unit-aware).

contradicted list[dict]

Claims naming a recorded scalar that holds another value.

unsourced list[dict]

Claims nothing in the session computed — including those the model itself marked literature or asserted, which are never contradicted.

unknown_ids list[str]

Parameter ids the answer quotes that exist in no catalogue.

recipe_flags list[dict]

Exports that look like a calibrated recipe's output while the recipe was never loaded, or loaded and never called.

figure_reviews list[str]

The vision verdicts on the run's figures.

prose dict | None

The regex-based provenance report on the free text, or None when the session computed nothing or the text states no number.

Source code in helioai/runtime/validator.py
@dataclass
class Verdict:
    """How an answer holds up, in every respect the runtime can check.

    Attributes:
        matched: Claims a ledger entry of the same name states (unit-aware).
        contradicted: Claims naming a recorded scalar that holds another value.
        unsourced: Claims nothing in the session computed — including those the model
            itself marked `literature` or `asserted`, which are never contradicted.
        unknown_ids: Parameter ids the answer quotes that exist in no catalogue.
        recipe_flags: Exports that look like a calibrated recipe's output while the
            recipe was never loaded, or loaded and never called.
        figure_reviews: The vision verdicts on the run's figures.
        prose: The regex-based provenance report on the free text, or None when the
            session computed nothing or the text states no number.
    """

    matched: list[dict] = field(default_factory=list)
    contradicted: list[dict] = field(default_factory=list)
    unsourced: list[dict] = field(default_factory=list)
    unknown_ids: list[str] = field(default_factory=list)
    recipe_flags: list[dict] = field(default_factory=list)
    figure_reviews: list[str] = field(default_factory=list)
    prose: dict | None = None

    @property
    def has_claims(self) -> bool:
        """Whether the answer named its numbers at all."""
        return bool(self.matched or self.contradicted or self.unsourced)

    def as_event(self) -> dict:
        """The `verdict` event payload: counts up front, every detail behind them."""
        return {
            "matched": len(self.matched),
            "contradicted": len(self.contradicted),
            "unsourced": len(self.unsourced),
            "unknown_ids": list(self.unknown_ids),
            "recipe_flags": list(self.recipe_flags),
            "figure_reviews": list(self.figure_reviews),
            "claims": [
                {"status": status, **c}
                for status, group in (
                    ("contradicted", self.contradicted),
                    ("unsourced", self.unsourced),
                    ("matched", self.matched),
                )
                for c in group
            ],
        }

has_claims property

has_claims: bool

Whether the answer named its numbers at all.

as_event

as_event() -> dict

The verdict event payload: counts up front, every detail behind them.

Source code in helioai/runtime/validator.py
def as_event(self) -> dict:
    """The `verdict` event payload: counts up front, every detail behind them."""
    return {
        "matched": len(self.matched),
        "contradicted": len(self.contradicted),
        "unsourced": len(self.unsourced),
        "unknown_ids": list(self.unknown_ids),
        "recipe_flags": list(self.recipe_flags),
        "figure_reviews": list(self.figure_reviews),
        "claims": [
            {"status": status, **c}
            for status, group in (
                ("contradicted", self.contradicted),
                ("unsourced", self.unsourced),
                ("matched", self.matched),
            )
            for c in group
        ],
    }

validate

validate(text: str, claims: list[dict], *, history: list, artifacts: list[dict], session_dir: Path, figure_reviews: list[str] | None = None) -> tuple[str, Verdict]

Judge a finished answer.

Parameters:

Name Type Description Default
text str

The answer, as the model wrote it.

required
claims list[dict]

The numbers it named (final_answer), each {name, value, units, source}; empty for a prose answer.

required
history list

The turn's messages, read for evidence a recipe was loaded.

required
artifacts list[dict]

What the run exported, to tell a computed value from a quoted one.

required
session_dir Path

Whose ledger to check against.

required
figure_reviews list[str] | None

The vision verdicts collected during the run.

None

Returns:

Type Description
str

The text — annotated when a check appends a correction, exactly as the two

Verdict

checks always annotated it — and the verdict.

Source code in helioai/runtime/validator.py
def validate(
    text: str,
    claims: list[dict],
    *,
    history: list,
    artifacts: list[dict],
    session_dir: Path,
    figure_reviews: list[str] | None = None,
) -> tuple[str, Verdict]:
    """Judge a finished answer.

    Args:
        text: The answer, as the model wrote it.
        claims: The numbers it named (`final_answer`), each `{name, value, units,
            source}`; empty for a prose answer.
        history: The turn's messages, read for evidence a recipe was loaded.
        artifacts: What the run exported, to tell a computed value from a quoted one.
        session_dir: Whose ledger to check against.
        figure_reviews: The vision verdicts collected during the run.

    Returns:
        The text — annotated when a check appends a correction, exactly as the two
        checks always annotated it — and the verdict.
    """
    text, unknown_ids = _flag_unknown_ids(text)
    text, recipe_flags = _flag_recipe_bypass(text, history, artifacts)
    verdict = Verdict(
        unknown_ids=unknown_ids,
        recipe_flags=recipe_flags,
        figure_reviews=list(figure_reviews or []),
        prose=_prose_report(text, session_dir),
    )
    if claims:
        entries = provenance.read_ledger(session_dir).get("values") or []
        for claim in claims:
            status, detail = judge_claim(claim, entries)
            getattr(verdict, status).append(detail)
        verdict.prose = _without_claimed_numbers(verdict.prose, claims)
    return text, verdict

judge_claim

judge_claim(claim: dict, entries: list[dict]) -> tuple[str, dict]

Place one claim against the ledger: matched, contradicted or unsourced.

The entries considered are those whose name is the claim's source or name — the export it says it came from. Any run that produced the value sources it (a session that exported compression_ratio twice, 3.045 then 2.538, computed both), so the entries are read newest first and the first that states the value wins, within the prose checker's tolerance (RTOL) or the rounding the claim shows. A named scalar that states another value contradicts the claim; an entry holding several values cannot accuse, since a reply legitimately quotes one component of a vector; nor can a dimensioned scalar accuse a claim that gives no units — the claim may be another quantity of the same run, and on the first live run it was (a normal's components filed under the angle's export). A claim the model marked literature or asserted is never contradicted: it did not say the session computed it.

Parameters:

Name Type Description Default
claim dict

{name, value, units, source} as final_answer delivered it.

required
entries list[dict]

provenance.read_ledger(...)["values"].

required

Returns:

Type Description
tuple[str, dict]

The status and the detail dict shown to a reader.

Source code in helioai/runtime/validator.py
def judge_claim(claim: dict, entries: list[dict]) -> tuple[str, dict]:
    """Place one claim against the ledger: `matched`, `contradicted` or `unsourced`.

    The entries considered are those whose name is the claim's `source` or `name` — the
    export it says it came from. Any run that produced the value sources it (a session
    that exported `compression_ratio` twice, 3.045 then 2.538, computed both), so the
    entries are read newest first and the first that states the value wins, within the
    prose checker's tolerance (`RTOL`) or the rounding the claim shows. A named scalar
    that states another value contradicts the claim; an entry holding several values
    cannot accuse, since a reply legitimately quotes one component of a vector; nor can
    a dimensioned scalar accuse a claim that gives no units — the claim may be another
    quantity of the same run, and on the first live run it was (a normal's components
    filed under the angle's export). A claim the model marked `literature` or
    `asserted` is never contradicted: it did not say the session computed it.

    Args:
        claim: `{name, value, units, source}` as `final_answer` delivered it.
        entries: `provenance.read_ledger(...)["values"]`.

    Returns:
        The status and the detail dict shown to a reader.
    """
    detail = {
        "name": claim.get("name"),
        "value": claim.get("value"),
        "units": claim.get("units") or "",
        "source": claim.get("source") or "asserted",
    }
    value = _number(claim.get("value"))
    source = str(detail["source"])
    if value is None or source in _SESSION_SOURCES_ONLY:
        return "unsourced", detail
    names = {source, str(detail["name"] or "")}
    named = [e for e in reversed(entries) if e.get("name") in names]
    if not named:
        return "unsourced", detail
    comparable = False
    for entry in named:
        converted = _in_ledger_units(value, str(detail["units"]), str(entry.get("units") or ""))
        if converted is None:
            continue
        comparable = True
        if converted is _INCOMPATIBLE:
            continue
        if _states(entry, converted, RTOL) or _within_rounding(entry, converted):
            detail["ledger"] = _ledger_value(entry)
            detail["code_path"] = entry.get("code_path")
            return "matched", detail
    if not comparable:
        detail["note"] = "units could not be reconciled with the ledger's"
        return "unsourced", detail
    scalars = [e for e in named if _is_scalar(e)]
    if not scalars:
        return "unsourced", detail
    recorded = scalars[0]
    recorded_units = str(recorded.get("units") or "")
    if not str(detail["units"]).strip() and recorded_units:
        # The first live run named the shock normal's components `source: theta_bn` —
        # the run that printed them — with no units; the ledger's `theta_bn` is the
        # angle. "n_x stated -0.509, the session computed 54.85 deg" accuses nothing a
        # reader can act on: a claim without units may be another quantity of the same
        # run, and here it was. Only a claim that says what it measures can be wrong.
        detail["note"] = (
            f"no units given; the export named holds {_ledger_value(recorded)} {recorded_units}"
        )
        return "unsourced", detail
    detail["ledger"] = _ledger_value(recorded)
    detail["ledger_units"] = recorded_units
    detail["code_path"] = recorded.get("code_path")
    return "contradicted", detail

helioai.runtime.plan

The plan the model announced, kept as data, and how the run then followed it.

present_plan(title, steps) was a display: the loop forwarded it to the interfaces as a plan event and forgot it. What the model then did was left to the reader to compare with what it had said — three tool calls later nobody remembers step 2 named get_timeseries, and a plan that promised the theta_bn recipe and ran hand-written arithmetic instead looked exactly like one that was followed. Plan is the announced plan as data; adherence places the turn's tool calls against it and reports, in one plan_report event, which planned tools were used, which were not, and which tools were used without being planned. The report describes; it never blocks, corrects or retries.

Two things the first live run settled. A step's tool field is prose — the model wrote "search_parameters + get_timeseries" and "task (data_analyst)" — so it is read word by word against the tools the lead knows, not compared whole. And a planned tool the lead had a sub-agent run is a step done, not a deviation: the lead planned the analysis in its analyst's tools and then delegated every step, and the report said 0/4 used, unplanned: task. A sub-agent's calls therefore satisfy the plan; they are never unplanned — the second run planned in delegations alone (task (data_analyst)) and its analyst's six tools were all "unplanned" — because the task step covers whatever the role does, and the role's whitelist, not the lead's plan, governs it. What the lead does with its own hands is held to the plan, task excepted: delegating is how a step gets done. The scaffolding calls (the plan itself, the skills, search_tools, final_answer) count for nothing on either side.

Step dataclass

One announced step.

Attributes:

Name Type Description
description str

What the step does, as the model put it.

tool str | None

The tool it said it would use, or None when it named none.

Source code in helioai/runtime/plan.py
@dataclass(frozen=True)
class Step:
    """One announced step.

    Attributes:
        description: What the step does, as the model put it.
        tool: The tool it said it would use, or None when it named none.
    """

    description: str
    tool: str | None = None

Plan dataclass

A plan as present_plan delivered it, with the payload's looseness removed.

Attributes:

Name Type Description
title str

The plan's one-line title.

steps tuple[Step, ...]

The steps in order; a step the model wrote as a bare string is kept as a description without a tool.

Source code in helioai/runtime/plan.py
@dataclass(frozen=True)
class Plan:
    """A plan as `present_plan` delivered it, with the payload's looseness removed.

    Attributes:
        title: The plan's one-line title.
        steps: The steps in order; a step the model wrote as a bare string is kept as a
            description without a tool.
    """

    title: str
    steps: tuple[Step, ...] = ()

    @classmethod
    def from_payload(cls, data: dict) -> Plan:
        """Build a plan from a `present_plan` payload or a `plan` event's data.

        Args:
            data: `{"title": str, "steps": [{"description", "tool"} | str, ...]}`, any
                part of which the model may have left out.

        Returns:
            The plan, with malformed steps dropped rather than raised on.
        """
        steps: list[Step] = []
        for raw in data.get("steps") or []:
            if isinstance(raw, str):
                steps.append(Step(raw))
            elif isinstance(raw, dict):
                tool = str(raw.get("tool") or "").strip() or None
                steps.append(Step(str(raw.get("description") or ""), tool))
        return cls(str(data.get("title") or ""), tuple(steps))

    def tools(self, known: Collection[str] | None = None) -> list[str]:
        """The tools the plan names, once each, in the order of their first mention.

        Args:
            known: The tools the lead can call. A step's `tool` field is read word by
                word and only these words count; None keeps every word, for a caller
                with no registry at hand.

        Returns:
            Tool names, first mention first.
        """
        words = (w for s in self.steps if s.tool for w in _WORD.findall(s.tool))
        return list(dict.fromkeys(w for w in words if known is None or w in known))

from_payload classmethod

from_payload(data: dict) -> Plan

Build a plan from a present_plan payload or a plan event's data.

Parameters:

Name Type Description Default
data dict

{"title": str, "steps": [{"description", "tool"} | str, ...]}, any part of which the model may have left out.

required

Returns:

Type Description
Plan

The plan, with malformed steps dropped rather than raised on.

Source code in helioai/runtime/plan.py
@classmethod
def from_payload(cls, data: dict) -> Plan:
    """Build a plan from a `present_plan` payload or a `plan` event's data.

    Args:
        data: `{"title": str, "steps": [{"description", "tool"} | str, ...]}`, any
            part of which the model may have left out.

    Returns:
        The plan, with malformed steps dropped rather than raised on.
    """
    steps: list[Step] = []
    for raw in data.get("steps") or []:
        if isinstance(raw, str):
            steps.append(Step(raw))
        elif isinstance(raw, dict):
            tool = str(raw.get("tool") or "").strip() or None
            steps.append(Step(str(raw.get("description") or ""), tool))
    return cls(str(data.get("title") or ""), tuple(steps))

tools

tools(known: Collection[str] | None = None) -> list[str]

The tools the plan names, once each, in the order of their first mention.

Parameters:

Name Type Description Default
known Collection[str] | None

The tools the lead can call. A step's tool field is read word by word and only these words count; None keeps every word, for a caller with no registry at hand.

None

Returns:

Type Description
list[str]

Tool names, first mention first.

Source code in helioai/runtime/plan.py
def tools(self, known: Collection[str] | None = None) -> list[str]:
    """The tools the plan names, once each, in the order of their first mention.

    Args:
        known: The tools the lead can call. A step's `tool` field is read word by
            word and only these words count; None keeps every word, for a caller
            with no registry at hand.

    Returns:
        Tool names, first mention first.
    """
    words = (w for s in self.steps if s.tool for w in _WORD.findall(s.tool))
    return list(dict.fromkeys(w for w in words if known is None or w in known))

calls_made

calls_made(events: Iterable[dict]) -> tuple[list[str], list[str]]

The tools a turn called: the lead's own, and its sub-agents', each once in order.

Parameters:

Name Type Description Default
events Iterable[dict]

The turn's events; tool_call events carrying a sub_agent_ctx are a sub-agent's, the others the lead's. Scaffolding calls count for neither.

required

Returns:

Type Description
tuple[list[str], list[str]]

(own, delegated), tool names in the order of their first call.

Source code in helioai/runtime/plan.py
def calls_made(events: Iterable[dict]) -> tuple[list[str], list[str]]:
    """The tools a turn called: the lead's own, and its sub-agents', each once in order.

    Args:
        events: The turn's events; `tool_call` events carrying a `sub_agent_ctx` are a
            sub-agent's, the others the lead's. Scaffolding calls count for neither.

    Returns:
        `(own, delegated)`, tool names in the order of their first call.
    """
    own: dict[str, None] = {}
    delegated: dict[str, None] = {}
    for ev in events:
        if ev.get("event") != "tool_call":
            continue
        name = ev["data"].get("name")
        if not name or name in SCAFFOLDING:
            continue
        (delegated if "sub_agent_ctx" in ev["data"] else own).setdefault(name, None)
    return list(own), list(delegated)

delegations_made

delegations_made(events: Iterable[dict]) -> list[dict]

How each delegation of the turn ended: the role, its turns, whether it was capped.

A plan written in delegations alone ("task (data_analyst)") is always followed by delegating, so the tools say nothing; what a reader wants to know is whether the role finished. The run after the Runner extraction is the case: the librarian hit its four-turn cap, was re-delegated and finished in two — visible in the trace, invisible in a count of tools.

Parameters:

Name Type Description Default
events Iterable[dict]

The turn's events; the lead's own sub_agent_end re-emissions count (those without a sub_agent_ctx).

required

Returns:

Type Description
list[dict]

[{role, n_iterations, capped}] in order of completion.

Source code in helioai/runtime/plan.py
def delegations_made(events: Iterable[dict]) -> list[dict]:
    """How each delegation of the turn ended: the role, its turns, whether it was capped.

    A plan written in delegations alone ("task (data_analyst)") is always followed by
    delegating, so the tools say nothing; what a reader wants to know is whether the
    role finished. The run after the Runner extraction is the case: the librarian hit
    its four-turn cap, was re-delegated and finished in two — visible in the trace,
    invisible in a count of tools.

    Args:
        events: The turn's events; the lead's own `sub_agent_end` re-emissions count
            (those without a `sub_agent_ctx`).

    Returns:
        `[{role, n_iterations, capped}]` in order of completion.
    """
    return [
        {
            "role": ev["data"].get("role", ""),
            "n_iterations": int(ev["data"].get("n_iterations") or 0),
            "capped": bool(ev["data"].get("capped", False)),
        }
        for ev in events
        if ev.get("event") == "sub_agent_end" and "sub_agent_ctx" not in ev["data"]
    ]

adherence

adherence(plan: Plan, events: Iterable[dict], known: Collection[str] | None = None) -> dict

Compare what the run did with what the plan said.

Parameters:

Name Type Description Default
plan Plan

The plan the turn opened with.

required
events Iterable[dict]

The turn's events, as yielded.

required
known Collection[str] | None

The tools the lead can call, to read the plan's tool fields against.

None

Returns:

Type Description
dict

The plan_report payload — title, planned (the tools the plan named),

dict

executed (the tools the lead called itself), delegated (the tools its

dict

sub-agents called), delegations (each sub-agent run: role, turns, capped),

dict

unplanned_tools (the lead's own calls that were never planned, task

dict

excepted), missed_tools (planned, called by neither) and ratio, the share

dict

of planned tools that were called, or None when the plan named no tool and

dict

there is nothing to hold the run to.

Source code in helioai/runtime/plan.py
def adherence(plan: Plan, events: Iterable[dict], known: Collection[str] | None = None) -> dict:
    """Compare what the run did with what the plan said.

    Args:
        plan: The plan the turn opened with.
        events: The turn's events, as yielded.
        known: The tools the lead can call, to read the plan's `tool` fields against.

    Returns:
        The `plan_report` payload — `title`, `planned` (the tools the plan named),
        `executed` (the tools the lead called itself), `delegated` (the tools its
        sub-agents called), `delegations` (each sub-agent run: role, turns, capped),
        `unplanned_tools` (the lead's own calls that were never planned, `task`
        excepted), `missed_tools` (planned, called by neither) and `ratio`, the share
        of planned tools that were called, or None when the plan named no tool and
        there is nothing to hold the run to.
    """
    planned = plan.tools(known)
    own, delegated = calls_made(events)
    done = list(dict.fromkeys(own + delegated))
    followed = [t for t in planned if t in done]
    return {
        "title": plan.title,
        "planned": planned,
        "executed": own,
        "delegated": delegated,
        "delegations": delegations_made(events),
        "unplanned_tools": [t for t in own if t not in planned and t != DELEGATION],
        "missed_tools": [t for t in planned if t not in done],
        "ratio": round(len(followed) / len(planned), 2) if planned else None,
    }

Agent loop

helioai.core.agent_loop

The agent decision loop.

Given a user message and a session id, run the LLM in a tool-using loop: the LLM may emit tool calls, execute them via the ToolRegistry, feed the results back, and iterate until the LLM produces a final text reply (or we hit the safety cap).

Two consumption modes share the same generator core (stream_chat): - chat() → collects all events, returns a single ChatResult - stream_chat() → async generator, yields one event dict per step

Event kinds and their payloads are listed once, in core/events.py, and held to the emitters and the three renderers by tests/test_events_contract.py.

ChatResult dataclass

Final outcome of a non-streaming chat() call.

Source code in helioai/core/agent_loop.py
@dataclass
class ChatResult:
    """Final outcome of a non-streaming `chat()` call."""

    reply: str
    n_iterations: int
    artifacts: list[dict] = field(default_factory=list)
    events: list[dict] = field(default_factory=list)

build_lead_system_prompt

build_lead_system_prompt(restricted: bool, experiments: frozenset[str] | None = None) -> str

Return the lead agent system prompt.

Parameters:

Name Type Description Default
restricted bool

True (the public default) appends the scope guardrail, so the model refuses off-topic requests itself. False is reached only with a valid dev token and yields the base prompt.

required
experiments frozenset[str] | None

The experiments in force; settings.agent.experiments when None. deferred_tools adds the sentence that tells the model some tools appear on request. A prompt that describes a tool the model is not given, or the reverse, is the worst of both, so the text and the Policy are driven by the same set.

None

Returns:

Type Description
str

The full system prompt text.

Source code in helioai/core/agent_loop.py
def build_lead_system_prompt(restricted: bool, experiments: frozenset[str] | None = None) -> str:
    """Return the lead agent system prompt.

    Args:
        restricted: True (the public default) appends the scope guardrail, so the
            model refuses off-topic requests itself. False is reached only with a
            valid dev token and yields the base prompt.
        experiments: The experiments in force; `settings.agent.experiments` when None.
            `deferred_tools` adds the sentence that tells the model some tools appear
            on request. A prompt that describes a tool the model is not given, or the
            reverse, is the worst of both, so the text and the `Policy` are driven by
            the same set.

    Returns:
        The full system prompt text.
    """
    if experiments is None:
        experiments = settings.agent.experiments
    prompt = SYSTEM_PROMPT
    if "deferred_tools" in experiments:
        prompt = prompt.replace(
            _DEFERRED_TOOLS_ANCHOR, _DEFERRED_TOOLS_ANCHOR + DEFERRED_TOOLS_NOTE, 1
        )
    if restricted:
        return prompt + "\n\n" + SCOPE_GUARDRAIL
    return prompt

stream_chat async

stream_chat(llm_client: LLMClient, user_id: str, session_id: str, user_text: str, *, restricted: bool = True) -> AsyncIterator[dict]

Run one conversational turn of the agent and stream its progress as events.

This is the package's central API: every interface (CLI, web SSE, Jupyter, MCP) is a consumer of this generator. History is loaded from and persisted to the session store keyed by (user_id, session_id), so consecutive calls with the same ids continue the same conversation.

Parameters:

Name Type Description Default
llm_client LLMClient

Provider client from build_llm_client().

required
user_id str

Storage namespace — workspaces and profiles live under it.

required
session_id str

Conversation id; reuse it to continue, mint one to start fresh.

required
user_text str

The user's message for this turn.

required
restricted bool

True (default) appends the heliophysics scope guardrail; False (dev token) exposes the base prompt only.

True

Yields:

Type Description
AsyncIterator[dict]

Dicts with an "event" key and a "data" payload — every kind listed in

AsyncIterator[dict]

core/events.py; the turn opens with user and ends with done. Each one is

AsyncIterator[dict]

appended to the session's journal (SessionStore.append_event) before it is

AsyncIterator[dict]

yielded, so store.events() replays exactly what an interface was shown.

Example

llm = build_llm_client() async for ev in stream_chat(llm, "cli", "my-session", "IMF Bz at L1 today?"): ... if ev["event"] == "reply": ... print(ev["text"], end="")

Source code in helioai/core/agent_loop.py
async def stream_chat(
    llm_client: LLMClient,
    user_id: str,
    session_id: str,
    user_text: str,
    *,
    restricted: bool = True,
) -> AsyncIterator[dict]:
    """Run one conversational turn of the agent and stream its progress as events.

    This is the package's central API: every interface (CLI, web SSE, Jupyter,
    MCP) is a consumer of this generator. History is loaded from and persisted
    to the session store keyed by (user_id, session_id), so consecutive calls
    with the same ids continue the same conversation.

    Args:
        llm_client: Provider client from `build_llm_client()`.
        user_id: Storage namespace — workspaces and profiles live under it.
        session_id: Conversation id; reuse it to continue, mint one to start fresh.
        user_text: The user's message for this turn.
        restricted: True (default) appends the heliophysics scope guardrail;
            False (dev token) exposes the base prompt only.

    Yields:
        Dicts with an `"event"` key and a `"data"` payload — every kind listed in
        `core/events.py`; the turn opens with `user` and ends with `done`. Each one is
        appended to the session's journal (`SessionStore.append_event`) before it is
        yielded, so `store.events()` replays exactly what an interface was shown.

    Example:
        >>> llm = build_llm_client()
        >>> async for ev in stream_chat(llm, "cli", "my-session", "IMF Bz at L1 today?"):
        ...     if ev["event"] == "reply":
        ...         print(ev["text"], end="")
    """
    # One turn at a time per session. The store hands every caller the same history
    # list; two turns interleaving their appends persisted a corrupted transcript — two
    # browser tabs on one session were enough. The lock is held for the whole turn, so
    # a second request waits for the first to finish (the web layer answers 409 instead
    # of waiting, see app.chat_stream). Closing the inner generator explicitly is what
    # runs its cleanup now rather than at garbage collection.
    async with store.turn_lock(user_id, session_id):
        opening = make("user", text=user_text)
        store.append_event(user_id, session_id, opening)
        yield opening
        turn = _stream_turn(llm_client, user_id, session_id, user_text, restricted=restricted)
        try:
            async for ev in turn:
                if ev["event"] not in NOT_JOURNALED:
                    store.append_event(user_id, session_id, ev)
                yield ev
        finally:
            await turn.aclose()

chat async

chat(llm_client: LLMClient, user_id: str, session_id: str, user_text: str, *, restricted: bool = True) -> ChatResult

Run one agent turn to completion and return the final result.

Non-streaming wrapper over stream_chat — same arguments, same session semantics — for callers that want the answer, not the progress feed (Jupyter magic, scripts, tests).

Parameters:

Name Type Description Default
llm_client LLMClient

Provider client from build_llm_client.

required
user_id str

Storage owner; decides which tree under data/users/ is used.

required
session_id str

Conversation to append to. An unknown id starts a new one.

required
user_text str

The question.

required
restricted bool

Whether the scope guardrail is in force. See build_lead_system_prompt.

True

Returns:

Type Description
ChatResult

ChatResult with reply (final text), n_iterations (LLM turns used),

ChatResult

artifacts (figures, parameter cards, code) and events (full trace).

Example

result = await chat(build_llm_client(), "cli", "my-session", ... "Plasma beta for B=5 nT, n=10 cm^-3, T=20 eV?") print(result.reply) # final assistant text (model-dependent) len(result.artifacts) # figures / parameter cards / code produced

Source code in helioai/core/agent_loop.py
async def chat(
    llm_client: LLMClient,
    user_id: str,
    session_id: str,
    user_text: str,
    *,
    restricted: bool = True,
) -> ChatResult:
    """Run one agent turn to completion and return the final result.

    Non-streaming wrapper over `stream_chat` — same arguments, same session
    semantics — for callers that want the answer, not the progress feed
    (Jupyter magic, scripts, tests).

    Args:
        llm_client: Provider client from `build_llm_client`.
        user_id: Storage owner; decides which tree under `data/users/` is used.
        session_id: Conversation to append to. An unknown id starts a new one.
        user_text: The question.
        restricted: Whether the scope guardrail is in force. See
            `build_lead_system_prompt`.

    Returns:
        ChatResult with `reply` (final text), `n_iterations` (LLM turns used),
        `artifacts` (figures, parameter cards, code) and `events` (full trace).

    Example:
        >>> result = await chat(build_llm_client(), "cli", "my-session",
        ...                     "Plasma beta for B=5 nT, n=10 cm^-3, T=20 eV?")
        >>> print(result.reply)        # final assistant text (model-dependent)
        >>> len(result.artifacts)      # figures / parameter cards / code produced
    """
    artifacts: list[dict] = []
    events: list[dict] = []
    reply = ""
    n_iters = 0
    error_msg: str | None = None

    async for ev in stream_chat(llm_client, user_id, session_id, user_text, restricted=restricted):
        name, data = ev["event"], ev["data"]
        if name == "reply":
            reply = data.get("text", "")
        elif name == "done":
            n_iters = data.get("n_iterations", 0)
        elif name == "artifact":
            artifacts.append(data)
        elif name == "tool_call":
            events.append(
                {
                    "turn": data["turn"],
                    "type": "tool_call",
                    "tool": data["name"],
                    "arguments": data.get("arguments", {}),
                }
            )
        elif name == "tool_result":
            events.append(
                {
                    "turn": data["turn"],
                    "type": "tool_result",
                    "tool": data["name"],
                    "summary": data.get("summary", ""),
                }
            )
        elif name == "error":
            error_msg = data.get("message", "unknown agent error")

    if error_msg is not None:
        raise RuntimeError(error_msg)
    return ChatResult(reply=reply, n_iterations=n_iters, artifacts=artifacts, events=events)

Sub-agents

helioai.core.sub_agents

Sub-agents for delegating focused heliophysics subtasks.

Each role declares a tool whitelist, a system addon, and an optional set of skills auto-loaded into the sub's system prompt. Sub-agents run in isolation with a fresh context — the lead's history is invisible to them.

The events a sub-agent yields are the non-lead-only kinds of core/events.py.

SubAgentRole dataclass

A specialised agent: its prompt, its tool whitelist and its turn budget.

allowed_tools is enforced, not advisory — a role calling outside its set gets an error naming what it may use, and the tool is never dispatched.

search_budget is the role's lookup allowance before it must touch data; it only takes effect under the search_budget experiment, and is 0 (no limit) otherwise.

Source code in helioai/core/sub_agents.py
@dataclass(frozen=True)
class SubAgentRole:
    """A specialised agent: its prompt, its tool whitelist and its turn budget.

    `allowed_tools` is enforced, not advisory — a role calling outside its set
    gets an error naming what it may use, and the tool is never dispatched.

    `search_budget` is the role's lookup allowance before it must touch data; it only
    takes effect under the `search_budget` experiment, and is 0 (no limit) otherwise.
    """

    name: str
    description: str
    system_addon: str
    allowed_tools: tuple[str, ...]
    max_turns: int = 5
    auto_load_skills: tuple[str, ...] = ()
    sandbox_no_network: bool = False
    search_budget: int = 0

task_tool_def

task_tool_def() -> ToolDef

Build the task tool definition offered to the lead agent.

Deliberately not registered in the ToolRegistry: the agent loop intercepts task and spawns a sub-agent instead of dispatching a function.

Returns:

Type Description
ToolDef

A ToolDef whose agent_role enum lists every available role.

Source code in helioai/core/sub_agents.py
def task_tool_def() -> ToolDef:
    """Build the `task` tool definition offered to the lead agent.

    Deliberately not registered in the ToolRegistry: the agent loop intercepts
    `task` and spawns a sub-agent instead of dispatching a function.

    Returns:
        A ToolDef whose `agent_role` enum lists every available role.
    """
    role_lines = "\n".join(f"  - `{r.name}`: {r.description}" for r in AGENT_ROLES.values())
    return ToolDef(
        name=TASK_TOOL_NAME,
        # Routing rules live in the lead's system prompt, not here. Stating them in both
        # places let them drift into contradiction — this text used to call
        # `parameter_hunter` "required" for unknown ids while the prompt said never to run
        # it first, and the model resolved the conflict differently from one run to the
        # next: three delegations, three, then none, on the same notebook.
        description=(
            "Spawn a specialist sub-agent for ONE focused subtask. "
            "The sub runs in isolation (empty context) — pre-resolve every fact "
            "(param ids, ISO times, missions) inside `description`.\n\nRoles:\n" + role_lines
        ),
        parameters={
            "type": "object",
            "properties": {
                "description": {
                    "type": "string",
                    "description": "Self-contained task description (1-3 sentences) with all needed facts.",
                },
                "agent_role": {
                    "type": "string",
                    "enum": sorted(AGENT_ROLES.keys()),
                    "description": "Sub-agent role to spawn.",
                },
            },
            "required": ["description", "agent_role"],
        },
    )

stream_subagent async

stream_subagent(role: str, description: str, *, parent_session_id: str, user_id: str, llm_client: LLMClient, task_id: str | None = None, context: RunContext | None = None) -> AsyncIterator[dict]

Async generator that runs a sub-agent and yields progress events.

Yields the kinds of core/events.py that are not lead-only (tool_call, tool_result, skill_loaded, artifact, figure_review, invalid_ids, recipe_bypassed), each enriched with sub_agent_ctx={role, task_id}, then a final sub_agent_end carrying findings/summary/artifacts/n_iterations/error.

summary is the sub-agent's whole deliverable and is emitted in full: it becomes the lead's tool result, so anything cut here is a measurement the lead can no longer report and will be tempted to invent. Callers that display it truncate on their own side. Reaching the turn cap is reported as an error, not as a result, for the same reason — a lead handed a capped run has nothing to summarise.

Parameters:

Name Type Description Default
role str

One of the configured roles. Its tool whitelist and turn cap are enforced, not advisory.

required
description str

The task, in the lead's own words.

required
parent_session_id str

The lead's session. The sub-agent writes into that same workspace, which is how load_data() reaches what the lead already downloaded.

required
user_id str

Storage owner, inherited from the lead.

required
llm_client LLMClient

Provider client, shared with the lead.

required
task_id str | None

Correlation id echoed in every event, so a caller running several sub-agents can tell their streams apart.

None
context RunContext | None

The lead's run context; the sub-agent runs under a child of it, in the same session directory. None derives one from the ids and the bound label, for callers that predate contexts.

None

Yields:

Type Description
AsyncIterator[dict]

Progress events, then a final sub_agent_end.

Source code in helioai/core/sub_agents.py
async def stream_subagent(
    role: str,
    description: str,
    *,
    parent_session_id: str,
    user_id: str,
    llm_client: LLMClient,
    task_id: str | None = None,
    context: RunContext | None = None,
) -> AsyncIterator[dict]:
    """Async generator that runs a sub-agent and yields progress events.

    Yields the kinds of `core/events.py` that are not lead-only (tool_call, tool_result,
    skill_loaded, artifact, figure_review, invalid_ids, recipe_bypassed), each enriched
    with sub_agent_ctx={role, task_id}, then a final sub_agent_end carrying
    findings/summary/artifacts/n_iterations/error.

    `summary` is the sub-agent's whole deliverable and is emitted in full: it becomes
    the lead's tool result, so anything cut here is a measurement the lead can no
    longer report and will be tempted to invent. Callers that display it truncate on
    their own side. Reaching the turn cap is reported as an `error`, not as a result,
    for the same reason — a lead handed a capped run has nothing to summarise.

    Args:
        role: One of the configured roles. Its tool whitelist and turn cap are
            enforced, not advisory.
        description: The task, in the lead's own words.
        parent_session_id: The lead's session. The sub-agent writes into that
            same workspace, which is how `load_data()` reaches what the lead
            already downloaded.
        user_id: Storage owner, inherited from the lead.
        llm_client: Provider client, shared with the lead.
        task_id: Correlation id echoed in every event, so a caller running
            several sub-agents can tell their streams apart.
        context: The lead's run context; the sub-agent runs under a child of it, in the
            same session directory. None derives one from the ids and the bound label,
            for callers that predate contexts.

    Yields:
        Progress events, then a final `sub_agent_end`.
    """
    import helioai.workspace as _ws

    if task_id is None:
        task_id = uuid.uuid4().hex[:8]
    ctx = {"role": role, "task_id": task_id}
    if context is None:
        context = RunContext.for_session(
            user_id or _ws.current_user(), parent_session_id, label=_ws._current_label.get()
        )

    if role not in AGENT_ROLES:
        known = ", ".join(sorted(AGENT_ROLES))
        yield make(
            "sub_agent_end",
            task_id=task_id,
            role=role,
            summary="",
            n_iterations=0,
            error=f"unknown agent_role {role!r}. Known: {known}",
            artifacts=[],
            findings={},
            usage={},
            capped=False,
        )
        return

    role_cfg = AGENT_ROLES[role]

    structlog.contextvars.bind_contextvars(
        parent_session_id=parent_session_id,
        sub_role=role,
        sub_task_id=task_id,
    )

    runner: Runner | None = None
    own_client: LLMClient | None = None
    provider = settings.llm.provider
    try:
        system_prompt, skills_loaded = _build_system_prompt(role_cfg)
        role_model = settings.agent.role_models.get(role)
        if role_model is not None:
            provider, model_name = role_model
            own_client = build_llm_client(provider, model=model_name)
            llm_client = own_client
            log.info("subagent_own_model", role=role, provider=provider, model=model_name)
        for skill_name in skills_loaded:
            yield make("skill_loaded", name=skill_name, sub_agent_ctx=ctx)

        allowed = frozenset(role_cfg.allowed_tools)
        tools = tuple(registry.list_tool_defs(only=set(allowed)))

        log.info(
            "subagent_start",
            role=role,
            description=description[:200],
            allowed_tools=sorted(allowed),
            max_turns=role_cfg.max_turns,
        )

        history: list[Message] = [
            Message(role="user", content=_with_inventory(description, context.session_dir))
        ]
        policy = Policy(
            name=role,
            system_prompt=system_prompt,
            tools=tools,
            max_turns=role_cfg.max_turns,
            tool_choice_first="required",
            allowed=allowed,
            sandbox_no_network=role_cfg.sandbox_no_network,
            sub_agent_ctx=ctx,
            provider=provider,
            model=role_model[1] if role_model else None,
            search_budget=(
                role_cfg.search_budget if "search_budget" in settings.agent.experiments else 0
            ),
        )
        runner = Runner(
            policy,
            llm_client,
            registry=registry,
            ctx=context.child(agent=role, task_id=task_id, no_network=role_cfg.sandbox_no_network),
            intercept=functools.partial(_deny_outside_whitelist, role=role, allowed=allowed),
        )
        t0 = time.monotonic()

        end: RunEnd | None = None
        async with aclosing(runner.run(history)) as run:
            async for item in run:
                if isinstance(item, RunEnd):
                    end = item
                    break
                yield item
        assert end is not None

        capped = end.capped
        final_text = (
            f"(sub-agent {role!r} reached its {role_cfg.max_turns}-turn cap)"
            if capped
            else end.final_text or ""
        )
        log.info(
            "subagent_end",
            role=role,
            n_iterations=end.turns,
            duration_ms=int((time.monotonic() - t0) * 1000),
            capped=capped,
            n_artifacts=len(end.artifacts),
        )

        final_text, bogus, bypassed_recipes = check_answer(final_text, history, end.artifacts)
        if bogus:
            log.warning("subagent_invented_ids", role=role, ids=bogus)
            yield make("invalid_ids", ids=bogus, sub_agent_ctx=ctx)

        if bypassed_recipes:
            log.warning("subagent_recipe_bypassed", role=role, recipes=bypassed_recipes)
            yield make("recipe_bypassed", recipes=bypassed_recipes, sub_agent_ctx=ctx)

        yield make(
            "sub_agent_end",
            task_id=task_id,
            role=role,
            findings=_findings(end.artifacts),
            summary=final_text,
            n_iterations=end.turns,
            error=final_text if capped else None,
            artifacts=end.artifacts,
            usage={**end.usage, "provider": provider},
            capped=capped,
        )

    except Exception as e:
        log.exception("subagent_error", role=role, task_id=task_id)
        artifacts = runner.artifacts if runner is not None else []
        yield make(
            "sub_agent_end",
            task_id=task_id,
            role=role,
            findings=_findings(artifacts),
            summary="",
            n_iterations=runner.turns if runner is not None else 0,
            error=str(e),
            artifacts=[],
            usage={**(runner.usage if runner is not None else _zero_usage()), "provider": provider},
            capped=False,
        )

    finally:
        if own_client is not None:
            await own_client.aclose()
        structlog.contextvars.unbind_contextvars("parent_session_id", "sub_role", "sub_task_id")

Tool execution

The tool-call mechanics the runner uses.

helioai.core.tool_exec

Shared tool-execution helpers for the agent loops.

Both the lead loop (agent_loop.stream_chat) and the sub-agent loop (sub_agents.stream_subagent) run the same tool-call mechanics: inject the sandbox run dir for run_python, summarise the result, detect skill loads, and extract renderable artifacts. Keeping that logic here — imported by both loops — prevents the two copies from drifting apart (a real bug source, see the session-13 _extract_artifact list/dict regression).

This module imports neither agent_loop nor sub_agents, so there is no cycle.

compact_history

compact_history(messages: list, keep_full: int = 2) -> list

Return a copy of messages where tool-result messages older than the last keep_full are summarized. The most recent results stay verbatim (the next LLM call usually needs them); older ones — already consumed — are trimmed so context does not grow unbounded over a long session. Persisted history is left untouched; only the per-call payload shrinks.

A loaded recipe is never summarized. It is the one result the model keeps writing code against for the rest of the task, and two run_python calls later it had been cut to its first line: in two live runs the analyst reloaded it, and once rewrote the formula from memory rather than call the function it could no longer see. A recipe is a few kilobytes; keeping it costs less than the extra turn.

The cap per stale result is 1 500 characters, up from 300. Measured on the live runs of 00_quickstart: every tool result of a four-turn session together is 14 kB, less than the system prompt and tool schemas re-sent on every call, and the numbers are exempt from the cap anyway. What 300 bought was a lead that could not recall the previous cell's method and an analyst that had lost its own exports.

ponytail: fixed window N=2; widen keep_full if a case regresses on stale results.

Parameters:

Name Type Description Default
messages list

The per-call payload, newest last.

required
keep_full int

How many of the most recent tool results to leave verbatim. Zero summarises every one of them, which is the degraded retry a context-length failure would want.

2

Returns:

Type Description
list

A new list. messages is not modified, because it is also what gets

list

persisted.

Source code in helioai/core/tool_exec.py
def compact_history(messages: list, keep_full: int = 2) -> list:
    """Return a copy of `messages` where tool-result messages older than the last
    `keep_full` are summarized. The most recent results stay verbatim (the next LLM
    call usually needs them); older ones — already consumed — are trimmed so context
    does not grow unbounded over a long session. Persisted history is left untouched;
    only the per-call payload shrinks.

    A loaded recipe is never summarized. It is the one result the model keeps writing
    code against for the rest of the task, and two `run_python` calls later it had been
    cut to its first line: in two live runs the analyst reloaded it, and once rewrote the
    formula from memory rather than call the function it could no longer see. A recipe
    is a few kilobytes; keeping it costs less than the extra turn.

    The cap per stale result is 1 500 characters, up from 300. Measured on the live
    runs of 00_quickstart: every tool result of a four-turn session together is 14 kB,
    less than the system prompt and tool schemas re-sent on every call, and the numbers
    are exempt from the cap anyway. What 300 bought was a lead that could not recall
    the previous cell's method and an analyst that had lost its own exports.

    ponytail: fixed window N=2; widen keep_full if a case regresses on stale results.

    Args:
        messages: The per-call payload, newest last.
        keep_full: How many of the most recent tool results to leave verbatim.
            Zero summarises every one of them, which is the degraded retry a
            context-length failure would want.

    Returns:
        A new list. `messages` is not modified, because it is also what gets
        persisted.
    """
    tool_idx = [i for i, m in enumerate(messages) if getattr(m, "role", None) == "tool"]
    if len(tool_idx) <= keep_full:
        return messages
    stale = set(tool_idx[:-keep_full])
    return [
        replace(m, content=_summarize_tool_result(m.content, max_chars=_STALE_RESULT_CHARS))
        if i in stale and m.content and not _is_recipe_source(m)
        else m
        for i, m in enumerate(messages)
    ]

trusted_args

trusted_args(name: str, ctx: RunContext | None = None, *, no_network: bool = False) -> dict

The framework-injected arguments of the tools that write to disk.

Passed via call_tool(..., trusted=...), so they bypass the private-argument guard that rejects model- or MCP-supplied _* overrides. run_python and run_recipe get their workspace, run index and network flag; get_timeseries and get_events_timeseries the session's data directory; save_catalog the user's catalogue directory. Every other tool gets nothing. Reading them off the context rather than off the workspace contextvars is what lets a tool called from a test, or over MCP, write where its caller said and nowhere else.

Parameters:

Name Type Description Default
name str

Tool about to be called.

required
ctx RunContext | None

The run's context. None falls back to the bound contextvars, the way the loops resolved the session before contexts existed.

None
no_network bool

Deny the sandbox a network namespace; ctx.no_network also does.

False

Returns:

Type Description
dict

The trusted argument dict, empty for a tool that writes nothing.

Source code in helioai/core/tool_exec.py
def trusted_args(name: str, ctx: RunContext | None = None, *, no_network: bool = False) -> dict:
    """The framework-injected arguments of the tools that write to disk.

    Passed via `call_tool(..., trusted=...)`, so they bypass the private-argument guard
    that rejects model- or MCP-supplied `_*` overrides. `run_python` and `run_recipe` get
    their workspace, run index and network flag; `get_timeseries` and
    `get_events_timeseries` the session's data directory; `save_catalog` the user's
    catalogue directory. Every other tool gets nothing. Reading them off the context
    rather than off the workspace contextvars is what lets a tool called from a test, or
    over MCP, write where its caller said and nowhere else.

    Args:
        name: Tool about to be called.
        ctx: The run's context. None falls back to the bound contextvars, the way the
            loops resolved the session before contexts existed.
        no_network: Deny the sandbox a network namespace; `ctx.no_network` also does.

    Returns:
        The trusted argument dict, empty for a tool that writes nothing.
    """
    if name not in _WRITING_TOOLS:
        return {}
    if ctx is None:
        from helioai.runtime.context import RunContext

        ctx = RunContext.current()
    if name == "save_catalog":
        return {"_catalogs_dir": str(ctx.catalogs_dir)} if ctx else {}
    if name in ("get_timeseries", "get_events_timeseries"):
        return {"_data_dir": str(ctx.data_dir)} if ctx else {}
    import helioai.workspace as _ws

    sdir = ctx.session_dir if ctx else _ws.get_session_dir()
    args = {"_plot_dir": str(sdir), "_run_idx": _ws.get_next_run_idx(sdir)}
    if no_network or (ctx is not None and ctx.no_network):
        args["_no_net"] = True
    return args

start_tool_calls

start_tool_calls(tool_calls: list[ToolCall] | None, *, allowed: set[str] | frozenset[str] | None = None, registry: ToolRegistry | None = None, ctx: RunContext | None = None) -> dict[str, asyncio.Task]

Start every parallel-safe registry call of a turn at once, keyed by call id.

The prompt asks the model to batch its downloads in one turn, and the loops then ran them one after the other: a data_analyst's first turn — three or four get_timeseries — took the sum of their durations. Started here, they overlap; the caller still awaits each result in the model's order, so every event and every tool message keeps the order it had when the calls were sequential.

Skipped, and left to the caller's sequential path: run_python (see _SEQUENTIAL_TOOLS), anything not in the registry (the task tool, the internal tools), and — for a sub-agent — anything outside its whitelist, which the caller refuses without dispatching.

Parameters:

Name Type Description Default
tool_calls list[ToolCall] | None

The assistant's tool calls for this turn.

required
allowed set[str] | frozenset[str] | None

A sub-agent's whitelist; None for the lead, who may call anything.

None
registry ToolRegistry | None

Where the calls are dispatched; the process-wide registry by default.

None
ctx RunContext | None

The run's context, from which the data tools get their directory as a trusted argument — the same trusted_args the sequential path passes.

None

Returns:

Type Description
dict[str, Task]

{tool_call.id: task} for the calls that were started; each task resolves to

dict[str, Task]

a ToolResult. registry.call_tool never raises — a failure is a failed

dict[str, Task]

result — so awaiting a task is safe.

Source code in helioai/core/tool_exec.py
def start_tool_calls(
    tool_calls: list[ToolCall] | None,
    *,
    allowed: set[str] | frozenset[str] | None = None,
    registry: ToolRegistry | None = None,
    ctx: RunContext | None = None,
) -> dict[str, asyncio.Task]:
    """Start every parallel-safe registry call of a turn at once, keyed by call id.

    The prompt asks the model to batch its downloads in one turn, and the loops then
    ran them one after the other: a data_analyst's first turn — three or four
    `get_timeseries` — took the sum of their durations. Started here, they overlap;
    the caller still awaits each result in the model's order, so every event and
    every `tool` message keeps the order it had when the calls were sequential.

    Skipped, and left to the caller's sequential path: `run_python` (see
    `_SEQUENTIAL_TOOLS`), anything not in the registry (the `task` tool, the internal
    tools), and — for a sub-agent — anything outside its whitelist, which the caller
    refuses without dispatching.

    Args:
        tool_calls: The assistant's tool calls for this turn.
        allowed: A sub-agent's whitelist; None for the lead, who may call anything.
        registry: Where the calls are dispatched; the process-wide registry by default.
        ctx: The run's context, from which the data tools get their directory as a
            trusted argument — the same `trusted_args` the sequential path passes.

    Returns:
        `{tool_call.id: task}` for the calls that were started; each task resolves to
        a `ToolResult`. `registry.call_tool` never raises — a failure is a failed
        result — so awaiting a task is safe.
    """
    if registry is None:
        from helioai.tools.registry import registry as _default

        registry = _default

    started: dict[str, asyncio.Task] = {}
    for tc in tool_calls or []:
        if tc.name in _SEQUENTIAL_TOOLS or tc.name not in registry:
            continue
        if allowed is not None and tc.name not in allowed:
            continue
        started[tc.id] = asyncio.create_task(
            registry.call_tool(tc.name, tc.arguments, trusted=trusted_args(tc.name, ctx))
        )
    return started

cancel_pending

cancel_pending(started: dict[str, Task]) -> None

Cancel the calls a turn started and never awaited — a cancelled turn must not leave downloads running for nobody.

Source code in helioai/core/tool_exec.py
def cancel_pending(started: dict[str, asyncio.Task]) -> None:
    """Cancel the calls a turn started and never awaited — a cancelled turn must not
    leave downloads running for nobody."""
    for task in started.values():
        if not task.done():
            task.cancel()

emit_post_tool_events

emit_post_tool_events(name: str, result: ToolResult, *, tool_result_extra: dict | None = None, common_extra: dict | None = None) -> Iterator[dict]

Yield the events that follow a completed tool call.

Order is tool_result → (skill_loaded if load_skill) → artifact(s), matching what both loops emitted before this was factored out.

Parameters:

Name Type Description Default
name str

The tool that just ran.

required
result ToolResult

Its result; the payload is read for artifacts, the model's text for the summary the event carries.

required
tool_result_extra dict | None

Merged into the tool_result event data, e.g. {turn}.

None
common_extra dict | None

Merged into skill_loaded and artifact event data, e.g. {sub_agent_ctx} when a sub-agent is the caller.

None

Yields:

Type Description
dict

The events, in the order both loops emitted them before this was

dict

factored out: tool_result, then skill_loaded, then artifacts.

Source code in helioai/core/tool_exec.py
def emit_post_tool_events(
    name: str,
    result: ToolResult,
    *,
    tool_result_extra: dict | None = None,
    common_extra: dict | None = None,
) -> Iterator[dict]:
    """Yield the events that follow a completed tool call.

    Order is `tool_result` → (`skill_loaded` if load_skill) → `artifact`(s),
    matching what both loops emitted before this was factored out.

    Args:
        name: The tool that just ran.
        result: Its result; the payload is read for artifacts, the model's text for
            the summary the event carries.
        tool_result_extra: Merged into the `tool_result` event data, e.g. `{turn}`.
        common_extra: Merged into `skill_loaded` and `artifact` event data, e.g.
            `{sub_agent_ctx}` when a sub-agent is the caller.

    Yields:
        The events, in the order both loops emitted them before this was
        factored out: `tool_result`, then `skill_loaded`, then artifacts.
    """
    tool_result_extra = tool_result_extra or {}
    common_extra = common_extra or {}
    text = result.for_llm()
    payload = result.payload

    # `summary` is written for the model and stays untouched. `display` is the same
    # event told to a person, computed here so the CLI, the Jupyter magic and the
    # browser show identical words without three copies of the logic.
    yield make(
        "tool_result",
        name=name,
        summary=_summarize_tool_result(text),
        display=describe_tool_result(name, text),
        **tool_result_extra,
    )

    if name == "load_skill" and isinstance(payload, dict):
        if payload.get("body") and not payload.get("error"):
            yield make("skill_loaded", name=payload.get("name", ""), **common_extra)

    for art in _extract_artifact(name, payload):
        if art.get("kind") == "exports":
            ctx = common_extra.get("sub_agent_ctx") or {}
            provenance.record(
                art.get("values") or {},
                code_path=_code_path(payload),
                agent=ctx.get("role") or "lead",
                task_id=ctx.get("task_id"),
                turn=tool_result_extra.get("turn"),
            )
        yield make("artifact", **art, **common_extra)

unknown_id_correction

unknown_id_correction(bogus: list[str]) -> str

The correction handed back when an answer quotes ids that are not in the catalogue.

Shared between the retry that buys the model another turn and the annotation left on an answer that has run out of turns, so both say exactly the same thing.

Parameters:

Name Type Description Default
bogus list[str]

The ids that are absent from the index.

required

Returns:

Type Description
str

The correction text, without the answer it refers to.

Source code in helioai/core/tool_exec.py
def unknown_id_correction(bogus: list[str]) -> str:
    """The correction handed back when an answer quotes ids that are not in the catalogue.

    Shared between the retry that buys the model another turn and the annotation left
    on an answer that has run out of turns, so both say exactly the same thing.

    Args:
        bogus: The ids that are absent from the index.

    Returns:
        The correction text, without the answer it refers to.
    """
    listed = "\n".join(f"  - {i}" for i in bogus)
    return (
        f"⚠️ AUTOMATED CORRECTION — the following ids are NOT in the catalogue and "
        f"must not be used:\n{listed}\n"
        f"They were not returned by any search. Call `search_parameters` again and "
        f"copy the ids from its output verbatim."
    )

recipe_available

recipe_available(tool_name: str, arguments: dict | None, result: ToolResult, history: list) -> ToolResult

Tell the model, on the run_python result it is about to read, that a shipped recipe covers what its code just computed.

_flag_recipe_bypass reads the same signals at the end of the turn and says so in the answer — after the model has stopped acting, for the reader. Here the same finding rides on the tool result itself, so the next turn can call run_recipe instead of defending a hand-written copy: on the fourth live run of 00_quickstart the analyst loaded theta_bn, rewrote the formula inline, exported theta_bn, and reported 54.85° from a window the recipe would not have chosen; nothing told it before its answer. Same constants and helpers, imported, so the two readings cannot disagree. Annotates, never blocks — the code ran, its exports stand, and the model may still argue. A run_recipe result is exempt: the shipped source ran.

Parameters:

Name Type Description Default
tool_name str

The tool that just ran.

required
arguments dict | None

Its arguments; the code of a run_python is what is read.

required
result ToolResult

Its result; the exports say what the code computed.

required
history list

The run so far, for the recipes it loaded or ran.

required

Returns:

Type Description
ToolResult

The result, with recipe_available — [{recipe, reason, run_with}] — added to

ToolResult

the payload when a signal fires; unchanged otherwise.

Source code in helioai/core/tool_exec.py
def recipe_available(
    tool_name: str, arguments: dict | None, result: ToolResult, history: list
) -> ToolResult:
    """Tell the model, on the `run_python` result it is about to read, that a shipped
    recipe covers what its code just computed.

    `_flag_recipe_bypass` reads the same signals at the end of the turn and says so in
    the answer — after the model has stopped acting, for the reader. Here the same
    finding rides on the tool result itself, so the next turn can call `run_recipe`
    instead of defending a hand-written copy: on the fourth live run of 00_quickstart
    the analyst loaded `theta_bn`, rewrote the formula inline, exported `theta_bn`, and
    reported 54.85° from a window the recipe would not have chosen; nothing told it
    before its answer. Same constants and helpers, imported, so the two readings cannot
    disagree. Annotates, never blocks — the code ran, its exports stand, and the model
    may still argue. A `run_recipe` result is exempt: the shipped source ran.

    Args:
        tool_name: The tool that just ran.
        arguments: Its arguments; the `code` of a `run_python` is what is read.
        result: Its result; the exports say what the code computed.
        history: The run so far, for the recipes it loaded or ran.

    Returns:
        The result, with `recipe_available` — `[{recipe, reason, run_with}]` — added to
        the payload when a signal fires; unchanged otherwise.
    """
    if tool_name != "run_python":
        return result
    data = result.payload
    if not isinstance(data, dict) or data.get("error"):
        return result
    exported = {
        name.lower()
        for name, stats in (data.get("exports") or {}).items()
        if isinstance(stats, dict) and not stats.get("error")
    }
    if not exported:
        return result
    from helioai.tools.recipes import HELPER_ALTERNATIVES, RECIPE_SIGNATURES, _recipe_path, run_with

    code = str((arguments or {}).get("code") or "")
    loaded, ran = _loaded_recipes(history)
    findings: list[dict] = []
    for name, source in loaded.items():
        if name in ran:
            continue
        functions = _recipe_functions(source)
        signatures = RECIPE_SIGNATURES.get(name, (name,))
        exports_it = any(sig in e for e in exported for sig in signatures)
        redefines_it = bool(functions & _recipe_functions(code))
        if functions and not _calls_any(code, functions) and (exports_it or redefines_it):
            findings.append(
                {"recipe": name, "reason": "not_called", "run_with": run_with(name, source)}
            )
    named = {f["recipe"] for f in findings}
    for name, signatures in RECIPE_SIGNATURES.items():
        if name in loaded or name in ran or name in named:
            continue
        if any(h in code for h in HELPER_ALTERNATIVES.get(name, ())):
            continue
        if any(sig in e for e in exported for sig in signatures):
            path = _recipe_path(name)
            source = path.read_text(encoding="utf-8") if path else ""
            findings.append(
                {"recipe": name, "reason": "not_loaded", "run_with": run_with(name, source)}
            )
    if not findings:
        return result
    return result.with_payload({**data, "recipe_available": findings})

check_answer

check_answer(text: str, history: list, artifacts: list[dict]) -> tuple[str, list[str], list[dict]]

Confront a finished answer with the catalogue and the recipe shelf.

Both loops call this, which is the whole point of it living here. Both checks were written inside the sub-agent loop and stayed there, so a lead agent that did the physics itself — Acts III and IV of the showcase notebook, on a run where load_recipe was called zero times all session — was never checked at all. The detectors were not silent because the run was clean; they were silent because nothing called them.

Parameters:

Name Type Description Default
text str

The finished answer.

required
history list

The turn's messages, read for evidence a recipe was loaded and then bypassed.

required
artifacts list[dict]

What the run exported, used to tell a computed value from a quoted one.

required

Returns:

Type Description
tuple[str, list[str], list[dict]]

The text (annotated when a check fires), the unknown ids, and the recipe flags.

Source code in helioai/core/tool_exec.py
def check_answer(
    text: str, history: list, artifacts: list[dict]
) -> tuple[str, list[str], list[dict]]:
    """Confront a finished answer with the catalogue and the recipe shelf.

    Both loops call this, which is the whole point of it living here. Both checks were
    written inside the sub-agent loop and stayed there, so a lead agent that did the
    physics itself — Acts III and IV of the showcase notebook, on a run where
    `load_recipe` was called zero times all session — was never checked at all. The
    detectors were not silent because the run was clean; they were silent because
    nothing called them.

    Args:
        text: The finished answer.
        history: The turn's messages, read for evidence a recipe was loaded and
            then bypassed.
        artifacts: What the run exported, used to tell a computed value from a
            quoted one.

    Returns:
        The text (annotated when a check fires), the unknown ids, and the recipe flags.
    """
    text, bogus = _flag_unknown_ids(text)
    text, bypassed = _flag_recipe_bypass(text, history, artifacts)
    return text, bogus, bypassed

Sessions

helioai.core.session

Conversation history keyed by (user_id, session_id), persisted to SQLite.

SessionStore

Conversation history keyed by (user_id, session_id), persisted to SQLite.

Histories are cached in memory per key and written back whole on save. Tests use a real database on tmp_path rather than a mock: a mocked store passed happily through a schema migration that broke production.

Example

store = SessionStore(tmp_path / "sessions.db") history = store.get_or_create("cli", "sess-1") # [] on first call history.append(Message(role="user", content="hello")) store.save("cli", "sess-1", history) store.get_or_create("cli", "sess-1")[0].role 'user'

Source code in helioai/core/session.py
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
class SessionStore:
    """Conversation history keyed by (user_id, session_id), persisted to SQLite.

    Histories are cached in memory per key and written back whole on `save`.
    Tests use a real database on `tmp_path` rather than a mock: a mocked store
    passed happily through a schema migration that broke production.

    Example:
        >>> store = SessionStore(tmp_path / "sessions.db")
        >>> history = store.get_or_create("cli", "sess-1")   # [] on first call
        >>> history.append(Message(role="user", content="hello"))
        >>> store.save("cli", "sess-1", history)
        >>> store.get_or_create("cli", "sess-1")[0].role
        'user'
    """

    def __init__(self, db_path: Path = DEFAULT_DB) -> None:
        self._db_path = db_path
        self._cache: dict[SessionKey, list[Message]] = {}
        self._lock = threading.Lock()
        self._turn_locks: dict[SessionKey, asyncio.Lock] = {}
        # Nothing touches the disk until the first query: the module-level `store`
        # is built at import, and creating directories on `import helioai` is the
        # kind of side effect a packager, a linter run or a docs build trips over.
        self._schema_ready = False
        # Its own lock: `save` already holds `_lock` when it connects, and `_lock` is
        # not re-entrant — taking it again here would deadlock the first save.
        self._schema_lock = threading.Lock()

    def _init_schema(self) -> None:
        self._db_path.parent.mkdir(parents=True, exist_ok=True)
        conn = sqlite3.connect(self._db_path, check_same_thread=False)
        try:
            conn.executescript(_SCHEMA)
            for statement in _MIGRATIONS:
                try:
                    conn.execute(statement)
                except Exception:
                    pass
            conn.commit()
        finally:
            conn.close()
        self._schema_ready = True

    @contextmanager
    def _connect(self) -> Iterator[sqlite3.Connection]:
        if not self._schema_ready:
            with self._schema_lock:
                if not self._schema_ready:
                    self._init_schema()
        conn = sqlite3.connect(self._db_path, check_same_thread=False)
        conn.execute("PRAGMA foreign_keys = ON")
        conn.execute("PRAGMA journal_mode = WAL")
        conn.execute("PRAGMA busy_timeout = 5000")
        try:
            yield conn
        finally:
            conn.close()

    def turn_lock(self, user_id: str, session_id: str) -> asyncio.Lock:
        """The lock a caller must hold for the whole of one conversational turn.

        `get_or_create` hands every caller the same list, and a turn is a
        read-modify-write of it that spans several awaits: without this, two turns on
        one session — two browser tabs, two MCP calls — interleave their appends and
        `save` persists the mix. `_lock` only serialises the SQL, never the turn.

        One `asyncio.Lock` per key, created on first use. A lock in Python ≥ 3.10 binds
        to an event loop only on its first *contended* acquisition, so the CLI and the
        Jupyter magic — a fresh `asyncio.run` per question, never two turns on one
        session at once — reuse it across loops safely, while the web and MCP servers
        run one loop. If that assumption ever breaks, asyncio raises a `RuntimeError`
        naming the loop mismatch instead of silently corrupting anything.

        Args:
            user_id: Owner of the session.
            session_id: Session identifier.

        Returns:
            The same lock object for the same key, for the life of the store.
        """
        key: SessionKey = (user_id, session_id)
        with self._lock:
            lock = self._turn_locks.get(key)
            if lock is None:
                lock = self._turn_locks[key] = asyncio.Lock()
            return lock

    def is_busy(self, user_id: str, session_id: str) -> bool:
        """Whether a turn is currently running for this session.

        Read without taking the lock — this is the web layer's fast refusal (409), not
        a guarantee; the guarantee is `turn_lock` itself.
        """
        lock = self._turn_locks.get((user_id, session_id))
        return bool(lock and lock.locked())

    def get_or_create(self, user_id: str, session_id: str) -> list[Message]:
        """Return the cached history for a session, loading it from disk if needed."""
        key: SessionKey = (user_id, session_id)
        with self._lock:
            if key in self._cache:
                return self._cache[key]
            history = self._load(user_id, session_id)
            self._cache[key] = history
            return history

    def _load(self, user_id: str, session_id: str) -> list[Message]:
        with self._connect() as conn:
            rows = conn.execute(
                "SELECT role, content, tool_calls, tool_call_id, origin, name, reasoning "
                "FROM messages WHERE user_id = ? AND session_id = ? ORDER BY seq",
                (user_id, session_id),
            ).fetchall()
        return [
            Message(
                role=role,
                content=content or "",
                tool_calls=_load_tool_calls(tool_calls),
                tool_call_id=tool_call_id,
                origin=origin,
                name=name,
                reasoning=reasoning,
            )
            for role, content, tool_calls, tool_call_id, origin, name, reasoning in rows
        ]

    def save(self, user_id: str, session_id: str, history: list[Message]) -> None:
        """Replace a session's stored history.

        Args:
            user_id: Owner of the session.
            session_id: Session identifier.
            history: Full message list; it replaces whatever was stored.
        """
        with self._lock, self._connect() as conn:
            conn.execute(
                "INSERT INTO sessions(user_id, session_id) VALUES(?, ?) "
                "ON CONFLICT(user_id, session_id) DO UPDATE SET updated_at = julianday('now')",
                (user_id, session_id),
            )
            conn.execute(
                "DELETE FROM messages WHERE user_id = ? AND session_id = ?", (user_id, session_id)
            )
            conn.executemany(
                "INSERT INTO messages(user_id, session_id, seq, role, content, tool_calls, "
                "tool_call_id, origin, name, reasoning) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
                [
                    (
                        user_id,
                        session_id,
                        i,
                        m.role,
                        m.content or "",
                        _dump_tool_calls(m.tool_calls),
                        m.tool_call_id,
                        m.origin,
                        m.name,
                        m.reasoning,
                    )
                    for i, m in enumerate(history)
                ],
            )
            conn.commit()

    def record_usage(
        self,
        user_id: str,
        session_id: str,
        *,
        turn: int | None,
        agent: str,
        provider: str,
        prompt_tokens: int,
        completion_tokens: int,
        cached_tokens: int = 0,
    ) -> None:
        """Append what one LLM call cost, as the provider reported it.

        The counts have ridden on `Message` for a while and were dropped at save time,
        so a session reloaded from disk reported no cost and nothing could say what a
        user had spent. Kept apart from `messages` on purpose: one row per call, never
        rewritten by `save`, so a compacted or reset history does not erase the bill.
        Zero counts are skipped — a provider that reports none leaves no row rather
        than a misleading zero.

        Args:
            user_id: Owner of the session.
            session_id: The conversation the call belonged to; a sub-agent's calls
                are charged to its parent session.
            turn: Turn index within the run, for ordering.
            agent: `"lead"` or the sub-agent role.
            provider: Provider name, since a session may switch providers.
            prompt_tokens: Input tokens billed.
            completion_tokens: Output tokens billed.
            cached_tokens: The part of the prompt served from the provider's cache.
        """
        if not (prompt_tokens or completion_tokens):
            return
        with self._lock, self._connect() as conn:
            conn.execute(
                "INSERT INTO usage(user_id, session_id, turn, agent, provider, "
                "prompt_tokens, completion_tokens, cached_tokens) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
                (
                    user_id,
                    session_id,
                    turn,
                    agent,
                    provider,
                    int(prompt_tokens),
                    int(completion_tokens),
                    int(cached_tokens),
                ),
            )
            conn.commit()

    def usage_totals(
        self, user_id: str, session_id: str | None = None, since_days: float | None = None
    ) -> dict:
        """Sum a user's token usage, optionally for one session or a recent window.

        Args:
            user_id: Whose usage.
            session_id: Restrict to one session; None for every session.
            since_days: Only calls in the last N days; None for all time.

        Returns:
            `{"prompt_tokens", "completion_tokens", "cached_tokens", "n_calls"}`, zeros
            when nothing was recorded.
        """
        clauses = ["user_id = ?"]
        params: list = [user_id]
        if session_id is not None:
            clauses.append("session_id = ?")
            params.append(session_id)
        if since_days is not None:
            clauses.append("recorded_at >= julianday('now') - ?")
            params.append(float(since_days))
        with self._connect() as conn:
            row = conn.execute(
                "SELECT COALESCE(SUM(prompt_tokens), 0), COALESCE(SUM(completion_tokens), 0), "
                "COALESCE(SUM(cached_tokens), 0), COUNT(*) FROM usage WHERE "
                + " AND ".join(clauses),
                params,
            ).fetchone()
        return {
            "prompt_tokens": int(row[0]),
            "completion_tokens": int(row[1]),
            "cached_tokens": int(row[2]),
            "n_calls": int(row[3]),
        }

    def append_event(self, user_id: str, session_id: str, event: dict) -> None:
        """Journal one event of a turn, in the order it was yielded.

        The journal is what a session replays from: the browser used to rebuild a
        past conversation by re-parsing the JSON of every `tool` message with a
        hundred lines of shape-sniffing, and lost the plan, the provenance verdict, the
        figure reviews and everything a sub-agent did on the way. Written *before* the
        event is handed to the interface, so a stream cut mid-turn still leaves what was
        shown. Its own table, never rewritten by `save`, so a compacted history does
        not erase the record.

        Args:
            user_id: Owner of the session.
            session_id: The conversation the event belongs to.
            event: `{"event": kind, "data": {...}}` as `events.make` builds it. Payloads
                are stored as JSON; unknown types are stringified, as they are for the
                model.
        """
        with self._lock, self._connect() as conn:
            conn.execute(
                "INSERT INTO events(user_id, session_id, seq, kind, data) VALUES (?, ?, "
                "(SELECT COALESCE(MAX(seq), -1) + 1 FROM events WHERE user_id = ? AND session_id = ?),"
                " ?, ?)",
                (
                    user_id,
                    session_id,
                    user_id,
                    session_id,
                    event["event"],
                    json.dumps(event.get("data") or {}, ensure_ascii=False, default=str),
                ),
            )
            conn.commit()

    def events(self, user_id: str, session_id: str) -> list[dict]:
        """A session's journal, oldest first, in the shape the loops yield.

        Args:
            user_id: Owner of the session.
            session_id: The conversation.

        Returns:
            `[{"event": kind, "data": {...}}, ...]` — exactly what `stream_chat` yielded,
            turn after turn, so a replay renders with the same code as the live stream.
            Empty for a session recorded before the journal existed.
        """
        with self._connect() as conn:
            rows = conn.execute(
                "SELECT kind, data FROM events WHERE user_id = ? AND session_id = ? ORDER BY seq",
                (user_id, session_id),
            ).fetchall()
        return [{"event": kind, "data": json.loads(data)} for kind, data in rows]

    def reset(self, user_id: str, session_id: str) -> None:
        """Delete a session, its messages, its usage rows and its journal, and drop it
        from the cache."""
        with self._lock, self._connect() as conn:
            conn.execute(
                "DELETE FROM sessions WHERE user_id = ? AND session_id = ?",
                (user_id, session_id),
            )
            for table in ("usage", "events"):
                conn.execute(
                    f"DELETE FROM {table} WHERE user_id = ? AND session_id = ?",  # noqa: S608
                    (user_id, session_id),
                )
            conn.commit()
        self._cache.pop((user_id, session_id), None)
        self._turn_locks.pop((user_id, session_id), None)

    def set_workspace_dir(self, user_id: str, session_id: str, workspace_dir: str) -> None:
        """Record which workspace directory a session's artifacts live in."""
        with self._lock, self._connect() as conn:
            conn.execute(
                "UPDATE sessions SET workspace_dir = ? WHERE user_id = ? AND session_id = ?",
                (workspace_dir, user_id, session_id),
            )
            conn.commit()

    def get_workspace_dir(self, user_id: str, session_id: str) -> str | None:
        """Return a session's workspace directory label, or None."""
        with self._connect() as conn:
            row = conn.execute(
                "SELECT workspace_dir FROM sessions WHERE user_id = ? AND session_id = ?",
                (user_id, session_id),
            ).fetchone()
        return row[0] if row else None

    def workspace_dirs(self, user_id: str) -> set[str]:
        """All workspace dir labels owned by a user (for path-ownership checks)."""
        with self._connect() as conn:
            rows = conn.execute(
                "SELECT workspace_dir FROM sessions "
                "WHERE user_id = ? AND workspace_dir IS NOT NULL",
                (user_id,),
            ).fetchall()
        return {r[0] for r in rows}

    def all_sessions(self, user_id: str) -> list[str]:
        """Return a user's session ids, most recently updated first."""
        with self._connect() as conn:
            rows = conn.execute(
                "SELECT session_id FROM sessions WHERE user_id = ? "
                "ORDER BY updated_at DESC, rowid DESC",
                (user_id,),
            ).fetchall()
        return [r[0] for r in rows]

    def list_summaries(self, user_id: str, limit: int = 50) -> list[dict]:
        """Summarise a user's recent sessions for the history view.

        Args:
            user_id: Owner of the sessions.
            limit: Maximum number of sessions to return.

        Returns:
            Dicts with session_id, updated_at, first_message, n_messages,
            workspace_dir and tokens (prompt + completion, all calls), most recent first.
        """
        from datetime import datetime

        with self._connect() as conn:
            rows = conn.execute(
                """
                SELECT s.session_id, s.updated_at AS jd,
                       (SELECT content FROM messages WHERE user_id = s.user_id
                          AND session_id = s.session_id AND role = 'user'
                          ORDER BY seq LIMIT 1) AS first_user,
                       (SELECT COUNT(*) FROM messages WHERE user_id = s.user_id
                          AND session_id = s.session_id) AS n_messages,
                       s.workspace_dir,
                       (SELECT COALESCE(SUM(prompt_tokens + completion_tokens), 0) FROM usage
                          WHERE user_id = s.user_id AND session_id = s.session_id) AS tokens
                FROM sessions s
                WHERE s.user_id = ?
                ORDER BY s.updated_at DESC, s.rowid DESC LIMIT ?
                """,
                (user_id, limit),
            ).fetchall()
        out: list[dict] = []
        for session_id, jd, first_user, n_messages, workspace_dir, tokens in rows:
            preview = (first_user or "").strip().replace("\n", " ")
            if len(preview) > 80:
                preview = preview[:77] + "..."
            unix_ts = (jd - 2440587.5) * 86400
            iso = datetime.fromtimestamp(unix_ts, tz=UTC).isoformat().replace("+00:00", "Z")
            out.append(
                {
                    "session_id": session_id,
                    "first_message": preview,
                    "n_messages": n_messages,
                    "updated_at": iso,
                    "workspace_dir": workspace_dir,
                    "tokens": int(tokens),
                }
            )
        return out

turn_lock

turn_lock(user_id: str, session_id: str) -> asyncio.Lock

The lock a caller must hold for the whole of one conversational turn.

get_or_create hands every caller the same list, and a turn is a read-modify-write of it that spans several awaits: without this, two turns on one session — two browser tabs, two MCP calls — interleave their appends and save persists the mix. _lock only serialises the SQL, never the turn.

One asyncio.Lock per key, created on first use. A lock in Python ≥ 3.10 binds to an event loop only on its first contended acquisition, so the CLI and the Jupyter magic — a fresh asyncio.run per question, never two turns on one session at once — reuse it across loops safely, while the web and MCP servers run one loop. If that assumption ever breaks, asyncio raises a RuntimeError naming the loop mismatch instead of silently corrupting anything.

Parameters:

Name Type Description Default
user_id str

Owner of the session.

required
session_id str

Session identifier.

required

Returns:

Type Description
Lock

The same lock object for the same key, for the life of the store.

Source code in helioai/core/session.py
def turn_lock(self, user_id: str, session_id: str) -> asyncio.Lock:
    """The lock a caller must hold for the whole of one conversational turn.

    `get_or_create` hands every caller the same list, and a turn is a
    read-modify-write of it that spans several awaits: without this, two turns on
    one session — two browser tabs, two MCP calls — interleave their appends and
    `save` persists the mix. `_lock` only serialises the SQL, never the turn.

    One `asyncio.Lock` per key, created on first use. A lock in Python ≥ 3.10 binds
    to an event loop only on its first *contended* acquisition, so the CLI and the
    Jupyter magic — a fresh `asyncio.run` per question, never two turns on one
    session at once — reuse it across loops safely, while the web and MCP servers
    run one loop. If that assumption ever breaks, asyncio raises a `RuntimeError`
    naming the loop mismatch instead of silently corrupting anything.

    Args:
        user_id: Owner of the session.
        session_id: Session identifier.

    Returns:
        The same lock object for the same key, for the life of the store.
    """
    key: SessionKey = (user_id, session_id)
    with self._lock:
        lock = self._turn_locks.get(key)
        if lock is None:
            lock = self._turn_locks[key] = asyncio.Lock()
        return lock

is_busy

is_busy(user_id: str, session_id: str) -> bool

Whether a turn is currently running for this session.

Read without taking the lock — this is the web layer's fast refusal (409), not a guarantee; the guarantee is turn_lock itself.

Source code in helioai/core/session.py
def is_busy(self, user_id: str, session_id: str) -> bool:
    """Whether a turn is currently running for this session.

    Read without taking the lock — this is the web layer's fast refusal (409), not
    a guarantee; the guarantee is `turn_lock` itself.
    """
    lock = self._turn_locks.get((user_id, session_id))
    return bool(lock and lock.locked())

get_or_create

get_or_create(user_id: str, session_id: str) -> list[Message]

Return the cached history for a session, loading it from disk if needed.

Source code in helioai/core/session.py
def get_or_create(self, user_id: str, session_id: str) -> list[Message]:
    """Return the cached history for a session, loading it from disk if needed."""
    key: SessionKey = (user_id, session_id)
    with self._lock:
        if key in self._cache:
            return self._cache[key]
        history = self._load(user_id, session_id)
        self._cache[key] = history
        return history

save

save(user_id: str, session_id: str, history: list[Message]) -> None

Replace a session's stored history.

Parameters:

Name Type Description Default
user_id str

Owner of the session.

required
session_id str

Session identifier.

required
history list[Message]

Full message list; it replaces whatever was stored.

required
Source code in helioai/core/session.py
def save(self, user_id: str, session_id: str, history: list[Message]) -> None:
    """Replace a session's stored history.

    Args:
        user_id: Owner of the session.
        session_id: Session identifier.
        history: Full message list; it replaces whatever was stored.
    """
    with self._lock, self._connect() as conn:
        conn.execute(
            "INSERT INTO sessions(user_id, session_id) VALUES(?, ?) "
            "ON CONFLICT(user_id, session_id) DO UPDATE SET updated_at = julianday('now')",
            (user_id, session_id),
        )
        conn.execute(
            "DELETE FROM messages WHERE user_id = ? AND session_id = ?", (user_id, session_id)
        )
        conn.executemany(
            "INSERT INTO messages(user_id, session_id, seq, role, content, tool_calls, "
            "tool_call_id, origin, name, reasoning) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
            [
                (
                    user_id,
                    session_id,
                    i,
                    m.role,
                    m.content or "",
                    _dump_tool_calls(m.tool_calls),
                    m.tool_call_id,
                    m.origin,
                    m.name,
                    m.reasoning,
                )
                for i, m in enumerate(history)
            ],
        )
        conn.commit()

record_usage

record_usage(user_id: str, session_id: str, *, turn: int | None, agent: str, provider: str, prompt_tokens: int, completion_tokens: int, cached_tokens: int = 0) -> None

Append what one LLM call cost, as the provider reported it.

The counts have ridden on Message for a while and were dropped at save time, so a session reloaded from disk reported no cost and nothing could say what a user had spent. Kept apart from messages on purpose: one row per call, never rewritten by save, so a compacted or reset history does not erase the bill. Zero counts are skipped — a provider that reports none leaves no row rather than a misleading zero.

Parameters:

Name Type Description Default
user_id str

Owner of the session.

required
session_id str

The conversation the call belonged to; a sub-agent's calls are charged to its parent session.

required
turn int | None

Turn index within the run, for ordering.

required
agent str

"lead" or the sub-agent role.

required
provider str

Provider name, since a session may switch providers.

required
prompt_tokens int

Input tokens billed.

required
completion_tokens int

Output tokens billed.

required
cached_tokens int

The part of the prompt served from the provider's cache.

0
Source code in helioai/core/session.py
def record_usage(
    self,
    user_id: str,
    session_id: str,
    *,
    turn: int | None,
    agent: str,
    provider: str,
    prompt_tokens: int,
    completion_tokens: int,
    cached_tokens: int = 0,
) -> None:
    """Append what one LLM call cost, as the provider reported it.

    The counts have ridden on `Message` for a while and were dropped at save time,
    so a session reloaded from disk reported no cost and nothing could say what a
    user had spent. Kept apart from `messages` on purpose: one row per call, never
    rewritten by `save`, so a compacted or reset history does not erase the bill.
    Zero counts are skipped — a provider that reports none leaves no row rather
    than a misleading zero.

    Args:
        user_id: Owner of the session.
        session_id: The conversation the call belonged to; a sub-agent's calls
            are charged to its parent session.
        turn: Turn index within the run, for ordering.
        agent: `"lead"` or the sub-agent role.
        provider: Provider name, since a session may switch providers.
        prompt_tokens: Input tokens billed.
        completion_tokens: Output tokens billed.
        cached_tokens: The part of the prompt served from the provider's cache.
    """
    if not (prompt_tokens or completion_tokens):
        return
    with self._lock, self._connect() as conn:
        conn.execute(
            "INSERT INTO usage(user_id, session_id, turn, agent, provider, "
            "prompt_tokens, completion_tokens, cached_tokens) VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
            (
                user_id,
                session_id,
                turn,
                agent,
                provider,
                int(prompt_tokens),
                int(completion_tokens),
                int(cached_tokens),
            ),
        )
        conn.commit()

usage_totals

usage_totals(user_id: str, session_id: str | None = None, since_days: float | None = None) -> dict

Sum a user's token usage, optionally for one session or a recent window.

Parameters:

Name Type Description Default
user_id str

Whose usage.

required
session_id str | None

Restrict to one session; None for every session.

None
since_days float | None

Only calls in the last N days; None for all time.

None

Returns:

Type Description
dict

{"prompt_tokens", "completion_tokens", "cached_tokens", "n_calls"}, zeros

dict

when nothing was recorded.

Source code in helioai/core/session.py
def usage_totals(
    self, user_id: str, session_id: str | None = None, since_days: float | None = None
) -> dict:
    """Sum a user's token usage, optionally for one session or a recent window.

    Args:
        user_id: Whose usage.
        session_id: Restrict to one session; None for every session.
        since_days: Only calls in the last N days; None for all time.

    Returns:
        `{"prompt_tokens", "completion_tokens", "cached_tokens", "n_calls"}`, zeros
        when nothing was recorded.
    """
    clauses = ["user_id = ?"]
    params: list = [user_id]
    if session_id is not None:
        clauses.append("session_id = ?")
        params.append(session_id)
    if since_days is not None:
        clauses.append("recorded_at >= julianday('now') - ?")
        params.append(float(since_days))
    with self._connect() as conn:
        row = conn.execute(
            "SELECT COALESCE(SUM(prompt_tokens), 0), COALESCE(SUM(completion_tokens), 0), "
            "COALESCE(SUM(cached_tokens), 0), COUNT(*) FROM usage WHERE "
            + " AND ".join(clauses),
            params,
        ).fetchone()
    return {
        "prompt_tokens": int(row[0]),
        "completion_tokens": int(row[1]),
        "cached_tokens": int(row[2]),
        "n_calls": int(row[3]),
    }

append_event

append_event(user_id: str, session_id: str, event: dict) -> None

Journal one event of a turn, in the order it was yielded.

The journal is what a session replays from: the browser used to rebuild a past conversation by re-parsing the JSON of every tool message with a hundred lines of shape-sniffing, and lost the plan, the provenance verdict, the figure reviews and everything a sub-agent did on the way. Written before the event is handed to the interface, so a stream cut mid-turn still leaves what was shown. Its own table, never rewritten by save, so a compacted history does not erase the record.

Parameters:

Name Type Description Default
user_id str

Owner of the session.

required
session_id str

The conversation the event belongs to.

required
event dict

{"event": kind, "data": {...}} as events.make builds it. Payloads are stored as JSON; unknown types are stringified, as they are for the model.

required
Source code in helioai/core/session.py
def append_event(self, user_id: str, session_id: str, event: dict) -> None:
    """Journal one event of a turn, in the order it was yielded.

    The journal is what a session replays from: the browser used to rebuild a
    past conversation by re-parsing the JSON of every `tool` message with a
    hundred lines of shape-sniffing, and lost the plan, the provenance verdict, the
    figure reviews and everything a sub-agent did on the way. Written *before* the
    event is handed to the interface, so a stream cut mid-turn still leaves what was
    shown. Its own table, never rewritten by `save`, so a compacted history does
    not erase the record.

    Args:
        user_id: Owner of the session.
        session_id: The conversation the event belongs to.
        event: `{"event": kind, "data": {...}}` as `events.make` builds it. Payloads
            are stored as JSON; unknown types are stringified, as they are for the
            model.
    """
    with self._lock, self._connect() as conn:
        conn.execute(
            "INSERT INTO events(user_id, session_id, seq, kind, data) VALUES (?, ?, "
            "(SELECT COALESCE(MAX(seq), -1) + 1 FROM events WHERE user_id = ? AND session_id = ?),"
            " ?, ?)",
            (
                user_id,
                session_id,
                user_id,
                session_id,
                event["event"],
                json.dumps(event.get("data") or {}, ensure_ascii=False, default=str),
            ),
        )
        conn.commit()

events

events(user_id: str, session_id: str) -> list[dict]

A session's journal, oldest first, in the shape the loops yield.

Parameters:

Name Type Description Default
user_id str

Owner of the session.

required
session_id str

The conversation.

required

Returns:

Type Description
list[dict]

[{"event": kind, "data": {...}}, ...] — exactly what stream_chat yielded,

list[dict]

turn after turn, so a replay renders with the same code as the live stream.

list[dict]

Empty for a session recorded before the journal existed.

Source code in helioai/core/session.py
def events(self, user_id: str, session_id: str) -> list[dict]:
    """A session's journal, oldest first, in the shape the loops yield.

    Args:
        user_id: Owner of the session.
        session_id: The conversation.

    Returns:
        `[{"event": kind, "data": {...}}, ...]` — exactly what `stream_chat` yielded,
        turn after turn, so a replay renders with the same code as the live stream.
        Empty for a session recorded before the journal existed.
    """
    with self._connect() as conn:
        rows = conn.execute(
            "SELECT kind, data FROM events WHERE user_id = ? AND session_id = ? ORDER BY seq",
            (user_id, session_id),
        ).fetchall()
    return [{"event": kind, "data": json.loads(data)} for kind, data in rows]

reset

reset(user_id: str, session_id: str) -> None

Delete a session, its messages, its usage rows and its journal, and drop it from the cache.

Source code in helioai/core/session.py
def reset(self, user_id: str, session_id: str) -> None:
    """Delete a session, its messages, its usage rows and its journal, and drop it
    from the cache."""
    with self._lock, self._connect() as conn:
        conn.execute(
            "DELETE FROM sessions WHERE user_id = ? AND session_id = ?",
            (user_id, session_id),
        )
        for table in ("usage", "events"):
            conn.execute(
                f"DELETE FROM {table} WHERE user_id = ? AND session_id = ?",  # noqa: S608
                (user_id, session_id),
            )
        conn.commit()
    self._cache.pop((user_id, session_id), None)
    self._turn_locks.pop((user_id, session_id), None)

set_workspace_dir

set_workspace_dir(user_id: str, session_id: str, workspace_dir: str) -> None

Record which workspace directory a session's artifacts live in.

Source code in helioai/core/session.py
def set_workspace_dir(self, user_id: str, session_id: str, workspace_dir: str) -> None:
    """Record which workspace directory a session's artifacts live in."""
    with self._lock, self._connect() as conn:
        conn.execute(
            "UPDATE sessions SET workspace_dir = ? WHERE user_id = ? AND session_id = ?",
            (workspace_dir, user_id, session_id),
        )
        conn.commit()

get_workspace_dir

get_workspace_dir(user_id: str, session_id: str) -> str | None

Return a session's workspace directory label, or None.

Source code in helioai/core/session.py
def get_workspace_dir(self, user_id: str, session_id: str) -> str | None:
    """Return a session's workspace directory label, or None."""
    with self._connect() as conn:
        row = conn.execute(
            "SELECT workspace_dir FROM sessions WHERE user_id = ? AND session_id = ?",
            (user_id, session_id),
        ).fetchone()
    return row[0] if row else None

workspace_dirs

workspace_dirs(user_id: str) -> set[str]

All workspace dir labels owned by a user (for path-ownership checks).

Source code in helioai/core/session.py
def workspace_dirs(self, user_id: str) -> set[str]:
    """All workspace dir labels owned by a user (for path-ownership checks)."""
    with self._connect() as conn:
        rows = conn.execute(
            "SELECT workspace_dir FROM sessions "
            "WHERE user_id = ? AND workspace_dir IS NOT NULL",
            (user_id,),
        ).fetchall()
    return {r[0] for r in rows}

all_sessions

all_sessions(user_id: str) -> list[str]

Return a user's session ids, most recently updated first.

Source code in helioai/core/session.py
def all_sessions(self, user_id: str) -> list[str]:
    """Return a user's session ids, most recently updated first."""
    with self._connect() as conn:
        rows = conn.execute(
            "SELECT session_id FROM sessions WHERE user_id = ? "
            "ORDER BY updated_at DESC, rowid DESC",
            (user_id,),
        ).fetchall()
    return [r[0] for r in rows]

list_summaries

list_summaries(user_id: str, limit: int = 50) -> list[dict]

Summarise a user's recent sessions for the history view.

Parameters:

Name Type Description Default
user_id str

Owner of the sessions.

required
limit int

Maximum number of sessions to return.

50

Returns:

Type Description
list[dict]

Dicts with session_id, updated_at, first_message, n_messages,

list[dict]

workspace_dir and tokens (prompt + completion, all calls), most recent first.

Source code in helioai/core/session.py
def list_summaries(self, user_id: str, limit: int = 50) -> list[dict]:
    """Summarise a user's recent sessions for the history view.

    Args:
        user_id: Owner of the sessions.
        limit: Maximum number of sessions to return.

    Returns:
        Dicts with session_id, updated_at, first_message, n_messages,
        workspace_dir and tokens (prompt + completion, all calls), most recent first.
    """
    from datetime import datetime

    with self._connect() as conn:
        rows = conn.execute(
            """
            SELECT s.session_id, s.updated_at AS jd,
                   (SELECT content FROM messages WHERE user_id = s.user_id
                      AND session_id = s.session_id AND role = 'user'
                      ORDER BY seq LIMIT 1) AS first_user,
                   (SELECT COUNT(*) FROM messages WHERE user_id = s.user_id
                      AND session_id = s.session_id) AS n_messages,
                   s.workspace_dir,
                   (SELECT COALESCE(SUM(prompt_tokens + completion_tokens), 0) FROM usage
                      WHERE user_id = s.user_id AND session_id = s.session_id) AS tokens
            FROM sessions s
            WHERE s.user_id = ?
            ORDER BY s.updated_at DESC, s.rowid DESC LIMIT ?
            """,
            (user_id, limit),
        ).fetchall()
    out: list[dict] = []
    for session_id, jd, first_user, n_messages, workspace_dir, tokens in rows:
        preview = (first_user or "").strip().replace("\n", " ")
        if len(preview) > 80:
            preview = preview[:77] + "..."
        unix_ts = (jd - 2440587.5) * 86400
        iso = datetime.fromtimestamp(unix_ts, tz=UTC).isoformat().replace("+00:00", "Z")
        out.append(
            {
                "session_id": session_id,
                "first_message": preview,
                "n_messages": n_messages,
                "updated_at": iso,
                "workspace_dir": workspace_dir,
                "tokens": int(tokens),
            }
        )
    return out

strip_orphan_tool_calls

strip_orphan_tool_calls(history: list[Message]) -> list[Message]

Remove assistant tool_calls that have no matching tool response.

An interrupted generation (e.g. client disconnect mid-tool) can leave an assistant message with tool_calls but no corresponding tool messages in the history. Sending such a sequence to the LLM API causes a 400 error.

For each orphaned tool_call id: - If the assistant message has content too, keep the message but drop the orphaned tool_calls list entry (or clear it entirely if all are orphaned). - If the assistant message has no content and all its tool_calls are orphaned, drop the message entirely.

Parameters:

Name Type Description Default
history list[Message]

Messages in order, as loaded from the session store.

required

Returns:

Type Description
list[Message]

A new list; the input is not modified. Messages without tool calls pass

list[Message]

through untouched, so a clean history is returned equal to its input.

Source code in helioai/core/session.py
def strip_orphan_tool_calls(history: list[Message]) -> list[Message]:
    """Remove assistant tool_calls that have no matching tool response.

    An interrupted generation (e.g. client disconnect mid-tool) can leave an
    assistant message with tool_calls but no corresponding tool messages in the
    history.  Sending such a sequence to the LLM API causes a 400 error.

    For each orphaned tool_call id:
    - If the assistant message has content too, keep the message but drop the
      orphaned tool_calls list entry (or clear it entirely if all are orphaned).
    - If the assistant message has no content and all its tool_calls are
      orphaned, drop the message entirely.

    Args:
        history: Messages in order, as loaded from the session store.

    Returns:
        A new list; the input is not modified. Messages without tool calls pass
        through untouched, so a clean history is returned equal to its input.
    """
    answered: set[str] = {m.tool_call_id for m in history if m.tool_call_id}
    cleaned: list[Message] = []
    for m in history:
        if not m.tool_calls:
            cleaned.append(m)
            continue
        live_tcs = [tc for tc in m.tool_calls if tc.id in answered]
        if len(live_tcs) == len(m.tool_calls):
            cleaned.append(m)
        elif live_tcs:
            cleaned.append(replace(m, tool_calls=live_tcs))
        elif m.content:
            cleaned.append(replace(m, tool_calls=None))
        # else: drop the message entirely (no content, no answered tool_calls)
    return cleaned

Skills

helioai.core.skills_loader

Discover and serve markdown-defined skills to the agent.

Each skill lives in skills//SKILL.md with YAML frontmatter: --- name: parameter_hunter description: Find speasy parameter ids from a natural language query. when_to_use: User asks about a parameter without giving its exact id. allowed_tools: [search_parameters] --- # body of the procedure

SkillError

Bases: RuntimeError

Raised when a skill is missing, unreadable, or fails its path check.

Source code in helioai/core/skills_loader.py
class SkillError(RuntimeError):
    """Raised when a skill is missing, unreadable, or fails its path check."""

    pass

SkillMeta dataclass

Header of a skill, as listed to the agent before it loads the body.

Source code in helioai/core/skills_loader.py
@dataclass(frozen=True)
class SkillMeta:
    """Header of a skill, as listed to the agent before it loads the body."""

    name: str
    description: str
    when_to_use: str
    allowed_tools: tuple[str, ...]
    path: Path

load_index

load_index() -> str

Return the markdown index of available skills, for the agent to browse.

Source code in helioai/core/skills_loader.py
def load_index() -> str:
    """Return the markdown index of available skills, for the agent to browse."""
    skills = _discover()
    if not skills:
        return "(no skills available)"
    lines = ["| Skill | When to use |", "|---|---|"]
    for meta in skills.values():
        when = meta.when_to_use.replace("|", "/").replace("\n", " ")
        lines.append(f"| `{meta.name}` | {when} |")
    return "\n".join(lines)

load_skill

load_skill(name: str) -> str

Return a skill's full markdown body.

Parameters:

Name Type Description Default
name str

Skill name as listed by list_skill_names.

required

Returns:

Type Description
str

The skill body, ready to append to a system prompt.

Raises:

Type Description
SkillError

If the skill does not exist or the name escapes the skills directory.

Source code in helioai/core/skills_loader.py
def load_skill(name: str) -> str:
    """Return a skill's full markdown body.

    Args:
        name: Skill name as listed by `list_skill_names`.

    Returns:
        The skill body, ready to append to a system prompt.

    Raises:
        SkillError: If the skill does not exist or the name escapes the skills
            directory.
    """
    skills = _discover()
    if name not in skills:
        known = ", ".join(sorted(skills)) or "(none)"
        raise SkillError(f"unknown skill {name!r}. Known: {known}")
    text = skills[name].path.read_text(encoding="utf-8")
    _, body = _split_frontmatter(text)
    return body

list_skill_names

list_skill_names() -> list[str]

Return the names of all available skills.

Source code in helioai/core/skills_loader.py
def list_skill_names() -> list[str]:
    """Return the names of all available skills."""
    return list(_discover().keys())

list_skills

list_skills() -> list[SkillMeta]

Return metadata for every discovered skill, for listings like MCP resources.

Source code in helioai/core/skills_loader.py
def list_skills() -> list[SkillMeta]:
    """Return metadata for every discovered skill, for listings like MCP resources."""
    return list(_discover().values())

Judgment

The seam for a System One judge beside the loop, and the joins that read its contract. See The judgment layer.

helioai.core.judgment

One seam between HelioAI and a System One judge — every question and every threshold.

Why one file: HelioAI is PyHC-listed and open source, and a reviewer who asks "what does this agent delegate to a proprietary model?" must be able to answer by reading one module. Every site that consults the judge builds its questions here-shaped (Noul, Choice), calls ask, and reads an Answers — or None.

Why abstention is None and not a sentinel: a sentinel can be compared, summed and sorted into a plausible wrong answer; None > 0.5 raises. The judge abstains when the backend is null (the default), when the site's experiment is off, when the call fails or times out, and per question when the answer does not clear that question's own threshold. Every call site reads if answers is None: <today's code> — so a build without a key, without the extra, or with the default backend is bit-for-bit the loop that exists today. That is the property the tests hold this module to.

Why observation only: the 2026-09-15 lesson — nothing that changes what the model sees ships without a replay on both sides of the same question — applies to every answer here. So ask records one JSON line per call under the session workspace (judgment.jsonl), with the state in full: a disagreement that cannot be adjudicated later without the state is not a measurement. What a site does with an answer is the site's business, and until a site has earned it, that is: annotate, never correct.

Why the model call is bounded: a judgment that can hold a turn is a judgment that can hang it. The whole call sits under settings.judgment.timeout_s (round trip measured from France on 2026-09-22: median 266 ms, p95 373 ms), and a timeout abstains.

DELIVERABLES module-attribute

DELIVERABLES: tuple[str, ...] = ('value', 'figure', 'catalogue', 'procedure')

What a request may want delivered, each asked as its own yes/no: one Choice over the same set abstained on every request that wanted two things at once — "plot |B| and compute θ_Bn" is a figure and a value — and abstaining on the commonest shape of a question is not a contract. A request that wants none of the four is an explanation.

Noul dataclass

A yes/no question.

The judge returns a probability; the answer is True at or above 0.5 + margin, False at or below 0.5 - margin, and None in between. The default margin makes 0.7 the least confident "yes" — a site that needs more certainty raises it.

Attributes:

Name Type Description
instructions str

The question, in plain English, as the judge reads it.

margin float

Half-width of the abstention band around 0.5.

Source code in helioai/core/judgment.py
@dataclass(frozen=True)
class Noul:
    """A yes/no question.

    The judge returns a probability; the answer is `True` at or above `0.5 + margin`,
    `False` at or below `0.5 - margin`, and `None` in between. The default margin makes
    0.7 the least confident "yes" — a site that needs more certainty raises it.

    Attributes:
        instructions: The question, in plain English, as the judge reads it.
        margin: Half-width of the abstention band around 0.5.
    """

    instructions: str
    margin: float = 0.2

Choice dataclass

One option out of a closed set.

The judge returns the chosen option and a confidence; below floor the answer is None. The options are the vocabulary the site already speaks — for a measurement type, the index's own measurement_type values — so an answer is usable as an exact filter, never as prose to interpret.

Attributes:

Name Type Description
instructions str

The question, in plain English.

options tuple[str, ...]

The closed set, in the order the site wants them presented.

floor float

Minimum confidence for the choice to count as an answer.

Source code in helioai/core/judgment.py
@dataclass(frozen=True)
class Choice:
    """One option out of a closed set.

    The judge returns the chosen option and a confidence; below `floor` the answer is
    `None`. The options are the vocabulary the site already speaks — for a measurement
    type, the index's own `measurement_type` values — so an answer is usable as an exact
    filter, never as prose to interpret.

    Attributes:
        instructions: The question, in plain English.
        options: The closed set, in the order the site wants them presented.
        floor: Minimum confidence for the choice to count as an answer.
    """

    instructions: str
    options: tuple[str, ...]
    floor: float = 0.5

Answers dataclass

What the judge said for one call, already decided question by question.

values holds bool for a Noul, str for a Choice, None where the answer did not clear its threshold; raw keeps the probabilities so the record can be re-judged at another threshold later without another call.

Attributes:

Name Type Description
values Mapping[str, bool | str | None]

Decided answers, keyed like the questions.

raw Mapping[str, Any]

The judge's probabilities, keyed like the questions.

model str

The judge model that answered.

latency_ms float

Wall time of the call as this process saw it.

request_id str | None

The provider's id for the call, for support and audit.

Source code in helioai/core/judgment.py
@dataclass(frozen=True)
class Answers:
    """What the judge said for one call, already decided question by question.

    `values` holds `bool` for a `Noul`, `str` for a `Choice`, `None` where the answer did
    not clear its threshold; `raw` keeps the probabilities so the record can be re-judged
    at another threshold later without another call.

    Attributes:
        values: Decided answers, keyed like the questions.
        raw: The judge's probabilities, keyed like the questions.
        model: The judge model that answered.
        latency_ms: Wall time of the call as this process saw it.
        request_id: The provider's id for the call, for support and audit.
    """

    values: Mapping[str, bool | str | None]
    raw: Mapping[str, Any]
    model: str
    latency_ms: float
    request_id: str | None = None

    def __getitem__(self, name: str) -> bool | str | None:
        return self.values[name]

    def get(self, name: str, default: Any = None) -> Any:
        """The decided answer for `name`, or `default` when absent or abstained."""
        value = self.values.get(name)
        return default if value is None else value

get

get(name: str, default: Any = None) -> Any

The decided answer for name, or default when absent or abstained.

Source code in helioai/core/judgment.py
def get(self, name: str, default: Any = None) -> Any:
    """The decided answer for `name`, or `default` when absent or abstained."""
    value = self.values.get(name)
    return default if value is None else value

enabled

enabled(site: str) -> bool

Whether site may ask at all: a judging backend AND the site's experiment name.

Two axes on purpose. With the backend alone, "Jev does not help at this site" and "the layer costs something" could not be told apart; with the experiment alone, a site could be on with nobody to answer.

Source code in helioai/core/judgment.py
def enabled(site: str) -> bool:
    """Whether `site` may ask at all: a judging backend AND the site's experiment name.

    Two axes on purpose. With the backend alone, "Jev does not help at this site" and "the
    layer costs something" could not be told apart; with the experiment alone, a site
    could be on with nobody to answer.
    """
    return settings.judgment.backend != "null" and f"judgment_{site}" in settings.agent.experiments

ask async

ask(site: str, state: Mapping[str, Any], questions: Mapping[str, Question], *, decided: Any = None) -> Answers | None

Ask the judge; None means abstain, and the caller runs today's code.

Parameters:

Name Type Description Default
site str

The call site's name; judgment_<site> is its experiment.

required
state Mapping[str, Any]

What the judge reads — JSON-serialisable, and recorded in full.

required
questions Mapping[str, Question]

The questions, keyed by the names the site will read back.

required
decided Any

What the deterministic path decided, when the site knows it before asking; recorded beside the answer so the two can be compared offline.

None

Returns:

Type Description
Answers | None

The decided answers, or None when the judge abstained as a whole.

Source code in helioai/core/judgment.py
async def ask(
    site: str,
    state: Mapping[str, Any],
    questions: Mapping[str, Question],
    *,
    decided: Any = None,
) -> Answers | None:
    """Ask the judge; `None` means abstain, and the caller runs today's code.

    Args:
        site: The call site's name; `judgment_<site>` is its experiment.
        state: What the judge reads — JSON-serialisable, and recorded in full.
        questions: The questions, keyed by the names the site will read back.
        decided: What the deterministic path decided, when the site knows it before
            asking; recorded beside the answer so the two can be compared offline.

    Returns:
        The decided answers, or `None` when the judge abstained as a whole.
    """
    if not enabled(site):
        return None
    t0 = time.perf_counter()
    try:
        response = await asyncio.wait_for(
            _backend().ask(dict(state), dict(questions)), timeout=settings.judgment.timeout_s
        )
    except Exception as e:
        latency = (time.perf_counter() - t0) * 1000
        _warn_once(site, e)
        _record(site, state, questions, None, decided, latency, error=f"{type(e).__name__}: {e}")
        return None
    latency = (time.perf_counter() - t0) * 1000
    answers = Answers(
        values=_decide(questions, response.answers),
        raw=_raw(questions, response.answers),
        model=response.model,
        latency_ms=latency,
        request_id=response.request_id,
    )
    _record(site, state, questions, answers, decided, latency)
    return answers

batch async

batch(site: str, states: list[Mapping[str, Any]], questions: Mapping[str, Question], *, concurrency: int = 8, record_to: Path | None = None, keys: Sequence[str] | None = None) -> list[Answers | None]

Ask the same questions of many states, off the agent loop — for a job, not a turn.

The runtime's ask is gated by a site's experiment name because it runs inside a conversation nobody asked to be judged. A job such as helioai index --classify is an explicit request: it needs only a judging backend, and it records every call to a file of its own (record_to) rather than to a session that does not exist. The same thresholds decide the answers, so what a job writes and what a turn would read agree.

Parameters:

Name Type Description Default
site str

The job's name, in the records.

required
states list[Mapping[str, Any]]

One state per item, JSON-serialisable.

required
questions Mapping[str, Question]

The questions, asked of every state.

required
concurrency int

How many calls in flight at once (37 ms per item at eight, measured).

8
record_to Path | None

The JSON-lines file every call is appended to; None records nothing.

None
keys Sequence[str] | None

One name per state, written into its record as key — the product id, so a record can be found again without matching its text. The 82 266 records of the first indexing pass carried none and had to be joined back by text.

None

Returns:

Type Description
list[Answers | None]

One Answers or None (abstained, failed, timed out) per state, in order.

Source code in helioai/core/judgment.py
async def batch(
    site: str,
    states: list[Mapping[str, Any]],
    questions: Mapping[str, Question],
    *,
    concurrency: int = 8,
    record_to: Path | None = None,
    keys: Sequence[str] | None = None,
) -> list[Answers | None]:
    """Ask the same questions of many states, off the agent loop — for a job, not a turn.

    The runtime's `ask` is gated by a site's experiment name because it runs inside a
    conversation nobody asked to be judged. A job such as `helioai index --classify` is an
    explicit request: it needs only a judging backend, and it records every call to a file
    of its own (`record_to`) rather than to a session that does not exist. The same
    thresholds decide the answers, so what a job writes and what a turn would read agree.

    Args:
        site: The job's name, in the records.
        states: One state per item, JSON-serialisable.
        questions: The questions, asked of every state.
        concurrency: How many calls in flight at once (37 ms per item at eight, measured).
        record_to: The JSON-lines file every call is appended to; None records nothing.
        keys: One name per state, written into its record as `key` — the product id, so
            a record can be found again without matching its text. The 82 266 records of
            the first indexing pass carried none and had to be joined back by text.

    Returns:
        One `Answers` or `None` (abstained, failed, timed out) per state, in order.
    """
    if settings.judgment.backend == "null" or not states:
        return [None] * len(states)
    if keys is not None and len(keys) != len(states):
        raise ValueError("one key per state")
    sem = asyncio.Semaphore(max(1, concurrency))

    async def one(state: Mapping[str, Any], key: str | None) -> Answers | None:
        async with sem:
            t0 = time.perf_counter()
            try:
                response = await asyncio.wait_for(
                    _backend().ask(dict(state), dict(questions)),
                    timeout=settings.judgment.timeout_s,
                )
            except Exception as e:
                latency = (time.perf_counter() - t0) * 1000
                _warn_once(site, e)
                _record(
                    site,
                    state,
                    questions,
                    None,
                    None,
                    latency,
                    error=str(e),
                    path=record_to,
                    key=key,
                )
                return None
            latency = (time.perf_counter() - t0) * 1000
            answers = Answers(
                values=_decide(questions, response.answers),
                raw=_raw(questions, response.answers),
                model=response.model,
                latency_ms=latency,
                request_id=response.request_id,
            )
            _record(site, state, questions, answers, None, latency, path=record_to, key=key)
            return answers

    return list(
        await asyncio.gather(
            *(one(st, keys[i] if keys is not None else None) for i, st in enumerate(states))
        )
    )

intent_contract async

intent_contract(query: str) -> Answers | None

The contract for one question, or None — see INTENT_QUESTIONS.

Parameters:

Name Type Description Default
query str

The user's question, verbatim.

required
Source code in helioai/core/judgment.py
async def intent_contract(query: str) -> Answers | None:
    """The contract for one question, or None — see `INTENT_QUESTIONS`.

    Args:
        query: The user's question, verbatim.
    """
    if not query or not query.strip():
        return None
    return await ask("intent", {"question": query}, INTENT_QUESTIONS)

contract_fields

contract_fields(answers: Answers) -> dict[str, Any]

The intent event's payload: the decided answers, the date assembled by the code.

A Choice of none, or an abstention, is None in the payload — the reader must not mistake "the judge did not say" for "the request named nothing". date is the ISO day, month or year the components make, with date_precision saying which; deliverable is the kinds wanted joined with + ("value+figure"), explanation when all four were decided no, None when any of them was left undecided and none was yes.

Source code in helioai/core/judgment.py
def contract_fields(answers: Answers) -> dict[str, Any]:
    """The `intent` event's payload: the decided answers, the date assembled by the code.

    A Choice of `none`, or an abstention, is `None` in the payload — the reader must not
    mistake "the judge did not say" for "the request named nothing". `date` is the ISO day,
    month or year the components make, with `date_precision` saying which; `deliverable`
    is the kinds wanted joined with `+` ("value+figure"), `explanation` when all four were
    decided no, `None` when any of them was left undecided and none was yes.
    """
    out: dict[str, Any] = {}
    for name in INTENT_QUESTIONS:
        value = answers.get(name)
        out[name] = None if value in (None, NONE) else value
    wanted = [kind for kind in DELIVERABLES if out.get(f"wants_{kind}") is True]
    if wanted:
        out["deliverable"] = "+".join(wanted)
    elif all(out.get(f"wants_{kind}") is False for kind in DELIVERABLES):
        out["deliverable"] = "explanation"
    else:
        out["deliverable"] = None
    year, month, day = out.pop("year"), out.pop("month"), out.pop("day")
    if year and month and day:
        out["date"], out["date_precision"] = f"{year}-{month}-{day}", "day"
    elif year and month:
        out["date"], out["date_precision"] = f"{year}-{month}", "month"
    elif year:
        out["date"], out["date_precision"] = year, "year"
    else:
        out["date"], out["date_precision"] = None, None
    out["model"] = answers.model
    out["latency_ms"] = round(answers.latency_ms, 1)
    return out

collect async

collect(task: Task) -> dict[str, Any] | None

Read a contract task the turn started; abstain if it has not answered in time.

The task ran concurrently with the model; by the time the answer is out it has almost always finished. A turn never waits on the judge longer than one judgment budget.

Source code in helioai/core/judgment.py
async def collect(task: asyncio.Task) -> dict[str, Any] | None:
    """Read a contract task the turn started; abstain if it has not answered in time.

    The task ran concurrently with the model; by the time the answer is out it has almost
    always finished. A turn never waits on the judge longer than one judgment budget.
    """
    try:
        answers = await asyncio.wait_for(task, timeout=settings.judgment.timeout_s)
    except Exception as e:
        _warn_once("intent", e)
        return None
    return contract_fields(answers) if answers is not None else None

to_sdk

to_sdk(questions: Mapping[str, Question]) -> dict[str, Any]

Our question shapes as typesafe_sdk ones — the only place the SDK's names appear.

Options of a Choice become criteria without descriptions: the vocabulary is the site's own (an index field, a closed list of regions), and describing it in prose would put a second, unmeasured question in front of the judge.

Source code in helioai/core/judgment.py
def to_sdk(questions: Mapping[str, Question]) -> dict[str, Any]:
    """Our question shapes as `typesafe_sdk` ones — the only place the SDK's names appear.

    Options of a `Choice` become criteria without descriptions: the vocabulary is the
    site's own (an index field, a closed list of regions), and describing it in prose
    would put a second, unmeasured question in front of the judge.
    """
    from typesafe_sdk import Choice as SdkChoice
    from typesafe_sdk import Noul as SdkNoul

    out: dict[str, Any] = {}
    for name, q in questions.items():
        if isinstance(q, Noul):
            out[name] = SdkNoul(instructions=q.instructions)
        else:
            out[name] = SdkChoice(
                instructions=q.instructions, criteria=dict.fromkeys(q.options, None)
            )
    return out

aclose async

aclose() -> None

Close the judge's connection pool at shutdown; safe when it was never opened.

Source code in helioai/core/judgment.py
async def aclose() -> None:
    """Close the judge's connection pool at shutdown; safe when it was never opened."""
    global _jev
    if _jev is not None and _jev._client is not None:
        from helioai.core.llm.base import close_sdk_client

        await close_sdk_client(_jev._client)
    _jev = None

helioai.core.joins

What the turn did about what was asked — four joins, no model, observation only.

The intent contract (judgment.intent_contract) says what a question committed the answer to. This module places that contract against what the turn actually produced — the parameter cards of every product loaded, the claims of final_answer, the figures — and records the comparison on the intent event. Each join is an exact operation on typed fields: a set membership, an interval intersection, a count. None calls a model, none reads prose, and none changes the answer; the intent event is what a reader of the journal, or a later correction path behind its own experiment, will consult.

The four joins are the failures the benches actually produced, one each:

  • frame — the frame the question named against the coord_sys of the cards the sandbox filled from the archive's metadata: "GSM asked, BGSE plotted" was caught by nothing.
  • window — the date the question named against the bounds the downloads obtained, and the cards whose series stopped short of the window asked for (coverage_note). The 2026-09-18 θ_Bn of 12° for a 54° shock had a window that ended at the shock; every check downstream was green because none looked at the bounds.
  • quantity — the measurement type the question named against the indexed type of the products loaded: twelve searches and a confident answer built on no product that measured the quantity (2026-09-15) is a count, once the field is filled.
  • responsiveness — what the question required against what the answer delivered: an uncertainty asked and no claim about a spread, a figure asked and none produced.

A join that has nothing to compare — no frame named, no card with a bound, no type on any loaded product — reports None, never a verdict: "the join could not say" and "the turn got it right" are kept apart, as they are in the contract itself.

checks

checks(contract: dict, artifacts: list[dict], claims: list[dict]) -> dict[str, Any]

The four joins for one turn, as the checks field of the intent event.

Parameters:

Name Type Description Default
contract dict

judgment.contract_fields output — the decided intent, None where the judge abstained or the request named nothing.

required
artifacts list[dict]

RunEnd.artifacts — every parameter card, figure, export and catalogue preview the lead produced — plus the artifacts its sub-agents streamed. RunEnd.artifacts alone never held the latter: six live turns with the judge on read "figure missing" on three of the four figures a sub-agent drew, and joined no frame or window at all. validate() keeps the lead's list by design — its recipe check reads exports against the lead's own history, and a delegated run_recipe would read as bypassed.

required
claims list[dict]

RunEnd.claims — the numbers the answer named through final_answer.

required
Source code in helioai/core/joins.py
def checks(contract: dict, artifacts: list[dict], claims: list[dict]) -> dict[str, Any]:
    """The four joins for one turn, as the `checks` field of the `intent` event.

    Args:
        contract: `judgment.contract_fields` output — the decided intent, `None` where
            the judge abstained or the request named nothing.
        artifacts: `RunEnd.artifacts` — every parameter card, figure, export and
            catalogue preview the lead produced — plus the artifacts its sub-agents
            streamed. `RunEnd.artifacts` alone never held the latter: six live turns
            with the judge on read "figure missing" on three of the four figures a
            sub-agent drew, and joined no frame or window at all. `validate()` keeps
            the lead's list by design — its recipe check reads exports against the
            lead's own history, and a delegated `run_recipe` would read as bypassed.
        claims: `RunEnd.claims` — the numbers the answer named through `final_answer`.
    """
    cards = [a for a in artifacts if a.get("kind") == "parameter_card"]
    return {
        "frame": _frame(contract.get("frame"), cards),
        "window": _window(contract.get("date"), contract.get("date_precision"), cards),
        "quantity": _quantity(contract.get("quantity"), cards),
        "responsiveness": _responsiveness(contract, artifacts, claims),
    }

summary

summary(check: dict[str, Any] | None) -> list[str]

The mismatches, as short phrases for a renderer; empty when every join agreed or had nothing to compare.

Source code in helioai/core/joins.py
def summary(check: dict[str, Any] | None) -> list[str]:
    """The mismatches, as short phrases for a renderer; empty when every join agreed or
    had nothing to compare."""
    if not check:
        return []
    out: list[str] = []
    frame = check.get("frame")
    if frame and frame.get("match") is False:
        out.append(f"frame {frame['asked']} asked, {'/'.join(frame['loaded'])} loaded")
    window = check.get("window")
    if window:
        if window.get("covered") is False:
            out.append(f"{window['date']} not in any loaded series")
        if window.get("short"):
            out.append(f"{len(window['short'])} series short of the window asked")
    quantity = check.get("quantity")
    if quantity and quantity.get("match") is False:
        loaded = ", ".join(sorted(quantity["loaded"]))
        out.append(f"{quantity['asked']} asked, loaded {loaded}")
    resp = check.get("responsiveness") or {}
    unc = resp.get("uncertainty")
    if unc and unc.get("required") and not unc.get("claimed"):
        out.append("uncertainty asked, none claimed")
    deliv = resp.get("deliverable")
    for kind in (deliv or {}).get("missing") or []:
        out.append(f"{kind} asked, none produced")
    return out

Figure review

helioai.core.vision

Stateless vision side-call: review sandbox figures after run_python.

The image is sent once, outside the conversation; only the short text verdict enters the tool result (and thus the history), so the cost stays a few hundred tokens per figure instead of re-sending images every turn. Never blocks the loop: any failure logs a warning and returns the result unchanged.

maybe_review async

maybe_review(tool_name: str, result: ToolResult) -> tuple[ToolResult, str | None]

Attach a vision verdict to a run_python result carrying figures.

No-op unless HELIOAI_VISION_ENABLED is set and the tool is run_python.

Parameters:

Name Type Description Default
tool_name str

Name of the tool that just ran.

required
result ToolResult

Its result; the payload is read for figure_paths.

required

Returns:

Type Description
ToolResult

(possibly amended result, verdict text or None) — the verdict is a

str | None

stateless side-call; only its text enters the history, never the image.

Source code in helioai/core/vision.py
async def maybe_review(tool_name: str, result: ToolResult) -> tuple[ToolResult, str | None]:
    """Attach a vision verdict to a run_python result carrying figures.

    No-op unless `HELIOAI_VISION_ENABLED` is set and the tool is `run_python`.

    Args:
        tool_name: Name of the tool that just ran.
        result: Its result; the payload is read for `figure_paths`.

    Returns:
        (possibly amended result, verdict text or None) — the verdict is a
        stateless side-call; only its text enters the history, never the image.
    """
    if not settings.vision.enabled or tool_name != "run_python":
        return result, None
    data = result.payload
    if not isinstance(data, dict) or data.get("error") or not data.get("figure_paths"):
        return result, None
    verdict = await _review(data["figure_paths"])
    if not verdict:
        return result, None
    return result.with_payload({**data, "figure_review": verdict}), verdict

LLM clients

helioai.core.llm.base

Provider-neutral message model and LLMClient interface.

ToolCall dataclass

A tool invocation requested by the model.

Attributes:

Name Type Description
id str

Provider-assigned identifier, echoed back on the matching tool result. Gemini has no native ids, so its client synthesises name::hex and parses the name back out.

name str

Registered tool name.

arguments dict

Decoded JSON arguments. Empty when the model emitted malformed JSON — a bad tool call must not kill the loop.

Source code in helioai/core/llm/base.py
@dataclass
class ToolCall:
    """A tool invocation requested by the model.

    Attributes:
        id: Provider-assigned identifier, echoed back on the matching tool
            result. Gemini has no native ids, so its client synthesises
            `name::hex` and parses the name back out.
        name: Registered tool name.
        arguments: Decoded JSON arguments. Empty when the model emitted
            malformed JSON — a bad tool call must not kill the loop.
    """

    id: str
    name: str
    arguments: dict

Message dataclass

One turn of conversation, in a provider-neutral form.

Every client converts to and from this shape, so the agent loop, the session store and the interfaces never see a provider's wire format.

Attributes:

Name Type Description
role Literal['system', 'user', 'assistant', 'tool']

Who produced the turn.

content str

Text content. Present alongside tool_calls when the model narrated what it was about to do.

tool_calls list[ToolCall] | None

Tools the assistant wants invoked, when it requested any.

tool_call_id str | None

For tool messages, the ToolCall.id being answered.

prompt_tokens int

Input tokens the provider billed for this reply, 0 when it reported none.

completion_tokens int

Output tokens the provider billed for this reply.

cached_tokens int

The subset of prompt_tokens served from the provider's prompt cache. Priced differently by every provider, so it is kept apart from the total rather than folded into it.

origin str | None

Who really wrote a user message when it was not the person. None for the human and the model; "correction" for the automated note the loop injects when an answer quotes ids the catalogue does not have. The provider clients only forward user and assistant turns from history, so a synthetic message must keep the user role to be seen at all — this field is what lets the replay and the notebook export tell it from a question. Persisted by SessionStore.

name str | None

For tool messages, the tool that produced the result. The wire formats carry only the call id, so every reader that needed to know which tool a result came from sniffed the JSON's shape — load_recipe was recognised by having name, code and metadata keys at once. Recorded here instead, and persisted, so a reader asks the message. Not sent to any provider.

reasoning str | None

For assistant messages, the chain of thought the provider returned beside the content (reasoning_content on DeepSeek's thinking mode). Never shown, but sent back: DeepSeek rejects with a 400 any request carrying tools whose history lacks an earlier turn's reasoning, and enforces it intermittently, so dropping it cost a lead mid-question on 2 of ~51 calls. Persisted by SessionStore for the same reason — a reloaded session is a history too. None when the provider returned none, which keeps it off the wire for every provider that never emits it.

Source code in helioai/core/llm/base.py
@dataclass
class Message:
    """One turn of conversation, in a provider-neutral form.

    Every client converts to and from this shape, so the agent loop, the session
    store and the interfaces never see a provider's wire format.

    Attributes:
        role: Who produced the turn.
        content: Text content. Present alongside `tool_calls` when the model
            narrated what it was about to do.
        tool_calls: Tools the assistant wants invoked, when it requested any.
        tool_call_id: For `tool` messages, the `ToolCall.id` being answered.
        prompt_tokens: Input tokens the provider billed for this reply, 0 when it
            reported none.
        completion_tokens: Output tokens the provider billed for this reply.
        cached_tokens: The subset of `prompt_tokens` served from the provider's
            prompt cache. Priced differently by every provider, so it is kept
            apart from the total rather than folded into it.
        origin: Who really wrote a `user` message when it was not the person.
            `None` for the human and the model; `"correction"` for the automated
            note the loop injects when an answer quotes ids the catalogue does
            not have. The provider clients only forward `user` and `assistant`
            turns from history, so a synthetic message must keep the user role to
            be seen at all — this field is what lets the replay and the notebook
            export tell it from a question. Persisted by `SessionStore`.
        name: For `tool` messages, the tool that produced the result. The wire
            formats carry only the call id, so every reader that needed to know
            which tool a result came from sniffed the JSON's shape — `load_recipe`
            was recognised by having `name`, `code` and `metadata` keys at once.
            Recorded here instead, and persisted, so a reader asks the message.
            Not sent to any provider.
        reasoning: For `assistant` messages, the chain of thought the provider
            returned beside the content (`reasoning_content` on DeepSeek's thinking
            mode). Never shown, but sent back: DeepSeek rejects with a 400 any request
            carrying `tools` whose history lacks an earlier turn's reasoning, and
            enforces it intermittently, so dropping it cost a lead mid-question on 2
            of ~51 calls. Persisted by `SessionStore` for the same reason — a reloaded
            session is a history too. `None` when the provider returned none, which
            keeps it off the wire for every provider that never emits it.
    """

    role: Literal["system", "user", "assistant", "tool"]
    content: str = ""
    tool_calls: list[ToolCall] | None = None
    tool_call_id: str | None = None
    origin: str | None = None
    name: str | None = None
    reasoning: str | None = None
    # Telemetry about the response, not part of the conversation: `SessionStore`
    # deliberately does not persist these, so a reloaded history reports no cost.
    prompt_tokens: int = 0
    completion_tokens: int = 0
    cached_tokens: int = 0

ToolDef dataclass

A tool as advertised to the model.

Attributes:

Name Type Description
name str

Tool name the model will call.

description str

What the tool does and when to reach for it — the model's only clue about applicability.

parameters dict

JSON Schema object describing the accepted arguments.

Source code in helioai/core/llm/base.py
@dataclass
class ToolDef:
    """A tool as advertised to the model.

    Attributes:
        name: Tool name the model will call.
        description: What the tool does and when to reach for it — the model's
            only clue about applicability.
        parameters: JSON Schema object describing the accepted arguments.
    """

    name: str
    description: str
    parameters: dict = field(default_factory=dict)

LLMClient

Bases: ABC

Interface every provider client implements.

One method is required, deliberately: the agent loop only ever needs a single completion with optional tool calling. Streaming happens at the loop level.

Source code in helioai/core/llm/base.py
class LLMClient(ABC):
    """Interface every provider client implements.

    One method is required, deliberately: the agent loop only ever needs a single
    completion with optional tool calling. Streaming happens at the loop level.
    """

    async def aclose(self) -> None:
        """Release the underlying HTTP connection pool.

        Callers that build a client per request — the CLI, the Jupyter magic and
        the web endpoints all do — must await this before their event loop ends.
        An async pool binds to the loop that used it, so a client left to the
        garbage collector schedules its own teardown after `asyncio.run` has
        closed that loop, and asyncio reports an unretrieved
        `RuntimeError: Event loop is closed` while the sockets stay open.

        The default is a no-op so a client without a pool needs no override.
        """
        return None

    @abstractmethod
    async def chat(
        self,
        messages: list[Message],
        tools: list[ToolDef],
        system_prompt: str | None = None,
        tool_choice: str = "auto",
    ) -> Message:
        """Send one turn and return the assistant's reply.

        Args:
            messages: Conversation history.
            tools: Tools the model may call this turn.
            system_prompt: Instructions placed before the history.
            tool_choice: `auto` to let the model decide, `required` to force a
                tool call — used on a sub-agent's first turn so it cannot answer
                from memory without looking anything up.

        Returns:
            The assistant reply, carrying `tool_calls` when the model requested any.
        """
        raise NotImplementedError

    async def stream_chat(
        self,
        messages: list[Message],
        tools: list[ToolDef],
        system_prompt: str | None = None,
        tool_choice: str = "auto",
    ) -> AsyncIterator[str | Message]:
        """Send one turn and yield the reply's text as it is generated, then the reply.

        The default is `chat()` in one piece: a client with no streaming support yields
        the finished `Message` and nothing before it, so a caller that streams works
        unchanged against every provider. Clients that stream yield text deltas (`str`)
        and end with the same `Message` `chat()` would have returned — same content,
        same tool calls, same usage.

        Args:
            messages: Conversation history.
            tools: Tools the model may call this turn.
            system_prompt: Instructions placed before the history.
            tool_choice: As for `chat`. Passed on only when it is not the default, so
                a client whose `chat` predates the argument still works.

        Yields:
            Text deltas, then the final `Message`.
        """
        kwargs: dict[str, Any] = {"system_prompt": system_prompt}
        if tool_choice != "auto":
            kwargs["tool_choice"] = tool_choice
        yield await self.chat(messages, tools, **kwargs)

aclose async

aclose() -> None

Release the underlying HTTP connection pool.

Callers that build a client per request — the CLI, the Jupyter magic and the web endpoints all do — must await this before their event loop ends. An async pool binds to the loop that used it, so a client left to the garbage collector schedules its own teardown after asyncio.run has closed that loop, and asyncio reports an unretrieved RuntimeError: Event loop is closed while the sockets stay open.

The default is a no-op so a client without a pool needs no override.

Source code in helioai/core/llm/base.py
async def aclose(self) -> None:
    """Release the underlying HTTP connection pool.

    Callers that build a client per request — the CLI, the Jupyter magic and
    the web endpoints all do — must await this before their event loop ends.
    An async pool binds to the loop that used it, so a client left to the
    garbage collector schedules its own teardown after `asyncio.run` has
    closed that loop, and asyncio reports an unretrieved
    `RuntimeError: Event loop is closed` while the sockets stay open.

    The default is a no-op so a client without a pool needs no override.
    """
    return None

chat abstractmethod async

chat(messages: list[Message], tools: list[ToolDef], system_prompt: str | None = None, tool_choice: str = 'auto') -> Message

Send one turn and return the assistant's reply.

Parameters:

Name Type Description Default
messages list[Message]

Conversation history.

required
tools list[ToolDef]

Tools the model may call this turn.

required
system_prompt str | None

Instructions placed before the history.

None
tool_choice str

auto to let the model decide, required to force a tool call — used on a sub-agent's first turn so it cannot answer from memory without looking anything up.

'auto'

Returns:

Type Description
Message

The assistant reply, carrying tool_calls when the model requested any.

Source code in helioai/core/llm/base.py
@abstractmethod
async def chat(
    self,
    messages: list[Message],
    tools: list[ToolDef],
    system_prompt: str | None = None,
    tool_choice: str = "auto",
) -> Message:
    """Send one turn and return the assistant's reply.

    Args:
        messages: Conversation history.
        tools: Tools the model may call this turn.
        system_prompt: Instructions placed before the history.
        tool_choice: `auto` to let the model decide, `required` to force a
            tool call — used on a sub-agent's first turn so it cannot answer
            from memory without looking anything up.

    Returns:
        The assistant reply, carrying `tool_calls` when the model requested any.
    """
    raise NotImplementedError

stream_chat async

stream_chat(messages: list[Message], tools: list[ToolDef], system_prompt: str | None = None, tool_choice: str = 'auto') -> AsyncIterator[str | Message]

Send one turn and yield the reply's text as it is generated, then the reply.

The default is chat() in one piece: a client with no streaming support yields the finished Message and nothing before it, so a caller that streams works unchanged against every provider. Clients that stream yield text deltas (str) and end with the same Message chat() would have returned — same content, same tool calls, same usage.

Parameters:

Name Type Description Default
messages list[Message]

Conversation history.

required
tools list[ToolDef]

Tools the model may call this turn.

required
system_prompt str | None

Instructions placed before the history.

None
tool_choice str

As for chat. Passed on only when it is not the default, so a client whose chat predates the argument still works.

'auto'

Yields:

Type Description
AsyncIterator[str | Message]

Text deltas, then the final Message.

Source code in helioai/core/llm/base.py
async def stream_chat(
    self,
    messages: list[Message],
    tools: list[ToolDef],
    system_prompt: str | None = None,
    tool_choice: str = "auto",
) -> AsyncIterator[str | Message]:
    """Send one turn and yield the reply's text as it is generated, then the reply.

    The default is `chat()` in one piece: a client with no streaming support yields
    the finished `Message` and nothing before it, so a caller that streams works
    unchanged against every provider. Clients that stream yield text deltas (`str`)
    and end with the same `Message` `chat()` would have returned — same content,
    same tool calls, same usage.

    Args:
        messages: Conversation history.
        tools: Tools the model may call this turn.
        system_prompt: Instructions placed before the history.
        tool_choice: As for `chat`. Passed on only when it is not the default, so
            a client whose `chat` predates the argument still works.

    Yields:
        Text deltas, then the final `Message`.
    """
    kwargs: dict[str, Any] = {"system_prompt": system_prompt}
    if tool_choice != "auto":
        kwargs["tool_choice"] = tool_choice
    yield await self.chat(messages, tools, **kwargs)

call_with_retry async

call_with_retry(fn: Callable[[], Awaitable[Any]], *, attempts: int = 4, base_delay: float = 1.0, max_delay: float = 60.0) -> Any

Call async fn with backoff on retryable HTTP errors.

A server-supplied Retry-After wins over the exponential backoff: waiting the window a rate limiter asks for costs less than losing a session that is already several minutes and several downloads deep.

But only when the wait is one we would actually sit through. A per-minute limiter asks for seconds; an exhausted daily quota answers Retry-After: 35464 — nearly ten hours. Capping that to max_delay and retrying anyway just buys four minutes of silence before the same failure, so a window longer than max_delay fails at once and says how long it really is.

Non-retryable errors and exhausted attempts are re-raised immediately.

Parameters:

Name Type Description Default
fn Callable[[], Awaitable[Any]]

Zero-argument coroutine function performing one attempt. Taking a callable rather than a coroutine is what lets it be re-invoked.

required
attempts int

Total tries, including the first.

4
base_delay float

First backoff, doubled per attempt.

1.0
max_delay float

Backoff ceiling, and the threshold above which a server's own Retry-After is refused rather than capped.

60.0

Returns:

Type Description
Any

Whatever fn returns on its first success.

Raises:

Type Description
Exception

The last error, re-raised once attempts run out or when the error is not retryable.

Source code in helioai/core/llm/base.py
async def call_with_retry(
    fn: Callable[[], Awaitable[Any]],
    *,
    attempts: int = 4,
    base_delay: float = 1.0,
    max_delay: float = 60.0,
) -> Any:
    """Call async fn with backoff on retryable HTTP errors.

    A server-supplied Retry-After wins over the exponential backoff: waiting the
    window a rate limiter asks for costs less than losing a session that is already
    several minutes and several downloads deep.

    But only when the wait is one we would actually sit through. A per-minute limiter
    asks for seconds; an exhausted daily quota answers `Retry-After: 35464` — nearly ten
    hours. Capping that to `max_delay` and retrying anyway just buys four minutes of
    silence before the same failure, so a window longer than `max_delay` fails at once
    and says how long it really is.

    Non-retryable errors and exhausted attempts are re-raised immediately.

    Args:
        fn: Zero-argument coroutine function performing one attempt. Taking a
            callable rather than a coroutine is what lets it be re-invoked.
        attempts: Total tries, including the first.
        base_delay: First backoff, doubled per attempt.
        max_delay: Backoff ceiling, and the threshold above which a server's
            own `Retry-After` is refused rather than capped.

    Returns:
        Whatever `fn` returns on its first success.

    Raises:
        Exception: The last error, re-raised once attempts run out or when the
            error is not retryable.
    """
    for attempt in range(attempts):
        try:
            return await fn()
        except Exception as exc:
            status = _error_status(exc)
            if status not in RETRYABLE_STATUS:
                raise
            hinted = _retry_after(exc)
            if hinted is not None and hinted > max_delay:
                log.warning(
                    "llm_retry_window_too_long",
                    extra={"status": status, "retry_after_s": hinted},
                )
                raise
            if attempt == attempts - 1:
                raise
            delay = min(
                hinted
                if hinted is not None
                else base_delay * (2**attempt) + random.uniform(0, 0.5),
                max_delay,
            )
            log.info(
                "llm_retry",
                extra={
                    "attempt": attempt + 1,
                    "status": status,
                    "delay": round(delay, 2),
                    "server_hinted": hinted is not None,
                },
            )
            await asyncio.sleep(delay)

close_sdk_client async

close_sdk_client(client: Any) -> None

Close an SDK client's connection pool, whether its close() is sync or async.

openai.AsyncOpenAI.close is a coroutine; google.genai.Client.close is not. Failures are swallowed: this only ever runs while tearing down, and a pool that will not close is not worth crashing a finished analysis over.

Parameters:

Name Type Description Default
client Any

Any SDK client. One exposing neither close nor aclose is accepted and ignored, so a stub in a test needs no teardown method.

required
Source code in helioai/core/llm/base.py
async def close_sdk_client(client: Any) -> None:
    """Close an SDK client's connection pool, whether its close() is sync or async.

    `openai.AsyncOpenAI.close` is a coroutine; `google.genai.Client.close` is not.
    Failures are swallowed: this only ever runs while tearing down, and a pool
    that will not close is not worth crashing a finished analysis over.

    Args:
        client: Any SDK client. One exposing neither `close` nor `aclose` is
            accepted and ignored, so a stub in a test needs no teardown method.
    """
    import inspect

    closer = getattr(client, "close", None) or getattr(client, "aclose", None)
    if closer is None:
        return
    try:
        result = closer()
        if inspect.isawaitable(result):
            await result
    except Exception as e:
        log.debug("sdk_client_close_failed: %s", e)

helioai.core.llm.openai_compat

Single client for every provider that speaks the OpenAI chat-completions wire format.

Groq, Ollama, Azure OpenAI and OpenAI itself all accept the same request shape, so they share one implementation here instead of one near-identical class each. A provider is a base_url plus a couple of dialect flags, not a subclass.

Azure is the one exception that still needs its own SDK client object (deployment routing and api-version), so it subclasses this to swap the constructor only.

OpenAICompatClient

Bases: LLMClient

Chat client for any OpenAI-compatible endpoint.

Parameters:

Name Type Description Default
model str

Model name sent in the request (the deployment name on Azure).

required
api_key str

Provider API key. Local endpoints such as Ollama ignore it, but the SDK requires a non-empty value.

''
base_url str | None

Endpoint root. None targets OpenAI itself.

None
system_role str

Role used for the system prompt — developer on Azure and the o-series, system everywhere else.

'system'
max_output_tokens int

Cap on generated tokens.

4096
temperature float | None

Sampling temperature. None omits the field entirely, which reasoning models require.

0.2
provider str

Name used to label log messages.

'openai'
client Any

Pre-built SDK client. Injected by tests and by subclasses.

None
Example

client = OpenAICompatClient( ... model="llama-3.3-70b-versatile", ... api_key="gsk_...", ... base_url="https://api.groq.com/openai/v1", ... ) reply = await client.chat([Message(role="user", content="hi")], tools=[])

Source code in helioai/core/llm/openai_compat.py
class OpenAICompatClient(LLMClient):
    """Chat client for any OpenAI-compatible endpoint.

    Args:
        model: Model name sent in the request (the deployment name on Azure).
        api_key: Provider API key. Local endpoints such as Ollama ignore it, but
            the SDK requires a non-empty value.
        base_url: Endpoint root. `None` targets OpenAI itself.
        system_role: Role used for the system prompt — `developer` on Azure and
            the o-series, `system` everywhere else.
        max_output_tokens: Cap on generated tokens.
        temperature: Sampling temperature. `None` omits the field entirely, which
            reasoning models require.
        provider: Name used to label log messages.
        client: Pre-built SDK client. Injected by tests and by subclasses.

    Example:
        >>> client = OpenAICompatClient(
        ...     model="llama-3.3-70b-versatile",
        ...     api_key="gsk_...",
        ...     base_url="https://api.groq.com/openai/v1",
        ... )
        >>> reply = await client.chat([Message(role="user", content="hi")], tools=[])
    """

    def __init__(
        self,
        *,
        model: str,
        api_key: str = "",
        base_url: str | None = None,
        system_role: str = "system",
        max_output_tokens: int = 4096,
        temperature: float | None = 0.2,
        provider: str = "openai",
        client: Any = None,
        default_headers: dict[str, str] | None = None,
    ):
        self._client = client or AsyncOpenAI(
            api_key=api_key or "unused",
            base_url=base_url,
            max_retries=0,
            default_headers=default_headers or None,
        )
        self._model = model
        self._system_role = system_role
        self._max_output_tokens = max_output_tokens
        self._temperature = temperature
        self._provider = provider
        self._reasoning_required = False

    async def aclose(self) -> None:
        """Close the httpx pool held by the SDK client."""
        await close_sdk_client(self._client)

    async def chat(
        self,
        messages: list[Message],
        tools: list[ToolDef],
        system_prompt: str | None = None,
        tool_choice: str = "auto",
    ) -> Message:
        """Send one chat turn and return the assistant's reply.

        Streamed underneath, with the deltas discarded. The OpenCode gateway drops
        `reasoning_content` from every non-streamed DeepSeek reply (0 of 24 turns on
        2026-09-29, 24 of 24 streamed), and DeepSeek then rejects a later request that
        does not send it back — so the sub-agents, which do not stream to anyone, died
        where the lead did not. Streaming is the one shape every provider here already
        serves to the lead, so asking it of every call costs nothing.

        Args:
            messages: Conversation history.
            tools: Tools the model may call.
            system_prompt: Instructions prepended as the first message.
            tool_choice: `auto` to let the model decide, `required` to force a call.

        Returns:
            The assistant reply, carrying `tool_calls` when the model requested any.
        """
        reply = None
        async for item in self.stream_chat(messages, tools, system_prompt, tool_choice):
            if isinstance(item, Message):
                reply = item
        assert reply is not None, "stream_chat must end with the reply"
        return reply

    async def _chat_unstreamed(
        self,
        messages: list[Message],
        tools: list[ToolDef],
        system_prompt: str | None,
        tool_choice: str,
    ) -> Message:
        kwargs = self._request(messages, tools, system_prompt, tool_choice)
        response = await self._create(kwargs)
        reply = from_openai_response(response, self._provider)
        if not (reply.content or "").strip() and not reply.tool_calls:
            log.warning("%s empty turn, retrying once: %s", self._provider, self._model)
            response = await call_with_retry(lambda: self._client.chat.completions.create(**kwargs))
            reply = from_openai_response(response, self._provider)
        return reply

    def _request(
        self,
        messages: list[Message],
        tools: list[ToolDef],
        system_prompt: str | None,
        tool_choice: str,
    ) -> dict:
        openai_messages: list[dict] = []
        if system_prompt:
            openai_messages.append({"role": self._system_role, "content": system_prompt})
        openai_messages.extend(to_openai_messages(messages))
        if self._reasoning_required:
            _carry_reasoning(openai_messages)

        kwargs: dict = {
            "model": self._model,
            "messages": openai_messages,
            "max_tokens": self._max_output_tokens,
        }
        if self._temperature is not None:
            kwargs["temperature"] = self._temperature
        if tools:
            kwargs["tools"] = to_openai_tools(tools)
            kwargs["tool_choice"] = tool_choice
        return kwargs

    async def _create(self, kwargs: dict) -> Any:
        while True:
            try:
                return await call_with_retry(lambda: self._client.chat.completions.create(**kwargs))
            except BadRequestError as e:
                # DeepSeek in thinking mode wants `reasoning_content` on every earlier
                # assistant turn of a request with tools, and an empty string satisfies
                # it. HelioAI sends back what it received, but the OpenCode gateway
                # routes a model to backends that disagree: one streams no reasoning at
                # all (0 of 24 turns in one batch on 2026-09-29, 23–24 of 24 in five
                # others), the next rejects the history that follows — a lead died on its
                # fourth call with both earlier turns empty. Asked of the provider, not
                # tabulated: the first such refusal fills the missing fields and replays,
                # and this client fills them from then on.
                if "reasoning_content" in str(e) and not self._reasoning_required:
                    log.warning("reasoning_content_required: %s", self._model)
                    self._reasoning_required = True
                    _carry_reasoning(kwargs["messages"])
                    continue
                # Forcing a tool call is a preference, never worth losing the turn over.
                # DeepSeek v4 in thinking mode rejects `required` outright ("Thinking
                # mode does not support this tool_choice"), which killed a sub-agent on
                # its very first turn. Whether a model accepts it depends on the model
                # and its reasoning mode, not on the provider, so it is asked rather
                # than tabulated.
                if kwargs.get("tool_choice") == "required" and "tool_choice" in str(e):
                    log.warning("tool_choice_rejected_falling_back_to_auto: %s", self._model)
                    kwargs["tool_choice"] = "auto"
                    continue
                raise

    async def stream_chat(
        self,
        messages: list[Message],
        tools: list[ToolDef],
        system_prompt: str | None = None,
        tool_choice: str = "auto",
    ) -> AsyncIterator[str | Message]:
        """Send one chat turn, yielding the reply's text as it arrives, then the reply.

        Text deltas are yielded as the provider sends them, with an inline
        `<think>` block held back until it closes, and a separate `reasoning_content`
        kept whole on `Message.reasoning`, never yielded. Tool-call fragments are
        reassembled by index and parsed as a finished call is; the usage comes from
        the final chunk (`stream_options.include_usage`). An endpoint that rejects the
        streaming request is asked again without it. A turn that ends with neither
        text nor a tool call is retried once, streamed again: on a reasoning model the
        whole allowance can go into hidden reasoning, and the identical request,
        replayed, came back with two tool calls — the loop above treats an empty turn
        as fatal, so without the retry the question was abandoned.

        Args:
            messages: Conversation history.
            tools: Tools the model may call.
            system_prompt: Instructions prepended as the first message.
            tool_choice: `auto` to let the model decide, `required` to force a call.

        Yields:
            Text deltas, then the final `Message`.
        """
        for attempt in range(2):
            kwargs = self._request(messages, tools, system_prompt, tool_choice)
            kwargs["stream"] = True
            kwargs["stream_options"] = {"include_usage": True}
            try:
                stream = await self._create(kwargs)
            except BadRequestError as e:
                if not _REFUSES_STREAMING.search(str(e)):
                    raise
                log.warning("%s rejected streaming, falling back to chat(): %s", self._provider, e)
                yield await self._chat_unstreamed(messages, tools, system_prompt, tool_choice)
                return

            raw = ""
            reasoning = ""
            sent = 0
            calls: dict[int, dict] = {}
            usage: dict = {"prompt_tokens": 0, "completion_tokens": 0, "cached_tokens": 0}
            finish_reason = None
            async for chunk in stream:
                if getattr(chunk, "usage", None):
                    usage = _usage(chunk)
                choices = getattr(chunk, "choices", None) or []
                if not choices:
                    continue
                choice = choices[0]
                finish_reason = getattr(choice, "finish_reason", None) or finish_reason
                delta = getattr(choice, "delta", None)
                if delta is None:
                    continue
                if getattr(delta, "reasoning_content", None):
                    reasoning += delta.reasoning_content
                if getattr(delta, "content", None):
                    raw += delta.content
                    visible = _visible_so_far(raw)
                    if len(visible) > sent:
                        yield visible[sent:]
                        sent = len(visible)
                for frag in getattr(delta, "tool_calls", None) or []:
                    slot = calls.setdefault(
                        getattr(frag, "index", 0) or 0, {"id": None, "name": None, "arguments": ""}
                    )
                    if getattr(frag, "id", None):
                        slot["id"] = frag.id
                    fn = getattr(frag, "function", None)
                    if fn is not None:
                        if getattr(fn, "name", None):
                            slot["name"] = fn.name
                        if getattr(fn, "arguments", None):
                            slot["arguments"] += fn.arguments

            content = _strip_reasoning(raw)
            tool_calls = [
                _parse_tool_call(
                    c["id"] or f"call_{i}",
                    c["name"] or "",
                    c["arguments"],
                    finish_reason,
                    self._provider,
                )
                for i, c in sorted(calls.items())
            ]
            if content.strip() or tool_calls or attempt:
                break
            log.warning("%s empty streamed turn, retrying once: %s", self._provider, self._model)
        yield Message(
            role="assistant",
            content=content,
            tool_calls=tool_calls or None,
            reasoning=reasoning or None,
            **usage,
        )

aclose async

aclose() -> None

Close the httpx pool held by the SDK client.

Source code in helioai/core/llm/openai_compat.py
async def aclose(self) -> None:
    """Close the httpx pool held by the SDK client."""
    await close_sdk_client(self._client)

chat async

chat(messages: list[Message], tools: list[ToolDef], system_prompt: str | None = None, tool_choice: str = 'auto') -> Message

Send one chat turn and return the assistant's reply.

Streamed underneath, with the deltas discarded. The OpenCode gateway drops reasoning_content from every non-streamed DeepSeek reply (0 of 24 turns on 2026-09-29, 24 of 24 streamed), and DeepSeek then rejects a later request that does not send it back — so the sub-agents, which do not stream to anyone, died where the lead did not. Streaming is the one shape every provider here already serves to the lead, so asking it of every call costs nothing.

Parameters:

Name Type Description Default
messages list[Message]

Conversation history.

required
tools list[ToolDef]

Tools the model may call.

required
system_prompt str | None

Instructions prepended as the first message.

None
tool_choice str

auto to let the model decide, required to force a call.

'auto'

Returns:

Type Description
Message

The assistant reply, carrying tool_calls when the model requested any.

Source code in helioai/core/llm/openai_compat.py
async def chat(
    self,
    messages: list[Message],
    tools: list[ToolDef],
    system_prompt: str | None = None,
    tool_choice: str = "auto",
) -> Message:
    """Send one chat turn and return the assistant's reply.

    Streamed underneath, with the deltas discarded. The OpenCode gateway drops
    `reasoning_content` from every non-streamed DeepSeek reply (0 of 24 turns on
    2026-09-29, 24 of 24 streamed), and DeepSeek then rejects a later request that
    does not send it back — so the sub-agents, which do not stream to anyone, died
    where the lead did not. Streaming is the one shape every provider here already
    serves to the lead, so asking it of every call costs nothing.

    Args:
        messages: Conversation history.
        tools: Tools the model may call.
        system_prompt: Instructions prepended as the first message.
        tool_choice: `auto` to let the model decide, `required` to force a call.

    Returns:
        The assistant reply, carrying `tool_calls` when the model requested any.
    """
    reply = None
    async for item in self.stream_chat(messages, tools, system_prompt, tool_choice):
        if isinstance(item, Message):
            reply = item
    assert reply is not None, "stream_chat must end with the reply"
    return reply

stream_chat async

stream_chat(messages: list[Message], tools: list[ToolDef], system_prompt: str | None = None, tool_choice: str = 'auto') -> AsyncIterator[str | Message]

Send one chat turn, yielding the reply's text as it arrives, then the reply.

Text deltas are yielded as the provider sends them, with an inline <think> block held back until it closes, and a separate reasoning_content kept whole on Message.reasoning, never yielded. Tool-call fragments are reassembled by index and parsed as a finished call is; the usage comes from the final chunk (stream_options.include_usage). An endpoint that rejects the streaming request is asked again without it. A turn that ends with neither text nor a tool call is retried once, streamed again: on a reasoning model the whole allowance can go into hidden reasoning, and the identical request, replayed, came back with two tool calls — the loop above treats an empty turn as fatal, so without the retry the question was abandoned.

Parameters:

Name Type Description Default
messages list[Message]

Conversation history.

required
tools list[ToolDef]

Tools the model may call.

required
system_prompt str | None

Instructions prepended as the first message.

None
tool_choice str

auto to let the model decide, required to force a call.

'auto'

Yields:

Type Description
AsyncIterator[str | Message]

Text deltas, then the final Message.

Source code in helioai/core/llm/openai_compat.py
async def stream_chat(
    self,
    messages: list[Message],
    tools: list[ToolDef],
    system_prompt: str | None = None,
    tool_choice: str = "auto",
) -> AsyncIterator[str | Message]:
    """Send one chat turn, yielding the reply's text as it arrives, then the reply.

    Text deltas are yielded as the provider sends them, with an inline
    `<think>` block held back until it closes, and a separate `reasoning_content`
    kept whole on `Message.reasoning`, never yielded. Tool-call fragments are
    reassembled by index and parsed as a finished call is; the usage comes from
    the final chunk (`stream_options.include_usage`). An endpoint that rejects the
    streaming request is asked again without it. A turn that ends with neither
    text nor a tool call is retried once, streamed again: on a reasoning model the
    whole allowance can go into hidden reasoning, and the identical request,
    replayed, came back with two tool calls — the loop above treats an empty turn
    as fatal, so without the retry the question was abandoned.

    Args:
        messages: Conversation history.
        tools: Tools the model may call.
        system_prompt: Instructions prepended as the first message.
        tool_choice: `auto` to let the model decide, `required` to force a call.

    Yields:
        Text deltas, then the final `Message`.
    """
    for attempt in range(2):
        kwargs = self._request(messages, tools, system_prompt, tool_choice)
        kwargs["stream"] = True
        kwargs["stream_options"] = {"include_usage": True}
        try:
            stream = await self._create(kwargs)
        except BadRequestError as e:
            if not _REFUSES_STREAMING.search(str(e)):
                raise
            log.warning("%s rejected streaming, falling back to chat(): %s", self._provider, e)
            yield await self._chat_unstreamed(messages, tools, system_prompt, tool_choice)
            return

        raw = ""
        reasoning = ""
        sent = 0
        calls: dict[int, dict] = {}
        usage: dict = {"prompt_tokens": 0, "completion_tokens": 0, "cached_tokens": 0}
        finish_reason = None
        async for chunk in stream:
            if getattr(chunk, "usage", None):
                usage = _usage(chunk)
            choices = getattr(chunk, "choices", None) or []
            if not choices:
                continue
            choice = choices[0]
            finish_reason = getattr(choice, "finish_reason", None) or finish_reason
            delta = getattr(choice, "delta", None)
            if delta is None:
                continue
            if getattr(delta, "reasoning_content", None):
                reasoning += delta.reasoning_content
            if getattr(delta, "content", None):
                raw += delta.content
                visible = _visible_so_far(raw)
                if len(visible) > sent:
                    yield visible[sent:]
                    sent = len(visible)
            for frag in getattr(delta, "tool_calls", None) or []:
                slot = calls.setdefault(
                    getattr(frag, "index", 0) or 0, {"id": None, "name": None, "arguments": ""}
                )
                if getattr(frag, "id", None):
                    slot["id"] = frag.id
                fn = getattr(frag, "function", None)
                if fn is not None:
                    if getattr(fn, "name", None):
                        slot["name"] = fn.name
                    if getattr(fn, "arguments", None):
                        slot["arguments"] += fn.arguments

        content = _strip_reasoning(raw)
        tool_calls = [
            _parse_tool_call(
                c["id"] or f"call_{i}",
                c["name"] or "",
                c["arguments"],
                finish_reason,
                self._provider,
            )
            for i, c in sorted(calls.items())
        ]
        if content.strip() or tool_calls or attempt:
            break
        log.warning("%s empty streamed turn, retrying once: %s", self._provider, self._model)
    yield Message(
        role="assistant",
        content=content,
        tool_calls=tool_calls or None,
        reasoning=reasoning or None,
        **usage,
    )

to_openai_messages

to_openai_messages(messages: list[Message]) -> list[dict]

Convert neutral messages to the OpenAI wire format.

System messages already in the history are dropped: the system prompt is passed separately by chat() so it always lands first.

Parameters:

Name Type Description Default
messages list[Message]

Conversation history in HelioAI's provider-neutral form.

required

Returns:

Type Description
list[dict]

Message dicts ready to send as the messages request field.

Source code in helioai/core/llm/openai_compat.py
def to_openai_messages(messages: list[Message]) -> list[dict]:
    """Convert neutral messages to the OpenAI wire format.

    System messages already in the history are dropped: the system prompt is
    passed separately by `chat()` so it always lands first.

    Args:
        messages: Conversation history in HelioAI's provider-neutral form.

    Returns:
        Message dicts ready to send as the `messages` request field.
    """
    out: list[dict] = []
    for msg in messages:
        if msg.role == "system":
            continue
        if msg.role == "user":
            out.append({"role": "user", "content": msg.content})
        elif msg.role == "assistant":
            if msg.tool_calls:
                wire = {
                    "role": "assistant",
                    "content": msg.content or None,
                    "tool_calls": [
                        {
                            "id": tc.id,
                            "type": "function",
                            "function": {
                                "name": tc.name,
                                "arguments": json.dumps(tc.arguments or {}),
                            },
                        }
                        for tc in msg.tool_calls
                    ],
                }
            else:
                wire = {"role": "assistant", "content": msg.content}
            if msg.reasoning:
                wire["reasoning_content"] = msg.reasoning
            out.append(wire)
        elif msg.role == "tool":
            out.append(
                {
                    "role": "tool",
                    "tool_call_id": msg.tool_call_id or "",
                    "content": msg.content,
                }
            )
    return out

to_openai_tools

to_openai_tools(tools: list[ToolDef]) -> list[dict]

Convert tool definitions to OpenAI function-calling schemas.

Parameters:

Name Type Description Default
tools list[ToolDef]

Tools the agent may call this turn.

required

Returns:

Type Description
list[dict]

Function schemas ready to send as the tools request field.

Source code in helioai/core/llm/openai_compat.py
def to_openai_tools(tools: list[ToolDef]) -> list[dict]:
    """Convert tool definitions to OpenAI function-calling schemas.

    Args:
        tools: Tools the agent may call this turn.

    Returns:
        Function schemas ready to send as the `tools` request field.
    """
    return [
        {
            "type": "function",
            "function": {
                "name": t.name,
                "description": t.description,
                "parameters": t.parameters or {"type": "object", "properties": {}},
            },
        }
        for t in tools
    ]

from_openai_response

from_openai_response(response: Any, provider: str = 'openai') -> Message

Convert an OpenAI chat-completions response to a neutral message.

Text content is preserved even when tool calls are present, and a tool call whose arguments are not valid JSON degrades to {} with a warning rather than raising — a malformed model output must not kill the agent loop. An inline <think>...</think> reasoning block, when a provider emits one, is stripped from the content before it reaches the agent loop or the user; a separate reasoning_content field is kept on Message.reasoning, because DeepSeek requires it back on every later request that carries tools.

Parameters:

Name Type Description Default
response Any

The SDK response object.

required
provider str

Provider name, used only to label log messages.

'openai'

Returns:

Type Description
Message

The assistant's reply, with tool_calls set when the model requested any.

Source code in helioai/core/llm/openai_compat.py
def from_openai_response(response: Any, provider: str = "openai") -> Message:
    """Convert an OpenAI chat-completions response to a neutral message.

    Text content is preserved even when tool calls are present, and a tool call
    whose arguments are not valid JSON degrades to `{}` with a warning rather
    than raising — a malformed model output must not kill the agent loop. An
    inline `<think>...</think>` reasoning block, when a provider emits one, is
    stripped from the content before it reaches the agent loop or the user; a
    separate `reasoning_content` field is kept on `Message.reasoning`, because
    DeepSeek requires it back on every later request that carries tools.

    Args:
        response: The SDK response object.
        provider: Provider name, used only to label log messages.

    Returns:
        The assistant's reply, with `tool_calls` set when the model requested any.
    """
    choice = response.choices[0]
    msg = choice.message
    usage = _usage(response)
    raw_content = msg.content or ""
    content = _strip_reasoning(raw_content)
    tool_calls_raw = getattr(msg, "tool_calls", None) or []
    finish_reason = getattr(choice, "finish_reason", None)
    reasoning = getattr(msg, "reasoning_content", None) or None

    if not tool_calls_raw:
        if not content.strip():
            # A turn that produced nothing has exactly three causes and they need
            # different fixes: the budget ran out mid-generation (raise it), the model
            # spent the whole turn inside <think> and closed with nothing (shorten the
            # question), or it genuinely returned an empty completion (retry/provider).
            # The agent loop used to state the first as fact for all three. It is
            # visible here and nowhere else, so it is recorded here.
            log.warning(
                "%s empty completion: finish_reason=%s, %d raw chars, %d after stripping reasoning",
                provider,
                finish_reason,
                len(raw_content),
                len(content),
            )
        return Message(role="assistant", content=content, reasoning=reasoning, **usage)

    tool_calls = [
        _parse_tool_call(tc.id, tc.function.name, tc.function.arguments, finish_reason, provider)
        for tc in tool_calls_raw
    ]
    return Message(
        role="assistant", content=content, tool_calls=tool_calls, reasoning=reasoning, **usage
    )

helioai.core.llm.azure_openai

Azure-hosted OpenAI.

Same wire format as every other OpenAI-compatible provider, so the conversion logic lives in openai_compat. Azure only differs in how the client is built — the endpoint embeds the deployment name and an api-version — plus two dialect details the base class already exposes as parameters: the system prompt is sent with the developer role, and temperature is omitted when unset because GPT-5 and the o-series reject it.

AzureOpenAIClient

Bases: OpenAICompatClient

Chat client for an Azure OpenAI deployment.

Parameters:

Name Type Description Default
api_key str

Azure OpenAI API key.

required
endpoint str

Resource root, e.g. https://myresource.openai.azure.com.

required
api_version str

Azure API version, e.g. 2024-12-01-preview.

required
deployment str

Deployment name, sent as the request's model.

required
max_output_tokens int

Cap on generated tokens.

4096
temperature float | None

Sampling temperature; None omits the field.

None
Source code in helioai/core/llm/azure_openai.py
class AzureOpenAIClient(OpenAICompatClient):
    """Chat client for an Azure OpenAI deployment.

    Args:
        api_key: Azure OpenAI API key.
        endpoint: Resource root, e.g. `https://myresource.openai.azure.com`.
        api_version: Azure API version, e.g. `2024-12-01-preview`.
        deployment: Deployment name, sent as the request's `model`.
        max_output_tokens: Cap on generated tokens.
        temperature: Sampling temperature; `None` omits the field.
    """

    def __init__(
        self,
        api_key: str,
        endpoint: str,
        api_version: str,
        deployment: str,
        max_output_tokens: int = 4096,
        temperature: float | None = None,
    ):
        super().__init__(
            client=AsyncAzureOpenAI(
                api_key=api_key,
                azure_endpoint=endpoint,
                api_version=api_version,
                max_retries=0,
            ),
            model=deployment,
            system_role="developer",
            max_output_tokens=max_output_tokens,
            temperature=temperature,
            provider="azure",
        )

helioai.core.llm.gemini

GeminiClient: Google genai SDK → neutral Message model.

GeminiClient

Bases: LLMClient

Chat client for Google Gemini, using the native google-genai SDK.

Kept separate from OpenAICompatClient because the wire format genuinely differs: turns are Content objects with typed parts, the assistant role is called model, and tool results are matched by function name rather than by id. Since Gemini issues no call ids, this client synthesises name::hex and parses the name back out on the way in — a format that is persisted in existing sessions, so it must stay readable.

Parameters:

Name Type Description Default
api_key str

Gemini API key.

required
model str

Model name, e.g. gemini-2.5-flash.

required
max_output_tokens int

Cap on generated tokens.

4096
temperature float

Sampling temperature.

0.2
Source code in helioai/core/llm/gemini.py
class GeminiClient(LLMClient):
    """Chat client for Google Gemini, using the native `google-genai` SDK.

    Kept separate from `OpenAICompatClient` because the wire format genuinely
    differs: turns are `Content` objects with typed parts, the assistant role is
    called `model`, and tool results are matched by function *name* rather than
    by id. Since Gemini issues no call ids, this client synthesises `name::hex`
    and parses the name back out on the way in — a format that is persisted in
    existing sessions, so it must stay readable.

    Args:
        api_key: Gemini API key.
        model: Model name, e.g. `gemini-2.5-flash`.
        max_output_tokens: Cap on generated tokens.
        temperature: Sampling temperature.
    """

    def __init__(
        self, api_key: str, model: str, max_output_tokens: int = 4096, temperature: float = 0.2
    ):
        self._client = genai.Client(api_key=api_key)
        self._model = model
        self._max_output_tokens = max_output_tokens
        self._temperature = temperature

    async def aclose(self) -> None:
        """Close the transport held by the genai client (its close() is sync)."""
        await close_sdk_client(self._client)

    async def chat(
        self,
        messages: list[Message],
        tools: list[ToolDef],
        system_prompt: str | None = None,
        tool_choice: str = "auto",
    ) -> Message:
        """Send one turn to Gemini and return the assistant's reply.

        When `system_prompt` is not given, it is recovered from any system message
        left in the history. `tool_choice="required"` maps to Gemini's ANY mode.
        """
        contents = self._to_gemini_contents(messages)
        gemini_tools = self._to_gemini_tools(tools) if tools else None

        if system_prompt is None:
            system_prompt = next((m.content for m in messages if m.role == "system"), None)

        tool_config = None
        if gemini_tools and tool_choice == "required":
            tool_config = gt.ToolConfig(
                function_calling_config=gt.FunctionCallingConfig(mode="ANY")
            )

        config = gt.GenerateContentConfig(
            max_output_tokens=self._max_output_tokens,
            temperature=self._temperature,
            tools=gemini_tools,
            tool_config=tool_config,
            system_instruction=system_prompt,
        )

        response = await call_with_retry(
            lambda: self._client.aio.models.generate_content(
                model=self._model,
                contents=contents,
                config=config,
            )
        )
        return self._from_gemini_response(response)

    @staticmethod
    def _to_gemini_contents(messages: list[Message]) -> list[gt.Content]:
        contents: list[gt.Content] = []
        for msg in messages:
            if msg.role == "system":
                continue
            if msg.role == "user":
                contents.append(gt.Content(role="user", parts=[gt.Part(text=msg.content)]))
            elif msg.role == "assistant":
                if msg.tool_calls:
                    parts = [
                        gt.Part(function_call=gt.FunctionCall(name=tc.name, args=tc.arguments))
                        for tc in msg.tool_calls
                    ]
                    contents.append(gt.Content(role="model", parts=parts))
                else:
                    contents.append(gt.Content(role="model", parts=[gt.Part(text=msg.content)]))
            elif msg.role == "tool":
                name, _, _ = (msg.tool_call_id or "::").partition("::")
                try:
                    response_payload = json.loads(msg.content)
                    if not isinstance(response_payload, dict):
                        response_payload = {"result": response_payload}
                except json.JSONDecodeError:
                    response_payload = {"result": msg.content}
                contents.append(
                    gt.Content(
                        role="user",
                        parts=[
                            gt.Part(
                                function_response=gt.FunctionResponse(
                                    name=name,
                                    response=response_payload,
                                )
                            )
                        ],
                    )
                )
        return contents

    @staticmethod
    def _to_gemini_tools(tools: list[ToolDef]) -> list[gt.Tool]:
        declarations = [
            gt.FunctionDeclaration(
                name=t.name,
                description=t.description,
                parameters=t.parameters or None,
            )
            for t in tools
        ]
        return [gt.Tool(function_declarations=declarations)]

    @staticmethod
    def _from_gemini_response(response) -> Message:
        tool_calls: list[ToolCall] = []
        text_chunks: list[str] = []
        candidate = response.candidates[0] if response.candidates else None
        parts = candidate.content.parts if candidate and candidate.content else []
        for part in parts or []:
            fc = getattr(part, "function_call", None)
            if fc and fc.name:
                args = dict(fc.args) if fc.args else {}
                call_id = f"{fc.name}::{uuid.uuid4().hex[:8]}"
                tool_calls.append(ToolCall(id=call_id, name=fc.name, arguments=args))
                continue
            text = getattr(part, "text", None)
            if text:
                text_chunks.append(text)
        # google-genai names every count differently from the OpenAI wire format,
        # which is the whole reason this client is not the compat one.
        um = getattr(response, "usage_metadata", None)
        usage = {
            "prompt_tokens": (getattr(um, "prompt_token_count", 0) or 0) if um else 0,
            "completion_tokens": (getattr(um, "candidates_token_count", 0) or 0) if um else 0,
            "cached_tokens": (getattr(um, "cached_content_token_count", 0) or 0) if um else 0,
        }
        if tool_calls:
            return Message(
                role="assistant",
                content="".join(text_chunks),
                tool_calls=tool_calls,
                **usage,
            )
        return Message(role="assistant", content="".join(text_chunks), **usage)

aclose async

aclose() -> None

Close the transport held by the genai client (its close() is sync).

Source code in helioai/core/llm/gemini.py
async def aclose(self) -> None:
    """Close the transport held by the genai client (its close() is sync)."""
    await close_sdk_client(self._client)

chat async

chat(messages: list[Message], tools: list[ToolDef], system_prompt: str | None = None, tool_choice: str = 'auto') -> Message

Send one turn to Gemini and return the assistant's reply.

When system_prompt is not given, it is recovered from any system message left in the history. tool_choice="required" maps to Gemini's ANY mode.

Source code in helioai/core/llm/gemini.py
async def chat(
    self,
    messages: list[Message],
    tools: list[ToolDef],
    system_prompt: str | None = None,
    tool_choice: str = "auto",
) -> Message:
    """Send one turn to Gemini and return the assistant's reply.

    When `system_prompt` is not given, it is recovered from any system message
    left in the history. `tool_choice="required"` maps to Gemini's ANY mode.
    """
    contents = self._to_gemini_contents(messages)
    gemini_tools = self._to_gemini_tools(tools) if tools else None

    if system_prompt is None:
        system_prompt = next((m.content for m in messages if m.role == "system"), None)

    tool_config = None
    if gemini_tools and tool_choice == "required":
        tool_config = gt.ToolConfig(
            function_calling_config=gt.FunctionCallingConfig(mode="ANY")
        )

    config = gt.GenerateContentConfig(
        max_output_tokens=self._max_output_tokens,
        temperature=self._temperature,
        tools=gemini_tools,
        tool_config=tool_config,
        system_instruction=system_prompt,
    )

    response = await call_with_retry(
        lambda: self._client.aio.models.generate_content(
            model=self._model,
            contents=contents,
            config=config,
        )
    )
    return self._from_gemini_response(response)

helioai.core.llm.factory

Build the configured LLM client.

Providers that speak the OpenAI wire format are table entries, not classes — see OPENAI_COMPAT. Azure and Gemini need their own SDK client objects and stay explicit below.

build_llm_client

build_llm_client(provider: str | None = None, model: str | None = None) -> LLMClient

Return a client for the requested provider.

Parameters:

Name Type Description Default
provider str | None

Provider name. Defaults to HELIOAI_LLM_PROVIDER.

None
model str | None

Model (or, on Azure, deployment) to use instead of the provider's configured one — how a delegated role runs on a smaller model than the lead (HELIOAI_ROLE_MODELS). None keeps the configured model.

None

Returns:

Type Description
LLMClient

A ready-to-use client.

Raises:

Type Description
RuntimeError

If the provider is unknown or its API key is missing.

Example

llm = build_llm_client("groq") type(llm).name 'OpenAICompatClient'

Source code in helioai/core/llm/factory.py
def build_llm_client(provider: str | None = None, model: str | None = None) -> LLMClient:
    """Return a client for the requested provider.

    Args:
        provider: Provider name. Defaults to `HELIOAI_LLM_PROVIDER`.
        model: Model (or, on Azure, deployment) to use instead of the provider's
            configured one — how a delegated role runs on a smaller model than the lead
            (`HELIOAI_ROLE_MODELS`). None keeps the configured model.

    Returns:
        A ready-to-use client.

    Raises:
        RuntimeError: If the provider is unknown or its API key is missing.

    Example:
        >>> llm = build_llm_client("groq")
        >>> type(llm).__name__
        'OpenAICompatClient'
    """
    p = (provider or settings.llm.provider).lower()
    validate_experiments()
    validate_judgment()

    if p == "azure":
        from helioai.core.llm.azure_openai import AzureOpenAIClient

        cfg = settings.llm.azure
        if not cfg.api_key:
            raise RuntimeError("AZURE_OPENAI_API_KEY is not set in .env")
        if not cfg.endpoint:
            raise RuntimeError("AZURE_OPENAI_ENDPOINT is not set in .env")
        return AzureOpenAIClient(
            api_key=cfg.api_key,
            endpoint=cfg.endpoint,
            api_version=cfg.api_version,
            deployment=model or cfg.deployment,
            max_output_tokens=cfg.max_output_tokens,
            temperature=cfg.temperature,
        )

    if p == "gemini":
        from helioai.core.llm.gemini import GeminiClient

        cfg = settings.llm.gemini
        if not cfg.api_key:
            raise RuntimeError(
                "GEMINI_API_KEY is not set in .env (https://aistudio.google.com/apikey)"
            )
        return GeminiClient(
            api_key=cfg.api_key,
            model=model or cfg.model,
            max_output_tokens=cfg.max_output_tokens,
            temperature=cfg.temperature,
        )

    if p in OPENAI_COMPAT:
        from helioai.core.llm.openai_compat import OpenAICompatClient

        spec = OPENAI_COMPAT[p]
        cfg = getattr(settings.llm, spec["config"])
        api_key = getattr(cfg, "api_key", "")
        if spec["key_env"] and not api_key:
            raise RuntimeError(f"{spec['key_env']} is not set in .env{spec.get('key_url', '')}")
        base_url = spec["base_url"] or f"{getattr(cfg, 'base_url', '').rstrip('/')}/v1"
        return OpenAICompatClient(
            provider=p,
            model=model or cfg.model,
            api_key=api_key,
            base_url=base_url,
            max_output_tokens=cfg.max_output_tokens,
            temperature=cfg.temperature,
            default_headers=_resolve_headers(getattr(cfg, "headers", None)),
        )

    known = "|".join(["azure", "gemini", *OPENAI_COMPAT])
    raise RuntimeError(f"Unknown LLM provider: {p!r}. Use {known}")

Storage

helioai.datastore

Session datastore — persist downloaded timeseries as .npz for reuse in run_python.

/data/.npz — compressed numpy archive

/data/manifest.json — index: name → metadata

Datasets are accessible in the sandbox via load_data("name"). The persisted data is the full-resolution download (before any downsampling). All I/O errors are silently swallowed — persistence must never break a tool call.

fill_mask

fill_mask(values, fillval: float | list | None = None)

Boolean mask of samples that carry no measurement.

Three conventions have to be caught at once, which is why every caller shares this one function instead of applying its own threshold:

  • non-finite (NaN/inf) — already unusable;
  • the ~1e31 magnitude convention, used by ACE among others;
  • the value the dataset declares in its CDF FILLVAL, which is the only way to catch Wind/SWE's 99999.9. A blanket "reject >= 99999" rule is wrong: OMNI carries a real proton temperature of 99093 K in the 2003 Halloween window.

Parameters:

Name Type Description Default
fillval float | list | None

the declared FILLVAL, when the provider exposes it. A list as often as a scalar — Wind/SWE declares [99999.8984375], and a bare float(fillval) on it raises. Compared with rtol=1e-6 to survive a float32 round-trip while staying far tighter than the ~1% gap to real nearby values.

None
Example

fill_mask(np.array([1.2, -1e31, 3.4]), [-1e31]).tolist() [False, True, False]

Source code in helioai/datastore.py
def fill_mask(values, fillval: float | list | None = None):
    """Boolean mask of samples that carry no measurement.

    Three conventions have to be caught at once, which is why every caller shares
    this one function instead of applying its own threshold:

    - non-finite (NaN/inf) — already unusable;
    - the ~1e31 magnitude convention, used by ACE among others;
    - the value the dataset *declares* in its CDF FILLVAL, which is the only way
      to catch Wind/SWE's 99999.9. A blanket "reject >= 99999" rule is wrong: OMNI
      carries a real proton temperature of 99093 K in the 2003 Halloween window.

    Args:
        fillval: the declared FILLVAL, when the provider exposes it. A list as often
            as a scalar — Wind/SWE declares `[99999.8984375]`, and a bare
            `float(fillval)` on it raises. Compared with rtol=1e-6 to survive a
            float32 round-trip while staying far tighter than the ~1% gap to real
            nearby values.

    Example:
        >>> fill_mask(np.array([1.2, -1e31, 3.4]), [-1e31]).tolist()
        [False, True, False]
    """
    import numpy as np

    bad = ~np.isfinite(values) | (np.abs(values) >= 1e30)
    if fillval is not None:
        try:
            fv = float(np.asarray(fillval).ravel()[0])
            if np.isfinite(fv):
                bad = bad | np.isclose(values, fv, rtol=1e-6)
        except (TypeError, ValueError, IndexError):
            pass
    return bad

blank_fill

blank_fill(values: Any, fillval: Any = None) -> tuple[Any, Any | None]

Return (values with fill blanked to NaN, mask), or (values, None) if not numeric.

Applied at every point where downloaded data is persisted, so that anything reading it back — the sandbox, the exported notebook — sees NaN for "no measurement" rather than a sentinel that looks like a plausible reading. Leaving it to the reader meant remembering to call clean(), which cannot see FILLVAL, so one forgotten call put a 99999.9 "speed" into a plot and a mean.

Non-numeric parameters (string labels, epochs) have no fill convention and are passed through untouched.

Parameters:

Name Type Description Default
values Any

Array as the provider delivered it.

required
fillval Any

Sentinel(s) the archive declares, read from FILLVAL rather than guessed — a hard-coded threshold blanks real measurements.

None

Returns:

Type Description
Any

(values with fill replaced by NaN, boolean mask of what was blanked),

Any | None

or (values, None) when the array is not numeric.

Example

vals, mask = blank_fill(np.array([1.2, -1e31, 3.4]), [-1e31]) vals.tolist(), mask.tolist() ([1.2, nan, 3.4], [False, True, False])

Source code in helioai/datastore.py
def blank_fill(values: Any, fillval: Any = None) -> tuple[Any, Any | None]:
    """Return (values with fill blanked to NaN, mask), or (values, None) if not numeric.

    Applied at every point where downloaded data is persisted, so that anything
    reading it back — the sandbox, the exported notebook — sees NaN for "no
    measurement" rather than a sentinel that looks like a plausible reading.
    Leaving it to the reader meant remembering to call clean(), which cannot see
    FILLVAL, so one forgotten call put a 99999.9 "speed" into a plot and a mean.

    Non-numeric parameters (string labels, epochs) have no fill convention and are
    passed through untouched.

    Args:
        values: Array as the provider delivered it.
        fillval: Sentinel(s) the archive declares, read from `FILLVAL` rather
            than guessed — a hard-coded threshold blanks real measurements.

    Returns:
        `(values with fill replaced by NaN, boolean mask of what was blanked)`,
        or `(values, None)` when the array is not numeric.

    Example:
        >>> vals, mask = blank_fill(np.array([1.2, -1e31, 3.4]), [-1e31])
        >>> vals.tolist(), mask.tolist()
        ([1.2, nan, 3.4], [False, True, False])
    """
    import numpy as np

    try:
        numeric = np.array(values, dtype="float64")
    except (TypeError, ValueError):
        return values, None
    mask = fill_mask(numeric, fillval)
    numeric[mask] = np.nan
    return numeric, mask

dir_lock

dir_lock(path: Path) -> threading.Lock

The lock to hold across a read-modify-write of the files under path.

A threading.Lock, not an asyncio.Lock: these writers are synchronous, called from coroutines today and from worker threads tomorrow, and a thread lock is correct in both places without binding to any event loop. One lock per directory string, created on first use and never dropped — a few dozen entries per process at most.

Parameters:

Name Type Description Default
path Path

The directory whose index files the caller is about to rewrite.

required

Returns:

Type Description
Lock

The same lock for the same directory, for the life of the process.

Source code in helioai/datastore.py
def dir_lock(path: Path) -> threading.Lock:
    """The lock to hold across a read-modify-write of the files under `path`.

    A `threading.Lock`, not an `asyncio.Lock`: these writers are synchronous, called
    from coroutines today and from worker threads tomorrow, and a thread lock is
    correct in both places without binding to any event loop. One lock per
    directory string, created on first use and never dropped — a few dozen entries
    per process at most.

    Args:
        path: The directory whose index files the caller is about to rewrite.

    Returns:
        The same lock for the same directory, for the life of the process.
    """
    key = str(path)
    with _dir_locks_guard:
        lock = _dir_locks.get(key)
        if lock is None:
            lock = _dir_locks[key] = threading.Lock()
        return lock

write_json_atomic

write_json_atomic(path: Path, payload: Any) -> None

Write payload as JSON so that path is never observed half-written.

The text goes to a sibling .tmp file first and is swapped in with os.replace, atomic on POSIX and on Windows when the target exists. A crash mid-write used to leave a truncated manifest.json and with it an unreadable session.

Parameters:

Name Type Description Default
path Path

Destination file.

required
payload Any

Anything json.dumps accepts.

required
Source code in helioai/datastore.py
def write_json_atomic(path: Path, payload: Any) -> None:
    """Write `payload` as JSON so that `path` is never observed half-written.

    The text goes to a sibling `.tmp` file first and is swapped in with `os.replace`,
    atomic on POSIX and on Windows when the target exists. A crash mid-write used to
    leave a truncated `manifest.json` and with it an unreadable session.

    Args:
        path: Destination file.
        payload: Anything `json.dumps` accepts.
    """
    tmp = path.with_name(path.name + ".tmp")
    tmp.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
    try:
        os.replace(tmp, path)
    except OSError:
        tmp.unlink(missing_ok=True)
        raise

read_manifest

read_manifest(session_dir: Path) -> dict

Return the manifest dict for a given session directory.

Parameters:

Name Type Description Default
session_dir Path

A session workspace directory (contains data/).

required

Returns:

Type Description
dict

{"datasets": {name: {kind, param_id, start, stop, ...}}} — empty

dict

datasets dict when no manifest exists.

Source code in helioai/datastore.py
def read_manifest(session_dir: Path) -> dict:
    """Return the manifest dict for a given session directory.

    Args:
        session_dir: A session workspace directory (contains `data/`).

    Returns:
        {"datasets": {name: {kind, param_id, start, stop, ...}}} — empty
        datasets dict when no manifest exists.
    """
    return _read_manifest_file(session_dir / DATA_SUBDIR)

find_existing

find_existing(param_id: str, start: str, stop: str, data_dir: Path | None = None) -> str | None

Name of an already-persisted dataset for this exact param and window, or None.

The same matching _unique_name uses to reuse a slot — but consulted before the download rather than after, so a repeat request costs a dict lookup instead of a network round-trip. The prompt has always said "download each parameter ONCE"; a real run still re-fetched the same Wind field three times across three turns, with the dataset name sitting in plain sight in its own history. Discipline the model does not reliably apply belongs in the tool.

Parameters:

Name Type Description Default
param_id str

Full speasy id, e.g. cda/WI_H0_MFI/BGSM.

required
start str

ISO start of the window, matched exactly.

required
stop str

ISO stop of the window, matched exactly.

required
data_dir Path | None

The session's data directory; None resolves the bound session.

None

Returns:

Type Description
str | None

The dataset name to load, or None. The match is exact: an overlapping

str | None

but different window is a miss, because returning a shorter series than

str | None

asked for would be silently wrong.

Source code in helioai/datastore.py
def find_existing(param_id: str, start: str, stop: str, data_dir: Path | None = None) -> str | None:
    """Name of an already-persisted dataset for this exact param and window, or None.

    The same matching `_unique_name` uses to reuse a slot — but consulted *before* the
    download rather than after, so a repeat request costs a dict lookup instead of a
    network round-trip. The prompt has always said "download each parameter ONCE"; a
    real run still re-fetched the same Wind field three times across three turns, with
    the dataset name sitting in plain sight in its own history. Discipline the model
    does not reliably apply belongs in the tool.

    Args:
        param_id: Full speasy id, e.g. `cda/WI_H0_MFI/BGSM`.
        start: ISO start of the window, matched exactly.
        stop: ISO stop of the window, matched exactly.
        data_dir: The session's data directory; None resolves the bound session.

    Returns:
        The dataset name to load, or None. The match is exact: an overlapping
        but different window is a miss, because returning a shorter series than
        asked for would be silently wrong.
    """
    try:
        data_dir = _session_data_dir(data_dir)
        if data_dir is None:
            return None
        for name, entry in _read_manifest_file(data_dir).get("datasets", {}).items():
            if (
                entry.get("param_id") == param_id
                and entry.get("start") == start
                and entry.get("stop") == stop
                and entry.get("kind") == "timeseries"
            ):
                return name
    except Exception as e:  # noqa: BLE001 — a cache miss must never block a download
        log.debug("find_existing failed: %s", e)
    return None

save_timeseries

save_timeseries(name_hint: str, *, time, values, param_id: str, units: str, start: str, stop: str, columns, source: str, data_dir: Path | None = None) -> dict | None

Persist a timeseries download as npz + a manifest entry.

Parameters:

Name Type Description Default
name_hint str

Basis for the dataset name (slugged, collision-suffixed).

required
time array - like

Time axis as returned by speasy.

required
values ndarray

Data array (fill already blanked).

required
param_id str

Speasy id, recorded in the manifest for the export rewrite.

required
units str

Physical units string.

required
start str

ISO window start, recorded for the export rewrite.

required
stop str

ISO window stop.

required
columns list[str]

Component names.

required
source str

Which tool produced the download.

required
data_dir Path | None

The session's data directory, as the run's context names it; None resolves the bound session and says so in the log.

None

Returns:

Type Description
dict | None

{"dataset": } to reference in load_data(), or None when

dict | None

persisting failed (the download result is still usable in-memory).

Source code in helioai/datastore.py
def save_timeseries(
    name_hint: str,
    *,
    time,
    values,
    param_id: str,
    units: str,
    start: str,
    stop: str,
    columns,
    source: str,
    data_dir: Path | None = None,
) -> dict | None:
    """Persist a timeseries download as npz + a manifest entry.

    Args:
        name_hint: Basis for the dataset name (slugged, collision-suffixed).
        time (array-like): Time axis as returned by speasy.
        values (numpy.ndarray): Data array (fill already blanked).
        param_id: Speasy id, recorded in the manifest for the export rewrite.
        units: Physical units string.
        start: ISO window start, recorded for the export rewrite.
        stop: ISO window stop.
        columns (list[str]): Component names.
        source: Which tool produced the download.
        data_dir: The session's data directory, as the run's context names it; None
            resolves the bound session and says so in the log.

    Returns:
        {"dataset": <final name>} to reference in `load_data()`, or None when
        persisting failed (the download result is still usable in-memory).
    """
    try:
        import time as _time

        import numpy as np

        data_dir = _session_data_dir(data_dir)
        if data_dir is None:
            return None

        time_arr = np.asarray(time)
        values_arr = np.asarray(values, dtype=float)

        if values_arr.nbytes > _MAX_BYTES:
            log.warning(
                "datastore: skipping %r — %d MB exceeds cap",
                param_id,
                values_arr.nbytes // (1024 * 1024),
            )
            return None

        with dir_lock(data_dir):
            manifest = _read_manifest_file(data_dir)
            base = _slug(name_hint or param_id)
            name = _unique_name(manifest, base, param_id, start, stop)
            fname = f"{name}.npz"

            np.savez_compressed(data_dir / fname, time=time_arr, values=values_arr)

            cols = list(columns) if isinstance(columns, (list, tuple)) else []
            # Derived from what was actually written rather than passed in, so it can
            # never drift from the file: get_timeseries blanks fill values to NaN
            # before saving, so this is the fraction with no measurement.
            missing_pct = (
                round(100 * float(np.isnan(values_arr).mean()), 1) if values_arr.size else 0.0
            )
            manifest.setdefault("datasets", {})[name] = {
                "kind": "timeseries",
                "file": fname,
                "param_id": param_id,
                "units": units,
                "start": start,
                "stop": stop,
                "shape": list(values_arr.shape),
                "columns": cols,
                "missing_pct": missing_pct,
                "source": source,
                "created": str(int(_time.time())),
            }
            _write_manifest_file(data_dir, manifest)
            return {"dataset": name}
    except Exception as e:
        log.warning("datastore: save_timeseries failed for %r: %s", param_id, e)
        return None

save_event_collection

save_event_collection(name_hint: str, *, series: list[tuple], param_id: str, units: str, source: str, data_dir: Path | None = None) -> dict | None

Persist a batch of per-event timeseries under one dataset name.

Parameters:

Name Type Description Default
name_hint str

Preferred dataset name; a suffix is added if it is taken.

required
series list[tuple]

[(event_start, event_stop, timeseries_or_None), ...]. A None entry records an event with no data rather than dropping it, so the count of events surveyed stays honest.

required
param_id str

The parameter every event was sampled from.

required
units str

Units as the provider reports them.

required
source str

Provenance string carried into the manifest.

required
data_dir Path | None

The session's data directory; None resolves the bound session.

None

Returns:

Type Description
dict | None

{"dataset": name}, or None when nothing could be written.

Source code in helioai/datastore.py
def save_event_collection(
    name_hint: str,
    *,
    series: list[tuple],
    param_id: str,
    units: str,
    source: str,
    data_dir: Path | None = None,
) -> dict | None:
    """Persist a batch of per-event timeseries under one dataset name.

    Args:
        name_hint: Preferred dataset name; a suffix is added if it is taken.
        series: `[(event_start, event_stop, timeseries_or_None), ...]`. A None
            entry records an event with no data rather than dropping it, so the
            count of events surveyed stays honest.
        param_id: The parameter every event was sampled from.
        units: Units as the provider reports them.
        source: Provenance string carried into the manifest.
        data_dir: The session's data directory; None resolves the bound session.

    Returns:
        `{"dataset": name}`, or None when nothing could be written.
    """
    try:
        import time as _time

        import numpy as np

        data_dir = _session_data_dir(data_dir)
        if data_dir is None:
            return None

        arrays: dict[str, any] = {}
        events_meta: list[dict] = []
        total_bytes = 0
        capped = False

        for i, (ev_start, ev_stop, ts) in enumerate(series):
            if ts is None or not hasattr(ts, "time") or not hasattr(ts, "values"):
                events_meta.append(
                    {"idx": i, "start": ev_start, "stop": ev_stop, "status": "no_data"}
                )
                continue
            try:
                t_arr = np.asarray(ts.time)
                v_arr, _ = blank_fill(ts.values, (getattr(ts, "meta", {}) or {}).get("FILLVAL"))
                v_arr = np.asarray(v_arr, dtype=float)
                total_bytes += t_arr.nbytes + v_arr.nbytes
                if total_bytes > _MAX_BYTES:
                    if not capped:
                        log.warning(
                            "datastore: event collection for %r exceeds 100 MB cap, truncating",
                            param_id,
                        )
                        capped = True
                    events_meta.append(
                        {"idx": i, "start": ev_start, "stop": ev_stop, "status": "truncated"}
                    )
                    continue
                arrays[f"t{i}"] = t_arr
                arrays[f"v{i}"] = v_arr
                events_meta.append({"idx": i, "start": ev_start, "stop": ev_stop, "status": "ok"})
            except Exception:
                events_meta.append(
                    {"idx": i, "start": ev_start, "stop": ev_stop, "status": "no_data"}
                )

        if not arrays:
            return None

        with dir_lock(data_dir):
            manifest = _read_manifest_file(data_dir)
            base = _slug(name_hint or param_id) + "_events"
            name = base
            if name in manifest.get("datasets", {}):
                i2 = 2
                while f"{base}_{i2}" in manifest.get("datasets", {}):
                    i2 += 1
                name = f"{base}_{i2}"
            fname = f"{name}.npz"

            np.savez_compressed(data_dir / fname, **arrays)

            manifest.setdefault("datasets", {})[name] = {
                "kind": "event_collection",
                "file": fname,
                "param_id": param_id,
                "units": units,
                "n_events": len(series),
                "events": events_meta,
                "source": source,
                "created": str(int(_time.time())),
            }
            _write_manifest_file(data_dir, manifest)
            return {"dataset": name}
    except Exception as e:
        log.warning("datastore: save_event_collection failed for %r: %s", param_id, e)
        return None

helioai.workspace

Workspace — stable output directory for sandbox figures and data.

Figures go to workspace//fig_N_M.png Code files go to workspace//code_N.py

The session label is a human-readable slug derived from the first user message, propagated via a contextvar set by stream_chat at the start of each request.

current_user

current_user() -> str

Return the user owning the current context, or the default user.

Returns:

Type Description
str

The bound user id, or DEFAULT_USER outside any bound context — the CLI

str

and the Jupyter magic never bind one, so unowned callers still get a

str

real storage home rather than an error.

Source code in helioai/workspace.py
def current_user() -> str:
    """Return the user owning the current context, or the default user.

    Returns:
        The bound user id, or `DEFAULT_USER` outside any bound context — the CLI
        and the Jupyter magic never bind one, so unowned callers still get a
        real storage home rather than an error.
    """
    return _current_user.get() or DEFAULT_USER

set_user

set_user(user_id: str) -> object

Bind the user that owns storage for the current context.

A contextvar rather than a global: every task spawned from here inherits the binding, which is what keeps a sub-agent writing into the same user's tree as the lead that spawned it, while a concurrent web request keeps its own.

Parameters:

Name Type Description Default
user_id str

Owner of every path derived until the binding is reset.

required

Returns:

Type Description
object

An opaque token to hand back to reset_user. Resetting by token rather

object

than by re-setting a previous value is what makes nesting safe.

Source code in helioai/workspace.py
def set_user(user_id: str) -> object:
    """Bind the user that owns storage for the current context.

    A contextvar rather than a global: every task spawned from here inherits the
    binding, which is what keeps a sub-agent writing into the same user's tree as
    the lead that spawned it, while a concurrent web request keeps its own.

    Args:
        user_id: Owner of every path derived until the binding is reset.

    Returns:
        An opaque token to hand back to `reset_user`. Resetting by token rather
        than by re-setting a previous value is what makes nesting safe.
    """
    return _current_user.set(user_id)

reset_user

reset_user(token: object) -> None

Restore the user binding that was in force before set_user.

Parameters:

Name Type Description Default
token object

The value set_user returned. A token from another contextvar, or one already reset, raises rather than silently rebinding.

required
Source code in helioai/workspace.py
def reset_user(token: object) -> None:
    """Restore the user binding that was in force before `set_user`.

    Args:
        token: The value `set_user` returned. A token from another contextvar,
            or one already reset, raises rather than silently rebinding.
    """
    _current_user.reset(token)  # type: ignore[arg-type]

user_home

user_home(user: str) -> Path

A user's private storage home: <data>/users/<user>/.

The directory is not created — callers that only need to read must not materialise a home for a user that does not exist.

Parameters:

Name Type Description Default
user str

Owner id. Not sanitised here; callers that accept it from a request pass it through safe_id first.

required

Returns:

Type Description
Path

The path, whether or not anything exists at it.

Example

user_home("cli") PosixPath('.../data/users/cli')

Source code in helioai/workspace.py
def user_home(user: str) -> Path:
    """A user's private storage home: `<data>/users/<user>/`.

    The directory is not created — callers that only need to read must not
    materialise a home for a user that does not exist.

    Args:
        user: Owner id. Not sanitised here; callers that accept it from a
            request pass it through `safe_id` first.

    Returns:
        The path, whether or not anything exists at it.

    Example:
        >>> user_home("cli")
        PosixPath('.../data/users/cli')
    """
    return _users_root() / user

set_session

set_session(session_id: str) -> object

Bind the session whose workspace the current context writes into.

Parameters:

Name Type Description Default
session_id str

Session id, used as a directory name when no label is bound.

required

Returns:

Type Description
object

An opaque token for reset_session.

Source code in helioai/workspace.py
def set_session(session_id: str) -> object:
    """Bind the session whose workspace the current context writes into.

    Args:
        session_id: Session id, used as a directory name when no label is bound.

    Returns:
        An opaque token for `reset_session`.
    """
    return _current_session.set(session_id)

reset_session

reset_session(token: object) -> None

Restore the session binding in force before set_session.

Parameters:

Name Type Description Default
token object

The value set_session returned.

required
Source code in helioai/workspace.py
def reset_session(token: object) -> None:
    """Restore the session binding in force before `set_session`.

    Args:
        token: The value `set_session` returned.
    """
    _current_session.reset(token)  # type: ignore[arg-type]

set_label

set_label(label: str) -> object

Bind the human-readable folder name for the current session.

Takes precedence over the session id in get_session_dir, so a workspace is findable by what was asked rather than by a uuid.

Parameters:

Name Type Description Default
label str

Directory name, sanitised on use rather than here.

required

Returns:

Type Description
object

An opaque token for reset_label.

Source code in helioai/workspace.py
def set_label(label: str) -> object:
    """Bind the human-readable folder name for the current session.

    Takes precedence over the session id in `get_session_dir`, so a workspace is
    findable by what was asked rather than by a uuid.

    Args:
        label: Directory name, sanitised on use rather than here.

    Returns:
        An opaque token for `reset_label`.
    """
    return _current_label.set(label)

reset_label

reset_label(token: object) -> None

Restore the label binding in force before set_label.

Parameters:

Name Type Description Default
token object

The value set_label returned.

required
Source code in helioai/workspace.py
def reset_label(token: object) -> None:
    """Restore the label binding in force before `set_label`.

    Args:
        token: The value `set_label` returned.
    """
    _current_label.reset(token)  # type: ignore[arg-type]

safe_id

safe_id(value: str, fallback: str = 'session') -> str

Reduce an identifier to something that cannot escape its parent directory.

Session ids are caller-supplied — a web request body, an MCP client, a CLI flag — and end up as path components here, in the export filename, and in the rmtree behind DELETE /api/sessions/{id}. Everything the project mints is a uuid4, so stripping to [A-Za-z0-9_-] is lossless in practice and turns ../.. into the fallback rather than a parent directory.

Parameters:

Name Type Description Default
value str

Caller-supplied identifier, trusted for nothing.

required
fallback str

Returned when nothing survives the filter, so the result is never the empty string — which would resolve to the parent directory.

'session'

Returns:

Type Description
str

At most 64 characters of [A-Za-z0-9_-], or fallback.

Example

safe_id("../../etc/passwd"), safe_id("sess-abc-123456") ('etcpasswd', 'sess-abc-123456')

Source code in helioai/workspace.py
def safe_id(value: str, fallback: str = "session") -> str:
    """Reduce an identifier to something that cannot escape its parent directory.

    Session ids are caller-supplied — a web request body, an MCP client, a CLI
    flag — and end up as path components here, in the export filename, and in the
    `rmtree` behind `DELETE /api/sessions/{id}`. Everything the project mints is a
    uuid4, so stripping to `[A-Za-z0-9_-]` is lossless in practice and turns
    `../..` into the fallback rather than a parent directory.

    Args:
        value: Caller-supplied identifier, trusted for nothing.
        fallback: Returned when nothing survives the filter, so the result is
            never the empty string — which would resolve to the parent directory.

    Returns:
        At most 64 characters of `[A-Za-z0-9_-]`, or `fallback`.

    Example:
        >>> safe_id("../../etc/passwd"), safe_id("sess-abc-123456")
        ('etcpasswd', 'sess-abc-123456')
    """
    cleaned = re.sub(r"[^A-Za-z0-9_-]", "", value)[:64]
    return cleaned or fallback

make_session_label

make_session_label(first_message: str, session_id: str, taken: Collection[str] = ()) -> str

Build a human-readable slug for the session workspace folder.

The session id suffix is what keeps two identical questions from sharing one directory; the words are only there so a human can find it. Six characters of the id are enough when ids are random, and not when a client chooses them: the web API accepts any ^[A-Za-z0-9_-]{1,64}$, so session-001 and session-002 — or a benchmark's bench-<question>-<hex> — collapsed onto one directory, and the second session found the first one's downloads in its inventory. taken is the set of labels the user already owns; when the six-character label is among them the suffix grows until it is not, and falls back to a digest of the whole id.

Parameters:

Name Type Description Default
first_message str

The question that opened the session.

required
session_id str

Session id; its first six characters are the usual discriminator.

required
taken Collection[str]

Labels already assigned to this user's other sessions.

()

Returns:

Type Description
str

A slug of at most four words plus the id suffix, distinct from every taken.

Example

make_session_label("Plot IMF Bz from ACE", "abc123def456") 'plot-imf-bz-from_abc123' make_session_label("Plot IMF Bz from ACE", "abc123def456", {"plot-imf-bz-from_abc123"}) 'plot-imf-bz-from_abc123de'

Source code in helioai/workspace.py
def make_session_label(first_message: str, session_id: str, taken: Collection[str] = ()) -> str:
    """Build a human-readable slug for the session workspace folder.

    The session id suffix is what keeps two identical questions from sharing one
    directory; the words are only there so a human can find it. Six characters of the
    id are enough when ids are random, and not when a client chooses them: the web API
    accepts any `^[A-Za-z0-9_-]{1,64}$`, so `session-001` and `session-002` — or a
    benchmark's `bench-<question>-<hex>` — collapsed onto one directory, and the second
    session found the first one's downloads in its inventory. `taken` is the set of
    labels the user already owns; when the six-character label is among them the suffix
    grows until it is not, and falls back to a digest of the whole id.

    Args:
        first_message: The question that opened the session.
        session_id: Session id; its first six characters are the usual discriminator.
        taken: Labels already assigned to this user's other sessions.

    Returns:
        A slug of at most four words plus the id suffix, distinct from every `taken`.

    Example:
        >>> make_session_label("Plot IMF Bz from ACE", "abc123def456")
        'plot-imf-bz-from_abc123'
        >>> make_session_label("Plot IMF Bz from ACE", "abc123def456", {"plot-imf-bz-from_abc123"})
        'plot-imf-bz-from_abc123de'
    """
    words = re.sub(r"[^a-z0-9\s]", "", first_message.lower().strip()).split()
    slug = "-".join(words[:4]) if words else "session"
    sid = safe_id(session_id)
    label = f"{slug[:25]}_{sid[:6]}"
    if label not in taken:
        return label
    for n in range(8, len(sid) + 1, 2):
        candidate = f"{slug[:25]}_{sid[:n]}"
        if candidate not in taken:
            return candidate
    digest = hashlib.blake2b(session_id.encode(), digest_size=3).hexdigest()
    return f"{slug[:25]}_{sid}-{digest}"

session_dir_for

session_dir_for(user: str, session_id: str, label: str | None = None) -> Path

The workspace directory of a session, from its ids rather than from the ambience.

The same rule get_session_dir applies to the bound contextvars — the label when the session has one, the id otherwise — so a RunContext built from ids and a caller reading the contextvars land in the same directory.

Parameters:

Name Type Description Default
user str

Owner of the session.

required
session_id str

The conversation.

required
label str | None

Its human-readable directory name, when already minted.

None

Returns:

Type Description
Path

The directory, created if missing.

Source code in helioai/workspace.py
def session_dir_for(user: str, session_id: str, label: str | None = None) -> Path:
    """The workspace directory of a session, from its ids rather than from the ambience.

    The same rule `get_session_dir` applies to the bound contextvars — the label when
    the session has one, the id otherwise — so a `RunContext` built from ids and a
    caller reading the contextvars land in the same directory.

    Args:
        user: Owner of the session.
        session_id: The conversation.
        label: Its human-readable directory name, when already minted.

    Returns:
        The directory, created if missing.
    """
    root = user_home(user) / "workspace"
    d = root / safe_id(label or session_id)
    d.mkdir(parents=True, exist_ok=True)
    return d

get_session_dir

get_session_dir() -> Path

Return the workspace directory for the current session.

Prefers the bound label, falls back to the bound session id, and lands in a temporary directory when neither is bound — an unbound caller still gets a writable place rather than an exception, because run_python must not fail for want of a session.

Returns:

Type Description
Path

The directory, created if missing.

Source code in helioai/workspace.py
def get_session_dir() -> Path:
    """Return the workspace directory for the current session.

    Prefers the bound label, falls back to the bound session id, and lands in a
    temporary directory when neither is bound — an unbound caller still gets a
    writable place rather than an exception, because `run_python` must not fail
    for want of a session.

    Returns:
        The directory, created if missing.
    """
    label = _current_label.get()
    if label:
        d = _root() / safe_id(label)
        d.mkdir(parents=True, exist_ok=True)
        return d
    session_id = _current_session.get()
    if session_id:
        d = _root() / safe_id(session_id)
        d.mkdir(parents=True, exist_ok=True)
        return d
    return _no_session_dir()

get_next_run_idx

get_next_run_idx(session_dir: Path) -> int

Return the next available run index for a session directory.

Derived from what is on disk rather than from a counter, so the numbering survives a restart and stays right when a session is resumed.

Parameters:

Name Type Description Default
session_dir Path

Directory holding the code_N.py files of past runs.

required

Returns:

Type Description
int

max(N) + 1, or 0 when the session has run nothing yet.

Source code in helioai/workspace.py
def get_next_run_idx(session_dir: Path) -> int:
    """Return the next available run index for a session directory.

    Derived from what is on disk rather than from a counter, so the numbering
    survives a restart and stays right when a session is resumed.

    Args:
        session_dir: Directory holding the `code_N.py` files of past runs.

    Returns:
        `max(N) + 1`, or 0 when the session has run nothing yet.
    """
    existing = list(session_dir.glob("code_*.py"))
    if not existing:
        return 0
    indices = []
    for p in existing:
        parts = p.stem.split("_")
        if len(parts) == 2 and parts[1].isdigit():
            indices.append(int(parts[1]))
    return max(indices) + 1 if indices else 0

get_run_dir_for_sandbox

get_run_dir_for_sandbox() -> str

The current session directory as a string, for the sandbox fallback.

Returns:

Type Description
str

get_session_dir() as a plain string — the fallback path builds a bwrap

str

argument list, which takes no Path.

Source code in helioai/workspace.py
def get_run_dir_for_sandbox() -> str:
    """The current session directory as a string, for the sandbox fallback.

    Returns:
        `get_session_dir()` as a plain string — the fallback path builds a bwrap
        argument list, which takes no `Path`.
    """
    return str(get_session_dir())

is_under_workspace

is_under_workspace(path: str | Path) -> bool

True if path is safely under the per-user storage root (no traversal).

is_relative_to rather than a string prefix: comparing against str(root) + "/" hard-coded the POSIX separator, so on Windows the check never matched and /figure and /code returned 404 for every legitimate path. Fail-closed, so it was a dead web UI rather than a hole — but dead all the same.

Parameters:

Name Type Description Default
path str | Path

Candidate path, resolved before comparison so symlinks and .. cannot point outside from within.

required

Returns:

Type Description
bool

True when the resolved path sits under the users root. False on any

bool

resolution error, which keeps an unreadable path from being served.

Source code in helioai/workspace.py
def is_under_workspace(path: str | Path) -> bool:
    """True if path is safely under the per-user storage root (no traversal).

    `is_relative_to` rather than a string prefix: comparing against `str(root) + "/"`
    hard-coded the POSIX separator, so on Windows the check never matched and `/figure`
    and `/code` returned 404 for every legitimate path. Fail-closed, so it was a dead
    web UI rather than a hole — but dead all the same.

    Args:
        path: Candidate path, resolved before comparison so symlinks and `..`
            cannot point outside from within.

    Returns:
        True when the resolved path sits under the users root. False on any
        resolution error, which keeps an unreadable path from being served.
    """
    try:
        p = Path(path).resolve()
        return p.is_relative_to(_users_root().resolve())
    except (ValueError, OSError):
        return False

cleanup_old_runs

cleanup_old_runs(ttl_seconds: int | None = None) -> int

Purge session directories older than the TTL, for every user.

Called on CLI startup, when the notebook magic loads and when the MCP server starts, and every hour by cleanup_periodically under the web server — the one process that never restarts, and so never reached this until it did. Removal failures are ignored rather than raised: housekeeping must not stop a user from asking a question. A user's home (profile, catalogs, speasy seed) is never touched; only what is under workspace/.

Parameters:

Name Type Description Default
ttl_seconds int | None

Age above which a session directory is deleted, measured on its mtime. Defaults to settings.workspace.ttl_seconds.

None

Returns:

Type Description
int

How many directories were removed.

Source code in helioai/workspace.py
def cleanup_old_runs(ttl_seconds: int | None = None) -> int:
    """Purge session directories older than the TTL, for every user.

    Called on CLI startup, when the notebook magic loads and when the MCP server
    starts, and every hour by `cleanup_periodically` under the web server — the one
    process that never restarts, and so never reached this until it did. Removal
    failures are ignored rather than raised: housekeeping must not stop a user from
    asking a question. A user's home (profile, catalogs, speasy seed) is never touched;
    only what is under `workspace/`.

    Args:
        ttl_seconds: Age above which a session directory is deleted, measured on
            its mtime. Defaults to `settings.workspace.ttl_seconds`.

    Returns:
        How many directories were removed.
    """
    from helioai.config import settings

    if ttl_seconds is None:
        ttl_seconds = settings.workspace.ttl_seconds
    users_root = _users_root()
    if not users_root.exists():
        return 0
    cutoff = time.time() - ttl_seconds
    removed = 0
    for home in users_root.iterdir():
        ws = home / "workspace"
        if not ws.is_dir():
            continue
        for session_dir in ws.iterdir():
            if session_dir.is_dir() and session_dir.stat().st_mtime < cutoff:
                shutil.rmtree(session_dir, ignore_errors=True)
                removed += 1
    return removed

cleanup_periodically async

cleanup_periodically(period_seconds: float = 3600.0) -> None

Run cleanup_old_runs every period_seconds until the task is cancelled.

The web server used to clean up once, at startup, and then run for weeks: a demo machine had 76 session directories of 269 MB each, every one of them past the TTL. The sweep runs in a worker thread so a slow disk never stalls a stream, and a failing sweep is logged and retried at the next tick rather than ending the task.

Parameters:

Name Type Description Default
period_seconds float

Time between sweeps; the first sweep happens after one period, since the caller already swept at startup.

3600.0
Source code in helioai/workspace.py
async def cleanup_periodically(period_seconds: float = 3600.0) -> None:
    """Run `cleanup_old_runs` every `period_seconds` until the task is cancelled.

    The web server used to clean up once, at startup, and then run for weeks: a demo
    machine had 76 session directories of 269 MB each, every one of them past the
    TTL. The sweep runs in a worker thread so a slow disk never stalls a stream, and
    a failing sweep is logged and retried at the next tick rather than ending the task.

    Args:
        period_seconds: Time between sweeps; the first sweep happens after one period,
            since the caller already swept at startup.
    """
    import asyncio

    from helioai.logging_config import get_logger

    while True:
        await asyncio.sleep(period_seconds)
        try:
            removed = await asyncio.to_thread(cleanup_old_runs)
            if removed:
                get_logger(__name__).info("workspace_cleanup", removed=removed)
        except Exception:
            get_logger(__name__).warning("workspace_cleanup_failed", exc_info=True)

Export

helioai.export

Export a session as a reproducible Jupyter notebook.

A research result that cannot be re-run is worthless. The agent already saves every sandbox run as code_N.py in the session workspace; this module bundles those runs plus the conversation into a self-contained, re-executable .ipynb with a provenance header (parameter ids, time, library versions).

The saved runs are rewritten to standalone code (to_standalone): load_data() becomes a fetch_series(...) call that wraps spz.get_data and blanks the declared fill value exactly as the session did, the agent-only param_card()/ document_method() calls are dropped, and a minimal header supplies the imports plus real clean()/export()/magnitude()/interp_to() helpers — so each cell runs in a plain Jupyter kernel with no HelioAI sandbox around it.

export_session_notebook

export_session_notebook(user_id: str, session_id: str, out_path: Path | None = None) -> Path

Write the session as a .ipynb and return its path.

Parameters:

Name Type Description Default
user_id str

Storage owner.

required
session_id str

Session to export.

required
out_path Path | None

Destination. Defaults to <workspace>/<label>.ipynb, falling back to the session id when the workspace has no label.

None

Returns:

Type Description
Path

The path written.

Example

export_session_notebook("cli", "8f3aa012-...") PosixPath('.../data/users/cli/workspace/plot-imf-bz_8f3aa0/plot-imf-bz_8f3aa0.ipynb')

The notebook opens with a provenance header (parameter ids, library versions), then one runnable cell per saved sandbox run, rewritten to standalone speasy calls.

Source code in helioai/export.py
def export_session_notebook(user_id: str, session_id: str, out_path: Path | None = None) -> Path:
    """Write the session as a .ipynb and return its path.

    Args:
        user_id: Storage owner.
        session_id: Session to export.
        out_path: Destination. Defaults to `<workspace>/<label>.ipynb`, falling
            back to the session id when the workspace has no label.

    Returns:
        The path written.

    Example:
        >>> export_session_notebook("cli", "8f3aa012-...")
        PosixPath('.../data/users/cli/workspace/plot-imf-bz_8f3aa0/plot-imf-bz_8f3aa0.ipynb')

        The notebook opens with a provenance header (parameter ids, library versions),
        then one runnable cell per saved sandbox run, rewritten to standalone speasy calls.
    """
    import nbformat as nbf

    nb = build_notebook(user_id, session_id)

    if out_path is None:
        from helioai.workspace import safe_id, user_home

        label = store.get_workspace_dir(user_id, session_id) or session_id
        root = user_home(user_id) / "workspace"
        root.mkdir(parents=True, exist_ok=True)
        out_path = root / f"{safe_id(label)}.ipynb"
    out_path = Path(out_path)
    nbf.write(nb, str(out_path))
    return out_path

to_standalone

to_standalone(code_src: str, manifest: dict, *, with_header: bool = True) -> str

Turn a saved sandbox run into standalone, re-executable code.

Strips agent-only calls, rewrites load_data() → fetch_series()/fetch_events(), and (unless embedded in a notebook that already has a setup cell) prepends the imports plus the real helpers it needs.

Parameters:

Name Type Description Default
code_src str

A code_N.py saved by the sandbox.

required
manifest dict

The session manifest from read_manifest() — provides the param_id and window behind each load_data() name.

required
with_header bool

Prepend the standalone imports/helpers header.

True
Example

manifest = {"datasets": {"imf_gsm": {"kind": "timeseries", ... "param_id": "amda/imf_gsm", "start": "2005-01-17", "stop": "2005-01-18"}}} to_standalone('d = load_data("imf_gsm")\nparam_card(d, "amda/imf_gsm")\n', ... manifest, with_header=False) 'd = fetch_series("amda/imf_gsm", "2005-01-17", "2005-01-18")'

Source code in helioai/export.py
def to_standalone(code_src: str, manifest: dict, *, with_header: bool = True) -> str:
    """Turn a saved sandbox run into standalone, re-executable code.

    Strips agent-only calls, rewrites load_data() → fetch_series()/fetch_events(), and
    (unless embedded in a notebook that already has a setup cell) prepends the imports
    plus the real helpers it needs.

    Args:
        code_src: A `code_N.py` saved by the sandbox.
        manifest: The session manifest from `read_manifest()` — provides the
            param_id and window behind each `load_data()` name.
        with_header: Prepend the standalone imports/helpers header.

    Example:
        >>> manifest = {"datasets": {"imf_gsm": {"kind": "timeseries",
        ...     "param_id": "amda/imf_gsm", "start": "2005-01-17", "stop": "2005-01-18"}}}
        >>> to_standalone('d = load_data("imf_gsm")\\nparam_card(d, "amda/imf_gsm")\\n',
        ...               manifest, with_header=False)
        'd = fetch_series("amda/imf_gsm", "2005-01-17", "2005-01-18")'
    """
    body = _strip_known_imports(
        _rewrite_load_data_calls(_strip_agent_only_calls(code_src), manifest)
    )
    if not with_header:
        return body
    header = _standalone_header(body)
    return f"{header}\n\n\n{body}" if header else body

Configuration

helioai.config

Centralized configuration — loads .env once at startup.

settings is a module-level singleton imported everywhere.

Importing this module never requires an API key. Credentials are validated where they are used, by llm.factory.build_llm_client, which already raised the same errors — the import-time copy only meant that surfaces needing no LLM at all could not start. The MCP server is exactly that: a pure tool provider whose client brings its own model, and it died on AZURE_OPENAI_API_KEY is not set before serving a single tool.

EXPERIMENTS module-attribute

EXPERIMENTS: frozenset[str] = frozenset({'deferred_tools', 'search_budget', 'search_variables', 'judgment_intent'})

The behaviours that change what the model sees or is told, each off by default.

Every one of them was committed on the strength of a single live run and never measured against the loop it replaced. They stay in the code as named experiments so that each can be switched on alone and compared on the same questions, N runs each (scripts/bench_live.py):

  • deferred_tools: the lead sees the formulary and catalogue tools only after asking for them with search_tools, and its prompt says so.
  • search_budget: a data_analyst past three lookups (a plasma_physicist past two) with nothing downloaded receives a correction listing the ids it already has.
  • search_variables: the top hit of a parameter search lists every variable of its dataset.
  • judgment_intent: with a judging backend (HELIOAI_JUDGMENT_BACKEND=jev), the user's question is read once by the judge, concurrently with the first model call, into an intent contract — deliverable, quantity, date, whether an uncertainty, a method or two spacecraft were asked for — placed against what the turn loaded and delivered (core/joins.py: frame, window, quantity, responsiveness) and emitted as an intent event after the answer. Observation: nothing in the loop acts on it.

final_answer was one of them and is now the default: on the third bench (34 clean runs, four configurations) it cost nothing on any question, re-enabled the claim verdict — zero contradictions over nine — and named rejected candidates as often as the loop it replaced. The two search experiments showed no value there (eight id-resolution runs correct without them); they stay off until a question that fails without them is recorded.

JUDGMENT_BACKENDS module-attribute

JUDGMENT_BACKENDS: frozenset[str] = frozenset({'null', 'jev'})

Who answers the judgment questions of helioai.core.judgment.

null abstains on every question, and it is the default: every call site runs the deterministic code it runs today when the answer is None, so a build with this backend is the loop as it was before the module existed. jev is TypeSafe's System One model through the optional judgment extra. There is no third "local" backend on purpose — the deterministic path is not a judge, it is the fallback, and naming it as a backend would suggest the two could disagree.

AzureOpenAIConfig dataclass

Azure OpenAI deployment settings.

Azure routes by deployment name rather than model name, and reasoning models (GPT-5, o-series) reject an explicit temperature — hence temperature=None by default, which omits the field entirely.

Source code in helioai/config.py
@dataclass
class AzureOpenAIConfig:
    """Azure OpenAI deployment settings.

    Azure routes by deployment name rather than model name, and reasoning models
    (GPT-5, o-series) reject an explicit temperature — hence `temperature=None`
    by default, which omits the field entirely.
    """

    deployment: str = "models-gpt-53-chat"
    api_version: str = "2024-12-01-preview"
    # 8192, not the 2048 this used to be and not the 4096 the other providers use.
    # Azure draws reasoning tokens from this same allowance, so a reasoning
    # deployment can spend the whole budget thinking and return an empty message —
    # no text, no tool call. At 2048 that happened on any request that generates a
    # file: the standalone-script export in examples/02 produced nothing at all.
    max_output_tokens: int = 8192
    temperature: float | None = None
    api_key: str = ""
    endpoint: str = ""

GeminiConfig dataclass

Google Gemini settings, used with the native google-genai client.

Source code in helioai/config.py
@dataclass
class GeminiConfig:
    """Google Gemini settings, used with the native `google-genai` client."""

    model: str = "gemini-2.5-flash"
    max_output_tokens: int = 4096
    temperature: float = 0.2
    api_key: str = ""

GroqConfig dataclass

Groq settings. Reached through the shared OpenAI-compatible client.

Source code in helioai/config.py
@dataclass
class GroqConfig:
    """Groq settings. Reached through the shared OpenAI-compatible client."""

    model: str = "llama-3.3-70b-versatile"
    max_output_tokens: int = 4096
    temperature: float = 0.2
    api_key: str = ""
    headers: dict[str, str] = field(default_factory=dict)

OpenCodeConfig dataclass

OpenCode's Zen gateway — OpenAI-compatible, whichever way you reach it: the flat-rate Go subscription, a BYOK-routed key, or any other model Zen hosts.

base_url defaults to the Go-plan endpoint, since that flat-rate tier is what most accounts actually have. It is a DIFFERENT catalogue from the general Zen endpoint (.../zen/v1, no /go/) — that one serves premium/BYOK-only models (Claude, ...) a Go subscription cannot reach, confirmed by querying both /models endpoints directly. Override HELIOAI_OPENCODE_URL to the plain Zen path if your access is not the Go plan.

No default model: what is reachable depends on your plan/BYOK setup and Zen's rotating catalogue (GLM, Kimi, DeepSeek, Qwen, MiniMax...). Set HELIOAI_OPENCODE_MODEL to the exact id from your dashboard — an empty string fails at the API with a clear "unknown model" rather than silently routing to a guessed default that may not exist on your plan.

16384 output tokens because everything this gateway serves is a reasoning model, and reasoning, prose AND tool-call arguments all draw on the same budget. At 4096, DeepSeek v4's run_python calls (~12k chars of JSON once a plot script is in them) were cut mid-string: the model saw "missing 1 required positional argument: 'code'", could not know why, and burned five turns re-sending the same truncated call. Same lesson as Azure's 2048→8192, one provider later.

Source code in helioai/config.py
@dataclass
class OpenCodeConfig:
    """OpenCode's Zen gateway — OpenAI-compatible, whichever way you reach it: the
    flat-rate Go subscription, a BYOK-routed key, or any other model Zen hosts.

    `base_url` defaults to the Go-plan endpoint, since that flat-rate tier is what
    most accounts actually have. It is a DIFFERENT catalogue from the general Zen
    endpoint (`.../zen/v1`, no `/go/`) — that one serves premium/BYOK-only models
    (Claude, ...) a Go subscription cannot reach, confirmed by querying both
    `/models` endpoints directly. Override `HELIOAI_OPENCODE_URL` to the plain Zen
    path if your access is not the Go plan.

    No default model: what is reachable depends on your plan/BYOK setup and Zen's
    rotating catalogue (GLM, Kimi, DeepSeek, Qwen, MiniMax...). Set
    `HELIOAI_OPENCODE_MODEL` to the exact id from your dashboard — an empty string
    fails at the API with a clear "unknown model" rather than silently routing to a
    guessed default that may not exist on your plan.

    16384 output tokens because everything this gateway serves is a reasoning model,
    and reasoning, prose AND tool-call arguments all draw on the same budget. At 4096,
    DeepSeek v4's run_python calls (~12k chars of JSON once a plot script is in them)
    were cut mid-string: the model saw "missing 1 required positional argument: 'code'",
    could not know why, and burned five turns re-sending the same truncated call.
    Same lesson as Azure's 2048→8192, one provider later.
    """

    base_url: str = "https://opencode.ai/zen/go"
    model: str = ""
    max_output_tokens: int = 16384
    temperature: float = 0.2
    api_key: str = ""
    headers: dict[str, str] = field(default_factory=dict)

OllamaConfig dataclass

Local Ollama settings.

Ollama serves an OpenAI-compatible API on /v1, so it needs no client of its own and no API key. Point base_url elsewhere for any other local endpoint.

Source code in helioai/config.py
@dataclass
class OllamaConfig:
    """Local Ollama settings.

    Ollama serves an OpenAI-compatible API on `/v1`, so it needs no client of its
    own and no API key. Point `base_url` elsewhere for any other local endpoint.
    """

    base_url: str = "http://localhost:11434"
    model: str = "qwen2.5:14b-instruct"
    max_output_tokens: int = 4096
    temperature: float = 0.2
    api_key: str = ""
    headers: dict[str, str] = field(default_factory=dict)

LLMConfig dataclass

Which provider to use, and the settings for each one.

opencode is the default because it is what the README tells a new user to set up: one key reaches hosted reasoning models. The default was azure until 0.4.0, an enterprise deployment no newcomer has, so the first error anyone saw named a key they had never heard of.

Source code in helioai/config.py
@dataclass
class LLMConfig:
    """Which provider to use, and the settings for each one.

    `opencode` is the default because it is what the README tells a new user to set up:
    one key reaches hosted reasoning models. The default was
    `azure` until 0.4.0, an enterprise deployment no newcomer has, so the first error
    anyone saw named a key they had never heard of.
    """

    provider: str = "opencode"
    azure: AzureOpenAIConfig = field(default_factory=AzureOpenAIConfig)
    gemini: GeminiConfig = field(default_factory=GeminiConfig)
    groq: GroqConfig = field(default_factory=GroqConfig)
    opencode: OpenCodeConfig = field(default_factory=OpenCodeConfig)
    ollama: OllamaConfig = field(default_factory=OllamaConfig)

AgentConfig dataclass

Agent loop limits, and which model each delegated role runs on.

max_iterations caps how many tool-calling rounds one question may take before the loop gives up, bounding both runtime and token spend.

role_models maps a sub-agent role to (provider, model). A parameter_hunter resolves ids from search results and needs no frontier model; a data_analyst writes the physics and does. Left empty, every role runs on the lead's client, as it always did. Parsed from HELIOAI_ROLE_MODELS="parameter_hunter=groq:llama-3.3-70b- versatile,data_analyst=opencode" — the model part is optional and defaults to the provider's configured model.

experiments names the behaviours of EXPERIMENTS that are switched on, from HELIOAI_EXPERIMENTS. Empty — the default — is the loop as it behaved before any of them existed.

Source code in helioai/config.py
@dataclass
class AgentConfig:
    """Agent loop limits, and which model each delegated role runs on.

    `max_iterations` caps how many tool-calling rounds one question may take
    before the loop gives up, bounding both runtime and token spend.

    `role_models` maps a sub-agent role to `(provider, model)`. A `parameter_hunter`
    resolves ids from search results and needs no frontier model; a `data_analyst`
    writes the physics and does. Left empty, every role runs on the lead's client, as
    it always did. Parsed from `HELIOAI_ROLE_MODELS="parameter_hunter=groq:llama-3.3-70b-
    versatile,data_analyst=opencode"` — the model part is optional and defaults to the
    provider's configured model.

    `experiments` names the behaviours of `EXPERIMENTS` that are switched on, from
    `HELIOAI_EXPERIMENTS`. Empty — the default — is the loop as it behaved before any of
    them existed.
    """

    max_iterations: int = 10
    role_models: dict[str, tuple[str, str | None]] = field(default_factory=dict)
    experiments: frozenset[str] = frozenset()

RAGConfig dataclass

Parameter search settings.

Retrieval is hybrid: dense embeddings for descriptions, BM25 for exact tokens like BGSEc, fused by Reciprocal Rank Fusion with parameter rrf_k.

There is no cross-encoder reranking stage, and that is a measured decision, not an omission: a generic MS MARCO cross-encoder (ms-marco-MiniLM-L-6-v2) was tried over the fused candidates and degraded results — trained on web prose, it discards the dense+sparse consensus that makes exact-code matching work. The plumbing sat disabled for a year and was removed; only a domain-tuned reranker would be worth adding back.

hybrid_fetch_k stays at 50 for the same kind of reason. Raising it to 100 or 200, or the BM25 cap alone to 100, was measured on 2026-09-22 against the 30 HelioBench n1 queries: MRR fell from 0.736 to 0.724, 0.712 and 0.729, because a wider pool lets an unpenalised candidate from deep in both lists climb over a relevant one that carries a small penalty, and two accepted products dropped out of the top-k altogether. The product the raise was meant to rescue (ssc/mms1, rank 54 in BM25) is reached through the dense channel instead, once ef_search looks wide enough.

index_repo is the Hugging Face dataset helioai index fetches a prebuilt index from when the local one is empty (helioai.index_snapshot). It is a setting so that a lab can point at its own mirror, and an empty value turns fetching off.

Source code in helioai/config.py
@dataclass
class RAGConfig:
    """Parameter search settings.

    Retrieval is hybrid: dense embeddings for descriptions, BM25 for exact tokens
    like `BGSEc`, fused by Reciprocal Rank Fusion with parameter `rrf_k`.

    There is no cross-encoder reranking stage, and that is a measured decision, not
    an omission: a generic MS MARCO cross-encoder (`ms-marco-MiniLM-L-6-v2`) was tried
    over the fused candidates and *degraded* results — trained on web prose, it
    discards the dense+sparse consensus that makes exact-code matching work. The
    plumbing sat disabled for a year and was removed; only a domain-tuned reranker
    would be worth adding back.

    `hybrid_fetch_k` stays at 50 for the same kind of reason. Raising it to 100 or 200,
    or the BM25 cap alone to 100, was measured on 2026-09-22 against the 30 HelioBench
    n1 queries: MRR fell from 0.736 to 0.724, 0.712 and 0.729, because a wider pool lets
    an unpenalised candidate from deep in both lists climb over a relevant one that
    carries a small penalty, and two accepted products dropped out of the top-k
    altogether. The product the raise was meant to rescue (`ssc/mms1`, rank 54 in BM25)
    is reached through the dense channel instead, once `ef_search` looks wide enough.

    `index_repo` is the Hugging Face dataset `helioai index` fetches a prebuilt index
    from when the local one is empty (`helioai.index_snapshot`). It is a setting so that
    a lab can point at its own mirror, and an empty value turns fetching off.
    """

    chroma_dir: Path = field(default_factory=lambda: _DATA / "chroma")
    collection_name: str = "speasy_catalog"
    catalogs_collection_name: str = "speasy_catalogs"
    embed_model: str = "sentence-transformers/all-MiniLM-L6-v2"
    hybrid_enabled: bool = True
    hybrid_fetch_k: int = 50
    rrf_k: int = 60
    index_repo: str = "erdoganfurkan/helioai-speasy-index"

WorkspaceConfig dataclass

Retention of per-session working directories.

Their location is not a setting: workspace.user_home derives it from data_dir per user. A workspace_dir field (and HELIOAI_WORKSPACE) used to sit here, read from the environment and consumed by nothing since storage became per-user.

Source code in helioai/config.py
@dataclass
class WorkspaceConfig:
    """Retention of per-session working directories.

    Their location is not a setting: `workspace.user_home` derives it from `data_dir`
    per user. A `workspace_dir` field (and `HELIOAI_WORKSPACE`) used to sit here, read
    from the environment and consumed by nothing since storage became per-user.
    """

    ttl_seconds: int = 86400 * 7  # 7 days

ProfileConfig dataclass

Legacy location of the single-user profile, kept for helioai migrate-storage.

The agent reads users/<user>/profile.md (workspace.user_home) since storage became per user; helioai profile, %helioai_profile and the web UI all edit that file. Nothing injects this path any more, so HELIOAI_PROFILE — which only moved it — was a knob that did nothing, and is gone. The path stays as the place the migration looks for a profile written by an older install.

Source code in helioai/config.py
@dataclass
class ProfileConfig:
    """Legacy location of the single-user profile, kept for `helioai migrate-storage`.

    The agent reads `users/<user>/profile.md` (`workspace.user_home`) since storage
    became per user; `helioai profile`, `%helioai_profile` and the web UI all edit that
    file. Nothing injects this path any more, so `HELIOAI_PROFILE` — which only moved
    it — was a knob that did nothing, and is gone. The path stays as the place the
    migration looks for a profile written by an older install.
    """

    profile_path: Path = field(default_factory=lambda: _DATA / "profile.md")

RecipesConfig dataclass

Where scientific recipes are loaded from.

Defaults to the copy shipped inside the package so pip install works; override with HELIOAI_RECIPES_DIR to use your own set.

Source code in helioai/config.py
@dataclass
class RecipesConfig:
    """Where scientific recipes are loaded from.

    Defaults to the copy shipped inside the package so `pip install` works;
    override with `HELIOAI_RECIPES_DIR` to use your own set.
    """

    recipes_dir: Path = field(default_factory=lambda: _PKG_RECIPES)

CatalogsConfig dataclass

Where user-saved event catalogs are written, in speasy format.

Source code in helioai/config.py
@dataclass
class CatalogsConfig:
    """Where user-saved event catalogs are written, in speasy format."""

    catalogs_dir: Path = field(default_factory=lambda: _DATA / "catalogs")

LiteratureConfig dataclass

NASA ADS credentials for find_papers. Free token, no key means no tool.

Source code in helioai/config.py
@dataclass
class LiteratureConfig:
    """NASA ADS credentials for `find_papers`. Free token, no key means no tool."""

    ads_token: str = ""

MCPConfig dataclass

Remote MCP servers to mount, plus auth for HelioAI's own MCP HTTP transport.

Source code in helioai/config.py
@dataclass
class MCPConfig:
    """Remote MCP servers to mount, plus auth for HelioAI's own MCP HTTP transport."""

    servers_json: str = ""
    # Shared secret required as `Authorization: Bearer <token>` on the HTTP transport
    # (stdio needs none — the client owns the process). Empty (default) → no auth
    # required on loopback; a non-loopback bind then refuses to start (mcp_server.main).
    token: str = ""

VisionConfig dataclass

Multimodal review of generated figures.

A stateless side-call outside the agent loop: the image is downscaled, sent once, and only the text verdict enters the history — never the image, which would otherwise be resent on every subsequent turn. Off by default.

Source code in helioai/config.py
@dataclass
class VisionConfig:
    """Multimodal review of generated figures.

    A stateless side-call outside the agent loop: the image is downscaled, sent
    once, and only the text verdict enters the history — never the image, which
    would otherwise be resent on every subsequent turn. Off by default.
    """

    # Reviews sandbox figures with a multimodal side-call; only the text
    # verdict enters the history, never the image.
    enabled: bool = False
    provider: str = "azure"
    model: str = ""
    timeout_s: float = 20.0

JudgmentConfig dataclass

A System One judge beside the loop, in observation.

Every question HelioAI asks it lives in helioai.core.judgment, so what is delegated to a proprietary model can be audited in one file. The backend is one axis; which sites may ask is the other, one EXPERIMENTS name per site — so "the judge does not help here" and "the layer costs something" can be told apart. Off by default, and nothing it answers corrects the model: it annotates and it is recorded.

Source code in helioai/config.py
@dataclass
class JudgmentConfig:
    """A System One judge beside the loop, in observation.

    Every question HelioAI asks it lives in `helioai.core.judgment`, so what is delegated
    to a proprietary model can be audited in one file. The backend is one axis; which
    sites may ask is the other, one `EXPERIMENTS` name per site — so "the judge does not
    help here" and "the layer costs something" can be told apart. Off by default, and
    nothing it answers corrects the model: it annotates and it is recorded.
    """

    backend: str = "null"
    model: str = "jev-latest"
    timeout_s: float = 2.0
    api_key: str = ""

DevConfig dataclass

Shared secret unlocking unrestricted mode past the heliophysics guardrail.

Empty by default, which means no token is valid and every request stays scoped. Compared in constant time.

Source code in helioai/config.py
@dataclass
class DevConfig:
    """Shared secret unlocking unrestricted mode past the heliophysics guardrail.

    Empty by default, which means no token is valid and every request stays
    scoped. Compared in constant time.
    """

    # Shared-secret that unlocks unrestricted LLM access (bypasses scope guardrail).
    # Empty (default) → no token is valid → all requests stay restricted.
    token: str = ""

WebAuthConfig dataclass

Nominative tokens for the web UI, parsed from HELIOAI_USERS.

Empty means no authentication and a single local user, which is the intended behaviour for local development only.

Source code in helioai/config.py
@dataclass
class WebAuthConfig:
    """Nominative tokens for the web UI, parsed from `HELIOAI_USERS`.

    Empty means no authentication and a single local user, which is the intended
    behaviour for local development only.
    """

    # Nominative tokens for the web UI: {token: user_id}. Parsed from
    # HELIOAI_USERS="tok1:vincent,tok2:alice". Empty → no auth, single local user.
    # ponytail: env-driven map, fine for a handful of researchers; move to a DB
    # table if tokens must be added/revoked at runtime.
    users: dict[str, str] = field(default_factory=dict)
    # `serve --web` refuses a non-loopback bind with no users configured, the way the
    # MCP HTTP server refuses one without a token: run_python is arbitrary code
    # execution. The one legitimate exception is a container, which must bind 0.0.0.0
    # inside its own network namespace while the host publishes the port on loopback —
    # docker-compose.yml sets HELIOAI_ALLOW_UNAUTHENTICATED_PUBLIC=1 for exactly that.
    allow_unauthenticated_public: bool = False

Settings dataclass

Root settings object.

Imported as the module-level settings singleton and read everywhere; built once at import by _load(). Credentials are not checked here — importing must work with no key at all — but in llm.factory.build_llm_client, where the selected provider is actually used.

Source code in helioai/config.py
@dataclass
class Settings:
    """Root settings object.

    Imported as the module-level `settings` singleton and read everywhere; built
    once at import by `_load()`. Credentials are not checked here — importing must
    work with no key at all — but in `llm.factory.build_llm_client`, where the
    selected provider is actually used.
    """

    data_dir: Path = field(default_factory=lambda: _DATA)
    llm: LLMConfig = field(default_factory=LLMConfig)
    agent: AgentConfig = field(default_factory=AgentConfig)
    rag: RAGConfig = field(default_factory=RAGConfig)
    workspace: WorkspaceConfig = field(default_factory=WorkspaceConfig)
    profile: ProfileConfig = field(default_factory=ProfileConfig)
    recipes: RecipesConfig = field(default_factory=RecipesConfig)
    catalogs: CatalogsConfig = field(default_factory=CatalogsConfig)
    literature: LiteratureConfig = field(default_factory=LiteratureConfig)
    mcp: MCPConfig = field(default_factory=MCPConfig)
    vision: VisionConfig = field(default_factory=VisionConfig)
    judgment: JudgmentConfig = field(default_factory=JudgmentConfig)
    dev: DevConfig = field(default_factory=DevConfig)
    web_auth: WebAuthConfig = field(default_factory=WebAuthConfig)

validate_experiments

validate_experiments(experiments: frozenset[str] | None = None) -> frozenset[str]

Refuse unknown experiment names — an error, not a warning.

An experiment that silently does nothing would be measured as if it did, and the comparison would be wrong without anyone knowing.

Parameters:

Name Type Description Default
experiments frozenset[str] | None

The set to check; settings.agent.experiments when None.

None

Returns:

Type Description
frozenset[str]

The same set, when every name is known.

Raises:

Type Description
RuntimeError

Naming the unknown entries and the known ones.

Source code in helioai/config.py
def validate_experiments(experiments: frozenset[str] | None = None) -> frozenset[str]:
    """Refuse unknown experiment names — an error, not a warning.

    An experiment that silently does nothing would be measured as if it did, and the
    comparison would be wrong without anyone knowing.

    Args:
        experiments: The set to check; `settings.agent.experiments` when None.

    Returns:
        The same set, when every name is known.

    Raises:
        RuntimeError: Naming the unknown entries and the known ones.
    """
    if experiments is None:
        experiments = settings.agent.experiments
    unknown = experiments - EXPERIMENTS
    if unknown:
        raise RuntimeError(
            f"HELIOAI_EXPERIMENTS: unknown {sorted(unknown)}; known: {sorted(EXPERIMENTS)}"
        )
    return experiments

validate_judgment

validate_judgment(cfg: JudgmentConfig | None = None) -> str

Refuse an unknown judgment backend — an error, not a warning.

Same reasoning as validate_experiments: a backend name that silently fell back to abstaining would be measured as if it had judged. Called where the API key is already checked, never at import.

Parameters:

Name Type Description Default
cfg JudgmentConfig | None

The config to check; settings.judgment when None.

None

Returns:

Type Description
str

The backend name, when known.

Raises:

Type Description
RuntimeError

Naming the unknown backend and the known ones.

Source code in helioai/config.py
def validate_judgment(cfg: JudgmentConfig | None = None) -> str:
    """Refuse an unknown judgment backend — an error, not a warning.

    Same reasoning as `validate_experiments`: a backend name that silently fell back to
    abstaining would be measured as if it had judged. Called where the API key is
    already checked, never at import.

    Args:
        cfg: The config to check; `settings.judgment` when None.

    Returns:
        The backend name, when known.

    Raises:
        RuntimeError: Naming the unknown backend and the known ones.
    """
    backend = (cfg or settings.judgment).backend
    if backend not in JUDGMENT_BACKENDS:
        raise RuntimeError(
            f"HELIOAI_JUDGMENT_BACKEND: unknown {backend!r}; known: {sorted(JUDGMENT_BACKENDS)}"
        )
    return backend

dev_unlock

dev_unlock(supplied: str | None) -> bool

True iff the supplied token matches the configured dev secret.

Parameters:

Name Type Description Default
supplied str | None

Token offered by the caller — a --dev flag or a request header. Compared with hmac.compare_digest, so a wrong token costs the same time as a right one.

required

Returns:

Type Description
bool

True only when a dev token is configured and the supplied one matches.

bool

An unconfigured instance answers False for every input, including

bool

None and the empty string, so a blank secret cannot unlock anything.

Source code in helioai/config.py
def dev_unlock(supplied: str | None) -> bool:
    """True iff the supplied token matches the configured dev secret.

    Args:
        supplied: Token offered by the caller — a `--dev` flag or a request
            header. Compared with `hmac.compare_digest`, so a wrong token costs
            the same time as a right one.

    Returns:
        True only when a dev token is configured and the supplied one matches.
        An unconfigured instance answers False for every input, including
        `None` and the empty string, so a blank secret cannot unlock anything.
    """
    return (
        bool(settings.dev.token)
        and supplied is not None
        and hmac.compare_digest(supplied, settings.dev.token)
    )

Logging

helioai.logging_config

Structured logging via structlog.

Output format is selected by HELIOAI_LOG_FORMAT
  • console (default): human-friendly, colourised.
  • json: one JSON object per line.

setup_logging

setup_logging(level: str | int = 'INFO', *, tracebacks: bool = True) -> None

Configure structlog and the root logger.

Output format follows HELIOAI_LOG_FORMAT: console (default) or json. Safe to call more than once — every entry point calls it, and repeated calls replace the handler rather than stacking duplicates.

HELIOAI_LOG_LEVEL overrides level. Every entry point hardcodes its own, so without this there is no way to quiet a third party that logs at the same level — speasy's inventory probes warn loudly on a provider it then disables, which is noise in a recorded session or a demo. An unrecognised value is ignored rather than obeyed: a typo must not silently turn logging up.

tracebacks=False is for the terminal a person is reading: an error keeps its one log line, with the exception's type and message, but not the stack under it — the interface prints what to do instead. DEBUG brings the stack back, whatever the flag.

Parameters:

Name Type Description Default
level str | int

Log level name or numeric value. Unknown names fall back to INFO.

'INFO'
tracebacks bool

Render the stack of a logged exception.

True
Source code in helioai/logging_config.py
def setup_logging(level: str | int = "INFO", *, tracebacks: bool = True) -> None:
    """Configure structlog and the root logger.

    Output format follows `HELIOAI_LOG_FORMAT`: `console` (default) or `json`.
    Safe to call more than once — every entry point calls it, and repeated calls
    replace the handler rather than stacking duplicates.

    `HELIOAI_LOG_LEVEL` overrides `level`. Every entry point hardcodes its own,
    so without this there is no way to quiet a third party that logs at the same
    level — speasy's inventory probes warn loudly on a provider it then disables,
    which is noise in a recorded session or a demo. An unrecognised value is
    ignored rather than obeyed: a typo must not silently turn logging up.

    `tracebacks=False` is for the terminal a person is reading: an error keeps its one
    log line, with the exception's type and message, but not the stack under it — the
    interface prints what to do instead. DEBUG brings the stack back, whatever the flag.

    Args:
        level: Log level name or numeric value. Unknown names fall back to INFO.
        tracebacks: Render the stack of a logged exception.
    """
    override = os.environ.get("HELIOAI_LOG_LEVEL", "").strip().upper()
    if override and isinstance(getattr(logging, override, None), int):
        level = getattr(logging, override)

    if isinstance(level, str):
        level = getattr(logging, level.upper(), logging.INFO)

    fmt = _format_from_env()

    shared_processors: list[Any] = [
        structlog.contextvars.merge_contextvars,
        structlog.processors.add_log_level,
        structlog.processors.TimeStamper(fmt="iso", utc=True),
        structlog.processors.StackInfoRenderer(),
    ]
    if not tracebacks and level > logging.DEBUG:
        shared_processors.append(_exception_as_one_line)
    shared_processors.append(structlog.processors.format_exc_info)

    if fmt == "json":
        renderer: Any = structlog.processors.JSONRenderer()
    else:
        renderer = structlog.dev.ConsoleRenderer(colors=sys.stderr.isatty())

    structlog.configure(
        processors=[
            *shared_processors,
            structlog.stdlib.ProcessorFormatter.wrap_for_formatter,
        ],
        wrapper_class=structlog.make_filtering_bound_logger(level),
        logger_factory=structlog.stdlib.LoggerFactory(),
        cache_logger_on_first_use=True,
    )

    formatter = structlog.stdlib.ProcessorFormatter(
        foreign_pre_chain=shared_processors,
        processors=[
            structlog.stdlib.ProcessorFormatter.remove_processors_meta,
            renderer,
        ],
    )

    handler = logging.StreamHandler(sys.stderr)
    handler.setFormatter(formatter)
    root = logging.getLogger()
    root.handlers = [handler]
    root.setLevel(level)
    _quiet_third_party_advisories()

get_logger

get_logger(name: str | None = None) -> Any

Return a structlog logger, optionally bound to a module name.

Parameters:

Name Type Description Default
name str | None

Usually __name__. Omitted, the logger carries no module field.

None

Returns:

Type Description
Any

A structlog bound logger. Its output format follows

Any

HELIOAI_LOG_FORMAT (console or json), decided at configuration time

Any

rather than here.

Source code in helioai/logging_config.py
def get_logger(name: str | None = None) -> Any:
    """Return a structlog logger, optionally bound to a module name.

    Args:
        name: Usually `__name__`. Omitted, the logger carries no module field.

    Returns:
        A structlog bound logger. Its output format follows
        `HELIOAI_LOG_FORMAT` (console or json), decided at configuration time
        rather than here.
    """
    return structlog.get_logger(name) if name else structlog.get_logger()