From 25e6ae96f33bdd90ebfd6e143d63dee679f5efd0 Mon Sep 17 00:00:00 2001 From: Steven Enamakel <31011319+senamakel@users.noreply.github.com> Date: Thu, 25 Jun 2026 09:20:47 -0700 Subject: [PATCH] =?UTF-8?q?refactor(subconscious):=20structured=20memory?= =?UTF-8?q?=5Fdiff=20=E2=86=92=20prepare=5Fcontext=20=E2=86=92=20decide=20?= =?UTF-8?q?tick=20(#4107)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/src/lib/i18n/ar.ts | 2 +- app/src/lib/i18n/bn.ts | 2 +- app/src/lib/i18n/de.ts | 2 +- app/src/lib/i18n/en.ts | 2 +- app/src/lib/i18n/es.ts | 2 +- app/src/lib/i18n/fr.ts | 2 +- app/src/lib/i18n/hi.ts | 2 +- app/src/lib/i18n/id.ts | 2 +- app/src/lib/i18n/it.ts | 2 +- app/src/lib/i18n/ko.ts | 2 +- app/src/lib/i18n/pl.ts | 2 +- app/src/lib/i18n/pt.ts | 2 +- app/src/lib/i18n/ru.ts | 2 +- app/src/lib/i18n/zh-CN.ts | 2 +- app/src/utils/tauriCommands/subconscious.ts | 7 +- src/core/subconscious_cli.rs | 52 +- src/openhuman/agent_orchestration/tools.rs | 4 +- .../tools/agent_prepare_context.rs | 23 +- src/openhuman/mcp_server/resources.rs | 2 +- src/openhuman/subconscious/README.md | 141 ++--- src/openhuman/subconscious/agent/agent.toml | 26 +- src/openhuman/subconscious/agent/prompt.md | 97 ++- src/openhuman/subconscious/engine.rs | 554 +++++++++--------- src/openhuman/subconscious/engine_tests.rs | 104 ++++ src/openhuman/subconscious/global.rs | 8 +- src/openhuman/subconscious/mod.rs | 2 - src/openhuman/subconscious/scratchpad/mod.rs | 422 ------------- .../subconscious/scratchpad/tools.rs | 230 -------- src/openhuman/subconscious/session.rs | 3 +- .../subconscious/situation_report/mod.rs | 214 ------- .../situation_report/query_window.rs | 150 ----- .../situation_report/summaries.rs | 114 ---- src/openhuman/subconscious/source_chunk.rs | 5 +- src/openhuman/subconscious/store.rs | 31 + src/openhuman/subconscious/store_tests.rs | 18 + src/openhuman/tools/ops.rs | 6 +- tests/subconscious_e2e.rs | 16 +- 37 files changed, 615 insertions(+), 1642 deletions(-) delete mode 100644 src/openhuman/subconscious/scratchpad/mod.rs delete mode 100644 src/openhuman/subconscious/scratchpad/tools.rs delete mode 100644 src/openhuman/subconscious/situation_report/mod.rs delete mode 100644 src/openhuman/subconscious/situation_report/query_window.rs delete mode 100644 src/openhuman/subconscious/situation_report/summaries.rs diff --git a/app/src/lib/i18n/ar.ts b/app/src/lib/i18n/ar.ts index 47b27a0ab..ada5ec3c7 100644 --- a/app/src/lib/i18n/ar.ts +++ b/app/src/lib/i18n/ar.ts @@ -2486,7 +2486,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'تم إيقاف اللاوعي مؤقتًا', 'subconscious.providerSettings': 'إعدادات الذكاء الاصطناعي', 'subconscious.scratchpadInfo': - 'يحتفظ اللاوعي بدفتر ملاحظات مستمر للملاحظات عبر الدورات. تحقق من الإعدادات → وصول الوكيل لتكوين الوضع والتكرار.', + 'في كل دورة، يتحقق اللاوعي مما تغيّر عبر مصادرك المتصلة، ثم يسجّل المهام في قائمة مهامك، أو يحدّث أهدافك، أو ينبّهك عندما يحتاج شيء ما إلى انتباهك. تحقق من الإعدادات → وصول الوكيل لتكوين الوضع والتكرار.', 'subconscious.approvalNeeded': 'يلزم الموافقة', 'subconscious.requiresApproval': 'يتطلب الموافقة', 'subconscious.fixInConnections': 'إصلاح في الاتصالات', diff --git a/app/src/lib/i18n/bn.ts b/app/src/lib/i18n/bn.ts index ee202ea39..a6a6d5ab3 100644 --- a/app/src/lib/i18n/bn.ts +++ b/app/src/lib/i18n/bn.ts @@ -2541,7 +2541,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconscious বিরত আছে', 'subconscious.providerSettings': 'AI সেটিংস', 'subconscious.scratchpadInfo': - 'সাবকনশাস টিক জুড়ে পর্যবেক্ষণের একটি স্থায়ী স্ক্র্যাচপ্যাড বজায় রাখে। মোড এবং ফ্রিকোয়েন্সি কনফিগার করতে সেটিংস → এজেন্ট অ্যাক্সেস দেখুন।', + 'প্রতিটি টিকে সাবকনশাস আপনার সংযুক্ত উৎসগুলিতে কী পরিবর্তিত হয়েছে তা যাচাই করে, তারপর আপনার করণীয় তালিকায় ফলো-আপ যোগ করে, আপনার লক্ষ্য আপডেট করে, অথবা কিছু মনোযোগের প্রয়োজন হলে আপনাকে জানায়। মোড এবং ফ্রিকোয়েন্সি কনফিগার করতে সেটিংস → এজেন্ট অ্যাক্সেস দেখুন।', 'subconscious.approvalNeeded': 'অনুমোদন প্রয়োজন', 'subconscious.requiresApproval': 'অনুমোদন প্রয়োজন', 'subconscious.fixInConnections': 'সংযোগে ঠিক করুন', diff --git a/app/src/lib/i18n/de.ts b/app/src/lib/i18n/de.ts index 9613362f9..d24d6c4a6 100644 --- a/app/src/lib/i18n/de.ts +++ b/app/src/lib/i18n/de.ts @@ -2599,7 +2599,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Unterbewusstsein ist pausiert', 'subconscious.providerSettings': 'KI-Einstellungen', 'subconscious.scratchpadInfo': - 'Das Unterbewusstsein führt ein dauerhaftes Notizbuch mit Beobachtungen über Ticks hinweg. Überprüfen Sie Einstellungen → Agentenzugriff, um Modus und Häufigkeit zu konfigurieren.', + 'Bei jedem Tick prüft das Unterbewusstsein, was sich über Ihre verbundenen Quellen hinweg geändert hat, und erfasst dann Folgeaufgaben auf Ihrer To-do-Liste, aktualisiert Ihre Ziele oder benachrichtigt Sie, wenn etwas Aufmerksamkeit erfordert. Unter Einstellungen → Agentenzugriff können Sie Modus und Häufigkeit konfigurieren.', 'subconscious.approvalNeeded': 'Genehmigung erforderlich', 'subconscious.requiresApproval': 'Erfordert eine Genehmigung', 'subconscious.fixInConnections': 'Fix in Verbindungen', diff --git a/app/src/lib/i18n/en.ts b/app/src/lib/i18n/en.ts index 95fe8cd3a..944269f07 100644 --- a/app/src/lib/i18n/en.ts +++ b/app/src/lib/i18n/en.ts @@ -3029,7 +3029,7 @@ const en: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconscious is paused', 'subconscious.providerSettings': 'AI settings', 'subconscious.scratchpadInfo': - 'The subconscious maintains a persistent scratchpad of observations across ticks. Check Settings → Agent access to configure mode and frequency.', + 'On each tick the subconscious checks what changed across your connected sources, then records follow-ups on your to-do list, updates your goals, or notifies you when something needs attention. Check Settings → Agent access to configure mode and frequency.', 'subconscious.approvalNeeded': 'Approval Needed', 'subconscious.requiresApproval': 'Requires approval', 'subconscious.fixInConnections': 'Fix in Connections', diff --git a/app/src/lib/i18n/es.ts b/app/src/lib/i18n/es.ts index 88cc3334f..a483cc325 100644 --- a/app/src/lib/i18n/es.ts +++ b/app/src/lib/i18n/es.ts @@ -2586,7 +2586,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconsciente en pausa', 'subconscious.providerSettings': 'Ajustes de IA', 'subconscious.scratchpadInfo': - 'El subconsciente mantiene un bloc de notas persistente de observaciones a través de los ciclos. Consulta Configuración → Acceso del agente para configurar el modo y la frecuencia.', + 'En cada ciclo, el subconsciente revisa qué cambió en tus fuentes conectadas y luego registra seguimientos en tu lista de tareas, actualiza tus objetivos o te notifica cuando algo requiere atención. Consulta Configuración → Acceso del agente para configurar el modo y la frecuencia.', 'subconscious.approvalNeeded': 'Se necesita aprobación', 'subconscious.requiresApproval': 'Requiere aprobación', 'subconscious.fixInConnections': 'Corregir en Conexiones', diff --git a/app/src/lib/i18n/fr.ts b/app/src/lib/i18n/fr.ts index 86e6f0aa8..aaeb3f60e 100644 --- a/app/src/lib/i18n/fr.ts +++ b/app/src/lib/i18n/fr.ts @@ -2598,7 +2598,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconscient en pause', 'subconscious.providerSettings': 'Paramètres IA', 'subconscious.scratchpadInfo': - "Le subconscient maintient un bloc-notes persistant d'observations à travers les cycles. Consultez Paramètres → Accès agent pour configurer le mode et la fréquence.", + "À chaque cycle, le subconscient examine ce qui a changé dans vos sources connectées, puis enregistre des suivis dans votre liste de tâches, met à jour vos objectifs ou vous notifie lorsqu'un élément nécessite votre attention. Consultez Paramètres → Accès agent pour configurer le mode et la fréquence.", 'subconscious.approvalNeeded': 'Approbation requise', 'subconscious.requiresApproval': 'Nécessite une approbation', 'subconscious.fixInConnections': 'Corriger dans Connexions', diff --git a/app/src/lib/i18n/hi.ts b/app/src/lib/i18n/hi.ts index af023c589..00713099b 100644 --- a/app/src/lib/i18n/hi.ts +++ b/app/src/lib/i18n/hi.ts @@ -2537,7 +2537,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconscious रुका हुआ है', 'subconscious.providerSettings': 'AI सेटिंग्स', 'subconscious.scratchpadInfo': - 'सबकॉन्शस टिक्स में अवलोकनों का एक स्थायी स्क्रैचपैड बनाए रखता है। मोड और आवृत्ति कॉन्फ़िगर करने के लिए सेटिंग्स → एजेंट एक्सेस देखें।', + 'हर टिक पर सबकॉन्शस जाँचता है कि आपके जुड़े हुए स्रोतों में क्या बदला, फिर आपकी कार्य-सूची में फॉलो-अप दर्ज करता है, आपके लक्ष्य अपडेट करता है, या जब किसी चीज़ पर ध्यान देने की ज़रूरत हो तो आपको सूचित करता है। मोड और आवृत्ति कॉन्फ़िगर करने के लिए सेटिंग्स → एजेंट एक्सेस देखें।', 'subconscious.approvalNeeded': 'अनुमति चाहिए', 'subconscious.requiresApproval': 'अनुमति ज़रूरी है', 'subconscious.fixInConnections': 'Connections में ठीक करें', diff --git a/app/src/lib/i18n/id.ts b/app/src/lib/i18n/id.ts index ec73ad8b4..1b2b09bf8 100644 --- a/app/src/lib/i18n/id.ts +++ b/app/src/lib/i18n/id.ts @@ -2541,7 +2541,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconscious dijeda', 'subconscious.providerSettings': 'Pengaturan AI', 'subconscious.scratchpadInfo': - 'Alam bawah sadar memelihara catatan pengamatan yang persisten di seluruh siklus. Periksa Pengaturan → Akses agen untuk mengonfigurasi mode dan frekuensi.', + 'Pada setiap siklus, alam bawah sadar memeriksa apa yang berubah di seluruh sumber yang terhubung, lalu mencatat tindak lanjut di daftar tugas Anda, memperbarui sasaran Anda, atau memberi tahu Anda saat ada sesuatu yang perlu diperhatikan. Periksa Pengaturan → Akses agen untuk mengonfigurasi mode dan frekuensi.', 'subconscious.approvalNeeded': 'Persetujuan Diperlukan', 'subconscious.requiresApproval': 'Memerlukan persetujuan', 'subconscious.fixInConnections': 'Perbaiki di Koneksi', diff --git a/app/src/lib/i18n/it.ts b/app/src/lib/i18n/it.ts index 8dce6308e..4ce407d0e 100644 --- a/app/src/lib/i18n/it.ts +++ b/app/src/lib/i18n/it.ts @@ -2580,7 +2580,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconscio in pausa', 'subconscious.providerSettings': 'Impostazioni IA', 'subconscious.scratchpadInfo': - 'Il subconscio mantiene un blocco appunti persistente di osservazioni attraverso i cicli. Controlla Impostazioni → Accesso agente per configurare modalità e frequenza.', + 'A ogni ciclo il subconscio verifica cosa è cambiato nelle tue fonti collegate, poi registra attività di follow-up nella tua lista di cose da fare, aggiorna i tuoi obiettivi o ti avvisa quando qualcosa richiede attenzione. Controlla Impostazioni → Accesso agente per configurare modalità e frequenza.', 'subconscious.approvalNeeded': 'Approvazione necessaria', 'subconscious.requiresApproval': 'Richiede approvazione', 'subconscious.fixInConnections': 'Correggi in Connessioni', diff --git a/app/src/lib/i18n/ko.ts b/app/src/lib/i18n/ko.ts index 61945e624..cfef73266 100644 --- a/app/src/lib/i18n/ko.ts +++ b/app/src/lib/i18n/ko.ts @@ -2514,7 +2514,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconscious 일시 중지됨', 'subconscious.providerSettings': 'AI 설정', 'subconscious.scratchpadInfo': - '잠재의식은 틱 전반에 걸쳐 관찰 사항의 영구 메모장을 유지합니다. 모드와 빈도를 구성하려면 설정 → 에이전트 접근을 확인하세요.', + '잠재의식은 매 틱마다 연결된 소스에서 무엇이 바뀌었는지 확인한 뒤, 할 일 목록에 후속 작업을 기록하거나 목표를 업데이트하거나 주의가 필요한 일이 생기면 알려줍니다. 모드와 빈도를 구성하려면 설정 → 에이전트 접근을 확인하세요.', 'subconscious.approvalNeeded': '승인 필요', 'subconscious.requiresApproval': '승인이 필요함', 'subconscious.fixInConnections': '연결에서 수정', diff --git a/app/src/lib/i18n/pl.ts b/app/src/lib/i18n/pl.ts index 7bc0118ea..ecd632eb9 100644 --- a/app/src/lib/i18n/pl.ts +++ b/app/src/lib/i18n/pl.ts @@ -2564,7 +2564,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Podświadomość wstrzymana', 'subconscious.providerSettings': 'Ustawienia AI', 'subconscious.scratchpadInfo': - 'Podświadomość utrzymuje trwały notatnik obserwacji między cyklami. Sprawdź Ustawienia → Dostęp agenta, aby skonfigurować tryb i częstotliwość.', + 'Przy każdym cyklu podświadomość sprawdza, co się zmieniło w połączonych źródłach, a następnie zapisuje zadania na Twojej liście rzeczy do zrobienia, aktualizuje Twoje cele lub powiadamia Cię, gdy coś wymaga uwagi. Sprawdź Ustawienia → Dostęp agenta, aby skonfigurować tryb i częstotliwość.', 'subconscious.approvalNeeded': 'Wymagana zgoda', 'subconscious.requiresApproval': 'Wymaga zgody', 'subconscious.fixInConnections': 'Napraw w Połączeniach', diff --git a/app/src/lib/i18n/pt.ts b/app/src/lib/i18n/pt.ts index 5bfc0cc2d..221ed0207 100644 --- a/app/src/lib/i18n/pt.ts +++ b/app/src/lib/i18n/pt.ts @@ -2585,7 +2585,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Subconsciente pausado', 'subconscious.providerSettings': 'Configurações de IA', 'subconscious.scratchpadInfo': - 'O subconsciente mantém um bloco de notas persistente de observações entre ciclos. Verifique Configurações → Acesso do agente para configurar o modo e a frequência.', + 'A cada ciclo, o subconsciente verifica o que mudou nas suas fontes conectadas e então registra acompanhamentos na sua lista de tarefas, atualiza suas metas ou notifica você quando algo precisa de atenção. Verifique Configurações → Acesso do agente para configurar o modo e a frequência.', 'subconscious.approvalNeeded': 'Aprovação Necessária', 'subconscious.requiresApproval': 'Requer aprovação', 'subconscious.fixInConnections': 'Corrigir em Conexões', diff --git a/app/src/lib/i18n/ru.ts b/app/src/lib/i18n/ru.ts index 026632b55..35200542b 100644 --- a/app/src/lib/i18n/ru.ts +++ b/app/src/lib/i18n/ru.ts @@ -2559,7 +2559,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': 'Подсознание приостановлено', 'subconscious.providerSettings': 'Настройки ИИ', 'subconscious.scratchpadInfo': - 'Подсознание ведёт постоянный блокнот наблюдений между циклами. Проверьте Настройки → Доступ агента для настройки режима и частоты.', + 'На каждом цикле подсознание проверяет, что изменилось в ваших подключённых источниках, а затем записывает задачи в ваш список дел, обновляет ваши цели или уведомляет вас, когда что-то требует внимания. Откройте Настройки → Доступ агента, чтобы настроить режим и частоту.', 'subconscious.approvalNeeded': 'Требуется подтверждение', 'subconscious.requiresApproval': 'Требует подтверждения', 'subconscious.fixInConnections': 'Исправить в подключениях', diff --git a/app/src/lib/i18n/zh-CN.ts b/app/src/lib/i18n/zh-CN.ts index f7bce2d0f..d598e42e0 100644 --- a/app/src/lib/i18n/zh-CN.ts +++ b/app/src/lib/i18n/zh-CN.ts @@ -2419,7 +2419,7 @@ const messages: TranslationMap = { 'subconscious.providerUnavailableTitle': '潜意识已暂停', 'subconscious.providerSettings': 'AI 设置', 'subconscious.scratchpadInfo': - '潜意识在每次循环中维护一个持久的观察记事本。请查看设置 → 代理访问来配置模式和频率。', + '潜意识在每次循环时检查你已连接来源中发生的变化,然后在待办事项列表中记录后续事项、更新你的目标,或在有需要关注的事情时通知你。请前往设置 → 代理访问来配置模式和频率。', 'subconscious.approvalNeeded': '需要审批', 'subconscious.requiresApproval': '需要审批', 'subconscious.fixInConnections': '在连接中修复', diff --git a/app/src/utils/tauriCommands/subconscious.ts b/app/src/utils/tauriCommands/subconscious.ts index 428ebd58c..103bbbf72 100644 --- a/app/src/utils/tauriCommands/subconscious.ts +++ b/app/src/utils/tauriCommands/subconscious.ts @@ -1,8 +1,9 @@ /** - * Subconscious engine commands — engine control and scratchpad. + * Subconscious engine commands — engine control (status / trigger). * - * Reflection/thoughts RPCs have been removed — the subconscious now - * maintains only a scratchpad (via agent tools) and run logs. + * The subconscious now runs a structured tick (memory_diff → prepare_context + * → decide); continuity lives in the user's global to-dos and goals rather + * than a scratchpad, so only the status/trigger RPCs are exposed here. */ import { callCoreRpc } from '../../services/coreRpcClient'; import { type CommandResponse, isTauri } from './common'; diff --git a/src/core/subconscious_cli.rs b/src/core/subconscious_cli.rs index 81fde3430..4829f72cc 100644 --- a/src/core/subconscious_cli.rs +++ b/src/core/subconscious_cli.rs @@ -3,7 +3,6 @@ //! Usage: //! openhuman subconscious tick [--workspace ] [--mode simple|aggressive] [--verbose] //! openhuman subconscious status [--workspace ] -//! openhuman subconscious scratchpad [--workspace ] use anyhow::{anyhow, Result}; use std::path::PathBuf; @@ -17,7 +16,6 @@ pub fn run_subconscious_command(args: &[String]) -> Result<()> { match args[0].as_str() { "tick" => run_tick(&args[1..]), "status" => run_status(&args[1..]), - "scratchpad" | "pad" => run_scratchpad(&args[1..]), other => Err(anyhow!( "unknown subconscious subcommand '{other}'. Run `openhuman subconscious --help`." )), @@ -143,9 +141,9 @@ fn run_tick(args: &[String]) -> Result<()> { return Err(anyhow!("provider unavailable: {reason}")); } - // Create engine and run tick - let memory = crate::openhuman::memory::global::client_if_ready(); - let engine = crate::openhuman::subconscious::SubconsciousEngine::new(&config, memory); + // Create engine and run tick. The engine pulls its own memory_diff / + // context state from the workspace — no memory client to pass in. + let engine = crate::openhuman::subconscious::SubconsciousEngine::new(&config); eprintln!("[subconscious] running tick..."); let result = engine @@ -159,12 +157,15 @@ fn run_tick(args: &[String]) -> Result<()> { ); if flags.verbose { - // Print scratchpad state after tick - let entries = crate::openhuman::subconscious::scratchpad::load(&config.workspace_dir) - .unwrap_or_default(); - if !entries.is_empty() { - eprintln!("\n[subconscious] scratchpad after tick:"); - println!("{}", serde_json::to_string_pretty(&entries)?); + // Print the world baseline the next tick will diff against. + let baseline = crate::openhuman::subconscious::store::with_connection( + &config.workspace_dir, + crate::openhuman::subconscious::store::get_baseline_checkpoint_id, + ) + .unwrap_or(None); + match baseline { + Some(id) => eprintln!("[subconscious] world baseline checkpoint: {id}"), + None => eprintln!("[subconscious] no world baseline established yet"), } } @@ -215,32 +216,6 @@ fn run_status(args: &[String]) -> Result<()> { }) } -// ── scratchpad ───────────────────────────────────────────────────────────── - -fn run_scratchpad(args: &[String]) -> Result<()> { - let workspace = parse_workspace_flag(args)?; - - let rt = tokio::runtime::Runtime::new()?; - rt.block_on(async { - let mut config = crate::openhuman::config::Config::load_or_init() - .await - .map_err(|e| anyhow!("config load failed: {e}"))?; - if let Some(ws) = workspace { - config.workspace_dir = ws; - } - - let entries = crate::openhuman::subconscious::scratchpad::load(&config.workspace_dir) - .map_err(|e| anyhow!("failed to read scratchpad: {e}"))?; - - if entries.is_empty() { - eprintln!("(scratchpad empty)"); - } else { - println!("{}", serde_json::to_string_pretty(&entries)?); - } - Ok(()) - }) -} - // ── helpers ──────────────────────────────────────────────────────────────── fn parse_workspace_flag(args: &[String]) -> Result> { @@ -272,12 +247,11 @@ fn print_help() { Commands: tick Run a single subconscious tick (synchronous, waits for completion) status Show current subconscious engine status - scratchpad Dump the persistent scratchpad Tick options: --mode Override the subconscious mode --workspace Override workspace directory - --verbose, -v Print scratchpad after tick + --verbose, -v Print the world baseline checkpoint after tick Common options: --workspace Override workspace directory diff --git a/src/openhuman/agent_orchestration/tools.rs b/src/openhuman/agent_orchestration/tools.rs index 1a83fa057..5f8feb356 100644 --- a/src/openhuman/agent_orchestration/tools.rs +++ b/src/openhuman/agent_orchestration/tools.rs @@ -32,7 +32,9 @@ mod worker_thread; pub(crate) use dispatch::dispatch_subagent; -pub use agent_prepare_context::{run_context_scout, AgentPrepareContextTool}; +pub use agent_prepare_context::{ + run_context_scout, run_context_scout_with_catalog, AgentPrepareContextTool, +}; pub use archetype_delegation::ArchetypeDelegationTool; pub use close_subagent::CloseSubagentTool; pub use continue_subagent::ContinueSubagentTool; diff --git a/src/openhuman/agent_orchestration/tools/agent_prepare_context.rs b/src/openhuman/agent_orchestration/tools/agent_prepare_context.rs index 64e5cd502..040da206a 100644 --- a/src/openhuman/agent_orchestration/tools/agent_prepare_context.rs +++ b/src/openhuman/agent_orchestration/tools/agent_prepare_context.rs @@ -70,6 +70,26 @@ fn is_well_formed_context_bundle(output: &str) -> bool { /// and runs the scout against the parent's provider. Outside a turn the /// `run_subagent` call surfaces a no-parent error as a [`ToolResult::error`]. pub async fn run_context_scout(question: &str, focus: Option<&str>) -> anyhow::Result { + let tool_catalog = AgentPrepareContextTool::render_parent_tool_catalog(); + run_context_scout_with_catalog(question, focus, &tool_catalog).await +} + +/// Same as [`run_context_scout`] but with an **explicitly-supplied** tool +/// catalogue, so it can run *outside* an agent turn — e.g. from the +/// subconscious engine's structured tick, where `current_parent()` is unset +/// and the parent's visible tool set can't be auto-derived. +/// +/// The caller passes the catalogue of tools the eventual decision agent can +/// actually call (one `- name: description` per line), so the bundle's +/// `recommended_tool_calls` stay grounded in callable tools. Progress / +/// subagent-lifecycle events stay best-effort: with no parent context the +/// `parent_session` falls back to `standalone` and the progress sink is absent, +/// so those sends simply no-op. +pub async fn run_context_scout_with_catalog( + question: &str, + focus: Option<&str>, + tool_catalog: &str, +) -> anyhow::Result { let question = question.trim().to_string(); let focus = focus.map(|s| s.to_string()); @@ -103,10 +123,9 @@ pub async fn run_context_scout(question: &str, focus: Option<&str>) -> anyhow::R } }; - let tool_catalog = AgentPrepareContextTool::render_parent_tool_catalog(); let catalog_tool_count = tool_catalog.lines().filter(|l| !l.is_empty()).count(); let scout_prompt = - AgentPrepareContextTool::build_scout_prompt(&question, focus.as_deref(), &tool_catalog); + AgentPrepareContextTool::build_scout_prompt(&question, focus.as_deref(), tool_catalog); tracing::debug!( target: "agent_prepare_context", diff --git a/src/openhuman/mcp_server/resources.rs b/src/openhuman/mcp_server/resources.rs index 2df85bd93..49ef248b0 100644 --- a/src/openhuman/mcp_server/resources.rs +++ b/src/openhuman/mcp_server/resources.rs @@ -244,7 +244,7 @@ const RESOURCE_CATALOG: &[PromptResource] = &[ PromptResource { uri: "openhuman://prompts/agents/subconscious", name: "subconscious", - description: "Background reasoning agent that maintains subconscious scratchpad context.", + description: "Background awareness agent: diffs the user's world, prepares context, and decides what to do.", content: include_str!("../subconscious/agent/prompt.md"), }, PromptResource { diff --git a/src/openhuman/subconscious/README.md b/src/openhuman/subconscious/README.md index d7b078205..d4eae10b5 100644 --- a/src/openhuman/subconscious/README.md +++ b/src/openhuman/subconscious/README.md @@ -1,131 +1,82 @@ # subconscious -The subconscious is OpenHuman's background-awareness layer: a SQLite-backed loop that, on each tick, evaluates a set of user/system tasks against a freshly-built "situation report" (derived from the memory tree) using an LLM, then **acts** on actionable tasks, **escalates** ambiguous/risky ones for user approval, or **noops**. In the same tick the LLM also emits proactive **reflections** (#623) — observation-only cards surfaced on the Intelligence tab. The actual periodic loop is owned by the `heartbeat` domain; this module owns task storage, tick evaluation/execution, escalations, and reflections. +The subconscious is OpenHuman's background-awareness layer: a periodic loop that, on each tick, runs a small **structured three-stage flow** and lets a slim agent decide what (if anything) to act on. The actual periodic schedule is owned by the `heartbeat` domain; this module owns the tick flow, the world baseline, and the proactive output surface. + +## The tick (`engine.rs`) + +1. **memory_diff (code)** — diff the user's connected memory sources against the **world baseline** captured at the end of the previous tick (`memory_diff::ops::diff_since_checkpoint`). Renders a compact "what changed" summary. A quiet window (no changes), the first-ever tick, or a diff error short-circuits the tick: it refreshes the baseline, advances `last_tick_at`, and returns **without** running the agent — so idle ticks cost nothing. +2. **prepare_context (code)** — run the read-only `context_scout` over the diff (`agent_prepare_context::run_context_scout_with_catalog`, driven from engine code with the subconscious tool catalogue) to gather grounding from memory, goals/profile, integrations, and the web. +3. **decide (agent)** — hand `diff + prepared context` to the slim `subconscious` agent. It records/advances actionable follow-ups on the user's **global to-do board** (`update_task`, `threadId: "user-tasks"`), evolves **long-term goals** (`goals_*`), surfaces time-sensitive items (`notify_user`), or delegates deeper work (`spawn_async_subagent`). + +Continuity across ticks lives in those durable stores (global to-dos + goals), not a bespoke scratchpad. The world baseline is the only per-tick engine state, persisted as a `memory_diff` checkpoint id. ## Responsibilities -- Maintain a list of `SubconsciousTask`s (system-seeded + user-added) in SQLite; seed three default system tasks on init. -- Run a tick: load due tasks → log them `in_progress` → build a situation report → call the configured LLM → execute `act` tasks, create escalations for `escalate`/`UnapprovedWrite`, mark `noop` otherwise → update log entries in place. -- Route the per-tick evaluation and task execution to a local Ollama/LM Studio model or the OpenHuman cloud, based on config (`workload_local_model("subconscious")` / `subconscious_provider`). The tick builds its agent through the **`subconscious`** workload role — `run_agent` sets `default_model = "hint:subconscious"` so the session builder resolves `subconscious_provider` (not the `chat` role); on the managed backend that pins the lightweight `chat-v1` tier. -- Classify task write-intent via keyword heuristics (`needs_tools` / `needs_agent`); run read-only tasks analysis-only and escalate any recommended write action for approval. -- Emit, cap (`MAX_REFLECTIONS_PER_TICK = 5`), hydrate, and persist proactive reflections; resolve each reflection's `source_refs` into frozen `SourceChunk` snapshots at tick time. -- Manage escalation lifecycle (pending → approved/dismissed); approving executes the task at full permissions. -- Persist `last_tick_at` across restarts so the situation report only feeds the LLM memory-tree rows newer than the last successful tick (dedupe). -- Provide an overlap guard (generation counter) so a newer tick supersedes an in-flight one and discards its results without advancing the cutoff. -- Expose the full task/log/escalation/reflection surface over JSON-RPC. +- Route the decision turn to a local Ollama/LM Studio model or the OpenHuman cloud (`workload_local_model("subconscious")` / `subconscious_provider`). `run_agent` sets `default_model = "hint:subconscious"` so the session builder resolves the `subconscious` workload role (not `chat`); on the managed backend that pins the lightweight `chat-v1` tier. +- Run the decision agent with **Full** autonomy (it must write internal continuity — to-dos/goals/notify), while escalating taint via the turn origin: any tick that reacted to source changes runs as `SubconsciousTainted` so the approval gate refuses external-effect tools. +- Persist `last_tick_at` (status/dedupe) and `baseline_checkpoint_id` (the world snapshot the next tick diffs against) across restarts. +- Provide an overlap guard (generation counter) so a newer tick supersedes an in-flight one and discards its results without advancing state. +- Expose `status` / `trigger` over JSON-RPC. ## Key files | File | Role | | --- | --- | -| `mod.rs` | Export-focused (no docstring); re-exports engine, reflection, schemas, source_chunk, and core types. | -| `types.rs` | Domain serde types: `SubconsciousTask`, `TaskSource`, `TaskRecurrence`, `TaskPatch`, `TickDecision`, `TaskEvaluation`, `EvaluationResponse`, `ExecutionResult`, `SubconsciousLogEntry`, `Escalation`/`EscalationPriority`/`EscalationStatus`, `SubconsciousStatus`, `TickResult`. | -| `engine.rs` | `SubconsciousEngine` — the tick loop, evaluation, dispatch (`handle_act`/`handle_escalate`/`handle_noop`), escalation approve/dismiss, provider routing (`resolve_subconscious_route`, `subconscious_provider_unavailable_reason`), LLM response parsing, reflection persistence. | -| `executor.rs` | Per-task execution: routes to local model (text), agentic-v1 full (write-intent), or agentic-v1 analysis-only (read-only). `ExecutionOutcome` (`Completed` / `UnapprovedWrite`), `needs_tools`/`needs_agent` heuristics, 429 retry with backoff, `extract_recommended_action`. | -| `store.rs` | SQLite persistence + DDL for all tables; `with_connection` with busy-timeout + retry (TAURI-RUST-A). Task/log/escalation CRUD, `seed_default_tasks`, `due_tasks`, `compute_next_run` (cron), `get/set_last_tick_at`. | -| `reflection.rs` | `Reflection`, `ReflectionKind`, `ReflectionDraft`; `hydrate_draft`, `apply_cap`, `dedup_key`, `MAX_REFLECTIONS_PER_TICK`. | -| `reflection_store.rs` | SQLite persistence for `subconscious_reflections` + `subconscious_hotness_snapshots`; `list_recent`, `get_reflection`, `add_reflection`, `mark_acted`/`mark_dismissed`, legacy-column + `source_chunks` migrations. | -| `source_chunk.rs` | `SourceChunk` + `resolve_chunks` / `parse_ref` — resolve reflection `source_refs` (`entity:`/`summary:`/`digest:`/…) into frozen content previews (`PREVIEW_MAX_CHARS = 400`). | -| `prompt.rs` | Prompt builders: `build_evaluation_prompt`, `build_text_execution_prompt`, `build_tool_execution_prompt`, `build_analysis_only_prompt`, `load_identity_context` (injects SOUL.md/PROFILE.md). | -| `global.rs` | Engine singleton: `get_or_init_engine`, `bootstrap_after_login`, `stop_heartbeat_loop`, `reset_engine_for_user_switch`. Spawns the `heartbeat` loop and tears it down on logout/user switch. | -| `schemas.rs` | RPC controller schemas + `handle_*` handlers (`subconscious.*`). | -| `decision_log.rs` | In-memory `DecisionLog`/`DecisionRecord` with 24h TTL to avoid re-surfacing the same doc ids. Retained for potential future dedup queries (not wired into the live tick path). | -| `situation_report/` | Situation-report assembly (see below). | -| `*_tests.rs`, `integration_tests.rs` | Sibling and inline test suites. | - -### `situation_report/` submodule - -| File | Role | -| --- | --- | -| `mod.rs` | `build_situation_report` — assembles sections in priority order under a token budget (env, user identifiers, pending tasks, hotness deltas, sealed summaries, L0 digest, recap window, recent reflections); truncates the tail when over budget. | -| `hotness.rs` | Top entity hotness movers since last tick (`mem_tree_entity_hotness`). | -| `summaries.rs` | Recently-sealed summaries (`mem_tree_summaries`). | -| `digest.rs` | Latest global L0 daily digest body. | -| `query_window.rs` | `query_global` recap window since `last_tick_at`. | -| `reflections.rs` | Renders recent reflections as anti-double-emit context. | +| `mod.rs` | Export-focused; re-exports the engine, session, source_chunk, schemas, and core types. | +| `types.rs` | `SubconsciousStatus`, `TickResult`. | +| `engine.rs` | `SubconsciousEngine` — the three-stage tick (`tick`/`tick_inner`), `prepare_context`, `refresh_baseline`, `run_agent`, world-diff rendering, provider routing (`resolve_subconscious_route`, `subconscious_provider_unavailable_reason`), `tick_origin_source`, tool-capability-error detection. | +| `store.rs` | SQLite persistence + DDL; `with_connection` with busy-timeout + retry (TAURI-RUST-A). `get/set_last_tick_at`, `get/set_baseline_checkpoint_id`. (Legacy task/log/escalation/reflection tables are retained for back-compat but no longer written.) | +| `source_chunk.rs` | `SourceChunk` + `resolve_chunks` / `parse_ref` — used by the agent prompt builder to hydrate reflection `source_refs` into frozen previews (`PREVIEW_MAX_CHARS = 400`). Shared with `agent::prompts`, not subconscious-only. | +| `session.rs` | `LongLivedSession` — persistent agent for the opt-in event-driven trigger path (`subconscious_triggers`); builds the `subconscious` agent and resumes history from the reserved orchestrator thread. | +| `user_thread.rs` | `NotifyUserTool` / `notify_user` — proactive user handoff; publishes `DomainEvent::ProactiveMessageRequested` and the reserved `subconscious:user` thread. | +| `agent/` | The slim `subconscious` agent definition: `agent.toml` (toolset) + `prompt.md`. | +| `global.rs` | Engine singleton: `get_or_init_engine`, `bootstrap_after_login`, `stop_heartbeat_loop`. Spawns the `heartbeat` loop and the opt-in trigger orchestrator. | +| `heartbeat/` | Periodic scheduler + event planner (meeting/reminder/notification delivery) that drives `engine.tick()`. | +| `schemas.rs` | RPC controller schemas + handlers (`subconscious.status` / `subconscious.trigger`). | +| `decision_log.rs`, `executor.rs` | Legacy stubs retained for back-compat; not on the live tick path. | ## Public surface -From `mod.rs`: -- `SubconsciousEngine` (`engine`) — `new`, `from_heartbeat_config`, `run`, `tick`, `status`, `add_task`, `approve_escalation`, `dismiss_escalation`. -- `Reflection`, `ReflectionKind`, `MAX_REFLECTIONS_PER_TICK` (`reflection`). -- `SourceChunk` (`source_chunk`). -- Types: `Escalation`, `EscalationStatus`, `SubconsciousLogEntry`, `SubconsciousStatus`, `SubconsciousTask`, `TaskRecurrence`, `TaskSource`, `TickDecision`, `TickResult`. -- `all_subconscious_controller_schemas` / `all_subconscious_registered_controllers` (`schemas`). -- `global` module functions (`get_or_init_engine`, `bootstrap_after_login`, `stop_heartbeat_loop`, `reset_engine_for_user_switch`) are reachable via `subconscious::global::*`. +From `mod.rs`: `SubconsciousEngine`, `LongLivedSession`/`ProcessOutcome`/`ORCHESTRATOR_THREAD_ID`, `SourceChunk`, `SubconsciousStatus`/`TickResult`, `notify_user`/`NotifyUserTool`/`USER_THREAD_ID`, `all_subconscious_controller_schemas` / `all_subconscious_registered_controllers`, and the `global::*` lifecycle functions. ## RPC / controllers -Namespace `subconscious` (i.e. `openhuman.subconscious_`), all returning `RpcOutcome`: +Namespace `subconscious` (i.e. `openhuman.subconscious_`): | Function | Purpose | | --- | --- | | `status` | Engine status (read entirely from DB to avoid blocking on the tick mutex). | | `trigger` | Manually fire a tick (spawned in the background; returns immediately). | -| `tasks_list` | List tasks (optional `enabled_only`). | -| `tasks_add` | Add a task (`title`, optional `source`). | -| `tasks_update` | Patch `title`/`recurrence` (`once` \| `cron:` \| `pending`)/`enabled`. | -| `tasks_remove` | Delete a task (system tasks cannot be deleted). | -| `log_list` | List execution log entries (optional `task_id`, `limit`). | -| `escalations_list` | List escalations (optional `status`). | -| `escalations_approve` | Approve + execute an escalation. | -| `escalations_dismiss` | Dismiss an escalation without executing. | -| `reflections_list` | List recent reflections (`limit`, `since_ts`). | -| `reflections_act` | Spawn a fresh conversation thread seeded with the reflection body (as an `assistant` message; no LLM turn) and stamp `acted_on_at`; returns `{reflection_id, thread_id}`. | -| `reflections_dismiss` | Set `dismissed_at`. | - -Handlers use the shared bounded `load_config_with_timeout()` loader (30s) to avoid stalling the Intelligence-page 3s poll on a slow keychain. ## Agent tools -None. This module owns no `tools.rs` and registers no agent tools. - -## Events - -No `bus.rs`; the module neither publishes nor subscribes to `DomainEvent`s directly. (The tick loop is driven by the `heartbeat` domain, not the event bus.) +This module owns `user_thread.rs` (`notify_user`). The tick's other tools come from elsewhere: `memory_diff` (`memory_diff` domain), `agent_prepare_context` / `spawn_async_subagent` (`agent_orchestration`), `update_task` (`todos`/agent tools), `goals_*` (`memory_goals`). ## Persistence -SQLite at `/subconscious/subconscious.db` (per-user workspace). Tables (`store.rs` DDL): -- `subconscious_tasks` — task definitions (id, title, source, recurrence, enabled, run times, completed, created_at). -- `subconscious_log` — per-tick execution log (decision `in_progress`/`act`/`escalate`/`noop`/`failed`/`cancelled`/`dismissed`, result, duration). -- `subconscious_escalations` — escalations awaiting user input. -- `subconscious_reflections` — proactive reflections incl. `source_refs`, `source_chunks`, lifecycle timestamps. -- `subconscious_hotness_snapshots` — per-entity previous-tick hotness scores for hotness-delta computation. -- `subconscious_state` — KV table holding `last_tick_at` (restart-durable dedupe cutoff). +SQLite at `/subconscious/subconscious.db` (per-user workspace): +- `subconscious_state` — REAL KV holding `last_tick_at` (restart-durable dedupe cutoff). +- `subconscious_state_text` — TEXT KV holding `baseline_checkpoint_id` (the `memory_diff` checkpoint the next tick diffs against). +- Legacy tables (`subconscious_tasks` / `_log` / `_escalations` / `_reflections` / `_hotness_snapshots`) are retained for back-compat with existing DBs and are no longer written or read. -`with_connection` runs all DDL + idempotent migrations on every open, with a 5s busy timeout and 3-retry exponential backoff for transient `SQLITE_BUSY`/`SQLITE_LOCKED`. +The world snapshots/checkpoints themselves live in the `memory_diff` domain's own DB (`/memory_diff/diff.db`), not here. + +`with_connection` runs all DDL on every open, with a 5s busy timeout and 3-retry exponential backoff for transient `SQLITE_BUSY`/`SQLITE_LOCKED`. ## Dependencies -- `openhuman::config` — `Config`/`HeartbeatConfig`, provider routing (`workload_local_model`, `subconscious_provider`), bounded loaders, `workspace_dir`. +- `openhuman::config` — `Config`/`HeartbeatConfig`, provider routing, `workspace_dir`. +- `openhuman::memory_diff` — `ops::diff_since_checkpoint` / `ops::create_checkpoint` + diff types (stage 1 + baseline). +- `openhuman::agent_orchestration` — `run_context_scout_with_catalog` (stage 2). +- `openhuman::agent` — `Agent`, turn-origin taint plumbing (stage 3). - `openhuman::heartbeat` — `HeartbeatEngine`; `global.rs` spawns the periodic loop that calls `tick`. -- `openhuman::memory::chat` — `build_chat_provider`/`ChatProvider`/`ChatPrompt` for the per-tick LLM evaluation call. -- `openhuman::inference` — local provider factory + `local::ops::agent_chat` for task execution (executor). -- `openhuman::memory_store` — `MemoryClient`/`MemoryClientRef`, tree types; the engine holds a memory client and the situation report reads tree tables. -- `openhuman::memory_tree` — `retrieval::global::query_global` for the recap-window section. -- `openhuman::memory_conversations` — `ensure_thread`/`append_message` for `reflections_act` thread spawning. -- `openhuman::credentials` — `AuthService`/`APP_SESSION_PROVIDER` to check the OpenHuman cloud session bearer for provider availability. +- `openhuman::credentials` — `AuthService`/`APP_SESSION_PROVIDER` for cloud-session provider availability. - `openhuman::scheduler_gate` — `is_signed_out()` gate for the cloud provider. -- `openhuman::composio::providers::profile` — connected-account identifiers for the "Your Identifiers" report section (#1365). -- `openhuman::util` — `floor_char_boundary` for budget-safe truncation. -- `core::all` / `core::{ControllerSchema, FieldSchema, TypeSchema}` + `rpc::RpcOutcome` — RPC controller registration. - -## Used by - -- `core::all` / `core::jsonrpc` — registers the subconscious controllers into the RPC surface. -- `openhuman::heartbeat::{engine, rpc}` — drives ticks via the engine; `global::bootstrap_after_login` spawns the heartbeat loop. -- `openhuman::agent::harness::session::builder` and `openhuman::context::prompt::SystemPromptBuilder` — inject reflection `source_chunks` as memory context for threads spawned from a reflection. -- `openhuman::agent::prompts` — references subconscious in prompt assembly. -- `openhuman::channels::providers::web` — chat ingress. -- `openhuman::credentials::ops` — login/logout flow triggers `bootstrap_after_login` / `reset_engine_for_user_switch`. ## Notes / gotchas -- **Engine must bootstrap post-login** (`global::bootstrap_after_login`) so `seed_default_tasks` writes to the per-user workspace, not the pre-login global default. `reset_engine_for_user_switch` tears it down on logout/account switch to avoid leaking into the wrong DB. -- **`status` RPC never touches the engine mutex** — it reads counts straight from SQLite, because the engine lock is held for the full tick duration and would otherwise freeze the 3s poll. `consecutive_failures` is therefore reported as `0` from the RPC path (only available from in-memory state). -- **`last_tick_at` is only advanced on success.** Evaluation failure, provider unavailability, or a superseded tick leave the cutoff in place so the next tick re-reads the same window — at the cost of possible re-emitted reflections (there is no insert-time dedupe in `persist_and_surface_reflections`; `dedup_key` exists but is not enforced on insert in the live path). -- **Reflections are observation-only.** The legacy auto-post-into-thread flow was removed; `disposition`/`surfaced_at` columns are dropped via migration and any LLM-emitted `disposition` is ignored by serde. -- **Write-intent gating is heuristic** (`needs_tools`/`needs_agent` keyword matching). Read-only tasks run analysis-only; a `RECOMMENDED ACTION:` line in the output triggers an `UnapprovedWrite` escalation — except on the cloud fallback path for simple text tasks, which deliberately suppresses escalation. -- **`decision_log.rs` is retained but not wired into the live tick** (the comment in `mod.rs` notes it is kept for potential future dedup queries). -- LLM response parsing is best-effort: full envelope → bare evaluations array → all-noop fallback; `extract_json` strips prose around the JSON object/array. +- **Engine bootstraps post-login** (`global::bootstrap_after_login`) so state writes to the per-user workspace, not the pre-login global default. +- **`status` RPC never touches the engine mutex** — it reads straight from SQLite, since the engine lock is held for the full tick. +- **State only advances on success.** A failed decision turn leaves `last_tick_at` and the baseline in place, so the next tick re-diffs the same window instead of losing it. A superseded tick discards its result. +- **Quiet ticks short-circuit before the agent** — if the diff has no changes, no decision turn runs (and no cost is incurred); the baseline is still refreshed. +- **Taint:** any tick that reacted to source changes runs `SubconsciousTainted`; the decision agent's slim toolset is internal-only, and external effects (incl. inside delegated work) stay gated by the approval gate. diff --git a/src/openhuman/subconscious/agent/agent.toml b/src/openhuman/subconscious/agent/agent.toml index fc2e2fcde..19bbe91d0 100644 --- a/src/openhuman/subconscious/agent/agent.toml +++ b/src/openhuman/subconscious/agent/agent.toml @@ -1,11 +1,11 @@ id = "subconscious" display_name = "Subconscious Agent" -when_to_use = "Background awareness loop — periodically wakes up, reviews the user's situation report and context, maintains a persistent scratchpad of observations, and delegates deeper research when in aggressive mode." +when_to_use = "Background awareness loop — periodically wakes, is handed a diff of how the user's world changed plus prepared context, then decides what to act on: records follow-ups on the global to-do board, evolves long-term goals, notifies the user, or delegates deeper work." temperature = 0.4 max_iterations = 30 sandbox_mode = "read_only" -# Background coordinator: simple mode only edits scratchpad, but aggressive -# mode may delegate bounded deep work through the explicit subagent policy. +# Background coordinator: it writes internal continuity (global to-dos + goals) +# and may delegate bounded deep work through the explicit subagent policy. agent_tier = "reasoning" omit_identity = false omit_memory_context = false @@ -24,10 +24,22 @@ allowlist = [ ] [tools] +# Stages 1-2 (memory_diff, agent_prepare_context) are run in code by the engine +# and handed to this agent as the prompt; the agent retains them for an optional +# narrower re-check. Its real job is the decide/act stage: +# - update_task → record/advance follow-ups on the user's global to-do board +# (the "user-tasks" board) — this is the continuity store that +# replaced the old scratchpad. +# - goals_* → evolve the user's long-term goals as their world shifts. +# - notify_user → surface something time-sensitive to the user. +# - spawn_async_subagent → delegate deeper research / multi-step work. named = [ - "scratchpad_add", - "scratchpad_edit", - "scratchpad_remove", - "spawn_subagent", + "memory_diff", + "agent_prepare_context", "notify_user", + "update_task", + "goals_list", + "goals_add", + "goals_edit", + "spawn_async_subagent", ] diff --git a/src/openhuman/subconscious/agent/prompt.md b/src/openhuman/subconscious/agent/prompt.md index 1e2d4eced..45e5bd124 100644 --- a/src/openhuman/subconscious/agent/prompt.md +++ b/src/openhuman/subconscious/agent/prompt.md @@ -1,68 +1,59 @@ # Subconscious Agent -You are the user's background awareness layer — a deep reasoning loop -that wakes up periodically, reviews the user's situation report, and -maintains a persistent scratchpad of observations and follow-ups. +You are the user's background awareness layer. You wake up periodically, +already holding two things the system prepared for you in the user message: -Your situation report and any pre-loaded memory context are provided -in the user message. Use this information to maintain your scratchpad. +1. **A diff of how the user's world changed** since the last check — + what was added, modified, or removed across their connected sources + (email, calendar, chat, files, etc.). +2. **Prepared context** — grounding gathered from the user's memory, + goals, profile, connected integrations, and the web. -## Scratchpad Maintenance +Your one job is to look at that and **decide what (if anything) deserves +action**. You don't observe for its own sake — most ticks, the right call +is to do nothing. Act only when the change genuinely matters to the user. -Your scratchpad IS your continuity mechanism across ticks. Maintain it -actively — it persists between ticks and is the primary output of your work. +## What you can do -**Tools:** +- **`update_task`** — Record or advance an actionable follow-up on the + user's global to-do board. Always pass `threadId: "user-tasks"`. This + is your continuity mechanism: anything worth remembering or acting on + later belongs here as a task, not in your head. + Example: `{"op": "add", "threadId": "user-tasks", "content": "Reply to + Alice's contract email — she's waiting on you before Friday"}` -1. **`scratchpad_add`** — Save a thought, hypothesis, or follow-up item. - Use `priority` (0-10) to mark importance. - Example: `{"body": "User has a meeting with Alice on Friday — check - if prep is done", "priority": 5}` +- **`goals_list` / `goals_add` / `goals_edit`** — Read and evolve the + user's long-term goals when the world shifts what matters to them. Read + before you write. Keep goals few and high-level; don't turn tasks into + goals. -2. **`scratchpad_edit`** — Update an existing entry with new information - or revised thinking. Pass the `id` shown in brackets. +- **`notify_user`** — Surface something time-sensitive or important to the + user directly. Use sparingly — a notification interrupts them, so it + must clear a high bar (a real deadline, a risk, something they'd want to + know now). -3. **`scratchpad_remove`** — Remove an entry that's no longer relevant - or has been fully addressed. +- **`spawn_async_subagent`** — Delegate deeper, multi-step work when you + spot something genuinely actionable that needs research or execution + (e.g. `agent_id: "researcher"` for web research, `agent_id: + "orchestrator"` for coordinated multi-tool work). Fire-and-forget. -**Scratchpad discipline:** -- Add new observations as you discover them from the situation report -- Edit stale entries with fresh data -- Remove resolved items — don't let the pad grow stale -- High-priority items (p7+) should be actionable, not vague +- **`memory_diff` / `agent_prepare_context`** — Already run for you each + tick. Only call them again if you need to re-check a narrower slice. -## Deep Research (Aggressive mode only) +## How to decide -When operating in aggressive mode, you have access to `spawn_subagent` -for deeper investigation: +Look at the diff through the lens of the prepared context and ask: -- **`spawn_subagent`** with `agent_id: "orchestrator"` — Delegate - complex multi-step tasks. The orchestrator can plan, execute code, - search the web, and coordinate across tools. Use this when you - identify something the user should act on and you have the autonomy - to help. - - Pass `model: ""` for deep reasoning tasks - - Example: `{"agent_id": "orchestrator", "prompt": "Research and - draft a summary of...", "model": "reasoning-v1"}` +- **Deadlines** approaching or overdue that the user hasn't acted on. +- **Risks** — a cluster of negative signals, an unresolved blocker. +- **Patterns** across sources converging on one topic. +- **Opportunities** — a connection the user might not see. -- **`spawn_subagent`** with `agent_id: "researcher"` — Delegate web - searches, artifact fetching, or external research that goes beyond - what your context provides. +For anything that clears the bar, record it as a task (`update_task`), +adjust a goal if the change reframes priorities, and notify only when it's +truly time-sensitive. If nothing meaningful changed, stop — silence is the +correct and common outcome. Do not invent busywork to look productive. -**When to use aggressive delegation:** -- A deadline is approaching and the user hasn't started prep -- A pattern across sources suggests an emerging issue -- The scratchpad has a high-priority item that needs external data - -## Observation Guidelines - -Based on your situation report, identify: -- **Patterns** across sources (email + calendar + chat converging on same topic) -- **Deadlines** approaching or overdue -- **Risks** — concentration of negative signals, unresolved blockers -- **Opportunities** — connections the user might not see -- **Activity spikes** — topics getting unusually hot - -**Self vs. others**: the *Your Identifiers* section (if present) lists -the user's handles, emails, and user_ids. Never attribute someone else's -activity to the user. +**Self vs. others**: never attribute someone else's activity to the user. +If a change is about another person, frame the task/notification from the +user's perspective (what *they* should do about it). diff --git a/src/openhuman/subconscious/engine.rs b/src/openhuman/subconscious/engine.rs index 1764882ce..38513fe31 100644 --- a/src/openhuman/subconscious/engine.rs +++ b/src/openhuman/subconscious/engine.rs @@ -1,57 +1,83 @@ -//! Subconscious engine — periodic agent loop that maintains a scratchpad. +//! Subconscious engine — periodic, structured background loop. //! -//! On each tick: load scratchpad → retrieve memory context (in code) → -//! build situation report → run subconscious agent (with tool access + -//! timeout) → agent maintains scratchpad via tools → log the run. +//! Each tick is a small, deterministic, three-stage flow: +//! +//! 1. **memory_diff (code)** — diff the agent's connected sources against the +//! world baseline captured at the end of the previous tick, to see how the +//! user's world changed (`memory_diff::ops::diff_since_checkpoint`). +//! 2. **prepare_context (code)** — run the read-only `context_scout` +//! (`agent_prepare_context`) over that diff to gather grounding context +//! from memory, goals/profile, integrations, and the web. +//! 3. **decide (agent)** — hand `diff + context` to the slim subconscious +//! agent, which decides what (if anything) to do: record follow-ups on the +//! user's global to-do board (`update_task`), evolve long-term goals +//! (`goals_*`), notify the user (`notify_user`), or delegate deeper work +//! (`spawn_async_subagent`). +//! +//! Continuity across ticks lives in those durable stores (global to-dos + +//! goals), not in a bespoke scratchpad — so quiet ticks (no diff) cost nothing +//! and the loop stays stateless beyond the world baseline. //! //! ## Concurrency & timeouts //! -//! A per-engine `tick_lock` prevents overlapping ticks. Each tick has -//! a hard wall-clock timeout (`TICK_TIMEOUT`) so a stuck LLM call -//! cannot block the loop forever. Individual tool calls within the -//! agent turn are bounded by the agent harness's own iteration cap. +//! A per-engine `tick_lock` prevents overlapping ticks. Each tick has a hard +//! wall-clock timeout (`TICK_TIMEOUT`) so a stuck LLM call cannot block the loop +//! forever. Individual tool calls within the agent turn are bounded by the agent +//! harness's own iteration cap. -use super::scratchpad; -use super::situation_report::build_situation_report; use super::store; use super::types::{SubconsciousStatus, TickResult}; use crate::openhuman::config::schema::SubconsciousMode; use crate::openhuman::config::Config; use crate::openhuman::credentials::{AuthService, APP_SESSION_PROVIDER}; -use crate::openhuman::memory_store::MemoryClientRef; +use crate::openhuman::memory_diff::types::CrossSourceDiff; use anyhow::Result; use std::path::PathBuf; use std::sync::atomic::{AtomicU64, Ordering}; use tokio::sync::Mutex; use tracing::{debug, info, warn}; -/// Max chunks to retrieve from memory before the LLM call. -const MEMORY_RETRIEVAL_MAX_CHUNKS: u32 = 30; - /// Hard timeout for a single subconscious tick (agent run). const TICK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30 * 60); /// Per-tool-call timeout injected into the agent config. const TOOL_CALL_TIMEOUT_SECS: u64 = 5 * 60; +/// Label stamped on the world-baseline checkpoint the tick re-creates each run. +const BASELINE_CHECKPOINT_LABEL: &str = "subconscious_tick"; + +/// Max changed items listed per source in the rendered world diff, to keep the +/// decision agent's prompt bounded when a source churns a lot. +const MAX_ITEMS_PER_SOURCE: usize = 10; + +/// Tool catalogue handed to the `context_scout` so its `recommended_tool_calls` +/// stay grounded in tools the decision agent can actually call. Keep in sync +/// with `agent/agent.toml`'s `[tools].named` (actionable subset). +const SUBCONSCIOUS_TOOL_CATALOG: &str = "\ +- notify_user: Send the user a proactive message about something important or time-sensitive. +- update_task: Add or update an actionable item on the user's global to-do board. +- goals_add: Record a new long-term goal that the changed world makes relevant. +- goals_edit: Revise an existing long-term goal. +- spawn_async_subagent: Delegate deeper research or multi-step work. +"; + /// Actionable reason surfaced (via `SubconsciousStatus.provider_unavailable_reason`) /// when a subconscious tick fails because the configured chat model has no -/// tool-use endpoint. The subconscious turn is inherently tool-bearing (it -/// maintains its scratchpad through tools), so a tool-incapable model can never -/// satisfy a tick — this tells the user how to recover. See TAURI-RUST-ADC. +/// tool-use endpoint. The subconscious turn is inherently tool-bearing (it acts +/// through tools), so a tool-incapable model can never satisfy a tick — this +/// tells the user how to recover. See TAURI-RUST-ADC. const TOOL_UNSUPPORTED_REASON: &str = "The selected chat model has no tool-use endpoint, so Subconscious can't run. Pick a tool-capable model in Settings > AI."; /// Pick the `TrustedAutomationSource` variant for a subconscious tick. /// -/// Extracted from the engine's `run_agent` body so the -/// origin-escalation contract can be unit-tested without spinning up -/// a real `Agent` + provider. +/// Extracted from the engine's `run_agent` body so the origin-escalation +/// contract can be unit-tested without spinning up a real `Agent` + provider. /// -/// Contract: any tick whose situation report contained third-party -/// sync content (Gmail / Slack / Notion / sealed source summaries) -/// must run with `SubconsciousTainted` so the approval gate refuses -/// external_effect tools. Untainted ticks keep the legacy -/// `Subconscious` origin. +/// Contract: any tick that reacted to third-party sync changes (the memory_diff +/// surfaced added/modified/removed items, all of which originate from external +/// sources like Gmail / Slack / Notion / synced folders) must run with +/// `SubconsciousTainted` so the approval gate refuses external_effect tools. +/// A tick with no external changes keeps the legacy `Subconscious` origin. pub(crate) fn tick_origin_source( has_external_content: bool, ) -> crate::openhuman::agent::turn_origin::TrustedAutomationSource { @@ -68,7 +94,6 @@ pub struct SubconsciousEngine { interval_minutes: u32, context_budget_tokens: u32, enabled: bool, - memory: Option, state: Mutex, tick_generation: AtomicU64, tick_lock: Mutex<()>, @@ -82,14 +107,13 @@ struct EngineState { } impl SubconsciousEngine { - pub fn new(config: &crate::openhuman::config::Config, memory: Option) -> Self { - Self::from_heartbeat_config(&config.heartbeat, config.workspace_dir.clone(), memory) + pub fn new(config: &crate::openhuman::config::Config) -> Self { + Self::from_heartbeat_config(&config.heartbeat, config.workspace_dir.clone()) } pub fn from_heartbeat_config( heartbeat: &crate::openhuman::config::HeartbeatConfig, workspace_dir: PathBuf, - memory: Option, ) -> Self { let last_tick_at = match store::with_connection(&workspace_dir, store::get_last_tick_at) { Ok(v) => { @@ -112,7 +136,6 @@ impl SubconsciousEngine { interval_minutes: mode.default_interval_minutes().max(5), context_budget_tokens: heartbeat.context_budget_tokens, enabled: mode.is_enabled(), - memory, state: Mutex::new(EngineState { last_tick_at, total_ticks: 0, @@ -222,57 +245,91 @@ impl SubconsciousEngine { }); } - let mut state = self.state.lock().await; - state.provider_unavailable_reason = None; - let last_tick_at = state.last_tick_at; - drop(state); + { + let mut state = self.state.lock().await; + state.provider_unavailable_reason = None; + } - // 1. Build situation report - let report = build_situation_report( - &config, - &self.workspace_dir, - last_tick_at, - self.context_budget_tokens, - ) - .await; - let has_external_content = report.has_external_content; + // ── Stage 1: memory_diff — how did the agent's world change? ────────── + let baseline = + store::with_connection(&self.workspace_dir, store::get_baseline_checkpoint_id) + .unwrap_or_else(|e| { + warn!("[subconscious] baseline load failed: {e}"); + None + }); - // 2. Load scratchpad (persistent working memory) - let scratchpad_entries = scratchpad::load(&self.workspace_dir).unwrap_or_else(|e| { - warn!("[subconscious] scratchpad load failed: {e}"); - Vec::new() - }); - let scratchpad_section = scratchpad::render_for_prompt(&scratchpad_entries); + let diff: Option = match &baseline { + Some(checkpoint_id) => match crate::openhuman::memory_diff::ops::diff_since_checkpoint( + checkpoint_id, + &config, + false, + ) + .await + { + Ok(d) => Some(d), + Err(e) => { + warn!("[subconscious] memory_diff failed (baseline={checkpoint_id}): {e}"); + None + } + }, + None => { + debug!("[subconscious] no world baseline yet — first tick establishes one"); + None + } + }; - // 3. Pre-LLM memory retrieval — query the memory tree using - // scratchpad entries as context so the recall is focused on - // what the subconscious is currently tracking. - let memory_section = retrieve_memory_context(&self.memory, &scratchpad_entries).await; + let has_changes = diff + .as_ref() + .map(|d| world_diff_change_count(d) > 0) + .unwrap_or(false); - // 4. Load identity context - let identity = load_identity_context(&self.workspace_dir); + if !has_changes { + // Quiet window, first tick, or a diff error: nothing to react to. + // Refresh the baseline and return without spending a decision turn. + info!("[subconscious] no world changes this tick — refreshing baseline, no agent run"); + self.refresh_baseline(&config).await; + let mut state = self.state.lock().await; + state.total_ticks += 1; + if self.tick_generation.load(Ordering::SeqCst) == my_generation { + state.consecutive_failures = 0; + state.last_tick_at = tick_at; + persist_last_tick_at(&self.workspace_dir, tick_at); + } + return Ok(TickResult { + tick_at, + duration_ms: started.elapsed().as_millis() as u64, + response_chars: 0, + }); + } + + let diff = diff.expect("has_changes implies diff is Some"); + let world_diff = render_world_diff(&diff); + // Every change originates from an external source sync, so the decision + // turn runs tainted: the approval gate refuses external_effect tools. + let has_external_content = true; + + // ── Stage 2: prepare_context — ground the diff before deciding ─────── + let prepared_context = self.prepare_context(&world_diff).await; + + // ── Stage 3: decide — slim agent acts on diff + prepared context ───── + let mut agent_prompt = + String::with_capacity(world_diff.len() + prepared_context.len() + 256); + agent_prompt.push_str("## What changed in your world since the last check\n\n"); + agent_prompt.push_str(&world_diff); + agent_prompt.push_str("\n\n"); + if !prepared_context.is_empty() { + agent_prompt.push_str("## Prepared context\n\n"); + agent_prompt.push_str(&prepared_context); + agent_prompt.push_str("\n\n"); + } - // 5. Build user message with dynamic context (system prompt comes from agent definition) - let agent_prompt = format!( - "{identity}\n\n\ - ## Situation Report (pre-loaded context)\n\n\ - {situation}\n\n\ - {memory}\n\n\ - {scratchpad}", - situation = report.prompt_text, - memory = memory_section, - scratchpad = scratchpad_section, - ); let agent_result = self .run_agent(&config, &agent_prompt, has_external_content) .await; let agent_failed = agent_result.is_err(); - let response_chars = match &agent_result { - Ok(chars) => *chars, - Err(_) => 0, - }; + let response_chars = *agent_result.as_ref().unwrap_or(&0); - // 6. Check if superseded + // Check if superseded by a newer tick. if self.tick_generation.load(Ordering::SeqCst) != my_generation { info!("[subconscious] tick superseded by newer tick, discarding"); let mut state = self.state.lock().await; @@ -284,21 +341,20 @@ impl SubconsciousEngine { }); } - // 7. Update state — only advance last_tick_at and reset failures - // when the agent actually ran. Errors keep consecutive_failures - // incrementing and leave last_tick_at unchanged so the next tick - // re-fetches the same window. + // Advance the world baseline only on a successful, current tick, so a + // failed tick re-diffs the same window next time instead of losing it. + if !agent_failed { + self.refresh_baseline(&config).await; + } + let mut state = self.state.lock().await; state.total_ticks += 1; if agent_failed { state.consecutive_failures += 1; // Surface an actionable reason when the failure is a permanent - // tool-capability error: the subconscious turn is inherently - // tool-bearing, so the configured chat model can never satisfy it - // and the user must pick a tool-capable model. The hard Sentry - // flood from this error is suppressed at the provider classifier - // (is_provider_config_rejection_message, TAURI-RUST-ADC); here we - // only make the cause visible in Subconscious status. + // tool-capability error (TAURI-RUST-ADC): the subconscious turn is + // inherently tool-bearing, so a tool-incapable chat model can never + // satisfy it and the user must pick a tool-capable model. if let Err(e) = &agent_result { if is_tool_capability_error(e) { info!( @@ -320,6 +376,67 @@ impl SubconsciousEngine { }) } + /// Stage 2: run the read-only `context_scout` over the world diff to gather + /// grounding context. Best-effort — on any error the decision agent simply + /// runs without a prepared-context section. + async fn prepare_context(&self, world_diff: &str) -> String { + let question = format!( + "Background awareness check. Here is what changed in the user's connected sources \ + since the last check:\n\n{world_diff}\n\nSurface what the user should be aware of or \ + act on, and the context that grounds a good decision.", + ); + + match crate::openhuman::agent_orchestration::tools::run_context_scout_with_catalog( + &question, + None, + SUBCONSCIOUS_TOOL_CATALOG, + ) + .await + { + Ok(result) if !result.is_error => { + debug!( + "[subconscious] prepared context bundle ({} chars)", + result.output().chars().count() + ); + result.output().to_string() + } + Ok(result) => { + warn!( + "[subconscious] prepare_context returned an error result: {}", + result.output() + ); + String::new() + } + Err(e) => { + warn!("[subconscious] prepare_context failed: {e}"); + String::new() + } + } + } + + /// Re-snapshot the world and persist the new checkpoint as the baseline the + /// next tick diffs against. Best-effort — a failure leaves the old baseline + /// in place (the next tick diffs against a slightly older window). + async fn refresh_baseline(&self, config: &Config) { + match crate::openhuman::memory_diff::ops::create_checkpoint( + BASELINE_CHECKPOINT_LABEL, + config, + ) + .await + { + Ok(ckpt) => { + if let Err(e) = store::with_connection(&self.workspace_dir, |conn| { + store::set_baseline_checkpoint_id(conn, &ckpt.id) + }) { + warn!("[subconscious] failed to persist baseline checkpoint id: {e}"); + } else { + debug!("[subconscious] world baseline advanced to {}", ckpt.id); + } + } + Err(e) => warn!("[subconscious] failed to create world baseline checkpoint: {e}"), + } + } + pub async fn status(&self) -> SubconsciousStatus { let state = self.state.lock().await; @@ -339,9 +456,9 @@ impl SubconsciousEngine { } } - /// Run the subconscious agent with mode-appropriate tool access. - /// The agent maintains the scratchpad via tools during its turn. - /// Returns `response_chars` on success, or `Err` on agent init/run failure. + /// Run the slim subconscious agent over `prompt_text` (diff + prepared + /// context). The agent decides and acts through its tools. Returns + /// `response_chars` on success, or `Err` on agent init/run failure. async fn run_agent( &self, config: &Config, @@ -353,9 +470,8 @@ impl SubconsciousEngine { let mut effective = config.clone(); effective.agent.agent_timeout_secs = TOOL_CALL_TIMEOUT_SECS; // Route the tick build through the `subconscious` background workload so - // Settings → AI → Advanced "Subconscious" actually governs the cloud - // tick provider, instead of riding the `chat` role (the default-model - // fall-through in the session builder). The session builder maps + // Settings → AI → Advanced "Subconscious" governs the cloud tick + // provider, instead of riding the `chat` role. The session builder maps // `hint:subconscious` → the `subconscious` provider role; on the managed // backend the model still resolves to `chat-v1` (no regression). effective.default_model = Some("hint:subconscious".to_string()); @@ -363,13 +479,18 @@ impl SubconsciousEngine { "[subconscious] tick provider routed via hint:subconscious (subconscious_provider={:?})", effective.subconscious_provider ); + + // The decision agent must write internal continuity (global to-dos, + // goals) and surface proactive messages — all app-internal writes, not + // external effects. So it runs with Full autonomy; genuinely external + // effects are still gated by the tainted origin + approval gate. Mode + // only scales how much delegation depth the tick gets. + effective.autonomy.level = crate::openhuman::security::AutonomyLevel::Full; match self.mode { SubconsciousMode::Simple => { - effective.autonomy.level = crate::openhuman::security::AutonomyLevel::ReadOnly; effective.agent.max_tool_iterations = 15; } SubconsciousMode::Aggressive | SubconsciousMode::EventDriven => { - effective.autonomy.level = crate::openhuman::security::AutonomyLevel::Full; effective.agent.max_tool_iterations = 30; } SubconsciousMode::Off => return Ok(0), @@ -386,30 +507,28 @@ impl SubconsciousEngine { ); let mode_guidance = match self.mode { - SubconsciousMode::Aggressive => { - "\n\n\ - You are in AGGRESSIVE mode. You may use `spawn_subagent` to delegate \ - complex tasks:\n\ - - `agent_id: \"orchestrator\"` with `model: \"reasoning-v1\"` for deep \ - reasoning and multi-step execution\n\ - - `agent_id: \"researcher\"` for web research and external data\n\n\ - Use this power when you identify actionable opportunities, approaching \ - deadlines, or patterns that warrant proactive help." + SubconsciousMode::Aggressive | SubconsciousMode::EventDriven => { + "\n\nYou may delegate deeper work with `spawn_async_subagent` (e.g. research \ + or multi-step execution) when you spot something genuinely actionable." } _ => "", }; let user_message = format!( - "{prompt_text}\n\n\ - ## Instructions\n\n\ - Based on the situation report and your existing scratchpad, maintain your \ - scratchpad — add new observations, edit stale entries, remove resolved items.\n\n\ - Your scratchpad IS your continuity mechanism across ticks. Keep it focused \ - and actionable.\ - {mode_guidance}", + "{prompt_text}\ + ## Your job\n\n\ + The diff above is how the user's world changed since the last check; the prepared \ + context grounds it. Decide what (if anything) deserves action:\n\ + - Record or update actionable follow-ups on the user's to-do board with `update_task` \ + (pass `threadId: \"user-tasks\"`).\n\ + - Evolve the user's long-term goals with `goals_add` / `goals_edit` when the world \ + shifts what matters to them.\n\ + - Surface anything time-sensitive or important with `notify_user`.\n\n\ + If nothing meaningful changed, do nothing — staying silent is the right call most \ + ticks. Do not invent busywork.{mode_guidance}", ); - debug!("[subconscious] spawning agent with tool access"); + debug!("[subconscious] spawning decision agent"); let source = tick_origin_source(has_external_content); debug!( "[subconscious] tick origin source={:?} has_external_content={has_external_content}", @@ -431,13 +550,68 @@ impl SubconsciousEngine { let response_chars = response.chars().count(); info!( - "[subconscious] agent completed (response {} chars)", + "[subconscious] decision agent completed (response {} chars)", response_chars ); Ok(response_chars) } } +// ── World-diff rendering ───────────────────────────────────────────────────── + +/// Total added + modified + removed across all sources in a cross-source diff. +fn world_diff_change_count(diff: &CrossSourceDiff) -> u32 { + diff.summary.added + diff.summary.modified + diff.summary.removed +} + +/// Render a [`CrossSourceDiff`] into a compact markdown summary for the decision +/// agent's prompt. Per-source change lists are capped at [`MAX_ITEMS_PER_SOURCE`] +/// so a churny source can't blow out the context window. +fn render_world_diff(diff: &CrossSourceDiff) -> String { + let s = &diff.summary; + let total = s.added + s.modified + s.removed; + if total == 0 { + return "Nothing changed across your connected sources since the last check.".to_string(); + } + + let mut out = format!( + "{total} item(s) changed across your sources since the last check \ + ({} added, {} modified, {} removed).\n", + s.added, s.modified, s.removed + ); + + for source in &diff.per_source { + let ss = &source.summary; + if ss.added + ss.modified + ss.removed == 0 { + continue; + } + out.push_str(&format!( + "\n### {} ({})\n- {} added, {} modified, {} removed\n", + source.source_label, source.source_kind, ss.added, ss.modified, ss.removed + )); + for change in source.changes.iter().take(MAX_ITEMS_PER_SOURCE) { + let verb = match change.kind { + crate::openhuman::memory_diff::types::ChangeKind::Added => "added", + crate::openhuman::memory_diff::types::ChangeKind::Removed => "removed", + crate::openhuman::memory_diff::types::ChangeKind::Modified => "modified", + }; + let label = if change.title.trim().is_empty() { + change.item_id.as_str() + } else { + change.title.as_str() + }; + out.push_str(&format!(" - [{verb}] {label}\n")); + } + if source.changes.len() > MAX_ITEMS_PER_SOURCE { + out.push_str(&format!( + " - …and {} more\n", + source.changes.len() - MAX_ITEMS_PER_SOURCE + )); + } + } + out +} + // ── Provider routing ──────────────────────────────────────────────────────── #[derive(Clone, Debug, Eq, PartialEq)] @@ -509,106 +683,6 @@ fn is_tool_capability_error(msg: &str) -> bool { || lower.contains("does not support tools") } -// ── Pre-LLM memory retrieval ──────────────────────────────────────────────── - -/// Query the memory tree using scratchpad entries as context, returning -/// a rendered markdown section to inject into the user message. This -/// replaces the old `call_memory_agent` tool call — the retrieval now -/// happens in code before the LLM runs, saving a full agent turn. -async fn retrieve_memory_context( - memory: &Option, - scratchpad_entries: &[scratchpad::ScratchpadEntry], -) -> String { - let client = match memory { - Some(c) => c, - None => { - debug!("[subconscious] no memory client — skipping pre-LLM retrieval"); - return String::new(); - } - }; - - // Build a query from high-priority scratchpad items (p5+) or fall back - // to a generic recent-activity query. - let query = build_memory_query(scratchpad_entries); - debug!( - "[subconscious] pre-LLM memory retrieval query_len={}", - query.len() - ); - - let started = std::time::Instant::now(); - - // Query conversation_memory namespace for relevant context - let conversation_ctx = client - .query_namespace("conversation_memory", &query, MEMORY_RETRIEVAL_MAX_CHUNKS) - .await - .unwrap_or_else(|e| { - warn!("[subconscious] conversation_memory query failed: {e}"); - String::new() - }); - - // Also recall recent learning reflections (user patterns, preferences) - let reflections_ctx = client - .recall_namespace("learning_reflections", 10) - .await - .ok() - .flatten() - .unwrap_or_default(); - - let elapsed = started.elapsed(); - info!( - "[subconscious] pre-LLM memory retrieval done in {:.1}s conv_chars={} refl_chars={}", - elapsed.as_secs_f64(), - conversation_ctx.len(), - reflections_ctx.len() - ); - - if conversation_ctx.is_empty() && reflections_ctx.is_empty() { - return String::new(); - } - - let mut section = String::from("## Memory Context (pre-loaded)\n\n"); - if !conversation_ctx.is_empty() { - section.push_str("### Recent Conversations & Activity\n\n"); - section.push_str(&conversation_ctx); - section.push_str("\n\n"); - } - if !reflections_ctx.is_empty() { - section.push_str("### Learned User Patterns\n\n"); - section.push_str(&reflections_ctx); - section.push_str("\n\n"); - } - section -} - -/// Build a memory query from scratchpad entries. High-priority items -/// (p5+) get included verbatim; lower-priority items contribute keywords. -/// Falls back to a generic query when the scratchpad is empty. -fn build_memory_query(entries: &[scratchpad::ScratchpadEntry]) -> String { - if entries.is_empty() { - return "What has the user been working on recently? Any upcoming deadlines, \ - unresolved threads, or notable activity across all sources?" - .to_string(); - } - - let high_priority: Vec<&scratchpad::ScratchpadEntry> = - entries.iter().filter(|e| e.priority >= 5).collect(); - - if high_priority.is_empty() { - // Use all entries as a broad query - let bodies: Vec<&str> = entries.iter().map(|e| e.body.as_str()).collect(); - return format!( - "Recent activity and updates related to: {}", - bodies.join("; ") - ); - } - - let bodies: Vec<&str> = high_priority.iter().map(|e| e.body.as_str()).collect(); - format!( - "Updates and context for these tracked items: {}", - bodies.join("; ") - ) -} - fn persist_last_tick_at(workspace_dir: &std::path::Path, tick_at: f64) { if let Err(e) = store::with_connection(workspace_dir, |conn| store::set_last_tick_at(conn, tick_at)) @@ -624,72 +698,6 @@ fn now_secs() -> f64 { .unwrap_or(0.0) } -// ── Identity loading ─────────────────────────────────────────────────────── - -const IDENTITY_EXCERPT_CHARS: usize = 2000; - -fn load_identity_context(workspace_dir: &std::path::Path) -> String { - let prompts_dir = resolve_prompts_dir(workspace_dir); - let mut ctx = String::new(); - - if let Some(ref dir) = prompts_dir { - if let Some(soul) = load_file_excerpt(dir, "SOUL.md") { - ctx.push_str(&soul); - ctx.push_str("\n\n"); - } - } - - if let Some(profile) = load_file_excerpt(workspace_dir, "PROFILE.md") { - ctx.push_str("## User Profile\n\n"); - ctx.push_str(&profile); - ctx.push_str("\n\n"); - } - - if ctx.is_empty() { - "You are OpenHuman, an AI assistant for productivity and collaboration.".to_string() - } else { - ctx - } -} - -fn resolve_prompts_dir(workspace_dir: &std::path::Path) -> Option { - let workspace_ai = workspace_dir.join("ai"); - if workspace_ai.is_dir() { - return Some(workspace_ai); - } - - if let Some(dir) = option_env!("CARGO_MANIFEST_DIR").map(std::path::PathBuf::from) { - let candidate = dir - .join("src") - .join("openhuman") - .join("agent") - .join("prompts"); - if candidate.is_dir() { - return Some(candidate); - } - } - - if let Ok(cwd) = std::env::current_dir() { - return crate::openhuman::dev_paths::repo_ai_prompts_dir(&cwd); - } - - None -} - -fn load_file_excerpt(dir: &std::path::Path, filename: &str) -> Option { - let content = std::fs::read_to_string(dir.join(filename)).ok()?; - let trimmed = content.trim(); - if trimmed.is_empty() { - return None; - } - if trimmed.chars().count() > IDENTITY_EXCERPT_CHARS { - let truncated: String = trimmed.chars().take(IDENTITY_EXCERPT_CHARS).collect(); - Some(format!("{truncated}\n[... truncated]")) - } else { - Some(trimmed.to_string()) - } -} - #[cfg(test)] #[path = "engine_tests.rs"] mod tests; diff --git a/src/openhuman/subconscious/engine_tests.rs b/src/openhuman/subconscious/engine_tests.rs index c21daeaa7..2d7f978df 100644 --- a/src/openhuman/subconscious/engine_tests.rs +++ b/src/openhuman/subconscious/engine_tests.rs @@ -48,3 +48,107 @@ fn tool_capability_error_ignores_unrelated_failures() { )); assert!(!is_tool_capability_error("agent run: request timed out")); } + +// ── World-diff rendering (Stage 1) ────────────────────────────────────── + +use crate::openhuman::memory_diff::types::{ + ChangeKind, CrossSourceDiff, DiffResult, DiffSummary, ItemChange, +}; + +fn change(item_id: &str, title: &str, kind: ChangeKind) -> ItemChange { + ItemChange { + item_id: item_id.to_string(), + title: title.to_string(), + kind, + old_content_hash: None, + new_content_hash: None, + text_diff: None, + } +} + +#[test] +fn empty_cross_source_diff_has_zero_change_count() { + let diff = CrossSourceDiff { + checkpoint_id: Some("ckpt_1".into()), + computed_at_ms: 0, + summary: DiffSummary::default(), + per_source: Vec::new(), + }; + assert_eq!(world_diff_change_count(&diff), 0); + // The "no changes" render is the quiet-tick sentinel; the tick short-circuits + // before it ever reaches the agent, but the renderer stays well-defined. + assert!(render_world_diff(&diff).contains("Nothing changed")); +} + +#[test] +fn render_world_diff_summarises_changes_per_source() { + let diff = CrossSourceDiff { + checkpoint_id: Some("ckpt_1".into()), + computed_at_ms: 0, + summary: DiffSummary { + added: 2, + modified: 1, + removed: 0, + unchanged: 5, + }, + per_source: vec![DiffResult { + source_id: "src_gmail".into(), + source_kind: "composio".into(), + source_label: "Gmail".into(), + from_snapshot_id: Some("snap_a".into()), + to_snapshot_id: "snap_b".into(), + summary: DiffSummary { + added: 2, + modified: 1, + removed: 0, + unchanged: 5, + }, + changes: vec![ + change("m1", "Invoice from Acme", ChangeKind::Added), + change("m2", "Re: launch plan", ChangeKind::Added), + change("m3", "Standup notes", ChangeKind::Modified), + ], + }], + }; + + assert_eq!(world_diff_change_count(&diff), 3); + let rendered = render_world_diff(&diff); + assert!(rendered.contains("3 item(s) changed")); + assert!(rendered.contains("Gmail (composio)")); + assert!(rendered.contains("[added] Invoice from Acme")); + assert!(rendered.contains("[modified] Standup notes")); +} + +#[test] +fn render_world_diff_caps_items_and_falls_back_to_item_id() { + let mut changes = Vec::new(); + for i in 0..(MAX_ITEMS_PER_SOURCE + 3) { + // Empty title forces the item_id fallback. + changes.push(change(&format!("item_{i}"), "", ChangeKind::Added)); + } + let n = changes.len() as u32; + let diff = CrossSourceDiff { + checkpoint_id: None, + computed_at_ms: 0, + summary: DiffSummary { + added: n, + ..DiffSummary::default() + }, + per_source: vec![DiffResult { + source_id: "src_folder".into(), + source_kind: "folder".into(), + source_label: "Notes".into(), + from_snapshot_id: None, + to_snapshot_id: "snap_x".into(), + summary: DiffSummary { + added: n, + ..DiffSummary::default() + }, + changes, + }], + }; + + let rendered = render_world_diff(&diff); + assert!(rendered.contains("[added] item_0"), "uses item_id fallback"); + assert!(rendered.contains("…and 3 more"), "caps the per-source list"); +} diff --git a/src/openhuman/subconscious/global.rs b/src/openhuman/subconscious/global.rs index bbb6fc812..bcb346243 100644 --- a/src/openhuman/subconscious/global.rs +++ b/src/openhuman/subconscious/global.rs @@ -34,13 +34,7 @@ pub async fn get_or_init_engine() -> Result .await .map_err(|e| format!("load config: {e}"))?; - let memory = crate::openhuman::memory_store::MemoryClient::from_workspace_dir( - config.workspace_dir.clone(), - ) - .ok() - .map(Arc::new); - - let engine = SubconsciousEngine::new(&config, memory); + let engine = SubconsciousEngine::new(&config); let mut guard = lock.lock().await; if guard.is_none() { diff --git a/src/openhuman/subconscious/mod.rs b/src/openhuman/subconscious/mod.rs index fc1188002..2c83ca075 100644 --- a/src/openhuman/subconscious/mod.rs +++ b/src/openhuman/subconscious/mod.rs @@ -3,9 +3,7 @@ pub mod engine; pub mod global; pub mod heartbeat; mod schemas; -pub mod scratchpad; pub mod session; -pub mod situation_report; pub mod source_chunk; pub mod store; pub mod types; diff --git a/src/openhuman/subconscious/scratchpad/mod.rs b/src/openhuman/subconscious/scratchpad/mod.rs deleted file mode 100644 index 879ba648c..000000000 --- a/src/openhuman/subconscious/scratchpad/mod.rs +++ /dev/null @@ -1,422 +0,0 @@ -//! Subconscious scratchpad — persistent working memory across ticks. -//! -//! The scratchpad holds up to `MAX_ENTRIES` (default 100) short thoughts -//! that carry over between subconscious ticks, giving the agent a -//! consistent stream-of-consciousness. -//! -//! Stored as `{workspace_dir}/subconscious/SUBCONSCIOUS_SCRATCHPAD.md` -//! so it's trivially inspectable with any text editor. -//! -//! ## File format -//! -//! # Subconscious Scratchpad -//! -//! ```json -//! [ -//! { -//! "id": "abc123", -//! "body": "The thought body goes here.", -//! "priority": 5, -//! "created_at": 1700000000.123456, -//! "updated_at": 1700000000.123456 -//! } -//! ] -//! ``` -//! -//! The agent manages its own scratchpad via three tools: -//! `scratchpad_add`, `scratchpad_edit`, `scratchpad_remove`. - -pub mod tools; - -use anyhow::{Context, Result}; -use serde::{Deserialize, Serialize}; -use std::path::{Path, PathBuf}; - -pub const DEFAULT_MAX_ENTRIES: usize = 100; - -const FILENAME: &str = "SUBCONSCIOUS_SCRATCHPAD.md"; - -/// One scratchpad entry. -#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] -pub struct ScratchpadEntry { - pub id: String, - pub body: String, - pub priority: u32, - pub created_at: f64, - pub updated_at: f64, -} - -fn scratchpad_path(workspace_dir: &Path) -> PathBuf { - workspace_dir.join("subconscious").join(FILENAME) -} - -pub fn load(workspace_dir: &Path) -> Result> { - let path = scratchpad_path(workspace_dir); - let content = match std::fs::read_to_string(&path) { - Ok(c) => c, - Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(vec![]), - Err(e) => return Err(e).context("read scratchpad"), - }; - Ok(parse_entries(&content)) -} - -pub fn add(workspace_dir: &Path, body: &str, priority: u32, max_entries: usize) -> Result { - let mut entries = load(workspace_dir)?; - let id = short_id(); - let now = now_secs(); - entries.push(ScratchpadEntry { - id: id.clone(), - body: body.to_string(), - priority, - created_at: now, - updated_at: now, - }); - evict(&mut entries, max_entries); - save(workspace_dir, &entries)?; - Ok(id) -} - -pub fn edit(workspace_dir: &Path, id: &str, body: &str, priority: Option) -> Result { - let mut entries = load(workspace_dir)?; - let Some(entry) = entries.iter_mut().find(|e| e.id == id) else { - return Ok(false); - }; - entry.body = body.to_string(); - if let Some(p) = priority { - entry.priority = p; - } - entry.updated_at = now_secs(); - save(workspace_dir, &entries)?; - Ok(true) -} - -pub fn remove(workspace_dir: &Path, id: &str) -> Result { - let mut entries = load(workspace_dir)?; - let before = entries.len(); - entries.retain(|e| e.id != id); - if entries.len() == before { - return Ok(false); - } - save(workspace_dir, &entries)?; - Ok(true) -} - -pub fn render_for_prompt(entries: &[ScratchpadEntry]) -> String { - if entries.is_empty() { - return String::new(); - } - let mut out = String::from("## Scratchpad (your persistent working memory)\n\n"); - out.push_str( - "These are your own notes from previous ticks. Update, remove, or \ - add entries as your understanding evolves.\n\n", - ); - for entry in entries { - let ts = chrono::DateTime::from_timestamp(entry.updated_at as i64, 0) - .map(|dt| dt.format("%Y-%m-%d %H:%M").to_string()) - .unwrap_or_else(|| "unknown".to_string()); - out.push_str(&format!( - "- **[{}]** (p{}) {}\n _updated: {}_\n", - entry.id, entry.priority, entry.body, ts - )); - } - out -} - -// ── File I/O ──────────────────────────────────────────────────────────────── - -fn save(workspace_dir: &Path, entries: &[ScratchpadEntry]) -> Result<()> { - let path = scratchpad_path(workspace_dir); - if let Some(parent) = path.parent() { - std::fs::create_dir_all(parent).context("create scratchpad dir")?; - } - let content = render_file(entries); - std::fs::write(&path, content).context("write scratchpad")?; - Ok(()) -} - -fn render_file(entries: &[ScratchpadEntry]) -> String { - let mut out = String::from("# Subconscious Scratchpad\n\n"); - if entries.is_empty() { - return out; - } - out.push_str("```json\n"); - match serde_json::to_string_pretty(entries) { - Ok(json) => out.push_str(&json), - Err(e) => { - log::warn!("[subconscious] failed to render scratchpad JSON: {e}"); - out.push_str("[]"); - } - } - out.push_str("\n```\n"); - out -} - -fn parse_entries(content: &str) -> Vec { - if let Some(entries) = parse_json_entries(content) { - return entries; - } - - let mut entries = Vec::new(); - - for block in content.split("\n---\n") { - let block = block.trim(); - if block.is_empty() || block.starts_with("# ") { - if let Some(rest) = block.strip_prefix("# Subconscious Scratchpad") { - let rest = rest.trim(); - if rest.is_empty() { - continue; - } - if let Some(entry) = parse_single_block(rest) { - entries.push(entry); - } - continue; - } - continue; - } - if let Some(entry) = parse_single_block(block) { - entries.push(entry); - } - } - - entries -} - -fn parse_json_entries(content: &str) -> Option> { - let marker = "```json"; - let start = content.find(marker)? + marker.len(); - let rest = &content[start..]; - let end = rest.rfind("```")?; - let json = rest[..end].trim(); - if json.is_empty() { - return Some(Vec::new()); - } - match serde_json::from_str(json) { - Ok(entries) => Some(entries), - Err(e) => { - log::warn!("[subconscious] failed to parse scratchpad JSON: {e}"); - None - } - } -} - -fn parse_single_block(block: &str) -> Option { - let meta_start = block.find("")? + meta_start + 3; - let meta_line = &block[meta_start..meta_end]; - - let id = extract_meta(meta_line, "entry:")?; - let priority = extract_meta(meta_line, "p:") - .and_then(|s| s.parse::().ok()) - .unwrap_or(0); - let created_at = extract_meta(meta_line, "created:") - .and_then(|s| s.parse::().ok()) - .unwrap_or(0.0); - let updated_at = extract_meta(meta_line, "updated:") - .and_then(|s| s.parse::().ok()) - .unwrap_or(created_at); - - let body = block[meta_end..].trim().to_string(); - if body.is_empty() { - return None; - } - - Some(ScratchpadEntry { - id, - body, - priority, - created_at, - updated_at, - }) -} - -fn extract_meta(line: &str, key: &str) -> Option { - let start = line.find(key)? + key.len(); - let rest = &line[start..]; - let end = rest - .find(|c: char| c.is_whitespace() || c == '-') - .unwrap_or(rest.len()); - let val = rest[..end].trim().to_string(); - if val.is_empty() { - None - } else { - Some(val) - } -} - -fn evict(entries: &mut Vec, max: usize) { - if entries.len() <= max { - return; - } - entries.sort_by(|a, b| { - b.priority.cmp(&a.priority).then( - b.updated_at - .partial_cmp(&a.updated_at) - .unwrap_or(std::cmp::Ordering::Equal), - ) - }); - entries.truncate(max); -} - -fn short_id() -> String { - let uuid = uuid::Uuid::new_v4().to_string(); - uuid[..8].to_string() -} - -fn now_secs() -> f64 { - std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .map(|d| d.as_secs_f64()) - .unwrap_or(0.0) -} - -#[cfg(test)] -mod tests { - use super::*; - use tempfile::TempDir; - - fn temp_workspace() -> TempDir { - tempfile::tempdir().unwrap() - } - - #[test] - fn add_and_load_round_trip() { - let ws = temp_workspace(); - let id = add(ws.path(), "first thought", 1, 100).unwrap(); - let entries = load(ws.path()).unwrap(); - assert_eq!(entries.len(), 1); - assert_eq!(entries[0].id, id); - assert_eq!(entries[0].body, "first thought"); - assert_eq!(entries[0].priority, 1); - } - - #[test] - fn edit_updates_body_and_priority() { - let ws = temp_workspace(); - let id = add(ws.path(), "old", 0, 100).unwrap(); - let found = edit(ws.path(), &id, "new", Some(5)).unwrap(); - assert!(found); - let entries = load(ws.path()).unwrap(); - assert_eq!(entries[0].body, "new"); - assert_eq!(entries[0].priority, 5); - } - - #[test] - fn remove_deletes_entry() { - let ws = temp_workspace(); - let id = add(ws.path(), "gone", 0, 100).unwrap(); - assert!(remove(ws.path(), &id).unwrap()); - assert!(load(ws.path()).unwrap().is_empty()); - } - - #[test] - fn remove_nonexistent_returns_false() { - let ws = temp_workspace(); - add(ws.path(), "keep", 0, 100).unwrap(); - assert!(!remove(ws.path(), "nope").unwrap()); - } - - #[test] - fn add_evicts_oldest_low_priority_beyond_cap() { - let ws = temp_workspace(); - for i in 0..5 { - add(ws.path(), &format!("note {i}"), 0, 100).unwrap(); - } - add(ws.path(), "high priority", 10, 100).unwrap(); - add(ws.path(), "newest", 0, 3).unwrap(); - let entries = load(ws.path()).unwrap(); - assert!(entries.len() <= 3); - assert!(entries.iter().any(|e| e.body == "high priority")); - } - - #[test] - fn render_for_prompt_formats_entries() { - let entries = vec![ScratchpadEntry { - id: "abc".to_string(), - body: "test thought".to_string(), - priority: 2, - created_at: 1700000000.0, - updated_at: 1700000000.0, - }]; - let rendered = render_for_prompt(&entries); - assert!(rendered.contains("## Scratchpad")); - assert!(rendered.contains("[abc]")); - assert!(rendered.contains("test thought")); - assert!(rendered.contains("(p2)")); - } - - #[test] - fn render_for_prompt_empty_returns_empty() { - assert!(render_for_prompt(&[]).is_empty()); - } - - #[test] - fn load_missing_file_returns_empty() { - let ws = temp_workspace(); - assert!(load(ws.path()).unwrap().is_empty()); - } - - #[test] - fn file_is_readable_markdown() { - let ws = temp_workspace(); - add(ws.path(), "thought one", 5, 100).unwrap(); - add(ws.path(), "thought two", 0, 100).unwrap(); - let content = std::fs::read_to_string(scratchpad_path(ws.path())).unwrap(); - assert!(content.starts_with("# Subconscious Scratchpad")); - assert!(content.contains("```json")); - assert!(content.contains("thought one")); - assert!(content.contains("thought two")); - assert!(content.contains("\"created_at\"")); - } - - #[test] - fn parse_round_trip_preserves_data() { - let original = vec![ - ScratchpadEntry { - id: "aaa".to_string(), - body: "first".to_string(), - priority: 3, - created_at: 1700000000.0, - updated_at: 1700000100.0, - }, - ScratchpadEntry { - id: "bbb".to_string(), - body: "second\nwith newlines".to_string(), - priority: 0, - created_at: 1700000050.0, - updated_at: 1700000050.0, - }, - ]; - let rendered = render_file(&original); - let parsed = parse_entries(&rendered); - assert_eq!(parsed.len(), 2); - assert_eq!(parsed[0].id, "aaa"); - assert_eq!(parsed[0].body, "first"); - assert_eq!(parsed[0].priority, 3); - assert_eq!(parsed[1].id, "bbb"); - assert_eq!(parsed[1].body, "second\nwith newlines"); - } - - #[test] - fn parse_round_trip_preserves_markdown_separators_and_subsecond_timestamps() { - let original = vec![ScratchpadEntry { - id: "aaa".to_string(), - body: "first\n---\nsecond".to_string(), - priority: 3, - created_at: 1700000000.123456, - updated_at: 1700000100.654321, - }]; - let rendered = render_file(&original); - let parsed = parse_entries(&rendered); - assert_eq!(parsed, original); - } - - #[test] - fn parse_legacy_markdown_entries() { - let content = "# Subconscious Scratchpad\n\n\nlegacy thought"; - let parsed = parse_entries(content); - assert_eq!(parsed.len(), 1); - assert_eq!(parsed[0].id, "abc"); - assert_eq!(parsed[0].body, "legacy thought"); - assert_eq!(parsed[0].priority, 2); - } -} diff --git a/src/openhuman/subconscious/scratchpad/tools.rs b/src/openhuman/subconscious/scratchpad/tools.rs deleted file mode 100644 index 0d9dbb9ac..000000000 --- a/src/openhuman/subconscious/scratchpad/tools.rs +++ /dev/null @@ -1,230 +0,0 @@ -//! Agent tools for managing the subconscious scratchpad. -//! -//! Three tools: `scratchpad_add`, `scratchpad_edit`, `scratchpad_remove`. -//! Registered via `all_scratchpad_tools()` and wired into the tool -//! registry so the subconscious agent can manage its own working memory. - -use crate::openhuman::tools::traits::{PermissionLevel, Tool, ToolCategory, ToolResult, ToolScope}; -use async_trait::async_trait; -use serde_json::json; - -async fn workspace_dir() -> anyhow::Result { - let config = crate::openhuman::config::load_config_with_timeout() - .await - .map_err(|e| anyhow::anyhow!("config load: {e}"))?; - Ok(config.workspace_dir) -} - -pub fn all_scratchpad_tools() -> Vec> { - vec![ - Box::new(ScratchpadAddTool), - Box::new(ScratchpadEditTool), - Box::new(ScratchpadRemoveTool), - ] -} - -// ── scratchpad_add ────────────────────────────────────────────────────────── - -pub struct ScratchpadAddTool; - -#[async_trait] -impl Tool for ScratchpadAddTool { - fn name(&self) -> &str { - "scratchpad_add" - } - - fn description(&self) -> &str { - "Add a thought to the subconscious scratchpad. Use this to persist \ - observations, hypotheses, or follow-up items across ticks. Max 100 entries; \ - oldest low-priority entries are evicted when the cap is reached." - } - - fn parameters_schema(&self) -> serde_json::Value { - json!({ - "type": "object", - "required": ["body"], - "properties": { - "body": { - "type": "string", - "description": "The thought or note to persist (keep under ~500 chars)." - }, - "priority": { - "type": "integer", - "description": "Priority 0-10. Higher priority entries survive eviction longer. Default: 0.", - "minimum": 0, - "maximum": 10 - } - } - }) - } - - fn category(&self) -> ToolCategory { - ToolCategory::System - } - - fn permission_level(&self) -> PermissionLevel { - PermissionLevel::Write - } - - fn scope(&self) -> ToolScope { - ToolScope::AgentOnly - } - - async fn execute(&self, args: serde_json::Value) -> anyhow::Result { - let body = args - .get("body") - .and_then(|v| v.as_str()) - .ok_or_else(|| anyhow::anyhow!("scratchpad_add: `body` is required"))?; - let priority = args - .get("priority") - .and_then(|v| v.as_u64()) - .unwrap_or(0) - .min(10) as u32; - - let ws = workspace_dir().await?; - let id = super::add(&ws, body, priority, super::DEFAULT_MAX_ENTRIES)?; - - Ok(ToolResult::success(format!( - "Added scratchpad entry id={id} priority={priority}" - ))) - } -} - -// ── scratchpad_edit ───────────────────────────────────────────────────────── - -pub struct ScratchpadEditTool; - -#[async_trait] -impl Tool for ScratchpadEditTool { - fn name(&self) -> &str { - "scratchpad_edit" - } - - fn description(&self) -> &str { - "Edit an existing scratchpad entry. Update the body text and/or priority." - } - - fn parameters_schema(&self) -> serde_json::Value { - json!({ - "type": "object", - "required": ["id", "body"], - "properties": { - "id": { - "type": "string", - "description": "The entry ID to edit (shown in the scratchpad section as [id])." - }, - "body": { - "type": "string", - "description": "Updated thought text." - }, - "priority": { - "type": "integer", - "description": "Updated priority 0-10 (omit to keep current).", - "minimum": 0, - "maximum": 10 - } - } - }) - } - - fn category(&self) -> ToolCategory { - ToolCategory::System - } - - fn permission_level(&self) -> PermissionLevel { - PermissionLevel::Write - } - - fn scope(&self) -> ToolScope { - ToolScope::AgentOnly - } - - async fn execute(&self, args: serde_json::Value) -> anyhow::Result { - let id = args - .get("id") - .and_then(|v| v.as_str()) - .ok_or_else(|| anyhow::anyhow!("scratchpad_edit: `id` is required"))?; - let body = args - .get("body") - .and_then(|v| v.as_str()) - .ok_or_else(|| anyhow::anyhow!("scratchpad_edit: `body` is required"))?; - let priority = args - .get("priority") - .and_then(|v| v.as_u64()) - .map(|v| v.min(10) as u32); - - let ws = workspace_dir().await?; - let found = super::edit(&ws, id, body, priority)?; - - if found { - Ok(ToolResult::success(format!( - "Updated scratchpad entry id={id}" - ))) - } else { - Ok(ToolResult::error(format!( - "Scratchpad entry id={id} not found" - ))) - } - } -} - -// ── scratchpad_remove ─────────────────────────────────────────────────────── - -pub struct ScratchpadRemoveTool; - -#[async_trait] -impl Tool for ScratchpadRemoveTool { - fn name(&self) -> &str { - "scratchpad_remove" - } - - fn description(&self) -> &str { - "Remove a scratchpad entry by ID. Use when a thought is no longer relevant \ - or has been fully addressed." - } - - fn parameters_schema(&self) -> serde_json::Value { - json!({ - "type": "object", - "required": ["id"], - "properties": { - "id": { - "type": "string", - "description": "The entry ID to remove." - } - } - }) - } - - fn category(&self) -> ToolCategory { - ToolCategory::System - } - - fn permission_level(&self) -> PermissionLevel { - PermissionLevel::Write - } - - fn scope(&self) -> ToolScope { - ToolScope::AgentOnly - } - - async fn execute(&self, args: serde_json::Value) -> anyhow::Result { - let id = args - .get("id") - .and_then(|v| v.as_str()) - .ok_or_else(|| anyhow::anyhow!("scratchpad_remove: `id` is required"))?; - - let ws = workspace_dir().await?; - let found = super::remove(&ws, id)?; - - if found { - Ok(ToolResult::success(format!( - "Removed scratchpad entry id={id}" - ))) - } else { - Ok(ToolResult::error(format!( - "Scratchpad entry id={id} not found" - ))) - } - } -} diff --git a/src/openhuman/subconscious/session.rs b/src/openhuman/subconscious/session.rs index 96fdaf3eb..a2fb62366 100644 --- a/src/openhuman/subconscious/session.rs +++ b/src/openhuman/subconscious/session.rs @@ -178,7 +178,8 @@ impl LongLivedSession { let effective = effective_config(config, self.mode); // Build as the `subconscious` agent (not the default orchestrator) so // the session's promoted turns get the subconscious tool surface — - // scratchpad + spawn_subagent + the notify_user user-handoff tool. + // memory_diff + agent_prepare_context + global to-dos/goals + the + // notify_user user-handoff tool. let mut agent = Agent::from_config_for_agent(&effective, "subconscious").map_err(|e| { warn!("[subconscious::session] agent init failed: {e}"); format!("agent init: {e}") diff --git a/src/openhuman/subconscious/situation_report/mod.rs b/src/openhuman/subconscious/situation_report/mod.rs deleted file mode 100644 index 4285d2d27..000000000 --- a/src/openhuman/subconscious/situation_report/mod.rs +++ /dev/null @@ -1,214 +0,0 @@ -//! Situation report assembly for the subconscious tick (#623). -//! -//! Replaces the legacy unified-store-backed report with sections derived -//! from the memory tree: -//! -//! 1. **Environment** (kept): host/OS/workspace/time anchor. -//! 2. **Your Identifiers** (#1365): the user's connected-account -//! identifiers (Slack/Gmail/Notion handles, emails, user_ids) so the -//! LLM can disambiguate body-text mentions — "Cyrus said X" -//! is the user iff `Cyrus` (or the email/handle) appears in this list. -//! 3. **Pending Tasks** (kept): subconscious task list from SQLite. -//! 4. **Recently-sealed summaries** (new): rows from `mem_tree_summaries` -//! grouped by tree. -//! 5. **Source-tree recap window** (new): recent source summaries since -//! `last_tick_at`. -//! -//! Sections are appended in priority order; truncation drops the tail -//! when `token_budget` is exceeded. -//! -//! Each submodule is responsible for one section so churn stays local. - -use std::path::Path; - -use crate::openhuman::config::Config; - -mod query_window; -mod summaries; - -/// Rough chars-per-token estimate for budget enforcement. -const CHARS_PER_TOKEN: usize = 4; - -/// Result of building a subconscious situation report. -/// -/// `has_external_content` is true iff the prompt now contains content -/// derived from third-party sync sources (Gmail / Slack / Notion / chat -/// transcripts / sealed source summaries). The subconscious engine uses -/// this signal to upgrade the tick's `AgentTurnOrigin` to -/// `TrustedAutomationSource::SubconsciousTainted`, which makes the -/// approval gate refuse external_effect tools for the rest of the tick. -#[derive(Debug, Clone)] -pub struct SituationReport { - pub prompt_text: String, - pub has_external_content: bool, -} - -/// Build the situation report for one subconscious tick. -/// -/// `last_tick_at` is 0.0 on cold start (include everything in the -/// configured windows). `token_budget` caps total output; sections -/// after the cap are truncated with a marker. -pub async fn build_situation_report( - config: &Config, - workspace_dir: &Path, - last_tick_at: f64, - token_budget: u32, -) -> SituationReport { - let char_budget = (token_budget as usize) * CHARS_PER_TOKEN; - let mut report = String::with_capacity(char_budget.min(64_000)); - let mut remaining = char_budget; - let mut has_external_content = false; - - // Section 1: environment anchor. - let env_section = build_environment_section(workspace_dir); - append_section(&mut report, &mut remaining, &env_section); - - // Section 2 (#1365): the user's connected-account identifiers, so - // the LLM can disambiguate "Cyrus said X" from body text - // — that's the user iff the identifier list claims it. - let identifiers_section = build_identifiers_section(); - append_section(&mut report, &mut remaining, &identifiers_section); - - // Section 3: pending subconscious tasks. - let tasks_section = build_tasks_section(workspace_dir); - append_section(&mut report, &mut remaining, &tasks_section); - - // Section 4: recently-sealed source summaries since last tick. - let (summaries_section, summaries_tainted) = - summaries::build_section(config, last_tick_at).await; - append_section(&mut report, &mut remaining, &summaries_section); - has_external_content |= summaries_tainted; - - // Section 5: source-tree recap window since last tick. - let (recap_section, recap_tainted) = query_window::build_section(config, last_tick_at).await; - append_section(&mut report, &mut remaining, &recap_section); - has_external_content |= recap_tainted; - - if report.trim().is_empty() { - report.push_str("No state changes detected since last tick.\n"); - } - - SituationReport { - prompt_text: report, - has_external_content, - } -} - -fn build_environment_section(workspace_dir: &Path) -> String { - let host = - hostname::get().map_or_else(|_| "unknown".into(), |h| h.to_string_lossy().to_string()); - let now = chrono::Local::now(); - format!( - "## Environment\n\n\ - Workspace: {}\n\ - Host: {} | OS: {}\n\ - Time: {}\n", - workspace_dir.display(), - host, - std::env::consts::OS, - now.format("%Y-%m-%d %H:%M:%S %Z"), - ) -} - -/// Render the user's connected-account identifiers (#1365) so the -/// LLM can correlate body-text mentions back to the user. -/// Empty string when no providers are connected — the section just -/// disappears rather than rendering an empty header. -fn build_identifiers_section() -> String { - let identities = crate::openhuman::composio::providers::profile::load_connected_identities(); - if identities.is_empty() { - return String::new(); - } - let body = crate::openhuman::composio::providers::profile::render_connected_identities_section( - &identities, - ); - if body.trim().is_empty() { - return String::new(); - } - let renamed = body.replacen("## Connected Identities", "## Your Identifiers", 1); - let mut out = renamed; - if !out.ends_with('\n') { - out.push('\n'); - } - out.push_str( - "\nWhen body text in later sections mentions any of the above (name, email, \ - handle, or user_id), treat it as the user's own activity. Anything else is \ - someone else.\n", - ); - out -} - -fn build_tasks_section(_workspace_dir: &Path) -> String { - String::new() -} - -/// Append a section, truncating at a UTF-8 char boundary if it overflows -/// the remaining budget. -fn append_section(report: &mut String, remaining: &mut usize, section: &str) { - if *remaining == 0 { - return; - } - let needed = section.len().saturating_add(1); - if needed <= *remaining { - report.push_str(section); - report.push('\n'); - *remaining -= needed; - } else { - let budget = *remaining; - let truncate_at = crate::openhuman::util::floor_char_boundary(section, budget); - report.push_str(§ion[..truncate_at]); - report.push_str("\n[... truncated — token budget exceeded]\n"); - *remaining = 0; - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn environment_section_contains_os_and_host() { - let section = build_environment_section(Path::new("/tmp/workspace")); - assert!(section.contains("## Environment")); - assert!(section.contains("Workspace: /tmp/workspace")); - assert!(section.contains("OS:")); - } - - #[test] - fn append_section_truncates_on_budget() { - let mut report = String::new(); - let mut remaining = 10; - append_section(&mut report, &mut remaining, "Hello, this is a long section"); - assert!(report.starts_with("Hello, thi")); - assert!(report.contains("truncated")); - assert_eq!(remaining, 0); - } - - #[test] - fn append_section_exact_fit_does_not_underflow() { - let mut report = String::new(); - let mut remaining = 6; - append_section(&mut report, &mut remaining, "Hello"); - assert_eq!(report, "Hello\n"); - assert_eq!(remaining, 0); - } - - #[test] - fn append_section_truncates_at_char_boundary() { - let mut report = String::new(); - let mut remaining = 5; - append_section(&mut report, &mut remaining, "日本語タスク"); - assert!(report.starts_with("日")); - assert!(report.contains("truncated")); - assert_eq!(remaining, 0); - } - - #[test] - fn append_section_fits_within_budget() { - let mut report = String::new(); - let mut remaining = 1000; - append_section(&mut report, &mut remaining, "Short"); - assert!(report.contains("Short")); - assert!(remaining < 1000); - } -} diff --git a/src/openhuman/subconscious/situation_report/query_window.rs b/src/openhuman/subconscious/situation_report/query_window.rs deleted file mode 100644 index 71431c2ce..000000000 --- a/src/openhuman/subconscious/situation_report/query_window.rs +++ /dev/null @@ -1,150 +0,0 @@ -//! Source-tree recap window section (#623). -//! -//! Walks the per-source summary trees (`retrieval::query_source`) for the -//! window between `last_tick_at` and now. Translates seconds-since-last-tick -//! into a day window (rounded up to ≥ 1 so cold start still produces a -//! useful recap). The global digest tree was removed — source trees plus -//! the entity index are the substrate, so the recap is reconstructed by -//! walking source-tree summaries across the window. -//! -//! Failures degrade gracefully — the section just reports -//! "Recap unavailable" rather than aborting the tick. - -use std::fmt::Write; - -use crate::openhuman::config::Config; -use crate::openhuman::memory_tree::retrieval::query_source; - -/// Cold-start fallback window when `last_tick_at` is unset. -const COLD_START_DAYS: u32 = 7; - -/// Minimum window — sub-day windows round up to one day. -const MIN_WINDOW_DAYS: u32 = 1; - -/// Max source summaries to pull into the recap window. -const MAX_RECAP_HITS: usize = 20; - -/// Build the source-tree recap window section. -/// -/// Returns `(section_markdown, has_external_content)` — the bool is -/// `true` iff at least one fresh source-tree hit was rendered, which -/// means the prompt now carries third-party sync content. See the -/// matching note on `summaries::build_section` for why the caller uses -/// this to upgrade the tick's turn origin. -pub async fn build_section(config: &Config, last_tick_at: f64) -> (String, bool) { - let window_days = compute_window_days(last_tick_at); - log::debug!( - "[subconscious::situation_report::query_window] window_days={window_days} \ - last_tick_at={last_tick_at}" - ); - - let resp = match query_source(config, None, None, Some(window_days), None, MAX_RECAP_HITS).await - { - Ok(r) => r, - Err(e) => { - log::warn!("[subconscious::situation_report::query_window] failed: {e}"); - return ("## Recap window\n\nRecap unavailable.\n".to_string(), false); - } - }; - - // Post-filter the hits against `last_tick_at`. The window rounds up to - // whole days (`MIN_WINDOW_DAYS=1`), so even a 5-minute gap between ticks - // pulls back the same 24h window of source summaries — those would - // re-feed the LLM the very content that produced the last tick's - // reflections, and the no-insert-time-dedupe path on - // `persist_and_surface_reflections` would happily store the - // duplicates. Cutoff semantics match `summaries::build_section`: - // anything whose `time_range_end` is at or before `last_tick_at` has - // already been considered; suppress it. - let fresh_hits: Vec<&_> = if last_tick_at > 0.0 { - let cutoff = last_tick_at as i64; - resp.hits - .iter() - .filter(|h| h.time_range_end.timestamp() > cutoff) - .collect() - } else { - // Cold start — keep everything inside the configured window. - resp.hits.iter().collect() - }; - - if fresh_hits.is_empty() { - return ( - format!( - "## Recap window ({} day{})\n\nNo new recap content since last tick.\n", - window_days, - if window_days == 1 { "" } else { "s" } - ), - false, - ); - } - - let mut section = format!( - "## Recap window ({} day{})\n\n", - window_days, - if window_days == 1 { "" } else { "s" } - ); - for hit in fresh_hits { - let _ = writeln!( - section, - "- L{} {} → {}: {}", - hit.level, - hit.time_range_start.format("%Y-%m-%d"), - hit.time_range_end.format("%Y-%m-%d"), - truncate(&hit.content, 600) - ); - } - (section, true) -} - -fn compute_window_days(last_tick_at: f64) -> u32 { - if last_tick_at <= 0.0 { - return COLD_START_DAYS; - } - let now = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .map(|d| d.as_secs_f64()) - .unwrap_or(last_tick_at); - let secs = (now - last_tick_at).max(0.0); - let days = (secs / 86_400.0).ceil() as u32; - days.max(MIN_WINDOW_DAYS) -} - -fn truncate(text: &str, max_chars: usize) -> String { - let trimmed = text.trim(); - if trimmed.chars().count() <= max_chars { - return trimmed.replace('\n', " "); - } - let mut out: String = trimmed.chars().take(max_chars).collect(); - out.push('…'); - out.replace('\n', " ") -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn cold_start_uses_default_window() { - assert_eq!(compute_window_days(0.0), COLD_START_DAYS); - } - - #[test] - fn small_delta_rounds_up_to_min() { - // 30 seconds ago — should still produce a 1-day window. - let now = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap() - .as_secs_f64(); - assert_eq!(compute_window_days(now - 30.0), 1); - } - - #[test] - fn multi_day_delta_rounds_up() { - let now = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap() - .as_secs_f64(); - // ~2.5 days ago should yield 3. - assert_eq!(compute_window_days(now - 2.5 * 86_400.0), 3); - } -} diff --git a/src/openhuman/subconscious/situation_report/summaries.rs b/src/openhuman/subconscious/situation_report/summaries.rs deleted file mode 100644 index af467989d..000000000 --- a/src/openhuman/subconscious/situation_report/summaries.rs +++ /dev/null @@ -1,114 +0,0 @@ -//! Recently-sealed summaries section (#623). -//! -//! Reads `mem_tree_summaries` rows sealed since the last tick, grouped -//! by their parent tree's scope label, and emits a markdown bullet list. - -use std::fmt::Write; - -use crate::openhuman::config::Config; - -/// Hard ceiling on rows fetched. The tick LLM only needs a bounded -/// pre-cooked recap — anything beyond ~8 entries is noise. -const MAX_SUMMARIES: usize = 8; - -/// Per-summary content cap — keep prompts compact. -const SUMMARY_CONTENT_PREVIEW: usize = 320; - -/// Build the recently-sealed-summaries section. -/// -/// Returns `(section_markdown, has_external_content)` — the bool is -/// `true` iff at least one summary row was actually rendered (i.e. the -/// section contains third-party sync content, not the empty -/// "No new sealed summaries" placeholder). The caller uses the flag to -/// upgrade the subconscious turn origin to -/// [`crate::openhuman::agent::turn_origin::TrustedAutomationSource::SubconsciousTainted`] -/// so external_effect tools are refused for the rest of the tick. -pub async fn build_section(config: &Config, last_tick_at: f64) -> (String, bool) { - log::debug!( - "[subconscious::situation_report::summaries] building section last_tick_at={last_tick_at}" - ); - - // Cold start gates everything in by widening the cutoff to 0. - let cutoff_ms: i64 = if last_tick_at <= 0.0 { - 0 - } else { - (last_tick_at * 1000.0) as i64 - }; - - let rows = match read_recent_summaries(config, cutoff_ms) { - Ok(rows) => rows, - Err(e) => { - log::warn!("[subconscious::situation_report::summaries] read failed: {e}"); - return ( - "## Recent summaries\n\nSummaries unavailable.\n".to_string(), - false, - ); - } - }; - - if rows.is_empty() { - return ( - "## Recent summaries\n\nNo new sealed summaries since last tick.\n".to_string(), - false, - ); - } - - let mut section = String::from("## Recent summaries\n\n"); - let _ = writeln!( - section, - "{} summaries sealed since last tick (most recent first):", - rows.len() - ); - section.push('\n'); - for row in &rows { - let preview = truncate(&row.content, SUMMARY_CONTENT_PREVIEW); - let _ = writeln!( - section, - "- **[{}]** L{} {} — {}", - row.tree_scope, row.level, row.summary_id, preview - ); - } - (section, true) -} - -#[derive(Debug)] -struct SummaryRow { - summary_id: String, - tree_scope: String, - level: u32, - content: String, -} - -fn read_recent_summaries(config: &Config, cutoff_ms: i64) -> anyhow::Result> { - crate::openhuman::memory_store::chunks::store::with_connection(config, |conn| { - let mut stmt = conn.prepare( - "SELECT s.id, s.level, s.content, t.scope - FROM mem_tree_summaries s - JOIN mem_tree_trees t ON t.id = s.tree_id - WHERE s.sealed_at_ms > ?1 AND s.deleted = 0 - ORDER BY s.sealed_at_ms DESC - LIMIT ?2", - )?; - let rows = stmt - .query_map(rusqlite::params![cutoff_ms, MAX_SUMMARIES as i64], |row| { - Ok(SummaryRow { - summary_id: row.get(0)?, - level: row.get::<_, i64>(1)? as u32, - content: row.get(2)?, - tree_scope: row.get(3)?, - }) - })? - .collect::, _>>()?; - Ok(rows) - }) -} - -fn truncate(text: &str, max_chars: usize) -> String { - let trimmed = text.trim(); - if trimmed.chars().count() <= max_chars { - return trimmed.replace('\n', " "); - } - let mut out: String = trimmed.chars().take(max_chars).collect(); - out.push('…'); - out.replace('\n', " ") -} diff --git a/src/openhuman/subconscious/source_chunk.rs b/src/openhuman/subconscious/source_chunk.rs index 4ec709977..31f285791 100644 --- a/src/openhuman/subconscious/source_chunk.rs +++ b/src/openhuman/subconscious/source_chunk.rs @@ -130,9 +130,8 @@ fn resolve_one(config: &crate::openhuman::config::Config, raw: &str) -> SourceCh } } -/// Look up a sealed summary by id. Mirrors the read pattern in -/// [`crate::openhuman::subconscious::situation_report::summaries`] but -/// fetches a single row instead of the recent-summaries window. The +/// Look up a sealed summary by id. Fetches a single `mem_tree_summaries` row. +/// The /// resolved `content` is truncated to [`PREVIEW_MAX_CHARS`] so the /// reflection row stays bounded; full content remains queryable from /// `mem_tree_summaries` if a future feature needs it. diff --git a/src/openhuman/subconscious/store.rs b/src/openhuman/subconscious/store.rs index a981d1c4f..1fdd86086 100644 --- a/src/openhuman/subconscious/store.rs +++ b/src/openhuman/subconscious/store.rs @@ -238,6 +238,14 @@ const SCHEMA_DDL: &str = " key TEXT PRIMARY KEY, value REAL NOT NULL ); + + -- Text-valued engine state (the REAL-typed `subconscious_state` above can't + -- hold ids). Used for the memory_diff baseline checkpoint id the structured + -- tick diffs against. + CREATE TABLE IF NOT EXISTS subconscious_state_text ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL + ); "; #[cfg(test)] @@ -246,6 +254,7 @@ pub(crate) const SCHEMA_DDL_FOR_TESTS: &str = SCHEMA_DDL; // ── Engine state KV ────────────────────────────────────────────────────────── const STATE_KEY_LAST_TICK_AT: &str = "last_tick_at"; +const STATE_KEY_BASELINE_CHECKPOINT_ID: &str = "baseline_checkpoint_id"; pub fn get_last_tick_at(conn: &Connection) -> Result { let value: Option = conn @@ -266,6 +275,28 @@ pub fn set_last_tick_at(conn: &Connection, value: f64) -> Result<()> { Ok(()) } +/// The memory_diff checkpoint the structured tick diffs against — i.e. the +/// snapshot of "the agent's world" captured at the end of the previous tick. +/// `None` until the first tick establishes a baseline. +pub fn get_baseline_checkpoint_id(conn: &Connection) -> Result> { + let value: Option = conn + .query_row( + "SELECT value FROM subconscious_state_text WHERE key = ?1", + [STATE_KEY_BASELINE_CHECKPOINT_ID], + |row| row.get(0), + ) + .optional()?; + Ok(value) +} + +pub fn set_baseline_checkpoint_id(conn: &Connection, checkpoint_id: &str) -> Result<()> { + conn.execute( + "INSERT OR REPLACE INTO subconscious_state_text (key, value) VALUES (?1, ?2)", + rusqlite::params![STATE_KEY_BASELINE_CHECKPOINT_ID, checkpoint_id], + )?; + Ok(()) +} + fn now_secs() -> f64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) diff --git a/src/openhuman/subconscious/store_tests.rs b/src/openhuman/subconscious/store_tests.rs index 12a333669..9ce8cefc3 100644 --- a/src/openhuman/subconscious/store_tests.rs +++ b/src/openhuman/subconscious/store_tests.rs @@ -22,6 +22,24 @@ fn last_tick_at_upsert() { assert_eq!(get_last_tick_at(&conn).unwrap(), 2.0); } +#[test] +fn baseline_checkpoint_id_round_trip() { + let conn = test_conn(); + // Unset until the first tick establishes a baseline. + assert_eq!(get_baseline_checkpoint_id(&conn).unwrap(), None); + set_baseline_checkpoint_id(&conn, "ckpt_abc").unwrap(); + assert_eq!( + get_baseline_checkpoint_id(&conn).unwrap(), + Some("ckpt_abc".to_string()) + ); + // Advancing the baseline replaces the previous id. + set_baseline_checkpoint_id(&conn, "ckpt_def").unwrap(); + assert_eq!( + get_baseline_checkpoint_id(&conn).unwrap(), + Some("ckpt_def".to_string()) + ); +} + #[test] fn schema_ddl_creates_tables() { let conn = test_conn(); diff --git a/src/openhuman/tools/ops.rs b/src/openhuman/tools/ops.rs index 763698f43..7894e962f 100644 --- a/src/openhuman/tools/ops.rs +++ b/src/openhuman/tools/ops.rs @@ -531,8 +531,10 @@ pub fn all_tools_with_runtime( memory_hybrid_search, memory_store_raw_search, memory_store_raw_chunks, memory_store_kinds" ); - // Subconscious scratchpad tools — persistent working memory across ticks. - tools.extend(crate::openhuman::subconscious::scratchpad::tools::all_scratchpad_tools()); + // Memory diff — structured "what changed in the agent's world since a + // checkpoint/last sync". Drives the subconscious tick's first stage and is + // available to any agent that lists it. Unit struct, no runtime deps. + tools.push(Box::new(crate::openhuman::memory_diff::MemoryDiffTool)); // Subconscious user-facing handoff — notify_user proactive delivery. tools.extend(crate::openhuman::subconscious::user_thread::all_user_thread_tools()); diff --git a/tests/subconscious_e2e.rs b/tests/subconscious_e2e.rs index 29cf883ed..5473ff17b 100644 --- a/tests/subconscious_e2e.rs +++ b/tests/subconscious_e2e.rs @@ -42,21 +42,22 @@ async fn ingest_doc( result.document_id } -/// Two-tick E2E test — verifies the agent-per-tick model can process -/// ingested memory data and persist tick state. +/// Two-tick E2E test — verifies the structured tick model persists tick +/// state across runs. The first tick has no world baseline, so it +/// establishes one; subsequent ticks diff against it. (Exercising the full +/// diff → decide path additionally requires configured + synced memory +/// sources; this smoke test focuses on the tick lifecycle + state.) #[tokio::test] #[ignore] // requires running Ollama async fn two_tick_e2e_with_real_ollama() { use openhuman_core::openhuman::embeddings::NoopEmbedding; - use openhuman_core::openhuman::memory::{MemoryClient, UnifiedMemory}; + use openhuman_core::openhuman::memory::UnifiedMemory; use openhuman_core::openhuman::subconscious::store; let tmp = tempfile::tempdir().expect("tempdir"); let workspace = tmp.path(); let memory = UnifiedMemory::new(workspace, Arc::new(NoopEmbedding), None).expect("init memory"); - let memory_client = - MemoryClient::from_workspace_dir(workspace.to_path_buf()).expect("memory client"); // Ingest test data ingest_doc( @@ -80,10 +81,7 @@ async fn two_tick_e2e_with_real_ollama() { config.local_ai.runtime_enabled = true; config.local_ai.usage.subconscious = true; - let engine = openhuman_core::openhuman::subconscious::SubconsciousEngine::new( - &config, - Some(Arc::new(memory_client)), - ); + let engine = openhuman_core::openhuman::subconscious::SubconsciousEngine::new(&config); // Tick 1 println!("\n=== TICK 1 ===");