From 9665b319008c2bd64252573c2de8fb646646f289 Mon Sep 17 00:00:00 2001 From: Steven Enamakel <31011319+senamakel@users.noreply.github.com> Date: Mon, 13 Jul 2026 02:58:46 -0700 Subject: [PATCH] feat: expose paid Medulla orchestration (#4807) --- docs/TEST-COVERAGE-MATRIX.md | 1 + .../developing/architecture/orchestration.md | 17 +- src/openhuman/about_app/catalog_data.rs | 7 +- src/openhuman/config/mod.rs | 31 +- src/openhuman/config/schema/mod.rs | 5 +- src/openhuman/config/schema/orchestration.rs | 120 ++- src/openhuman/orchestration/medulla.rs | 699 ++++++++++++++++++ src/openhuman/orchestration/mod.rs | 1 + src/openhuman/orchestration/schemas.rs | 48 +- 9 files changed, 897 insertions(+), 32 deletions(-) create mode 100644 src/openhuman/orchestration/medulla.rs diff --git a/docs/TEST-COVERAGE-MATRIX.md b/docs/TEST-COVERAGE-MATRIX.md index bc9daf23c..422bbf066 100644 --- a/docs/TEST-COVERAGE-MATRIX.md +++ b/docs/TEST-COVERAGE-MATRIX.md @@ -501,6 +501,7 @@ End-to-end coverage of the agent harness via the web-chat RPC surface against an | ID | Feature | Layer | Test path(s) | Status | Notes | | ------ | ---------------------------------------------------------------- | ----- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | | 11.3.1 | Hosted-only orchestration (client = trigger + effects + render) | RI+VU | `tests/orchestration_hosted_client.rs`, `app/src/components/intelligence/TinyPlaceOrchestrationTab.test.tsx`, `app/src/components/orchestration/__tests__/AgentChatPanel.test.tsx` | ✅ | Local wake-graph brain retired (frontend_agent/graph/master_agent/reasoning_agent deleted). Client forwards events to the hosted brain (`POST /orchestration/v1/events`), uploads world-diffs, syncs hosted sessions/messages/steering into the render cache, and executes `send_dm`/`evict` socket effects; cloud-unreachable banner on outage. | +| 11.3.2 | Direct paid Medulla orchestration with local OpenHuman tools | RU | `src/openhuman/orchestration/medulla.rs`, `src/openhuman/orchestration/schemas.rs` | ✅ | `openhuman.orchestration_run` checks the active paid plan, starts a hosted Medulla cycle, executes requested contact/session/send tools locally, and continues pending/tool-use events to a final result. Mocked HTTP tests cover direct and tool-loop success plus plan, pending, backend-error, unknown-tool, and tool-failure paths. | --- diff --git a/gitbooks/developing/architecture/orchestration.md b/gitbooks/developing/architecture/orchestration.md index 8f22b2dbb..1d81c7109 100644 --- a/gitbooks/developing/architecture/orchestration.md +++ b/gitbooks/developing/architecture/orchestration.md @@ -82,6 +82,14 @@ chat message. The Brain → Orchestration tab (`TinyPlaceOrchestrationTab.tsx` + `useOrchestrationChats.ts`) reads real store classification, live-updates, and lets the owner steer the front-end agent from the Master composer. +The `openhuman.orchestration_run` renderer RPC is the direct request/response +entry point for the backend's paid Medulla engine. It performs a paid-plan +preflight, advertises OpenHuman's local contact, session-history, and +send-to-agent tools, and executes requested calls on the device. Tool results +are returned through `/orchestration/v1/run/continue` until Medulla produces a +final reply. The backend repeats authentication, entitlement, and resource +limit enforcement on every cycle. + ## Running unattended (stage 8) - **No message loss**: ingest dedupes by relay `message_id` *before* decrypt (the @@ -103,6 +111,9 @@ the owner steer the front-end agent from the Master composer. ## Configuration -The `[orchestration]` config block (`src/openhuman/config/schema/orchestration.rs`): -`enabled`, `debounce_ms`, `max_supersteps`, `message_window`, -`context_evict_threshold` (clamped 0.8 to 0.9), `subagent_concurrency`. +The `[orchestration]` config block (`src/openhuman/config/schema/orchestration.rs`) +uses `enabled` as its master switch. Optional direct-cycle overrides are grouped +under `[orchestration.medulla.prompt_overrides]`, +`[orchestration.medulla.config]`, and `[orchestration.medulla.limits]`. Omitted +values retain backend defaults, and the backend clamps all supplied graph and +resource limits. diff --git a/src/openhuman/about_app/catalog_data.rs b/src/openhuman/about_app/catalog_data.rs index 16b3dfcba..a5ce31534 100644 --- a/src/openhuman/about_app/catalog_data.rs +++ b/src/openhuman/about_app/catalog_data.rs @@ -1815,10 +1815,11 @@ pub(super) const CAPABILITIES: &[Capability] = &[ forwards session DMs to the hosted orchestration brain, which reasons, \ replies, and steers on its own cadence server-side. The device executes the \ effects the brain pushes back (send the reply, mirror an eviction into local \ - memory), renders the hosted read surface, and shows a notice when the cloud \ - brain is unreachable.", + memory), renders the hosted read surface, and can run the paid Medulla API \ + directly with its local contact, session-history, and send-to-agent tools.", how_to: "Intelligence > Orchestration (pair a wrapped session, then chat via the Master \ - window).", + window), or call openhuman.orchestration_run. Prompt, graph, and resource \ + overrides live under [orchestration.medulla] in config.toml.", status: CapabilityStatus::Beta, privacy: DERIVED_TO_BACKEND, }, diff --git a/src/openhuman/config/mod.rs b/src/openhuman/config/mod.rs index ddd859401..2964e508e 100644 --- a/src/openhuman/config/mod.rs +++ b/src/openhuman/config/mod.rs @@ -38,21 +38,22 @@ pub use schema::{ DockerRuntimeConfig, EmailConfig, EmbeddingRouteConfig, GitbooksConfig, HeartbeatConfig, HttpHeader, HttpRequestConfig, IMessageConfig, IntegrationToggle, IntegrationsConfig, LarkConfig, LearningConfig, LinqConfig, LlmBackend, LocalAiConfig, MatrixConfig, McpAuthConfig, - McpClientConfig, McpClientIdentityConfig, McpServerConfig, MeetConfig, MemoryConfig, - MemoryTreeConfig, ModelRouteConfig, MultimodalConfig, MultimodalFileConfig, - ObservabilityConfig, OrchestratorModelConfig, PolymarketClobCredentials, PolymarketConfig, - PrivacyConfig, PrivacyMode, ProxyConfig, ProxyScope, ReflectionSource, ReliabilityConfig, - ResourceLimitsConfig, RuntimeConfig, SandboxBackend, SandboxConfig, SchedulerConfig, - SchedulerGateConfig, SchedulerGateMode, ScreenIntelligenceConfig, SearchConfig, SearchEngine, - SearchEngineCredentials, SearxngConfig, SecretsConfig, SecurityConfig, ShellConfig, - SlackConfig, StorageConfig, StorageProviderConfig, StorageProviderSection, StreamMode, - TeamModelConfig, TelegramConfig, TokenjuiceConfig, UpdateConfig, UpdateRestartStrategy, - VoiceActivationMode, VoiceServerConfig, WebSearchConfig, WebhookConfig, YuanbaoConfig, - DEFAULT_CLOUD_LLM_MODEL, DEFAULT_MEMORY_SYNC_INTERVAL_SECS, DEFAULT_MODEL, - MEMORY_SYNC_INTERVAL_PRESETS_SECS, MODEL_AGENTIC_V1, MODEL_BURST_V1, MODEL_CHAT_V1, - MODEL_CODING_V1, MODEL_REASONING_QUICK_V1, MODEL_REASONING_V1, MODEL_SUMMARIZATION_V1, - MODEL_VISION_V1, SEARCH_ENGINE_BRAVE, SEARCH_ENGINE_DISABLED, SEARCH_ENGINE_MANAGED, - SEARCH_ENGINE_PARALLEL, SEARCH_ENGINE_QUERIT, + McpClientConfig, McpClientIdentityConfig, McpServerConfig, MedullaClientConfig, + MedullaCycleConfig, MedullaCycleLimits, MedullaPromptOverrides, MedullaVerification, + MeetConfig, MemoryConfig, MemoryTreeConfig, ModelRouteConfig, MultimodalConfig, + MultimodalFileConfig, ObservabilityConfig, OrchestratorModelConfig, PolymarketClobCredentials, + PolymarketConfig, PrivacyConfig, PrivacyMode, ProxyConfig, ProxyScope, ReflectionSource, + ReliabilityConfig, ResourceLimitsConfig, RuntimeConfig, SandboxBackend, SandboxConfig, + SchedulerConfig, SchedulerGateConfig, SchedulerGateMode, ScreenIntelligenceConfig, + SearchConfig, SearchEngine, SearchEngineCredentials, SearxngConfig, SecretsConfig, + SecurityConfig, ShellConfig, SlackConfig, StorageConfig, StorageProviderConfig, + StorageProviderSection, StreamMode, TeamModelConfig, TelegramConfig, TokenjuiceConfig, + UpdateConfig, UpdateRestartStrategy, VoiceActivationMode, VoiceServerConfig, WebSearchConfig, + WebhookConfig, YuanbaoConfig, DEFAULT_CLOUD_LLM_MODEL, DEFAULT_MEMORY_SYNC_INTERVAL_SECS, + DEFAULT_MODEL, MEMORY_SYNC_INTERVAL_PRESETS_SECS, MODEL_AGENTIC_V1, MODEL_BURST_V1, + MODEL_CHAT_V1, MODEL_CODING_V1, MODEL_REASONING_QUICK_V1, MODEL_REASONING_V1, + MODEL_SUMMARIZATION_V1, MODEL_VISION_V1, SEARCH_ENGINE_BRAVE, SEARCH_ENGINE_DISABLED, + SEARCH_ENGINE_MANAGED, SEARCH_ENGINE_PARALLEL, SEARCH_ENGINE_QUERIT, }; pub use schemas::{ all_controller_schemas as all_config_controller_schemas, diff --git a/src/openhuman/config/schema/mod.rs b/src/openhuman/config/schema/mod.rs index 0b5873456..2bf48e97c 100644 --- a/src/openhuman/config/schema/mod.rs +++ b/src/openhuman/config/schema/mod.rs @@ -71,7 +71,10 @@ pub use local_ai::{LocalAiConfig, LocalAiUsage}; pub use meet::{AutoJoinPolicy, AutoSummarizePolicy, CalendarProvider, MeetConfig}; pub use node::NodeConfig; pub use observability::{AgentTracingBackend, AgentTracingConfig, ObservabilityConfig}; -pub use orchestration::OrchestrationConfig; +pub use orchestration::{ + MedullaClientConfig, MedullaCycleConfig, MedullaCycleLimits, MedullaPromptOverrides, + MedullaVerification, OrchestrationConfig, +}; pub use privacy::{PrivacyConfig, PrivacyMode}; pub use proxy::{ apply_runtime_proxy_to_builder, build_runtime_proxy_client, diff --git a/src/openhuman/config/schema/orchestration.rs b/src/openhuman/config/schema/orchestration.rs index 738f1db3b..000b9c58f 100644 --- a/src/openhuman/config/schema/orchestration.rs +++ b/src/openhuman/config/schema/orchestration.rs @@ -1,8 +1,4 @@ -//! Orchestration configuration. -//! -//! The reasoning/wake graph runs server-side now, so the device-side config is a -//! single opt-out: whether to ingest tiny.place harness session DMs and forward -//! them to the hosted orchestration brain. +//! Hosted orchestration configuration. use schemars::JsonSchema; use serde::{Deserialize, Serialize}; @@ -14,18 +10,124 @@ fn default_enabled() -> bool { #[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] #[serde(default)] pub struct OrchestrationConfig { - /// Ingest inbound tiny.place harness session DMs, forward them to the hosted - /// brain (`POST /orchestration/v1/events`), run the device tail (effect - /// executor, world-diff uploader, health probe), and render the hosted read - /// surface. When `false` the device does none of this. Default: `true`. + /// Master switch for hosted orchestration ingest, rendering, and direct + /// Medulla runs. Default: `true`. #[serde(default = "default_enabled")] pub enabled: bool, + + /// Options forwarded to the paid `/orchestration/v1/run` surface. + #[serde(default)] + pub medulla: MedullaClientConfig, } impl Default for OrchestrationConfig { fn default() -> Self { Self { enabled: default_enabled(), + medulla: MedullaClientConfig::default(), } } } + +#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)] +#[serde(default)] +pub struct MedullaClientConfig { + /// Optional prompt replacements. Empty fields retain backend defaults. + #[serde(default)] + pub prompt_overrides: MedullaPromptOverrides, + + /// Optional graph tuning. The backend clamps every supplied value. + #[serde(default)] + pub config: MedullaCycleConfig, + + /// Optional per-cycle resource limits. The backend clamps every value. + #[serde(default)] + pub limits: MedullaCycleLimits, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)] +#[serde(default)] +pub struct MedullaPromptOverrides { + pub orchestrate_system: Option, + pub reasoning_execute_system: Option, + pub orchestrate_rlm_system: Option, + pub compress_system: Option, + pub frontend_gate_system: Option, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)] +#[serde(default)] +pub struct MedullaCycleConfig { + pub max_passes: Option, + pub max_steps: Option, + pub max_depth: Option, + pub context_window_tokens: Option, + pub verification: Option, +} + +#[derive(Debug, Clone, Copy, Serialize, Deserialize, JsonSchema)] +#[serde(rename_all = "snake_case")] +pub enum MedullaVerification { + Remind, + Off, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema)] +#[serde(default)] +pub struct MedullaCycleLimits { + pub max_concurrency: Option, + pub max_tokens: Option, + pub deadline_ms: Option, + pub max_tasks_per_delegate: Option, + pub max_depth: Option, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn defaults_keep_orchestration_enabled_without_overrides() { + let config = OrchestrationConfig::default(); + assert!(config.enabled); + assert!(config.medulla.config.max_passes.is_none()); + assert!(config.medulla.limits.deadline_ms.is_none()); + assert!(config.medulla.prompt_overrides.orchestrate_system.is_none()); + } + + #[test] + fn medulla_nested_toml_deserializes() { + let config: OrchestrationConfig = toml::from_str( + r#" +enabled = true + +[medulla.prompt_overrides] +orchestrate_system = "Plan carefully" + +[medulla.config] +max_passes = 3 +verification = "remind" + +[medulla.limits] +max_concurrency = 8 +deadline_ms = 90000 +"#, + ) + .unwrap(); + + assert_eq!(config.medulla.config.max_passes, Some(3)); + assert!(matches!( + config.medulla.config.verification, + Some(MedullaVerification::Remind) + )); + assert_eq!(config.medulla.limits.max_concurrency, Some(8)); + assert_eq!( + config + .medulla + .prompt_overrides + .orchestrate_system + .as_deref(), + Some("Plan carefully") + ); + } +} diff --git a/src/openhuman/orchestration/medulla.rs b/src/openhuman/orchestration/medulla.rs new file mode 100644 index 000000000..534d9af20 --- /dev/null +++ b/src/openhuman/orchestration/medulla.rs @@ -0,0 +1,699 @@ +//! Client for the backend's paid Medulla request/response surface. +//! +//! OpenHuman supplies its existing local orchestration tools to +//! `/orchestration/v1/run`, executes requested calls on-device, and returns the +//! results through `/run/continue` until the cycle ends. + +use std::sync::Arc; + +use reqwest::Method; +use serde::{Deserialize, Serialize}; +use serde_json::{json, Map, Value}; + +use crate::api::config::effective_backend_api_url; +use crate::api::BackendOAuthClient; +use crate::openhuman::config::{Config, MedullaClientConfig}; +use crate::openhuman::tools::Tool; + +const LOG: &str = "orchestration"; +const RUN_PATH: &str = "/orchestration/v1/run"; +const CONTINUE_PATH: &str = "/orchestration/v1/run/continue"; +const PLAN_PATH: &str = "/payments/stripe/currentPlan"; +const MAX_LOOP_EVENTS: usize = 128; + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct MedullaRunResult { + pub reply: String, + #[serde(default)] + pub pass_count: u32, + #[serde(default)] + pub compressed_history: Vec, + #[serde(default)] + pub escalations: Vec, + pub session_id: String, + pub cycle_id: String, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct ToolCall { + id: String, + name: String, + #[serde(default)] + args: Value, +} + +#[derive(Debug, Deserialize)] +#[serde(tag = "stop", rename_all = "snake_case")] +enum LoopEvent { + ToolUse { + #[serde(rename = "cycleId")] + cycle_id: String, + #[serde(rename = "sessionId")] + session_id: String, + #[serde(rename = "toolCalls")] + tool_calls: Vec, + }, + End { + #[serde(flatten)] + result: MedullaRunResult, + }, + Pending { + #[serde(rename = "cycleId")] + cycle_id: String, + #[serde(rename = "sessionId")] + session_id: String, + }, + Error { + #[serde(rename = "cycleId")] + cycle_id: String, + error: String, + }, +} + +/// Run a Medulla cycle using the signed-in user's paid backend plan. The local +/// plan lookup is a fail-fast UX guard; the backend remains authoritative and +/// repeats the paid-plan check on both run endpoints. +pub async fn run( + config: &Config, + input: &str, + session_id: Option<&str>, +) -> Result { + if !config.orchestration.enabled { + return Err("hosted orchestration is disabled in config".to_string()); + } + let input = input.trim(); + if input.is_empty() { + return Err("input is required".to_string()); + } + + let token = crate::openhuman::credentials::session_support::require_live_session_token(config)?; + let api_url = effective_backend_api_url(&config.api_url); + let client = BackendOAuthClient::new(&api_url).map_err(|err| err.to_string())?; + ensure_paid_plan(&client, &token).await?; + + let tuning = config.orchestration.medulla.clone(); + let config = Arc::new(config.clone()); + let tools: Vec> = vec![ + Box::new(super::tools::ListContactsTool), + Box::new(super::tools::ListSessionsTool::new(Arc::clone(&config))), + Box::new(super::tools::ReadSessionTool::new(Arc::clone(&config))), + Box::new(super::tools::SendToAgentTool::new(config)), + ]; + + run_with_client(&client, &token, input, session_id, &tools, &tuning).await +} + +async fn ensure_paid_plan(client: &BackendOAuthClient, token: &str) -> Result<(), String> { + let data = client + .authed_json(token, Method::GET, PLAN_PATH, None) + .await + .map_err(crate::api::flatten_authed_error)?; + if paid_plan_active(&data) { + return Ok(()); + } + Err("Medulla orchestration requires an active Basic or Pro plan".to_string()) +} + +fn paid_plan_active(data: &Value) -> bool { + let plan = data.get("plan").and_then(Value::as_str).unwrap_or_default(); + let active = data + .get("hasActiveSubscription") + .and_then(Value::as_bool) + .unwrap_or(false); + active && matches!(plan.to_ascii_uppercase().as_str(), "BASIC" | "PRO") +} + +async fn run_with_client( + client: &BackendOAuthClient, + token: &str, + input: &str, + session_id: Option<&str>, + tools: &[Box], + tuning: &MedullaClientConfig, +) -> Result { + let mut body = Map::new(); + body.insert("input".to_string(), Value::String(input.to_string())); + if let Some(session_id) = session_id.map(str::trim).filter(|id| !id.is_empty()) { + body.insert( + "sessionId".to_string(), + Value::String(session_id.to_string()), + ); + } + if !tools.is_empty() { + body.insert( + "tools".to_string(), + Value::Array(tools.iter().map(|tool| tool_spec(tool.as_ref())).collect()), + ); + } + if let Some(options) = tuning_value(tuning) { + body.insert("options".to_string(), options); + } + + tracing::debug!( + input_bytes = input.len(), + tool_count = tools.len(), + has_session = session_id.is_some(), + "[orchestration] medulla.run.start" + ); + let first = post(client, token, RUN_PATH, Value::Object(body)).await?; + if tools.is_empty() { + let result = serde_json::from_value(first) + .map_err(|err| format!("parse Medulla run response: {err}"))?; + tracing::debug!("[orchestration] medulla.run.end direct=true"); + return Ok(result); + } + + drive_tool_loop(client, token, tools, first).await +} + +async fn drive_tool_loop( + client: &BackendOAuthClient, + token: &str, + tools: &[Box], + mut value: Value, +) -> Result { + for event_index in 0..MAX_LOOP_EVENTS { + let event: LoopEvent = serde_json::from_value(value) + .map_err(|err| format!("parse Medulla loop event: {err}"))?; + match event { + LoopEvent::End { result } => { + tracing::debug!( + cycle_id = %result.cycle_id, + pass_count = result.pass_count, + events = event_index + 1, + "[orchestration] medulla.run.end" + ); + return Ok(result); + } + LoopEvent::Error { cycle_id, error } => { + tracing::warn!(cycle_id = %cycle_id, "[orchestration] medulla.run.error"); + return Err(format!("Medulla cycle {cycle_id} failed: {error}")); + } + LoopEvent::Pending { + cycle_id, + session_id, + } => { + tracing::debug!( + cycle_id = %cycle_id, + session_id = %session_id, + "[orchestration] medulla.run.pending" + ); + value = continue_run(client, token, &cycle_id, Vec::new()).await?; + } + LoopEvent::ToolUse { + cycle_id, + session_id, + tool_calls, + } => { + tracing::debug!( + cycle_id = %cycle_id, + session_id = %session_id, + call_count = tool_calls.len(), + "[orchestration] medulla.run.tool_use" + ); + let mut results = Vec::with_capacity(tool_calls.len()); + for call in tool_calls { + results.push(execute_tool_call(tools, call).await); + } + value = continue_run(client, token, &cycle_id, results).await?; + } + } + } + Err(format!( + "Medulla tool loop exceeded {MAX_LOOP_EVENTS} events" + )) +} + +async fn execute_tool_call(tools: &[Box], call: ToolCall) -> Value { + let Some(tool) = tools.iter().find(|tool| tool.name() == call.name) else { + tracing::warn!(tool = %call.name, "[orchestration] medulla.tool.unknown"); + return json!({ + "id": call.id, + "ok": false, + "error": format!("unknown OpenHuman tool: {}", call.name), + }); + }; + + tracing::debug!(tool = %call.name, call_id = %call.id, "[orchestration] medulla.tool.start"); + match tool.execute(call.args).await { + Ok(result) if !result.is_error => { + tracing::debug!(tool = %call.name, call_id = %call.id, "[orchestration] medulla.tool.end"); + json!({ "id": call.id, "ok": true, "result": result.output_for_llm(true) }) + } + Ok(result) => { + tracing::warn!(tool = %call.name, call_id = %call.id, "[orchestration] medulla.tool.failed"); + json!({ "id": call.id, "ok": false, "error": result.output() }) + } + Err(err) => { + tracing::warn!(tool = %call.name, call_id = %call.id, error = %err, "[orchestration] medulla.tool.failed"); + json!({ "id": call.id, "ok": false, "error": err.to_string() }) + } + } +} + +async fn continue_run( + client: &BackendOAuthClient, + token: &str, + cycle_id: &str, + tool_results: Vec, +) -> Result { + post( + client, + token, + CONTINUE_PATH, + json!({ "cycleId": cycle_id, "toolResults": tool_results }), + ) + .await +} + +async fn post( + client: &BackendOAuthClient, + token: &str, + path: &str, + body: Value, +) -> Result { + client + .authed_json(token, Method::POST, path, Some(body)) + .await + .map_err(crate::api::flatten_authed_error) +} + +fn tool_spec(tool: &dyn Tool) -> Value { + json!({ + "name": tool.name(), + "description": tool.description(), + "parameters": tool.parameters_schema(), + }) +} + +fn tuning_value(tuning: &MedullaClientConfig) -> Option { + let prompt_overrides = &tuning.prompt_overrides; + let prompt_overrides = json_object([ + ( + "ORCHESTRATE_SYSTEM", + prompt_overrides + .orchestrate_system + .as_ref() + .map(|v| json!(v)), + ), + ( + "REASONING_EXECUTE_SYSTEM", + prompt_overrides + .reasoning_execute_system + .as_ref() + .map(|v| json!(v)), + ), + ( + "ORCHESTRATE_RLM_SYSTEM", + prompt_overrides + .orchestrate_rlm_system + .as_ref() + .map(|v| json!(v)), + ), + ( + "COMPRESS_SYSTEM", + prompt_overrides.compress_system.as_ref().map(|v| json!(v)), + ), + ( + "FRONTEND_GATE_SYSTEM", + prompt_overrides + .frontend_gate_system + .as_ref() + .map(|v| json!(v)), + ), + ]); + let config = &tuning.config; + let config = json_object([ + ("maxPasses", config.max_passes.map(|v| json!(v))), + ("maxSteps", config.max_steps.map(|v| json!(v))), + ("maxDepth", config.max_depth.map(|v| json!(v))), + ( + "contextWindowTokens", + config.context_window_tokens.map(|v| json!(v)), + ), + ( + "verification", + config.verification.map(|v| { + json!(match v { + crate::openhuman::config::MedullaVerification::Remind => "remind", + crate::openhuman::config::MedullaVerification::Off => "off", + }) + }), + ), + ]); + let limits = &tuning.limits; + let limits = json_object([ + ("maxConcurrency", limits.max_concurrency.map(|v| json!(v))), + ("maxTokens", limits.max_tokens.map(|v| json!(v))), + ("deadlineMs", limits.deadline_ms.map(|v| json!(v))), + ( + "maxTasksPerDelegate", + limits.max_tasks_per_delegate.map(|v| json!(v)), + ), + ("maxDepth", limits.max_depth.map(|v| json!(v))), + ]); + let options = json_object([ + ("promptOverrides", prompt_overrides.map(Value::Object)), + ("config", config.map(Value::Object)), + ("limits", limits.map(Value::Object)), + ]); + options.map(Value::Object) +} + +fn json_object(entries: [(&str, Option); N]) -> Option> { + let object: Map = entries + .into_iter() + .filter_map(|(key, value)| value.map(|value| (key.to_string(), value))) + .collect(); + (!object.is_empty()).then_some(object) +} + +#[cfg(test)] +mod tests { + use super::*; + use async_trait::async_trait; + use wiremock::matchers::{header, method, path}; + use wiremock::{Mock, MockServer, ResponseTemplate}; + + use crate::openhuman::config::{MedullaCycleConfig, MedullaPromptOverrides}; + use crate::openhuman::tools::ToolResult; + + struct EchoTool; + + #[async_trait] + impl Tool for EchoTool { + fn name(&self) -> &str { + "echo" + } + + fn description(&self) -> &str { + "Echo text" + } + + fn parameters_schema(&self) -> Value { + json!({"type":"object","properties":{"text":{"type":"string"}},"required":["text"]}) + } + + async fn execute(&self, args: Value) -> anyhow::Result { + Ok(ToolResult::success( + args.get("text").and_then(Value::as_str).unwrap_or_default(), + )) + } + } + + struct FailingTool; + + #[async_trait] + impl Tool for FailingTool { + fn name(&self) -> &str { + "fail" + } + + fn description(&self) -> &str { + "Return a tool-level failure" + } + + fn parameters_schema(&self) -> Value { + json!({"type":"object"}) + } + + async fn execute(&self, _args: Value) -> anyhow::Result { + Ok(ToolResult::error("expected failure")) + } + } + + struct ExplodingTool; + + #[async_trait] + impl Tool for ExplodingTool { + fn name(&self) -> &str { + "explode" + } + + fn description(&self) -> &str { + "Return an execution error" + } + + fn parameters_schema(&self) -> Value { + json!({"type":"object"}) + } + + async fn execute(&self, _args: Value) -> anyhow::Result { + anyhow::bail!("execution exploded") + } + } + + fn end_event() -> Value { + json!({ + "stop":"end", + "reply":"done", + "passCount":2, + "compressedHistory":[], + "escalations":[], + "cycleId":"cycle-1", + "sessionId":"session-1" + }) + } + + fn envelope(data: Value) -> Value { + json!({"success": true, "data": data}) + } + + #[test] + fn paid_plan_requires_active_basic_or_pro() { + assert!(paid_plan_active( + &json!({"plan":"PRO","hasActiveSubscription":true}) + )); + assert!(paid_plan_active( + &json!({"plan":"basic","hasActiveSubscription":true}) + )); + assert!(!paid_plan_active( + &json!({"plan":"FREE","hasActiveSubscription":true}) + )); + assert!(!paid_plan_active( + &json!({"plan":"PRO","hasActiveSubscription":false}) + )); + } + + #[test] + fn tuning_uses_backend_field_names_and_omits_empty_sections() { + assert!(tuning_value(&MedullaClientConfig::default()).is_none()); + let tuning = MedullaClientConfig { + prompt_overrides: MedullaPromptOverrides { + orchestrate_system: Some("custom".to_string()), + ..Default::default() + }, + config: MedullaCycleConfig { + max_passes: Some(3), + ..Default::default() + }, + ..Default::default() + }; + assert_eq!( + tuning_value(&tuning), + Some(json!({ + "promptOverrides": {"ORCHESTRATE_SYSTEM":"custom"}, + "config": {"maxPasses":3} + })) + ); + } + + #[tokio::test] + async fn paid_plan_check_accepts_paid_and_rejects_free() { + for (plan, active, should_pass) in [("PRO", true, true), ("FREE", true, false)] { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path(PLAN_PATH)) + .and(header("authorization", "Bearer test-token")) + .respond_with(ResponseTemplate::new(200).set_body_json(envelope(json!({ + "plan": plan, + "hasActiveSubscription": active + })))) + .mount(&server) + .await; + + let client = BackendOAuthClient::new(&server.uri()).unwrap(); + let result = ensure_paid_plan(&client, "test-token").await; + assert_eq!(result.is_ok(), should_pass); + } + } + + #[tokio::test] + async fn public_run_validates_enabled_and_input_before_credentials() { + let mut config = Config::default(); + config.orchestration.enabled = false; + assert_eq!( + run(&config, "task", None).await.unwrap_err(), + "hosted orchestration is disabled in config" + ); + + config.orchestration.enabled = true; + assert_eq!( + run(&config, " ", None).await.unwrap_err(), + "input is required" + ); + } + + #[tokio::test] + async fn run_without_tools_returns_direct_result_and_forwards_options() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path(RUN_PATH)) + .and(header("authorization", "Bearer test-token")) + .respond_with(ResponseTemplate::new(200).set_body_json(envelope(end_event()))) + .mount(&server) + .await; + + let client = BackendOAuthClient::new(&server.uri()).unwrap(); + let tuning = MedullaClientConfig { + config: MedullaCycleConfig { + max_passes: Some(4), + ..Default::default() + }, + ..Default::default() + }; + let result = run_with_client( + &client, + "test-token", + "direct", + Some(" session-1 "), + &[], + &tuning, + ) + .await + .unwrap(); + + assert_eq!(result.reply, "done"); + let requests = server.received_requests().await.unwrap(); + let body: Value = serde_json::from_slice(&requests[0].body).unwrap(); + assert_eq!(body["sessionId"], "session-1"); + assert_eq!(body["options"]["config"]["maxPasses"], 4); + assert!(body.get("tools").is_none()); + } + + #[tokio::test] + async fn pending_event_polls_until_end() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path(CONTINUE_PATH)) + .respond_with(ResponseTemplate::new(200).set_body_json(envelope(end_event()))) + .mount(&server) + .await; + let client = BackendOAuthClient::new(&server.uri()).unwrap(); + + let result = drive_tool_loop( + &client, + "test-token", + &[Box::new(EchoTool)], + json!({"stop":"pending","cycleId":"cycle-1","sessionId":"session-1"}), + ) + .await + .unwrap(); + + assert_eq!(result.reply, "done"); + let requests = server.received_requests().await.unwrap(); + let body: Value = serde_json::from_slice(&requests[0].body).unwrap(); + assert_eq!(body, json!({"cycleId":"cycle-1","toolResults":[]})); + } + + #[tokio::test] + async fn error_event_is_returned_with_cycle_context() { + let server = MockServer::start().await; + let client = BackendOAuthClient::new(&server.uri()).unwrap(); + let error = drive_tool_loop( + &client, + "test-token", + &[], + json!({"stop":"error","cycleId":"cycle-bad","error":"budget exceeded"}), + ) + .await + .unwrap_err(); + + assert_eq!(error, "Medulla cycle cycle-bad failed: budget exceeded"); + } + + #[tokio::test] + async fn unknown_and_failed_tools_are_reported_to_backend() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path(CONTINUE_PATH)) + .respond_with(ResponseTemplate::new(200).set_body_json(envelope(end_event()))) + .mount(&server) + .await; + let client = BackendOAuthClient::new(&server.uri()).unwrap(); + let tools: Vec> = vec![Box::new(FailingTool), Box::new(ExplodingTool)]; + + drive_tool_loop( + &client, + "test-token", + &tools, + json!({ + "stop":"tool_use", + "cycleId":"cycle-1", + "sessionId":"session-1", + "toolCalls":[ + {"id":"call-unknown","name":"missing","args":{}}, + {"id":"call-failed","name":"fail","args":{}}, + {"id":"call-exploded","name":"explode","args":{}} + ] + }), + ) + .await + .unwrap(); + + let requests = server.received_requests().await.unwrap(); + let body: Value = serde_json::from_slice(&requests[0].body).unwrap(); + assert_eq!(body["toolResults"][0]["ok"], false); + assert_eq!( + body["toolResults"][0]["error"], + "unknown OpenHuman tool: missing" + ); + assert_eq!(body["toolResults"][1]["error"], "expected failure"); + assert_eq!(body["toolResults"][2]["error"], "execution exploded"); + } + + #[tokio::test] + async fn tool_loop_executes_locally_and_continues_to_end() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path(RUN_PATH)) + .and(header("authorization", "Bearer test-token")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "success": true, + "data": { + "stop":"tool_use", + "cycleId":"cycle-1", + "sessionId":"session-1", + "toolCalls":[{"id":"call-1","name":"echo","args":{"text":"hello"}}] + } + }))) + .mount(&server) + .await; + Mock::given(method("POST")) + .and(path(CONTINUE_PATH)) + .respond_with(ResponseTemplate::new(200).set_body_json(envelope(end_event()))) + .mount(&server) + .await; + + let client = BackendOAuthClient::new(&server.uri()).unwrap(); + let tools: Vec> = vec![Box::new(EchoTool)]; + let result = run_with_client( + &client, + "test-token", + "use echo", + None, + &tools, + &MedullaClientConfig::default(), + ) + .await + .unwrap(); + + assert_eq!(result.reply, "done"); + assert_eq!(result.pass_count, 2); + let requests = server.received_requests().await.unwrap(); + let continued: Value = serde_json::from_slice(&requests[1].body).unwrap(); + assert_eq!(continued["toolResults"][0]["result"], "hello"); + } +} diff --git a/src/openhuman/orchestration/mod.rs b/src/openhuman/orchestration/mod.rs index b0d320d9f..15df8f312 100644 --- a/src/openhuman/orchestration/mod.rs +++ b/src/openhuman/orchestration/mod.rs @@ -19,6 +19,7 @@ pub mod cloud; pub mod effect_executor; pub mod exec_gate; pub mod ingest; +pub mod medulla; pub mod migrate_history; pub mod ops; pub mod presence; diff --git a/src/openhuman/orchestration/schemas.rs b/src/openhuman/orchestration/schemas.rs index 56624c8ce..0d0b44d8b 100644 --- a/src/openhuman/orchestration/schemas.rs +++ b/src/openhuman/orchestration/schemas.rs @@ -36,6 +36,7 @@ pub fn all_controller_schemas() -> Vec { schema_for("orchestration_self_identity"), schema_for("orchestration_publish_identity"), schema_for("orchestration_relay_info"), + schema_for("orchestration_run"), ] } @@ -81,6 +82,10 @@ pub fn all_registered_controllers() -> Vec { schema: schema_for("orchestration_relay_info"), handler: handle_relay_info, }, + RegisteredController { + schema: schema_for("orchestration_run"), + handler: handle_medulla_run, + }, ] } @@ -178,6 +183,19 @@ fn schema_for(function: &str) -> ControllerSchema { inputs: vec![], outputs: vec![json_output("result", "{ baseUrl, network }.")], }, + "orchestration_run" => ControllerSchema { + namespace: "orchestration", + function: "run", + description: "Run the paid hosted Medulla engine with OpenHuman's local contact/session/send tools. Tool calls execute on this device and are returned through the backend continuation loop until a final reply is available.", + inputs: vec![ + required_str("input", "The task or prompt for Medulla to orchestrate."), + optional_str("sessionId", "Optional Medulla session id to continue."), + ], + outputs: vec![json_output( + "result", + "{ reply, passCount, compressedHistory, escalations, sessionId, cycleId }.", + )], + }, other => unreachable!("unknown orchestration schema: {other}"), } } @@ -1030,6 +1048,30 @@ fn handle_relay_info(_params: Map) -> ControllerFuture { }) } +/// Direct paid Medulla entry point. The backend remains the authority for both +/// authentication and plan enforcement; `medulla::run` also performs a local +/// plan preflight so free users receive an immediate actionable error. +fn handle_medulla_run(params: Map) -> ControllerFuture { + Box::pin(async move { + let input = required_param(¶ms, "input")?.trim().to_string(); + let session_id = params + .get("sessionId") + .and_then(Value::as_str) + .map(str::trim) + .filter(|value| !value.is_empty()) + .map(str::to_string); + let config = load_config("medulla_run").await?; + log::debug!( + target: LOG, + "[orchestration_rpc] medulla_run.entry input_bytes={} has_session={}", + input.len(), + session_id.is_some(), + ); + let result = super::medulla::run(&config, &input, session_id.as_deref()).await?; + to_json(result) + }) +} + // ── helpers ───────────────────────────────────────────────────────────────── async fn load_config(action: &str) -> Result { @@ -1087,7 +1129,7 @@ mod tests { #[test] fn schemas_use_orchestration_namespace() { let schemas = all_controller_schemas(); - assert_eq!(schemas.len(), 10); + assert_eq!(schemas.len(), 11); assert!(schemas.iter().all(|s| s.namespace == "orchestration")); assert_eq!(schema_for("orchestration_attention").function, "attention"); assert_eq!( @@ -1115,6 +1157,10 @@ mod tests { schema_for("orchestration_sessions_create").function, "sessions_create" ); + assert_eq!(schema_for("orchestration_run").function, "run"); + assert!(all_registered_controllers() + .iter() + .any(|controller| controller.schema.function == "run")); } #[test]