ive harnessed the harness
Something went wrong. Try again.
136 kB · 3453 lines
Rust
at main
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454//! One authoritative SQLite connection, used only inside bounded blocking jobs.//! Network calls, Python execution and hooks never run inside a transaction.use crate::{ behavior::Release, hooks::{DeliveryDecision, ReviewedDelivery}, protocol::*,};use anyhow::{anyhow, bail, ensure, Context, Result};use rusqlite::{params, Connection, OptionalExtension, Transaction, TransactionBehavior};use serde::{Deserialize, Serialize};use serde_json::{json, Value};use std::{ fs::{self, File}, path::{Path, PathBuf}, sync::{Arc, Mutex}, time::{Duration, SystemTime, UNIX_EPOCH},};use tokio::sync::{watch, OwnedSemaphorePermit, Semaphore};
use crate::work::{ ClaimResult, ConsequenceMetadata, Impact, ObserveResult, ObserverCaseRecord, OccurrenceId, OriginKind, Provenance, Readiness, RetireResult, RunCause, WorkItem, WorkItemId, WorkStatus,};
/// Durable revisit gate attached to a deferred work item: the actual/// dependency that keeps it blocked and/or a bounded revisit deadline.#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize)]pub struct RevisitGate { /// The actual dependency, as a durable reason string. When set the /// item reads as blocked until it is released. pub dependency: Option<String>, /// Bounded revisit deadline in wall-clock ms. Eligibility resumes /// only at or after this time. pub deadline_ms: Option<i64>,}
impl RevisitGate { #[must_use] pub fn blocked_on(dependency: impl Into<String>) -> Self { Self { dependency: Some(dependency.into()), deadline_ms: None, } }
#[must_use] pub fn at(deadline_ms: i64) -> Self { Self { dependency: None, deadline_ms: Some(deadline_ms), } }
fn readiness_str_and_reason(&self) -> (&'static str, Option<String>) { match &self.dependency { Some(d) => ("blocked", Some(d.clone())), None => ("ready", None), } }}
/// Durable revisit state persisted on a work item.#[derive(Clone, Debug, PartialEq, Eq, Serialize)]pub struct RevisitState { pub deadline_ms: Option<i64>, pub reason: Option<String>,}
const APPLICATION_ID: i64 = 1263288882;pub const CURRENT_SCHEMA_VERSION: i64 = 15;pub const MIGRATION: &str = include_str!("../migrations/001_runtime.sql");pub const MIGRATION_002: &str = include_str!("../migrations/002_lifecycle.sql");pub const MIGRATION_003: &str = include_str!("../migrations/003_work.sql");pub const MIGRATION_004: &str = include_str!("../migrations/004_observations.sql");pub const MIGRATION_005: &str = include_str!("../migrations/005_originations.sql");pub const MIGRATION_006: &str = include_str!("../migrations/006_assessment_integrity.sql");pub const MIGRATION_007: &str = include_str!("../migrations/007_placements.sql");pub const MIGRATION_008: &str = include_str!("../migrations/008_placement_events.sql");pub const MIGRATION_009: &str = include_str!("../migrations/009_unknown_executions.sql");pub const MIGRATION_010: &str = include_str!("../migrations/010_peers.sql");pub const MIGRATION_011: &str = include_str!("../migrations/011_orientations.sql");pub const MIGRATION_012: &str = include_str!("../migrations/012_budget.sql");pub const MIGRATION_013: &str = include_str!("../migrations/013_frontier.sql");pub const MIGRATION_014: &str = include_str!("../migrations/014_event_kinds.sql");pub const MIGRATION_015: &str = include_str!("../migrations/015_prompt_header.sql");pub const MIGRATION_016: &str = include_str!("../migrations/016_retire_inbox.sql");
const MIGRATIONS: &[(i64, &str)] = &[ (1, MIGRATION), (2, MIGRATION_002), (3, MIGRATION_003), (4, MIGRATION_004), (5, MIGRATION_005), (6, MIGRATION_006), (7, MIGRATION_007), (8, MIGRATION_008), (9, MIGRATION_009), (10, MIGRATION_010), (11, MIGRATION_011), (12, MIGRATION_012), (13, MIGRATION_013), (14, MIGRATION_014), (15, MIGRATION_015), (16, MIGRATION_016),];
mod orientation;
/// Bounded concurrent blocking jobs sharing the one connection. Each job takes a/// permit AND a share of the ownership lock for its whole transaction, so the/// permit count is also the quiescence barrier a restarting owner waits on.const BLOCKING_JOBS: usize = 8;
/// One blocking job's share of the ownership lock, held together with its permit.////// The pairing is the invariant: fields drop in declaration order, so the lock/// share is dropped before the permit frees and `Store::quiesce` may read a free/// permit as "no job still holds the database". Declaration order is used rather/// than a pair of locals because it holds while unwinding from a panic too; two/// locals would drop in reverse and free the permit first.struct Job { _share: Arc<File>, _permit: OwnedSemaphorePermit,}
impl Job { fn holding(share: Arc<File>, permit: OwnedSemaphorePermit) -> Self { Self { _share: share, _permit: permit, } }}
type EffectRow = (String, String, String, Option<i64>, Option<String>);
#[derive(Clone)]pub struct Store { connection: Arc<Mutex<Connection>>, gate: Arc<Semaphore>, changed: watch::Sender<()>, _owner: Arc<File>,}
/// A live `Store` has no useful printable state: the connection is not `Debug` and/// the path is deliberately not retained. This exists only so that a struct which/// holds a `Store` (for example a restored deployment) can carry `Debug` without/// either dropping that field from its own output or leaking the database's/// location. It prints no path and no credentials.impl std::fmt::Debug for Store { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.write_str("Store { .. }") }}#[derive(Clone, Debug, Serialize, Deserialize)]pub struct Event { pub id: i64, pub kind: String, pub payload: Value, pub created_ms: i64,}#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]pub struct Receipt { pub effect_id: EffectId, pub status: ReceiptStatus, pub event_id: Option<i64>, pub reason: Option<String>,}
#[derive(Clone, Debug, Serialize, Deserialize)]pub struct ExecutionEffect { pub effect_id: EffectId, pub text: String, pub status: ReceiptStatus, pub event_id: Option<i64>, pub reason: Option<String>,}#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]#[serde(rename_all = "snake_case")]pub enum ReceiptStatus { Confirmed, Held,}pub enum PreparedEffect { New(EffectId), Existing(Receipt),}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]pub enum StoreAdmit { Inserted(i64), Duplicate(i64), Conflict(String), Error(String),}
impl StoreAdmit { pub fn seq(&self) -> Option<i64> { match self { Self::Inserted(seq) | Self::Duplicate(seq) => Some(*seq), _ => None, } }
pub fn is_inserted(&self) -> bool { matches!(self, Self::Inserted(_)) }
pub fn is_duplicate(&self) -> bool { matches!(self, Self::Duplicate(_)) }
pub fn is_conflict(&self) -> bool { matches!(self, Self::Conflict(_)) }
pub fn is_error(&self) -> bool { matches!(self, Self::Error(_)) }}#[derive(Clone, Debug)]pub struct Activation { pub epoch: i64, pub revision: String, pub path: PathBuf,}
pub fn now_ms() -> Result<i64> { Ok(SystemTime::now() .duration_since(UNIX_EPOCH)? .as_millis() .try_into()?)}
fn event(tx: &Transaction<'_>, session: &str, kind: &str, payload: &Value) -> Result<i64> { tx.execute( "INSERT INTO events(session_id,kind,payload,created_ms) VALUES(?1,?2,?3,?4)", params![session, kind, serde_json::to_string(payload)?, now_ms()?], )?; Ok(tx.last_insert_rowid())}
fn live(tx: &Transaction<'_>, session: &str, execution: &str, generation: &str) -> Result<()> { let count: i64 = tx.query_row("SELECT count(*) FROM executions WHERE id=?1 AND session_id=?2 AND generation=?3 AND state='running' AND completed_event IS NULL", params![execution, session, generation], |row| row.get(0))?; ensure!( count == 1, "execution authority has ended or belongs to another generation" ); Ok(())}
fn apply_migrations(connection: &mut Connection, starting_version: i64) -> Result<()> { for &(version, migration) in MIGRATIONS { if version <= starting_version || version > 10 { continue; } if version == 2 { connection.pragma_update(None, "foreign_keys", "OFF")?; } connection.execute_batch(migration)?; if version == 2 { connection.pragma_update(None, "foreign_keys", "ON")?; } }
for &(version, migration) in MIGRATIONS { if version <= 10 { continue; } let migrated_version: i64 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; if migrated_version == version - 1 { connection.execute_batch(migration)?; } } Ok(())}
fn retire_preflight( connection: &Connection, id: &str, attempt_id: &str,) -> Result<Option<RetireResult>> { let current: Option<(String, Option<String>)> = connection .query_row( "SELECT status, attempt_id FROM work_items WHERE id = ?1", params![id], |row| Ok((row.get(0)?, row.get(1)?)), ) .optional()?;
let Some((status, holder)) = current else { return Ok(Some(RetireResult::NotFound)); };
match status.as_str() { "completed" | "failed" => Ok(Some(RetireResult::AlreadyTerminal)), "claimed" if holder.as_deref() != Some(attempt_id) => { let expected = holder .map(|value| value.parse().unwrap_or_else(|_| OperationId::fresh())) .unwrap_or_else(OperationId::fresh); Ok(Some(RetireResult::WrongOwner { expected })) } "claimed" => Ok(None), _ => Ok(Some(RetireResult::NotClaimed)), }}
fn merge_evidence_refs(existing_json: &str, new_refs: &[String]) -> (Vec<String>, bool) { let mut refs: Vec<String> = serde_json::from_str(existing_json).unwrap_or_default(); let mut changed = false; for reference in new_refs { if !refs.contains(reference) { refs.push(reference.clone()); changed = true; } } (refs, changed)}
impl Store { /// Open only the new runtime format. This refuses an unrelated populated DB. /// Call once per daemon; clones share its connection and admission bound. pub fn open(path: impl AsRef<Path>) -> Result<Self> { let path = path.as_ref(); ensure!( path != Path::new(":memory:"), "runtime store requires a durable file" ); let parent = path .parent() .filter(|p| !p.as_os_str().is_empty()) .unwrap_or(Path::new(".")); fs::create_dir_all(parent)?; let path = parent .canonicalize()? .join(path.file_name().context("database filename required")?); if let Ok(metadata) = fs::symlink_metadata(&path) { ensure!( !metadata.file_type().is_symlink(), "database symlink aliases are not supported" ); } let mut lock_path = path.as_os_str().to_os_string(); lock_path.push(".owner.lock"); let owner = File::options() .read(true) .write(true) .create(true) .truncate(false) .open(PathBuf::from(lock_path))?; owner.try_lock().map_err(|error| { anyhow!( "database {} is owned elsewhere or locking failed: {error:?}; \ a live daemon, or an in-flight write that has not quiesced, still holds {}.owner.lock", path.display(), path.display() ) })?; // Refuse an unsupported SQLite runtime before modifying even a new DB. let v = rusqlite::version_number(); ensure!( v >= 3_051_003 || (3_050_007..3_051_000).contains(&v) || (3_044_006..3_045_000).contains(&v), "SQLite {} lacks a verified WAL-reset fix", rusqlite::version() ); let mut connection = Connection::open(path)?; connection.busy_timeout(Duration::from_secs(5))?; connection.pragma_update(None, "foreign_keys", "ON")?; let app: i64 = connection.pragma_query_value(None, "application_id", |row| row.get(0))?; let version: i64 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?; let starting_version = if app == 0 && version == 0 { let tables: i64 = connection.query_row( "SELECT count(*) FROM sqlite_master WHERE name NOT LIKE 'sqlite_%'", [], |r| r.get(0), )?; ensure!( tables == 0, "refusing to initialize over a legacy or unrelated database" ); 0 } else { ensure!( app == APPLICATION_ID && (1..=CURRENT_SCHEMA_VERSION).contains(&version), "unrecognized runtime database version" ); version }; apply_migrations(&mut connection, starting_version)?;
// Runner-level integrity gate for the whole migration chain: the // chain runs with enforcement toggled per rebuild (migration 008 // drops referenced tables, so enforcement is OFF inside those // transactions); no migration may commit an orphaned reference. let orphans: i64 = connection.query_row("SELECT COUNT(*) FROM pragma_foreign_key_check", [], |r| { r.get(0) })?; ensure!( orphans == 0, "runtime database has {orphans} orphaned foreign-key references after migration" ); connection.pragma_update(None, "journal_mode", "WAL")?; connection.pragma_update(None, "synchronous", "FULL")?; let tx = connection.transaction_with_behavior(TransactionBehavior::Immediate)?; let unclosed: Vec<(String, String, String)> = { let mut query = tx.prepare( "SELECT id,session_id,state FROM executions WHERE completed_event IS NULL", )?; let rows = query .query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))? .collect::<rusqlite::Result<_>>()?; rows }; for (id, session, prior) in unclosed { let (kind, state, reason, effect_reason) = if prior == "yielded" { ( "execution.completed", "yielded", "host restarted after a yielded result; publication was not observed", "host restarted before local publication", ) } else { ( "execution.unknown", "unknown", "host restarted before completion; code may have run and its outcome was not observed", "execution outcome unknown; publication withheld", ) }; let seq = event( &tx, &session, kind, &json!({"execution_id":id,"state":state,"reason":reason}), )?; tx.execute( "UPDATE executions SET state=?2,completed_event=?3 WHERE id=?1", params![id, state, seq], )?; tx.execute( "UPDATE effects SET status='held',reason=?2 WHERE execution_id=?1 AND status='proposed'", params![id, effect_reason], )?; } let now_startup = now_ms().unwrap_or_else(|_| now_ms().unwrap_or(0)); tx.execute( "UPDATE work_items SET status = 'failed', error = 'host restarted; attempt orphaned', failed_ms = ?1 WHERE status = 'claimed'", params![now_startup], )?; tx.commit()?; Ok(Self { connection: Arc::new(Mutex::new(connection)), gate: Arc::new(Semaphore::new(BLOCKING_JOBS)), changed: watch::channel(()).0, _owner: Arc::new(owner), }) }
/// A coalesced wake hint, never a second history source. Readers resume by durable sequence. pub fn changes(&self) -> watch::Receiver<()> { self.changed.subscribe() }
async fn with<T: Send + 'static>( &self, run: impl FnOnce(&mut Connection) -> Result<T> + Send + 'static, ) -> Result<T> { let permit = self.gate.clone().acquire_owned().await?; let connection = self.connection.clone(); let owner = self._owner.clone(); let changed = self.changed.clone(); tokio::task::spawn_blocking(move || { // Cancellation drops interest, not an in-flight database write: the // share is held until the transaction finishes, and released before // the permit, however this closure ends. let _job = Job::holding(owner, permit); let mut connection = connection .lock() .map_err(|_| anyhow!("runtime database lock poisoned"))?; let before = connection.total_changes(); let result = run(&mut connection); let wrote = connection.total_changes() != before; drop(connection); // Run has committed or rolled back and released its lock. A rollback may wake // a reader, but only committed rows can ever become client records. if wrote { changed.send_modify(|_| {}); } result }) .await .context("runtime database task panicked")? }
/// Take ownership of a database whose previous owner is a *stopped* process or /// a dropped handle, tolerating the moment it needs to release the file. A /// live owner is still refused: it never releases, and the wait expires into /// the same error `open` reports immediately. /// /// Distinct from `open` on purpose. `open` answers "is this mine, now?" and is /// what a second live daemon must be refused by; this answers "take this over /// once the previous owner has let go". pub async fn open_on_restart(path: impl AsRef<Path>, wait: Duration) -> Result<Self> { let path = path.as_ref(); let deadline = std::time::Instant::now() + wait; loop { match Self::open(path) { Ok(store) => return Ok(store), Err(error) => { if std::time::Instant::now() >= deadline { return Err(error.context(format!( "no previous owner released {} within {wait:?}", path.display() ))); } tokio::time::sleep(Duration::from_millis(5)).await; } } } }
/// Wait until no blocking job holds the connection or its share of the /// ownership lock. An aborted task does NOT stop an in-flight write, so /// teardown that returns before this can hand a restarting owner a database /// still locked by a previous job. Bounded, because a wedged job must be /// reported rather than hang the service forever. /// /// This covers jobs, which is what a barrier over their permits can decide; /// it does not release ownership itself. The database is handed over when the /// last `Store` handle drops - for a daemon, when it exits. pub async fn quiesce(&self) -> Result<()> { const DEADLINE: Duration = Duration::from_secs(30); let barrier = self.gate.clone().acquire_many_owned(BLOCKING_JOBS as u32); match tokio::time::timeout(DEADLINE, barrier).await { Ok(Ok(all)) => { drop(all); Ok(()) } Ok(Err(error)) => Err(error.into()), Err(_) => bail!("runtime store did not quiesce within {DEADLINE:?}"), } }
pub async fn ensure_session(&self, session: &SessionId) -> Result<()> { let session = session.to_string(); self.with(move |db| { db.execute( "INSERT INTO sessions(id) VALUES(?1) ON CONFLICT DO NOTHING", [session], )?; Ok(()) }) .await }
/// Returns only after the input and its work item commit together. The channel/scheduler is not the /// source of truth. Equal content with different source keys stays distinct. pub async fn admit_with_status( &self, session: &SessionId, source_key: String, payload: Value, ) -> Result<StoreAdmit> { ensure!( !source_key.is_empty() && source_key.len() <= 512, "invalid source key" ); ensure!( serde_json::to_vec(&payload)?.len() <= MAX_CODE, "input payload too large" ); let session = session.to_string(); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; let current_state: Option<String> = tx .query_row("SELECT state FROM sessions WHERE id=?1", [&session], |r| r.get(0)) .optional()?; if let Some(st) = ¤t_state { if st == "stopped" || st == "draining" { bail!("cannot admit input: session {session} is in state {st}"); } } let previous: Option<(i64,String)> = tx.query_row("SELECT seq,payload FROM events WHERE session_id=?1 AND source_key=?2", params![session,source_key], |r| Ok((r.get(0)?,r.get(1)?))).optional()?; if let Some((id, prior)) = previous { let prior_val: Value = serde_json::from_str(&prior)?; if prior_val != payload { return Ok(StoreAdmit::Conflict("source key reused with different content".into())); } return Ok(StoreAdmit::Duplicate(id)); } let now = now_ms()?; tx.execute("INSERT INTO events(session_id,kind,source_key,payload,created_ms) VALUES(?1,'input',?2,?3,?4)", params![session,source_key,serde_json::to_string(&payload)?,now])?; let id = tx.last_insert_rowid(); let work_id = format!("work-input-{}", id); let cause_payload = serde_json::to_string(&json!({"event_ids": [id]}))?; let evidence_refs = serde_json::to_string(&json!([format!("event:{}", id)]))?; tx.execute( "INSERT INTO work_items ( id, session_id, purpose, cause_kind, cause_payload, readiness, evidence_refs, impact, aging_ms, reporter_id, origin_kind, status, created_ms ) VALUES (?1, ?2, 'user_input', 'input_frame', ?3, 'ready', ?4, 2, 0, 'ingress', 'human', 'pending', ?5)", params![work_id, session, cause_payload, evidence_refs, now], )?; let changed = tx.execute("UPDATE sessions SET state='ready',wake_at_ms=NULL WHERE id=?1 AND state='waiting'", [&session])?; if changed == 1 { event(&tx, &session, "session.woke", &json!({"reason":"input","event_id":id}))?; } tx.commit()?; Ok(StoreAdmit::Inserted(id)) }).await }
/// Returns only after the input and its work item commit together. The channel/scheduler is not the /// source of truth. Equal content with different source keys stays distinct. pub async fn admit( &self, session: &SessionId, source_key: String, payload: Value, ) -> Result<i64> { match self.admit_with_status(session, source_key, payload).await? { StoreAdmit::Inserted(id) | StoreAdmit::Duplicate(id) => Ok(id), StoreAdmit::Conflict(err) => bail!("conflict: {err}"), StoreAdmit::Error(err) => bail!("error: {err}"), } }
/// Persist this session's workbench placement, keyed by the namespace /// generation the binding was issued under. A placement is bound once /// per generation: re-writing the SAME descriptor is idempotent /// (restart), but re-writing a DIFFERENT descriptor under a generation /// that already holds one is refused — a relocate must mint a fresh /// generation so a target change is never an in-place mutation that /// implies the old namespace survived. pub async fn set_session_placement( &self, session: &SessionId, placement: &crate::target::SessionPlacement, ) -> Result<()> { let session_str = session.to_string(); let generation = placement.namespace_generation.as_str().to_owned(); let placement_json = serde_json::to_string(&placement)?; self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let existing: Option<String> = tx .query_row( "SELECT placement FROM session_placements WHERE session_id = ?1 AND namespace_generation = ?2", params![session_str, generation], |r| r.get(0), ) .optional()?; if let Some(prior) = existing { if prior != placement_json { bail!( "namespace generation {generation} is already bound to a different placement; mint a new namespace generation for the relocation" ); } tx.commit()?; return Ok(()); } tx.execute( "INSERT INTO session_placements (session_id, namespace_generation, placement, created_ms) VALUES (?1, ?2, ?3, ?4)", params![session_str, generation, placement_json, now_ms()?], )?; tx.commit()?; Ok(()) }) .await }
/// The session's current placement: the most recently bound namespace /// generation's descriptor, type-checked through the target vocabulary /// at read time. A hand-edited or corrupt row is a typed error, never a /// silent fall back to Local. None only before the first binding. pub async fn session_placement( &self, session: &SessionId, ) -> Result<Option<crate::target::SessionPlacement>> { let session_str = session.to_string(); self.with(move |conn| { let json: Option<String> = conn .query_row( "SELECT placement FROM session_placements WHERE session_id = ?1 ORDER BY created_ms DESC, namespace_generation DESC LIMIT 1", params![session_str], |r| r.get(0), ) .optional()?; match json { Some(text) => { let value: serde_json::Value = serde_json::from_str(&text)?; Ok(Some(crate::target::SessionPlacement::from_json(&value)?)) } None => Ok(None), } }) .await }
pub async fn session_state( &self, session: &SessionId, ) -> Result<Option<(String, Option<i64>)>> { let session = session.to_string(); self.with(move |db| { let res = db .query_row( "SELECT state, wake_at_ms FROM sessions WHERE id = ?1", [&session], |r| Ok((r.get(0)?, r.get(1)?)), ) .optional()?; Ok(res) }) .await }
pub async fn set_session_state(&self, session: &SessionId, state: &str) -> Result<()> { ensure!( matches!(state, "ready" | "waiting" | "draining" | "stopped"), "invalid session state" ); let session = session.to_string(); let state = state.to_string(); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; if state == "waiting" { tx.execute( "UPDATE sessions SET state=?1 WHERE id=?2", params![state, session], )?; } else { tx.execute( "UPDATE sessions SET state=?1, wake_at_ms=NULL WHERE id=?2", params![state, session], )?; } tx.commit()?; Ok(()) }) .await }
pub async fn get_event(&self, session: &SessionId, seq: i64) -> Result<Option<Event>> { let session = session.to_string(); self.with(move |db| { let res = db .query_row( "SELECT seq, kind, payload, created_ms FROM events WHERE session_id = ?1 AND seq = ?2", params![session, seq], |r| { Ok(Event { id: r.get(0)?, kind: r.get(1)?, payload: serde_json::from_str(&r.get::<_, String>(2)?).map_err(|e| { rusqlite::Error::FromSqlConversionFailure( 2, rusqlite::types::Type::Text, Box::new(e), ) })?, created_ms: r.get(3)?, }) }, ) .optional()?; Ok(res) }) .await }
pub async fn find_event_by_source_key( &self, session: &SessionId, source_key: &str, ) -> Result<Option<i64>> { let session = session.to_string(); let source_key = source_key.to_string(); self.with(move |db| { let res: Option<i64> = db .query_row( "SELECT seq FROM events WHERE session_id = ?1 AND source_key = ?2", params![session, source_key], |r| r.get(0), ) .optional()?; Ok(res) }) .await }
/// The session's unclaimed input events, oldest first. /// /// The queue is `work_items`, not a second table: a claim on an input is durable there, so an /// input claimed by an attempt that died is recovered rather than silently becoming /// unclaimed again in this process's memory. pub async fn pending(&self, session: &SessionId, limit: usize) -> Result<Vec<i64>> { ensure!((1..=1024).contains(&limit), "invalid input query limit"); let session = session.to_string(); self.with(move |db| { let mut query = db.prepare( "SELECT cause_payload FROM work_items WHERE session_id=?1 AND cause_kind='input_frame' AND status='pending' ORDER BY created_ms, CAST(substr(id, 12) AS INTEGER) LIMIT ?2", )?; let payloads = query .query_map(params![session, limit as i64], |r| r.get::<_, String>(0))? .collect::<rusqlite::Result<Vec<_>>>()?; let mut ids = Vec::new(); for payload in payloads { let value: Value = serde_json::from_str(&payload)?; if let Some(id) = value .get("event_ids") .and_then(|ids| ids.as_array()) .and_then(|ids| ids.first()) .and_then(Value::as_i64) { ids.push(id); } } Ok(ids) }) .await }
/// Claim up to `limit` of the session's pending inputs for one attempt. pub async fn claim_pending( &self, session: &SessionId, attempt: &OperationId, limit: usize, ) -> Result<Vec<i64>> { let mut claimed = Vec::new(); for id in self.pending(session, limit).await? { match self .claim_work_item(&format!("work-input-{id}"), attempt, now_ms()?) .await { Ok(result) if result.held_by(attempt) => claimed.push(id), Ok(_) => {} Err(error) => { // Nothing is half-claimed: inputs this attempt did take go back to pending. for done in &claimed { let _ = self .yield_work_item(&format!("work-input-{done}"), attempt) .await; } return Err(error); } } } Ok(claimed) }
/// Give claimed inputs back, so the next attempt can take them. pub async fn release_inputs(&self, attempt: &OperationId, ids: &[i64]) -> Result<()> { for id in ids { if let Err(error) = self .yield_work_item(&format!("work-input-{id}"), attempt) .await { eprintln!("input claim {id} not released: {error}"); } } Ok(()) }
/// Host-only acknowledgement; not exposed to the kernel. A future turn /// coordinator should call this only after persisting its input consumption. pub async fn acknowledge_inputs(&self, session: &SessionId, ids: Vec<i64>) -> Result<()> { ensure!(ids.len() <= 1024, "too many input acknowledgements"); let session_str = session.to_string(); let ids_clone = ids.clone(); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; let now = now_ms().unwrap_or_else(|_| now_ms().unwrap_or(0)); for id in ids_clone { let exists: i64 = tx.query_row( "SELECT count(*) FROM events WHERE seq=?1 AND session_id=?2", params![id, session_str], |r| r.get(0), )?; ensure!( exists == 1, "acknowledgement is stale or belongs to another session" ); let work_id = format!("work-input-{}", id); tx.execute( "UPDATE work_items SET status='completed', disposition='consumed_by_turn', completed_ms=?1, attempt_id=NULL, claimed_ms=NULL WHERE id=?2 AND status IN ('pending','claimed','deferred')", params![now, work_id], )?; } tx.commit()?; Ok(()) }).await }
pub async fn scratchpad(&self, session: &SessionId, text: String) -> Result<i64> { ensure!(text.len() <= MAX_CODE, "scratchpad too large"); self.append(session, "scratchpad", json!({"text":text})) .await }
/// The session's active context view: the newest installed frontier, or none when the /// session has never had one installed. pub async fn active_frontier( &self, session: &SessionId, ) -> Result<Option<crate::frontier::Frontier>> { let session = session.to_string(); let raw: Option<String> = self .with(move |db| { let mut query = db.prepare( "SELECT payload FROM events WHERE session_id=?1 AND kind=?2 ORDER BY seq DESC LIMIT 1", )?; Ok(query .query_row(params![session, crate::frontier::FRONTIER_KIND], |row| { row.get::<_, String>(0) }) .optional()?) }) .await?; raw.map(|text| { let payload: Value = serde_json::from_str(&text)?; crate::frontier::Frontier::from_payload(&payload) }) .transpose() }
/// Installs a frontier IF AND ONLY IF it is the next view. /// /// Reading the active epoch and inserting the new record are one transaction on one connection, /// so two writers cannot both believe they are next: the loser is refused, naming the epoch it /// should have planned from, and the winner's view stands. A blind append would let a second /// writer replace a view another task had just committed, or install one derived from a view /// that no longer exists. /// /// The append is still the commit point: the previous view is active until this commits and the /// new one is active after it, so there is no window in which the session has neither. pub async fn install_frontier( &self, session: &SessionId, frontier: &crate::frontier::Frontier, ) -> Result<i64> { let session = session.to_string(); let payload = frontier.to_payload(); let expected_epoch = frontier.epoch; self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; let current: Option<String> = tx .query_row( "SELECT payload FROM events WHERE session_id=?1 AND kind=?2 ORDER BY seq DESC LIMIT 1", params![session, crate::frontier::FRONTIER_KIND], |row| row.get(0), ) .optional()?; let current_epoch = match current { Some(raw) => { crate::frontier::Frontier::from_payload(&serde_json::from_str::<Value>(&raw)?)? .epoch } None => 0, }; ensure!( current_epoch + 1 == expected_epoch, "refusing to install frontier epoch {expected_epoch}: the active view is epoch {current_epoch}" ); let id = event(&tx, &session, crate::frontier::FRONTIER_KIND, &payload)?; tx.commit()?; Ok(id) }) .await }
pub async fn append( &self, session: &SessionId, kind: &'static str, payload: Value, ) -> Result<i64> { let session = session.to_string(); self.with(move |db| { let tx = db.transaction()?; let id = event(&tx, &session, kind, &payload)?; tx.commit()?; Ok(id) }) .await }
pub async fn history( &self, session: &SessionId, after: i64, limit: usize, ) -> Result<Vec<Event>> { ensure!( after >= 0 && (1..=1024).contains(&limit), "invalid history page" ); let session = session.to_string(); self.with(move |db| { let mut query = db.prepare("SELECT seq,kind,payload,created_ms FROM events WHERE session_id=?1 AND seq>?2 ORDER BY seq LIMIT ?3")?; let raw: Vec<(i64,String,String,i64)> = query.query_map(params![session,after,limit as i64], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?)))?.collect::<rusqlite::Result<_>>()?; raw.into_iter().map(|(id,kind,payload,created_ms)| Ok(Event {id,kind,payload:serde_json::from_str(&payload)?,created_ms})).collect() }).await }
/// The oldest and newest sequence a session holds, or `None` when it holds nothing. /// /// This is what lets a daemon answer "how much of this session can I serve" without reading the /// session: a window and a floor are both decided from these two numbers. pub async fn seq_bounds(&self, session: &SessionId) -> Result<Option<(i64, i64)>> { let session = session.to_string(); self.with(move |db| { let bounds: (Option<i64>, Option<i64>) = db.query_row( "SELECT MIN(seq),MAX(seq) FROM events WHERE session_id=?1", [session], |r| Ok((r.get(0)?, r.get(1)?)), )?; Ok(match bounds { (Some(low), Some(high)) => Some((low, high)), _ => None, }) }) .await }
/// The newest `limit` records strictly before `before`, oldest first: one page of the past. /// /// The mirror of `history`, and what a reader asks for when they press "load earlier history": the /// same records read from the other end, returned in the order they happened. pub async fn history_before( &self, session: &SessionId, before: i64, limit: usize, ) -> Result<Vec<Event>> { ensure!( before >= 0 && (1..=1024).contains(&limit), "invalid history page" ); let session = session.to_string(); self.with(move |db| { let mut query = db.prepare("SELECT seq,kind,payload,created_ms FROM events WHERE session_id=?1 AND seq<?2 ORDER BY seq DESC LIMIT ?3")?; let raw: Vec<(i64,String,String,i64)> = query.query_map(params![session,before,limit as i64], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?)))?.collect::<rusqlite::Result<_>>()?; let mut events: Vec<Event> = raw.into_iter().map(|(id,kind,payload,created_ms)| Ok(Event {id,kind,payload:serde_json::from_str(&payload)?,created_ms})).collect::<Result<_>>()?; events.reverse(); Ok(events) }).await }
/// Records after `after` and at or before `through`, oldest first: one page of a bounded replay. /// /// A replay states the range it is sending, and this is how that range is kept true while the session /// is still being written to: without the upper bound a page would carry events that arrived after the /// range was announced. pub async fn history_through( &self, session: &SessionId, after: i64, through: i64, limit: usize, ) -> Result<Vec<Event>> { ensure!( after >= 0 && (1..=1024).contains(&limit), "invalid history page" ); let session = session.to_string(); self.with(move |db| { let mut query = db.prepare("SELECT seq,kind,payload,created_ms FROM events WHERE session_id=?1 AND seq>?2 AND seq<=?3 ORDER BY seq LIMIT ?4")?; let raw: Vec<(i64,String,String,i64)> = query.query_map(params![session,after,through,limit as i64], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?)))?.collect::<rusqlite::Result<_>>()?; raw.into_iter().map(|(id,kind,payload,created_ms)| Ok(Event {id,kind,payload:serde_json::from_str(&payload)?,created_ms})).collect() }).await }
pub async fn activation(&self, session: &SessionId) -> Result<Option<Activation>> { let session = session.to_string(); self.with(move |db| { Ok(db.query_row("SELECT s.hook_epoch,r.id,r.path FROM sessions s JOIN releases r ON r.id=s.hook_revision WHERE s.id=?1", [session], |r| Ok(Activation {epoch:r.get(0)?,revision:r.get(1)?,path:PathBuf::from(r.get::<_,String>(2)?)})).optional()?) }).await }
pub async fn activate( &self, session: &SessionId, expected_epoch: i64, release: Release, ) -> Result<i64> { self.activate_with_work(session, expected_epoch, release, None) .await }
pub async fn activate_with_work( &self, session: &SessionId, expected_epoch: i64, release: Release, work_item_id: Option<String>, ) -> Result<i64> { let session = session.to_string(); let path = release .path .to_str() .context("release path is not UTF-8")? .to_owned(); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; tx.execute("INSERT INTO releases(id,path) VALUES(?1,?2) ON CONFLICT DO NOTHING", params![release.id,path])?; let changed = tx.execute("UPDATE sessions SET hook_revision=?1,hook_epoch=hook_epoch+1 WHERE id=?2 AND hook_epoch=?3", params![release.id,session,expected_epoch])?; ensure!(changed == 1, "activation conflict: expected epoch is stale"); let mut payload = json!({"revision": release.id, "epoch": expected_epoch + 1}); if let Some(ref wid) = work_item_id { payload["work_item_id"] = json!(wid); } event(&tx, &session, "hooks.activated", &payload)?; tx.commit()?; Ok(expected_epoch + 1) }).await }
pub async fn begin( &self, session: &SessionId, id: &OperationId, generation: &Generation, revision: String, code: String, ) -> Result<()> { ensure!(code.len() <= MAX_CODE, "cell source exceeds limit"); let (session, id, generation) = (session.to_string(), id.to_string(), generation.to_string()); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; let seq = event(&tx,&session,"execution.started",&json!({"execution_id":id,"generation":generation,"revision":revision,"code":code}))?; tx.execute("INSERT INTO executions(id,session_id,generation,hook_revision,state,started_event) VALUES(?1,?2,?3,?4,'running',?5)", params![id,session,generation,revision,seq])?; tx.commit()?; Ok(()) }).await }
pub async fn finish(&self, id: &OperationId, outcome: Outcome) -> Result<()> { if let Outcome::Unknown { reason } = outcome { return self.finish_unknown(id, reason).await; } let id = id.to_string(); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; let (session,prior,completed): (String,String,Option<i64>) = tx.query_row("SELECT session_id,state,completed_event FROM executions WHERE id=?1", [&id], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?)))?; ensure!(completed.is_none(), "execution already completed"); let state = if prior == "yielded" { "yielded" } else { match &outcome { Outcome::Cell {status:CellStatus::Ok,..} => "ok", Outcome::Cell {status:CellStatus::Yielded,..} => bail!("worker claimed an uncommitted yield"), Outcome::Cell {status:CellStatus::Error,..} => "error", _ => "interrupted", } }; let seq = event(&tx,&session,"execution.completed",&json!({"execution_id":id,"state":state,"outcome":outcome}))?; tx.execute("UPDATE executions SET state=?2,completed_event=?3 WHERE id=?1", params![id,state,seq])?; tx.execute("UPDATE effects SET status='held',reason='execution ended before publication' WHERE execution_id=?1 AND status='proposed'", [&id])?; tx.commit()?; Ok(()) }).await }
/// Commits an execution whose remote side may have run but whose result /// was not observed. Repeating the same operation/reason is idempotent; /// no uncertain operation is replayed or silently treated as a failure. pub async fn finish_unknown(&self, id: &OperationId, reason: impl Into<String>) -> Result<()> { let id = id.to_string(); let reason = reason.into(); ensure!( !reason.trim().is_empty() && reason.len() <= 65_536, "unknown execution requires a bounded reason" ); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; let (session, prior, completed): (String, String, Option<i64>) = tx.query_row( "SELECT session_id,state,completed_event FROM executions WHERE id=?1", [&id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), )?; if let Some(completed) = completed { if prior == "unknown" { let payload: Value = tx.query_row( "SELECT payload FROM events WHERE seq=?1", [completed], |r| serde_json::from_str(&r.get::<_, String>(0)?).map_err(|error| { rusqlite::Error::FromSqlConversionFailure( 0, rusqlite::types::Type::Text, Box::new(error), ) }), )?; if payload.get("reason").and_then(Value::as_str) == Some(reason.as_str()) { tx.commit()?; return Ok(()); } } bail!("execution already completed") } let seq = event( &tx, &session, "execution.unknown", &json!({"execution_id": id, "state": "unknown", "reason": reason}), )?; tx.execute( "UPDATE executions SET state='unknown',completed_event=?2 WHERE id=?1", params![id, seq], )?; tx.execute( "UPDATE effects SET status='held',reason='execution outcome unknown; publication withheld' WHERE execution_id=?1 AND status='proposed'", [&id], )?; tx.commit()?; Ok(()) }) .await }
pub async fn prepare_send( &self, session: &SessionId, execution: &OperationId, generation: &Generation, request: RequestId, text: String, ) -> Result<PreparedEffect> { ensure!( !text.trim().is_empty() && text.len() <= MAX_OUTPUT, "invalid local message" ); let (session, execution, generation) = ( session.to_string(), execution.to_string(), generation.to_string(), ); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; live(&tx,&session,&execution,&generation)?; let prior: Option<EffectRow> = tx.query_row("SELECT id,text,status,event_id,reason FROM effects WHERE execution_id=?1 AND request_id=?2", params![execution,request.as_str()], |r| Ok((r.get(0)?,r.get(1)?,r.get(2)?,r.get(3)?,r.get(4)?))).optional()?; if let Some((id,prior_text,status,event_id,reason)) = prior { ensure!(prior_text == text, "effect identity reused with changed text"); let status = match status.as_str() { "confirmed" => ReceiptStatus::Confirmed, "held" => ReceiptStatus::Held, _ => bail!("effect review is already in progress") }; return Ok(PreparedEffect::Existing(Receipt {effect_id:id.try_into().map_err(|e:String| anyhow!(e))?,status,event_id,reason})); } let id = EffectId::fresh(); tx.execute("INSERT INTO effects(id,execution_id,request_id,text,status) VALUES(?1,?2,?3,?4,'proposed')", params![id.as_str(),execution,request.as_str(),text])?; tx.commit()?; Ok(PreparedEffect::New(id)) }).await }
pub async fn publish_send( &self, session: &SessionId, execution: &OperationId, generation: &Generation, effect: EffectId, revision: String, review: ReviewedDelivery, ) -> Result<Receipt> { let (session, execution, generation) = ( session.to_string(), execution.to_string(), generation.to_string(), ); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; live(&tx, &session, &execution, &generation)?; let expected: String = tx.query_row( "SELECT hook_revision FROM executions WHERE id=?1", [&execution], |r| r.get(0), )?; ensure!( expected == revision, "delivery policy does not match the execution's pinned revision" ); let text: String = tx.query_row( "SELECT text FROM effects WHERE id=?1 AND execution_id=?2 AND status='proposed'", params![effect.as_str(), execution], |r| r.get(0), )?; tx.execute( "INSERT INTO hook_traces(effect_id,revision,decision,fault) VALUES(?1,?2,?3,?4)", params![ effect.as_str(), revision, serde_json::to_string(&review.decision)?, review.fault ], )?; let receipt = match review.decision { DeliveryDecision::Allow => { let id = event( &tx, &session, "local.sent", &json!({"effect_id":effect,"text":text}), )?; tx.execute( "UPDATE effects SET status='confirmed',event_id=?2 WHERE id=?1", params![effect.as_str(), id], )?; Receipt { effect_id: effect, status: ReceiptStatus::Confirmed, event_id: Some(id), reason: None, } } DeliveryDecision::Hold { reason } => { event( &tx, &session, "delivery.held", &json!({"effect_id":effect,"reason":reason}), )?; tx.execute( "UPDATE effects SET status='held',reason=?2 WHERE id=?1", params![effect.as_str(), reason], )?; Receipt { effect_id: effect, status: ReceiptStatus::Held, event_id: None, reason: Some(reason), } } }; tx.commit()?; Ok(receipt) }) .await }
pub async fn execution_effects(&self, execution: &OperationId) -> Result<Vec<ExecutionEffect>> { let execution = execution.to_string(); self.with(move |db| { let mut query = db.prepare( "SELECT id, text, status, event_id, reason FROM effects WHERE execution_id = ?1 ORDER BY rowid", )?; let rows = query.query_map([execution], |r| { let id_str: String = r.get(0)?; let status_str: String = r.get(2)?; let status = match status_str.as_str() { "confirmed" => ReceiptStatus::Confirmed, "held" => ReceiptStatus::Held, other => return Err(rusqlite::Error::FromSqlConversionFailure( 2, rusqlite::types::Type::Text, Box::new(std::io::Error::new( std::io::ErrorKind::InvalidData, format!("unexpected status {other}"), )), )), }; let effect_id = match EffectId::try_from(id_str) { Ok(id) => id, Err(e) => return Err(rusqlite::Error::FromSqlConversionFailure( 0, rusqlite::types::Type::Text, Box::new(std::io::Error::new(std::io::ErrorKind::InvalidData, e)), )), }; Ok(ExecutionEffect { effect_id, text: r.get(1)?, status, event_id: r.get(3)?, reason: r.get(4)?, }) })?.collect::<rusqlite::Result<Vec<_>>>()?; Ok(rows) }).await }
pub async fn execution_yield( &self, execution: &OperationId, ) -> Result<Option<(Option<f64>, String)>> { let execution = execution.to_string(); self.with(move |db| { let mut query = db.prepare( "SELECT e.payload FROM events e JOIN executions x ON e.session_id = x.session_id AND e.seq > x.started_event AND (x.completed_event IS NULL OR e.seq <= x.completed_event) WHERE x.id = ?1 AND e.kind = 'session.waited' ORDER BY e.seq DESC LIMIT 1", )?; let payload_opt: Option<String> = query.query_row([execution], |r| r.get(0)).optional()?; if let Some(payload_str) = payload_opt { let val: Value = serde_json::from_str(&payload_str) .map_err(|e| rusqlite::Error::FromSqlConversionFailure( 0, rusqlite::types::Type::Text, Box::new(e), ))?; let reason = val.get("reason").and_then(|v| v.as_str()).unwrap_or("").to_string(); let seconds = val.get("seconds").and_then(|v| v.as_f64()); Ok(Some((seconds, reason))) } else { Ok(None) } }).await }
pub async fn wait_session(&self, session: &SessionId, wake_at_ms: Option<i64>) -> Result<()> { let session = session.to_string(); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; tx.execute( "UPDATE sessions SET state='waiting', wake_at_ms=?1 WHERE id=?2", params![wake_at_ms, session], )?; tx.commit()?; Ok(()) }) .await }
pub async fn pending_count(&self, session: &SessionId) -> Result<i64> { let session = session.to_string(); self.with(move |db| { let count: i64 = db.query_row( "SELECT count(*) FROM work_items WHERE session_id=?1 AND cause_kind='input_frame' AND status='pending'", [&session], |r| r.get(0), )?; Ok(count) }) .await }
pub async fn wait( &self, session: &SessionId, execution: &OperationId, generation: &Generation, seconds: Option<f64>, reason: String, ) -> Result<Value> { ensure!( !reason.trim().is_empty() && reason.len() <= 2048, "invalid wait reason" ); ensure!( seconds.is_none_or(|s| s.is_finite() && s > 0.0 && s <= 600.0), "invalid wait duration" ); let (session_str, execution_str, generation_str) = ( session.to_string(), execution.to_string(), generation.to_string(), ); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; live(&tx,&session_str,&execution_str,&generation_str)?; let pending = { // Unclaimed durable inputs only: the claim is the work item's status, so this // count is the same before and after a restart. let mut query = tx.prepare( "SELECT count(*) FROM work_items WHERE session_id=?1 AND cause_kind='input_frame' AND status='pending'", )?; let count = query.query_row([&session_str], |r| r.get::<_, i64>(0))?; count as usize }; let deadline = seconds.map(|s| now_ms().map(|n| n + (s * 1000.0).ceil() as i64)).transpose()?; tx.execute("UPDATE sessions SET state=?2,wake_at_ms=?3 WHERE id=?1", params![session_str,if pending > 0 {"ready"} else {"waiting"},if pending > 0 {None} else {deadline}])?; tx.execute("UPDATE executions SET state='yielded' WHERE id=?1", [&execution_str])?; let result = json!({"pending_input":pending > 0,"wake_at_ms":if pending > 0 {None} else {deadline},"reason":reason,"seconds":seconds}); event(&tx,&session_str,"session.waited",&result)?; tx.commit()?; Ok(result) }).await }
/// A scheduler will call this on startup and at the next deadline. Committing /// the ready state and wake event together makes repeated ticks idempotent. pub async fn wake_due(&self, now: i64) -> Result<Vec<String>> { ensure!(now >= 0, "invalid clock"); self.with(move |db| { let tx = db.transaction_with_behavior(TransactionBehavior::Immediate)?; let sessions: Vec<String> = { let mut query = tx.prepare("SELECT id FROM sessions WHERE state='waiting' AND wake_at_ms<=?1")?; let rows = query .query_map([now], |r| r.get(0))? .collect::<rusqlite::Result<_>>()?; rows }; for session in &sessions { tx.execute( "UPDATE sessions SET state='ready',wake_at_ms=NULL WHERE id=?1", [session], )?; event(&tx, session, "session.woke", &json!({"reason":"deadline"}))?; let work_id = format!("wake-{}-{}", session, now); let cause_payload = serde_json::to_string(&json!({"deadline_ms": now, "reason": "deadline"}))?; let _ = tx.execute( "INSERT OR IGNORE INTO work_items ( id, session_id, purpose, cause_kind, cause_payload, readiness, evidence_refs, impact, urgency_ms, aging_ms, reporter_id, origin_kind, status, created_ms ) VALUES (?1, ?2, 'deadline_wake', 'deadline_wake', ?3, 'ready', '[]', 2, ?4, 0, 'runtime', 'system', 'pending', ?5)", params![work_id, session, cause_payload, now, now], ); } tx.commit()?; Ok(sessions) }) .await }
/// Check if a session has a wake event that has not yet had an execution started. pub async fn has_unserviced_wake(&self, session: &SessionId) -> Result<bool> { let session_str = session.to_string(); self.with(move |db| { let mut stmt = db.prepare( "SELECT 1 FROM events WHERE session_id=?1 AND kind='session.woke' AND seq > COALESCE((SELECT MAX(seq) FROM events WHERE session_id=?1 AND kind='execution.started'), 0) LIMIT 1", )?; let exists = stmt.exists([&session_str])?; Ok(exists) }) .await }
pub async fn insert_work_item(&self, item: &WorkItem) -> Result<bool> { let item = item.clone(); self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let inserted = insert_work_item_tx(&tx, &item)?; tx.commit()?; Ok(inserted) }) .await }
pub async fn get_work_item(&self, id: &str) -> Result<Option<WorkItem>> { let id = id.to_string(); self.with(move |conn| { let mut stmt = conn.prepare( "SELECT id, session_id, purpose, cause_kind, cause_payload, readiness, blocked_reason, evidence_refs, impact, urgency_ms, aging_ms, reporter_id, origin_kind, status, attempt_id, claimed_ms, completed_ms, disposition, failed_ms, error, created_ms FROM work_items WHERE id = ?1", )?; let item = stmt.query_row(params![id], parse_work_item_row).optional()?; Ok(item) }) .await }
pub async fn get_work_item_disposition(&self, id: &str) -> Result<Option<String>> { let id = id.to_string(); self.with(move |conn| { let disp: Option<String> = conn .query_row( "SELECT disposition FROM work_items WHERE id = ?1", params![id], |r| r.get(0), ) .optional()?; Ok(disp) }) .await }
/// The one eligibility contract for runnable work, as a SQL predicate. /// `alias` qualifies the columns for joined queries; `now` is the bound /// wall-clock parameter string. The per-session selector, the driver's /// session query, and claim_work_item all share this contract; changing /// eligibility means changing this function and nothing else. fn eligibility_where(alias: Option<&str>, now: &str) -> String { let q = |col: &str| match alias { Some(a) => format!("{a}.{col}"), None => col.to_owned(), }; format!( "{} IN ('pending', 'deferred') AND {} = 'ready' \ AND ({} IS NULL OR {} <= {now})", q("status"), q("readiness"), q("revisit_deadline_ms"), q("revisit_deadline_ms"), ) }
const WORK_ITEM_COLUMNS: &str = "id, session_id, purpose, cause_kind, cause_payload, readiness, blocked_reason, evidence_refs, impact, urgency_ms, aging_ms, reporter_id, origin_kind, status, attempt_id, claimed_ms, completed_ms, disposition, failed_ms, error, created_ms";
/// Every work item satisfying the single eligibility contract at `now_ms`. /// `session = None` selects across all sessions (the driver's scope). pub async fn eligible_work_items( &self, session: Option<&SessionId>, now_ms: i64, ) -> Result<Vec<WorkItem>> { let session = session.map(|s| s.to_string()); self.with(move |conn| { let (sql, params) = match &session { Some(s) => ( format!( "SELECT {} FROM work_items WHERE session_id = ?1 AND ({}) ORDER BY created_ms ASC", Self::WORK_ITEM_COLUMNS, Self::eligibility_where(None, "?2") ), (s.clone(), now_ms), ), None => ( format!( "SELECT {} FROM work_items WHERE ({}) ORDER BY created_ms ASC", Self::WORK_ITEM_COLUMNS, Self::eligibility_where(None, "?1") ), (String::new(), now_ms), ), }; let mut stmt = conn.prepare(&sql)?; let rows = if session.is_some() { stmt.query_map(params![params.0, params.1], parse_work_item_row)? .collect::<rusqlite::Result<Vec<_>>>()? } else { stmt.query_map(params![params.1], parse_work_item_row)? .collect::<rusqlite::Result<Vec<_>>>()? }; Ok(rows) }) .await }
pub async fn pending_work_items(&self, session: &SessionId) -> Result<Vec<WorkItem>> { let now = now_ms()?; self.eligible_work_items(Some(session), now).await }
pub async fn all_pending_work_items(&self) -> Result<Vec<WorkItem>> { let now = now_ms()?; self.eligible_work_items(None, now).await }
pub async fn sessions_with_pending_work(&self) -> Result<Vec<String>> { let where_clause = Self::eligibility_where(Some("w"), "?1"); let now = now_ms()?; self.with(move |conn| { let mut stmt = conn.prepare(&format!( "SELECT DISTINCT w.session_id FROM work_items w JOIN sessions s ON s.id = w.session_id WHERE ({where_clause}) AND s.state = 'ready'" ))?; let rows = stmt .query_map(params![now], |r| r.get(0))? .collect::<rusqlite::Result<Vec<_>>>()?; Ok(rows) }) .await }
pub async fn all_sessions(&self) -> Result<Vec<SessionId>> { self.with(|conn| { let mut stmt = conn.prepare("SELECT id FROM sessions ORDER BY id")?; let ids = stmt .query_map([], |row| row.get::<_, String>(0))? .collect::<rusqlite::Result<Vec<_>>>()?; ids.into_iter() .map(|id| SessionId::try_from(id).map_err(|e| anyhow!(e))) .collect() }) .await }
pub async fn claim_work_item( &self, id: &str, attempt_id: &OperationId, now_ms: i64, ) -> Result<ClaimResult> { let id = id.to_string(); let attempt_id = attempt_id.clone(); let attempt_str = attempt_id.to_string(); self.with(move |conn| { // First check current state, including the durable revisit gate let current: Option<(String, Option<String>, String, Option<i64>)> = conn .query_row( "SELECT status, attempt_id, readiness, revisit_deadline_ms FROM work_items WHERE id = ?1", params![id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)), ) .optional()?;
let Some((status, holder, readiness, revisit)) = current else { return Ok(ClaimResult::NotFound); };
match status.as_str() { "completed" => return Ok(ClaimResult::AlreadyCompleted), "failed" => return Ok(ClaimResult::AlreadyFailed), "claimed" => { let holder_id = holder .map(|s| s.parse().unwrap_or_else(|_| OperationId::fresh())) .unwrap_or_else(OperationId::fresh); return Ok(ClaimResult::AlreadyClaimed { holder: holder_id }); } "pending" | "deferred" => { if readiness != "ready" { return Ok(ClaimResult::NotReady); } if revisit.is_some_and(|d| d > now_ms) { // A bounded revisit deadline gates the item until it is due. return Ok(ClaimResult::NotReady); } } _ => return Ok(ClaimResult::NotReady), }
let rows = conn.execute( &format!( "UPDATE work_items SET status = 'claimed', attempt_id = ?2, claimed_ms = ?3 WHERE id = ?1 AND ({})", Self::eligibility_where(None, "?4") ), params![id, attempt_str, now_ms, now_ms], )?;
if rows > 0 { Ok(ClaimResult::Claimed { attempt_id }) } else { Ok(ClaimResult::NotReady) } }) .await }
pub async fn complete_work_item( &self, id: &str, attempt_id: &OperationId, disposition: &str, now_ms: i64, ) -> Result<RetireResult> { let id = id.to_string(); let attempt_str = attempt_id.to_string(); let disp = disposition.to_string(); self.with(move |conn| { // Check current state if let Some(result) = retire_preflight(conn, &id, &attempt_str)? { return Ok(result); }
let rows = conn.execute( "UPDATE work_items SET status = 'completed', disposition = ?3, completed_ms = ?4 WHERE id = ?1 AND status = 'claimed' AND attempt_id = ?2", params![id, attempt_str, disp, now_ms], )?;
if rows > 0 { Ok(RetireResult::Retired) } else { Ok(RetireResult::NotClaimed) } }) .await }
pub async fn fail_work_item( &self, id: &str, attempt_id: &OperationId, error: &str, now_ms: i64, ) -> Result<RetireResult> { let id = id.to_string(); let attempt_str = attempt_id.to_string(); let err = error.to_string(); self.with(move |conn| { // Check current state if let Some(result) = retire_preflight(conn, &id, &attempt_str)? { return Ok(result); }
let rows = conn.execute( "UPDATE work_items SET status = 'failed', error = ?3, failed_ms = ?4 WHERE id = ?1 AND status = 'claimed' AND attempt_id = ?2", params![id, attempt_str, err, now_ms], )?;
if rows > 0 { Ok(RetireResult::Retired) } else { Ok(RetireResult::NotClaimed) } }) .await }
/// Defer a work item while preserving the attempt history. The item becomes /// re-claimable after a wake event. Use when maintenance yields mid-turn. pub async fn defer_work_item( &self, id: &str, attempt_id: &OperationId, reason: &str, ) -> Result<RetireResult> { let id = id.to_string(); let attempt_str = attempt_id.to_string(); let reason = reason.to_string(); self.with(move |conn| { // Check current state if let Some(result) = retire_preflight(conn, &id, &attempt_str)? { return Ok(result); }
// Transition to deferred, clearing the claim but keeping readiness=ready // so it can be re-claimed after the wake let rows = conn.execute( "UPDATE work_items SET status = 'deferred', error = ?3, attempt_id = NULL, claimed_ms = NULL WHERE id = ?1 AND status = 'claimed' AND attempt_id = ?2", params![id, attempt_str, reason], )?;
if rows > 0 { Ok(RetireResult::Retired) } else { Ok(RetireResult::NotClaimed) } }) .await }
/// Retire a claim whose slice ran out so the continuation can take the same /// work. The item returns to `pending`, not to `deferred`: the scheduler /// selects pending work, and a task that merely ran out of steps is ready to /// run again. Unlike `defer_work_item` nothing went wrong, so no error is /// recorded - the yield reason lives in the durable slice record. pub async fn yield_work_item( &self, id: &str, attempt_id: &OperationId, ) -> Result<RetireResult> { let id = id.to_string(); let attempt_str = attempt_id.to_string(); self.with(move |conn| { if let Some(result) = retire_preflight(conn, &id, &attempt_str)? { return Ok(result); }
let rows = conn.execute( "UPDATE work_items SET status = 'pending', attempt_id = NULL, claimed_ms = NULL WHERE id = ?1 AND status = 'claimed' AND attempt_id = ?2", params![id, attempt_str], )?;
if rows > 0 { Ok(RetireResult::Retired) } else { Ok(RetireResult::NotClaimed) } }) .await }
/// Defer a work item by retiring its claim and persisting the truthful /// revisit gate: blocked readiness when a dependency is named, and a /// bounded revisit deadline when one is declared. A deferred item is /// re-eligible only through the shared eligibility contract. pub async fn defer_work_item_with_revisit( &self, id: &str, attempt_id: &OperationId, reason: &str, gate: &RevisitGate, ) -> Result<RetireResult> { if let Some(d) = gate.deadline_ms { ensure!( d > 0, "revisit deadline must be a positive wall-clock ms value" ); } let id = id.to_string(); let attempt_str = attempt_id.to_string(); let reason = reason.to_string(); let (readiness_str, blocked_reason) = gate.readiness_str_and_reason(); let deadline = gate.deadline_ms; let revisit_reason = reason.clone(); self.with(move |conn| { if let Some(result) = retire_preflight(conn, &id, &attempt_str)? { return Ok(result); }
let rows = conn.execute( "UPDATE work_items SET status = 'deferred', error = ?3, attempt_id = NULL, claimed_ms = NULL, readiness = ?4, blocked_reason = ?5, revisit_deadline_ms = ?6, revisit_reason = ?7 WHERE id = ?1 AND status = 'claimed' AND attempt_id = ?2", params![ id, attempt_str, reason, readiness_str, blocked_reason, deadline, revisit_reason ], )?;
if rows > 0 { Ok(RetireResult::Retired) } else { Ok(RetireResult::NotClaimed) } }) .await }
/// Persisted revisit state for a work item: its deadline and reason. pub async fn get_work_item_revisit(&self, id: &str) -> Result<Option<RevisitState>> { let id = id.to_string(); self.with(move |conn| { let state: Option<(Option<i64>, Option<String>)> = conn .query_row( "SELECT revisit_deadline_ms, revisit_reason FROM work_items WHERE id = ?1", params![id], |r| Ok((r.get(0)?, r.get(1)?)), ) .optional()?; Ok(state.map(|(deadline_ms, reason)| RevisitState { deadline_ms, reason, })) }) .await }
/// Release a dependency-blocked work item: merge new evidence refs and /// set readiness back to ready. Any persisted revisit deadline still /// gates eligibility through the shared contract. Returns false when the /// item does not exist or is not currently blocked. pub async fn release_blocked_work_item( &self, id: &str, evidence_refs: &[String], ) -> Result<bool> { let id = id.to_string(); let new_refs = evidence_refs.to_vec(); self.with(move |conn| { let existing: Option<(String, Option<String>)> = conn .query_row( "SELECT readiness, blocked_reason FROM work_items WHERE id = ?1", params![id], |r| Ok((r.get(0)?, r.get(1)?)), ) .optional()?;
let Some((readiness, _)) = existing else { return Ok(false); }; if readiness != "blocked" { return Ok(false); }
let refs_str: String = conn.query_row( "SELECT evidence_refs FROM work_items WHERE id = ?1", params![id], |r| r.get(0), )?; let (refs, _) = merge_evidence_refs(&refs_str, &new_refs); let merged = serde_json::to_string(&refs)?;
let rows = conn.execute( "UPDATE work_items SET readiness = 'ready', blocked_reason = NULL, evidence_refs = ?2 WHERE id = ?1 AND readiness = 'blocked'", params![id, merged], )?; Ok(rows > 0) }) .await }
/// Earliest not-yet-due revisit deadline among this session's pending or /// deferred work, for scheduling a wake. None when no deadline pends. pub async fn earliest_revisit_deadline( &self, session: &SessionId, now_ms: i64, ) -> Result<Option<i64>> { let session = session.to_string(); self.with(move |conn| { let earliest: Option<Option<i64>> = conn .query_row( "SELECT MIN(revisit_deadline_ms) FROM work_items WHERE session_id = ?1 AND revisit_deadline_ms > ?2 AND status IN ('pending', 'deferred')", params![session, now_ms], |r| r.get(0), ) .optional()?; Ok(earliest.flatten()) }) .await }
/// Release every work item currently claimed by `attempt_id` back to /// 'pending' with readiness unchanged, clearing attempt_id and claimed_ms, /// recording `reason` in `error`. Returns the released work item ids. /// A claimed item whose attempt_id differs must not be touched. pub async fn release_attempt_claims( &self, attempt_id: &OperationId, reason: &str, ) -> Result<Vec<String>> { let attempt_str = attempt_id.to_string(); let reason = reason.to_string(); self.with(move |conn| { let ids: Vec<String> = { let mut stmt = conn.prepare( "SELECT id FROM work_items WHERE attempt_id = ?1 AND status = 'claimed'", )?; let rows = stmt .query_map(params![attempt_str], |r| r.get(0))? .collect::<rusqlite::Result<Vec<_>>>()?; rows }; if ids.is_empty() { return Ok(Vec::new()); } let placeholders = (1..=ids.len()) .map(|i| format!("?{i}")) .collect::<Vec<_>>() .join(", "); let mut bind_params: Vec<&dyn rusqlite::ToSql> = Vec::new(); for id in &ids { bind_params.push(id); } bind_params.push(&reason); bind_params.push(&attempt_str); let sql = format!( "UPDATE work_items SET status = 'pending', attempt_id = NULL, claimed_ms = NULL, error = ?{} WHERE status = 'claimed' AND attempt_id = ?{} AND id IN ({placeholders})", ids.len() + 1, ids.len() + 2 ); let released = conn.execute(&sql, bind_params.as_slice())?; debug_assert_eq!( released, ids.len(), "attempt claims changed between select and update" ); Ok(ids) }) .await }
/// Attach additional evidence refs to an existing work item without touching /// its status, ownership, or completion state. Safe to call on claimed/completed items. pub async fn attach_evidence_to_work_item( &self, work_item_id: &str, new_refs: &[String], ) -> Result<bool> { let id = work_item_id.to_string(); let new_refs = new_refs.to_vec(); self.with(move |conn| { let existing: Option<String> = conn .query_row( "SELECT evidence_refs FROM work_items WHERE id = ?1", params![id], |r| r.get(0), ) .optional()?;
let Some(existing_str) = existing else { return Ok(false); };
let (refs, changed) = merge_evidence_refs(&existing_str, &new_refs);
if changed { let merged = serde_json::to_string(&refs)?; conn.execute( "UPDATE work_items SET evidence_refs = ?2 WHERE id = ?1", params![id, merged], )?; } Ok(true) }) .await }
/// Atomically create an origination event AND its runnable work item in one transaction. /// Returns (work_item_id, event_sequence) on success. /// The operation_id is used as idempotency key bound to session+execution. /// Persists the canonical request -> work association durably. pub async fn commit_origination_with_work( &self, session_id: &SessionId, operation_id: &OperationId, work_item: WorkItem, origination_payload: Value, ) -> Result<(WorkItemId, i64)> { ensure!( work_item.session_id == session_id.as_str(), "work item session {} does not match request session {}", work_item.session_id, session_id ); let session = session_id.to_string(); let op_str = operation_id.to_string(); let source_key = format!("origination:{op_str}"); self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let now = now_ms()?;
// Check if this request identity has already been recorded let existing_record: Option<(String, i64, String)> = tx .query_row( "SELECT work_item_id, event_seq, payload FROM origination_requests WHERE session_id = ?1 AND operation_id = ?2", params![session, op_str], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), ) .optional()?;
if let Some((recorded_work_id, recorded_seq, recorded_payload_str)) = existing_record { // Verify payload matches (changed payload on reused request identity is a conflict) let recorded_val: Value = serde_json::from_str(&recorded_payload_str)?; if recorded_val != origination_payload { bail!( "request identity reused with changed payload: expected {:?}, got {:?}", recorded_val, origination_payload ); }
// Verify that recorded work item exists and belongs to this session let item_session: Option<String> = tx .query_row( "SELECT session_id FROM work_items WHERE id = ?1", params![recorded_work_id], |r| r.get(0), ) .optional()?;
let item_session = item_session.ok_or_else(|| { anyhow!( "acknowledged work identity {} does not exist in work_items", recorded_work_id ) })?;
ensure!( item_session == session, "recorded work item {} belongs to unrelated session {}", recorded_work_id, item_session );
// Verify that the origination event exists for this session let event_exists: bool = tx .query_row( "SELECT 1 FROM events WHERE seq = ?1 AND session_id = ?2", params![recorded_seq, session], |_| Ok(true), ) .optional()? .unwrap_or(false);
ensure!( event_exists, "recorded event seq {} does not exist for session {}", recorded_seq, session );
tx.commit()?; return Ok((WorkItemId(recorded_work_id), recorded_seq)); }
// Also check if an event with this source_key exists let existing_event: Option<i64> = tx .query_row( "SELECT seq FROM events WHERE session_id = ?1 AND source_key = ?2", params![session, source_key], |r| r.get(0), ) .optional()?;
if let Some(seq) = existing_event { let prior_payload_str: String = tx.query_row( "SELECT payload FROM events WHERE seq = ?1", params![seq], |r| r.get(0), )?; let prior_val: Value = serde_json::from_str(&prior_payload_str)?; if prior_val != origination_payload { bail!( "event source key reused with changed payload: expected {:?}, got {:?}", prior_val, origination_payload ); }
let existing_id: String = tx .query_row( "SELECT id FROM work_items WHERE id = ?1 AND session_id = ?2", params![work_item.id, session], |r| r.get(0), ) .map_err(|_| anyhow!("acknowledged work identity {} does not exist in work_items", work_item.id))?;
tx.execute( "INSERT INTO origination_requests (session_id, operation_id, work_item_id, event_seq, payload, created_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", params![session, op_str, existing_id, seq, prior_payload_str, now], )?;
tx.commit()?; return Ok((WorkItemId(existing_id), seq)); }
// Disallow reusing an existing work item id for a brand new origination request let existing_work: Option<(String, String)> = tx .query_row( "SELECT id, session_id FROM work_items WHERE id = ?1", params![work_item.id], |r| Ok((r.get(0)?, r.get(1)?)), ) .optional()?;
if let Some((existing_id, existing_sess)) = existing_work { bail!( "work item {} already exists in session {}; refusing to reuse work identity for new origination request", existing_id, existing_sess ); }
insert_work_item_tx(&tx, &work_item)?;
let payload_str = serde_json::to_string(&origination_payload)?; tx.execute( "INSERT INTO events (session_id, kind, source_key, payload, created_ms) VALUES (?1, 'scratchpad', ?2, ?3, ?4)", params![session, source_key, payload_str, now], )?; let seq = tx.last_insert_rowid();
tx.execute( "INSERT INTO origination_requests (session_id, operation_id, work_item_id, event_seq, payload, created_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", params![session, op_str, work_item.id, seq, payload_str, now], )?;
tx.commit()?; Ok((WorkItemId(work_item.id), seq)) }) .await }
pub async fn commit_maintenance_assessment( &self, session_id: &SessionId, operation_id: &OperationId, work_item_id: &str, outcome: &str, assessment_payload: Value, ) -> Result<i64> { // Closed outcome parse at the durable boundary: unknown tags, missing // variant fields or contradictory references are typed errors. let assessment = MaintenanceAssessment::parse(outcome, &assessment_payload)?; let session = session_id.to_string(); let op_str = operation_id.to_string(); let wid = work_item_id.to_string(); let source_key = format!("assessment:{op_str}"); self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let now = now_ms()?;
let existing: Option<(i64, String)> = tx .query_row( "SELECT seq, payload FROM events WHERE session_id = ?1 AND source_key = ?2", params![session, source_key], |r| Ok((r.get(0)?, r.get(1)?)), ) .optional()?;
if let Some((seq, prior_payload_str)) = existing { let prior_val: Value = serde_json::from_str(&prior_payload_str)?; let canonical = Self::canonical_assessment_payload( &session, &wid, prior_val .get("attempt_id") .and_then(|v| v.as_str()) .unwrap_or_default(), &assessment, ); if prior_val != canonical { bail!( "assessment request identity reused with changed payload: expected {:?}, got {:?}", prior_val, canonical ); } tx.commit()?; return Ok(seq); }
let binding: Option<(String, String, Option<String>)> = tx .query_row( "SELECT session_id, status, attempt_id FROM work_items WHERE id = ?1", params![wid], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), ) .optional()?;
let Some((item_sess, status, holder)) = binding else { bail!("work item {} does not exist for assessment", wid); };
ensure!( item_sess == session, "work item {} belongs to unrelated session {}", wid, item_sess );
// A submission binds to the work item AND its current owning // attempt. An assessment without an owning attempt cannot be // tied to the evidence the attempt observed. let attempt_str = holder.ok_or_else(|| { anyhow!( "work item {} is not currently claimed by an owning attempt (status {}); assessment must be submitted during the owning attempt", wid, status ) })?; ensure!( status == "claimed", "work item {} has no active owning attempt (status {})", wid, status );
if let MaintenanceAssessment::Activated { candidate_revision, check_receipts, } = &assessment { Self::validate_activated_assessment( &tx, &session, &wid, candidate_revision, check_receipts, )?; }
let canonical = Self::canonical_assessment_payload(&session, &wid, &attempt_str, &assessment); let canonical_str = serde_json::to_string(&canonical)?; let assessment_str = serde_json::to_string(&assessment)?;
tx.execute( "INSERT INTO events (session_id, kind, source_key, payload, created_ms) VALUES (?1, 'scratchpad', ?2, ?3, ?4)", params![session, source_key, canonical_str, now], )?; let seq = tx.last_insert_rowid();
tx.execute( "UPDATE work_items SET disposition = ?2 WHERE id = ?1", params![wid, assessment_str], )?;
tx.execute( "INSERT INTO maintenance_assessments (session_id, operation_id, work_item_id, attempt_id, assessment, event_seq, created_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)", params![session, op_str, wid, attempt_str, assessment_str, seq, now], )?;
tx.commit()?; Ok(seq) }) .await }
fn canonical_assessment_payload( session: &str, work_item_id: &str, attempt_id: &str, assessment: &MaintenanceAssessment, ) -> Value { let mut canonical = serde_json::to_value(assessment).unwrap_or_else(|_| json!({})); canonical["type"] = json!("maintenance.assessed"); canonical["session_id"] = json!(session); canonical["work_item_id"] = json!(work_item_id); canonical["attempt_id"] = json!(attempt_id); canonical }
/// An Activated assessment is a deployment fact, not a claim: the /// candidate must be a real release, a host activation record must bind /// that revision to THIS work item, and every check receipt must resolve /// to a committed event. Proposed checks are not executed receipts. fn validate_activated_assessment( tx: &Transaction<'_>, session: &str, work_item_id: &str, candidate_revision: &str, check_receipts: &[String], ) -> Result<()> { ensure!( !check_receipts.is_empty(), "activated assessment for work item {work_item_id} requires at least one executed check receipt" ); let release_exists: Option<i64> = tx .query_row( "SELECT 1 FROM releases WHERE id = ?1", params![candidate_revision], |r| r.get(0), ) .optional()?; ensure!( release_exists.is_some(), "candidate revision {candidate_revision:?} is not a real release known to this host" );
let activations: Vec<(i64, Value)> = { let mut stmt = tx.prepare( "SELECT seq, payload FROM events WHERE session_id = ?1 AND kind = 'hooks.activated' ORDER BY seq", )?; let rows = stmt .query_map(params![session], |r| { let seq: i64 = r.get(0)?; let payload_str: String = r.get(1)?; let payload: Value = serde_json::from_str(&payload_str).map_err(|e| { rusqlite::Error::FromSqlConversionFailure( 1, rusqlite::types::Type::Text, e.into(), ) })?; Ok((seq, payload)) })? .collect::<rusqlite::Result<Vec<_>>>()?; rows }; let record = activations .iter() .rev() .find(|(_, payload)| { payload.get("work_item_id").and_then(|v| v.as_str()) == Some(work_item_id) }) .ok_or_else(|| { anyhow!( "no host activation record binds work item {work_item_id} in session {session}; an activation without this link is not proof" ) })?; let record_revision = record .1 .get("revision") .and_then(|v| v.as_str()) .unwrap_or_default(); ensure!( record_revision == candidate_revision, "activation record for work item {work_item_id} has revision {record_revision:?}, not candidate revision {candidate_revision:?}" );
for receipt in check_receipts { let seq = receipt .strip_prefix("event:") .and_then(|rest| rest.parse::<i64>().ok()) .ok_or_else(|| { anyhow!( "check receipt {receipt:?} is not a resolvable event reference (event:<seq>); proposed checks are not executed receipts" ) })?; let exists: Option<i64> = tx .query_row( "SELECT 1 FROM events WHERE seq = ?1 AND session_id = ?2", params![seq, session], |r| r.get(0), ) .optional()?; ensure!( exists.is_some(), "check receipt {receipt:?} does not resolve to a committed event in this session" ); } Ok(()) }
pub async fn reopen_work_item( &self, id: &str, reason: &str, new_refs: &[String], ) -> Result<bool> { let id = id.to_string(); let reason = reason.to_string(); let new_refs = new_refs.to_vec(); self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let existing_refs_str: Option<String> = tx .query_row( "SELECT evidence_refs FROM work_items WHERE id = ?1 AND status IN ('completed', 'failed', 'deferred')", params![id], |r| r.get(0), ) .optional()?;
let Some(refs_str) = existing_refs_str else { return Ok(false); };
let (refs, _) = merge_evidence_refs(&refs_str, &new_refs); let merged_refs = serde_json::to_string(&refs)?;
let rows = tx.execute( "UPDATE work_items SET status = 'pending', readiness = 'ready', attempt_id = NULL, claimed_ms = NULL, completed_ms = NULL, disposition = NULL, failed_ms = NULL, error = ?2, evidence_refs = ?3 WHERE id = ?1 AND status IN ('completed', 'failed', 'deferred')", params![id, reason, merged_refs], )?; tx.commit()?; Ok(rows > 0) }) .await }
pub async fn create_follow_up_work( &self, parent_work_id: &str, follow_up: WorkItem, ) -> Result<WorkItemId> { let parent_id = parent_work_id.to_string(); self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let parent_exists: bool = tx .query_row( "SELECT 1 FROM work_items WHERE id = ?1", params![parent_id], |_| Ok(true), ) .optional()? .unwrap_or(false);
ensure!(parent_exists, "parent work item not found");
insert_work_item_tx(&tx, &follow_up)?; tx.commit()?; Ok(WorkItemId(follow_up.id)) }) .await }
pub async fn recover_startup_work_items(&self, now_ms: i64) -> Result<Vec<String>> { self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; tx.execute( "UPDATE work_items SET status = 'failed', error = 'host restarted; attempt orphaned', failed_ms = ?1 WHERE status = 'claimed'", params![now_ms], )?; let mut stmt = tx.prepare(&format!( "SELECT DISTINCT session_id FROM work_items WHERE ({})", Self::eligibility_where(None, "?1") ))?; let rows = stmt.query_map(params![now_ms], |r| r.get(0))?; let mut sids = Vec::new(); for r in rows { sids.push(r?); } drop(stmt); tx.commit()?; Ok(sids) }) .await }
/// Hook revision recorded for a committed execution, so historical /// provenance resolves from the linked source execution instead of the /// current environment. None when the execution id is unknown. pub async fn get_execution_revision(&self, execution_id: &str) -> Result<Option<String>> { let exec_id = execution_id.to_string(); self.with(move |db| { let rev: Option<String> = db .query_row( "SELECT hook_revision FROM executions WHERE id = ?1", params![exec_id], |r| r.get(0), ) .optional()?; Ok(rev) }) .await }
/// Sessions whose committed events reach past the recovery observer /// cursor, for startup recovery and the periodic driver. pub async fn sessions_with_unobserved_events(&self) -> Result<Vec<SessionId>> { self.with(move |db| { let mut stmt = db.prepare( "SELECT DISTINCT e.session_id FROM events e LEFT JOIN observer_cursors c ON c.partition = ('observer:recovered:' || e.session_id) WHERE e.seq > COALESCE(c.last_seq, 0)", )?; let rows = stmt .query_map([], |r| { let id_str: String = r.get(0)?; SessionId::try_from(id_str).map_err(|e| { rusqlite::Error::FromSqlConversionFailure( 0, rusqlite::types::Type::Text, e.into(), ) }) })? .collect::<rusqlite::Result<Vec<_>>>()?; Ok(rows) }) .await }
pub async fn get_observer_cursor(&self, partition: &str) -> Result<i64> { let partition = partition.to_string(); self.with(move |conn| { let last_seq: Option<i64> = conn .query_row( "SELECT last_seq FROM observer_cursors WHERE partition = ?1", params![partition], |r| r.get(0), ) .optional()?; Ok(last_seq.unwrap_or(0)) }) .await }
pub async fn advance_observer_cursor(&self, partition: &str, new_seq: i64) -> Result<()> { let partition = partition.to_string(); self.with(move |conn| { let now = now_ms().unwrap_or(0); conn.execute( "INSERT INTO observer_cursors (partition, last_seq, updated_ms) VALUES (?1, ?2, ?3) ON CONFLICT(partition) DO UPDATE SET last_seq = CASE WHEN excluded.last_seq > observer_cursors.last_seq THEN excluded.last_seq ELSE observer_cursors.last_seq END, updated_ms = ?3", params![partition, new_seq, now], )?; Ok(()) }) .await }
pub async fn get_observer_case(&self, id: &str) -> Result<Option<ObserverCaseRecord>> { let id = id.to_string(); self.with(move |conn| { let res = conn .query_row( "SELECT id, session_id, operation, component_revision, signature, workspace, count, first_seen_ms, last_seen_ms, evidence_refs, work_item_id FROM observer_cases WHERE id = ?1", params![id], |r| { let refs_str: String = r.get(9)?; let refs: Vec<String> = serde_json::from_str(&refs_str).unwrap_or_default(); Ok(ObserverCaseRecord { id: r.get(0)?, session_id: r.get(1)?, operation: r.get(2)?, component_revision: r.get(3)?, signature: r.get(4)?, workspace: r.get(5)?, count: r.get(6)?, first_seen_ms: r.get(7)?, last_seen_ms: r.get(8)?, evidence_refs: refs, work_item_id: r.get(10)?, }) }, ) .optional()?; Ok(res) }) .await }
pub async fn commit_observation_case_with_cursor( &self, partition: &str, new_cursor_seq: i64, occurrence_id: &OccurrenceId, case: ObserverCaseRecord, work_item: Option<WorkItem>, ) -> Result<ObserveResult> { let partition = partition.to_string(); let occurrence = occurrence_id.0.clone(); self.with(move |conn| { let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; let now = now_ms().unwrap_or_else(|_| now_ms().unwrap_or(0));
// Advance cursor (max semantics - never go backward) tx.execute( "INSERT INTO observer_cursors (partition, last_seq, updated_ms) VALUES (?1, ?2, ?3) ON CONFLICT(partition) DO UPDATE SET last_seq = CASE WHEN excluded.last_seq > observer_cursors.last_seq THEN excluded.last_seq ELSE observer_cursors.last_seq END, updated_ms = ?3", params![partition, new_cursor_seq, now], )?;
// Check if this occurrence was already applied to this case let already_applied: bool = tx .query_row( "SELECT 1 FROM observed_occurrences WHERE occurrence_id = ?1 AND case_id = ?2", params![occurrence, case.id], |_| Ok(true), ) .optional()? .unwrap_or(false);
if already_applied { tx.commit()?; return Ok(ObserveResult::AlreadyObserved); }
// Insert work item if provided (INSERT OR IGNORE - won't clobber existing) if let Some(ref item) = work_item { insert_work_item_tx(&tx, item)?; }
// Check if case already exists let existing: Option<(i64, String)> = tx .query_row( "SELECT count, evidence_refs FROM observer_cases WHERE id = ?1", params![case.id], |r| Ok((r.get(0)?, r.get(1)?)), ) .optional()?;
let is_new_case; let work_item_id;
if let Some((_cnt, existing_refs_str)) = existing { is_new_case = false; // Merge evidence refs (deduplicated) let (existing_refs, _) = merge_evidence_refs(&existing_refs_str, &case.evidence_refs); let merged_refs = serde_json::to_string(&existing_refs)?; // Increment count and update evidence tx.execute( "UPDATE observer_cases SET count = count + 1, last_seen_ms = ?2, evidence_refs = ?3, work_item_id = COALESCE(work_item_id, ?4) WHERE id = ?1", params![case.id, case.last_seen_ms, merged_refs, case.work_item_id], )?; work_item_id = tx .query_row( "SELECT work_item_id FROM observer_cases WHERE id = ?1", params![case.id], |r| r.get::<_, Option<String>>(0), ) .optional()? .flatten();
if let Some(ref wid) = work_item_id { let w_refs: Option<String> = tx .query_row( "SELECT evidence_refs FROM work_items WHERE id = ?1", params![wid], |r| r.get(0), ) .optional()?; if let Some(w_refs_str) = w_refs { let (item_refs, w_changed) = merge_evidence_refs(&w_refs_str, &case.evidence_refs); if w_changed { let merged_w = serde_json::to_string(&item_refs)?; tx.execute( "UPDATE work_items SET evidence_refs = ?2 WHERE id = ?1", params![wid, merged_w], )?; } } } } else { is_new_case = true; let refs_json = serde_json::to_string(&case.evidence_refs)?; tx.execute( "INSERT INTO observer_cases ( id, session_id, operation, component_revision, signature, workspace, count, first_seen_ms, last_seen_ms, evidence_refs, work_item_id ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)", params![ case.id, case.session_id, case.operation, case.component_revision, case.signature, case.workspace, case.count, case.first_seen_ms, case.last_seen_ms, refs_json, case.work_item_id ], )?; work_item_id = case.work_item_id.clone(); }
// Record that this occurrence has been applied tx.execute( "INSERT INTO observed_occurrences (occurrence_id, case_id, applied_ms) VALUES (?1, ?2, ?3)", params![occurrence, case.id, now], )?;
tx.commit()?; Ok(ObserveResult::Applied { work_item_id: WorkItemId(work_item_id.unwrap_or_else(|| case.id.clone())), is_new_case, }) }) .await }
pub async fn record_peer_session(&self, peer: PeerSessionRecord) -> Result<()> { self.with(move |db| { db.execute( "INSERT INTO peer_sessions (peer_id, owner_session_id, init_mode, profile, workspace, kernel_target, budget, branch_parent_revision, created_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9) ON CONFLICT(peer_id) DO UPDATE SET profile = coalesce(excluded.profile, peer_sessions.profile), workspace = coalesce(excluded.workspace, peer_sessions.workspace)", rusqlite::params![ peer.peer_id, peer.owner_session_id, peer.init_mode, peer.profile, peer.workspace, peer.kernel_target, peer.budget, peer.branch_parent_revision, peer.created_ms, ], )?; Ok(()) }) .await }
pub async fn get_peer_session(&self, peer_id: &str) -> Result<Option<PeerSessionRecord>> { let peer_id = peer_id.to_string(); self.with(move |db| { let mut stmt = db.prepare( "SELECT peer_id, owner_session_id, init_mode, profile, workspace, kernel_target, budget, branch_parent_revision, created_ms FROM peer_sessions WHERE peer_id = ?1", )?; let mut rows = stmt.query([&peer_id])?; if let Some(row) = rows.next()? { Ok(Some(decode_peer_session(row)?)) } else { Ok(None) } }) .await }
pub async fn list_peer_sessions( &self, owner_session_id: &str, ) -> Result<Vec<PeerSessionRecord>> { let owner_session_id = owner_session_id.to_string(); self.with(move |db| { let mut stmt = db.prepare( "SELECT peer_id, owner_session_id, init_mode, profile, workspace, kernel_target, budget, branch_parent_revision, created_ms FROM peer_sessions WHERE owner_session_id = ?1 ORDER BY created_ms ASC", )?; let rows = stmt.query_map([&owner_session_id], decode_peer_session)?; let mut list = Vec::new(); for r in rows { list.push(r?); } Ok(list) }) .await }
pub async fn admit_peer_request(&self, req: PeerRequestRecord) -> Result<PeerAdmitOutcome> { self.with(move |db| { let mut stmt = db.prepare( "SELECT request_id, sender_session_id, peer_id, text, evidence_refs, status, admitted_at_ms, completed_at_ms, result, error FROM peer_requests WHERE request_id = ?1", )?; let mut rows = stmt.query([&req.request_id])?; if let Some(row) = rows.next()? { return Ok(PeerAdmitOutcome::AlreadyAdmitted(decode_peer_request(row)?)); } drop(rows); drop(stmt);
let evidence_json = serde_json::to_string(&req.evidence_refs)?; db.execute( "INSERT INTO peer_requests (request_id, sender_session_id, peer_id, text, evidence_refs, status, admitted_at_ms) VALUES (?1, ?2, ?3, ?4, ?5, 'admitted', ?6)", rusqlite::params![ req.request_id, req.sender_session_id, req.peer_id, req.text, evidence_json, req.admitted_at_ms, ], )?; Ok(PeerAdmitOutcome::Admitted(req)) }) .await }
pub async fn get_peer_request(&self, request_id: &str) -> Result<Option<PeerRequestRecord>> { let request_id = request_id.to_string(); self.with(move |db| { let mut stmt = db.prepare( "SELECT request_id, sender_session_id, peer_id, text, evidence_refs, status, admitted_at_ms, completed_at_ms, result, error FROM peer_requests WHERE request_id = ?1", )?; let mut rows = stmt.query([&request_id])?; if let Some(row) = rows.next()? { Ok(Some(decode_peer_request(row)?)) } else { Ok(None) } }) .await }
/// Peer requests that are admitted but not yet answered. This is the dispatch /// and startup-recovery worklist: a request leaves it only through the status /// CAS in `complete_peer_request`/`fail_peer_request`, so re-running the /// dispatch pass cannot deliver the same request to a recipient twice. pub async fn admitted_peer_requests(&self) -> Result<Vec<PeerRequestRecord>> { self.with(|db| { let mut stmt = db.prepare( "SELECT request_id, sender_session_id, peer_id, text, evidence_refs, status, admitted_at_ms, completed_at_ms, result, error FROM peer_requests WHERE status = 'admitted' ORDER BY admitted_at_ms, request_id", )?; let mut rows = stmt.query([])?; let mut out = Vec::new(); while let Some(row) = rows.next()? { out.push(decode_peer_request(row)?); } Ok(out) }) .await }
pub async fn complete_peer_request( &self, request_id: &str, result: Value, completed_at_ms: i64, ) -> Result<bool> { let request_id = request_id.to_string(); let result_json = serde_json::to_string(&result)?; self.with(move |db| { let changed = db.execute( "UPDATE peer_requests SET status = 'completed', result = ?1, completed_at_ms = ?2 WHERE request_id = ?3 AND status = 'admitted'", rusqlite::params![result_json, completed_at_ms, request_id], )?; Ok(changed > 0) }) .await }
pub async fn cancel_peer_request(&self, request_id: &str) -> Result<bool> { let request_id = request_id.to_string(); self.with(move |db| { let changed = db.execute( "UPDATE peer_requests SET status = 'cancelled' WHERE request_id = ?1 AND status = 'admitted'", [&request_id], )?; Ok(changed > 0) }) .await }
pub async fn fail_peer_request( &self, request_id: &str, error: &str, completed_at_ms: i64, ) -> Result<bool> { let request_id = request_id.to_string(); let error = error.to_string(); self.with(move |db| { let changed = db.execute( "UPDATE peer_requests SET status = 'failed', error = ?1, completed_at_ms = ?2 WHERE request_id = ?3 AND status = 'admitted'", rusqlite::params![error, completed_at_ms, request_id], )?; Ok(changed > 0) }) .await }}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]pub struct PeerSessionRecord { pub peer_id: String, pub owner_session_id: String, pub init_mode: String, pub profile: Option<String>, pub workspace: Option<String>, pub kernel_target: Option<String>, pub budget: Option<String>, pub branch_parent_revision: Option<String>, pub created_ms: i64,}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]pub struct PeerRequestRecord { pub request_id: String, pub sender_session_id: String, pub peer_id: String, pub text: String, pub evidence_refs: Vec<String>, pub status: String, pub admitted_at_ms: i64, pub completed_at_ms: Option<i64>, pub result: Option<Value>, pub error: Option<String>,}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]pub enum PeerAdmitOutcome { Admitted(PeerRequestRecord), AlreadyAdmitted(PeerRequestRecord),}
fn decode_peer_session(row: &rusqlite::Row<'_>) -> rusqlite::Result<PeerSessionRecord> { Ok(PeerSessionRecord { peer_id: row.get(0)?, owner_session_id: row.get(1)?, init_mode: row.get(2)?, profile: row.get(3)?, workspace: row.get(4)?, kernel_target: row.get(5)?, budget: row.get(6)?, branch_parent_revision: row.get(7)?, created_ms: row.get(8)?, })}
fn decode_peer_request(row: &rusqlite::Row<'_>) -> rusqlite::Result<PeerRequestRecord> { let evidence_json: String = row.get(4)?; let evidence_refs: Vec<String> = serde_json::from_str(&evidence_json).unwrap_or_default(); let result_json: Option<String> = row.get(8)?; let result: Option<Value> = result_json.and_then(|json| serde_json::from_str(&json).ok()); Ok(PeerRequestRecord { request_id: row.get(0)?, sender_session_id: row.get(1)?, peer_id: row.get(2)?, text: row.get(3)?, evidence_refs, status: row.get(5)?, admitted_at_ms: row.get(6)?, completed_at_ms: row.get(7)?, result, error: row.get(9)?, })}
fn insert_work_item_tx(tx: &Transaction<'_>, item: &WorkItem) -> Result<bool> { let cause_payload = match &item.cause { RunCause::InputFrame { event_ids } => json!({"event_ids": event_ids}), RunCause::DeadlineWake { deadline_ms, reason, } => { json!({"deadline_ms": deadline_ms, "reason": reason}) } RunCause::MaintenanceReady { item_id, purpose } => { json!({"item_id": item_id, "purpose": purpose}) } RunCause::CheckDue { check_id, due_ms } => { json!({"check_id": check_id, "due_ms": due_ms}) } RunCause::JobCompleted { job_id, outcome } => { json!({"job_id": job_id, "outcome": outcome}) } }; let cause_payload_str = serde_json::to_string(&cause_payload)?; let (readiness_str, blocked_reason) = match &item.consequence.readiness { Readiness::Ready => ("ready", None), Readiness::Blocked { reason } => ("blocked", Some(reason.clone())), Readiness::Paused { reason } => ("paused", Some(reason.clone())), }; let evidence_refs_str = serde_json::to_string(&item.provenance.source_refs)?; let (status_str, attempt_id, claimed_ms, completed_ms, disposition, failed_ms, error) = match &item.status { WorkStatus::Pending => ("pending", None, None, None, None, None, None), WorkStatus::Claimed { attempt_id, claimed_ms, } => ( "claimed", Some(attempt_id.to_string()), Some(*claimed_ms), None, None, None, None, ), WorkStatus::Completed { completed_ms, disposition, } => ( "completed", None, None, Some(*completed_ms), Some(disposition.clone()), None, None, ), WorkStatus::Failed { failed_ms, error } => ( "failed", None, None, None, None, Some(*failed_ms), Some(error.clone()), ), WorkStatus::Deferred { reason } => ( "deferred", None, None, None, None, None, Some(reason.clone()), ), };
let existing: bool = tx .query_row( "SELECT 1 FROM work_items WHERE id = ?1", params![item.id], |_| Ok(true), ) .optional()? .unwrap_or(false);
if existing { return Ok(false); }
let rows = tx.execute( "INSERT INTO work_items ( id, session_id, purpose, cause_kind, cause_payload, readiness, blocked_reason, evidence_refs, impact, urgency_ms, aging_ms, reporter_id, origin_kind, status, attempt_id, claimed_ms, completed_ms, disposition, failed_ms, error, created_ms ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20, ?21)", params![ item.id, item.session_id, item.purpose, item.cause.kind_str(), cause_payload_str, readiness_str, blocked_reason, evidence_refs_str, item.consequence.impact.as_i64(), item.consequence.urgency_ms, item.consequence.aging_ms, item.provenance.reporter_id, item.provenance.origin_kind.as_str(), status_str, attempt_id, claimed_ms, completed_ms, disposition, failed_ms, error, item.created_ms, ], )?; Ok(rows > 0)}
fn parse_work_item_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<WorkItem> { let id: String = row.get(0)?; let session_id: String = row.get(1)?; let purpose: String = row.get(2)?; let cause_kind: String = row.get(3)?; let cause_payload_str: String = row.get(4)?; let readiness_str: String = row.get(5)?; let blocked_reason: Option<String> = row.get(6)?; let evidence_refs_str: String = row.get(7)?; let impact_num: i64 = row.get(8)?; let urgency_ms: Option<i64> = row.get(9)?; let aging_ms: i64 = row.get(10)?; let reporter_id: String = row.get(11)?; let origin_kind_str: String = row.get(12)?; let status_str: String = row.get(13)?; let attempt_id_str: Option<String> = row.get(14)?; let claimed_ms: Option<i64> = row.get(15)?; let completed_ms: Option<i64> = row.get(16)?; let disposition: Option<String> = row.get(17)?; let failed_ms: Option<i64> = row.get(18)?; let error: Option<String> = row.get(19)?; let created_ms: i64 = row.get(20)?;
let payload: serde_json::Value = serde_json::from_str(&cause_payload_str).unwrap_or(serde_json::Value::Null); let cause = match cause_kind.as_str() { "input_frame" => { let event_ids = payload .get("event_ids") .and_then(|v| v.as_array()) .map(|arr| arr.iter().filter_map(|x| x.as_i64()).collect()) .unwrap_or_default(); RunCause::InputFrame { event_ids } } "deadline_wake" => { let deadline_ms = payload .get("deadline_ms") .and_then(|v| v.as_i64()) .unwrap_or(0); let reason = payload .get("reason") .and_then(|v| v.as_str()) .map(|s| s.to_string()); RunCause::DeadlineWake { deadline_ms, reason, } } "maintenance_ready" => { let item_id = payload .get("item_id") .and_then(|v| v.as_str()) .unwrap_or(&id) .to_string(); let p = payload .get("purpose") .and_then(|v| v.as_str()) .unwrap_or(&purpose) .to_string(); RunCause::MaintenanceReady { item_id, purpose: p, } } "check_due" => { let check_id = payload .get("check_id") .and_then(|v| v.as_str()) .unwrap_or(&id) .to_string(); let due_ms = payload.get("due_ms").and_then(|v| v.as_i64()).unwrap_or(0); RunCause::CheckDue { check_id, due_ms } } "job_completed" => { let job_id = payload .get("job_id") .and_then(|v| v.as_str()) .unwrap_or(&id) .to_string(); let outcome = payload .get("outcome") .and_then(|v| v.as_str()) .unwrap_or("ok") .to_string(); RunCause::JobCompleted { job_id, outcome } } _ => RunCause::InputFrame { event_ids: vec![] }, };
let readiness = match readiness_str.as_str() { "blocked" => Readiness::Blocked { reason: blocked_reason.unwrap_or_default(), }, "paused" => Readiness::Paused { reason: blocked_reason.unwrap_or_default(), }, _ => Readiness::Ready, };
let consequence = ConsequenceMetadata { impact: Impact::from_i64(impact_num), urgency_ms, dependencies: Vec::new(), aging_ms, estimated_cost: None, readiness, };
let evidence_refs: Vec<String> = serde_json::from_str(&evidence_refs_str).unwrap_or_default(); let provenance = Provenance { reporter_id, origin_kind: OriginKind::from_str_lossy(&origin_kind_str), source_refs: evidence_refs, observed_ms: created_ms, };
let status = match status_str.as_str() { "claimed" => WorkStatus::Claimed { attempt_id: attempt_id_str .and_then(|s| OperationId::try_from(s).ok()) .unwrap_or_else(OperationId::fresh), claimed_ms: claimed_ms.unwrap_or(0), }, "completed" => WorkStatus::Completed { completed_ms: completed_ms.unwrap_or(0), disposition: disposition.unwrap_or_default(), }, "failed" => WorkStatus::Failed { failed_ms: failed_ms.unwrap_or(0), error: error.unwrap_or_default(), }, "deferred" => WorkStatus::Deferred { reason: error.unwrap_or_default(), }, _ => WorkStatus::Pending, };
Ok(WorkItem { id, session_id, purpose, cause, consequence, provenance, status, created_ms, })}
#[cfg(test)]mod ownership_tests { use super::*; use std::time::Instant;
/// Ownership is a share of an exclusive file lock, and `quiesce` reads a free /// permit as "no job is left holding it". Drive the REAL write path rather /// than a synthetic permit grab: a job that never started, or one whose lock /// share outlives its permit, both defeat the barrier this test claims. #[tokio::test] async fn quiesce_waits_for_a_real_in_flight_job_and_leaves_ownership_free() { const JOB: Duration = Duration::from_millis(300);
let dir = tempfile::tempdir().unwrap(); let path = dir.path().join("runtime.sqlite"); let store = Store::open(&path).unwrap();
let job = tokio::spawn({ let store = store.clone(); async move { store .with(|_| { std::thread::sleep(JOB); Ok(()) }) .await } });
// Wait for the job to actually take its share, then demand the refusal. let started = Instant::now(); while store.gate.available_permits() == BLOCKING_JOBS { assert!(started.elapsed() < JOB, "job never took a permit"); tokio::task::yield_now().await; } let held = store._owner.clone(); assert!( Store::open(&path).is_err(), "a live owner must refuse a second owner for the same database" );
let waiting = tokio::spawn({ let store = store.clone(); async move { store.quiesce().await } }); tokio::task::yield_now().await; assert!( !waiting.is_finished(), "quiesce must wait for the in-flight write instead of passing it" );
job.await.unwrap().unwrap(); waiting.await.unwrap().unwrap(); // quiesce returned only once nothing can still hold the lock, so releasing // the last handles hands the database over with no further wait. drop(held); drop(store); Store::open(&path).expect("a restarting owner takes the database after quiesce"); }}
#[cfg(test)]mod job_ownership_tests { use super::*;
/// A panicking job is the case a naive "drop the lock share at the end of the /// body" fix misses: unwinding runs the guard's drop, so the share is still /// released before the permit and the barrier stays honest. #[tokio::test] async fn a_panicking_job_still_releases_its_permit_and_lock_share() { let dir = tempfile::tempdir().unwrap(); let path = dir.path().join("runtime.sqlite"); let store = Store::open(&path).unwrap();
let panicked = store .with(|_: &mut Connection| -> Result<()> { panic!("job panicked mid-write") }) .await; assert!( panicked.is_err(), "a panicking job is an error, not a success" );
assert_eq!( store.gate.available_permits(), BLOCKING_JOBS, "a panicking job must not strand its permit" ); store .quiesce() .await .expect("the barrier completes once the panicking job unwound");
drop(store); Store::open(&path).expect("ownership transfers after a panicking job"); }}