diff --git a/Cargo.lock b/Cargo.lock index 728012f..3862bc8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1520,11 +1520,13 @@ dependencies = [ "futures-util", "hex", "hmac", + "rusqlite", "santi-core", "santi-provider", "serde", "serde_json", "sha2", + "tempfile", "tokio", "toml", "tower-http", diff --git a/crates/santi-api/Cargo.toml b/crates/santi-api/Cargo.toml index 67f091e..b146b44 100644 --- a/crates/santi-api/Cargo.toml +++ b/crates/santi-api/Cargo.toml @@ -21,3 +21,7 @@ tokio.workspace = true toml.workspace = true tower-http.workspace = true utoipa.workspace = true + +[dev-dependencies] +tempfile = "3.27.0" +rusqlite.workspace = true diff --git a/crates/santi-api/src/config.rs b/crates/santi-api/src/config.rs index 3034530..a7eaaaf 100644 --- a/crates/santi-api/src/config.rs +++ b/crates/santi-api/src/config.rs @@ -277,6 +277,32 @@ pub(crate) fn santi_home() -> PathBuf { PathBuf::from(base).join(".santi") } +/// The runtime data locations, all anchored on `santi_home()` unless an explicit +/// env overrides. Pure computation — creating the dirs is the caller's job (the +/// offline ops paths read without creating; `serve` creates). Shared by `serve`, +/// `santi doctor`, and `santi inbox seed` so they never drift. +#[derive(Debug, Clone)] +pub struct RuntimePaths { + pub database_path: PathBuf, + pub runtime_root: PathBuf, + pub execution_root: PathBuf, +} + +pub fn resolve_runtime_paths() -> RuntimePaths { + let home = santi_home(); + RuntimePaths { + database_path: optional_env("SANTI_DB") + .map(PathBuf::from) + .unwrap_or_else(|| home.join("runtime").join("db")), + runtime_root: optional_env("SANTI_RUNTIME_ROOT") + .map(PathBuf::from) + .unwrap_or_else(|| home.join("runtime")), + execution_root: optional_env("SANTI_EXECUTION_ROOT") + .map(PathBuf::from) + .unwrap_or_else(|| home.join("execution")), + } +} + fn expand_home(path: &str) -> PathBuf { if let Some(rest) = path.strip_prefix("~/") && let Ok(home) = env::var("HOME") diff --git a/crates/santi-api/src/lib.rs b/crates/santi-api/src/lib.rs index 213e374..afd8dff 100644 --- a/crates/santi-api/src/lib.rs +++ b/crates/santi-api/src/lib.rs @@ -1,4 +1,5 @@ pub mod config; +pub mod ops; pub mod provider; mod bucket; diff --git a/crates/santi-api/src/ops.rs b/crates/santi-api/src/ops.rs new file mode 100644 index 0000000..277d9d8 --- /dev/null +++ b/crates/santi-api/src/ops.rs @@ -0,0 +1,258 @@ +//! Local (non-HTTP) operator commands. Unlike the client commands, these do NOT +//! reach a running runtime over HTTP — they act directly on the on-disk store +//! and runtime files. They exist for the self-upgrade lifecycle (PHASE-07), +//! where the service is often stopped, and are grouped here so `santi-api` +//! stays the single owner of path resolution (`config::resolve_runtime_paths`). + +use std::fs; + +use serde::Serialize; + +use crate::config::{self, RuntimePaths}; + +/// A read-only pre-check of the on-disk store + default soul memory. Pure reads: +/// it never opens the store (which would migrate/wipe) and is safe against a +/// live or stopped service. The soul-deep health "come up + coherent" contract +/// is confirmed functionally elsewhere (PHASE-07 open #2); this is the cheap +/// deterministic half — "is the store at the expected schema and is the default +/// soul's memory readable". +#[derive(Debug, Clone, Serialize)] +pub struct DoctorReport { + pub database_path: String, + pub database_exists: bool, + /// The DB's `user_version`, or null when the DB file does not exist yet. + pub schema_version: Option, + pub expected_schema_version: u32, + /// The DB is present and already at the version this binary expects (so a + /// start would NOT wipe/migrate). + pub schema_ok: bool, + pub default_soul_id: String, + pub memory_path: String, + pub memory_present: bool, + pub memory_readable: bool, + pub memory_bytes: u64, + /// Overall gate: schema at the expected version AND (memory absent, which is + /// a fresh soul that falls back to the encoded default, OR memory readable). + pub ok: bool, +} + +/// Run the offline pre-check against the runtime paths resolved from env. +pub fn doctor() -> Result { + doctor_at(&config::resolve_runtime_paths()) +} + +/// The pure pre-check over explicit paths (env-free, so it is unit-testable). +pub fn doctor_at(paths: &RuntimePaths) -> Result { + let database_exists = paths.database_path.exists(); + let schema_version = santi_core::read_schema_version(&paths.database_path)?; + let schema_ok = schema_version == Some(santi_core::SCHEMA_VERSION); + + let memory_path = + santi_core::soul_memory_file(&paths.runtime_root, santi_core::DEFAULT_SOUL_ID); + let memory_present = memory_path.exists(); + let (memory_readable, memory_bytes) = match fs::read(&memory_path) { + Ok(bytes) => (true, bytes.len() as u64), + Err(_) => (false, 0), + }; + // Absent memory is fine (a fresh soul falls back to the encoded default); + // present-but-unreadable is the failure — the soul's continuity would break. + let memory_ok = !memory_present || memory_readable; + + Ok(DoctorReport { + database_path: paths.database_path.display().to_string(), + database_exists, + schema_version, + expected_schema_version: santi_core::SCHEMA_VERSION, + schema_ok, + default_soul_id: santi_core::DEFAULT_SOUL_ID.to_string(), + memory_path: memory_path.display().to_string(), + memory_present, + memory_readable, + memory_bytes, + ok: schema_ok && memory_ok, + }) +} + +/// The result of an offline inbox seed. +#[derive(Debug, Clone, Serialize)] +pub struct SeedReport { + pub strand_id: String, + /// Durably enqueued (false ⟺ the inbox gate rejected it — the strand is far + /// behind; the caller should treat this as a failure). + pub accepted: bool, + pub reason: Option, +} + +/// Enqueue one `santi_system` record into a strand's durable inbox WITHOUT a +/// running service (a direct MQ producer). Used by the self-upgrade flow to seed +/// the "you were upgrading — come look" record before starting the final version, +/// so boot recovery drains it and the soul wakes into the result (PHASE-07). +/// +/// Opens the store with THIS binary (so it migrates to this binary's schema — +/// seed with the FINAL version). The target strand MUST already exist: enqueuing +/// into an unknown strand would leave an inbox row that boot recovery can never +/// turn into a turn, so we reject instead of writing an orphan. +pub fn inbox_seed(strand_id: &str, text: &str) -> Result { + inbox_seed_at(&config::resolve_runtime_paths(), strand_id, text) +} + +pub fn inbox_seed_at( + paths: &RuntimePaths, + strand_id: &str, + text: &str, +) -> Result { + let store = santi_core::SantiStore::open(&paths.database_path)?; + if store.strand(strand_id)?.is_none() { + return Err(format!("unknown strand: {strand_id}")); + } + let outcome = store.enqueue_inbox( + strand_id, + santi_core::MessageKind::SantiSystem, + santi_core::MessageContent::text(text), + )?; + Ok(match outcome { + santi_core::IngestOutcome::Accepted { strand_id } => SeedReport { + strand_id, + accepted: true, + reason: None, + }, + santi_core::IngestOutcome::Rejected { reason } => SeedReport { + strand_id: strand_id.to_string(), + accepted: false, + reason: Some(reason), + }, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::path::PathBuf; + + fn paths_under(root: &std::path::Path) -> RuntimePaths { + RuntimePaths { + database_path: root.join("runtime").join("db"), + runtime_root: root.join("runtime"), + execution_root: root.join("execution"), + } + } + + #[test] + fn healthy_when_migrated_and_memory_readable() { + let temp = tempfile::tempdir().expect("temp dir"); + let paths = paths_under(temp.path()); + // Open (migrates to the current schema) then write the default soul memory. + santi_core::SantiStore::open(&paths.database_path).expect("open store"); + let memory = santi_core::soul_memory_file(&paths.runtime_root, santi_core::DEFAULT_SOUL_ID); + std::fs::create_dir_all(memory.parent().unwrap()).unwrap(); + std::fs::write(&memory, "# I am the first secretary").unwrap(); + + let report = doctor_at(&paths).expect("doctor"); + assert!(report.ok, "expected healthy: {report:?}"); + assert!(report.schema_ok); + assert_eq!(report.schema_version, Some(santi_core::SCHEMA_VERSION)); + assert!(report.memory_present && report.memory_readable); + assert!(report.memory_bytes > 0); + } + + #[test] + fn unhealthy_when_schema_stale() { + let temp = tempfile::tempdir().expect("temp dir"); + let paths = paths_under(temp.path()); + std::fs::create_dir_all(paths.database_path.parent().unwrap()).unwrap(); + // A DB stamped at a stale version — a start WOULD wipe/migrate it. + let conn = rusqlite::Connection::open(&paths.database_path).unwrap(); + conn.pragma_update(None, "user_version", 5u32).unwrap(); + drop(conn); + + let report = doctor_at(&paths).expect("doctor"); + assert!(!report.ok); + assert!(!report.schema_ok); + assert_eq!(report.schema_version, Some(5)); + } + + #[test] + fn absent_memory_is_ok_but_absent_db_is_not() { + let temp = tempfile::tempdir().expect("temp dir"); + let paths = paths_under(temp.path()); + santi_core::SantiStore::open(&paths.database_path).expect("open store"); + // No memory file written: a fresh soul falls back to the encoded default, + // so memory-absent is healthy as long as the schema is current. + let report = doctor_at(&paths).expect("doctor"); + assert!(report.ok, "absent memory should be fine: {report:?}"); + assert!(!report.memory_present); + + // But a DB that does not exist at all is not healthy (schema unknown). + let missing = RuntimePaths { + database_path: temp.path().join("void").join("db"), + ..paths + }; + let report = doctor_at(&missing).expect("doctor"); + assert!(!report.ok); + assert_eq!(report.schema_version, None); + } + + #[test] + fn report_serializes_to_json() { + let temp = tempfile::tempdir().expect("temp dir"); + let paths = paths_under(temp.path()); + santi_core::SantiStore::open(&paths.database_path).expect("open store"); + let report = doctor_at(&paths).expect("doctor"); + let json = serde_json::to_string(&report).expect("serialize"); + assert!(json.contains("\"schema_ok\"")); + let _ = PathBuf::from(&report.database_path); + } + + #[test] + fn seed_into_existing_strand_is_drainable_on_boot() { + let temp = tempfile::tempdir().expect("temp dir"); + let paths = paths_under(temp.path()); + let strand_id = { + let store = santi_core::SantiStore::open(&paths.database_path).expect("open"); + store.create_strand().expect("create strand").id + }; + + let report = inbox_seed_at(&paths, &strand_id, "you were upgrading — come look").unwrap(); + assert!(report.accepted); + assert_eq!(report.strand_id, strand_id); + + // The seed is exactly what boot recovery scans for, and try_start_turn + // (what a boot poke calls) drains it into a real turn carrying our text. + let store = santi_core::SantiStore::open(&paths.database_path).expect("reopen"); + assert!( + store + .strands_with_pending_requests() + .unwrap() + .contains(&strand_id), + "boot recovery would re-drive this strand" + ); + // "strand_send" is exactly the trigger boot recovery uses (resume_pending). + let started = store + .try_start_turn(&strand_id, "strand_send", None) + .unwrap() + .expect("a turn starts by draining the seed"); + assert_eq!(started.drained_messages.len(), 1); + assert_eq!( + started.drained_messages[0].content_text, + "you were upgrading — come look" + ); + assert_eq!( + started.drained_messages[0].message.message_kind, + santi_core::MessageKind::SantiSystem + ); + } + + #[test] + fn seed_into_unknown_strand_errors_without_writing_orphan() { + let temp = tempfile::tempdir().expect("temp dir"); + let paths = paths_under(temp.path()); + santi_core::SantiStore::open(&paths.database_path).expect("open"); + + let err = inbox_seed_at(&paths, "ss_does_not_exist", "x").unwrap_err(); + assert!(err.contains("unknown strand"), "got: {err}"); + + // Nothing was enqueued — boot recovery finds no orphan to choke on. + let store = santi_core::SantiStore::open(&paths.database_path).expect("reopen"); + assert!(store.strands_with_pending_requests().unwrap().is_empty()); + } +} diff --git a/crates/santi-api/src/server.rs b/crates/santi-api/src/server.rs index b178119..d3168ba 100644 --- a/crates/santi-api/src/server.rs +++ b/crates/santi-api/src/server.rs @@ -1,4 +1,4 @@ -use std::{convert::Infallible, env, fs, net::SocketAddr, path::PathBuf, sync::Arc}; +use std::{convert::Infallible, env, fs, net::SocketAddr, sync::Arc}; use crate::{ config, provider, @@ -36,25 +36,20 @@ pub fn export_openapi_json() -> Result { pub async fn serve(config: config::ConfigService) -> Result<(), String> { let provider = provider::from_config(config.provider_config()?); - // Defaults anchor on the santi home (`SANTI_HOME`, else `~/.santi`); explicit - // env always overrides. The data dirs are created so a zero-config run works. - let home = config::santi_home(); - let database_path = env::var("SANTI_DB") - .unwrap_or_else(|_| home.join("runtime").join("db").display().to_string()); - let runtime_root = env::var("SANTI_RUNTIME_ROOT") - .unwrap_or_else(|_| home.join("runtime").display().to_string()); - let execution_root = env::var("SANTI_EXECUTION_ROOT") - .unwrap_or_else(|_| home.join("execution").display().to_string()); - if let Some(parent) = PathBuf::from(&database_path).parent() { + // Paths anchor on the santi home (`SANTI_HOME`, else `~/.santi`); explicit env + // always overrides (see `resolve_runtime_paths`). The data dirs are created + // here so a zero-config run works (the offline ops paths only read). + let paths = config::resolve_runtime_paths(); + if let Some(parent) = paths.database_path.parent() { fs::create_dir_all(parent).map_err(|error| error.to_string())?; } - fs::create_dir_all(&runtime_root).map_err(|error| error.to_string())?; - fs::create_dir_all(&execution_root).map_err(|error| error.to_string())?; + fs::create_dir_all(&paths.runtime_root).map_err(|error| error.to_string())?; + fs::create_dir_all(&paths.execution_root).map_err(|error| error.to_string())?; let service = SantiService::open( SantiServiceConfig { - database_path, - runtime_root, - execution_root, + database_path: paths.database_path.display().to_string(), + runtime_root: paths.runtime_root.display().to_string(), + execution_root: paths.execution_root.display().to_string(), bind_addr: Some(bind_addr_string()), }, provider, @@ -78,9 +73,61 @@ pub async fn serve(config: config::ConfigService) -> Result<(), String> { // Liveness: re-drive any requests stranded by a previous crash. service.resume_pending(); println!("santi-api listening on http://{address}"); + // Graceful shutdown (PHASE-07): on SIGTERM/Ctrl-C, latch the service so no + // new turns start (inbox consumption pauses; ingest still enqueues durably), + // let axum drain in-flight HTTP, then wait out the in-flight turn before + // exiting. The external upgrade flow owns the hard bound (SIGKILL after its + // timeout); this is the cooperative half. + let shutdown_signal = { + let service = service.clone(); + async move { + wait_for_shutdown_signal().await; + println!("santi-api: shutdown signal received — quiescing (no new turns)"); + service.begin_shutdown(); + } + }; + let drainer = service.clone(); axum::serve(listener, router(service, api_key)) + .with_graceful_shutdown(shutdown_signal) .await - .map_err(|error| error.to_string()) + .map_err(|error| error.to_string())?; + drainer.drain_running_turns(shutdown_grace()).await; + println!("santi-api: drained; exiting"); + Ok(()) +} + +/// Resolve on the shutdown signal: SIGTERM (systemd/`systemctl stop`) or Ctrl-C. +async fn wait_for_shutdown_signal() { + #[cfg(unix)] + { + use tokio::signal::unix::{SignalKind, signal}; + let mut term = match signal(SignalKind::terminate()) { + Ok(term) => term, + Err(error) => { + eprintln!("santi-api: cannot install SIGTERM handler: {error}"); + return; + } + }; + tokio::select! { + _ = term.recv() => {} + _ = tokio::signal::ctrl_c() => {} + } + } + #[cfg(not(unix))] + { + let _ = tokio::signal::ctrl_c().await; + } +} + +/// How long the service waits for the in-flight turn to finish on shutdown. +/// `SANTI_SHUTDOWN_GRACE_SECS`, default 600s (turns can run minutes). The systemd +/// unit's `TimeoutStopSec` must be at least this so systemd does not SIGKILL first. +fn shutdown_grace() -> std::time::Duration { + let secs = env::var("SANTI_SHUTDOWN_GRACE_SECS") + .ok() + .and_then(|value| value.parse::().ok()) + .unwrap_or(600); + std::time::Duration::from_secs(secs) } fn bind_addr_string() -> String { diff --git a/crates/santi-core/src/lib.rs b/crates/santi-core/src/lib.rs index 8060d66..f8aad1d 100644 --- a/crates/santi-core/src/lib.rs +++ b/crates/santi-core/src/lib.rs @@ -11,7 +11,9 @@ pub use model::*; pub use object_store::{LocalObjectStore, ObjectBucket, ObjectMeta, ObjectPayload, ObjectUri}; pub use santi_provider::ProviderItem; pub use service::{SantiService, SantiServiceConfig}; -pub use store::SantiStore; +pub use store::{ + DEFAULT_SOUL_ID, SCHEMA_VERSION, SantiStore, read_schema_version, soul_memory_file, +}; pub use workspace_uri::{ MEMORY_FILE, SOUL_WORKSPACE_URI, STRAND_WORKSPACE_URI, WorkspaceRoot, WorkspaceUri, parse_workspace_uri, soul_memory_uri, strand_memory_uri, workspace_uri, diff --git a/crates/santi-core/src/service.rs b/crates/santi-core/src/service.rs index 2944e43..2138347 100644 --- a/crates/santi-core/src/service.rs +++ b/crates/santi-core/src/service.rs @@ -9,7 +9,11 @@ use futures_util::StreamExt; use santi_provider::{ProviderClient, ProviderEvent, ProviderRequest}; use std::{ collections::HashMap, - sync::{Arc, Mutex}, + sync::{ + Arc, Mutex, + atomic::{AtomicBool, Ordering}, + }, + time::{Duration, Instant}, }; use tokio::sync::broadcast; @@ -34,6 +38,11 @@ pub struct SantiService { pub(crate) config: SantiServiceConfig, material_cache: Arc>>, stream_events: broadcast::Sender, + /// Graceful-shutdown latch (PHASE-07): once set, `poke` refuses to START new + /// turns, so inbox CONSUMPTION pauses while ingest keeps durably enqueuing + /// (the inbox is an MQ — we stop consuming, never producing). The in-flight + /// turn is left to finish; `drain_running_turns` waits it out. + shutting_down: Arc, } type MaterialCacheKey = (String, MaterialKind); @@ -67,9 +76,48 @@ impl SantiService { config, material_cache: Arc::new(Mutex::new(HashMap::new())), stream_events: broadcast::channel(1024).0, + shutting_down: Arc::new(AtomicBool::new(false)), }) } + /// Begin a graceful shutdown: stop consuming the inbox (no new turns start). + /// Idempotent. Ingest still durably enqueues; the in-flight turn finishes. + pub fn begin_shutdown(&self) { + self.shutting_down.store(true, Ordering::SeqCst); + } + + pub fn is_shutting_down(&self) -> bool { + self.shutting_down.load(Ordering::SeqCst) + } + + /// Wait until no turn is `running` (the in-flight turn finished), or until + /// `cap` elapses. Called after the HTTP server has stopped accepting, so no + /// new turns can appear once shutdown has begun. On cap-timeout it returns + /// anyway: the still-running turn will be reconciled to `interrupted` on the + /// next boot (honest occurrence), and the external upgrade flow's own bound + /// (SIGKILL) is the hard stop. + pub async fn drain_running_turns(&self, cap: Duration) { + let start = Instant::now(); + loop { + match self.store.running_turn_count() { + Ok(0) => return, + Ok(remaining) => { + if start.elapsed() >= cap { + eprintln!( + "santi: shutdown drain cap reached with {remaining} turn(s) still running; leaving them to boot-recovery" + ); + return; + } + tokio::time::sleep(Duration::from_millis(200)).await; + } + Err(error) => { + eprintln!("santi: shutdown drain scan failed: {error}"); + return; + } + } + } + } + /// Re-drive strands left "behind" by a crash (their inbox durably holds /// content nobody ever drained). Liveness only — no retry of /// attempted/failed turns. Call once at server startup (inside the tokio @@ -340,6 +388,13 @@ impl SantiService { /// coalesces) or there is nothing pending. The atomic guard in /// `try_start_turn` keeps "one present per thread of experience". fn poke(&self, strand_id: &str, trigger_type: &str) -> DrivenTurn { + // Graceful shutdown: pause CONSUMPTION. The content stays durably in the + // inbox (ingest already enqueued it) and boot recovery drains it on the + // next start. This also stops the completion re-poke from spawning a + // follow-on turn, so an in-flight turn can finish and the strand quiesce. + if self.is_shutting_down() { + return None; + } match self.store.try_start_turn(strand_id, trigger_type, None) { Ok(Some(started)) => { for message in started.drained_messages.iter().cloned() { diff --git a/crates/santi-core/src/service/tools.rs b/crates/santi-core/src/service/tools.rs index dbc7616..02e010b 100644 --- a/crates/santi-core/src/service/tools.rs +++ b/crates/santi-core/src/service/tools.rs @@ -128,7 +128,9 @@ impl SantiService { } pub(super) fn soul_memory_file(&self, soul_id: &str) -> PathBuf { - self.soul_memory_dir(soul_id).join("MEMORY.md") + // Delegate to the free function so offline ops (`santi doctor`) and the + // running service always resolve the same path. + crate::store::soul_memory_file(self.runtime_root(), soul_id) } pub(super) fn strand_memory_dir(&self, strand_id: &str) -> PathBuf { diff --git a/crates/santi-core/src/store.rs b/crates/santi-core/src/store.rs index 97c8015..5e3bf06 100644 --- a/crates/santi-core/src/store.rs +++ b/crates/santi-core/src/store.rs @@ -21,8 +21,14 @@ use db::*; use rows::{actor_type_db, collect_rows, map_webhook_row, message_state_db}; use schema::SCHEMA; -const SANTI_SCHEMA_VERSION: u32 = 19; -const DEFAULT_SOUL_ID: &str = "soul_default"; +/// The schema version this binary expects. On open, a DB at any other version +/// is wiped + rebuilt (beta: no back-compat migrations yet — see PHASE-07 crux #5). +/// Public so ops paths (`santi doctor`) can compare a DB's `user_version` to it +/// WITHOUT opening the store (which would migrate/wipe). +pub const SCHEMA_VERSION: u32 = 19; +/// The default soul's id. Public so offline ops (doctor/seed) can address it +/// without a running service. +pub const DEFAULT_SOUL_ID: &str = "soul_default"; /// The runtime's one system actor identity. No account/user: every non-soul /// actor speaks as `system`, whether it's a runtime-authored notice (kind /// santi_system) or opaque world-inbound content (kind text) — the sender's @@ -52,6 +58,37 @@ pub struct StartedTurn { pub drained_messages: Vec, } +/// Read a DB's `user_version` WITHOUT opening the store (which would migrate +/// and, on a version mismatch, WIPE). `Ok(None)` when the file does not exist +/// yet (a fresh instance). Read-only: safe to run against a live service (WAL) +/// or a stopped one. Used by the offline pre-check `santi doctor`. +pub fn read_schema_version(path: impl AsRef) -> Result, String> { + if !path.as_ref().exists() { + return Ok(None); + } + let conn = Connection::open_with_flags( + path, + rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_URI, + ) + .map_err(|error| error.to_string())?; + conn.query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)) + .map(Some) + .map_err(|error| error.to_string()) +} + +/// The default soul's memory file, given a runtime root: +/// `/souls//memory/MEMORY.md`. A free function so offline +/// ops can compute the path without a `SantiService`; the running service's +/// `soul_memory_file` (service/tools.rs) delegates here to stay in lockstep. +pub fn soul_memory_file(runtime_root: impl AsRef, soul_id: &str) -> std::path::PathBuf { + runtime_root + .as_ref() + .join("souls") + .join(soul_id) + .join("memory") + .join(crate::workspace_uri::MEMORY_FILE) +} + impl SantiStore { pub fn open(path: impl AsRef) -> Result { if let Some(parent) = path.as_ref().parent() { @@ -71,7 +108,7 @@ impl SantiStore { let version = conn .query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)) .map_err(|error| error.to_string())?; - if version != SANTI_SCHEMA_VERSION { + if version != SCHEMA_VERSION { conn.execute_batch( r#" DROP TABLE IF EXISTS response_stream_deltas; @@ -107,7 +144,7 @@ impl SantiStore { } conn.execute_batch(SCHEMA) .map_err(|error| error.to_string())?; - conn.pragma_update(None, "user_version", SANTI_SCHEMA_VERSION) + conn.pragma_update(None, "user_version", SCHEMA_VERSION) .map_err(|error| error.to_string())?; Ok(()) } diff --git a/crates/santi-core/src/store/runtime.rs b/crates/santi-core/src/store/runtime.rs index 9121d7a..c782ed2 100644 --- a/crates/santi-core/src/store/runtime.rs +++ b/crates/santi-core/src/store/runtime.rs @@ -308,6 +308,19 @@ impl SantiStore { .map_err(|error| error.to_string()) } + /// How many turns are currently `running`. Used by graceful shutdown to + /// wait for the in-flight turn to finish (0 ⟺ the strand has quiesced). + pub fn running_turn_count(&self) -> Result { + let conn = self.conn.lock().unwrap(); + conn.query_row( + "SELECT COUNT(*) FROM turns WHERE status = 'running'", + [], + |row| row.get::<_, i64>(0), + ) + .map(|count| count as usize) + .map_err(|error| error.to_string()) + } + /// Strands that are "behind" (their inbox is non-empty). Used on boot to /// re-drive durable requests stranded by a crash (liveness) — the inbox /// itself is durable, so this is exactly "which strands still have diff --git a/crates/santi-core/tests/service.rs b/crates/santi-core/tests/service.rs index 1a26b17..c554809 100644 --- a/crates/santi-core/tests/service.rs +++ b/crates/santi-core/tests/service.rs @@ -418,6 +418,68 @@ async fn boot_recovery_drains_stranded_inbox_entries() { ); } +/// Graceful shutdown pauses inbox CONSUMPTION (no new turns start) while ingest +/// keeps PRODUCING durably; a later fresh boot then drains what queued up. This +/// is the enabling behavior for self-upgrade: quiesce → stop → swap → start → +/// boot recovery wakes the soul on whatever queued during the window (PHASE-07). +#[tokio::test] +async fn graceful_shutdown_pauses_consumption_but_not_production() { + let temp = tempfile::tempdir().expect("temp dir"); + let config = SantiServiceConfig { + database_path: temp.path().join("santi.sqlite").display().to_string(), + runtime_root: temp.path().join("runtime").display().to_string(), + execution_root: temp.path().join("execution").display().to_string(), + bind_addr: Some("127.0.0.1:0".to_string()), + }; + let provider = Arc::new(FakeProvider::default()); + + // A quiescing service: it accepts (durably enqueues) but starts NO turn. + let strand_id = { + let service = SantiService::open(config.clone(), provider.clone()).expect("open service"); + service.begin_shutdown(); + assert!(service.is_shutting_down()); + let outcome = service + .ingest_external_event( + "soul_default", + "shutdown:quiesce", + "arrived while quiescing".to_string(), + ) + .expect("ingest during shutdown"); + match outcome { + santi_core::IngestOutcome::Accepted { strand_id } => strand_id, + other => panic!("expected accepted, got {other:?}"), + } + }; + + // Consumption paused: no turn was started. Production intact: the record is + // durably queued (exactly what boot recovery scans for). + let store = SantiStore::open(&config.database_path).expect("open store directly"); + assert_eq!( + store.running_turn_count().expect("count"), + 0, + "shutdown must not start a turn" + ); + assert!( + store + .strands_with_pending_requests() + .expect("pending") + .contains(&strand_id), + "the ingested record must still be durably queued" + ); + drop(store); + + // A fresh service (not shutting down) drains the backlog on boot. + let service = SantiService::open(config, provider.clone()).expect("reopen service"); + service.resume_pending(); + let runtime = wait_for_any_completed_turn(&service, &strand_id).await; + assert!( + runtime + .messages + .iter() + .any(|message| message.content_text == "arrived while quiescing") + ); +} + #[tokio::test] async fn send_strand_targets_the_strands_own_soul() { let temp = tempfile::tempdir().expect("temp dir"); diff --git a/crates/santi-core/tests/store.rs b/crates/santi-core/tests/store.rs index 3aa9e91..63ba9a7 100644 --- a/crates/santi-core/tests/store.rs +++ b/crates/santi-core/tests/store.rs @@ -539,3 +539,52 @@ fn create_soul_and_label_anchoring() { .is_err() ); } + +#[test] +fn read_schema_version_none_when_db_absent() { + let temp = tempfile::tempdir().expect("temp dir"); + let missing = temp.path().join("nope.sqlite"); + assert_eq!( + santi_core::read_schema_version(&missing).expect("read"), + None + ); +} + +#[test] +fn read_schema_version_is_readonly_and_matches_after_open() { + let temp = tempfile::tempdir().expect("temp dir"); + let db = temp.path().join("santi.sqlite"); + + // A DB stamped at a stale version: the probe reports it AS-IS and, crucially, + // does NOT migrate/wipe it (unlike SantiStore::open). + { + let conn = Connection::open(&db).expect("open sqlite"); + conn.pragma_update(None, "user_version", 5u32) + .expect("stamp version"); + } + assert_eq!( + santi_core::read_schema_version(&db).expect("read"), + Some(5), + "probe must report the stored version, not migrate it" + ); + assert_eq!( + santi_core::read_schema_version(&db).expect("read again"), + Some(5), + "a second probe still sees the stale version — the first was read-only" + ); + + // Opening the store DOES migrate to the runtime's version. + let store = SantiStore::open(&db).expect("open store"); + drop(store); + assert_eq!( + santi_core::read_schema_version(&db).expect("read post-open"), + Some(santi_core::SCHEMA_VERSION) + ); +} + +#[test] +fn soul_memory_file_composes_under_runtime_root() { + let path = santi_core::soul_memory_file("/srv/santi/runtime", "soul_default"); + assert!(path.ends_with("souls/soul_default/memory/MEMORY.md")); + assert!(path.starts_with("/srv/santi/runtime")); +} diff --git a/crates/santi/src/main.rs b/crates/santi/src/main.rs index 286b605..59774ae 100644 --- a/crates/santi/src/main.rs +++ b/crates/santi/src/main.rs @@ -64,6 +64,13 @@ enum Command { #[arg(trailing_var_arg = true, allow_hyphen_values = true)] args: Vec, }, + /// Offline pre-check of the on-disk store + default soul memory (read-only). + /// A local ops command (NOT an HTTP client): exits non-zero when unhealthy, + /// so the upgrade flow can gate on it. See PHASE-07. + Doctor, + /// Offline store-level ops (act directly on the DB, no running service). + #[command(subcommand)] + Inbox(InboxCommand), /// GET /api/v1/health Health, /// Strand resources under /api/v1/strands @@ -109,6 +116,18 @@ enum CompactCommand { }, } +#[derive(Subcommand)] +enum InboxCommand { + /// Enqueue one `santi_system` record into a strand's durable inbox WITHOUT a + /// running service (a direct MQ producer). The strand comes from + /// --strand/SANTI_STRAND_ID and must already exist. Used by the self-upgrade + /// flow to seed the "come look" record before starting the final version. + Seed { + /// The message text (the "come look" occurrence). + text: String, + }, +} + #[derive(Subcommand)] enum StrandCommand { /// POST /api/v1/strands @@ -145,6 +164,8 @@ async fn main() -> Result<()> { let cli = Cli::parse(); match cli.command { Command::Service { args } => run_service(args).await, + Command::Doctor => run_doctor(), + Command::Inbox(inbox) => run_inbox(inbox, cli.strand), other => { let defaults = ClientDefaults { strand: cli.strand, @@ -185,6 +206,38 @@ impl ClientDefaults { } } +/// Offline pre-check (local ops, no HTTP). Prints the report as JSON to stdout +/// and exits non-zero when unhealthy, so a caller (the upgrade flow) can gate. +fn run_doctor() -> Result<()> { + let report = santi_api::ops::doctor().map_err(|error| anyhow::anyhow!(error))?; + println!("{}", serde_json::to_string_pretty(&report)?); + if !report.ok { + anyhow::bail!("doctor: unhealthy (see report above)"); + } + Ok(()) +} + +/// Offline inbox producer (local ops, no HTTP). Resolves the strand from +/// --strand/SANTI_STRAND_ID and seeds a durable record; exits non-zero if the +/// inbox gate rejects it, so the upgrade flow notices a badly-behind strand. +fn run_inbox(command: InboxCommand, default_strand: Option) -> Result<()> { + match command { + InboxCommand::Seed { text } => { + let strand_id = default_strand + .map(|id| id.trim().to_string()) + .filter(|id| !id.is_empty()) + .ok_or_else(|| anyhow::anyhow!("no strand id: set --strand / SANTI_STRAND_ID"))?; + let report = santi_api::ops::inbox_seed(&strand_id, &text) + .map_err(|error| anyhow::anyhow!(error))?; + println!("{}", serde_json::to_string_pretty(&report)?); + if !report.accepted { + anyhow::bail!("inbox seed rejected: {}", report.reason.unwrap_or_default()); + } + Ok(()) + } + } +} + /// Run the runtime server in-process via `santi-api`. async fn run_service(args: Vec) -> Result<()> { let argv = std::iter::once("santi".to_string()).chain(args); @@ -214,6 +267,8 @@ async fn run_client( let base = base_url.trim_end_matches('/').to_string(); match command { Command::Service { .. } => unreachable!("service is handled before the client path"), + Command::Doctor => unreachable!("doctor is handled before the client path"), + Command::Inbox(_) => unreachable!("inbox is handled before the client path"), Command::Health => get(&client, &format!("{base}/api/v1/health")).await, Command::Strand(StrandCommand::Create) => { post(&client, &format!("{base}/api/v1/strands"), None).await