//! Durable run causes, origin-neutral consequence scheduling, and work metadata. //! //! # Architecture & Contracts //! - Provenance (author/origin) is strictly evidence metadata, NEVER a priority or permission rank. //! - Priority and selection are pure projections over consequences: impact, urgency, dependencies, aging. //! - Tagged run causes (`InputFrame`, `DeadlineWake`, `MaintenanceReady`, `CheckDue`) share one runtime mechanism. use crate::protocol::OperationId; use serde::{Deserialize, Serialize}; use std::fmt; /// Identity for a single event delivery and idempotency key. #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(transparent)] pub struct OccurrenceId(pub String); impl OccurrenceId { #[must_use] pub fn from_event(event_id: i64) -> Self { Self(format!("event:{event_id}")) } #[must_use] pub fn as_str(&self) -> &str { &self.0 } } impl fmt::Display for OccurrenceId { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.write_str(&self.0) } } impl From for OccurrenceId { fn from(s: String) -> Self { Self(s) } } impl From<&str> for OccurrenceId { fn from(s: &str) -> Self { Self(s.to_string()) } } impl From for String { fn from(id: OccurrenceId) -> Self { id.0 } } impl AsRef for OccurrenceId { fn as_ref(&self) -> &str { &self.0 } } /// Identity for a recurring failure signature group. #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(transparent)] pub struct CaseId(pub String); impl CaseId { #[must_use] pub fn as_str(&self) -> &str { &self.0 } } impl fmt::Display for CaseId { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.write_str(&self.0) } } impl From for CaseId { fn from(s: String) -> Self { Self(s) } } impl From<&str> for CaseId { fn from(s: &str) -> Self { Self(s.to_string()) } } impl From for String { fn from(id: CaseId) -> Self { id.0 } } impl AsRef for CaseId { fn as_ref(&self) -> &str { &self.0 } } /// Identity for a runnable unit of work. #[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(transparent)] pub struct WorkItemId(pub String); impl WorkItemId { #[must_use] pub fn as_str(&self) -> &str { &self.0 } } impl fmt::Display for WorkItemId { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.write_str(&self.0) } } impl From for WorkItemId { fn from(s: String) -> Self { Self(s) } } impl From<&str> for WorkItemId { fn from(s: &str) -> Self { Self(s.to_string()) } } impl From for String { fn from(id: WorkItemId) -> Self { id.0 } } impl AsRef for WorkItemId { fn as_ref(&self) -> &str { &self.0 } } /// Alias for attempt identity. Each claim of a work item gets a fresh attempt. pub type AttemptId = OperationId; /// Result of attempting to claim a work item. #[derive(Clone, Debug, PartialEq, Eq)] pub enum ClaimResult { /// Successfully claimed the work item. Claimed { attempt_id: AttemptId }, /// Work item does not exist. NotFound, /// Work item is already claimed by another attempt. AlreadyClaimed { holder: AttemptId }, /// Work item is already completed. AlreadyCompleted, /// Work item is already failed. AlreadyFailed, /// Work item is not ready (blocked/paused). NotReady, } impl ClaimResult { #[must_use] pub fn is_claimed(&self) -> bool { matches!(self, Self::Claimed { .. }) } /// Whether this attempt holds the item: a claim it just took, or one it already had. /// /// A turn claims the work item it is running as well as the inputs it serves, and those can be /// the same row, so "already claimed by me" is success here and "already claimed by someone /// else" is not. #[must_use] pub fn held_by(&self, attempt: &AttemptId) -> bool { match self { Self::Claimed { attempt_id } => attempt_id == attempt, Self::AlreadyClaimed { holder } => holder == attempt, Self::NotFound | Self::AlreadyCompleted | Self::AlreadyFailed | Self::NotReady => false, } } } impl From for bool { fn from(r: ClaimResult) -> bool { r.is_claimed() } } impl std::ops::Not for ClaimResult { type Output = bool; fn not(self) -> bool { !self.is_claimed() } } impl std::ops::Not for &ClaimResult { type Output = bool; fn not(self) -> bool { !self.is_claimed() } } /// Result of attempting to retire a work item (complete or fail). #[derive(Clone, Debug, PartialEq, Eq)] pub enum RetireResult { /// Successfully retired the work item. Retired, /// Work item does not exist. NotFound, /// Attempt ID does not match the active claim holder. WrongOwner { expected: AttemptId }, /// Work item is not in claimed status. NotClaimed, /// Work item was already completed/failed. AlreadyTerminal, } impl RetireResult { #[must_use] pub fn is_retired(&self) -> bool { matches!(self, Self::Retired) } } impl From for bool { fn from(r: RetireResult) -> bool { r.is_retired() } } impl std::ops::Not for RetireResult { type Output = bool; fn not(self) -> bool { !self.is_retired() } } impl std::ops::Not for &RetireResult { type Output = bool; fn not(self) -> bool { !self.is_retired() } } /// Result of observing an occurrence. #[derive(Clone, Debug, PartialEq, Eq)] pub enum ObserveResult { /// This is a new occurrence; work item created or evidence attached. Applied { work_item_id: WorkItemId, is_new_case: bool, }, /// This occurrence was already processed; no-op. AlreadyObserved, /// The occurrence could not be processed. Error(String), } /// Provenance metadata. Origin identifies who reported or authored the work. /// Contracts mandate that `origin_kind` is NEVER used for authorization, approval gates, /// or priority ranking. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct Provenance { pub reporter_id: String, pub origin_kind: OriginKind, pub source_refs: Vec, pub observed_ms: i64, } impl Provenance { #[must_use] pub fn new(reporter_id: impl Into, origin_kind: OriginKind, observed_ms: i64) -> Self { Self { reporter_id: reporter_id.into(), origin_kind, source_refs: Vec::new(), observed_ms, } } #[must_use] pub fn human(reporter_id: impl Into, observed_ms: i64) -> Self { Self::new(reporter_id, OriginKind::Human, observed_ms) } #[must_use] pub fn agent(reporter_id: impl Into, observed_ms: i64) -> Self { Self::new(reporter_id, OriginKind::Agent, observed_ms) } #[must_use] pub fn system(reporter_id: impl Into, observed_ms: i64) -> Self { Self::new(reporter_id, OriginKind::System, observed_ms) } } #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum OriginKind { Human, Agent, System, External, } impl OriginKind { #[must_use] pub const fn as_str(self) -> &'static str { match self { Self::Human => "human", Self::Agent => "agent", Self::System => "system", Self::External => "external", } } pub fn from_str_lossy(s: &str) -> Self { match s { "human" => Self::Human, "agent" => Self::Agent, "external" => Self::External, _ => Self::System, } } } /// Operational impact level of a work item. #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum Impact { Low = 1, Normal = 2, High = 3, Critical = 4, } impl Impact { #[must_use] pub const fn base_weight(self) -> i64 { match self { Self::Low => 10, Self::Normal => 100, Self::High => 500, Self::Critical => 1000, } } #[must_use] pub fn from_i64(val: i64) -> Self { match val { 1 => Self::Low, 3 => Self::High, 4 => Self::Critical, _ => Self::Normal, } } #[must_use] pub const fn as_i64(self) -> i64 { self as i64 } } #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "status", rename_all = "snake_case")] pub enum Readiness { Ready, Blocked { reason: String }, Paused { reason: String }, } impl Readiness { #[must_use] pub fn is_ready(&self) -> bool { matches!(self, Self::Ready) } } /// Consequence metadata used for origin-neutral scheduling. #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] pub struct ConsequenceMetadata { pub impact: Impact, pub urgency_ms: Option, pub dependencies: Vec, pub aging_ms: i64, pub estimated_cost: Option, pub readiness: Readiness, } impl ConsequenceMetadata { #[must_use] pub fn new(impact: Impact) -> Self { Self { impact, urgency_ms: None, dependencies: Vec::new(), aging_ms: 0, estimated_cost: None, readiness: Readiness::Ready, } } } /// Tagged causes that can trigger a session turn or autonomous work. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum RunCause { InputFrame { event_ids: Vec, }, DeadlineWake { deadline_ms: i64, reason: Option, }, MaintenanceReady { item_id: String, purpose: String, }, CheckDue { check_id: String, due_ms: i64, }, JobCompleted { job_id: String, outcome: String, }, } impl RunCause { /// Returns true if this cause is autonomous and runs even without pending inputs. #[must_use] pub fn runs_without_pending_input(&self) -> bool { matches!( self, Self::DeadlineWake { .. } | Self::MaintenanceReady { .. } | Self::CheckDue { .. } | Self::JobCompleted { .. } ) } #[must_use] pub fn kind_str(&self) -> &'static str { match self { Self::InputFrame { .. } => "input_frame", Self::DeadlineWake { .. } => "deadline_wake", Self::MaintenanceReady { .. } => "maintenance_ready", Self::CheckDue { .. } => "check_due", Self::JobCompleted { .. } => "job_completed", } } } #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "status", rename_all = "snake_case")] pub enum WorkStatus { Pending, Claimed { attempt_id: OperationId, claimed_ms: i64, }, Completed { completed_ms: i64, disposition: String, }, Failed { failed_ms: i64, error: String, }, Deferred { reason: String, }, } /// A general unit of runnable work, encompassing operator commands, maintenance tasks, and checks. #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] pub struct WorkItem { pub id: String, pub session_id: String, pub purpose: String, pub cause: RunCause, pub consequence: ConsequenceMetadata, pub provenance: Provenance, pub status: WorkStatus, pub created_ms: i64, } /// Durable observation case grouped by failure signature and workspace. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct ObserverCaseRecord { pub id: String, pub session_id: String, pub operation: String, pub component_revision: String, pub signature: String, pub workspace: String, pub count: i64, pub first_seen_ms: i64, pub last_seen_ms: i64, pub evidence_refs: Vec, pub work_item_id: Option, } impl WorkItem { #[must_use] pub fn new( id: impl Into, purpose: impl Into, cause: RunCause, consequence: ConsequenceMetadata, provenance: Provenance, created_ms: i64, ) -> Self { Self { id: id.into(), session_id: "default".into(), purpose: purpose.into(), cause, consequence, provenance, status: WorkStatus::Pending, created_ms, } } #[must_use] pub fn with_session(mut self, session_id: impl Into) -> Self { self.session_id = session_id.into(); self } #[must_use] pub fn with_evidence_refs(mut self, refs: Vec) -> Self { self.provenance.source_refs = refs; self } #[must_use] pub fn with_readiness(mut self, readiness: Readiness) -> Self { self.consequence.readiness = readiness; self } #[must_use] pub fn with_impact(mut self, impact: Impact) -> Self { self.consequence.impact = impact; self } #[must_use] pub fn with_urgency(mut self, urgency_ms: i64) -> Self { self.consequence.urgency_ms = Some(urgency_ms); self } #[must_use] pub fn is_eligible(&self) -> bool { matches!(self.status, WorkStatus::Pending) && self.consequence.readiness.is_ready() } /// Claim this work item under a specific attempt identity. pub fn claim(&mut self, attempt_id: OperationId, now_ms: i64) -> Result<(), &'static str> { if !matches!( self.status, WorkStatus::Pending | WorkStatus::Deferred { .. } ) { return Err("work item is not in claimable status"); } if !self.consequence.readiness.is_ready() { return Err("work item is not ready"); } self.status = WorkStatus::Claimed { attempt_id, claimed_ms: now_ms, }; Ok(()) } /// Retire this attempt with a terminal completion. pub fn complete( &mut self, attempt_id: &OperationId, disposition: impl Into, now_ms: i64, ) -> Result<(), &'static str> { match &self.status { WorkStatus::Claimed { attempt_id: active, .. } if active == attempt_id => { self.status = WorkStatus::Completed { completed_ms: now_ms, disposition: disposition.into(), }; Ok(()) } WorkStatus::Claimed { .. } => Err("mismatched attempt ID cannot retire work"), _ => Err("work item is not in claimed status"), } } /// Retire this attempt with a terminal failure. pub fn fail( &mut self, attempt_id: &OperationId, error: impl Into, now_ms: i64, ) -> Result<(), &'static str> { match &self.status { WorkStatus::Claimed { attempt_id: active, .. } if active == attempt_id => { self.status = WorkStatus::Failed { failed_ms: now_ms, error: error.into(), }; Ok(()) } WorkStatus::Claimed { .. } => Err("mismatched attempt ID cannot fail work"), _ => Err("work item is not in claimed status"), } } } /// Compute priority score for origin-neutral scheduling. /// /// INVARIANT: `provenance.origin_kind` is NOT inspected. /// Only consequence facts (impact, urgency, aging, readiness) determine priority. #[must_use] pub fn project_priority_score(item: &WorkItem, now_ms: i64) -> i64 { if !item.consequence.readiness.is_ready() { return i64::MIN; } let base = item.consequence.impact.base_weight(); // Aging: 1 point per 1000ms waited (prevents starvation of lower-impact work under sustained load - G08) let age_ms = (now_ms - item.created_ms).max(0) + item.consequence.aging_ms; let age_boost = age_ms / 1_000; // Urgency: boost if approaching or past deadline let urgency_boost = match item.consequence.urgency_ms { Some(target) => { let delta = target - now_ms; if delta <= 0 { 1000 // Overdue } else if delta <= 5_000 { 500 // Approaching deadline } else { 100 } } None => 0, }; base + age_boost + urgency_boost } /// Select the next runnable work item based purely on consequence priority. /// /// Proves G05 (origin invariance), G06 (critical repair outranks feature), /// and G08 (aging prevents indefinite starvation). #[must_use] pub fn select_next_work_item(items: &[WorkItem], now_ms: i64) -> Option<&WorkItem> { items .iter() .filter(|item| { matches!(item.status, WorkStatus::Pending) && item.consequence.readiness.is_ready() }) .max_by_key(|item| project_priority_score(item, now_ms)) } #[cfg(test)] mod tests { use super::*; #[test] fn g05_origin_invariance_identical_decisions_for_human_and_agent() { let now = 10_000; let c = ConsequenceMetadata::new(Impact::High); let human_work = WorkItem::new( "task-1", "fix parser crash", RunCause::MaintenanceReady { item_id: "fix-1".into(), purpose: "repair".into(), }, c.clone(), Provenance::human("operator", now), now, ); let agent_work = WorkItem::new( "task-2", "fix parser crash", RunCause::MaintenanceReady { item_id: "fix-1".into(), purpose: "repair".into(), }, c, Provenance::agent("agent-eval", now), now, ); // Scoring must be bit-identical regardless of origin let score_h = project_priority_score(&human_work, now); let score_a = project_priority_score(&agent_work, now); assert_eq!( score_h, score_a, "G05: human and agent origin receive identical scheduling priority" ); } #[test] fn g06_reliability_repair_outranks_feature_regardless_of_origin() { let now = 10_000; // Human requested a nonurgent feature (Normal impact) let feature = WorkItem::new( "feat-1", "add dark mode", RunCause::InputFrame { event_ids: vec![1] }, ConsequenceMetadata::new(Impact::Normal), Provenance::human("operator", now), now, ); // Agent self-diagnosed a critical capability repair (Critical impact) let repair = WorkItem::new( "repair-1", "restore context delivery", RunCause::MaintenanceReady { item_id: "m-1".into(), purpose: "restore delivery".into(), }, ConsequenceMetadata::new(Impact::Critical), Provenance::agent("self-check", now), now, ); let items = [feature.clone(), repair.clone()]; let selected = select_next_work_item(&items, now); assert_eq!( selected.map(|w| w.id.as_str()), Some("repair-1"), "G06: critical repair outranks normal feature" ); // Swap origins: human requests the repair, agent suggests the feature let mut feature_agent = feature; feature_agent.provenance = Provenance::agent("agent-suggestion", now); let mut repair_human = repair; repair_human.provenance = Provenance::human("operator", now); let items_swapped = [feature_agent, repair_human]; let selected_swapped = select_next_work_item(&items_swapped, now); assert_eq!( selected_swapped.map(|w| w.id.as_str()), Some("repair-1"), "G06: swapped origins produce identical selection" ); } #[test] fn g08_aging_prevents_indefinite_maintenance_starvation() { let start = 1_000; // Low impact maintenance item created at start let maintenance = WorkItem::new( "maint-1", "clean temp scratchpad", RunCause::MaintenanceReady { item_id: "m-clean".into(), purpose: "cleanup".into(), }, ConsequenceMetadata::new(Impact::Low), Provenance::agent("janitor", start), start, ); // Continuous stream of Normal human work arrives later at now = 100_000 let now = 100_000; // 99 seconds later: aging adds 99 points let fresh_normal = WorkItem::new( "user-1", "hello", RunCause::InputFrame { event_ids: vec![2] }, ConsequenceMetadata::new(Impact::Normal), // base 100 Provenance::human("user", now), now, ); // At now = 100_000: // maintenance score: base 10 + age 99 = 109 // fresh_normal score: base 100 + age 0 = 100 // Maintenance has aged sufficiently to be selected! let items = [fresh_normal, maintenance]; let selected = select_next_work_item(&items, now); assert_eq!( selected.map(|w| w.id.as_str()), Some("maint-1"), "G08: aged maintenance is not starved" ); } #[test] fn single_attempt_ownership_and_identity_checked_retirement() { let now = 1_000; let mut item = WorkItem::new( "item-1", "test item", RunCause::CheckDue { check_id: "c1".into(), due_ms: now, }, ConsequenceMetadata::new(Impact::Normal), Provenance::system("monitor", now), now, ); let attempt_1 = OperationId::fresh(); let attempt_2 = OperationId::fresh(); assert!(item.claim(attempt_1.clone(), now).is_ok()); // Second claim while claimed must fail assert!(item.claim(attempt_2.clone(), now).is_err()); // Attempt 2 cannot retire Attempt 1's work assert!(item .complete(&attempt_2, "stale runner completion", now + 10) .is_err()); assert!(matches!(item.status, WorkStatus::Claimed { .. })); // Attempt 1 can retire its own work assert!(item .complete(&attempt_1, "completed successfully", now + 20) .is_ok()); assert!(matches!(item.status, WorkStatus::Completed { .. })); } #[test] fn durable_jobs_have_recorded_recovery_policy() { let job1 = DurableJob::new( "job-1", "session-a", "klbr_routines.compact", "rev-1", JobRecoveryPolicy::retry_safe(), 1_000, ); assert_eq!(job1.recovery_policy, JobRecoveryPolicy::RetrySafe); assert!(job1.recovery_policy.is_retry_safe()); let job2 = DurableJob::new( "job-2", "session-a", "klbr_routines.index", "rev-1", JobRecoveryPolicy::resume_from_checkpoint("cursor:150"), 1_000, ); assert_eq!(job2.recovery_policy.checkpoint(), Some("cursor:150")); let job3 = DurableJob::new( "job-3", "session-a", "klbr_routines.mutate", "rev-1", JobRecoveryPolicy::manual_unknown("uncommitted external API write"), 1_000, ); assert!(!job3.recovery_policy.is_retry_safe()); let serialized = serde_json::to_string(&job1).unwrap(); let deserialized: DurableJob = serde_json::from_str(&serialized).unwrap(); assert_eq!(deserialized.recovery_policy, JobRecoveryPolicy::RetrySafe); let serialized2 = serde_json::to_string(&job2).unwrap(); let deserialized2: DurableJob = serde_json::from_str(&serialized2).unwrap(); assert_eq!( deserialized2.recovery_policy.checkpoint(), Some("cursor:150") ); } } /// Explicit recovery policy for durable jobs / routines. #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[serde(tag = "policy", rename_all = "snake_case")] pub enum JobRecoveryPolicy { /// Safe to retry from the beginning upon restart or failure. RetrySafe, /// Resumes execution from the last recorded checkpoint. ResumeFromCheckpoint { checkpoint: String }, /// Requires manual triage; state after interruption is unknown. ManualUnknown { reason: String }, } impl JobRecoveryPolicy { #[must_use] pub fn retry_safe() -> Self { Self::RetrySafe } #[must_use] pub fn resume_from_checkpoint(checkpoint: impl Into) -> Self { Self::ResumeFromCheckpoint { checkpoint: checkpoint.into(), } } #[must_use] pub fn manual_unknown(reason: impl Into) -> Self { Self::ManualUnknown { reason: reason.into(), } } #[must_use] pub fn is_retry_safe(&self) -> bool { matches!(self, Self::RetrySafe) } #[must_use] pub fn checkpoint(&self) -> Option<&str> { match self { Self::ResumeFromCheckpoint { checkpoint } => Some(checkpoint.as_str()), _ => None, } } } impl Default for JobRecoveryPolicy { fn default() -> Self { Self::ManualUnknown { reason: "interrupted job requires explicit recovery triage".into(), } } } /// A durable job registered with the runtime, having an explicit recorded recovery policy. /// /// Unlike ephemeral background tasks which vanish on kernel reset or process end, /// a durable job has recorded identity, routine entrypoint, behavior revision, /// and recovery policy. #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] pub struct DurableJob { pub id: String, pub session_id: String, pub routine: String, pub revision: String, pub recovery_policy: JobRecoveryPolicy, pub status: WorkStatus, pub created_ms: i64, } impl DurableJob { #[must_use] pub fn new( id: impl Into, session_id: impl Into, routine: impl Into, revision: impl Into, recovery_policy: JobRecoveryPolicy, created_ms: i64, ) -> Self { Self { id: id.into(), session_id: session_id.into(), routine: routine.into(), revision: revision.into(), recovery_policy, status: WorkStatus::Pending, created_ms, } } #[must_use] pub fn to_run_cause(&self) -> RunCause { RunCause::JobCompleted { job_id: self.id.clone(), outcome: match &self.status { WorkStatus::Completed { disposition, .. } => disposition.clone(), WorkStatus::Failed { error, .. } => format!("failed: {error}"), WorkStatus::Deferred { reason } => format!("deferred: {reason}"), _ => "pending".into(), }, } } }