From 753349700b39cf5d3fd605f584176315c97b2c79 Mon Sep 17 00:00:00 2001 From: Steven Enamakel <31011319+senamakel@users.noreply.github.com> Date: Mon, 6 Jul 2026 18:24:00 -0700 Subject: [PATCH] feat(inference): migrate inference callers onto tinyagents ChatModel (Phase 1 + 3a + 5 groundwork) (#4614) --- docs/tinyagents-inference-migration-plan.md | 158 +++++++++++ src/openhuman/accessibility/automate.rs | 11 +- .../agent/harness/session/builder/factory.rs | 39 +-- src/openhuman/agent_meetings/summary.rs | 19 +- src/openhuman/inference/provider/factory.rs | 71 +++++ .../inference/provider/factory_tests.rs | 41 +++ src/openhuman/inference/provider/mod.rs | 2 + src/openhuman/learning/linkedin_enrichment.rs | 24 +- src/openhuman/learning/reflection.rs | 30 +- src/openhuman/learning/reflection_tests.rs | 16 +- src/openhuman/memory/chat.rs | 65 +++-- .../memory_tree/tree_runtime/engine.rs | 49 ++-- .../memory_tree/tree_runtime/engine_tests.rs | 66 +++-- src/openhuman/memory_tree/tree_runtime/ops.rs | 29 +- src/openhuman/orchestration/ops.rs | 20 +- src/openhuman/tinyagents/mod.rs | 71 +++-- src/openhuman/tinyagents/model.rs | 263 +++++++++++++++++- src/openhuman/tinyagents/routes.rs | 8 +- src/openhuman/tinyagents/tests.rs | 10 +- 19 files changed, 810 insertions(+), 182 deletions(-) create mode 100644 docs/tinyagents-inference-migration-plan.md diff --git a/docs/tinyagents-inference-migration-plan.md b/docs/tinyagents-inference-migration-plan.md new file mode 100644 index 000000000..fbe03e3fe --- /dev/null +++ b/docs/tinyagents-inference-migration-plan.md @@ -0,0 +1,158 @@ +# Migrating `src/openhuman/inference/` onto `vendor/tinyagents/` + +**Status:** plan (not started) +**Relates to:** #4249 (tinyagents migration), `docs/tinyagents-migration-spec.md` (Phase 2 — model registry / profiles / fallback), `docs/tinyagents-full-migration-plan/`, `docs/tinyagents-drift-ledger.md` +**tinyagents:** 1.7.1, vendored as a git submodule at `vendor/tinyagents/` and path-patched over the crates.io pin (root `Cargo.toml`: `tinyagents = { version = "1.7", features = ["sqlite", "repl"] }` + `[patch.crates-io] tinyagents = { path = "vendor/tinyagents" }`). + +--- + +## 1. Why + +The agent loop already runs on tinyagents; `run_turn_via_tinyagents_shared` drives every turn. But the **model layer underneath it is still entirely in-house**: the harness reaches a crate `ChatModel` only through the `ProviderModel` adapter (`src/openhuman/tinyagents/model.rs`), which wraps openhuman's own `Box` stack from `src/openhuman/inference/provider/` — ~29.6k lines that re-implement what tinyagents 1.7 now ships natively: + +| openhuman (`inference/provider/`, etc.) | tinyagents 1.7 equivalent | +| ------------------------------------------------------------------------ | ------------------------------------------------------------------------------------------------------ | +| `Provider` trait, `ChatMessage`/`ChatRequest`/`ChatResponse`/`ProviderDelta` (`traits.rs`) | `harness::model::{ChatModel, ModelRequest, ModelResponse, ModelStream, ModelStreamItem}` + `harness::message::*` | +| `compatible*.rs` (~15 files: OpenAI-compat client, SSE streaming, parse, dump, repeat, timeout) | `harness::providers::openai` (convert / sse / transport / types) — serves every OpenAI-compatible endpoint | +| `reliable.rs` retry/backoff wrapper | `harness::retry::{RetryPolicy, FallbackPolicy, RateLimiter, is_retryable}` | +| `router.rs` `RouterProvider` hint-table multi-model routing | `harness::model::ModelRegistry` (per-depth model resolution) | +| `model_context.rs` context-window table | `MODEL_CONTEXT_PATTERNS` fallback + `registry::ModelCatalog` (offline snapshot: windows, pricing, capability flags) | +| `error_classify.rs` / `error_code.rs` (retryability, HTTP status parsing) | `harness::retry::is_retryable` + `harness::model::ProviderError` | +| provider-string → provider construction (`factory.rs`, partially) | `harness::providers::{ProviderKind, ProviderSpec}` factory (`ProviderKind::infer`, `Compatible`) | +| embeddings dispatch (`ops.rs` → local/cloud) | `harness::embeddings` traits (+ openai impl) | +| `provider/openai_codex.rs`, `openhuman_backend.rs`, `claude_agent_sdk/` | no equivalent — stay as host `ChatModel` impls | + +Maintaining both stacks means every fix (streaming edge case, retry policy, context-window entry, provider quirk) lands twice, and the adapter seam (`ProviderModel` + `ThinkingForwarder` + usage-carry plumbing) exists only to translate between two isomorphic type systems. Since `vendor/tinyagents` is our own crate (same GPL-3.0 license, same org), host-agnostic gaps get **upstreamed into the crate**; only genuinely openhuman-specific glue stays. + +**End state:** the crate's `ChatModel`/`ModelRequest` is the native model interface everywhere in openhuman; `inference/` keeps only host concerns (RPC surface, config, local runtime management, voice, OAuth, `/v1` endpoint); `ProviderModel` and the entire `Provider` trait stack are deleted. + +### Blast radius + +- 170 files outside `inference/` import `openhuman::inference`; 151 import `inference::provider` specifically. Top consumers: agent harness/session/tools/triage, `context`, `voice`, `routing`, `memory_tree`, `learning`, `channels`, `embeddings`, `subconscious`, `screen_intelligence`, `threads`, `migrations`, `config/schema`. +- Module sizes: `provider/` ≈ 29.6k lines, `local/` ≈ 13.7k, `voice/` + `http/` + `openai_oauth/` ≈ 4.4k, root files ≈ 5.5k. Only `provider/` + parts of the root files migrate; the rest stays. + +--- + +## 2. Disposition map + +### Migrates (replaced by crate, or upstreamed into crate) + +| Component | Destination | Notes | +| --- | --- | --- | +| `provider/traits.rs` — `Provider` trait + request/response/delta types | **delete**, use crate `ChatModel` + message types | The big inversion (Phase 1). `UsageInfo` → crate `Usage` (see gap G1 on USD/cached tokens). | +| `provider/compatible*.rs` — OpenAI-compat client | **delete**, use crate `providers::openai` | Gap audit first (Phase 2): request-dump debugging, repeat-detection, per-request timeout, BYOK auth styles. | +| `provider/reliable.rs` | **delete**, use crate `RetryPolicy` (+ `FallbackPolicy`) | Also resolves the known double-retry (reliable.rs *and* harness-level retry both fire today). | +| `provider/router.rs` (`RouterProvider`) | **delete**, use crate `ModelRegistry` | Hint table (`reasoning-v1`, `agentic-v1`, …) becomes registry aliases. | +| `model_context.rs` (`context_window_for_model`) | **upstream** entries into crate `ModelCatalog` snapshot / `MODEL_CONTEXT_PATTERNS`; host keeps a thin lookup that prefers config overrides | Crate catalog also carries pricing + capability flags — feeds the cost tracker. | +| `provider/error_classify.rs`, `error_code.rs` | **delete**, use crate `is_retryable` / `ProviderError` | Upstream any status-classification the crate misses. | +| `provider/temperature.rs` (`@` suffix) | host parses the suffix in the factory; value rides `ModelRequest` params | Grammar stays host-side; the plumbing type goes away. | +| `provider/config_rejection.rs`, `billing_error.rs` | **split**: generic classification upstreamed as crate error kinds; openhuman semantics (Sentry demotion, budget messaging) stay as a host classifier over `TinyAgentsError` | | +| `provider/factory.rs` | **shrinks, stays**: resolves openhuman provider strings (`openhuman`, `cloud`, `ollama:`, `:[@temp]`) + config + credentials → crate `ProviderSpec`/`Arc` | This is the host↔crate boundary after migration. `BYOK_INCOMPLETE_SENTINEL` stays. | +| `provider/ops.rs` (`list_configured_models`, SessionExpired publishing) | **stays**, retargeted to crate types | SessionExpired needs an auth-failure signal from the crate client (gap G3). | +| embeddings dispatch | crate `harness::embeddings` traits (seam `tinyagents/embeddings.rs` already exists — finish it) | Local (Ollama) embedding stays a host impl of the crate trait. | +| `provider/thread_context.rs`, `resolved_route.rs`, `auth_error_registry.rs` | **re-home** into `src/openhuman/tinyagents/` (they're seam concerns, not provider concerns) | `thread_context` task-locals already consumed by `model.rs`. | + +### Stays in `inference/` (host concerns, out of scope for the crate) + +- **RPC surface**: `schemas.rs`, `ops.rs`, `local/schemas.rs` — all `inference.*` controllers, legacy aliases. +- **Local runtime management** (`local/`): Ollama/LM Studio detect/spawn/adopt, Whisper/Piper install, download progress, model artifacts, context floor. The *chat client* to Ollama/LM Studio migrates to crate `ProviderKind::Ollama`/`Compatible`; process lifecycle stays. +- **`voice/`**: STT/TTS inference impls (whisper-cpp bindings, Piper) — not LLM-shaped; unchanged. +- **`openai_oauth/`**: Codex OAuth PKCE + encrypted token store (credentials domain integration). +- **`http/`**: the `/v1/chat/completions` OpenAI-compat *server* endpoint. (Later option: crate 1.3+ has OpenAI-compat runtime model listing; revisit after Phase 5.) +- **`device.rs`, `presets.rs`, `model_ids.rs`, `paths.rs`, `parse.rs`**: hardware profiles, preset tiers, config-derived model-id resolution, artifact paths. +- **`sentiment.rs`**: stays as a host op, but its model call is rewritten onto `ChatModel` (+ crate structured output) in Phase 6. +- **Bespoke providers** (`openhuman_backend.rs` session-JWT managed backend, `claude_agent_sdk/` subprocess, `openai_codex.rs` OAuth-token compat variant): stay in-repo, reimplemented as crate `ChatModel` impls (Phase 4). + +--- + +## 3. Known crate gaps to close first (upstream work in `vendor/tinyagents`) + +Verified against 1.7.1 source; re-audit at Phase 0 since the crate moves fast. + +- **G1 — Usage fidelity**: crate `Usage` has no `charged_amount_usd` and needs verifying for cache-read/cache-write token fields. The $0-cost-turn bug (fixed host-side via `cost::catalog::estimate_cost_usd`) shows exactly what breaks when this is lossy. Either upstream optional cost/cached fields on `Usage`, or keep host-side estimation keyed off the crate `ModelCatalog` pricing. +- **G2 — Tool-call start metadata**: crate `ToolDelta` carries `call_id`/`content` but no `tool_name`; the UI timeline's tool-start event is still forwarded out-of-band (`model.rs` forwarder). Upstream `tool_name` on the first `ToolDelta` (or a dedicated start item) so the forwarder can die with `ProviderModel`. +- **G3 — Auth-failure signal**: openhuman publishes `DomainEvent::SessionExpired` when a chat attempt fails auth. The crate client must classify 401/expired distinctly (as a `ProviderError` kind) so the host factory can hook it without string-sniffing. +- **G4 — Request dump / wire observability**: `compatible_dump.rs` writes raw request/response dumps for debugging. Upstream a transport-level hook (or confirm the crate's observability exporters cover it) before deleting. +- **G5 — Per-request timeout policy**: `compatible_timeout.rs` semantics vs. what the crate transport exposes. Upstream a per-call timeout on `ModelRequest`/`ProviderSpec` if missing. +- **G6 — BYOK auth styles**: the provider preset catalog (`docs/inference-provider-catalog.md`) supports multiple `AuthStyle`s (headers etc.). Confirm crate `ProviderSpec` can express every style in the catalog; upstream what's missing. +- **G7 — Repeat-output guard**: `compatible_repeat.rs` (degenerate-repetition detection). Decide: upstream as an optional stream guard, or accept the loss (note #4463 already tracks deleted repeat guards). + +Upstream flow: change in `vendor/tinyagents` (submodule working tree) → PR to `tinyhumansai/tinyagents` → publish → bump the crates.io pin in **both** Cargo worlds (root + `app/src-tauri`) and the submodule ref in lockstep. Nothing in openhuman may depend on unpublished vendored-only API at a merge point. + +--- + +## 4. Phases + +Each phase compiles green in both Cargo worlds, keeps ≥80% diff coverage, and lands as its own PR-sized slice. Provider-string grammar, RPC names, and observable UI behavior (streaming, cost footer, tool timeline) are parity-locked throughout. + +### Phase 0 — Inventory & gap re-audit +- Enumerate every consumer of `Provider` / `ChatRequest` / `ChatResponse` / `ChatMessage` outside `inference/` (151 files) and bucket them: (a) goes through the seam already, (b) direct one-shot `provider.chat(...)` callers (learning, memory, subconscious, screen_intelligence, sentiment, triage…), (c) type-only imports. +- Re-verify §3 gaps against current crate HEAD; file crate issues; update `docs/tinyagents-sdk-gaps.md`. +- Golden-transcript capture: record request/response wire dumps for the BYOK catalog matrix + Ollama + openhuman backend on the current stack, as fixtures for Phase 2 parity. + +**Exit:** disposition table confirmed per-file; crate gap PRs filed. + +### Phase 1 — Model-layer inversion (the pivot) +- Introduce `create_chat_model(...) -> Arc` in `factory.rs` alongside `create_chat_provider`, initially wrapping the existing stack via `ProviderModel` (zero behavior change). +- Migrate consumers bucket-by-bucket from `Box` to `Arc`: one-shot callers first (they use `chat()` once — mechanical: `ChatRequest` → `ModelRequest::new(...)`), then the seam (`run_turn_via_tinyagents_shared` takes the model directly — delete the wrap at the call sites), then streaming consumers. +- Keep a temporary reverse adapter (`ChatModel` → `Provider`) only if a consumer can't move in one slice; delete it before phase exit. + +**Exit:** no caller outside `inference/provider/` names the `Provider` trait; `ProviderModel` is constructed in exactly one place (factory). + +### Phase 2 — OpenAI-compatible client swap +- Behind the factory, construct crate `providers::openai` clients (via `ProviderSpec`) instead of `CompatibleProvider` for: BYOK slugs, Ollama, LM Studio, the `cloud` slug. +- Parity-test against Phase 0 golden dumps: request shape (tools JSON, temperature, multimodal blocks), SSE streaming (text, reasoning, tool-arg deltas), usage extraction, error mapping. Note recent regression surface: tool-calling defaults to JSON with P-Format opt-in (9b84f9684) — the crate path must honor the same default. +- Delete `compatible*.rs` (15 files, the bulk of the 29.6k lines) once all construction sites are switched. + +**Exit:** no `compatible*.rs` left; wire parity fixtures green; `pnpm test:rust` + `json_rpc_e2e` green. + +### Phase 3 — Reliability, routing, model metadata +- Replace `ReliableProvider` layering (`session/builder/factory.rs` re-layered it in the P1 parity fix) with crate `RetryPolicy` at the client level; audit and remove the double-retry. +- Replace `RouterProvider` with `ModelRegistry`: abstract tier names (`reasoning-v1`, `coding-v1`, …) become registry entries resolved per call; `provider_for_role`/workload resolution feeds the registry instead of building a router provider. +- Upstream openhuman's context-window table entries into the crate `ModelCatalog` snapshot; `context_window_for_model` becomes a host shim: config override → catalog → crate pattern fallback. Wire catalog pricing into `cost::catalog` (replacing the hand-rolled rate table, or seeding it). + +**Exit:** `reliable.rs`, `router.rs`, `model_context.rs` deleted; retry fires exactly once per layer by design. + +### Phase 4 — Bespoke providers as `ChatModel` impls +- Rewrite `openhuman_backend.rs` (managed backend, session JWT + SessionExpired publishing via G3), `openai_codex.rs` (Codex OAuth token source over the crate openai client), and `claude_agent_sdk/` (subprocess protocol) as direct `ChatModel` implementations in their current homes. +- Local runtime: `local/` keeps process lifecycle; its chat/vision/embed entrypoints call crate clients pointed at the local base URL. + +**Exit:** `Provider` trait has zero implementations → delete `traits.rs` and the trait itself. + +### Phase 5 — Seam shrink +- Delete `ProviderModel`, `ThinkingForwarder` remnants, `ProviderUsageCarry`, and the `ChatMessage`↔crate-message conversion layer in `tinyagents/convert.rs` (the harness now receives crate types natively). +- Re-home `thread_context.rs` / `resolved_route.rs` / `auth_error_registry.rs` into `src/openhuman/tinyagents/`. +- `inference/provider/` collapses to: `factory.rs` (string grammar → `ChatModel`), bespoke impls, host error classifier, `ops.rs`, `schemas.rs`. + +**Exit:** `src/openhuman/tinyagents/model.rs` deleted; adapter inventory test updated. + +### Phase 6 — One-shot inference ops onto the crate +- `sentiment.rs`, `should_react`, `summarize`, vision prompts, triage-style single calls: rewrite onto `ChatModel` + crate structured output (`harness/structured`) instead of hand-rolled parse (`parse.rs` shrinks or dies). +- Embeddings: finish `tinyagents/embeddings.rs` — cloud via crate openai embeddings, local as a host impl of the crate trait; `LocalAiEmbeddingResult` maps from crate types. + +**Exit:** no ad-hoc prompt/parse loops outside the crate surface. + +### Phase 7 — Cleanup, docs, deletion ledger +- Update `inference/README.md`, `gitbooks/developing/architecture/agent-harness.md`, `docs/inference-provider-catalog.md`; add deletions to `docs/tinyagents-full-migration-plan/99-deletion-ledger.md`; refresh `docs/tinyagents-drift-ledger.md`. +- Remove dead re-exports from `inference/mod.rs`; keep temporary `pub use` shims only where a follow-up PR is already open. +- Sweep for stale doc-comments naming `Provider`/`CompatibleProvider`. + +--- + +## 5. Risks & gotchas + +- **Wire-format regressions are silent until a provider hiccups.** The compat client encodes years of quirk handling (SSE edge cases, malformed tool-arg fragments, providers that omit usage). Mitigation: Phase 0 golden dumps + per-provider live smoke tests gated on env keys, run against the real BYOK matrix before each deletion. +- **Cost accounting**: the event bridge + `record_unobserved_turn_usage` fallback were hard-won ($0-turn bug). Any `Usage` shape change must keep cached-token and USD flow intact end-to-end (dashboard + footer). +- **Streaming UI parity**: tool-start events (G2) and post-hoc reasoning still ride the out-of-band forwarder; deleting it before the crate gap closes breaks the tool timeline. +- **Two Cargo worlds**: root and `app/src-tauri` pin tinyagents independently — every crate bump lands in both lockfiles plus the submodule ref, same commit. +- **Vendored-crate discipline**: the path patch means local vendor edits silently take effect; CI and other clones need the submodule at the matching ref. Never merge openhuman code that requires unpublished crate API. +- **Sentry noise contract**: `ops.rs` deliberately demotes provider/user-config failures to `warn!`. The new host classifier over `TinyAgentsError` must preserve `expected_error_kind` behavior or Sentry floods. +- **Test serialization**: everything runs under `inference_test_guard()` (process-global mutex over the runtime singleton + config); new tests must too. +- **Open regressions in the same area** (#4451–#4469, esp. #4460 streamed calls losing thread_id task-locals, #4463 repeat guards): coordinate so this migration doesn't re-break or mask those fixes; thread-context task-locals move in Phase 5 — verify #4460's fix survives the re-home. +- **`/v1` server endpoint** reuses provider types for its request/response DTOs (`http/types.rs`) — it must keep its external wire shape while internals switch to crate types. + +## 6. Explicit non-goals + +- No change to the `inference.*` RPC surface, provider-string grammar, or the Settings > AI preset catalog UX. +- No migration of `local/` process management, `voice/` STT/TTS engines, `openai_oauth/` flows, or `device`/`presets`/`paths` — they are not LLM-driven in the crate's sense. +- No sub-agent execution changes (P5 of the harness migration already declined crate `SubAgentTool`). diff --git a/src/openhuman/accessibility/automate.rs b/src/openhuman/accessibility/automate.rs index ce20f2b3c..374a994b5 100644 --- a/src/openhuman/accessibility/automate.rs +++ b/src/openhuman/accessibility/automate.rs @@ -663,14 +663,19 @@ impl AutomateBackend for RealBackend { // `automation` provider knob is a follow-up (see plan §5); routing through // `summarization` keeps M1 free of Config-schema churn while still keeping // the chat model out of the loop. - let (provider, model) = crate::openhuman::inference::provider::create_chat_provider( + use tinyagents::harness::message::Message; + use tinyagents::harness::model::{ChatModel, ModelRequest}; + let model = crate::openhuman::inference::provider::create_chat_model( "summarization", &self.config, + 0.0, ) .map_err(|e| format!("fast-model provider unavailable: {e}"))?; - provider - .chat_with_system(Some(system), user, &model, 0.0) + let request = ModelRequest::new(vec![Message::system(system), Message::user(user)]); + model + .invoke(&(), request) .await + .map(|response| response.text()) .map_err(|e| format!("fast-model call failed: {e}")) } diff --git a/src/openhuman/agent/harness/session/builder/factory.rs b/src/openhuman/agent/harness/session/builder/factory.rs index 056104dc4..ad2c944fe 100644 --- a/src/openhuman/agent/harness/session/builder/factory.rs +++ b/src/openhuman/agent/harness/session/builder/factory.rs @@ -469,23 +469,15 @@ impl Agent { // `chat_provider` selection. Subagents still set their own role // through `ModelSpec::Hint(...)` in the subagent runner. let provider_role = provider_role_for(agent_id, config.default_model.as_deref()); - let (raw_provider, mut model_name): (Box, String) = + // Retry/backoff is now owned by the crate `RetryPolicy` at the harness + // model call (issue #4249, Phase 3a) — see `tinyagents::run_policy_for`. + // The turn path therefore no longer wraps the resolved provider in + // `ReliableProvider`; wrapping it here plus the crate retry would + // double-retry every transient error. Cross-route fallback is likewise + // the crate registry `FallbackPolicy`, so `config.reliability.*` no longer + // layers on the turn path (it still governs the non-seam provider paths). + let (provider, mut model_name): (Box, String) = crate::openhuman::inference::provider::create_chat_provider(provider_role, config)?; - // Re-layer the ReliableProvider retry/backoff + model-fallback wrapper on - // top of the factory's resolved backend (issue #4249, 1c). The migration to - // `create_chat_provider` dropped this; restore it so rate-limit/5xx retries - // and the user's `model_fallbacks` apply to the main chat turn exactly as - // the legacy `create_intelligent_routing_provider` path did. Capability - // probes (`supports_native_tools` / `supports_vision`) forward to the inner - // backend, so downstream dispatcher/vision selection is unchanged. - let provider: Box = Box::new( - crate::openhuman::inference::provider::reliable::ReliableProvider::new( - vec![(provider_role.to_string(), raw_provider)], - config.reliability.provider_retries, - config.reliability.provider_backoff_ms, - ) - .with_model_fallbacks(config.reliability.model_fallbacks.clone()), - ); log::info!( "[session-builder] agent_id={} provider_role={} resolved_model={} supports_native_tools={}", agent_id, @@ -697,18 +689,27 @@ impl Agent { // For cloud reflection, wrap the provider in an Arc. // For local, no provider needed. let reflection_provider: Option< - Arc, + Arc>, > = if config.learning.reflection_source == crate::openhuman::config::ReflectionSource::Cloud { - Some(Arc::from(provider::create_routed_provider( + // Reflection always calls with the `hint:reasoning` route + + // 0.3 temperature (formerly `simple_chat(prompt, + // "hint:reasoning", 0.3)`), so bake both into the wrapped + // model. The routed provider still resolves the hint per call. + let routed = provider::create_routed_provider( config.inference_url.as_deref(), config.api_url.as_deref(), config.api_key.as_deref(), &config.reliability, &config.model_routes, &model_name, - )?)) + )?; + Some(provider::chat_model_from_provider( + routed, + "hint:reasoning".to_string(), + 0.3, + )) } else { None }; diff --git a/src/openhuman/agent_meetings/summary.rs b/src/openhuman/agent_meetings/summary.rs index 3cf6fc10e..45dfbf862 100644 --- a/src/openhuman/agent_meetings/summary.rs +++ b/src/openhuman/agent_meetings/summary.rs @@ -15,7 +15,9 @@ use serde::Deserialize; use crate::core::event_bus::BackendMeetTurn; use crate::openhuman::config::{AutoSummarizePolicy, Config}; -use crate::openhuman::inference::provider::create_chat_provider; +use crate::openhuman::inference::provider::create_chat_model; +use tinyagents::harness::message::Message; +use tinyagents::harness::model::{ChatModel, ModelRequest}; use super::types::{ActionItem, ActionItemKind, MeetingSummary}; @@ -142,13 +144,20 @@ pub async fn generate_meeting_summary( let config = Config::load_or_init() .await .map_err(|e| format!("config load failed: {e}"))?; - let (provider, model) = create_chat_provider(SUMMARIZATION_ROLE, &config) + let model = create_chat_model(SUMMARIZATION_ROLE, &config, 0.3) .map_err(|e| format!("summarisation provider init failed: {e}"))?; - let reply = provider - .chat_with_system(Some(SYSTEM_PROMPT), &transcript, &model, 0.3) + // One user turn under a fixed system prompt — the crate model interface's + // native shape. Equivalent to the former `chat_with_system` one-shot call. + let request = ModelRequest::new(vec![ + Message::system(SYSTEM_PROMPT), + Message::user(transcript), + ]); + let reply = model + .invoke(&(), request) .await - .map_err(|e| format!("summarisation LLM call failed: {e}"))?; + .map_err(|e| format!("summarisation LLM call failed: {e}"))? + .text(); let raw = parse_summary_json(&reply) .ok_or_else(|| "model reply did not contain parseable summary JSON".to_string())?; diff --git a/src/openhuman/inference/provider/factory.rs b/src/openhuman/inference/provider/factory.rs index c0e7a1f8c..abf7f99fe 100644 --- a/src/openhuman/inference/provider/factory.rs +++ b/src/openhuman/inference/provider/factory.rs @@ -41,6 +41,8 @@ use crate::openhuman::inference::provider::openai_codex::{ use crate::openhuman::inference::provider::openhuman_backend::OpenHumanBackendProvider; use crate::openhuman::inference::provider::traits::Provider; use crate::openhuman::inference::provider::ProviderRuntimeOptions; +use std::sync::Arc; +use tinyagents::harness::model::ChatModel; /// Sentinel meaning "use the OpenHuman backend session JWT". pub const PROVIDER_OPENHUMAN: &str = "openhuman"; @@ -891,6 +893,75 @@ pub fn create_chat_provider_from_string( ) } +/// Build an `Arc` for the given workload role. +/// +/// Phase 1 of the tinyagents inference migration (#4249): the crate +/// [`ChatModel`] is the model interface the harness and one-shot inference +/// callers target. Today this wraps the existing openhuman [`Provider`] stack +/// via [`ProviderModel`](crate::openhuman::tinyagents::model) — a **zero +/// behaviour change** shim — so callers can move off `Box` +/// incrementally while the provider stack is dismantled underneath. `temperature` +/// is pinned onto the returned model because the crate model interface bakes +/// sampling into the model rather than the per-call request. +pub fn create_chat_model( + role: &str, + config: &Config, + temperature: f64, +) -> anyhow::Result>> { + Ok(create_chat_model_with_model_id(role, config, temperature)?.0) +} + +/// Like [`create_chat_model`], but also returns the resolved model id. +/// +/// One-shot callers that persist or log the concrete model (e.g. the memory +/// summarise audit) need the id the role resolved to; the plain +/// [`create_chat_model`] drops it. +pub fn create_chat_model_with_model_id( + role: &str, + config: &Config, + temperature: f64, +) -> anyhow::Result<(Arc>, String)> { + let (provider, model) = create_chat_provider(role, config)?; + let chat = chat_model_from_provider(provider, model.clone(), temperature); + Ok((chat, model)) +} + +/// Build an `Arc` from an explicit provider string and config. +/// +/// The [`ChatModel`] counterpart of [`create_chat_provider_from_string`]; see +/// [`create_chat_model`] for the migration rationale. +pub fn create_chat_model_from_string( + role: &str, + provider: &str, + config: &Config, + temperature: f64, +) -> anyhow::Result>> { + let (provider, model) = create_chat_provider_from_string(role, provider, config)?; + Ok(chat_model_from_provider(provider, model, temperature)) +} + +/// Wrap an owned [`Provider`] as an `Arc` pinned to +/// `model`/`temperature`. +/// +/// The single seam where a boxed provider becomes the crate model interface. +/// As consumers migrate off `Box` this stays the conversion point, +/// shrinking toward the Phase 1 exit criterion (`ProviderModel` constructed in +/// exactly one place) and, ultimately, the `Provider` trait's deletion in +/// Phase 4. Exposed `pub(crate)` so a caller that must build a specific provider +/// itself (e.g. the LinkedIn enrichment path, which deliberately forces the +/// managed backend) can still hand back a `ChatModel` without naming the trait. +pub(crate) fn chat_model_from_provider( + provider: Box, + model: String, + temperature: f64, +) -> Arc> { + crate::openhuman::tinyagents::model::provider_chat_model( + Arc::from(provider), + model, + temperature, + ) +} + /// Build a local-runtime provider without applying the custom-provider session gate. /// /// Used by setup/probe flows that need to validate an endpoint before the diff --git a/src/openhuman/inference/provider/factory_tests.rs b/src/openhuman/inference/provider/factory_tests.rs index 678af277a..e388e430e 100644 --- a/src/openhuman/inference/provider/factory_tests.rs +++ b/src/openhuman/inference/provider/factory_tests.rs @@ -2710,3 +2710,44 @@ fn enforce_local_only_inference_errors_on_external_when_local_only() { crate::openhuman::security::live_policy::reload_privacy(PrivacyMode::Standard) .expect("policy installed"); } + +// ── Phase 1 (#4249): `create_chat_model` seam ────────────────────────────── +// The crate `ChatModel` factory wraps the resolved `Provider` via +// `ProviderModel` with zero behaviour change; a one-shot `invoke` must +// round-trip through the underlying provider. +#[tokio::test] +async fn create_chat_model_wraps_provider_and_round_trips() { + use crate::openhuman::inference::provider::traits::Provider; + use async_trait::async_trait; + use std::sync::Arc; + use tinyagents::harness::message::Message; + use tinyagents::harness::model::{ChatModel, ModelRequest}; + + let _guard = crate::openhuman::inference::inference_test_guard(); + + struct CannedProvider; + #[async_trait] + impl Provider for CannedProvider { + async fn chat_with_system( + &self, + _system_prompt: Option<&str>, + message: &str, + _model: &str, + _temperature: f64, + ) -> anyhow::Result { + Ok(format!("echo: {message}")) + } + } + + // The factory consults this override under cfg(test), so `create_chat_model` + // resolves to the mock without needing configured cloud providers. + let _override = test_provider_override::install(Arc::new(CannedProvider)); + let config = Config::default(); + + let model = create_chat_model("chat", &config, 0.3).expect("create_chat_model must build"); + let response = model + .invoke(&(), ModelRequest::new(vec![Message::user("hi there")])) + .await + .expect("invoke must succeed"); + assert_eq!(response.text(), "echo: hi there"); +} diff --git a/src/openhuman/inference/provider/mod.rs b/src/openhuman/inference/provider/mod.rs index 5db66f733..faa3e0f3a 100644 --- a/src/openhuman/inference/provider/mod.rs +++ b/src/openhuman/inference/provider/mod.rs @@ -45,7 +45,9 @@ pub use error_code::{ is_backend_malformed_bad_request, is_managed_backend_envelope, managed_error_skips_sentry, BackendErrorCode, }; +pub(crate) use factory::chat_model_from_provider; pub use factory::{ + create_chat_model, create_chat_model_from_string, create_chat_model_with_model_id, create_chat_provider, provider_for_role, role_for_model_tier, BYOK_INCOMPLETE_SENTINEL, }; pub use ops::*; diff --git a/src/openhuman/learning/linkedin_enrichment.rs b/src/openhuman/learning/linkedin_enrichment.rs index 5cee7ffa0..933ddbc27 100644 --- a/src/openhuman/learning/linkedin_enrichment.rs +++ b/src/openhuman/learning/linkedin_enrichment.rs @@ -334,6 +334,15 @@ pub async fn summarise_profile_with_llm(config: &Config, raw_md: &str) -> anyhow config.api_key.as_deref(), &options, )?; + // This path deliberately forces the managed backend (not role routing), so + // wrap the built provider directly rather than resolving via + // `create_chat_model`. Behaviour-preserving; only the call surface moves to + // the crate `ChatModel`. + let model_chat = crate::openhuman::inference::provider::chat_model_from_provider( + provider, + "summarization-v1".to_string(), + 0.3, + ); let system = "\ You are a profile analyst. You will receive a user's LinkedIn profile in Markdown format. \ @@ -359,9 +368,18 @@ Rules:\n\ "[linkedin_enrichment] sending profile to LLM for summarisation" ); - let summary = provider - .chat_with_system(Some(system), raw_md, model, 0.3) - .await?; + use tinyagents::harness::message::Message; + use tinyagents::harness::model::{ChatModel, ModelRequest}; + let summary = model_chat + .invoke( + &(), + ModelRequest::new(vec![ + Message::system(system), + Message::user(raw_md.to_string()), + ]), + ) + .await? + .text(); tracing::debug!( output_len = summary.len(), diff --git a/src/openhuman/learning/reflection.rs b/src/openhuman/learning/reflection.rs index 4ef32d538..a46a0e240 100644 --- a/src/openhuman/learning/reflection.rs +++ b/src/openhuman/learning/reflection.rs @@ -44,17 +44,38 @@ pub struct ReflectionHook { config: LearningConfig, full_config: Arc, memory: Arc, - provider: Option>, + provider: Option>>, /// Per-session reflection counts for throttling. Key is session_id (or "__global__"). session_counts: Mutex>, } +/// Run one cloud reflection call through the crate model interface. +/// +/// The reflection model is built with the `hint:reasoning` route + 0.3 +/// temperature already baked in (see the session builder), so this is a single +/// user turn with no per-request overrides — the [`ChatModel`] equivalent of the +/// former `simple_chat(prompt, "hint:reasoning", 0.3)`. +async fn invoke_cloud_reflection( + provider: &Arc>, + prompt: &str, +) -> anyhow::Result { + use tinyagents::harness::message::Message; + use tinyagents::harness::model::{ChatModel, ModelRequest}; + Ok(provider + .invoke( + &(), + ModelRequest::new(vec![Message::user(prompt.to_string())]), + ) + .await? + .text()) +} + impl ReflectionHook { pub fn new( config: LearningConfig, full_config: Arc, memory: Arc, - provider: Option>, + provider: Option>>, ) -> Self { Self { config, @@ -171,8 +192,7 @@ impl ReflectionHook { "[learning::reflection] local_ai.usage.learning_reflection not enabled — \ falling back to cloud provider" ); - return provider - .simple_chat(prompt, "hint:reasoning", 0.3) + return invoke_cloud_reflection(provider, prompt) .await .map_err(|e| anyhow::anyhow!("cloud reflection fallback failed: {e}")); } @@ -201,7 +221,7 @@ impl ReflectionHook { let provider = self.provider.as_ref().ok_or_else(|| { anyhow::anyhow!("no cloud provider configured for reflection") })?; - provider.simple_chat(prompt, "hint:reasoning", 0.3).await + invoke_cloud_reflection(provider, prompt).await } } } diff --git a/src/openhuman/learning/reflection_tests.rs b/src/openhuman/learning/reflection_tests.rs index b86c22208..6bb296139 100644 --- a/src/openhuman/learning/reflection_tests.rs +++ b/src/openhuman/learning/reflection_tests.rs @@ -378,11 +378,18 @@ async fn on_turn_complete_dedupes_reflections_across_heuristic_and_llm_paths() { let memory_impl = Arc::new(MockMemory::default()); let memory: Arc = memory_impl.clone(); + // Wrap the stub provider as a `ChatModel` (the field's type after the Phase 1 + // migration); the hint/temperature baked here are inert for the stub. + let stub_model = crate::openhuman::inference::provider::chat_model_from_provider( + Box::new(StubProvider), + "hint:reasoning".to_string(), + 0.3, + ); let hook = ReflectionHook::new( reflection_config(), Arc::new(Config::default()), memory, - Some(Arc::new(StubProvider)), + Some(stub_model), ); let turn = TurnContext { @@ -567,11 +574,16 @@ async fn on_turn_complete_emits_style_candidates_from_llm_preferences() { } let memory: Arc = Arc::new(MockMemory::default()); + let stub_model = crate::openhuman::inference::provider::chat_model_from_provider( + Box::new(StubPrefProvider), + "hint:reasoning".to_string(), + 0.3, + ); let hook = ReflectionHook::new( reflection_config(), Arc::new(Config::default()), memory, - Some(Arc::new(StubPrefProvider)), + Some(stub_model), ); let before = candidate::global().len(); diff --git a/src/openhuman/memory/chat.rs b/src/openhuman/memory/chat.rs index 834a0424e..1beea7f4d 100644 --- a/src/openhuman/memory/chat.rs +++ b/src/openhuman/memory/chat.rs @@ -13,8 +13,10 @@ use async_trait::async_trait; use crate::openhuman::config::{Config, DEFAULT_CLOUD_LLM_MODEL}; use crate::openhuman::inference::provider::{ - create_chat_provider, provider_for_role, ChatMessage, ChatRequest, Provider, UsageInfo, + create_chat_model_with_model_id, provider_for_role, UsageInfo, }; +use tinyagents::harness::message::Message; +use tinyagents::harness::model::{ChatModel, ModelRequest}; /// One pair of prompt messages handed to the memory LLM backend. #[derive(Debug, Clone)] @@ -62,17 +64,17 @@ pub trait ChatProvider: Send + Sync { } struct InferenceChatProvider { - inner: Box, - model: String, + inner: Arc>, + model_id: String, display: String, } impl InferenceChatProvider { - fn new(inner: Box, model: String) -> Self { - let display = format!("inference:{model}"); + fn new(inner: Arc>, model_id: String) -> Self { + let display = format!("inference:{model_id}"); Self { inner, - model, + model_id, display, } } @@ -93,41 +95,42 @@ impl InferenceChatProvider { "[memory::chat] provider={} kind={} model={} sys_chars={} user_chars={}", self.display, prompt.kind, - self.model, + self.model_id, prompt.system.len(), prompt.user.len() ); - let messages = vec![ - ChatMessage::system(prompt.system.clone()), - ChatMessage::user(prompt.user.clone()), - ]; + // One system + one user turn — the crate model interface's native shape. + // Temperature and the output cap ride the request (the shared model is + // reused across memory prompts of differing temperature/budget), and the + // adapter honors both per-request values. + let mut request = ModelRequest::new(vec![ + Message::system(prompt.system.clone()), + Message::user(prompt.user.clone()), + ]) + .with_temperature(prompt.temperature); + if let Some(cap) = prompt.max_tokens { + request = request.with_max_tokens(cap); + } - let request = ChatRequest { - messages: &messages, - tools: None, - stream: None, - max_tokens: prompt.max_tokens, - }; - - let response = self - .inner - .chat(request, &self.model, prompt.temperature) - .await?; + let response = self.inner.invoke(&(), request).await?; // Fail fast on a missing body rather than masking it as an empty // string: an empty summary would still be ingested (and, post-#3110, // counted against the run's real charge) as if it were valid output. // The caller's fallback path (`fallback_summary`) is the correct // recovery for a silent provider, and it only runs on `Err`. - let Some(text) = response.text else { + let text = response.text(); + if text.is_empty() { anyhow::bail!( "inference provider '{}' returned no text for {} summarise request", self.display, prompt.kind ); - }; - let usage = response.usage; + } + // Recover the full host usage (real token counts + backend-charged USD + + // context window) the adapter round-tripped through the response (G1). + let usage = crate::openhuman::tinyagents::model::usage_info_from_response(&response); log::debug!( "[memory::chat] provider={} kind={} response_chars={} usage_present={} input_tokens={} output_tokens={} charged_usd={}", @@ -187,17 +190,21 @@ pub fn build_chat_runtime(config: &Config) -> Result<(Arc, Str // per-caller `default_model` pre-routing is needed here. BYOK/local routes // carry their own model in the provider string. let resolved_provider = provider_for_role("summarization", config); - let (provider, model) = create_chat_provider("summarization", config)?; + // Temperature is applied per-prompt via `ModelRequest::with_temperature` + // (each memory `ChatPrompt` carries its own), so the construction temperature + // is just a default the per-call value overrides. + let (model, model_id) = + create_chat_model_with_model_id("summarization", config, config.default_temperature)?; log::debug!( "[memory::chat] built provider route={} model={}", resolved_provider, - model + model_id ); Ok(( - Arc::new(InferenceChatProvider::new(provider, model.clone())), - model, + Arc::new(InferenceChatProvider::new(model, model_id.clone())), + model_id, )) } diff --git a/src/openhuman/memory_tree/tree_runtime/engine.rs b/src/openhuman/memory_tree/tree_runtime/engine.rs index 59f79264f..85681dcb0 100644 --- a/src/openhuman/memory_tree/tree_runtime/engine.rs +++ b/src/openhuman/memory_tree/tree_runtime/engine.rs @@ -7,12 +7,14 @@ use std::collections::BTreeMap; use crate::core::event_bus::{publish_global, DomainEvent}; use crate::openhuman::config::Config; -use crate::openhuman::inference::provider::traits::Provider; use crate::openhuman::memory_tree::tree_runtime::store; use crate::openhuman::memory_tree::tree_runtime::types::{ derive_node_ids, derive_parent_id, estimate_tokens, level_from_node_id, NodeLevel, TreeNode, TreeStatus, }; +use std::sync::Arc; +use tinyagents::harness::message::Message; +use tinyagents::harness::model::{ChatModel, ModelRequest}; const SUMMARIZATION_TEMP: f64 = 0.3; @@ -32,8 +34,7 @@ const MAX_SUMMARY_CHARS: usize = 20_000 * 4; /// Returns the last hour leaf node created, or `None` if the buffer was empty. pub async fn run_summarization( config: &Config, - provider: &dyn Provider, - model: &str, + provider: &dyn ChatModel<()>, namespace: &str, _ts: DateTime, ) -> Result> { @@ -86,7 +87,6 @@ pub async fn run_summarization( NodeLevel::Hour.max_tokens(), "hour", hour_id, - model, ) .await .context("summarize hour leaf")?; @@ -147,7 +147,7 @@ pub async fn run_summarization( ] { for (node_id, node_level) in &all_propagation_ids { if *node_level == level && seen.insert(node_id.clone()) { - match propagate_node(config, provider, namespace, node_id, level, model).await { + match propagate_node(config, provider, namespace, node_id, level).await { Ok(()) => propagated += 1, Err(e) => { log::warn!( @@ -193,8 +193,7 @@ pub async fn run_summarization( /// Preserves buffered data that hasn't been summarized yet. pub async fn rebuild_tree( config: &Config, - provider: &dyn Provider, - model: &str, + provider: &dyn ChatModel<()>, namespace: &str, ) -> Result { tracing::debug!("[tree_summarizer] rebuilding tree for namespace '{namespace}'"); @@ -278,9 +277,7 @@ pub async fn rebuild_tree( let mut propagate = |id: &str, level: NodeLevel| { let node_id = id.to_string(); async move { - if let Err(e) = - propagate_node(config, provider, namespace, &node_id, level, model).await - { + if let Err(e) = propagate_node(config, provider, namespace, &node_id, level).await { log::warn!( "[tree_summarizer] rebuild propagate failed (continuing) namespace='{namespace}' \ node={node_id} level={}: {e:#}", @@ -338,11 +335,10 @@ pub async fn rebuild_tree( /// Re-summarize a single non-leaf node from its children. pub(crate) async fn propagate_node( config: &Config, - provider: &dyn Provider, + provider: &dyn ChatModel<()>, namespace: &str, node_id: &str, level: NodeLevel, - model: &str, ) -> Result<()> { let children = store::read_children(config, namespace, node_id)?; if children.is_empty() { @@ -380,15 +376,7 @@ pub(crate) async fn propagate_node( combined_tokens, max_tokens ); - summarize_to_limit( - provider, - &combined, - max_tokens, - level.as_str(), - node_id, - model, - ) - .await? + summarize_to_limit(provider, &combined, max_tokens, level.as_str(), node_id).await? }; let now = Utc::now(); @@ -429,12 +417,11 @@ pub(crate) async fn propagate_node( /// Summarize text to fit within a token limit using the LLM provider. /// Enforces a hard character limit on the response to prevent runaway output. async fn summarize_to_limit( - provider: &dyn Provider, + provider: &dyn ChatModel<()>, content: &str, max_tokens: u32, level_name: &str, node_id: &str, - model: &str, ) -> Result { let max_chars = (max_tokens as usize) * 4; let system_prompt = format!( @@ -449,11 +436,19 @@ async fn summarize_to_limit( ); let response = provider - .chat_with_system(Some(&system_prompt), content, model, SUMMARIZATION_TEMP) + .invoke( + &(), + ModelRequest::new(vec![ + Message::system(system_prompt), + Message::user(content.to_string()), + ]) + .with_temperature(SUMMARIZATION_TEMP), + ) .await .with_context(|| { format!("LLM summarization failed for node {node_id} (level={level_name})") - })?; + })? + .text(); // Enforce hard character limit on LLM response (use the stricter of the two limits) let char_limit = max_chars.min(MAX_SUMMARY_CHARS); @@ -587,7 +582,7 @@ fn collect_hour_leaves_recursive( /// /// This should be called once at application startup. The task runs /// indefinitely, sleeping until the next hour boundary. -pub async fn run_hourly_loop(config: Config, provider: Box, model: String) { +pub async fn run_hourly_loop(config: Config, provider: Arc>) { tracing::debug!("[tree_summarizer] hourly loop started"); loop { @@ -615,7 +610,7 @@ pub async fn run_hourly_loop(config: Config, provider: Box, model: let ts = Utc::now(); let namespaces = discover_active_namespaces(&config); for ns in &namespaces { - match run_summarization(&config, provider.as_ref(), &model, ns, ts).await { + match run_summarization(&config, provider.as_ref(), ns, ts).await { Ok(Some(node)) => { tracing::debug!( "[tree_summarizer] hourly job completed for '{}': node {} ({} tokens)", diff --git a/src/openhuman/memory_tree/tree_runtime/engine_tests.rs b/src/openhuman/memory_tree/tree_runtime/engine_tests.rs index d42614349..866eb41f6 100644 --- a/src/openhuman/memory_tree/tree_runtime/engine_tests.rs +++ b/src/openhuman/memory_tree/tree_runtime/engine_tests.rs @@ -69,6 +69,18 @@ impl Provider for FailAtLevelProvider { } } +/// Wrap a stub [`Provider`] as the crate `ChatModel` the engine now consumes. +/// The baked model/temperature are inert for the stubs (they ignore both). +fn as_model( + provider: impl Provider + 'static, +) -> std::sync::Arc> { + crate::openhuman::inference::provider::chat_model_from_provider( + Box::new(provider), + "test-model".to_string(), + 0.3, + ) +} + // ── group_by_hour ──────────────────────────────────────────────────────── #[test] @@ -191,16 +203,14 @@ async fn propagate_node_with_no_children_is_noop() { let tmp = TempDir::new().unwrap(); let cfg = test_config(&tmp); std::fs::create_dir_all(&cfg.workspace_dir).unwrap(); - let provider = StubProvider::with_reply("should not be called"); - let model = "stub-model"; + let provider = as_model(StubProvider::with_reply("should not be called")); // No children exist for "2024/03/15" — propagate must succeed silently. let result = propagate_node( &cfg, - &provider, + provider.as_ref(), "test-ns", "2024/03/15", NodeLevel::Day, - model, ) .await; assert!(result.is_ok(), "empty children must not error: {result:?}"); @@ -247,8 +257,8 @@ async fn propagate_node_day_from_hour_children_fits_budget() { .unwrap(); // StubProvider reply should NOT be used when content fits budget. - let provider = StubProvider::with_reply("SHOULD_NOT_APPEAR"); - propagate_node(&cfg, &provider, ns, "2024/03/15", NodeLevel::Day, "m") + let provider = as_model(StubProvider::with_reply("SHOULD_NOT_APPEAR")); + propagate_node(&cfg, provider.as_ref(), ns, "2024/03/15", NodeLevel::Day) .await .unwrap(); @@ -289,8 +299,8 @@ async fn propagate_node_month_from_day_children() { store::write_node(&cfg, &make_day("2024/03/14", "Day 14 recap.")).unwrap(); store::write_node(&cfg, &make_day("2024/03/15", "Day 15 recap.")).unwrap(); - let provider = StubProvider::with_reply("Month summary from LLM."); - propagate_node(&cfg, &provider, ns, "2024/03", NodeLevel::Month, "m") + let provider = as_model(StubProvider::with_reply("Month summary from LLM.")); + propagate_node(&cfg, provider.as_ref(), ns, "2024/03", NodeLevel::Month) .await .unwrap(); @@ -346,8 +356,8 @@ async fn propagate_node_preserves_created_at_on_update() { }; store::write_node(&cfg, &child).unwrap(); - let provider = StubProvider::with_reply("updated summary"); - propagate_node(&cfg, &provider, ns, "2024/03/15", NodeLevel::Day, "m") + let provider = as_model(StubProvider::with_reply("updated summary")); + propagate_node(&cfg, provider.as_ref(), ns, "2024/03/15", NodeLevel::Day) .await .unwrap(); @@ -369,9 +379,9 @@ async fn run_summarization_empty_buffer_returns_none() { let tmp = TempDir::new().unwrap(); let cfg = test_config(&tmp); std::fs::create_dir_all(&cfg.workspace_dir).unwrap(); - let provider = StubProvider::with_reply("should not be called"); + let provider = as_model(StubProvider::with_reply("should not be called")); let ts = Utc::now(); - let result = run_summarization(&cfg, &provider, "test-model", "test-ns", ts) + let result = run_summarization(&cfg, provider.as_ref(), "test-ns", ts) .await .unwrap(); assert!(result.is_none(), "empty buffer must return None"); @@ -389,8 +399,8 @@ async fn run_summarization_drains_buffer_and_writes_hour_node() { store::buffer_write(&cfg, ns, "entry one", &ts, None).unwrap(); store::buffer_write(&cfg, ns, "entry two", &ts, None).unwrap(); - let provider = StubProvider::with_reply("hour leaf summary from LLM"); - let last_node = run_summarization(&cfg, &provider, "test-model", ns, ts) + let provider = as_model(StubProvider::with_reply("hour leaf summary from LLM")); + let last_node = run_summarization(&cfg, provider.as_ref(), ns, ts) .await .unwrap(); @@ -425,8 +435,8 @@ async fn run_summarization_builds_ancestor_chain() { store::buffer_write(&cfg, ns, "test content", &ts, None).unwrap(); - let provider = StubProvider::with_reply("summary text"); - run_summarization(&cfg, &provider, "test-model", ns, ts) + let provider = as_model(StubProvider::with_reply("summary text")); + run_summarization(&cfg, provider.as_ref(), ns, ts) .await .unwrap(); @@ -463,8 +473,8 @@ async fn run_summarization_multi_hour_groups_produce_multiple_hour_leaves() { store::buffer_write(&cfg, ns, "morning entry", &ts_h08, None).unwrap(); store::buffer_write(&cfg, ns, "afternoon entry", &ts_h14, None).unwrap(); - let provider = StubProvider::with_reply("grouped summary"); - run_summarization(&cfg, &provider, "test-model", ns, ts_h14) + let provider = as_model(StubProvider::with_reply("grouped summary")); + run_summarization(&cfg, provider.as_ref(), ns, ts_h14) .await .unwrap(); @@ -537,11 +547,9 @@ async fn rebuild_tree_restores_buffer_and_rewrites_ancestors() { // Preserve unsummarized buffer content across rebuild. store::buffer_write(&cfg, ns, "pending buffer item", &ts, None).unwrap(); - let provider = StubProvider::with_reply("rebuilt summary"); + let provider = as_model(StubProvider::with_reply("rebuilt summary")); - let status = rebuild_tree(&cfg, &provider, "test-model", ns) - .await - .unwrap(); + let status = rebuild_tree(&cfg, provider.as_ref(), ns).await.unwrap(); assert!(status.total_nodes >= 5, "expected leaf + ancestor chain"); let restored_buffer = store::buffer_read(&cfg, ns).unwrap(); @@ -594,13 +602,13 @@ async fn rebuild_tree_partial_success_when_one_level_fails() { store::write_node(&cfg, &make_hour("2024/03/15/10")).unwrap(); store::write_node(&cfg, &make_hour("2024/03/15/11")).unwrap(); - let provider = FailAtLevelProvider { + let provider = as_model(FailAtLevelProvider { fail_level: "day", reply: "ok summary".to_string(), - }; + }); // Must NOT error despite the day-level summarization failing. - let status = rebuild_tree(&cfg, &provider, "test-model", ns) + let status = rebuild_tree(&cfg, provider.as_ref(), ns) .await .expect("partial failure must not abort the rebuild"); @@ -623,8 +631,8 @@ async fn rebuild_tree_on_empty_namespace_is_noop() { let cfg = test_config(&tmp); std::fs::create_dir_all(&cfg.workspace_dir).unwrap(); - let provider = StubProvider::with_reply("unused"); - let status = rebuild_tree(&cfg, &provider, "test-model", "empty-rebuild") + let provider = as_model(StubProvider::with_reply("unused")); + let status = rebuild_tree(&cfg, provider.as_ref(), "empty-rebuild") .await .unwrap(); assert_eq!(status.total_nodes, 0); @@ -633,8 +641,8 @@ async fn rebuild_tree_on_empty_namespace_is_noop() { #[tokio::test] async fn summarize_to_limit_truncates_overlong_provider_output() { - let provider = StubProvider::with_reply("x".repeat(MAX_SUMMARY_CHARS + 50)); - let summary = summarize_to_limit(&provider, "short input", 10, "day", "2024/03/15", "m") + let provider = as_model(StubProvider::with_reply("x".repeat(MAX_SUMMARY_CHARS + 50))); + let summary = summarize_to_limit(provider.as_ref(), "short input", 10, "day", "2024/03/15") .await .unwrap(); diff --git a/src/openhuman/memory_tree/tree_runtime/ops.rs b/src/openhuman/memory_tree/tree_runtime/ops.rs index 82d3adf6a..2192bea37 100644 --- a/src/openhuman/memory_tree/tree_runtime/ops.rs +++ b/src/openhuman/memory_tree/tree_runtime/ops.rs @@ -44,10 +44,10 @@ pub async fn tree_summarizer_run( ) -> Result, String> { store::validate_namespace(namespace)?; - let (provider, model) = create_provider(config)?; + let (provider, _model) = create_provider(config)?; let ts = Utc::now(); - match engine::run_summarization(config, provider.as_ref(), &model, namespace.trim(), ts).await { + match engine::run_summarization(config, provider.as_ref(), namespace.trim(), ts).await { Ok(Some(node)) => Ok(RpcOutcome::single_log( serde_json::to_value(&node).map_err(|e| e.to_string())?, format!( @@ -126,9 +126,9 @@ pub async fn tree_summarizer_rebuild( ) -> Result, String> { store::validate_namespace(namespace)?; - let (provider, model) = create_provider(config)?; + let (provider, _model) = create_provider(config)?; - let status = engine::rebuild_tree(config, provider.as_ref(), &model, namespace.trim()) + let status = engine::rebuild_tree(config, provider.as_ref(), namespace.trim()) .await .map_err(|e| format!("rebuild failed: {e:#}"))?; @@ -166,15 +166,24 @@ fn create_provider( config: &Config, ) -> Result< ( - Box, + std::sync::Arc>, String, ), String, > { + // The summarizer applies its own temperature per request + // (`SUMMARIZATION_TEMP` in `engine`), so the construction temperature here is + // just a default the per-call value overrides. if config.local_ai.runtime_enabled { // Local path: Ollama + the user's local chat model. let provider = create_local_ai_provider(config)?; - return Ok((provider, config.local_ai.chat_model_id.clone())); + let model = config.local_ai.chat_model_id.clone(); + let chat = crate::openhuman::inference::provider::chat_model_from_provider( + provider, + model.clone(), + config.default_temperature, + ); + return Ok((chat, model)); } if !config.memory_tree.cloud_summarization_opt_in { @@ -185,8 +194,12 @@ fn create_provider( // Cloud path — user has explicitly opted in. Build the configured // provider for the summarization role (`memory_provider` hint). - crate::openhuman::inference::provider::factory::create_chat_provider("summarization", config) - .map_err(|e| format!("tree summarizer: failed to build cloud provider: {e:#}")) + crate::openhuman::inference::provider::create_chat_model_with_model_id( + "summarization", + config, + config.default_temperature, + ) + .map_err(|e| format!("tree summarizer: failed to build cloud provider: {e:#}")) } /// Whether a summarization provider can be resolved for "Build Summary Trees" diff --git a/src/openhuman/orchestration/ops.rs b/src/openhuman/orchestration/ops.rs index 0ea71b396..bbd87d510 100644 --- a/src/openhuman/orchestration/ops.rs +++ b/src/openhuman/orchestration/ops.rs @@ -398,11 +398,13 @@ async fn synthesize_steering( use crate::openhuman::agent::turn_origin::{ with_origin, AgentTurnOrigin, TrustedAutomationSource, }; - use crate::openhuman::inference::provider::create_chat_provider; + use crate::openhuman::inference::provider::create_chat_model; + use tinyagents::harness::message::Message; + use tinyagents::harness::model::{ChatModel, ModelRequest}; for attempt in 1..=2 { - let (provider, model) = match create_chat_provider("subconscious", config) { - Ok(pm) => pm, + let model = match create_chat_model("subconscious", config, 0.3) { + Ok(m) => m, Err(e) => { log::warn!(target: LOG, "[orchestration] review.provider_unavailable: {e}"); return None; @@ -412,11 +414,13 @@ async fn synthesize_steering( job_id: tick_id.to_string(), source: TrustedAutomationSource::SubconsciousTainted, }; - match with_origin( - origin, - provider.chat_with_system(Some(STEERING_SYNTH_SYSTEM), prompt, &model, 0.3), - ) - .await + let request = ModelRequest::new(vec![ + Message::system(STEERING_SYNTH_SYSTEM), + Message::user(prompt), + ]); + match with_origin(origin, model.invoke(&(), request)) + .await + .map(|response| response.text()) { Ok(text) => { if let Some(parsed) = parse_steering_output(&text) { diff --git a/src/openhuman/tinyagents/mod.rs b/src/openhuman/tinyagents/mod.rs index 9828cde99..73867bb71 100644 --- a/src/openhuman/tinyagents/mod.rs +++ b/src/openhuman/tinyagents/mod.rs @@ -54,6 +54,7 @@ use tinyagents::harness::middleware::{ ToolPolicyMiddleware as TaToolPolicyMiddleware, }; use tinyagents::harness::model::CapabilitySet; +use tinyagents::harness::retry::RetryPolicy; use tinyagents::harness::runtime::{AgentHarness, RunPolicy, UnknownToolPolicy}; use tinyagents::harness::steering::SteeringHandle; use tinyagents::harness::store::StoreRegistry; @@ -114,14 +115,24 @@ pub(crate) struct ToolPolicyEnforcement { /// the tinyagents default of 25 — far more than openhuman's `max_iterations`. /// The recursion depth cap is also set here so TinyAgents uses OpenHuman's /// existing sub-agent spawn depth instead of the SDK default. -/// Retry is set to a single attempt: the openhuman [`Provider`] already does its -/// own internal retry/backoff (via the still-wrapped `ReliableProvider`), so a -/// second harness-level retry layer would double-retry transient errors and, -/// worse, swallow a deterministic provider error when a mock/test provider yields -/// a different result on the retry. This pin stays until `ReliableProvider` is -/// un-wrapped in the 02.2 conformance pass (Workstream 11); the crate's -/// exp-backoff [`RetryPolicy`](tinyagents::harness::retry::RetryPolicy) fields -/// stay at the default schedule so raising `max_attempts` later is a one-line flip. +/// Retry is now owned by the crate [`RetryPolicy`] (issue #4249, Phase 3a): the +/// turn path no longer wraps its provider in `ReliableProvider` (removed in +/// `session/builder/factory.rs`), so the single retry layer is here, at the +/// harness model call. The schedule mirrors the former `ReliableProvider` +/// defaults — 2 retries (3 attempts) with 500 ms exponential backoff — so +/// transient 429/5xx behavior is preserved. Retryability is decided by the crate +/// `is_retryable`, which the [`ProviderModel`](super::model) adapter feeds +/// correctly: a permanent config/auth/quota/context error is mapped to a +/// non-retryable `TinyAgentsError::Validation`, a transient blip to a retryable +/// `Model` error. The crate caps `max_attempts` at +/// `RunLimits::max_retries_per_call + 1` (default 3 retries), so this stays +/// within the loop's own bound. +/// +/// (Config parity note: the former `config.reliability.provider_retries` / +/// `provider_backoff_ms` / `model_fallbacks` no longer drive the turn path — +/// retry is the fixed schedule below and cross-route fallback is the crate +/// registry `FallbackPolicy` from [`routes::route_fallback_policy`]. Those config +/// knobs still apply to the non-seam `ReliableProvider` paths.) /// /// Cross-route **fallback** (`RunPolicy.fallback`) is orthogonal to retry and is /// populated per-turn by the caller ([`assemble_turn_harness`] via @@ -133,7 +144,17 @@ fn run_policy_for(max_iterations: usize, response_cache_enabled: bool) -> RunPol policy.limits.max_model_calls = max_iterations; policy.limits.max_tool_calls = max_iterations.saturating_mul(8).max(8); policy.limits.max_depth = MAX_SPAWN_DEPTH; - policy.retry.max_attempts = 1; + // Crate-owned retry (Phase 3a): mirror the former `ReliableProvider` schedule + // (2 retries, 500 ms exponential backoff). `backoff_sleep` is on so a + // transient 429/5xx actually waits before retrying, as it did before. + policy.retry = RetryPolicy { + max_attempts: 3, + initial_backoff_ms: 500, + max_backoff_ms: 30_000, + multiplier: 2.0, + jitter: false, + backoff_sleep: true, + }; // Unknown-tool recovery (01.2 / C3): the crate policy owns this end to end — // the `__openhuman_unknown_tool__` sentinel tool + `UnknownToolRewriteMiddleware` // were already deleted. We deliberately keep `ReturnToolError` rather than @@ -494,7 +515,6 @@ pub(crate) async fn run_turn_via_tinyagents_shared( subagent_scope.clone(), context_window, early_exit_tools, - max_output_tokens, context_mw, tool_policy, routes::turn_required_capabilities(model), @@ -545,7 +565,7 @@ pub(crate) async fn run_turn_via_tinyagents_shared( ); } - let config = RunConfig::new("agent_turn") + let mut config = RunConfig::new("agent_turn") .with_max_model_calls(max_iterations) .with_max_tool_calls(max_iterations.saturating_mul(8).max(8)) .with_max_depth(MAX_SPAWN_DEPTH) @@ -560,6 +580,13 @@ pub(crate) async fn run_turn_via_tinyagents_shared( } else { "unobserved" }); + // Per-turn output cap rides RunConfig now (Phase 5 groundwork): the loop + // stamps it onto every `ModelRequest.max_tokens` and the ProviderModel + // adapter honors it, so the cap no longer bakes into the primary + route + // models. Mirrors the legacy `AGENT_TURN_MAX_OUTPUT_TOKENS` / sub-agent cap. + if let Some(cap) = max_output_tokens { + config = config.with_max_turn_output_tokens(cap); + } tracing::info!( model, @@ -1160,7 +1187,6 @@ fn assemble_turn_harness( subagent_scope: Option, context_window: Option, early_exit_tools: &[&str], - max_output_tokens: Option, context_mw: TurnContextMiddleware, tool_policy: Option, required_capabilities: Option, @@ -1212,13 +1238,10 @@ fn assemble_turn_harness( let summary_provider = provider.clone(); let mut provider_model = ProviderModel::new(provider, model, temperature) .with_usage_carry(provider_usage_carry.clone()); - // Cap the model's per-call output budget (parity with the legacy engine, - // which bounded the main agent at `AGENT_TURN_MAX_OUTPUT_TOKENS` and each - // sub-agent at its `max_turn_output_tokens`). Without this the tinyagents - // path ran the provider uncapped. - if let Some(cap) = max_output_tokens { - provider_model = provider_model.with_max_tokens(cap); - } + // The per-call output cap now rides `RunConfig.max_turn_output_tokens` + // (Phase 5 groundwork), set by the caller: the loop stamps it onto every + // `ModelRequest` and the adapter honors `request.max_tokens`, so the cap no + // longer needs to be baked into the primary model or each route model. // Record the model's context window on its capability profile (issue #4249, // Phase 2) so the crate can validate input capacity before dispatch. if let Some(window) = context_window.filter(|w| *w > 0) { @@ -1249,13 +1272,9 @@ fn assemble_turn_harness( // the retained provider handle (the other clone was consumed into the // primary `ProviderModel`); `build_route_models` clones it per route and // skips the turn's own model so we don't shadow the default. - for route in routes::build_route_models( - &summary_provider, - temperature, - model, - max_output_tokens, - &provider_usage_carry, - ) { + for route in + routes::build_route_models(&summary_provider, temperature, model, &provider_usage_carry) + { let routes::RouteModel { name, model: route_model, diff --git a/src/openhuman/tinyagents/model.rs b/src/openhuman/tinyagents/model.rs index acf6eb964..f7fea1972 100644 --- a/src/openhuman/tinyagents/model.rs +++ b/src/openhuman/tinyagents/model.rs @@ -25,7 +25,7 @@ use crate::openhuman::agent::progress::AgentProgress; use crate::openhuman::inference::provider::thread_context::{current_thread_id, with_thread_id}; use crate::openhuman::inference::provider::{ current_route_slot, with_route_slot, ChatMessage, ChatRequest, ChatResponse, Provider, - ProviderDelta, + ProviderDelta, UsageInfo, }; use crate::openhuman::tools::ToolSpec; @@ -245,13 +245,14 @@ fn response_to_model_response( content.push(block); } let usage = response.usage.as_ref().map(|u| { - // Carry the provider's cached-prefix input count through the crate - // `Usage` (it has a `cache_read_tokens` field) so downstream cost - // accounting can price it at the cached rate. `Usage::new` seeds - // input/output/total; set the cache field on top. (`charged_amount_usd` - // has no crate home; the event bridge estimates cost from token counts.) + // Carry every token breakdown the crate `Usage` can express so a + // standalone `invoke` is usage-faithful (gap G1): cache reads/writes and + // reasoning tokens all have crate homes as of tinyagents 1.7. `Usage::new` + // seeds input/output/total; set the detail fields on top. let mut usage = Usage::new(u.input_tokens, u.output_tokens); usage.cache_read_tokens = u.cached_input_tokens; + usage.cache_creation_tokens = u.cache_creation_tokens; + usage.reasoning_tokens = u.reasoning_tokens; usage }); let finish_reason = if tool_calls.is_empty() { @@ -268,11 +269,78 @@ fn response_to_model_response( }, usage, finish_reason: Some(finish_reason.to_string()), - raw: None, + // The crate `Usage` has no field for the provider's **charged USD** or the + // model's **context window**, but `ModelResponse.raw` is exactly "raw + // provider metadata preserved for callers who need it" — so stash them + // there (gap G1). A standalone `invoke` then round-trips the OpenHuman + // managed backend's charged amount + window via + // [`usage_info_from_response`], no crate change required. Omitted when the + // provider reported neither (keeps non-managed responses byte-clean). + raw: openhuman_usage_meta_raw(response.usage.as_ref()), resolved_model: None, } } +/// JSON key under which the model adapter stashes the provider-reported +/// billing/context metadata that the crate [`Usage`] has no field for +/// (gap G1). Consumed by [`usage_info_from_response`]. +const OPENHUMAN_USAGE_META_KEY: &str = "openhuman_usage_meta"; + +/// The two host [`UsageInfo`] fields with no crate [`Usage`] home, ferried +/// through [`ModelResponse::raw`] so a standalone `invoke` stays usage-faithful. +#[derive(Debug, Clone, Copy, Default, serde::Serialize, serde::Deserialize)] +struct OpenhumanUsageMeta { + /// Provider-charged amount in USD (`UsageInfo::charged_amount_usd`). + #[serde(default)] + charged_amount_usd: f64, + /// Model context window in tokens (`UsageInfo::context_window`). + #[serde(default)] + context_window: u64, +} + +/// Build the `ModelResponse.raw` value carrying charged-USD + context-window +/// metadata, or `None` when the provider reported neither (so responses from +/// providers that don't surface billing stay `raw: None`). +fn openhuman_usage_meta_raw(usage: Option<&UsageInfo>) -> Option { + let u = usage?; + if u.charged_amount_usd <= 0.0 && u.context_window == 0 { + return None; + } + let meta = OpenhumanUsageMeta { + charged_amount_usd: u.charged_amount_usd, + context_window: u.context_window, + }; + Some(serde_json::json!({ OPENHUMAN_USAGE_META_KEY: meta })) +} + +/// Reconstruct a host [`UsageInfo`] from a crate [`ModelResponse`], recovering +/// the provider-charged USD + context window the adapter stashed in +/// [`ModelResponse::raw`] (gap G1). Returns `None` when the response carried no +/// usage at all. +/// +/// This is the seam one-shot inference callers use once they move off +/// `Box` (`chat` → `UsageInfo`) onto `Arc` +/// (`invoke` → `ModelResponse`): the full host usage record — real token +/// counts *and* backend-charged USD — survives the crossing. +pub(crate) fn usage_info_from_response(response: &ModelResponse) -> Option { + let usage = response.usage.as_ref()?; + let meta = response + .raw + .as_ref() + .and_then(|v| v.get(OPENHUMAN_USAGE_META_KEY)) + .and_then(|v| serde_json::from_value::(v.clone()).ok()) + .unwrap_or_default(); + Some(UsageInfo { + input_tokens: usage.input_tokens, + output_tokens: usage.output_tokens, + context_window: meta.context_window, + cached_input_tokens: usage.cache_read_tokens, + cache_creation_tokens: usage.cache_creation_tokens, + reasoning_tokens: usage.reasoning_tokens, + charged_amount_usd: meta.charged_amount_usd, + }) +} + /// Forward one openhuman [`ProviderDelta`]. Visible text, reasoning, and /// tool-call **argument** fragments all become harness [`ModelStreamItem`]s (so /// the [`OpenhumanEventBridge`](super::OpenhumanEventBridge) mirrors them as @@ -501,6 +569,12 @@ impl ChatModel<()> for ProviderModel { ) -> tinyagents::Result { let native = self.provider.supports_native_tools(); let (messages, specs) = build_chat_inputs(&request, native); + // Honor a per-request temperature when the caller sets one (e.g. one-shot + // inference callers that reuse a single model across prompts of differing + // temperature), else fall back to the temperature pinned at construction. + // The agent-loop seam leaves `request.temperature` `None`, so the pinned + // value still governs every turn — behaviour-neutral for the harness path. + let temperature = request.temperature.unwrap_or(self.temperature); // Positional layouts for the text-mode P-Format fallback (issue #4465); // empty (and thus behaviour-neutral) when no tools are advertised. let pformat_registry = pformat_registry_from_request(&request); @@ -512,7 +586,12 @@ impl ChatModel<()> for ProviderModel { // payload would defeat the opt-out and get rejected/ignored. tools: (native && !specs.is_empty()).then_some(&specs), stream: None, - max_tokens: self.max_tokens, + // Prefer a per-request output cap when the caller set one, else the + // cap pinned at construction. The agent-loop seam pins via + // `with_max_tokens` and leaves `request.max_tokens` `None` + // (openhuman never sets the crate `RunConfig.max_turn_output_tokens`), + // so the pinned cap still governs every turn. + max_tokens: request.max_tokens.or(self.max_tokens), }; tracing::debug!( @@ -524,7 +603,7 @@ impl ChatModel<()> for ProviderModel { let response = match self .provider - .chat(chat_request, &self.model, self.temperature) + .chat(chat_request, &self.model, temperature) .await { Ok(response) => { @@ -611,8 +690,11 @@ impl ChatModel<()> for ProviderModel { let pformat_registry = pformat_registry_from_request(&request); let provider = self.provider.clone(); let model = self.model.clone(); - let temperature = self.temperature; - let max_tokens = self.max_tokens; + // Per-request temperature when set (see `invoke`), else the pinned value; + // the agent-loop seam leaves it `None`, so streamed turns are unchanged. + let temperature = request.temperature.unwrap_or(self.temperature); + // Same precedence for the output cap (see `invoke`). + let max_tokens = request.max_tokens.or(self.max_tokens); let thinking = self.thinking.clone(); let error_slot = self.error_slot.clone(); // Captured for the spawned producer (task-locals/`self` do not cross the @@ -765,3 +847,162 @@ impl ChatModel<()> for ProviderModel { Ok(Box::pin(stream)) } } + +#[cfg(test)] +mod g1_usage_tests { + //! Gap G1: a standalone `invoke` must stay usage-faithful — token + //! breakdowns ride the crate `Usage`, and the two host fields with no crate + //! home (charged USD + context window) ride `ModelResponse.raw` and + //! reconstruct exactly via [`usage_info_from_response`]. + use super::*; + + fn empty_registry() -> crate::openhuman::agent::pformat::PFormatRegistry { + crate::openhuman::agent::pformat::PFormatRegistry::default() + } + + #[test] + fn usage_round_trips_charged_usd_and_all_token_breakdowns() { + let chat = ChatResponse { + text: Some("hi".to_string()), + tool_calls: Vec::new(), + usage: Some(UsageInfo { + input_tokens: 100, + output_tokens: 20, + context_window: 128_000, + cached_input_tokens: 40, + cache_creation_tokens: 10, + reasoning_tokens: 7, + charged_amount_usd: 0.0123, + }), + reasoning_content: None, + }; + let model_response = response_to_model_response(&chat, &empty_registry()); + + // Crate Usage carries every token breakdown natively. + let usage = model_response.usage.expect("usage present"); + assert_eq!(usage.input_tokens, 100); + assert_eq!(usage.output_tokens, 20); + assert_eq!(usage.cache_read_tokens, 40); + assert_eq!(usage.cache_creation_tokens, 10); + assert_eq!(usage.reasoning_tokens, 7); + + // Charged USD + context window ride raw and reconstruct exactly. + let recovered = usage_info_from_response(&model_response).expect("usage info"); + assert_eq!(recovered.input_tokens, 100); + assert_eq!(recovered.output_tokens, 20); + assert_eq!(recovered.context_window, 128_000); + assert_eq!(recovered.cached_input_tokens, 40); + assert_eq!(recovered.cache_creation_tokens, 10); + assert_eq!(recovered.reasoning_tokens, 7); + assert!((recovered.charged_amount_usd - 0.0123).abs() < 1e-9); + } + + #[test] + fn no_billing_metadata_leaves_raw_clean() { + let chat = ChatResponse { + text: Some("hi".to_string()), + tool_calls: Vec::new(), + usage: Some(UsageInfo { + input_tokens: 5, + output_tokens: 3, + ..Default::default() + }), + reasoning_content: None, + }; + let model_response = response_to_model_response(&chat, &empty_registry()); + assert!( + model_response.raw.is_none(), + "no charged USD / window ⇒ raw stays None" + ); + let recovered = usage_info_from_response(&model_response).expect("usage info"); + assert_eq!(recovered.charged_amount_usd, 0.0); + assert_eq!(recovered.context_window, 0); + assert_eq!(recovered.input_tokens, 5); + } + + #[test] + fn no_usage_reconstructs_to_none() { + let chat = ChatResponse { + text: Some("hi".to_string()), + tool_calls: Vec::new(), + usage: None, + reasoning_content: None, + }; + let model_response = response_to_model_response(&chat, &empty_registry()); + assert!(usage_info_from_response(&model_response).is_none()); + } +} + +#[cfg(test)] +mod adapter_param_tests { + //! The adapter honors a per-request temperature / output cap when the caller + //! sets one (one-shot callers reuse a model across differing prompts), and + //! otherwise the value pinned at construction (the agent-loop seam path). + use super::*; + use tinyagents::harness::message::Message; + + #[derive(Default)] + struct CaptureProvider { + seen: Arc)>>>, + } + + #[async_trait] + impl Provider for CaptureProvider { + async fn chat_with_system( + &self, + _system_prompt: Option<&str>, + _message: &str, + _model: &str, + _temperature: f64, + ) -> anyhow::Result { + unreachable!("chat() is overridden") + } + + async fn chat( + &self, + request: ChatRequest<'_>, + _model: &str, + temperature: f64, + ) -> anyhow::Result { + self.seen + .lock() + .unwrap() + .push((temperature, request.max_tokens)); + Ok(ChatResponse { + text: Some("ok".to_string()), + ..Default::default() + }) + } + } + + #[tokio::test] + async fn per_request_overrides_win_else_pinned() { + let seen: Arc)>>> = Arc::default(); + let provider: Arc = Arc::new(CaptureProvider { seen: seen.clone() }); + let model = ProviderModel::new(provider, "m", 0.7).with_max_tokens(100); + + // Request carries its own temperature + cap → those win. + model + .invoke( + &(), + ModelRequest::new(vec![Message::user("x")]) + .with_temperature(0.1) + .with_max_tokens(42), + ) + .await + .unwrap(); + // Request leaves both unset → pinned construction values apply. + model + .invoke(&(), ModelRequest::new(vec![Message::user("x")])) + .await + .unwrap(); + + let seen = seen.lock().unwrap(); + assert_eq!( + seen[0], + (0.1, Some(42)), + "per-request temperature + cap win" + ); + assert_eq!(seen[1], (0.7, Some(100)), "unset falls back to pinned"); + } +} diff --git a/src/openhuman/tinyagents/routes.rs b/src/openhuman/tinyagents/routes.rs index a9d49c540..b077d4f76 100644 --- a/src/openhuman/tinyagents/routes.rs +++ b/src/openhuman/tinyagents/routes.rs @@ -88,7 +88,6 @@ pub(super) fn build_route_models( provider: &Arc, temperature: f64, skip_model: &str, - max_output_tokens: Option, // Shared provider-usage carry (#4467): fallback route models must drain the // SAME carry the bridge reads, or a successful fallback call never feeds its // backend-charged USD / context-window / cache-creation-reasoning breakdown @@ -115,9 +114,10 @@ pub(super) fn build_route_models( .with_vision(vision) .with_reasoning(reasoning) .with_usage_carry(usage_carry.clone()); - if let Some(cap) = max_output_tokens { - model = model.with_max_tokens(cap); - } + // The per-turn output cap now rides `RunConfig.max_turn_output_tokens` + // (Phase 5 groundwork): the loop stamps it onto every `ModelRequest`, so + // route models no longer bake it in — they carry only model identity + + // capability profile. if let Some(window) = window.filter(|w| *w > 0) { model = model.with_context_window(window); } diff --git a/src/openhuman/tinyagents/tests.rs b/src/openhuman/tinyagents/tests.rs index e4fcac2fc..facd39d7e 100644 --- a/src/openhuman/tinyagents/tests.rs +++ b/src/openhuman/tinyagents/tests.rs @@ -446,7 +446,6 @@ fn adapter_inventory_registers_model_tools_and_middleware() { None, // subagent_scope: top-level turn Some(200_000), // known context window → compression + trim install &["ask_user_clarification"], - Some(1024), TurnContextMiddleware::defaults(), None, // no builder tool policy on this path None, // no per-turn required capabilities @@ -553,7 +552,13 @@ fn adapter_inventory_registers_model_tools_and_middleware() { assert!(profile.tool_calling, "EchoThenDone supports native tools"); assert!(!profile.modalities.image_in, "no vision on the mock"); assert_eq!(profile.max_input_tokens, Some(200_000), "context window"); - assert_eq!(profile.max_output_tokens, Some(1024), "output cap"); + // The per-turn output cap now rides `RunConfig.max_turn_output_tokens` + // (Phase 5 groundwork), not the model profile, so the profile carries no + // output cap. + assert_eq!( + profile.max_output_tokens, None, + "output cap rides RunConfig" + ); } /// The context-management middlewares gate on a known context window: without @@ -578,7 +583,6 @@ fn adapter_inventory_gates_context_middleware_on_window() { None, None, // unknown context window &[], // no early-exit tools - None, TurnContextMiddleware::defaults(), None, None, // no per-turn required capabilities