From 58e0f46c97deccf5751df1d42f7b2fd52477932f Mon Sep 17 00:00:00 2001 From: vitaliytv Date: Mon, 10 Aug 2026 13:48:42 +0300 Subject: [PATCH 1/2] =?UTF-8?q?feat(agent-server):=20orchestrator-=D1=80?= =?UTF-8?q?=D0=BE=D0=BB=D1=8C=20=D1=96=D0=B7=20=D0=BF=D0=BE=D0=B4=D1=96?= =?UTF-8?q?=D1=94=D0=B2=D0=B8=D0=BC=20wake?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Хвиля 3. runtime.md, «Wake: push замість polling»: базовий MT прокидався cron-ом кожні 5 хвилин; цільова архітектура має три джерела пробудження, і на кожне виконує ту саму трійку mt watch-логіки. Реалізовано всі три джерела: - relay push «є нові події у задачі X» -> AppState::wake_orchestrator, тобто хост ресканить негайно, а не чекає таймера; - touch .mt/wake від post-merge git hook; - періодичний tick як FALLBACK — щоб при недоступному relay система працювала як базовий MT, а не зупинялась. Один привід — одне прокидання: сигнал споживається, mtime файлу запам'ятовується. Інакше кожна ітерація циклу перезапускала б dispatch на тій самій події. tick(): dispatch -> алерти -> GC. Порядок не випадковий: dispatch першим, бо заради нього хост і прокидається; алерти після нього, щоб вузол, який щойно став unresolvable у цьому ж прогоні, потрапив у той самий звіт; GC останнім — він прибирає за тим, що вже завершилось. Алерти дедуплікуються за станом процесу, а не файлом-міткою: алерт це доставка, а не артефакт графа, тож перезапуск хоста має право нагадати про вузол, який усе ще чекає людину. Опитування .mt/wake замість inotify — свідомий компроміс: інотифай додав би залежність і платформні гілки заради того, що на масштабі одного репозиторію коштує один stat на 200 мс. Це закриває і хвіст хвилі 2: алерт при unresolvable тепер має того, хто його помічає. Тести: 5 — алерт раз на вузол, алерти доходять до вкладених вузлів, явний сигнал будить негайно і споживається, touch файлу будить (а сама наявність файлу — ні), fallback повертає без події. 360 passed, clippy чистий. Co-Authored-By: Claude Fable 5 --- crates/agent-server/src/lib.rs | 1 + crates/agent-server/src/orchestrator.rs | 277 ++++++++++++++++++++++++ crates/agent-server/src/relay_client.rs | 3 + crates/agent-server/src/ws.rs | 24 ++ docs/conformance.md | 8 +- 5 files changed, 309 insertions(+), 4 deletions(-) create mode 100644 crates/agent-server/src/orchestrator.rs diff --git a/crates/agent-server/src/lib.rs b/crates/agent-server/src/lib.rs index ee23efd..9766220 100644 --- a/crates/agent-server/src/lib.rs +++ b/crates/agent-server/src/lib.rs @@ -12,6 +12,7 @@ pub mod approvals_gate; pub mod discovery; pub mod graph; +pub mod orchestrator; pub mod relay_client; pub mod runner; pub mod session; diff --git a/crates/agent-server/src/orchestrator.rs b/crates/agent-server/src/orchestrator.rs new file mode 100644 index 0000000..0dfec2a --- /dev/null +++ b/crates/agent-server/src/orchestrator.rs @@ -0,0 +1,277 @@ +//! Orchestrator-роль хоста (runtime.md, «Wake: push замість polling»). +//! +//! Базовий MT прокидався cron-ом кожні 5 хвилин. Тут — подієвий wake із +//! трьома джерелами: +//! +//! 1. relay push «є нові події у задачі X» → [`Wake::signal`]; +//! 2. `post-merge` git hook → `touch .mt/wake`; +//! 3. періодичний tick — **fallback**, щоб система працювала як базовий MT, +//! коли relay недоступний. +//! +//! На кожен wake виконується та сама трійка `mt watch`-логіки: **dispatch** +//! (запуск готових вузлів), **алерти** (вузли, що стали `unresolvable`) і +//! **GC** (прибирання відпрацьованих worktree). + +use std::collections::HashSet; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::time::{Duration, SystemTime}; + +use serde::{Deserialize, Serialize}; + +/// Підсумок одного прокидання — те, що хост має відзвітувати назовні. +#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)] +pub struct TickReport { + /// Вузли, запущені цим прокиданням, і їхні результати. + pub dispatched: Vec<(String, String)>, + /// Вузли, що вперше стали `unresolvable` — алерт власнику. + pub alerts: Vec, + /// Прибрані worktree. + pub pruned: Vec, + /// Помилки, які не мають валити цикл (наступний wake спробує знову). + pub errors: Vec, +} + +/// Джерело пробуджень: явний сигнал (relay push), файл-мітка `.mt/wake` +/// (git hook) і періодичний fallback. +pub struct Wake { + wake_file: PathBuf, + last_seen: Option, + signalled: Arc, + interval: Duration, +} + +impl Wake { + pub fn new(project_root: &Path, interval: Duration) -> Self { + let wake_file = project_root.join(".mt/wake"); + Self { + last_seen: file_mtime(&wake_file), + wake_file, + signalled: Arc::new(AtomicBool::new(false)), + interval, + } + } + + /// Ручка для relay-клієнта: push «є нові події у задачі X» будить хост + /// негайно, не чекаючи періодичного tick-а. + pub fn signaller(&self) -> Arc { + Arc::clone(&self.signalled) + } + + /// Чи є привід прокинутись **зараз** — без блокування. + /// + /// Споживає сигнал і оновлює позначку часу файлу, тож повторний виклик + /// без нової події поверне `false`: один привід — одне прокидання. + pub fn should_wake_now(&mut self) -> bool { + if self.signalled.swap(false, Ordering::SeqCst) { + return true; + } + let current = file_mtime(&self.wake_file); + if current.is_some() && current != self.last_seen { + self.last_seen = current; + return true; + } + false + } + + /// Чекає на привід не довше за `interval`; повертає `true`, якщо + /// прокинулись за подією, і `false`, якщо спрацював fallback-таймер. + /// + /// Опитування файлу — свідомий компроміс: інотифай додав би залежність + /// і платформні гілки заради того, що на масштабі одного репозиторію + /// коштує один `stat` на 200 мс. + pub fn wait(&mut self) -> bool { + let deadline = SystemTime::now() + self.interval; + loop { + if self.should_wake_now() { + return true; + } + if SystemTime::now() >= deadline { + return false; + } + std::thread::sleep(Duration::from_millis(200).min(self.interval)); + } + } +} + +fn file_mtime(path: &Path) -> Option { + std::fs::metadata(path).ok()?.modified().ok() +} + +/// Стан orchestrator-ролі між прокиданнями. +pub struct Orchestrator { + tasks_dir: String, + project_root: PathBuf, + concurrency: usize, + /// Вузли, за які алерт уже відправлено — щоб кожне прокидання не + /// повторювало той самий алерт про той самий термінальний вузол. + alerted: HashSet, +} + +impl Orchestrator { + pub fn new(tasks_dir: impl Into, concurrency: usize) -> Self { + let tasks_dir = tasks_dir.into(); + let project_root = Path::new(&tasks_dir) + .parent() + .unwrap_or(Path::new(".")) + .to_path_buf(); + Self { + tasks_dir, + project_root, + concurrency: concurrency.max(1), + alerted: HashSet::new(), + } + } + + /// Вузли, що потребують алерту цього прокидання: у стані `unresolvable` + /// і ще не оголошені. + /// + /// Дедуплікація за станом процесу, а не за файлом-міткою: алерт — це + /// доставка, а не артефакт графа, і перезапуск хоста має право нагадати + /// про вузол, який усе ще чекає людину. + fn pending_alerts(&mut self, nodes: &[mt_core::TaskNode]) -> Vec { + let mut out = Vec::new(); + let mut stack: Vec<&mt_core::TaskNode> = nodes.iter().collect(); + while let Some(node) = stack.pop() { + if node.state == mt_core::TaskState::Unresolvable + && !self.alerted.contains(&node.path) + { + self.alerted.insert(node.path.clone()); + out.push(node.path.clone()); + } + stack.extend(node.children.iter()); + } + out.sort(); + out + } + + /// Одне прокидання: dispatch → алерти → GC. + /// + /// Порядок не випадковий. Dispatch першим, бо саме заради нього хост і + /// прокидається; алерти після нього, щоб вузол, який щойно став + /// `unresolvable` у цьому ж прогоні, потрапив у той самий звіт; GC + /// останнім — він прибирає за тим, що вже завершилось. + pub fn tick(&mut self) -> TickReport { + let mut report = TickReport::default(); + + match mt_core::orchestrate::run_auto(&self.tasks_dir, self.concurrency) { + Ok(results) => { + for r in results { + if let Some(error) = r.error { + report.errors.push(format!("{}: {error}", r.path)); + } + report.dispatched.push((r.path, r.result)); + } + } + Err(error) => report.errors.push(format!("dispatch: {error}")), + } + + match mt_core::scan_tasks_with_claims(self.tasks_dir.clone(), Vec::new()) { + Ok(nodes) => report.alerts = self.pending_alerts(&nodes), + Err(error) => report.errors.push(format!("scan: {error}")), + } + + match mt_core::worktree::prune_worktrees(&self.project_root) { + Ok(output) => { + report.pruned = output + .lines() + .map(str::trim) + .filter(|l| !l.is_empty()) + .map(String::from) + .collect(); + } + Err(error) => report.errors.push(format!("gc: {error}")), + } + + report + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn node(path: &str, state: mt_core::TaskState) -> mt_core::TaskNode { + mt_core::TaskNode { + id: path.rsplit('/').next().unwrap_or(path).to_string(), + path: path.to_string(), + state, + deps: Vec::new(), + mode: "agent".to_string(), + budget_sec: None, + budget_hard_sec: None, + deadline: None, + hint: None, + created_at: None, + children: Vec::new(), + is_composite: false, + warnings: Vec::new(), + } + } + + #[test] + fn alerts_fire_once_per_node() { + let mut orch = Orchestrator::new("mt", 1); + let nodes = vec![ + node("solo", mt_core::TaskState::Unresolvable), + node("other", mt_core::TaskState::Waiting), + ]; + assert_eq!(orch.pending_alerts(&nodes), ["solo"]); + // Наступне прокидання не повторює алерт про той самий вузол. + assert!(orch.pending_alerts(&nodes).is_empty()); + } + + #[test] + fn alerts_reach_nested_nodes() { + let mut orch = Orchestrator::new("mt", 1); + let mut parent = node("parent", mt_core::TaskState::Spawned); + parent + .children + .push(node("parent/child", mt_core::TaskState::Unresolvable)); + assert_eq!(orch.pending_alerts(&[parent]), ["parent/child"]); + } + + #[test] + fn explicit_signal_wakes_immediately() { + let tmp = tempfile::tempdir().unwrap(); + let mut wake = Wake::new(tmp.path(), Duration::from_secs(3600)); + assert!(!wake.should_wake_now(), "без події — не будимо"); + + wake.signaller().store(true, Ordering::SeqCst); + assert!(wake.should_wake_now(), "relay push будить негайно"); + // Сигнал спожито: один привід — одне прокидання. + assert!(!wake.should_wake_now()); + } + + #[test] + fn touching_wake_file_wakes_host() { + let tmp = tempfile::tempdir().unwrap(); + std::fs::create_dir_all(tmp.path().join(".mt")).unwrap(); + let wake_file = tmp.path().join(".mt/wake"); + std::fs::write(&wake_file, "").unwrap(); + + let mut wake = Wake::new(tmp.path(), Duration::from_secs(3600)); + assert!(!wake.should_wake_now(), "наявний файл сам по собі — не подія"); + + // git hook торкається файлу після мерджу. + std::thread::sleep(Duration::from_millis(10)); + filetime_touch(&wake_file); + assert!(wake.should_wake_now()); + assert!(!wake.should_wake_now(), "повторно за тим самим mtime — ні"); + } + + #[test] + fn periodic_fallback_returns_without_event() { + let tmp = tempfile::tempdir().unwrap(); + let mut wake = Wake::new(tmp.path(), Duration::from_millis(50)); + // Relay недоступний, hook мовчить — прокидаємось за таймером. + assert!(!wake.wait(), "fallback, а не подія"); + } + + fn filetime_touch(path: &Path) { + let now = std::time::SystemTime::now(); + let file = std::fs::OpenOptions::new().write(true).open(path).unwrap(); + file.set_modified(now + Duration::from_secs(1)).unwrap(); + } +} diff --git a/crates/agent-server/src/relay_client.rs b/crates/agent-server/src/relay_client.rs index de74530..622083c 100644 --- a/crates/agent-server/src/relay_client.rs +++ b/crates/agent-server/src/relay_client.rs @@ -123,6 +123,9 @@ async fn handle_incoming(state: &Arc, text: &str) { let Some(envelope) = frame.get("envelope") else { return; }; + // Push «є нові події у задачі X» (runtime.md): хост ресканить негайно, + // не чекаючи fallback-таймера. + state.wake_orchestrator(); let device_id = envelope .get("device_id") .and_then(Value::as_str) diff --git a/crates/agent-server/src/ws.rs b/crates/agent-server/src/ws.rs index 61b7575..988b5a1 100644 --- a/crates/agent-server/src/ws.rs +++ b/crates/agent-server/src/ws.rs @@ -39,6 +39,11 @@ pub struct AppState { /// Гейт підписаних approvals (access.md); pubkey-кеш наповнює relay-міст. pub approvals: Arc, graph: Option, + /// Ручка пробудження orchestrator-ролі: relay push «є нові події у + /// задачі X» має будити хост негайно, а не чекати fallback-таймера + /// (runtime.md, «Wake: push замість polling»). `None` — orchestrator + /// у цьому процесі не запущений. + wake: std::sync::Mutex>>, /// Активні інтерактивні run-и за node-ключем кімнати. Git-операції /// швидкі й локальні — виконуються під локом (spawn_blocking — TODO /// разом із віддаленими remote). @@ -55,6 +60,24 @@ impl AppState { ) } + /// Реєструє ручку пробудження orchestrator-ролі. + pub fn set_wake(&self, signaller: std::sync::Arc) { + if let Ok(mut slot) = self.wake.lock() { + *slot = Some(signaller); + } + } + + /// Будить orchestrator-роль, якщо вона є в цьому процесі. Тихо + /// нічого не робить інакше — relay-міст не має знати, чи хост + /// виконує ще й оркестрацію. + pub fn wake_orchestrator(&self) { + if let Ok(slot) = self.wake.lock() { + if let Some(flag) = slot.as_ref() { + flag.store(true, std::sync::atomic::Ordering::SeqCst); + } + } + } + /// Конструктор зі спільними частинами — коли sessions/gate потрібні /// runner-фабриці ДО створення AppState (approval-гейт тулів). pub fn from_parts( @@ -68,6 +91,7 @@ impl AppState { runner, token, approvals, + wake: std::sync::Mutex::new(None), graph: None, runs: tokio::sync::Mutex::new(HashMap::new()), } diff --git a/docs/conformance.md b/docs/conformance.md index 88e092a..f533d7d 100644 --- a/docs/conformance.md +++ b/docs/conformance.md @@ -46,13 +46,13 @@ | Composite-агрегація вгору | РЕАЛІЗОВАНО | `signal.rs` `propagate_composite` | — | | Протокол spawn | РЕАЛІЗОВАНО | `spawn.rs` `spawn_approve`/`publish_spawn`, `publish.rs` `publish_lifecycle` | — (`plan_reject_max` закрито через `unresolvable`) | | Git-протокол `invalidate`/`kill` + re-run семантика | РЕАЛІЗОВАНО | `lifecycle.rs` `publish_mutation`/`stop`/`reconcile_after_rerun`; CLI — `mt stop` | — | -| Оркестрація `run --auto` | ЧАСТКОВО | `orchestrate.rs` `run_auto` | Continuous backfill і rescan на кожній події — є; лишається wake (подієвий запуск) — orchestrator-роль | +| Оркестрація `run --auto` | РЕАЛІЗОВАНО | `orchestrate.rs` `run_auto`; подієвий запуск — `agent-server/orchestrator.rs` | — | | Worktree lifecycle | РЕАЛІЗОВАНО | `worktree.rs` | — | | Git-межа (`gix` + вузький shim) | РЕАЛІЗОВАНО | `git/` | — | | Аудит-цикл: вердикт, clarification, amend, `audit_failed_streak` | РЕАЛІЗОВАНО | `audit.rs`; CLI — `mt verdict`/`mt clarify`/`mt amend` | — | -| Аудитор як актор (`mt run --actor auditor`, `audit_model`) | ЧАСТКОВО | `audit.rs` `run_auditor`/`build_auditor_prompt`; `runner.rs` `run_single_phase` | Тригери `audit_schedule_days`/`audit_on_patch` — потребують orchestrator-ролі (хвиля 3) | +| Аудитор як актор (`mt run --actor auditor`, `audit_model`) | ЧАСТКОВО | `audit.rs` `run_auditor`/`build_auditor_prompt`; `runner.rs` `run_single_phase` | Тригери `audit_schedule_days`/`audit_on_patch` — черга оркестратора наступним кроком | | EngineerAgent | РЕАЛІЗОВАНО | `runner.rs` `Actor`/`build_engineer_prompt`/`full_run_history`; CLI — `mt run --actor engineer` | — (GraphPatch реалізовано як дозволені втручання через штатні команди, окремого артефакту спека не задає) | -| `unresolvable` (3 тригери + алерт) | ЧАСТКОВО | `lib.rs` `unresolvable_reason`/`write_unresolvable`; тригери — `runner.rs` (перед комітом), `spawn.rs` `spawn_reject` | Алерт власнику (relay push) — потребує orchestrator-ролі й relay-шляху (хвиля 3) | +| `unresolvable` (3 тригери + алерт) | РЕАЛІЗОВАНО | `lib.rs` `unresolvable_reason`/`write_unresolvable`; алерт — `agent-server/orchestrator.rs` `pending_alerts` | — (доставка алерту назовні — push-транспорт M2) | | Recurrence | ВІДСУТНЄ | — | Уся глава `recurrence.md` | | Secrets broker / sandbox `skill_profiles` | ВІДСУТНЄ | — | `a.md.secrets` не інжектиться, allowlist немає | | `.mt.json` — дефолти для реалізованого | РЕАЛІЗОВАНО | `config.rs` `config_defaults` | — (ключі нереалізованих фіч свідомо без дефолтів, див. «Закриті питання») | @@ -71,7 +71,7 @@ | Handoff між хостами | ЧАСТКОВО | `agent-server/graph.rs`, `ws.rs` | Немає `HandoffRequest` як події й доставки через relay; немає checkpoint-режиму | | Approvals-гейт mid-run | ЧАСТКОВО | `agent-server/approvals_gate.rs` | **Немає матеріалізації підпису в `## Approvals`** — це блокує demo-критерій M2 | | ACP-транспорт | РЕАЛІЗОВАНО | `agent-core/acp.rs` | `mcpServers` жорстко порожній | -| Orchestrator-роль у agent-server + wake | ВІДСУТНЄ | — | Скан/dispatch/GC/алерти при wake | +| Orchestrator-роль у agent-server + wake | РЕАЛІЗОВАНО | `agent-server/orchestrator.rs` `Wake`/`Orchestrator::tick`; relay push → `AppState::wake_orchestrator` | — (тригери аудиту за розкладом — окремий рядок) | | `client_kind: mt-dashboard` | ВІДСУТНЄ | — | Типи подій є, ніхто не емітить | | Surface-профілі, MCP, preview, `ContextSelected` | ВІДСУТНЄ | — | Уся глава `surfaces.md` | From f13bd0b72d43def8993fed737b3104c458eec515 Mon Sep 17 00:00:00 2001 From: vitaliytv Date: Tue, 11 Aug 2026 09:04:47 +0300 Subject: [PATCH 2/2] =?UTF-8?q?feat:=20=D1=82=D1=80=D0=B8=D0=B3=D0=B5?= =?UTF-8?q?=D1=80=D0=B8=20=D0=B0=D1=83=D0=B4=D0=B8=D1=82=D1=83=20=E2=80=94?= =?UTF-8?q?=20=D1=87=D0=B5=D1=80=D0=B3=D0=B0=20=D0=BE=D1=80=D0=BA=D0=B5?= =?UTF-8?q?=D1=81=D1=82=D1=80=D0=B0=D1=82=D0=BE=D1=80=D0=B0=20=D1=96=20aud?= =?UTF-8?q?it=5Fon=5Fpatch?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Хвиля 3, продовження orchestrator-ролі. graph.md називає три тригери аудиту: mt audit, audit_schedule_days, audit_on_patch. Був лише перший. Черга (audit_schedule_days). Оркестратор на кожному прокиданні розбирає вузли у стані pending-audit і віддає їх аудиторові. Це той самий тригер, але подієвий, а не за таймером: цикл, відкритий сигналом mt audit, не має чекати наступної доби, щоб хтось на нього подивився — саме через це аудит у спеці названо async ЧЕРГОЮ, а не періодичною задачею. Порядок черги детермінований (за шляхом): два хости, які прокинулись одночасно, беруться за неї однаково, а не змагаються хаотично. audit_on_patch. Вузол, який пропатчили й перезапустили, тепер проходить аудит навіть якщо в контракті стоїть optional: патч змінює саме те, за чим оцінювали результат, тож попереднє «зійшло і так» не переноситься. Слід патчу — архів history/*-invalidate/. Виняток: audit: off патчем НЕ піднімається. Вимкнений аудит — свідоме рішення автора контракту, і система не має його переголошувати. Тести: 4 — черга збирає лише відкриті цикли (включно з вкладеними) у стабільному порядку, порожня черга без циклів, патч піднімає optional до required, off лишається off. 364 passed, clippy чистий. Co-Authored-By: Claude Fable 5 --- crates/agent-server/src/orchestrator.rs | 62 +++++++++++++++++++- crates/mt-core/src/signal.rs | 75 ++++++++++++++++++++++++- docs/conformance.md | 6 +- 3 files changed, 137 insertions(+), 6 deletions(-) diff --git a/crates/agent-server/src/orchestrator.rs b/crates/agent-server/src/orchestrator.rs index 0dfec2a..1ab1554 100644 --- a/crates/agent-server/src/orchestrator.rs +++ b/crates/agent-server/src/orchestrator.rs @@ -27,6 +27,8 @@ pub struct TickReport { pub dispatched: Vec<(String, String)>, /// Вузли, що вперше стали `unresolvable` — алерт власнику. pub alerts: Vec, + /// Вузли з відкритим аудит-циклом, які цей tick віддав аудиторові. + pub audited: Vec, /// Прибрані worktree. pub pruned: Vec, /// Помилки, які не мають валити цикл (наступний wake спробує знову). @@ -95,6 +97,26 @@ impl Wake { } } +/// Черга аудиту (graph.md, «Аудит (async черга)»): вузли у стані +/// `pending-audit`. Саме її розбирає orchestrator на кожному прокиданні — +/// це і є тригер `audit_schedule_days`, тільки подієвий, а не за таймером: +/// цикл, відкритий сигналом `mt audit`, не має чекати наступної доби. +/// +/// Порядок детермінований (за шляхом) — щоб два хости, які прокинулись +/// одночасно, бралися за чергу однаково, а не змагались хаотично. +pub fn audit_queue(nodes: &[mt_core::TaskNode]) -> Vec { + let mut out = Vec::new(); + let mut stack: Vec<&mt_core::TaskNode> = nodes.iter().collect(); + while let Some(node) = stack.pop() { + if node.state == mt_core::TaskState::PendingAudit { + out.push(node.path.clone()); + } + stack.extend(node.children.iter()); + } + out.sort(); + out +} + fn file_mtime(path: &Path) -> Option { std::fs::metadata(path).ok()?.modified().ok() } @@ -168,7 +190,20 @@ impl Orchestrator { } match mt_core::scan_tasks_with_claims(self.tasks_dir.clone(), Vec::new()) { - Ok(nodes) => report.alerts = self.pending_alerts(&nodes), + Ok(nodes) => { + report.alerts = self.pending_alerts(&nodes); + let queue = audit_queue(&nodes); + for path in queue { + match mt_core::audit::run_auditor( + &self.tasks_dir, + &path, + &mt_core::config::agent_cli_env_from_process(), + ) { + Ok(_) => report.audited.push(path), + Err(error) => report.errors.push(format!("audit {path}: {error}")), + } + } + } Err(error) => report.errors.push(format!("scan: {error}")), } @@ -232,6 +267,31 @@ mod tests { assert_eq!(orch.pending_alerts(&[parent]), ["parent/child"]); } + #[test] + fn audit_queue_collects_open_cycles_deterministically() { + let mut parent = node("b-parent", mt_core::TaskState::Spawned); + parent + .children + .push(node("b-parent/child", mt_core::TaskState::PendingAudit)); + let nodes = vec![ + node("a-solo", mt_core::TaskState::PendingAudit), + parent, + node("c-done", mt_core::TaskState::Resolved), + ]; + // Лише відкриті цикли, і в стабільному порядку — щоб два хости, + // які прокинулись разом, бралися за чергу однаково. + assert_eq!(audit_queue(&nodes), ["a-solo", "b-parent/child"]); + } + + #[test] + fn audit_queue_is_empty_without_open_cycles() { + let nodes = vec![ + node("solo", mt_core::TaskState::Waiting), + node("done", mt_core::TaskState::Resolved), + ]; + assert!(audit_queue(&nodes).is_empty()); + } + #[test] fn explicit_signal_wakes_immediately() { let tmp = tempfile::tempdir().unwrap(); diff --git a/crates/mt-core/src/signal.rs b/crates/mt-core/src/signal.rs index e28f01a..3a10614 100644 --- a/crates/mt-core/src/signal.rs +++ b/crates/mt-core/src/signal.rs @@ -245,7 +245,7 @@ fn check_contract_unchanged(dir: &Path) -> Result<(), String> { /// Політика аудиту вузла: frontmatter `audit:` task.md (required|optional|off). fn audit_policy(dir: &Path) -> String { - fs::read_to_string(dir.join("task.md")) + let declared = fs::read_to_string(dir.join("task.md")) .ok() .map(|c| parse_front_matter(&c)) .and_then(|fm| { @@ -253,7 +253,44 @@ fn audit_policy(dir: &Path) -> String { .and_then(serde_json::Value::as_str) .map(String::from) }) - .unwrap_or_else(|| "optional".to_string()) + .unwrap_or_else(|| "optional".to_string()); + + // `audit_on_patch` (graph.md, тригери аудиту): вузол, який пропатчили й + // перезапустили, проходить аудит навіть якщо в контракті стоїть + // `optional`. Патч змінює саме те, за чим оцінювали результат, тож + // довіра до попереднього «зійшло і так» не переноситься. + // `off` не піднімається: вимкнений аудит — свідоме рішення автора. + if declared == "optional" && was_patched(dir) && audit_on_patch_enabled(dir) { + return "required".to_string(); + } + declared +} + +/// Чи вузол колись інвалідували — у ньому є архів `history/*-invalidate/`. +fn was_patched(dir: &Path) -> bool { + fs::read_dir(dir.join("history")) + .map(|entries| { + entries.flatten().any(|e| { + e.file_name() + .to_string_lossy() + .ends_with("-invalidate") + }) + }) + .unwrap_or(false) +} + +fn audit_on_patch_enabled(dir: &Path) -> bool { + let Some(project_root) = dir.parent().and_then(Path::parent) else { + return false; + }; + crate::config::merge_config( + fs::read_to_string(project_root.join(".mt.json")) + .ok() + .as_deref(), + ) + .get("audit_on_patch") + .and_then(serde_json::Value::as_bool) + .unwrap_or(true) } fn signal_success( @@ -635,6 +672,40 @@ mod tests { assert!(!tmp.path().join("solo/run_001.md").exists()); } + #[test] + fn patched_node_requires_audit_even_when_optional() { + // graph.md, тригер `audit_on_patch`: патч змінює саме те, за чим + // оцінювали результат, тож попереднє «зійшло і так» не переноситься. + let tmp = tempfile::tempdir().unwrap(); + node(tmp.path(), "solo", TASK_WITH_CHECK); // audit не вказано → optional + let dir = tmp.path().join("solo"); + let root = tmp.path().to_string_lossy().into_owned(); + write_fact(&root, "solo", "Зроблено.", None).unwrap(); + assert!(done(&root, "solo", "agent").is_ok(), "до патчу — optional"); + + // Слід інвалідації робить аудит обов'язковим. + fs::create_dir_all(dir.join("history/20260810-120000-invalidate")).unwrap(); + write_fact(&root, "solo", "Перероблено.", None).unwrap(); + assert!(done(&root, "solo", "agent") + .unwrap_err() + .contains("приймає лише сигнал audit")); + assert!(audit(&root, "solo", "agent").is_ok()); + } + + #[test] + fn audit_off_is_not_raised_by_patch() { + // Вимкнений аудит — свідоме рішення автора контракту, патч його не + // скасовує. + let tmp = tempfile::tempdir().unwrap(); + let task = TASK_WITH_CHECK.replace("---\n\n## Task", "audit: off\n---\n\n## Task"); + node(tmp.path(), "solo", &task); + let dir = tmp.path().join("solo"); + fs::create_dir_all(dir.join("history/20260810-120000-invalidate")).unwrap(); + let root = tmp.path().to_string_lossy().into_owned(); + write_fact(&root, "solo", "Зроблено.", None).unwrap(); + assert!(done(&root, "solo", "agent").is_ok()); + } + #[test] fn audit_opens_cycle_and_required_blocks_done() { let tmp = tempfile::tempdir().unwrap(); diff --git a/docs/conformance.md b/docs/conformance.md index f533d7d..02f2508 100644 --- a/docs/conformance.md +++ b/docs/conformance.md @@ -17,7 +17,7 @@ | Мілстоун | Стан | Головне, чого бракує | | --- | --- | --- | | M0 — dogfood ядра | цикл замкнено | recurrence; secrets broker; телеметрія вартості | -| M1 — agent-server | значною мірою | orchestrator-роль і wake, backpressure за спекою, глибокий реплей, злиття `agent-cli` у `mt serve\|attach` | +| M1 — agent-server | значною мірою | backpressure за спекою, глибокий реплей, злиття `agent-cli` у `mt serve\|attach` | | M2 — mission control | частково | матеріалізація підпису в `## Approvals`, персистентний store і auth, реальний push-транспорт, handoff між машинами через relay, presence | | M3 — dashboard і поверхні | не починався | surface-профілі, MCP-сервери, preview/`ContextSelected`, `client_kind: mt-dashboard` | | M4 — файловий шар i18n | не починався | `refs/mt/i18n`, worktree-матеріалізація, write path у base, lazy-мови (`layers/` — суміжна задача, інший контейнер і конфіг) | @@ -50,7 +50,7 @@ | Worktree lifecycle | РЕАЛІЗОВАНО | `worktree.rs` | — | | Git-межа (`gix` + вузький shim) | РЕАЛІЗОВАНО | `git/` | — | | Аудит-цикл: вердикт, clarification, amend, `audit_failed_streak` | РЕАЛІЗОВАНО | `audit.rs`; CLI — `mt verdict`/`mt clarify`/`mt amend` | — | -| Аудитор як актор (`mt run --actor auditor`, `audit_model`) | ЧАСТКОВО | `audit.rs` `run_auditor`/`build_auditor_prompt`; `runner.rs` `run_single_phase` | Тригери `audit_schedule_days`/`audit_on_patch` — черга оркестратора наступним кроком | +| Аудитор як актор + тригери аудиту | РЕАЛІЗОВАНО | `audit.rs` `run_auditor`; черга — `orchestrator.rs` `audit_queue`; `audit_on_patch` — `signal.rs` `audit_policy` | — | | EngineerAgent | РЕАЛІЗОВАНО | `runner.rs` `Actor`/`build_engineer_prompt`/`full_run_history`; CLI — `mt run --actor engineer` | — (GraphPatch реалізовано як дозволені втручання через штатні команди, окремого артефакту спека не задає) | | `unresolvable` (3 тригери + алерт) | РЕАЛІЗОВАНО | `lib.rs` `unresolvable_reason`/`write_unresolvable`; алерт — `agent-server/orchestrator.rs` `pending_alerts` | — (доставка алерту назовні — push-транспорт M2) | | Recurrence | ВІДСУТНЄ | — | Уся глава `recurrence.md` | @@ -71,7 +71,7 @@ | Handoff між хостами | ЧАСТКОВО | `agent-server/graph.rs`, `ws.rs` | Немає `HandoffRequest` як події й доставки через relay; немає checkpoint-режиму | | Approvals-гейт mid-run | ЧАСТКОВО | `agent-server/approvals_gate.rs` | **Немає матеріалізації підпису в `## Approvals`** — це блокує demo-критерій M2 | | ACP-транспорт | РЕАЛІЗОВАНО | `agent-core/acp.rs` | `mcpServers` жорстко порожній | -| Orchestrator-роль у agent-server + wake | РЕАЛІЗОВАНО | `agent-server/orchestrator.rs` `Wake`/`Orchestrator::tick`; relay push → `AppState::wake_orchestrator` | — (тригери аудиту за розкладом — окремий рядок) | +| Orchestrator-роль у agent-server + wake | РЕАЛІЗОВАНО | `agent-server/orchestrator.rs` `Wake`/`Orchestrator::tick`; relay push → `AppState::wake_orchestrator` | — | | `client_kind: mt-dashboard` | ВІДСУТНЄ | — | Типи подій є, ніхто не емітить | | Surface-профілі, MCP, preview, `ContextSelected` | ВІДСУТНЄ | — | Уся глава `surfaces.md` |