diff --git a/app/src-tauri/Cargo.lock b/app/src-tauri/Cargo.lock index b232d5e75a..309596244c 100644 --- a/app/src-tauri/Cargo.lock +++ b/app/src-tauri/Cargo.lock @@ -789,9 +789,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.4.2" +version = "1.4.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5d262e149917187838d5b42777c8253bcb64500067342904e7d429499a6f277e" +checksum = "0ad534f4357a5264cce5019c989cf66a4f0dc4e0d1b1d15f8aacec0ff7360273" dependencies = [ "find-msvc-tools", "jobserver", @@ -1735,9 +1735,9 @@ dependencies = [ [[package]] name = "either" -version = "1.17.0" +version = "1.18.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e5e8f6c15a24b9a3ee5efec809ccd006d3b30e8b3bb63c39af737c7f87daa1d" +checksum = "252afb9ae5eaa683babdc6a068b3f5726eb19e05070c731f9b2a23a7c3e8ed34" [[package]] name = "email-encoding" @@ -1877,9 +1877,9 @@ dependencies = [ [[package]] name = "error-code" -version = "3.3.2" +version = "3.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dea2df4cf52843e0452895c455a1a2cfbb842a1e7329671acf418fdc53ed4c59" +checksum = "0b5343afd4a8365a643ac588dab4cf234a190c7f6c88c9f6dd6ffe00837661b7" [[package]] name = "event-listener" @@ -1968,9 +1968,9 @@ dependencies = [ [[package]] name = "find-msvc-tools" -version = "0.1.10" +version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26b73573e6edcd2af0cdf47bd6cb58f0b3839491263c314eaad1ccf24430e1de" +checksum = "d45db016d36b838f563236e9193d0ee6ce38f3f68b6c94e914b4929c96bbb890" [[package]] name = "findshlibs" @@ -2507,9 +2507,9 @@ dependencies = [ [[package]] name = "h2" -version = "0.4.15" +version = "0.4.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6cb093c84e8bd9b188d4c4a8cb6579fc016968d14c99882163cd3ff402a4f155" +checksum = "839c0e8a181239723652be9062bb56ca5bf5f64011f73b623f6f4fc59086a228" dependencies = [ "atomic-waker", "bytes", @@ -2798,7 +2798,7 @@ dependencies = [ "js-sys", "log", "wasm-bindgen", - "windows-core 0.57.0", + "windows-core 0.62.2", ] [[package]] @@ -2822,9 +2822,9 @@ dependencies = [ [[package]] name = "icu_collections" -version = "2.2.0" +version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2984d1cd16c883d7935b9e07e44071dca8d917fd52ecc02c04d5fa0b5a3f191c" +checksum = "fa68d21081c4a05d5a901a1c62add574c77048b6a1c67be3b50ce0b60d4ca513" dependencies = [ "displaydoc", "potential_utf", @@ -2836,9 +2836,9 @@ dependencies = [ [[package]] name = "icu_locale_core" -version = "2.2.0" +version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "92219b62b3e2b4d88ac5119f8904c10f8f61bf7e95b640d25ba3075e6cac2c29" +checksum = "d56e28588da92eee5c3201a6eff33fabdd49b62269c8938d4ff050ce4d900deb" dependencies = [ "displaydoc", "litemap", @@ -2849,9 +2849,9 @@ dependencies = [ [[package]] name = "icu_normalizer" -version = "2.2.0" +version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c56e5ee99d6e3d33bd91c5d85458b6005a22140021cc324cea84dd0e72cff3b4" +checksum = "12f9cf5f235641ed274641dd81c3f28d870e276763d0797aeeab72317b1c646f" dependencies = [ "icu_collections", "icu_normalizer_data", @@ -2863,16 +2863,17 @@ dependencies = [ [[package]] name = "icu_normalizer_data" -version = "2.2.0" +version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "da3be0ae77ea334f4da67c12f149704f19f81d1adf7c51cf482943e84a2bad38" +checksum = "1563da1ed3e0b3bf3d74c9b85917ac9c56464d2f57242270c09c9e752f8021a0" [[package]] name = "icu_properties" -version = "2.2.0" +version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bee3b67d0ea5c2cca5003417989af8996f8604e34fb9ddf96208a033901e70de" +checksum = "7e7ca276ad3145661a65914e6daf131ca5120cd3dcee8f8f3214b8875184a148" dependencies = [ + "displaydoc", "icu_collections", "icu_locale_core", "icu_properties_data", @@ -2883,15 +2884,15 @@ dependencies = [ [[package]] name = "icu_properties_data" -version = "2.2.0" +version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8e2bbb201e0c04f7b4b3e14382af113e17ba4f63e2c9d2ee626b720cbce54a14" +checksum = "e590f038c1464a96894fd6d10127e90a8be4509f56ff7ecef851b15cee0b7caa" [[package]] name = "icu_provider" -version = "2.2.0" +version = "2.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "139c4cf31c8b5f33d7e199446eff9c1e02decfc2f0eec2c8d71f65befa45b421" +checksum = "d27bbb9d3abbefac45d55f647c9de1d44aafcd1186eb91879afef17c396c3e73" dependencies = [ "displaydoc", "icu_locale_core", @@ -3076,9 +3077,9 @@ dependencies = [ [[package]] name = "jaq-std" -version = "3.0.1" +version = "3.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8bdc5a74b0feeb5e6a1dc2dd08c34280a61e37668d10a6a3b27ad69d0fb9ce2e" +checksum = "7941c8de9c591052050550f228c62ef80d3ecbd84c330f5c454bc8ebb7a04089" dependencies = [ "bstr", "jaq-core", @@ -3388,9 +3389,9 @@ dependencies = [ [[package]] name = "libgit2-sys" -version = "0.18.7+1.9.6" +version = "0.18.8+1.9.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "23c7391e4b9f4ffab1a624223cc1d7385ff9a678f490768add717de7ea2f4d89" +checksum = "7f7c568b25d7489bc3fb2988ed69ab111d2944d2f5fec3d5c987fe545ea97b50" dependencies = [ "cc", "libc", @@ -3420,9 +3421,9 @@ dependencies = [ [[package]] name = "libredox" -version = "0.1.19" +version = "0.1.20" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2026a5056764a10b2bf5d56488cba40da507f5493a6a429340e2004d9ed085fa" +checksum = "28d0a00925a9f930d679b6789b721e3a7f9ed110f41b86d2497caa780c3a070a" dependencies = [ "libc", ] @@ -3468,9 +3469,9 @@ checksum = "32a66949e030da00e8c7d4434b251670a91556f4144941d37452769c25d58a53" [[package]] name = "litemap" -version = "0.8.2" +version = "0.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "92daf443525c4cce67b150400bc2316076100ce0b3686209eb8cf3c31612e6f0" +checksum = "47d9d19d1d6efa0109d2f65ff4c85cddd50bd572e5a00127ab10987290bcefae" [[package]] name = "lock_api" @@ -3527,9 +3528,9 @@ dependencies = [ [[package]] name = "mail-parser" -version = "0.11.6" +version = "0.11.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4084ec5c2f90b341d0c70990e92a23b128f75ca14fc1dd5edd8fd5c9b417da4d" +checksum = "5d1bb2f9fb98d69b0369be719dc8ab2f535a959e3761128c5c044fae97f68c12" dependencies = [ "hashify", ] @@ -4753,9 +4754,9 @@ dependencies = [ [[package]] name = "pkg-config" -version = "0.3.33" +version = "0.3.34" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "19f132c84eca552bf34cab8ec81f1c1dcc229b811638f9d283dceabe58c5569e" +checksum = "f6b464fbc74e149a392436b17d523f769e057cb6877f6a5c4618bc6f11800548" [[package]] name = "plist" @@ -4856,9 +4857,9 @@ dependencies = [ [[package]] name = "potential_utf" -version = "0.1.5" +version = "0.1.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0103b1cef7ec0cf76490e969665504990193874ea05c85ff9bab8b911d0a0564" +checksum = "d83eb9bc6d8e5cf568e7a1101d60ee05e81ed50ea106026f3d18deeb046d7661" dependencies = [ "zerovec", ] @@ -5021,9 +5022,9 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.16" +version = "0.11.17" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2f4bfc015262b9df63c8845072ce59068853ff5872180c2ce2f13038b970e560" +checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" dependencies = [ "bytes", "getrandom 0.4.3", @@ -5222,18 +5223,18 @@ dependencies = [ [[package]] name = "ref-cast" -version = "1.0.26" +version = "1.0.27" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "216e8f773d7923bcba9ceb86a86c93cabb3903a11872fc3f138c49630e50b96d" +checksum = "7e440fb4e4b4147295338efb76001ab9e4efc0e5839df2c47fc5ac2381d365c3" dependencies = [ "ref-cast-impl", ] [[package]] name = "ref-cast-impl" -version = "1.0.26" +version = "1.0.27" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2c9283685feec7d69af75fb0e858d5e7378f33fe4fc699383b2916ab9273e03c" +checksum = "92ecd8964f8453721699a1ed72037b0db49ce2f5a5138486ee89bed6f67cdf3a" dependencies = [ "proc-macro2", "quote", @@ -5537,9 +5538,9 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.14" +version = "0.103.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0527518605e68109d875e248ea259b6758801cf165e4b2c2733ae3b51f12535a" +checksum = "f3c3cf1d8b1e7d4927e2d154c3fcb02979afb9939629c62cd9048d4f07b60ac2" dependencies = [ "ring", "rustls-pki-types", @@ -6360,9 +6361,9 @@ checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" [[package]] name = "swift-rs" -version = "1.0.7" +version = "1.0.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4057c98e2e852d51fdcfca832aac7b571f6b351ad159f9eda5db1655f8d0c4d7" +checksum = "e45c444e496845d3f2a351146bff59aae4975b2280238df1dfaa0c7d1846f38e" dependencies = [ "base64 0.21.7", "serde", @@ -7408,9 +7409,9 @@ dependencies = [ [[package]] name = "tinystr" -version = "0.8.3" +version = "0.8.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c8323304221c2a851516f22236c5722a72eaa19749016521d6dff0824447d96d" +checksum = "b1e27c91459209c2986af3dcf603a5a74a4368754ce37414f59acc971167f643" dependencies = [ "displaydoc", "zerovec", @@ -8064,9 +8065,9 @@ checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be" [[package]] name = "uuid" -version = "1.24.0" +version = "1.24.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bf3923a6f5c4c6382e0b653c4117f48d631ea17f38ed86e2a828e6f7412f5239" +checksum = "2cefc03fd367c0c6d4305de1b312cf00248c4114f4a0418ce6a6af769e3b0bd9" dependencies = [ "getrandom 0.4.3", "js-sys", @@ -8244,9 +8245,9 @@ dependencies = [ [[package]] name = "wayland-backend" -version = "0.3.16" +version = "0.3.17" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "016ccf01d1c58b6f8999612813e17c9b2390f7d70671428869913310f83f54b8" +checksum = "38a91b4eaddff87b1cd1074985e3713da4af2c49742d1b356b2c01670a67a078" dependencies = [ "cc", "downcast-rs", @@ -8593,6 +8594,19 @@ dependencies = [ "windows-strings 0.4.2", ] +[[package]] +name = "windows-core" +version = "0.62.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8e83a14d34d0623b51dce9581199302a221863196a1dde71a7663a4c2be9deb" +dependencies = [ + "windows-implement 0.60.2", + "windows-interface 0.59.3", + "windows-link 0.2.1", + "windows-result 0.4.1", + "windows-strings 0.5.1", +] + [[package]] name = "windows-future" version = "0.2.1" @@ -8730,6 +8744,15 @@ dependencies = [ "windows-link 0.1.3", ] +[[package]] +name = "windows-result" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7781fa89eaf60850ac3d2da7af8e5242a5ea78d1a11c49bf2910bb5a73853eb5" +dependencies = [ + "windows-link 0.2.1", +] + [[package]] name = "windows-strings" version = "0.1.0" @@ -8749,6 +8772,15 @@ dependencies = [ "windows-link 0.1.3", ] +[[package]] +name = "windows-strings" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7837d08f69c77cf6b07689544538e017c1bfcf57e34b4c0ff58e6c2cd3b37091" +dependencies = [ + "windows-link 0.2.1", +] + [[package]] name = "windows-sys" version = "0.45.0" @@ -9106,9 +9138,9 @@ checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" [[package]] name = "writeable" -version = "0.6.3" +version = "0.6.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1ffae5123b2d3fc086436f8834ae3ab053a283cfac8fe0a0b8eaae044768a4c4" +checksum = "3ad82d2a33cdc9674dc7465672f271e096168fcdbe0f799d9e6db8c5892679dc" [[package]] name = "wry" @@ -9375,9 +9407,9 @@ dependencies = [ [[package]] name = "zerotrie" -version = "0.2.4" +version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0f9152d31db0792fa83f70fb2f83148effb5c1f5b8c7686c3459e361d9bc20bf" +checksum = "4ea269c3bd32f0a32c321907a2ae912ba6f4649bb0fc764a15627e99a7095a3f" dependencies = [ "displaydoc", "yoke", @@ -9386,9 +9418,9 @@ dependencies = [ [[package]] name = "zerovec" -version = "0.11.6" +version = "0.11.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "90f911cbc359ab6af17377d242225f4d75119aec87ea711a880987b18cd7b239" +checksum = "bb0464e17806c1d976d5cba29399c7f08e516e279e2ba493f63123b5fca67dd8" dependencies = [ "yoke", "zerofrom", @@ -9397,13 +9429,13 @@ dependencies = [ [[package]] name = "zerovec-derive" -version = "0.11.3" +version = "0.11.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "625dc425cab0dca6dc3c3319506e6593dcb08a9f387ea3b284dbd52a92c40555" +checksum = "34df6fc39dbd26ddc9c10e6a2984476e13acce22e64e4487636ef494369225da" dependencies = [ "proc-macro2", "quote", - "syn 2.0.119", + "syn 3.0.3", ] [[package]] @@ -9470,9 +9502,9 @@ dependencies = [ [[package]] name = "zvariant" -version = "5.14.0" +version = "5.15.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b5e28c25bd8bb8da5a1f3e7065d0c156b9ee9a7973adf78b0e35eaefdf3b1b5c" +checksum = "c1d34c27cc6cdd1f458427519dd6b8612f7b7e3f7b9a0b2355d041dda9869147" dependencies = [ "endi", "enumflags2", @@ -9486,9 +9518,9 @@ dependencies = [ [[package]] name = "zvariant_derive" -version = "5.14.0" +version = "5.15.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d496a145685283b67e232bd9e47377f6b60ad9d51e3601b23867f77c42477f96" +checksum = "864155e69b4352db0c7f374917bf45d1e0c8d17659c8b3dbf9795f3673f8c497" dependencies = [ "proc-macro-crate 3.5.0", "proc-macro2", @@ -9499,9 +9531,9 @@ dependencies = [ [[package]] name = "zvariant_utils" -version = "4.0.0" +version = "4.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "629d80ece222cad20fe0e8741be493c4ab166acf3b85341bdc2cdbcfd8f3c2d6" +checksum = "bad0294361a320b694a328460dc73add56c306150f5cb6bfafc44446120008a3" dependencies = [ "proc-macro2", "quote", diff --git a/src/core/runtime/services.rs b/src/core/runtime/services.rs index cc540efb73..70d9e0e365 100644 --- a/src/core/runtime/services.rs +++ b/src/core/runtime/services.rs @@ -375,8 +375,8 @@ async fn run_legacy_migrations(config: &Config) { // Idempotent copy of any task boards left in the retired // `{workspace}/agent_task_boards/*.json` file-JSON tree into the crate // `graph.todos` store, which is now authoritative. Idempotent and returns - // fast on an empty/absent legacy dir (the `*.runs.json` ledger stays local). - // As above, each core boot must inspect its own workspace. + // fast on an empty/absent legacy dir. As above, each core boot must inspect + // its own workspace. match crate::openhuman::agent::tinyagents::todos::migrate_legacy_task_boards( &config.workspace_dir, ) @@ -393,6 +393,26 @@ async fn run_legacy_migrations(config: &Config) { Ok(_) => {} Err(e) => log::warn!("[todos] legacy→crate task-board migration failed: {e}"), } + + // The `*.runs.json` claim/heartbeat ledgers that sat beside those boards + // move with them: run records now live in the crate `graph.todos.runs` + // store, so a board and its run log cannot drift apart across a restart. + // Left behind, an in-flight claim would be invisible to the reclaim sweep + // and its card would stay wedged at `in_progress` forever. + match crate::openhuman::threads::todos::runs::migrate_legacy_task_runs(&config.workspace_dir) + .await + { + Ok(report) if report.total > 0 => { + log::info!( + "[todos] legacy→crate run-ledger migration: total={} copied={} skipped={}", + report.total, + report.copied, + report.skipped + ); + } + Ok(_) => {} + Err(e) => log::warn!("[todos] legacy→crate run-ledger migration failed: {e}"), + } } /// Auto-connect Socket.IO to the backend when enabled by the service selection. diff --git a/src/openhuman/agent/task_dispatcher/dispatch.rs b/src/openhuman/agent/task_dispatcher/dispatch.rs index 140d185ef1..1db560bde1 100644 --- a/src/openhuman/agent/task_dispatcher/dispatch.rs +++ b/src/openhuman/agent/task_dispatcher/dispatch.rs @@ -107,7 +107,7 @@ pub async fn dispatch_card( "[task_dispatcher] card claimed (→in_progress), spawning autonomous run" ); - if let Err(e) = runs::create_run(&location, &run_id, &card_id, &executor.label) { + if let Err(e) = runs::create_run(&location, &run_id, &card_id, &executor.label).await { tracing::warn!( run_id = %run_id, card_id = %card_id, @@ -182,8 +182,8 @@ pub async fn dispatch_card( tid, ActiveRun { abort: join.abort_handle(), - hb_cancel: hb_cancel_tx, - location: reg_location, + heartbeat_cancel: hb_cancel_tx, + context: reg_location, card_id: reg_card_id, run_id: reg_run_id, }, diff --git a/src/openhuman/agent/task_dispatcher/executor.rs b/src/openhuman/agent/task_dispatcher/executor.rs index 69a7bed453..fe39cd40d3 100644 --- a/src/openhuman/agent/task_dispatcher/executor.rs +++ b/src/openhuman/agent/task_dispatcher/executor.rs @@ -364,7 +364,8 @@ pub(super) async fn write_back( Vec::new(), ), }; - if let Err(e) = runs::complete_run(location, run_id, run_outcome, run_error, run_evidence) { + if let Err(e) = runs::complete_run(location, run_id, run_outcome, run_error, run_evidence).await + { tracing::warn!( run_id = %run_id, error = %e, diff --git a/src/openhuman/agent/task_dispatcher/poller.rs b/src/openhuman/agent/task_dispatcher/poller.rs index 7aac349b80..7a749a360d 100644 --- a/src/openhuman/agent/task_dispatcher/poller.rs +++ b/src/openhuman/agent/task_dispatcher/poller.rs @@ -7,7 +7,9 @@ use std::sync::OnceLock; use std::time::Duration; -use crate::openhuman::agent::task_board::{TaskApprovalMode, TaskBoardCard, TaskCardStatus}; +use tinyagents::graph::todos::dispatch::select; + +use crate::openhuman::agent::task_board::{TaskApprovalMode, TaskBoardCard}; use crate::openhuman::config::Config; use crate::openhuman::threads::todos::ops::{self, BoardLocation, USER_TASKS_THREAD_ID}; use crate::openhuman::threads::todos::runs::{self, RunLimits}; @@ -30,27 +32,22 @@ const POLLER_MAX_BACKOFF_SECONDS: u64 = 15 * 60; /// immediately slow down. const POLLER_IDLE_GRACE_TICKS: u32 = 2; +/// The backoff curve itself lives in the crate +/// ([`select::PollCadence`](tinyagents::graph::todos::dispatch::select::PollCadence)); +/// this is OpenHuman's tuning of it (issue #4090). +const POLLER_CADENCE: select::PollCadence = select::PollCadence { + base: Duration::from_secs(POLLER_TICK_SECONDS), + max_backoff: Duration::from_secs(POLLER_MAX_BACKOFF_SECONDS), + grace_ticks: POLLER_IDLE_GRACE_TICKS, +}; + static POLLER_STARTED: OnceLock<()> = OnceLock::new(); -/// Compute the next sleep before a poll tick given how many consecutive idle -/// ticks have elapsed (issue #4090). Pure + deterministic so the backoff curve -/// is unit-testable without the real timer. -/// -/// - Fresh work (`idle_ticks == 0`) or within the grace window → base cadence. -/// - Beyond the grace window → exponential backoff (double per extra idle tick) -/// saturating at [`POLLER_MAX_BACKOFF_SECONDS`]. +/// How long to sleep before the next poll tick, given how many consecutive idle +/// ticks have elapsed: base cadence through the grace window, then doubling up +/// to the ceiling. fn next_poll_delay(idle_ticks: u32) -> Duration { - let over = idle_ticks.saturating_sub(POLLER_IDLE_GRACE_TICKS); - if over == 0 { - return Duration::from_secs(POLLER_TICK_SECONDS); - } - // Double per idle tick past the grace window, saturating at the cap. Clamp - // the shift so a long idle streak can't overflow the multiply. - let factor = 1u64.checked_shl(over.min(20)).unwrap_or(u64::MAX); - let secs = POLLER_TICK_SECONDS - .saturating_mul(factor) - .min(POLLER_MAX_BACKOFF_SECONDS); - Duration::from_secs(secs) + POLLER_CADENCE.next_delay(idle_ticks) } /// Spawn the board poller. Idempotent — only the first call installs the loop. @@ -200,11 +197,7 @@ async fn poll_board(location: &BoardLocation, agent_assigned_only: bool) -> Resu // `enforce_single_in_progress` caps the board at one running card, so if // one is already in progress there's nothing for this tick to claim. - if snapshot - .cards - .iter() - .any(|c| c.status == TaskCardStatus::InProgress) - { + if select::has_card_in_progress(&snapshot.cards) { return Ok(false); } @@ -230,61 +223,37 @@ async fn poll_board(location: &BoardLocation, agent_assigned_only: bool) -> Resu /// When `agent_assigned_only` is set, cards without an `assigned_agent` are /// excluded — used on the `user-tasks` board so the poller runs only /// agent-generated tasks and never picks up a human's manually-created card. +/// +/// The selection policy itself is +/// [`select::pick_next_card`](tinyagents::graph::todos::dispatch::select::pick_next_card). pub(super) fn pick_next_todo( cards: &[TaskBoardCard], agent_assigned_only: bool, ) -> Option { - cards - .iter() - .filter(|c| matches!(c.status, TaskCardStatus::Todo | TaskCardStatus::Ready)) - .filter(|c| { - !agent_assigned_only - || c.assigned_agent - .as_deref() - .map(|a| !a.trim().is_empty()) - .unwrap_or(false) - }) - .max_by(|a, b| { - card_urgency(a) - .partial_cmp(&card_urgency(b)) - .unwrap_or(std::cmp::Ordering::Equal) - // On equal urgency, prefer the lower `order` (earlier card): - // reversing the order comparison makes it the "greater" pick. - .then(b.order.cmp(&a.order)) - }) - .cloned() + select::pick_next_card(cards, agent_assigned_only) } /// Whether a card must be parked at `awaiting_approval` before it can run. /// /// Per-card `approval_mode` is authoritative when set; the global /// `require_task_plan_approval` setting is only the fallback for cards with no -/// explicit preference: -/// - `Required` → always park, **even when the global default is off**. The -/// interactive plan-review gate (WebChat turns, see -/// [`crate::openhuman::agent::tools::todo`]) stamps `Required`, and that -/// review must hold regardless of the global switch — otherwise an -/// interactive plan would execute before the user ever sees the review card. -/// - `NotRequired` → never park (already cleared human review, e.g. approved -/// out of the `task-sources` inbox onto `user-tasks`). -/// - unset → fall back to the global default. +/// explicit preference. In particular `Required` parks the card **even when the +/// global default is off**: the interactive plan-review gate (WebChat turns, +/// see [`crate::openhuman::agent::tools::todo`]) stamps `Required`, and that +/// review must hold regardless of the global switch — otherwise an interactive +/// plan would execute before the user ever saw the review card. +/// +/// The rule itself is +/// [`select::requires_plan_approval`](tinyagents::graph::todos::dispatch::select::requires_plan_approval). pub(super) fn requires_plan_approval( global_required: bool, approval_mode: Option<&TaskApprovalMode>, ) -> bool { - match approval_mode { - Some(TaskApprovalMode::Required) => true, - Some(TaskApprovalMode::NotRequired) => false, - None => global_required, - } + select::requires_plan_approval(global_required, approval_mode) } pub(super) fn card_urgency(card: &TaskBoardCard) -> f64 { - card.source_metadata - .as_ref() - .and_then(|m| m.get("urgency")) - .and_then(serde_json::Value::as_f64) - .unwrap_or(0.0) + select::card_urgency(card) } #[cfg(test)] diff --git a/src/openhuman/agent/task_dispatcher/prompt.rs b/src/openhuman/agent/task_dispatcher/prompt.rs index 48d95a36fd..bc3d52af44 100644 --- a/src/openhuman/agent/task_dispatcher/prompt.rs +++ b/src/openhuman/agent/task_dispatcher/prompt.rs @@ -1,121 +1,38 @@ -//! Task prompt construction helpers. +//! Task prompt construction — a thin binding of +//! [`tinyagents::graph::todos::dispatch::prompt`] to OpenHuman's tool names. //! -//! Builds the goal prompt handed to autonomous runs from a [`TaskBoardCard`], -//! and the live-progress instruction that keeps the card current while the -//! run works. +//! The crate owns the rendering (objective, plan, acceptance criteria, source +//! provenance, and the "block rather than guess" progress addendum). All this +//! module supplies is which tools the generated text should point the model at: +//! `memory_recall` for the ingested activity of a card's originating item, and +//! `update_task` for the card write-back. -use crate::openhuman::agent::task_board::TaskBoardCard; - -/// Render a card into the goal prompt handed to the autonomous run. -/// -/// The card's `content`/title is the display form; the prompt leads with the -/// clean `objective`, then any `plan` steps and `acceptance_criteria`, and a -/// pointer to the originating source so the agent can pull related context from -/// memory via its `memory_recall` tool (the GitHub/Notion/… activity for this -/// item is ingested into the summary tree by the memory-sources domain). -pub fn build_task_prompt(card: &TaskBoardCard) -> String { - let mut lines: Vec = Vec::new(); +use std::sync::LazyLock; - let objective = card - .objective - .as_deref() - .map(str::trim) - .filter(|s| !s.is_empty()) - .unwrap_or_else(|| card.title.trim()); - lines.push(format!( - "You are autonomously executing one task to completion. Objective:\n{objective}" - )); +use tinyagents::graph::todos::dispatch::prompt as crate_prompt; +use tinyagents::graph::todos::dispatch::TaskPromptTools; - if !card.plan.is_empty() { - lines.push("\nPlan:".to_string()); - for (i, step) in card.plan.iter().enumerate() { - lines.push(format!("{}. {}", i + 1, step.trim())); - } - } - - if !card.acceptance_criteria.is_empty() { - lines.push("\nAcceptance criteria (the task is done only when all hold):".to_string()); - for c in &card.acceptance_criteria { - lines.push(format!("- {}", c.trim())); - } - } - - if let Some(meta) = &card.source_metadata { - let provider = meta.get("provider").and_then(|v| v.as_str()); - let repo = meta.get("repo").and_then(|v| v.as_str()); - let external_id = meta.get("external_id").and_then(|v| v.as_str()); - let url = meta.get("url").and_then(|v| v.as_str()); - let mut origin = String::new(); - if let Some(p) = provider { - origin.push_str(p); - } - if let Some(r) = repo { - origin.push_str(&format!(" {r}")); - } - if let Some(id) = external_id { - origin.push_str(&format!("#{id}")); - } - // Gate on a known provider so the origin string is always meaningful - // (an id-only card would render "#123" with a leading space). - if provider.is_some() { - lines.push(format!( - "\nThis task originates from {}. Its activity has been ingested into memory — use \ - your memory_recall tool to pull related context (prior discussion, linked items) \ - before and while you work.", - origin.trim() - )); - } - if let Some(u) = url { - lines.push(format!("Source link: {u}")); - } - // G9b — agent-driven external write-back. When the upstream item is - // addressable (provider + id), instruct the agent to close the loop on - // the source itself via its integration tools. Runs under the - // connection's existing write scope (no extra approval gate); if it - // can't, it reports that instead of failing. - if provider.is_some() && external_id.is_some() { - lines.push(format!( - "\nWhen the task is complete, record the outcome on the upstream source ({}): use \ - your integration tools to add a comment summarising the resolution and, if the \ - work fully addresses it, close/resolve the item. If you lack the permission or \ - connection to do so, say so in your final summary instead of guessing.", - origin.trim() - )); - } - } +use crate::openhuman::agent::task_board::TaskBoardCard; - lines.push( - "\nWork the task to completion. Do not pick up unrelated work. When finished, your final \ - message should summarise what you did and the evidence (commits, PRs, results)." - .to_string(), - ); +/// OpenHuman's tool names for the two tools a task prompt references. +static TOOLS: LazyLock = LazyLock::new(|| TaskPromptTools { + memory_recall: Some("memory_recall".to_string()), + update_task: "update_task".to_string(), +}); - lines.join("\n") +/// Render a card into the goal prompt handed to the autonomous run. +pub fn build_task_prompt(card: &TaskBoardCard) -> String { + crate_prompt::build_task_prompt(card, &TOOLS) } /// Instruction appended to the run prompt so the autonomous turn keeps its own /// task card current via the `update_task` tool while it works. /// -/// The card is already `in_progress` (the dispatcher claimed it before -/// spawning the run), addressed by the exact card id + board the run owns -/// (without the explicit `threadId` the tool defaults to the `task-sources` -/// board and would miss a `user-tasks` card). Two things this asks for: -/// 1. *progress* updates (notes/evidence) as the run works, and -/// 2. an explicit `status: blocked` + `blocker` when the run needs a -/// decision/information from the user or cannot proceed — which -/// [`write_back`] now preserves rather than force-completing, so the task -/// pauses for the user instead of being silently marked done. +/// The card is addressed by exact id **and** board: without the explicit +/// `threadId` the tool defaults to the `task-sources` board and would miss a +/// `user-tasks` card. A run that blocks itself is preserved as blocked by +/// [`write_back`](super::executor::write_back) rather than force-completed, so +/// the task pauses for the user instead of being silently marked done. pub(super) fn build_progress_instruction(card_id: &str, thread_id: &str) -> String { - format!( - "\n\nThis task is tracked as card `{card_id}` on the `{thread_id}` board. As you work, \ - call the `update_task` tool (id `{card_id}`, threadId `{thread_id}`) to keep the card \ - current — append `notes`/`evidence` as you make progress.\n\nIf you need a decision or \ - information from the user, or you genuinely cannot proceed (missing access, ambiguous \ - requirement, an action that needs the user's confirmation), call `update_task` with \ - `status: blocked` and a `blocker` that states exactly what you need from the user. The \ - task will stay paused in that blocked state until the user responds — do NOT guess, \ - fabricate, or take a risky irreversible action just to avoid blocking. If instead you \ - finish the work, end with a summary of what you did and the evidence; completion is \ - recorded automatically." - ) + crate_prompt::build_progress_instruction(card_id, thread_id, &TOOLS) } diff --git a/src/openhuman/agent/task_dispatcher/registry.rs b/src/openhuman/agent/task_dispatcher/registry.rs index 855209c559..5bf2e560a2 100644 --- a/src/openhuman/agent/task_dispatcher/registry.rs +++ b/src/openhuman/agent/task_dispatcher/registry.rs @@ -3,33 +3,35 @@ //! Tracks active runs by session `thread_id` so the web-channel cancel path //! can abort them even though they are detached tokio tasks rather than //! web-channel turns. +//! +//! The map itself is +//! [`ActiveRunRegistry`](tinyagents::graph::todos::dispatch::ActiveRunRegistry), +//! which owns the race-free removal that decides who writes a run's terminal +//! card state. What stays here is the OpenHuman side of a cancel: the board +//! write-back and the terminal chat event. -use std::collections::HashMap; -use std::sync::{Mutex, OnceLock}; +use std::sync::OnceLock; -use super::types::ActiveRun; +use tinyagents::graph::todos::dispatch::ActiveRunRegistry; -static ACTIVE_RUNS: OnceLock>> = OnceLock::new(); +use crate::openhuman::threads::todos::ops::BoardLocation; -pub(super) fn active_runs() -> &'static Mutex> { - ACTIVE_RUNS.get_or_init(|| Mutex::new(HashMap::new())) +use super::types::ActiveRun; + +fn registry() -> &'static ActiveRunRegistry { + static ACTIVE_RUNS: OnceLock> = OnceLock::new(); + ACTIVE_RUNS.get_or_init(ActiveRunRegistry::new) } pub(super) fn register_active_run(thread_id: String, run: ActiveRun) { - active_runs() - .lock() - .expect("active_runs mutex poisoned") - .insert(thread_id, run); + registry().register(thread_id, run); } /// Remove and return the active-run entry for `thread_id`. The naturally /// completing run and a concurrent [`cancel_session`] race on this — whoever /// gets `Some` "owns" the terminal board write-back, so it happens exactly once. pub(super) fn take_active_run(thread_id: &str) -> Option { - active_runs() - .lock() - .expect("active_runs mutex poisoned") - .remove(thread_id) + registry().take(thread_id) } /// Atomically remove the active-run entry for `thread_id`, but only when it @@ -40,36 +42,11 @@ pub(super) fn take_active_run(thread_id: &str) -> Option { /// by a newer run before removal — the "stale cancel kills a newer turn" race a /// separate peek-then-`take_active_run` would reopen (#4760). A `None` /// `request_id` removes whatever run is on the thread (unscoped Stop / -/// teardown). -/// -/// Both scoped no-op cases (no active run, or a `run_id` mismatch from a -/// superseded/unrelated request) emit grep-friendly `debug` diagnostics so an +/// teardown). Both scoped no-op cases (no active run, or a `run_id` mismatch +/// from a superseded/unrelated request) are logged by the crate registry, so an /// intentional no-op cancel is still traceable. pub(super) fn take_active_run_if(thread_id: &str, request_id: Option<&str>) -> Option { - let mut guard = active_runs().lock().expect("active_runs mutex poisoned"); - if let Some(rid) = request_id { - match guard.get(thread_id) { - None => { - tracing::debug!( - thread_id = %thread_id, - request_id = %rid, - "[task_dispatcher] scoped cancel ignored: no active run on thread" - ); - return None; - } - Some(run) if run.run_id != rid => { - tracing::debug!( - thread_id = %thread_id, - request_id = %rid, - active_run_id = %run.run_id, - "[task_dispatcher] scoped cancel ignored: run_id mismatch (superseded/unrelated request)" - ); - return None; - } - _ => {} - } - } - guard.remove(thread_id) + registry().take_if(thread_id, request_id) } /// Cancel the in-flight autonomous run streaming into session `thread_id`. @@ -96,12 +73,11 @@ pub async fn cancel_session(thread_id: &str) -> bool { /// exact run it atomically removed via [`take_active_run_if`], rather than /// re-acquiring the lock and racing a replacement run (#4760). async fn cancel_taken_run(thread_id: &str, run: ActiveRun) { - run.abort.abort(); - let _ = run.hb_cancel.send(true); + run.cancel(); // The aborted task never reaches its own write-back — do it here so the // card lands in a terminal state instead of a stale `in_progress`. super::executor::write_back( - &run.location, + &run.context, &run.card_id, &run.run_id, Err("Cancelled by user".to_string()), diff --git a/src/openhuman/agent/task_dispatcher/tests.rs b/src/openhuman/agent/task_dispatcher/tests.rs index 2f70f8211f..accc2717f0 100644 --- a/src/openhuman/agent/task_dispatcher/tests.rs +++ b/src/openhuman/agent/task_dispatcher/tests.rs @@ -22,8 +22,8 @@ async fn active_run_registry_take_is_once() { key.to_string(), ActiveRun { abort: handle.abort_handle(), - hb_cancel: tx, - location: BoardLocation::Scratch, + heartbeat_cancel: tx, + context: BoardLocation::Scratch, card_id: "c1".to_string(), run_id: "r1".to_string(), }, @@ -55,8 +55,8 @@ async fn cancel_session_scoped_ignores_a_mismatched_request() { key.to_string(), ActiveRun { abort: handle.abort_handle(), - hb_cancel: tx, - location: BoardLocation::Scratch, + heartbeat_cancel: tx, + context: BoardLocation::Scratch, card_id: "c1".to_string(), run_id: "r1".to_string(), }, @@ -95,8 +95,8 @@ async fn cancel_session_scoped_aborts_the_run_when_the_request_matches() { key.to_string(), ActiveRun { abort: handle.abort_handle(), - hb_cancel: tx, - location: loc.clone(), + heartbeat_cancel: tx, + context: loc.clone(), card_id: id.clone(), run_id: "r1".to_string(), }, diff --git a/src/openhuman/agent/task_dispatcher/types.rs b/src/openhuman/agent/task_dispatcher/types.rs index 56464203e0..754730703a 100644 --- a/src/openhuman/agent/task_dispatcher/types.rs +++ b/src/openhuman/agent/task_dispatcher/types.rs @@ -7,15 +7,14 @@ use crate::openhuman::threads::todos::ops::BoardLocation; /// Autonomous runs are detached `tokio` tasks, not web-channel turns, so they /// are invisible to the web channel's own in-flight registry — which is why the /// chat **Cancel** button (which calls `channel_web_cancel`) couldn't stop them. -/// Registering the run's [`AbortHandle`](tokio::task::AbortHandle) here lets -/// [`cancel_session`] abort it from that same cancel path. -pub(super) struct ActiveRun { - pub(super) abort: tokio::task::AbortHandle, - pub(super) hb_cancel: tokio::sync::watch::Sender, - pub(super) location: BoardLocation, - pub(super) card_id: String, - pub(super) run_id: String, -} +/// Registering the run's [`AbortHandle`](tokio::task::AbortHandle) in the +/// crate's [`ActiveRunRegistry`](tinyagents::graph::todos::dispatch::ActiveRunRegistry) +/// lets [`cancel_session`](super::registry::cancel_session) abort it from that +/// same cancel path. +/// +/// The `context` the crate carries for us is the run's [`BoardLocation`], which +/// the canceller needs to write the card back to a terminal state. +pub(super) type ActiveRun = tinyagents::graph::todos::dispatch::ActiveRun; /// A resolved executor: which built-in agent definition to build, an optional /// system-prompt suffix carrying a personality identity or skill guidelines, diff --git a/src/openhuman/threads/goals/runtime.rs b/src/openhuman/threads/goals/runtime.rs index 36d7a1e5ae..56330af65d 100644 --- a/src/openhuman/threads/goals/runtime.rs +++ b/src/openhuman/threads/goals/runtime.rs @@ -21,6 +21,10 @@ use std::path::{Path, PathBuf}; use async_trait::async_trait; +use tinyagents::graph::goals::budget as crate_budget; +use tinyagents::graph::goals::{BudgetVerdict, GoalBudgetGuard}; + +use super::migration::goals_store; use super::store; use super::{ThreadGoal, ThreadGoalStatus}; use crate::core::bus::BUS; @@ -88,7 +92,7 @@ pub async fn pause_for_current_thread(workspace_dir: &Path) { /// The per-turn token total used for budget accounting (prompt + completion). fn turn_tokens(input: u64, output: u64) -> u64 { - input.saturating_add(output) + crate_budget::turn_tokens(input, output) } /// Whether the current turn is an autonomous goal-continuation (vs. a @@ -109,55 +113,43 @@ fn is_goal_continuation_turn() -> bool { /// Account a finished turn's usage against the ambient thread's goal. /// -/// Only **active** goals are charged (a paused/complete/budget-limited goal -/// doesn't accrue usage from incidental chat). Best-effort: a failure is logged -/// and swallowed so accounting never fails a user turn. Emits -/// `ThreadGoalUpdated` when the status changes (e.g. → `budget_limited`) so the -/// UI chip refreshes. +/// The accounting rules are the crate's +/// ([`crate_budget::account_turn`](tinyagents::graph::goals::account_turn)): +/// only **active** goals are charged, so a paused/complete/budget-limited goal +/// doesn't accrue usage from incidental chat, and a user-initiated turn clears +/// the one-shot continuation suppression (a continuation turn must not clear +/// its own, see [`super::continuation`]). +/// +/// What is OpenHuman's here: reading the ambient thread from the turn scope, +/// classifying the turn as user-initiated vs. continuation from its origin, and +/// emitting `ThreadGoalUpdated` when the status changes (e.g. → +/// `budget_limited`) so the UI chip refreshes. Best-effort throughout: a +/// failure is logged and swallowed so accounting never fails a user turn. pub async fn account_turn_against_goal(workspace_dir: &Path, input: u64, output: u64, secs: u64) { let Some(thread_id) = current_thread_id() else { return; }; - let goal = match store::get(workspace_dir, &thread_id).await { - Ok(Some(g)) => g, + let prev_status = match store::get(workspace_dir, &thread_id).await { + Ok(Some(goal)) => goal.status, Ok(None) => return, Err(e) => { tracing::debug!(thread_id = %thread_id, error = %e, "[thread_goals] account get failed"); return; } }; - if !goal.status.is_active() { - return; - } - // Reset the one-shot continuation suppression on user-initiated activity: a - // real turn in this thread means the user re-engaged, so a future idle - // period may auto-continue again. The continuation turn itself runs under a - // GoalContinuation origin and must NOT clear its own suppression. - if goal.continuation_suppressed && !is_goal_continuation_turn() { - if let Err(e) = - store::set_continuation_suppressed_if(workspace_dir, &thread_id, &goal.goal_id, false) - .await - { - tracing::debug!( - thread_id = %thread_id, - error = %e, - "[thread_goals] failed to clear continuation suppression" - ); - } - } - let delta = turn_tokens(input, output); - if delta == 0 && secs == 0 { - return; - } - let prev_status = goal.status; - match store::account_usage(workspace_dir, &thread_id, &goal.goal_id, delta, secs).await { + + let store = goals_store(workspace_dir); + let user_initiated = !is_goal_continuation_turn(); + match crate_budget::account_turn(&store, &thread_id, input, output, secs, user_initiated).await + { Ok(Some(updated)) => { tracing::debug!( thread_id = %thread_id, goal_id = %updated.goal_id, tokens_used = updated.tokens_used, status = updated.status.as_str(), - "[thread_goals] accounted turn usage (+{delta} tok, +{secs}s)" + "[thread_goals] accounted turn usage (+{} tok, +{secs}s)", + turn_tokens(input, output) ); if updated.status != prev_status { BUS.publish(DomainEvent::ThreadGoalUpdated { @@ -169,7 +161,7 @@ pub async fn account_turn_against_goal(workspace_dir: &Path, input: u64, output: } Ok(None) => {} Err(e) => { - tracing::debug!(thread_id = %thread_id, error = %e, "[thread_goals] account_usage failed"); + tracing::debug!(thread_id = %thread_id, error = %e, "[thread_goals] account_turn failed"); } } } @@ -178,32 +170,34 @@ pub async fn account_turn_against_goal(workspace_dir: &Path, input: u64, output: /// running usage (already-accounted tokens from prior turns + this turn's /// tokens so far) would meet or exceed its budget. /// -/// It only fires for goals that are still `Active` with a configured budget — +/// The decision is the crate's +/// [`GoalBudgetGuard`](tinyagents::graph::goals::GoalBudgetGuard); this is the +/// adapter that votes it into OpenHuman's [`StopHook`] chain. #4469 item 1: the +/// stop is a graceful *pause*, not an instantaneous abort — the vote fires in +/// the stop-hook middleware's `after_model`, and the harness drains the pause +/// at the **top of the next iteration**, so the tool round for the model call +/// that tripped the budget still runs and the turn's wrap-up summary may spend +/// one more model call before the partial transcript is returned. It bounds an +/// autonomous run to a small, deterministic overshoot past the ceiling rather +/// than a hard cut at the exact accounting point. +/// +/// The guard only arms for a goal that is `Active` with a configured budget, +/// and stands down if that goal is completed, replaced, or paused mid-turn — /// once a goal is `budget_limited`/`paused`/`complete` the user can still chat -/// freely (the injected context steers the model to summarise), so we never -/// hard-stop a user-present turn that isn't actively burning a live budget. +/// freely (the injected context steers the model to summarise), so a +/// user-present turn is never hard-stopped by a budget that is no longer live. #[derive(Debug, Clone)] pub struct GoalBudgetStopHook { workspace_dir: PathBuf, - thread_id: String, - /// The goal version this hook was armed for. Stops enforcing if the goal is - /// replaced mid-turn (a new objective mints a new id). - goal_id: String, - budget: u64, + guard: GoalBudgetGuard, } impl GoalBudgetStopHook { /// Build a hook for `goal` if it's active and has a budget; `None` otherwise. pub fn for_goal(workspace_dir: &Path, goal: &ThreadGoal) -> Option { - if !goal.status.is_active() { - return None; - } - let budget = goal.token_budget?; Some(Self { workspace_dir: workspace_dir.to_path_buf(), - thread_id: goal.thread_id.clone(), - goal_id: goal.goal_id.clone(), - budget, + guard: GoalBudgetGuard::for_goal(goal)?, }) } } @@ -215,27 +209,16 @@ impl StopHook for GoalBudgetStopHook { } async fn check(&self, ctx: &TurnState<'_>) -> StopDecision { - // Read the goal's already-accounted usage (prior turns). If it's gone, - // replaced, or no longer active, stop enforcing. - let goal = match store::get(&self.workspace_dir, &self.thread_id).await { - Ok(Some(g)) => g, - _ => return StopDecision::Continue, - }; - if goal.goal_id != self.goal_id || !goal.status.is_active() { - return StopDecision::Continue; - } - let projected = goal - .tokens_used - .saturating_add(turn_tokens(ctx.cost.input_tokens, ctx.cost.output_tokens)); - if projected >= self.budget { - StopDecision::Stop { - reason: format!( - "thread goal budget reached: {projected} tokens >= {} budget — stopping to summarise progress", - self.budget - ), + let store = goals_store(&self.workspace_dir); + let in_flight = turn_tokens(ctx.cost.input_tokens, ctx.cost.output_tokens); + match self.guard.check(&store, in_flight).await { + Ok(BudgetVerdict::Stop { reason }) => StopDecision::Stop { reason }, + Ok(BudgetVerdict::Continue) => StopDecision::Continue, + Err(e) => { + // An unreadable goal is not grounds for killing a live turn. + tracing::debug!(error = %e, "[thread_goals] budget check failed; continuing"); + StopDecision::Continue } - } else { - StopDecision::Continue } } } diff --git a/src/openhuman/threads/todos/README.md b/src/openhuman/threads/todos/README.md index 8e4b332cca..466e0043ac 100644 --- a/src/openhuman/threads/todos/README.md +++ b/src/openhuman/threads/todos/README.md @@ -14,10 +14,15 @@ OpenHuman keeps this module to preserve app-specific integration: `AgentProgress::TaskBoardUpdated`. - `schemas.rs` preserves the `openhuman.todos_*` JSON-RPC API. - `tools.rs` preserves the granular `todo_*` agent tools. -- `runs.rs` owns the OpenHuman autonomous-run ledger, which is separate from - task-board storage. +- `runs.rs` binds the TinyAgents autonomous-run ledger + (`tinyagents::graph::todos::runs`) to `BoardLocation` addressing, renders its + timestamps as RFC 3339 for the wire, publishes `TaskRunReclaimed`, and imports + the retired `agent_task_boards/.runs.json` ledgers. The run record, + heartbeat, staleness policy, and reclaim sweep are the crate's. `agent::task_board` re-exports the TinyAgents board types and keeps the legacy `TaskBoardStore` facade for existing callers. Legacy -`agent_task_boards/*.json` values are imported at startup through -`openhuman::agent::tinyagents::todos`; existing TinyAgents values are never replaced. +`agent_task_boards/*.json` boards are imported at startup through +`openhuman::agent::tinyagents::todos`, and the `*.runs.json` ledgers beside them +through `threads::todos::runs::migrate_legacy_task_runs`; existing TinyAgents +values are never replaced. diff --git a/src/openhuman/threads/todos/ops.rs b/src/openhuman/threads/todos/ops.rs index cb41da2eee..4139aa8e6a 100644 --- a/src/openhuman/threads/todos/ops.rs +++ b/src/openhuman/threads/todos/ops.rs @@ -50,7 +50,7 @@ impl BoardLocation { } } -fn target(location: &BoardLocation) -> (Arc, &str) { +pub(super) fn target(location: &BoardLocation) -> (Arc, &str) { match location { BoardLocation::Thread { workspace_dir, diff --git a/src/openhuman/threads/todos/runs.rs b/src/openhuman/threads/todos/runs.rs index ec94362cbf..6ec8a374d9 100644 --- a/src/openhuman/threads/todos/runs.rs +++ b/src/openhuman/threads/todos/runs.rs @@ -1,508 +1,253 @@ -//! Durable task-run records with heartbeat liveness and stale reclaim. +//! Compatibility facade over [`tinyagents::graph::todos::runs`]. //! -//! Each time the [`crate::openhuman::agent::task_dispatcher`] claims a card, -//! it creates a [`TaskRun`] that tracks: who claimed it, when, last heartbeat, -//! completion outcome, and error/evidence. A background heartbeat timer ticks -//! alongside the autonomous run so healthy long-running workers stay live while -//! wedged workers can be detected and reclaimed. +//! TinyAgents owns the durable task-run record, the heartbeat, the staleness +//! policy, and the reclaim sweep (card back to `todo`, or parked at `blocked` +//! once a card has burned through its reclaim budget). What stays here is +//! OpenHuman's own shape around it: [`BoardLocation`] addressing (including the +//! process-global scratch board), RFC 3339 timestamps on the wire, the +//! `TaskRunReclaimed` domain event, and the one-time import of the retired +//! `{workspace}/agent_task_boards/.runs.json` ledger. //! -//! Stale reclaim policy: a run whose heartbeat is older than -//! [`RunLimits::heartbeat_stale_secs`] **or** whose total age exceeds -//! [`RunLimits::claim_ttl_secs`] is eligible for reclaim. Reclaimed cards move -//! back to `todo` (re-dispatchable) unless they've been reclaimed more than -//! [`RunLimits::max_reclaim_count`] times, in which case they park as `blocked` -//! with a diagnostic blocker message. - -use chrono::{DateTime, Utc}; -use parking_lot::Mutex; -use serde::{Deserialize, Serialize}; -use std::collections::HashMap; -use std::fs; -use std::io::{Read, Write}; -use std::path::{Path, PathBuf}; -use std::sync::{Arc, OnceLock}; - -use crate::openhuman::agent::task_board::TaskCardStatus; - -use super::ops::{self, BoardLocation, CardPatch}; - -// ── Defaults ─────────────────────────────────────────────────────────── - -pub const DEFAULT_HEARTBEAT_STALE_SECS: u64 = 300; -pub const DEFAULT_CLAIM_TTL_SECS: u64 = 3600; -pub const DEFAULT_MAX_RECLAIM_COUNT: u32 = 3; -const HEARTBEAT_TICK_SECS: u64 = 30; - -// ── Types ────────────────────────────────────────────────────────────── - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -pub enum RunOutcome { - Success, - Failed, - Reclaimed, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct TaskRun { - pub run_id: String, - pub card_id: String, - pub claimed_by: String, - pub claim_token: String, - pub started_at: String, - pub last_heartbeat_at: String, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub completed_at: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub outcome: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - pub error: Option, - #[serde(default, skip_serializing_if = "Vec::is_empty")] - pub evidence: Vec, -} - -impl TaskRun { - pub fn is_active(&self) -> bool { - self.completed_at.is_none() - } -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct RunLimits { - pub heartbeat_stale_secs: u64, - pub claim_ttl_secs: u64, - pub max_reclaim_count: u32, -} +//! Run records live in the crate KV store beside the board itself +//! (`graph.todos.runs`), so a board and its run log can no longer drift apart +//! across a restart. -impl Default for RunLimits { - fn default() -> Self { - Self { - heartbeat_stale_secs: DEFAULT_HEARTBEAT_STALE_SECS, - claim_ttl_secs: DEFAULT_CLAIM_TTL_SECS, - max_reclaim_count: DEFAULT_MAX_RECLAIM_COUNT, - } - } -} +use std::path::Path; -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct ReclaimResult { - pub reclaimed_count: usize, - pub blocked_count: usize, - pub details: Vec, -} +use serde::{Deserialize, Serialize}; +use tinyagents::graph::todos::runs as crate_runs; -#[derive(Debug, Clone, Serialize, Deserialize)] -#[serde(rename_all = "camelCase")] -pub struct ReclaimDetail { - pub run_id: String, - pub card_id: String, - pub reason: String, - pub new_card_status: String, -} +pub use tinyagents::graph::todos::runs::{ + ReclaimDetail, ReclaimResult, RunLimits, RunOutcome, TaskRun, DEFAULT_CLAIM_TTL_SECS, + DEFAULT_HEARTBEAT_STALE_SECS, DEFAULT_MAX_RECLAIM_COUNT, +}; -// ── Per-board lock for run records ───────────────────────────────────── +use crate::openhuman::agent::task_board::normalize_timestamp_for_wire; -fn run_lock(location: &BoardLocation) -> Arc> { - static MAP: OnceLock>>>> = OnceLock::new(); - let map_mu = MAP.get_or_init(|| Mutex::new(HashMap::new())); - let key = match location { - BoardLocation::Thread { thread_id, .. } => format!("runs:{thread_id}"), - BoardLocation::Scratch => "runs:_scratch_".to_string(), - }; - map_mu.lock().entry(key).or_default().clone() -} +use super::ops::{target, BoardLocation}; -// ── Store ────────────────────────────────────────────────────────────── +/// Cadence of the background heartbeat spawned alongside an autonomous run. +const HEARTBEAT_TICK: std::time::Duration = crate_runs::DEFAULT_HEARTBEAT_TICK; +/// Legacy on-disk ledger the crate store replaced. const TASK_BOARD_DIR: &str = "agent_task_boards"; -fn runs_path(workspace_dir: &Path, thread_id: &str) -> PathBuf { - workspace_dir - .join(TASK_BOARD_DIR) - .join(format!("{}.runs.json", hex::encode(thread_id.as_bytes()))) -} - -fn load_runs(location: &BoardLocation) -> Result, String> { - let BoardLocation::Thread { - workspace_dir, - thread_id, - } = location - else { - return Ok(Vec::new()); - }; - let path = runs_path(workspace_dir, thread_id); - if !path.exists() { - return Ok(Vec::new()); - } - let mut buf = String::new(); - fs::File::open(&path) - .map_err(|e| format!("open runs {}: {e}", path.display()))? - .read_to_string(&mut buf) - .map_err(|e| format!("read runs {}: {e}", path.display()))?; - serde_json::from_str::>(&buf) - .map_err(|e| format!("parse runs {}: {e}", path.display())) +fn map_err(result: tinyagents::error::Result) -> Result { + result.map_err(|error| error.to_string()) } -fn save_runs(location: &BoardLocation, runs: &[TaskRun]) -> Result<(), String> { - let BoardLocation::Thread { - workspace_dir, - thread_id, - } = location - else { - return Ok(()); - }; - let dir = workspace_dir.join(TASK_BOARD_DIR); - fs::create_dir_all(&dir).map_err(|e| format!("create runs dir {}: {e}", dir.display()))?; - let path = runs_path(workspace_dir, thread_id); - let bytes = serde_json::to_vec_pretty(&runs).map_err(|e| format!("serialize runs: {e}"))?; - let mut tmp = - tempfile::NamedTempFile::new_in(&dir).map_err(|e| format!("create runs tempfile: {e}"))?; - tmp.write_all(&bytes) - .map_err(|e| format!("write runs tempfile: {e}"))?; - tmp.as_file() - .sync_all() - .map_err(|e| format!("fsync runs tempfile: {e}"))?; - tmp.persist(&path) - .map_err(|e| format!("persist runs {}: {e}", path.display()))?; - Ok(()) +/// Crate stamps are unix-epoch milliseconds; the `openhuman.todos_run_*` RPC +/// surface has always spoken RFC 3339, so translate on the way out. +fn for_wire(mut run: TaskRun) -> TaskRun { + run.started_at = normalize_timestamp_for_wire(&run.started_at); + run.last_heartbeat_at = normalize_timestamp_for_wire(&run.last_heartbeat_at); + run.completed_at = run + .completed_at + .as_deref() + .map(normalize_timestamp_for_wire); + run } -// ── Operations ───────────────────────────────────────────────────────── - -pub fn create_run( +pub async fn create_run( location: &BoardLocation, run_id: &str, card_id: &str, claimed_by: &str, ) -> Result { - let lock = run_lock(location); - let _guard = lock.lock(); - - let now = Utc::now().to_rfc3339(); - let claim_token = uuid::Uuid::new_v4().to_string(); - - tracing::debug!( - run_id = %run_id, - card_id = %card_id, - claimed_by = %claimed_by, - "[todos][runs] create_run entry" - ); - - let run = TaskRun { - run_id: run_id.to_string(), - card_id: card_id.to_string(), - claimed_by: claimed_by.to_string(), - claim_token: claim_token.clone(), - started_at: now.clone(), - last_heartbeat_at: now, - completed_at: None, - outcome: None, - error: None, - evidence: Vec::new(), - }; - - let mut runs = load_runs(location)?; - runs.push(run.clone()); - save_runs(location, &runs)?; - - tracing::info!( - run_id = %run_id, - card_id = %card_id, - claim_token = %claim_token, - "[todos][runs] create_run ok" - ); - Ok(run) + let (store, thread_id) = target(location); + let run = map_err( + crate_runs::create_run(&store, thread_id, Some(run_id), card_id, claimed_by).await, + )?; + Ok(for_wire(run)) } -pub fn update_heartbeat(location: &BoardLocation, run_id: &str) -> Result<(), String> { - let lock = run_lock(location); - let _guard = lock.lock(); - - let mut runs = load_runs(location)?; - let run = runs - .iter_mut() - .find(|r| r.run_id == run_id && r.is_active()) - .ok_or_else(|| format!("[todos][runs] active run '{run_id}' not found for heartbeat"))?; - - run.last_heartbeat_at = Utc::now().to_rfc3339(); - save_runs(location, &runs)?; - - tracing::trace!( - run_id = %run_id, - "[todos][runs] heartbeat updated" - ); - Ok(()) +pub async fn update_heartbeat(location: &BoardLocation, run_id: &str) -> Result<(), String> { + let (store, thread_id) = target(location); + map_err(crate_runs::update_heartbeat(&store, thread_id, run_id).await) } -pub fn complete_run( +pub async fn complete_run( location: &BoardLocation, run_id: &str, outcome: RunOutcome, error: Option, evidence: Vec, ) -> Result { - let lock = run_lock(location); - let _guard = lock.lock(); - - tracing::debug!( - run_id = %run_id, - outcome = ?outcome, - "[todos][runs] complete_run entry" - ); - - let mut runs = load_runs(location)?; - let run = runs - .iter_mut() - .find(|r| r.run_id == run_id && r.is_active()) - .ok_or_else(|| format!("[todos][runs] active run '{run_id}' not found for completion"))?; - - run.completed_at = Some(Utc::now().to_rfc3339()); - run.outcome = Some(outcome); - run.error = error; - run.evidence = evidence; - let completed = run.clone(); - - save_runs(location, &runs)?; - - tracing::info!( - run_id = %run_id, - outcome = ?completed.outcome, - "[todos][runs] complete_run ok" - ); - Ok(completed) + let (store, thread_id) = target(location); + let run = map_err( + crate_runs::complete_run(&store, thread_id, run_id, outcome, error, evidence).await, + )?; + Ok(for_wire(run)) } -pub fn list_runs(location: &BoardLocation, card_id: Option<&str>) -> Result, String> { - let lock = run_lock(location); - let _guard = lock.lock(); - - let runs = load_runs(location)?; - Ok(match card_id { - Some(cid) => runs.into_iter().filter(|r| r.card_id == cid).collect(), - None => runs, - }) +pub async fn list_runs( + location: &BoardLocation, + card_id: Option<&str>, +) -> Result, String> { + let (store, thread_id) = target(location); + let runs = map_err(crate_runs::list_runs(&store, thread_id, card_id).await)?; + Ok(runs.into_iter().map(for_wire).collect()) } -pub fn get_run(location: &BoardLocation, run_id: &str) -> Result, String> { - let lock = run_lock(location); - let _guard = lock.lock(); - - let runs = load_runs(location)?; - Ok(runs.into_iter().find(|r| r.run_id == run_id)) +pub async fn get_run(location: &BoardLocation, run_id: &str) -> Result, String> { + let (store, thread_id) = target(location); + let run = map_err(crate_runs::get_run(&store, thread_id, run_id).await)?; + Ok(run.map(for_wire)) } -pub fn find_stale_runs( +pub async fn find_stale_runs( location: &BoardLocation, limits: &RunLimits, ) -> Result, String> { - let lock = run_lock(location); - let _guard = lock.lock(); - - let runs = load_runs(location)?; - let now = Utc::now(); - let mut stale = Vec::new(); - - for run in &runs { - if !run.is_active() { - continue; - } - if let Some(reason) = check_staleness(run, &now, limits) { - stale.push((run.clone(), reason)); - } - } - Ok(stale) -} - -fn check_staleness(run: &TaskRun, now: &DateTime, limits: &RunLimits) -> Option { - let started: DateTime = run.started_at.parse().ok()?; - let last_hb: DateTime = run.last_heartbeat_at.parse().ok()?; - - let age_secs = (*now - started).num_seconds().max(0) as u64; - let hb_age_secs = (*now - last_hb).num_seconds().max(0) as u64; - - if age_secs > limits.claim_ttl_secs { - return Some(format!( - "claim TTL expired (age {age_secs}s > limit {}s)", - limits.claim_ttl_secs - )); - } - if hb_age_secs > limits.heartbeat_stale_secs { - return Some(format!( - "heartbeat stale (last heartbeat {hb_age_secs}s ago > limit {}s)", - limits.heartbeat_stale_secs - )); - } - None + let (store, thread_id) = target(location); + let stale = map_err(crate_runs::find_stale_runs(&store, thread_id, limits).await)?; + Ok(stale + .into_iter() + .map(|(run, reason)| (for_wire(run), reason)) + .collect()) } -/// Reclaim stale runs: mark the run as `Reclaimed`, then move the card -/// back to `todo` (re-dispatchable) or `blocked` (if reclaim count -/// exceeds `max_reclaim_count`). +/// Reclaim stale runs and publish a `TaskRunReclaimed` event per reclaimed +/// card, so the Tasks board UI sees a wedged card come back without a refresh. pub async fn reclaim_stale( location: &BoardLocation, limits: &RunLimits, ) -> Result { - tracing::debug!( - thread_id = ?location.thread_id(), - "[todos][runs] reclaim_stale entry" - ); - - let stale_runs = find_stale_runs(location, limits)?; - if stale_runs.is_empty() { - return Ok(ReclaimResult { - reclaimed_count: 0, - blocked_count: 0, - details: Vec::new(), - }); + let (store, thread_id) = target(location); + let result = map_err(crate_runs::reclaim_stale(&store, thread_id, limits).await)?; + + if let Some(thread_id) = location.thread_id() { + for detail in &result.details { + crate::core::bus::BUS.publish(crate::core::events::DomainEvent::TaskRunReclaimed { + run_id: detail.run_id.clone(), + card_id: detail.card_id.clone(), + thread_id: thread_id.to_string(), + reason: detail.reason.clone(), + }); + } } + Ok(result) +} - let mut reclaimed_count = 0usize; - let mut blocked_count = 0usize; - let mut details = Vec::new(); - - for (stale_run, reason) in &stale_runs { - if let Err(e) = complete_run( - location, - &stale_run.run_id, - RunOutcome::Reclaimed, - Some(reason.clone()), - Vec::new(), - ) { - tracing::warn!( - run_id = %stale_run.run_id, - error = %e, - "[todos][runs] failed to complete stale run" - ); - continue; - } +/// Tick the run's heartbeat in the background until it completes or `cancel` +/// fires. Board-location addressing is resolved once, here, so the crate task +/// carries only a store and a thread id. +pub fn spawn_heartbeat_task( + location: BoardLocation, + run_id: String, + cancel: tokio::sync::watch::Receiver, +) { + let (store, thread_id) = target(&location); + crate_runs::spawn_heartbeat_task(store, thread_id.to_string(), run_id, cancel, HEARTBEAT_TICK); +} - let prior_reclaims = count_reclaims_for_card(location, &stale_run.card_id).unwrap_or(0); +// ── Legacy ledger migration ──────────────────────────────────────────── - let (new_status, new_status_str) = if prior_reclaims >= limits.max_reclaim_count { - (TaskCardStatus::Blocked, "blocked") - } else { - (TaskCardStatus::Todo, "todo") - }; +/// Outcome of the one-time `.runs.json` import. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct TaskRunMigrationReport { + pub total: usize, + pub copied: usize, + pub skipped: usize, +} - let blocker_msg = if new_status == TaskCardStatus::Blocked { - Some(format!( - "Reclaimed {prior_reclaims} time(s), exceeding limit of {}. \ - Last reclaim reason: {reason}", - limits.max_reclaim_count - )) - } else { - None - }; +/// Copy any run ledgers left in the retired file tree into the crate store, +/// without replacing runs the crate already holds. +/// +/// A thread whose crate log is non-empty is skipped wholesale: the crate log is +/// authoritative, and merging two histories would double-count the reclaims the +/// sweep's `max_reclaim_count` budget is derived from. +pub async fn migrate_legacy_task_runs( + workspace_dir: &Path, +) -> Result { + let dir = workspace_dir.join(TASK_BOARD_DIR); + let mut entries = match tokio::fs::read_dir(&dir).await { + Ok(entries) => entries, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + return Ok(TaskRunMigrationReport::default()); + } + Err(error) => return Err(format!("read legacy runs dir {}: {error}", dir.display())), + }; - let patch = CardPatch { - status: Some(new_status), - blocker: blocker_msg, - ..Default::default() + let mut report = TaskRunMigrationReport::default(); + while let Some(entry) = entries + .next_entry() + .await + .map_err(|error| format!("iterate legacy runs dir: {error}"))? + { + let path = entry.path(); + let Some(thread_id) = legacy_thread_id(&path) else { + continue; }; - - match ops::edit(location, &stale_run.card_id, patch).await { - Ok(_) => { - tracing::info!( - run_id = %stale_run.run_id, - card_id = %stale_run.card_id, - new_status = new_status_str, - reason = %reason, - prior_reclaims, - "[todos][runs] card reclaimed" - ); - - if let Some(thread_id) = location.thread_id() { - crate::core::bus::BUS.publish( - crate::core::events::DomainEvent::TaskRunReclaimed { - run_id: stale_run.run_id.clone(), - card_id: stale_run.card_id.clone(), - thread_id: thread_id.to_string(), - reason: reason.clone(), - }, - ); + report.total += 1; + + let runs: Vec = match tokio::fs::read_to_string(&path).await { + Ok(body) => match serde_json::from_str(&body) { + Ok(runs) => runs, + Err(error) => { + tracing::warn!(path = %path.display(), %error, "skip invalid legacy run ledger"); + report.skipped += 1; + continue; } - - if new_status == TaskCardStatus::Blocked { - blocked_count += 1; - } else { - reclaimed_count += 1; - } - details.push(ReclaimDetail { - run_id: stale_run.run_id.clone(), - card_id: stale_run.card_id.clone(), - reason: reason.clone(), - new_card_status: new_status_str.to_string(), - }); + }, + Err(error) => { + tracing::warn!(path = %path.display(), %error, "skip unreadable legacy run ledger"); + report.skipped += 1; + continue; } - Err(e) => { - tracing::warn!( - run_id = %stale_run.run_id, - card_id = %stale_run.card_id, - error = %e, - "[todos][runs] failed to update card after reclaim" - ); + }; + + let location = BoardLocation::Thread { + workspace_dir: workspace_dir.to_path_buf(), + thread_id: thread_id.clone(), + }; + let (store, thread_id) = target(&location); + match map_err(crate_runs::import_if_absent(&store, thread_id, runs).await) { + Ok(true) => report.copied += 1, + Ok(false) => report.skipped += 1, + Err(error) => { + tracing::warn!(path = %path.display(), %error, "skip legacy run ledger: store write failed"); + report.skipped += 1; } } } - - tracing::info!( - reclaimed_count, - blocked_count, - "[todos][runs] reclaim_stale complete" - ); - - Ok(ReclaimResult { - reclaimed_count, - blocked_count, - details, - }) + Ok(report) } -fn count_reclaims_for_card(location: &BoardLocation, card_id: &str) -> Result { - let runs = load_runs(location)?; - let count = runs - .iter() - .filter(|r| r.card_id == card_id && r.outcome.as_ref() == Some(&RunOutcome::Reclaimed)) - .count(); - Ok(count as u32) +/// Decode the thread id encoded in a `.runs.json` file name. +/// +/// Only strict, ASCII, even-length lowercase hex is accepted. The inner two +/// bytes of every pair must both be hexadecimal digits (`0-9a-f`), so a signed +/// or malformed stem such as `+f` is rejected rather than accepted by +/// `u8::from_str_radix`. Decoding walks `hex.as_bytes()` in whole pairs, never +/// slicing a multi-byte UTF-8 character, so a non-ASCII stem like `aéb` returns +/// `None` instead of panicking mid-startup. +fn legacy_thread_id(path: &Path) -> Option { + let name = path.file_name()?.to_str()?; + let hex = name.strip_suffix(".runs.json")?; + if hex.is_empty() || hex.len() % 2 != 0 || !hex.is_ascii() { + return None; + } + let bytes: Vec = hex + .as_bytes() + .chunks_exact(2) + .map(|pair| { + let hi = lowercase_hex_nibble(pair[0])?; + let lo = lowercase_hex_nibble(pair[1])?; + Some(hi * 16 + lo) + }) + .collect::>>()?; + String::from_utf8(bytes).ok() } -// ── Heartbeat background task ────────────────────────────────────────── - -pub fn spawn_heartbeat_task( - location: BoardLocation, - run_id: String, - cancel: tokio::sync::watch::Receiver, -) { - tokio::spawn(async move { - let mut ticker = tokio::time::interval(std::time::Duration::from_secs(HEARTBEAT_TICK_SECS)); - let mut cancel = cancel; - ticker.tick().await; // skip the immediate fire - loop { - tokio::select! { - _ = ticker.tick() => { - if let Err(e) = update_heartbeat(&location, &run_id) { - tracing::debug!( - run_id = %run_id, - error = %e, - "[todos][runs] heartbeat tick failed (run may have completed)" - ); - break; - } - } - _ = cancel.changed() => { - tracing::debug!( - run_id = %run_id, - "[todos][runs] heartbeat cancelled (run completed)" - ); - break; - } - } - } - }); +/// Decode one ASCII byte as a lowercase hexadecimal nibble (`0-9a-f`), or +/// `None` for any other byte (including uppercase `A-F`). +fn lowercase_hex_nibble(byte: u8) -> Option { + match byte { + b'0'..=b'9' => Some(byte - b'0'), + b'a'..=b'f' => Some(byte - b'a' + 10), + _ => None, + } } // ── Tests ────────────────────────────────────────────────────────────── @@ -510,6 +255,8 @@ pub fn spawn_heartbeat_task( #[cfg(test)] mod tests { use super::*; + use crate::openhuman::agent::task_board::TaskCardStatus; + use crate::openhuman::threads::todos::ops::{self, CardPatch}; use tempfile::tempdir; fn thread_loc(dir: &Path, id: &str) -> BoardLocation { @@ -519,295 +266,344 @@ mod tests { } } - #[test] - fn create_and_list_run() { + /// Lowercase-hex encode a thread id, matching [`super::legacy_thread_id`]'s + /// decoder so test-built keys round-trip through the migration. + fn hex_key(id: &str) -> String { + id.as_bytes().iter().map(|b| format!("{b:02x}")).collect() + } + + #[tokio::test] + async fn create_and_list_run() { let dir = tempdir().unwrap(); let loc = thread_loc(dir.path(), "run-test-1"); - let run = create_run(&loc, "run-1", "card-1", "default").unwrap(); + let run = create_run(&loc, "run-1", "card-1", "default") + .await + .unwrap(); assert_eq!(run.run_id, "run-1"); assert_eq!(run.card_id, "card-1"); assert_eq!(run.claimed_by, "default"); assert!(run.is_active()); assert!(!run.claim_token.is_empty()); - let all = list_runs(&loc, None).unwrap(); + let all = list_runs(&loc, None).await.unwrap(); assert_eq!(all.len(), 1); + assert_eq!(list_runs(&loc, Some("card-1")).await.unwrap().len(), 1); + assert!(list_runs(&loc, Some("card-other")) + .await + .unwrap() + .is_empty()); + } + + #[tokio::test] + async fn timestamps_reach_the_wire_as_rfc3339() { + let dir = tempdir().unwrap(); + let loc = thread_loc(dir.path(), "wire-test"); + let run = create_run(&loc, "run-1", "card-1", "default") + .await + .unwrap(); - let by_card = list_runs(&loc, Some("card-1")).unwrap(); - assert_eq!(by_card.len(), 1); + // The RPC surface has always spoken RFC 3339; the crate stores millis. + assert!(chrono::DateTime::parse_from_rfc3339(&run.started_at).is_ok()); + assert!(chrono::DateTime::parse_from_rfc3339(&run.last_heartbeat_at).is_ok()); - let empty = list_runs(&loc, Some("card-other")).unwrap(); - assert!(empty.is_empty()); + let done = complete_run(&loc, "run-1", RunOutcome::Success, None, vec![]) + .await + .unwrap(); + let completed_at = done.completed_at.expect("completed stamp"); + assert!(chrono::DateTime::parse_from_rfc3339(&completed_at).is_ok()); } - #[test] - fn heartbeat_updates_timestamp() { + #[tokio::test] + async fn heartbeat_updates_timestamp() { let dir = tempdir().unwrap(); let loc = thread_loc(dir.path(), "hb-test-1"); - create_run(&loc, "run-hb", "card-1", "default").unwrap(); - let before = get_run(&loc, "run-hb").unwrap().unwrap(); + create_run(&loc, "run-hb", "card-1", "default") + .await + .unwrap(); + let before = get_run(&loc, "run-hb").await.unwrap().unwrap(); - std::thread::sleep(std::time::Duration::from_millis(10)); - update_heartbeat(&loc, "run-hb").unwrap(); + update_heartbeat(&loc, "run-hb").await.unwrap(); - let after = get_run(&loc, "run-hb").unwrap().unwrap(); + let after = get_run(&loc, "run-hb").await.unwrap().unwrap(); assert!(after.last_heartbeat_at >= before.last_heartbeat_at); } - #[test] - fn heartbeat_fails_for_completed_run() { + #[tokio::test] + async fn heartbeat_fails_for_completed_run() { let dir = tempdir().unwrap(); let loc = thread_loc(dir.path(), "hb-test-2"); - create_run(&loc, "run-done", "card-1", "default").unwrap(); - complete_run(&loc, "run-done", RunOutcome::Success, None, Vec::new()).unwrap(); + create_run(&loc, "run-done", "card-1", "default") + .await + .unwrap(); + complete_run(&loc, "run-done", RunOutcome::Success, None, vec![]) + .await + .unwrap(); - assert!(update_heartbeat(&loc, "run-done").is_err()); + assert!(update_heartbeat(&loc, "run-done").await.is_err()); } - #[test] - fn complete_run_sets_outcome() { + #[tokio::test] + async fn complete_run_sets_outcome() { let dir = tempdir().unwrap(); let loc = thread_loc(dir.path(), "complete-test"); - create_run(&loc, "run-c", "card-1", "default").unwrap(); - let completed = complete_run( + create_run(&loc, "run-ok", "card-1", "default") + .await + .unwrap(); + let done = complete_run( &loc, - "run-c", + "run-ok", RunOutcome::Success, None, - vec!["opened PR #5".to_string()], + vec!["evidence".to_string()], ) + .await .unwrap(); - assert!(!completed.is_active()); - assert_eq!(completed.outcome, Some(RunOutcome::Success)); - assert!(completed.completed_at.is_some()); - assert_eq!(completed.evidence, vec!["opened PR #5"]); + assert!(!done.is_active()); + assert_eq!(done.outcome, Some(RunOutcome::Success)); + assert_eq!(done.evidence, vec!["evidence".to_string()]); } - #[test] - fn complete_run_with_failure() { + #[tokio::test] + async fn complete_run_with_failure() { let dir = tempdir().unwrap(); let loc = thread_loc(dir.path(), "fail-test"); - create_run(&loc, "run-f", "card-1", "default").unwrap(); - let completed = complete_run( + create_run(&loc, "run-bad", "card-1", "default") + .await + .unwrap(); + let done = complete_run( &loc, - "run-f", + "run-bad", RunOutcome::Failed, - Some("agent build failed".to_string()), - Vec::new(), + Some("boom".to_string()), + vec![], ) + .await .unwrap(); - assert_eq!(completed.outcome, Some(RunOutcome::Failed)); - assert_eq!(completed.error.as_deref(), Some("agent build failed")); + assert_eq!(done.outcome, Some(RunOutcome::Failed)); + assert_eq!(done.error.as_deref(), Some("boom")); } - #[test] - fn get_run_returns_none_for_missing() { + #[tokio::test] + async fn get_run_returns_none_for_missing() { let dir = tempdir().unwrap(); - let loc = thread_loc(dir.path(), "get-test"); - - assert!(get_run(&loc, "no-such-run").unwrap().is_none()); - } - - #[test] - fn check_staleness_detects_expired_ttl() { - let now = Utc::now(); - let old = (now - chrono::Duration::seconds(7200)).to_rfc3339(); - let run = TaskRun { - run_id: "r1".into(), - card_id: "c1".into(), - claimed_by: "test".into(), - claim_token: "t".into(), - started_at: old.clone(), - last_heartbeat_at: now.to_rfc3339(), - completed_at: None, - outcome: None, - error: None, - evidence: Vec::new(), - }; - let limits = RunLimits::default(); - let reason = check_staleness(&run, &now, &limits); - assert!(reason.is_some()); - assert!(reason.unwrap().contains("TTL expired")); + let loc = thread_loc(dir.path(), "missing-test"); + assert!(get_run(&loc, "nope").await.unwrap().is_none()); } - #[test] - fn check_staleness_detects_stale_heartbeat() { - let now = Utc::now(); - let recent_start = (now - chrono::Duration::seconds(60)).to_rfc3339(); - let old_hb = (now - chrono::Duration::seconds(600)).to_rfc3339(); - let run = TaskRun { - run_id: "r2".into(), - card_id: "c1".into(), - claimed_by: "test".into(), - claim_token: "t".into(), - started_at: recent_start, - last_heartbeat_at: old_hb, - completed_at: None, - outcome: None, - error: None, - evidence: Vec::new(), - }; - let limits = RunLimits::default(); - let reason = check_staleness(&run, &now, &limits); - assert!(reason.is_some()); - assert!(reason.unwrap().contains("heartbeat stale")); + /// Age a run past every limit by rewriting its stamps in the crate store. + async fn wedge(loc: &BoardLocation, run_id: &str) { + let (store, thread_id) = target(loc); + let mut runs = crate_runs::list_runs(&store, thread_id, None) + .await + .unwrap(); + for run in runs.iter_mut().filter(|run| run.run_id == run_id) { + run.started_at = "0".to_string(); + run.last_heartbeat_at = "0".to_string(); + } + let key = hex_key(&thread_id); + store + .put( + crate_runs::RUNS_NAMESPACE, + &key, + serde_json::to_value(&runs).unwrap(), + ) + .await + .unwrap(); } - #[test] - fn check_staleness_passes_healthy_run() { - let now = Utc::now(); - let recent = (now - chrono::Duration::seconds(10)).to_rfc3339(); - let run = TaskRun { - run_id: "r3".into(), - card_id: "c1".into(), - claimed_by: "test".into(), - claim_token: "t".into(), - started_at: recent.clone(), - last_heartbeat_at: recent, - completed_at: None, - outcome: None, - error: None, - evidence: Vec::new(), - }; - let limits = RunLimits::default(); - assert!(check_staleness(&run, &now, &limits).is_none()); + async fn seed_in_progress_card(loc: &BoardLocation, title: &str) -> String { + let snapshot = ops::add(loc, title, CardPatch::default()).await.unwrap(); + let card_id = snapshot.cards.last().unwrap().id.clone(); + ops::edit( + loc, + &card_id, + CardPatch { + status: Some(TaskCardStatus::InProgress), + ..Default::default() + }, + ) + .await + .unwrap(); + card_id } #[tokio::test] async fn reclaim_stale_moves_card_to_todo() { let dir = tempdir().unwrap(); - let loc = thread_loc(dir.path(), "reclaim-test-1"); + let loc = thread_loc(dir.path(), "reclaim-test"); + let card_id = seed_in_progress_card(&loc, "wedged work").await; - let snap = ops::add(&loc, "reclaimable task", CardPatch::default()) + create_run(&loc, "run-stale", &card_id, "default") .await .unwrap(); - let card_id = snap.cards[0].id.clone(); - ops::update_status(&loc, &card_id, TaskCardStatus::InProgress) - .await - .unwrap(); - - // Create a run with an old heartbeat - { - let lock = run_lock(&loc); - let _guard = lock.lock(); - let old = (Utc::now() - chrono::Duration::seconds(600)).to_rfc3339(); - let run = TaskRun { - run_id: "stale-run".into(), - card_id: card_id.clone(), - claimed_by: "test".into(), - claim_token: "t".into(), - started_at: old.clone(), - last_heartbeat_at: old, - completed_at: None, - outcome: None, - error: None, - evidence: Vec::new(), - }; - save_runs(&loc, &[run]).unwrap(); - } + wedge(&loc, "run-stale").await; let result = reclaim_stale(&loc, &RunLimits::default()).await.unwrap(); assert_eq!(result.reclaimed_count, 1); assert_eq!(result.blocked_count, 0); assert_eq!(result.details[0].new_card_status, "todo"); - let snap = ops::list(&loc).await.unwrap(); - assert_eq!(snap.cards[0].status, TaskCardStatus::Todo); + let snapshot = ops::list(&loc).await.unwrap(); + assert_eq!(snapshot.cards[0].status, TaskCardStatus::Todo); } #[tokio::test] async fn reclaim_blocks_after_max_reclaims() { let dir = tempdir().unwrap(); - let loc = thread_loc(dir.path(), "reclaim-block-test"); - - let snap = ops::add(&loc, "troublesome task", CardPatch::default()) - .await - .unwrap(); - let card_id = snap.cards[0].id.clone(); - ops::update_status(&loc, &card_id, TaskCardStatus::InProgress) - .await - .unwrap(); + let loc = thread_loc(dir.path(), "reclaim-max-test"); + let card_id = seed_in_progress_card(&loc, "poison work").await; + // Two reclaims are tolerated; the one that reaches the limit parks it. + let limits = RunLimits { + max_reclaim_count: 2, + ..RunLimits::default() + }; - // Seed prior reclaimed runs (3 = at the limit) - { - let lock = run_lock(&loc); - let _guard = lock.lock(); - let old = (Utc::now() - chrono::Duration::seconds(600)).to_rfc3339(); - let mut runs = Vec::new(); - for i in 0..3 { - runs.push(TaskRun { - run_id: format!("prior-{i}"), - card_id: card_id.clone(), - claimed_by: "test".into(), - claim_token: format!("t{i}"), - started_at: old.clone(), - last_heartbeat_at: old.clone(), - completed_at: Some(old.clone()), - outcome: Some(RunOutcome::Reclaimed), - error: Some("stale".into()), - evidence: Vec::new(), - }); + for (attempt, expected) in ["todo", "blocked"].iter().enumerate() { + if attempt > 0 { + ops::edit( + &loc, + &card_id, + CardPatch { + status: Some(TaskCardStatus::InProgress), + ..Default::default() + }, + ) + .await + .unwrap(); } - // Active stale run - runs.push(TaskRun { - run_id: "current-stale".into(), - card_id: card_id.clone(), - claimed_by: "test".into(), - claim_token: "tc".into(), - started_at: old.clone(), - last_heartbeat_at: old, - completed_at: None, - outcome: None, - error: None, - evidence: Vec::new(), - }); - save_runs(&loc, &runs).unwrap(); + let run_id = format!("run-{attempt}"); + create_run(&loc, &run_id, &card_id, "default") + .await + .unwrap(); + wedge(&loc, &run_id).await; + let result = reclaim_stale(&loc, &limits).await.unwrap(); + assert_eq!(&result.details[0].new_card_status, expected); } - let result = reclaim_stale(&loc, &RunLimits::default()).await.unwrap(); - assert_eq!(result.reclaimed_count, 0); - assert_eq!(result.blocked_count, 1); - assert_eq!(result.details[0].new_card_status, "blocked"); - - let snap = ops::list(&loc).await.unwrap(); - assert_eq!(snap.cards[0].status, TaskCardStatus::Blocked); - assert!(snap.cards[0] + let snapshot = ops::list(&loc).await.unwrap(); + assert_eq!(snapshot.cards[0].status, TaskCardStatus::Blocked); + assert!(snapshot.cards[0] .blocker .as_deref() .unwrap_or_default() - .contains("Reclaimed")); + .contains("exceeding limit of 2")); } #[tokio::test] async fn reclaim_skips_healthy_runs() { let dir = tempdir().unwrap(); - let loc = thread_loc(dir.path(), "healthy-test"); - - let snap = ops::add(&loc, "healthy task", CardPatch::default()) + let loc = thread_loc(dir.path(), "reclaim-healthy-test"); + let card_id = seed_in_progress_card(&loc, "live work").await; + create_run(&loc, "run-live", &card_id, "default") .await .unwrap(); - let card_id = snap.cards[0].id.clone(); - ops::update_status(&loc, &card_id, TaskCardStatus::InProgress) - .await - .unwrap(); - - create_run(&loc, "healthy-run", &card_id, "default").unwrap(); let result = reclaim_stale(&loc, &RunLimits::default()).await.unwrap(); assert_eq!(result.reclaimed_count, 0); assert_eq!(result.blocked_count, 0); + + let snapshot = ops::list(&loc).await.unwrap(); + assert_eq!(snapshot.cards[0].status, TaskCardStatus::InProgress); } - #[test] - fn scratch_location_returns_empty_runs() { - let runs = list_runs(&BoardLocation::Scratch, None).unwrap(); + #[tokio::test] + async fn scratch_location_returns_empty_runs() { + // Serialize against the process-global scratch store shared with + // `todos::ops` / agent-tool tests (see `scratch_test_lock`). + let _guard = ops::scratch_test_lock(); + let runs = list_runs(&BoardLocation::Scratch, None).await.unwrap(); assert!(runs.is_empty()); } + + #[tokio::test] + async fn legacy_ledger_is_imported_once_and_never_replaces_crate_runs() { + let workspace = tempdir().unwrap(); + let legacy_dir = workspace.path().join(TASK_BOARD_DIR); + tokio::fs::create_dir_all(&legacy_dir).await.unwrap(); + + let thread_id = "legacy-thread"; + let hex = hex_key(thread_id); + let legacy = vec![TaskRun { + run_id: "legacy-run".to_string(), + card_id: "card-1".to_string(), + claimed_by: "default".to_string(), + claim_token: "token".to_string(), + started_at: "0".to_string(), + last_heartbeat_at: "0".to_string(), + completed_at: None, + outcome: None, + error: None, + evidence: Vec::new(), + }]; + tokio::fs::write( + legacy_dir.join(format!("{hex}.runs.json")), + serde_json::to_vec(&legacy).unwrap(), + ) + .await + .unwrap(); + // A board file alongside it must not be mistaken for a run ledger. + tokio::fs::write(legacy_dir.join(format!("{hex}.json")), b"{}") + .await + .unwrap(); + + let first = migrate_legacy_task_runs(workspace.path()).await.unwrap(); + assert_eq!( + first, + TaskRunMigrationReport { + total: 1, + copied: 1, + skipped: 0, + } + ); + + let loc = thread_loc(workspace.path(), thread_id); + assert_eq!(list_runs(&loc, None).await.unwrap().len(), 1); + + // Second pass: the crate log is authoritative and is left alone. + let second = migrate_legacy_task_runs(workspace.path()).await.unwrap(); + assert_eq!(second.copied, 0); + assert_eq!(second.skipped, 1); + assert_eq!(list_runs(&loc, None).await.unwrap().len(), 1); + } + + #[test] + fn legacy_file_names_decode_only_run_ledgers() { + let hex: String = "thread-1" + .as_bytes() + .iter() + .map(|b| format!("{b:02x}")) + .collect(); + assert_eq!( + legacy_thread_id(Path::new(&format!("/w/{hex}.runs.json"))).as_deref(), + Some("thread-1") + ); + assert!(legacy_thread_id(Path::new(&format!("/w/{hex}.json"))).is_none()); + assert!(legacy_thread_id(Path::new("/w/notes.txt")).is_none()); + assert!(legacy_thread_id(Path::new("/w/zz.runs.json")).is_none()); + } + + #[test] + fn legacy_file_names_reject_malformed_stems_without_panicking() { + // Multi-byte UTF-8 in the stem slices an odd index if decoded by byte + // pairs; it must be rejected, not panic. + assert!(legacy_thread_id(Path::new("/w/aéb.runs.json")).is_none()); + // from_str_radix would accept "+f" as 15; strict hex decoding rejects the sign. + assert!(legacy_thread_id(Path::new("/w/+f+f.runs.json")).is_none()); + assert!(legacy_thread_id(Path::new("/w/gg.runs.json")).is_none()); + assert!(legacy_thread_id(Path::new("/w/GG.runs.json")).is_none()); + // Uppercase hex is rejected even though to_digit(16) would accept it. + assert!(legacy_thread_id(Path::new("/w/4A.runs.json")).is_none()); + assert!(legacy_thread_id(Path::new("/w/4a.runs.json")).is_some()); + // A mixed pair with a valid second nibble but invalid first is rejected. + assert!(legacy_thread_id(Path::new("/w/0g.runs.json")).is_none()); + // Non-hex ASCII letters are rejected. + assert!(legacy_thread_id(Path::new("/w/zz.runs.json")).is_none()); + } } diff --git a/src/openhuman/threads/todos/schemas.rs b/src/openhuman/threads/todos/schemas.rs index 4c8fc75ad5..043521a759 100644 --- a/src/openhuman/threads/todos/schemas.rs +++ b/src/openhuman/threads/todos/schemas.rs @@ -655,7 +655,7 @@ fn handle_run_list(params: Map) -> ControllerFuture { card_id = ?p.card_id, "[rpc][todos] run_list entry" ); - let run_list = runs::list_runs(&loc, p.card_id.as_deref())?; + let run_list = runs::list_runs(&loc, p.card_id.as_deref()).await?; serde_json::to_value(&run_list).map_err(|e| format!("serialize runs: {e}")) }) } @@ -669,7 +669,7 @@ fn handle_run_get(params: Map) -> ControllerFuture { run_id = %p.run_id, "[rpc][todos] run_get entry" ); - let run = runs::get_run(&loc, &p.run_id)?; + let run = runs::get_run(&loc, &p.run_id).await?; serde_json::to_value(&run).map_err(|e| format!("serialize run: {e}")) }) } diff --git a/tests/raw_coverage/agent_archivist_debug_round21_raw_coverage_e2e.rs b/tests/raw_coverage/agent_archivist_debug_round21_raw_coverage_e2e.rs index f8412c7dce..5f449ac469 100644 --- a/tests/raw_coverage/agent_archivist_debug_round21_raw_coverage_e2e.rs +++ b/tests/raw_coverage/agent_archivist_debug_round21_raw_coverage_e2e.rs @@ -436,11 +436,11 @@ fn debug_dump_writer_sanitizes_names_and_writes_summary_sidecars() -> Result<()> workspace_dir: PathBuf::from("/tmp/round21-workspace"), text: "SYSTEM PROMPT\n".to_string(), tool_names: vec!["echo".to_string(), "search".to_string()], - skill_tool_count: 1, tool_specs: vec![ - json!({"name": "echo", "description": "echo", "parameters": {}}), - json!({"name": "search", "description": "search", "parameters": {}}), + json!({"name": "echo", "description": "echo back", "parameters": {}}), + json!({"name": "search", "description": "search docs", "parameters": {}}), ], + skill_tool_count: 1, }]; let summary = write_prompt_dumps(tmp.path(), &dumps)?; @@ -461,5 +461,14 @@ fn debug_dump_writer_sanitizes_names_and_writes_summary_sidecars() -> Result<()> let summary_text = std::fs::read_to_string(summary.summary_path)?; assert!(summary_text.contains("agent/with spaces@gmail:primary")); assert!(summary_text.contains("tools=2")); + // The per-dump tools sidecar carries the rendered tool schemas verbatim, + // one entry per tool in `tool_names` order. + let tools_json = std::fs::read_to_string( + tmp.path().join("1_agent_with_spaces_gmail_primary.tools.json"), + )?; + let specs: Vec = serde_json::from_str(&tools_json)?; + assert_eq!(specs.len(), 2); + assert_eq!(specs[0]["name"], "echo"); + assert_eq!(specs[1]["name"], "search"); Ok(()) } diff --git a/tests/raw_coverage/inference_agent_raw_coverage_e2e.rs b/tests/raw_coverage/inference_agent_raw_coverage_e2e.rs index 35f8b33bc2..d8b0c29dc1 100644 --- a/tests/raw_coverage/inference_agent_raw_coverage_e2e.rs +++ b/tests/raw_coverage/inference_agent_raw_coverage_e2e.rs @@ -3875,11 +3875,11 @@ async fn agent_debug_prompt_dump_and_identity_rendering_cover_file_layouts() { workspace_dir: workspace.path().join("ws"), text: "# planner\nbody\n".to_string(), tool_names: vec!["todo".to_string(), "delegate".to_string()], - skill_tool_count: 0, tool_specs: vec![ - json!({"name": "todo", "description": "todo", "parameters": {}}), - json!({"name": "delegate", "description": "delegate", "parameters": {}}), + json!({"name": "todo", "description": "manage todos", "parameters": {}}), + json!({"name": "delegate", "description": "delegate a task", "parameters": {}}), ], + skill_tool_count: 0, }, DumpedPrompt { agent_id: "integrations_agent".to_string(), @@ -3889,12 +3889,12 @@ async fn agent_debug_prompt_dump_and_identity_rendering_cover_file_layouts() { workspace_dir: workspace.path().join("ws"), text: "# integrations\nbody\n".to_string(), tool_names: vec!["GMAIL_SEND_EMAIL".to_string()], - skill_tool_count: 1, tool_specs: vec![json!({ "name": "GMAIL_SEND_EMAIL", - "description": "send email", + "description": "send an email", "parameters": {}, })], + skill_tool_count: 1, }, ]; @@ -3929,6 +3929,32 @@ async fn agent_debug_prompt_dump_and_identity_rendering_cover_file_layouts() { assert!(summary_text.contains("planner/coverage")); assert!(summary_text.contains("integrations_agent@gmail+calendar")); + // Each per-dump tools sidecar carries the rendered tool schemas verbatim, + // one entry per tool in `tool_names` order — compare the full payload + // (name, description and parameters), not just count and names. + let planner_tools: Vec = serde_json::from_str( + &std::fs::read_to_string( + workspace.path().join("1_planner_coverage.tools.json"), + ) + .expect("planner tools sidecar"), + ) + .expect("planner tools json"); + assert_eq!(planner_tools.as_slice(), dumps[0].tool_specs.as_slice()); + + let integrations_tools: Vec = serde_json::from_str( + &std::fs::read_to_string( + workspace + .path() + .join("2_integrations_agent_gmail_calendar.tools.json"), + ) + .expect("integrations tools sidecar"), + ) + .expect("integrations tools json"); + assert_eq!( + integrations_tools.as_slice(), + dumps[1].tool_specs.as_slice() + ); + let identities = openhuman_core::openhuman::agent::prompts::render_connected_identities(); assert_eq!(identities, ""); } diff --git a/tests/raw_coverage/tools_approval_channels_raw_coverage_e2e.rs b/tests/raw_coverage/tools_approval_channels_raw_coverage_e2e.rs index d11472336e..b6e8532f51 100644 --- a/tests/raw_coverage/tools_approval_channels_raw_coverage_e2e.rs +++ b/tests/raw_coverage/tools_approval_channels_raw_coverage_e2e.rs @@ -1620,19 +1620,14 @@ async fn orchestrator_tool_synthesis_covers_agent_and_integration_delegation_edg assert_eq!(names, vec!["research", "delegate_to_integrations_agent"]); let research = &tools[0]; - // Delegate-tool descriptions carry the target agent's `when_to_use` - // verbatim. The repeated "direct response/direct tools are insufficient" - // prefix was intentionally removed (the orchestrator's own prompt already - // carries that rule once), so assert it is gone rather than still required. - assert!(research - .description() - .contains("careful public-source research")); - assert!(research - .description() - .contains("Use for careful public-source research.")); - assert!(!research - .description() - .contains("direct tools are insufficient")); + // The delegation tool's description is the target agent's `when_to_use` + // verbatim (the "Use only when direct response/direct tools are + // insufficient." prefix was deliberately dropped — it is stated once in + // the orchestrator prompt instead of once per delegate schema per turn). + assert_eq!( + research.description(), + "Use for careful public-source research." + ); assert_eq!(research.permission_level(), PermissionLevel::Execute); assert_eq!(research.category(), ToolCategory::System); assert_eq!(