mirror of
https://github.com/tinyhumansai/openhuman.git
synced 2026-07-27 21:08:00 +00:00
fix(subconscious): route text-only tasks through cloud when local AI disabled (#1565)
Co-authored-by: Steven Enamakel <enamakel@tinyhumans.ai>
This commit is contained in:
co-authored by
Steven Enamakel
parent
3c542b1109
commit
00fda569e0
@@ -10,7 +10,35 @@ use super::prompt;
|
||||
use super::types::{ExecutionResult, SubconsciousTask};
|
||||
use tracing::{debug, info, warn};
|
||||
|
||||
#[cfg(test)]
|
||||
mod test_mocks {
|
||||
use std::sync::atomic::{AtomicU8, Ordering};
|
||||
|
||||
const MODE_REAL: u8 = 0;
|
||||
const MODE_LOCAL_FAIL: u8 = 1;
|
||||
const MODE_AGENT_FAIL: u8 = 2;
|
||||
|
||||
static MODE: AtomicU8 = AtomicU8::new(MODE_REAL);
|
||||
|
||||
pub fn mock_local() {
|
||||
MODE.store(MODE_LOCAL_FAIL, Ordering::Release);
|
||||
}
|
||||
pub fn mock_agent() {
|
||||
MODE.store(MODE_AGENT_FAIL, Ordering::Release);
|
||||
}
|
||||
pub fn reset() {
|
||||
MODE.store(MODE_REAL, Ordering::Release);
|
||||
}
|
||||
pub fn is_local_mocked() -> bool {
|
||||
MODE.load(Ordering::Acquire) == MODE_LOCAL_FAIL
|
||||
}
|
||||
pub fn is_agent_mocked() -> bool {
|
||||
MODE.load(Ordering::Acquire) == MODE_AGENT_FAIL
|
||||
}
|
||||
}
|
||||
|
||||
/// Outcome of executing a task — either completed or needs user approval.
|
||||
#[derive(Debug)]
|
||||
pub enum ExecutionOutcome {
|
||||
/// Task completed (either read-only analysis or approved write).
|
||||
Completed(ExecutionResult),
|
||||
@@ -31,14 +59,17 @@ pub async fn execute_task(
|
||||
) -> Result<ExecutionOutcome, String> {
|
||||
let started = std::time::Instant::now();
|
||||
let task_has_write_intent = needs_tools(&task.title);
|
||||
let mut config = crate::openhuman::config::Config::load_or_init()
|
||||
.await
|
||||
.map_err(|e| format!("config load: {e}"))?;
|
||||
|
||||
let result = if task_has_write_intent {
|
||||
// Task explicitly asks for a write action — run with full permissions.
|
||||
info!(
|
||||
"[subconscious:executor] write task: {} — agentic loop, full permissions",
|
||||
task.title
|
||||
"[subconscious:executor] write task: id={} — agentic loop, full permissions",
|
||||
task.id
|
||||
);
|
||||
execute_with_agent_full(task, situation_report, identity_context)
|
||||
execute_with_agent_full(&mut config, task, situation_report, identity_context)
|
||||
.await
|
||||
.map(|output| {
|
||||
ExecutionOutcome::Completed(ExecutionResult {
|
||||
@@ -50,10 +81,12 @@ pub async fn execute_task(
|
||||
} else if needs_agent(&task.title) {
|
||||
// Read-only task but needs deeper reasoning — run analysis-only.
|
||||
info!(
|
||||
"[subconscious:executor] read-only task escalated: {} — agentic loop, analysis only",
|
||||
task.title
|
||||
"[subconscious:executor] read-only task escalated: id={} — agentic loop, analysis only",
|
||||
task.id
|
||||
);
|
||||
let output = execute_with_agent_analysis(task, situation_report, identity_context).await?;
|
||||
let output =
|
||||
execute_with_agent_analysis(&mut config, task, situation_report, identity_context)
|
||||
.await?;
|
||||
let duration_ms = started.elapsed().as_millis() as u64;
|
||||
|
||||
if let Some(recommendation) = extract_recommended_action(&output) {
|
||||
@@ -70,24 +103,53 @@ pub async fn execute_task(
|
||||
}))
|
||||
}
|
||||
} else {
|
||||
// Simple text-only task — local model handles it.
|
||||
debug!(
|
||||
"[subconscious:executor] text task: {} — using local model",
|
||||
task.title
|
||||
);
|
||||
execute_with_local_model(task, situation_report, identity_context)
|
||||
.await
|
||||
.map(|output| {
|
||||
ExecutionOutcome::Completed(ExecutionResult {
|
||||
output,
|
||||
used_tools: false,
|
||||
duration_ms: started.elapsed().as_millis() as u64,
|
||||
// Simple text-only task. Use local model if configured for subconscious
|
||||
// tasks, otherwise fall back to the cloud agentic analysis path.
|
||||
if config.local_ai.use_local_for_subconscious() {
|
||||
debug!(
|
||||
"[subconscious:executor] text task: id={} — using local model",
|
||||
task.id
|
||||
);
|
||||
execute_with_local_model(&config, task, situation_report, identity_context)
|
||||
.await
|
||||
.map(|output| {
|
||||
ExecutionOutcome::Completed(ExecutionResult {
|
||||
output,
|
||||
used_tools: false,
|
||||
duration_ms: started.elapsed().as_millis() as u64,
|
||||
})
|
||||
})
|
||||
})
|
||||
} else {
|
||||
info!(
|
||||
"[subconscious:executor] text task: id={} — local AI disabled, using cloud fallback",
|
||||
task.id
|
||||
);
|
||||
let output =
|
||||
execute_with_agent_analysis(&mut config, task, situation_report, identity_context)
|
||||
.await
|
||||
.map_err(|e| format!("cloud fallback agent execution: {e}"))?;
|
||||
let duration_ms = started.elapsed().as_millis() as u64;
|
||||
debug!(
|
||||
"[subconscious:executor] text task cloud fallback complete: id={} — duration_ms={}",
|
||||
task.id, duration_ms
|
||||
);
|
||||
|
||||
// Suppress UnapprovedWrite: passive tasks that didn't trigger
|
||||
// needs_agent should never escalate even if the cloud model's
|
||||
// output contains RECOMMENDED ACTION. The write-intent gate is
|
||||
// needs_tools for active tasks and needs_agent for read-only
|
||||
// escalations; the cloud fallback is a passthrough for simple
|
||||
// text tasks and must not silently change the contract.
|
||||
Ok(ExecutionOutcome::Completed(ExecutionResult {
|
||||
output,
|
||||
used_tools: false,
|
||||
duration_ms,
|
||||
}))
|
||||
}
|
||||
};
|
||||
|
||||
if let Err(ref e) = result {
|
||||
warn!("[subconscious:executor] task '{}' failed: {e}", task.title);
|
||||
warn!("[subconscious:executor] task id={} failed: {e}", task.id);
|
||||
}
|
||||
|
||||
result
|
||||
@@ -95,13 +157,24 @@ pub async fn execute_task(
|
||||
|
||||
/// Execute an approved write action — called after user approves an escalation
|
||||
/// that originated from `UnapprovedWrite`.
|
||||
///
|
||||
/// Independent `Config::load_or_init()`: the task was originally routed under
|
||||
/// config_A in `execute_task`; now executes under config_B after user approval.
|
||||
/// If `use_local_for_subconscious()` toggled between the two calls, the approval
|
||||
/// was made under different assumptions. Risk is negligible in practice (config
|
||||
/// changes require a restart to take effect on most fields), but callers should
|
||||
/// be aware of this TOCTOU window.
|
||||
pub async fn execute_approved_write(
|
||||
task: &SubconsciousTask,
|
||||
situation_report: &str,
|
||||
identity_context: &str,
|
||||
) -> Result<ExecutionResult, String> {
|
||||
let started = std::time::Instant::now();
|
||||
let output = execute_with_agent_full(task, situation_report, identity_context).await?;
|
||||
let mut config = crate::openhuman::config::Config::load_or_init()
|
||||
.await
|
||||
.map_err(|e| format!("config load: {e}"))?;
|
||||
let output =
|
||||
execute_with_agent_full(&mut config, task, situation_report, identity_context).await?;
|
||||
Ok(ExecutionResult {
|
||||
output,
|
||||
used_tools: true,
|
||||
@@ -111,35 +184,18 @@ pub async fn execute_approved_write(
|
||||
|
||||
/// Execute a text-only task using the local Ollama model.
|
||||
///
|
||||
/// Gated by `local_ai.usage.subconscious`. When the flag is off (or
|
||||
/// `runtime_enabled` is off), returns `Err` so callers don't mistake a
|
||||
/// disabled subsystem for a successfully-completed empty execution.
|
||||
/// TODO: wire a cloud fallback here when use_local_for_subconscious is false.
|
||||
/// The caller MUST have already checked `config.local_ai.use_local_for_subconscious()`
|
||||
/// before calling this function.
|
||||
async fn execute_with_local_model(
|
||||
config: &crate::openhuman::config::Config,
|
||||
task: &SubconsciousTask,
|
||||
situation_report: &str,
|
||||
identity_context: &str,
|
||||
) -> Result<String, String> {
|
||||
let config = crate::openhuman::config::Config::load_or_init()
|
||||
.await
|
||||
.map_err(|e| format!("config load: {e}"))?;
|
||||
|
||||
if !config.local_ai.use_local_for_subconscious() {
|
||||
// Fail fast rather than returning Ok("") — upstream code uses an
|
||||
// empty string as a normal "task ran but produced no output"
|
||||
// signal, so a silent skip would mask a disabled subsystem as a
|
||||
// completed action. Surface the gate state so callers can
|
||||
// distinguish "skipped" from "succeeded with empty output".
|
||||
tracing::info!(
|
||||
"[subconscious:executor] local_ai.usage.subconscious not enabled — \
|
||||
refusing to execute task '{}' (no cloud fallback configured)",
|
||||
task.title
|
||||
);
|
||||
return Err(
|
||||
"local_ai.usage.subconscious not enabled (no cloud fallback configured)".to_string(),
|
||||
);
|
||||
#[cfg(test)]
|
||||
if test_mocks::is_local_mocked() {
|
||||
return Err("local model: mocked failure (test)".into());
|
||||
}
|
||||
|
||||
let prompt_text = prompt::build_text_execution_prompt(task, situation_report, identity_context);
|
||||
|
||||
let messages = vec![
|
||||
@@ -165,34 +221,32 @@ async fn execute_with_local_model(
|
||||
/// Retries up to 3 times with exponential backoff (2s, 4s, 8s) on 429 rate-limit
|
||||
/// errors from the agentic-v1 cloud model.
|
||||
async fn execute_with_agent_full(
|
||||
config: &mut crate::openhuman::config::Config,
|
||||
task: &SubconsciousTask,
|
||||
situation_report: &str,
|
||||
identity_context: &str,
|
||||
) -> Result<String, String> {
|
||||
let mut config = crate::openhuman::config::Config::load_or_init()
|
||||
.await
|
||||
.map_err(|e| format!("config load: {e}"))?;
|
||||
|
||||
let prompt_text = prompt::build_tool_execution_prompt(task, situation_report, identity_context);
|
||||
|
||||
agent_chat_with_retry(&mut config, &prompt_text).await
|
||||
agent_chat_with_retry(config, &prompt_text).await
|
||||
}
|
||||
|
||||
/// Execute with agentic-v1 in analysis-only mode (read-only tasks).
|
||||
///
|
||||
/// The prompt instructs the model to analyze but not execute write actions.
|
||||
async fn execute_with_agent_analysis(
|
||||
config: &mut crate::openhuman::config::Config,
|
||||
task: &SubconsciousTask,
|
||||
situation_report: &str,
|
||||
identity_context: &str,
|
||||
) -> Result<String, String> {
|
||||
let mut config = crate::openhuman::config::Config::load_or_init()
|
||||
.await
|
||||
.map_err(|e| format!("config load: {e}"))?;
|
||||
|
||||
#[cfg(test)]
|
||||
if test_mocks::is_agent_mocked() {
|
||||
return Err("cloud fallback: mocked failure (test)".into());
|
||||
}
|
||||
let prompt_text = prompt::build_analysis_only_prompt(task, situation_report, identity_context);
|
||||
|
||||
agent_chat_with_retry(&mut config, &prompt_text).await
|
||||
agent_chat_with_retry(config, &prompt_text).await
|
||||
}
|
||||
|
||||
/// Call agent_chat with rate-limit retry (429 only, up to 3 attempts).
|
||||
@@ -310,6 +364,121 @@ pub fn needs_tools(title: &str) -> bool {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use std::path::Path;
|
||||
use tempfile::tempdir;
|
||||
|
||||
/// Guard that sets an env var for the duration of the test and restores it on drop.
|
||||
struct EnvVarGuard {
|
||||
key: String,
|
||||
old: Option<String>,
|
||||
}
|
||||
|
||||
impl EnvVarGuard {
|
||||
fn set_to_path(key: &str, value: &Path) -> Self {
|
||||
let old = std::env::var(key).ok();
|
||||
std::env::set_var(key, value.to_str().expect("path is valid utf-8"));
|
||||
Self {
|
||||
key: key.to_string(),
|
||||
old,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for EnvVarGuard {
|
||||
fn drop(&mut self) {
|
||||
match &self.old {
|
||||
Some(value) => std::env::set_var(&self.key, value),
|
||||
None => std::env::remove_var(&self.key),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn write_subconscious_test_config(workspace_root: &Path, local_ai_enabled: bool) {
|
||||
let cfg = format!(
|
||||
r#"default_temperature = 0.7
|
||||
|
||||
[local_ai]
|
||||
runtime_enabled = {local_ai_enabled}
|
||||
provider = "ollama"
|
||||
|
||||
[local_ai.usage]
|
||||
subconscious = {local_ai_enabled}
|
||||
|
||||
[memory]
|
||||
backend = "sqlite"
|
||||
auto_save = true
|
||||
embedding_provider = "none"
|
||||
embedding_model = "none"
|
||||
embedding_dimensions = 0
|
||||
|
||||
[secrets]
|
||||
encrypt = false
|
||||
"#
|
||||
);
|
||||
std::fs::create_dir_all(workspace_root).expect("mkdir test workspace root");
|
||||
let config_path = workspace_root.join("config.toml");
|
||||
std::fs::write(&config_path, &cfg).expect("write test config");
|
||||
let _: crate::openhuman::config::Config =
|
||||
toml::from_str(&cfg).expect("test config should deserialize");
|
||||
}
|
||||
|
||||
fn make_text_task(title: &str) -> SubconsciousTask {
|
||||
SubconsciousTask {
|
||||
id: "test-id".into(),
|
||||
title: title.into(),
|
||||
source: super::super::types::TaskSource::User,
|
||||
recurrence: super::super::types::TaskRecurrence::Once,
|
||||
enabled: true,
|
||||
last_run_at: None,
|
||||
next_run_at: None,
|
||||
completed: false,
|
||||
created_at: 1700000000.0,
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_task_routes_to_cloud_when_local_disabled() {
|
||||
let _env_lock = crate::openhuman::config::TEST_ENV_LOCK
|
||||
.lock()
|
||||
.expect("test env lock");
|
||||
test_mocks::reset();
|
||||
let tmp = tempdir().expect("tempdir");
|
||||
let _workspace = EnvVarGuard::set_to_path("OPENHUMAN_WORKSPACE", tmp.path());
|
||||
write_subconscious_test_config(tmp.path(), false);
|
||||
|
||||
test_mocks::mock_agent();
|
||||
let task = make_text_task("Summarize unread emails");
|
||||
let result = execute_task(&task, "", "").await;
|
||||
|
||||
assert!(result.is_err(), "expected error (cloud path)");
|
||||
let err = result.unwrap_err();
|
||||
assert!(
|
||||
err.contains("cloud fallback"),
|
||||
"expected cloud fallback error, got: {err}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_task_routes_to_local_when_local_enabled() {
|
||||
let _env_lock = crate::openhuman::config::TEST_ENV_LOCK
|
||||
.lock()
|
||||
.expect("test env lock");
|
||||
test_mocks::reset();
|
||||
let tmp = tempdir().expect("tempdir");
|
||||
let _workspace = EnvVarGuard::set_to_path("OPENHUMAN_WORKSPACE", tmp.path());
|
||||
write_subconscious_test_config(tmp.path(), true);
|
||||
|
||||
test_mocks::mock_local();
|
||||
let task = make_text_task("Summarize unread emails");
|
||||
let result = execute_task(&task, "", "").await;
|
||||
|
||||
assert!(result.is_err(), "expected error (local path)");
|
||||
let err = result.unwrap_err();
|
||||
assert!(
|
||||
err.contains("local model"),
|
||||
"expected local model error, got: {err}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn needs_tools_detects_action_verbs() {
|
||||
|
||||
Reference in New Issue
Block a user