From d25cf1c7486d6c898040146e0f64a9d6e4299644 Mon Sep 17 00:00:00 2001 From: LiberteCode <291704231+LiberteCode@users.noreply.github.com> Date: Wed, 8 Jul 2026 19:39:54 +0800 Subject: [PATCH 1/2] feat: record inbox drain provenance --- crates/santi-api/src/ops.rs | 3 +- crates/santi-api/src/server.rs | 28 ++++-- crates/santi-api/src/webhook.rs | 31 +++++++ crates/santi-core/src/model.rs | 31 +++++++ crates/santi-core/src/service.rs | 40 ++++++-- crates/santi-core/src/service/im.rs | 11 ++- crates/santi-core/src/store.rs | 37 +++++++- crates/santi-core/src/store/db.rs | 124 ++++++++++++++++++++++--- crates/santi-core/src/store/rows.rs | 25 ++++- crates/santi-core/src/store/runtime.rs | 4 +- crates/santi-core/src/store/schema.rs | 9 +- crates/santi-core/tests/service.rs | 109 ++++++++++++++++++++++ crates/santi-core/tests/store.rs | 64 +++++++++++++ 13 files changed, 470 insertions(+), 46 deletions(-) diff --git a/crates/santi-api/src/ops.rs b/crates/santi-api/src/ops.rs index 013a7d6..953696e 100644 --- a/crates/santi-api/src/ops.rs +++ b/crates/santi-api/src/ops.rs @@ -128,10 +128,11 @@ fn inbox_seed_existing_strand( strand_id: &str, text: &str, ) -> Result { - let outcome = store.enqueue_inbox( + let outcome = store.enqueue_inbox_with_source( strand_id, santi_core::MessageKind::SantiSystem, santi_core::MessageContent::text(text), + Some(santi_core::InboxSource::new("offline_inbox_seed")), )?; Ok(match outcome { santi_core::IngestOutcome::Accepted { strand_id } => SeedReport { diff --git a/crates/santi-api/src/server.rs b/crates/santi-api/src/server.rs index 80ef28e..0339903 100644 --- a/crates/santi-api/src/server.rs +++ b/crates/santi-api/src/server.rs @@ -19,11 +19,12 @@ use futures_core::Stream; use santi_core::{ CompactExecRequest, CompactExecResponse, CompactQueryResponse, CreateSoulRequest, CreateStrandResponse, CreateWebhookRequest, ErrorResponse, ForkStrandResponse, HealthResponse, - ImInboxEntry, ImSendRequest, ImSendResponse, IngestOutcome, MaterialRequest, SantiService, - SantiServiceConfig, SantiStreamEvent, SantiStreamPayload, SendStrandAcceptedResponse, - SendStrandRequest, Soul, Strand, StrandDetail, StrandMaterial, StrandRuntimeSnapshot, - WebhookSubscription, prefixed_id, timestamp_now, + ImInboxEntry, ImSendRequest, ImSendResponse, InboxSource, IngestOutcome, MaterialRequest, + SantiService, SantiServiceConfig, SantiStreamEvent, SantiStreamPayload, + SendStrandAcceptedResponse, SendStrandRequest, Soul, Strand, StrandDetail, StrandMaterial, + StrandRuntimeSnapshot, WebhookSubscription, prefixed_id, timestamp_now, }; +use serde_json::json; use tower_http::{ cors::{Any, CorsLayer}, trace::TraceLayer, @@ -354,13 +355,28 @@ async fn ingest_webhook( let label = if subscription.strand_strategy == "single" { format!("{}:{}", subscription.adaptor, name) } else { - event.label + event.label.clone() }; + let source = InboxSource::new("webhook") + .with_ref(format!("{}:{name}", subscription.adaptor)) + .with_metadata(json!({ + "subscription": name, + "adaptor": subscription.adaptor, + "strand_strategy": subscription.strand_strategy, + "event_label": event.label, + "materialized_label": label, + "event": event.source_metadata, + })); // Rejection handling is the adaptor's own policy: a webhook silently drops // + logs (the sender has no way to retry a specific event) rather than // surfacing the inbox gate as an error. match service - .ingest_external_event(&subscription.soul_id, &label, event.santi_system_text) + .ingest_external_event_with_source( + &subscription.soul_id, + &label, + event.santi_system_text, + Some(source), + ) .map_err(ApiError::from_service)? { IngestOutcome::Accepted { .. } => {} diff --git a/crates/santi-api/src/webhook.rs b/crates/santi-api/src/webhook.rs index 9b228be..c70459c 100644 --- a/crates/santi-api/src/webhook.rs +++ b/crates/santi-api/src/webhook.rs @@ -28,6 +28,9 @@ pub(crate) struct NormalizedEvent { pub santi_system_text: String, /// The opaque external label that anchors the strand (per-thread identity). pub label: String, + /// Bounded adaptor-owned provenance for runtime diagnostics. This is NOT + /// appended to message content; it is carried only into drain evidence. + pub source_metadata: Option, /// Whether this event type is in scope for santi. Out-of-scope events (a /// GitHub `ping`, an unhandled action) verify fine but produce no turn. pub in_scope: bool, @@ -168,6 +171,7 @@ impl WebhookAdaptor for GithubAdaptor { return Ok(WebhookOutcome::Event(NormalizedEvent { santi_system_text: String::new(), label: format!("github:{webhook_name}:{event_type}"), + source_metadata: None, in_scope: false, self_authored: false, })); @@ -222,9 +226,22 @@ impl WebhookAdaptor for GithubAdaptor { // subscriptions never share a thread. let label = format!("github:{webhook_name}:issue:{repo}#{number}"); + let source_metadata = json!({ + "adaptor": "github", + "webhook_name": webhook_name, + "event_type": event_type, + "action": action, + "delivery": delivery, + "repo": repo, + "issue_number": number, + "url": url, + "label": label, + }); + Ok(WebhookOutcome::Event(NormalizedEvent { santi_system_text, label, + source_metadata: Some(source_metadata), in_scope, self_authored, })) @@ -404,6 +421,7 @@ fn feishu_normalize( return Ok(WebhookOutcome::Event(NormalizedEvent { santi_system_text: String::new(), label: format!("feishu:{webhook_name}:{event_type}"), + source_metadata: None, in_scope: false, self_authored: false, })); @@ -435,9 +453,22 @@ fn feishu_normalize( // One strand per feishu chat, scoped by subscription name. let label = format!("feishu:{webhook_name}:chat:{chat_id}"); + let source_metadata = json!({ + "adaptor": "feishu", + "webhook_name": webhook_name, + "event_type": event_type, + "event_id": event_id, + "chat_id": chat_id, + "chat_type": chat_type, + "message_id": message_id, + "sender_type": sender_type, + "label": label, + }); + Ok(WebhookOutcome::Event(NormalizedEvent { santi_system_text, label, + source_metadata: Some(source_metadata), in_scope, self_authored: false, })) diff --git a/crates/santi-core/src/model.rs b/crates/santi-core/src/model.rs index e266b48..b19e201 100644 --- a/crates/santi-core/src/model.rs +++ b/crates/santi-core/src/model.rs @@ -353,6 +353,36 @@ pub enum IngestOutcome { Rejected { reason: String }, } +/// Bounded provenance for an inbound item at the moment it is enqueued. This is +/// runtime evidence, not model-visible message content: provider assembly reads +/// `messages`, while this metadata is carried only into drain/audit diagnostics. +#[derive(Debug, Clone, Serialize, Deserialize, ToSchema)] +pub struct InboxSource { + pub source_type: String, + pub source_ref: Option, + pub metadata: Option, +} + +impl InboxSource { + pub fn new(source_type: impl Into) -> Self { + Self { + source_type: source_type.into(), + source_ref: None, + metadata: None, + } + } + + pub fn with_ref(mut self, source_ref: impl Into) -> Self { + self.source_ref = Some(source_ref.into()); + self + } + + pub fn with_metadata(mut self, metadata: Value) -> Self { + self.metadata = Some(metadata); + self + } +} + /// The external-label prefix that marks a strand as an IM conversation. The IM /// layer builds `im:` labels; the reply-routing correlation /// strips this back to the participant. Shared by the IM store, the API send @@ -483,6 +513,7 @@ pub enum SantiStreamPayload { pub struct StrandRuntimeSnapshot { pub strand: Strand, pub messages: Vec, + pub message_events: Vec, pub turns: Vec, pub thinking_spans: Vec, pub tool_calls: Vec, diff --git a/crates/santi-core/src/service.rs b/crates/santi-core/src/service.rs index 1e4b2d5..ecc0cd9 100644 --- a/crates/santi-core/src/service.rs +++ b/crates/santi-core/src/service.rs @@ -24,11 +24,11 @@ use crate::assembly::input::provider_input; use crate::service_prompt::provider_tools; use crate::{ CompactExecRequest, CompactExecResponse, CompactQueryResponse, CreateSoulRequest, - CreateStrandResponse, CreateWebhookRequest, IngestOutcome, MaterialKind, MessageContent, - MessageKind, SantiStore, SantiStreamEvent, SantiStreamPayload, SendStrandAcceptedResponse, - SendStrandRequest, Soul, Strand, StrandDetail, StrandMaterial, StrandMessage, - StrandRuntimeSnapshot, StrandSelector, ThinkingCompletionReason, ThinkingSpan, Turn, - TurnActivityState, WebhookSubscription, prefixed_id, timestamp_now, + CreateStrandResponse, CreateWebhookRequest, InboxSource, IngestOutcome, MaterialKind, + MessageContent, MessageKind, SantiStore, SantiStreamEvent, SantiStreamPayload, + SendStrandAcceptedResponse, SendStrandRequest, Soul, Strand, StrandDetail, StrandMaterial, + StrandMessage, StrandRuntimeSnapshot, StrandSelector, ThinkingCompletionReason, ThinkingSpan, + Turn, TurnActivityState, WebhookSubscription, prefixed_id, timestamp_now, }; use failure::ProviderTurnFailure; use runtime_notice::{ProviderInputObservation, RuntimeNoticeBus}; @@ -232,9 +232,20 @@ impl SantiService { content: MessageContent, kind: MessageKind, trigger_type: &str, + ) -> Result { + self.ingest_with_source(selector, content, kind, trigger_type, None) + } + + pub fn ingest_with_source( + &self, + selector: StrandSelector, + content: MessageContent, + kind: MessageKind, + trigger_type: &str, + source: Option, ) -> Result { let strand = self.store.resolve_strand_selector(&selector)?; - let (outcome, _driven) = self.ingest_into(&strand, content, kind, trigger_type)?; + let (outcome, _driven) = self.ingest_into(&strand, content, kind, trigger_type, source)?; Ok(outcome) } @@ -251,8 +262,11 @@ impl SantiService { content: MessageContent, kind: MessageKind, trigger_type: &str, + source: Option, ) -> Result<(IngestOutcome, DrivenTurn), String> { - let outcome = self.store.enqueue_inbox(&strand.id, kind, content)?; + let outcome = self + .store + .enqueue_inbox_with_source(&strand.id, kind, content, source)?; let driven = match outcome { IngestOutcome::Accepted { .. } => self.poke(&strand.id, trigger_type), IngestOutcome::Rejected { .. } => None, @@ -270,6 +284,16 @@ impl SantiService { soul_id: &str, label: &str, system_text: String, + ) -> Result { + self.ingest_external_event_with_source(soul_id, label, system_text, None) + } + + pub fn ingest_external_event_with_source( + &self, + soul_id: &str, + label: &str, + system_text: String, + source: Option, ) -> Result { let strand = self .store @@ -282,6 +306,7 @@ impl SantiService { MessageContent::text(system_text), MessageKind::SantiSystem, "system", + source, )?; Ok(outcome) } @@ -367,6 +392,7 @@ impl SantiService { }, MessageKind::Text, "strand_send", + Some(InboxSource::new("strand_send").with_ref(strand.id.clone())), )?; if let IngestOutcome::Rejected { reason } = outcome { return Err(reason); diff --git a/crates/santi-core/src/service/im.rs b/crates/santi-core/src/service/im.rs index 7ed7a30..07133c5 100644 --- a/crates/santi-core/src/service/im.rs +++ b/crates/santi-core/src/service/im.rs @@ -1,11 +1,13 @@ //! IM layer service methods — the thin seam between the plain IM and the runtime. //! The IM is conceptually ORTHOGONAL to the runtime (PHASE-08 CONVERGED MODEL v4): //! inbound reuses the source-less runtime primitive `ingest`, addressing the soul -//! by an `im:` conversation label; the participant address lives only -//! in the IM envelope (the label + the IM store), never in the runtime. +//! by an `im:` conversation label. Reply-routing authority lives in +//! the IM envelope (the label + the IM store); the runtime receives only bounded +//! diagnostic provenance for incident review, not a reply capability. use crate::{ - IM_LABEL_PREFIX, ImInboxEntry, IngestOutcome, MessageContent, MessageKind, StrandSelector, + IM_LABEL_PREFIX, ImInboxEntry, InboxSource, IngestOutcome, MessageContent, MessageKind, + StrandSelector, }; use super::SantiService; @@ -23,7 +25,7 @@ impl SantiService { ) -> Result { self.store.ensure_im_participant(participant_id, "human")?; let label = format!("{IM_LABEL_PREFIX}{participant_id}"); - self.ingest( + self.ingest_with_source( StrandSelector::ByLabel { soul_id: soul_id.to_string(), label, @@ -31,6 +33,7 @@ impl SantiService { MessageContent::text(content.to_string()), MessageKind::Text, "strand_send", + Some(InboxSource::new("im").with_ref(participant_id.to_string())), ) } diff --git a/crates/santi-core/src/store.rs b/crates/santi-core/src/store.rs index 1cb6323..89f8b51 100644 --- a/crates/santi-core/src/store.rs +++ b/crates/santi-core/src/store.rs @@ -6,8 +6,9 @@ use std::{ use rusqlite::{Connection, params}; use crate::{ - ActorType, IngestOutcome, MessageContent, MessageIntake, MessageKind, MessageState, Strand, - StrandMessage, StrandSelector, StrandTargetType, Turn, prefixed_id, timestamp_now, + ActorType, InboxSource, IngestOutcome, MessageContent, MessageIntake, MessageKind, + MessageState, Strand, StrandMessage, StrandSelector, StrandTargetType, Turn, prefixed_id, + timestamp_now, }; mod assembly; @@ -27,7 +28,7 @@ use schema::SCHEMA; /// is wiped + rebuilt (beta: no back-compat migrations yet — see PHASE-07 crux #5). /// Public so ops paths (`santi doctor`) can compare a DB's `user_version` to it /// WITHOUT opening the store (which would migrate/wipe). -pub const SCHEMA_VERSION: u32 = 21; +pub const SCHEMA_VERSION: u32 = 22; /// The default soul's id. Public so offline ops (doctor/seed) can address it /// without a running service. pub const DEFAULT_SOUL_ID: &str = "soul_default"; @@ -242,6 +243,7 @@ impl SantiStore { }; Ok(Some(crate::StrandRuntimeSnapshot { messages: strand_messages(&conn, strand_id)?, + message_events: message_events_for_strand(&conn, strand_id)?, turns: turns_for_strand(&conn, &strand.id)?, thinking_spans: soul_thinking_spans(&conn, &strand.id)?, tool_calls: soul_tool_calls(&conn, &strand.id)?, @@ -392,6 +394,16 @@ impl SantiStore { strand_id: &str, message_kind: MessageKind, content: MessageContent, + ) -> Result { + self.enqueue_inbox_with_source(strand_id, message_kind, content, None) + } + + pub fn enqueue_inbox_with_source( + &self, + strand_id: &str, + message_kind: MessageKind, + content: MessageContent, + source: Option, ) -> Result { let conn = self.conn.lock().unwrap(); let pending: i64 = conn @@ -411,16 +423,31 @@ impl SantiStore { let inbox_id = prefixed_id("inbox"); let now = timestamp_now(); let content_json = serde_json::to_string(&content).map_err(|error| error.to_string())?; + let source_type = source.as_ref().map(|source| source.source_type.as_str()); + let source_ref = source + .as_ref() + .and_then(|source| source.source_ref.as_deref()); + let source_metadata = source + .as_ref() + .and_then(|source| source.metadata.as_ref()) + .map(serde_json::to_string) + .transpose() + .map_err(|error| error.to_string())?; conn.execute( r#" - INSERT INTO strand_inbox (id, strand_id, message_kind, content, created_at) - VALUES (?1, ?2, ?3, ?4, ?5) + INSERT INTO strand_inbox ( + id, strand_id, message_kind, content, source_type, source_ref, source_metadata, created_at + ) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8) "#, params![ inbox_id, strand_id, rows::message_kind_db(&message_kind), content_json, + source_type, + source_ref, + source_metadata, now ], ) diff --git a/crates/santi-core/src/store/db.rs b/crates/santi-core/src/store/db.rs index 4988b88..2d029da 100644 --- a/crates/santi-core/src/store/db.rs +++ b/crates/santi-core/src/store/db.rs @@ -1,11 +1,12 @@ mod timeline; use rusqlite::{Connection, OptionalExtension, params}; +use serde_json::{Value, json}; use crate::{ - ActorType, Compact, MessageKind, Soul, Strand, StrandEffect, StrandEntry, StrandMessage, - StrandTargetType, ThinkingSpan, ToolCall, ToolResult, Turn, WebhookSubscription, prefixed_id, - timestamp_now, + ActorType, Compact, MessageEvent, MessageKind, Soul, Strand, StrandEffect, StrandEntry, + StrandMessage, StrandTargetType, ThinkingSpan, ToolCall, ToolResult, Turn, WebhookSubscription, + prefixed_id, timestamp_now, }; use super::rows::*; @@ -63,11 +64,13 @@ pub(super) fn append_entry_in_tx( pub(super) fn drain_inbox_in_tx( conn: &Connection, strand_id: &str, + committing_turn_id: &str, ) -> Result, String> { let mut stmt = conn .prepare( r#" - SELECT id, message_kind, content FROM strand_inbox + SELECT id, message_kind, content, source_type, source_ref, source_metadata, created_at + FROM strand_inbox WHERE strand_id = ?1 ORDER BY rowid ASC "#, @@ -75,11 +78,15 @@ pub(super) fn drain_inbox_in_tx( .map_err(|error| error.to_string())?; let pending = stmt .query_map(params![strand_id], |row| { - Ok(( - row.get::<_, String>(0)?, - row.get::<_, String>(1)?, - row.get::<_, String>(2)?, - )) + Ok(PendingInboxEntry { + id: row.get(0)?, + message_kind: row.get(1)?, + content: row.get(2)?, + source_type: row.get(3)?, + source_ref: row.get(4)?, + source_metadata: row.get(5)?, + enqueued_at: row.get(6)?, + }) }) .map_err(|error| error.to_string())? .collect::, _>>() @@ -88,7 +95,7 @@ pub(super) fn drain_inbox_in_tx( let now = timestamp_now(); let mut drained = Vec::with_capacity(pending.len()); - for (inbox_id, message_kind_db, content_json) in pending { + for pending_entry in pending { let message_id = prefixed_id("msg"); conn.execute( r#" @@ -101,15 +108,26 @@ pub(super) fn drain_inbox_in_tx( params![ message_id, super::SANTI_SYSTEM_ACTOR_ID, - message_kind_db, - content_json, + pending_entry.message_kind.as_str(), + pending_entry.content.as_str(), now ], ) .map_err(|error| error.to_string())?; - append_entry_in_tx(conn, strand_id, StrandTargetType::Message, &message_id)?; - conn.execute("DELETE FROM strand_inbox WHERE id = ?1", params![inbox_id]) - .map_err(|error| error.to_string())?; + let relation = append_entry_in_tx(conn, strand_id, StrandTargetType::Message, &message_id)?; + insert_inbox_drain_event_in_tx( + conn, + &pending_entry, + &message_id, + relation.strand_seq, + committing_turn_id, + &now, + )?; + conn.execute( + "DELETE FROM strand_inbox WHERE id = ?1", + params![pending_entry.id], + ) + .map_err(|error| error.to_string())?; drained.push( message_by_id(conn, &message_id)? .ok_or_else(|| "drained message missing".to_string())?, @@ -118,6 +136,82 @@ pub(super) fn drain_inbox_in_tx( Ok(drained) } +struct PendingInboxEntry { + id: String, + message_kind: String, + content: String, + source_type: Option, + source_ref: Option, + source_metadata: Option, + enqueued_at: String, +} + +fn insert_inbox_drain_event_in_tx( + conn: &Connection, + pending: &PendingInboxEntry, + message_id: &str, + strand_seq: i64, + turn_id: &str, + drained_at: &str, +) -> Result<(), String> { + let source_metadata = pending.source_metadata.as_deref().map(|raw| { + serde_json::from_str::(raw).unwrap_or_else(|_| json!({ "invalid_json": true })) + }); + let payload = json!({ + "kind": "inbox_drain", + "inbox_id": pending.id.as_str(), + "enqueued_at": pending.enqueued_at.as_str(), + "drained_at": drained_at, + "committing_turn_id": turn_id, + "message_id": message_id, + "strand_seq": strand_seq, + "source": { + "type": pending.source_type.as_deref(), + "ref": pending.source_ref.as_deref(), + "metadata": source_metadata, + } + }); + conn.execute( + r#" + INSERT INTO message_events ( + id, message_id, action, actor_type, actor_id, base_version, payload, created_at + ) + VALUES (?1, ?2, 'insert', 'system', ?3, 1, ?4, ?5) + "#, + params![ + prefixed_id("mev"), + message_id, + super::SANTI_SYSTEM_ACTOR_ID, + payload.to_string(), + drained_at + ], + ) + .map_err(|error| error.to_string())?; + Ok(()) +} + +pub(super) fn message_events_for_strand( + conn: &Connection, + strand_id: &str, +) -> Result, String> { + let mut stmt = conn + .prepare( + r#" + SELECT e.id, e.message_id, e.action, e.actor_type, e.actor_id, + e.base_version, e.payload, e.created_at + FROM message_events e + JOIN r_strand_entries r ON r.target_type = 'message' AND r.target_id = e.message_id + WHERE r.strand_id = ?1 + ORDER BY r.strand_seq ASC, e.created_at ASC, e.id ASC + "#, + ) + .map_err(|error| error.to_string())?; + let rows = stmt + .query_map(params![strand_id], map_message_event_row) + .map_err(|error| error.to_string())?; + collect_rows(rows) +} + pub(super) fn soul_by_id(conn: &Connection, soul_id: &str) -> Result, String> { conn.query_row( r#" diff --git a/crates/santi-core/src/store/rows.rs b/crates/santi-core/src/store/rows.rs index 615e584..b702667 100644 --- a/crates/santi-core/src/store/rows.rs +++ b/crates/santi-core/src/store/rows.rs @@ -2,10 +2,10 @@ use rusqlite::Row; use serde_json::Value; use crate::{ - ActorType, Compact, Message, MessageContent, MessageKind, MessageState, Soul, Strand, - StrandEffect, StrandMessage, StrandMessageRef, StrandTargetType, ThinkingCompletionReason, - ThinkingSpan, ThinkingSpanState, ToolCall, ToolResult, Turn, TurnStatus, TurnTriggerType, - WebhookSubscription, + ActorType, Compact, Message, MessageContent, MessageEvent, MessageKind, MessageState, Soul, + Strand, StrandEffect, StrandMessage, StrandMessageRef, StrandTargetType, + ThinkingCompletionReason, ThinkingSpan, ThinkingSpanState, ToolCall, ToolResult, Turn, + TurnStatus, TurnTriggerType, WebhookSubscription, }; pub(super) fn map_soul_row(row: &Row<'_>) -> rusqlite::Result { @@ -64,6 +64,23 @@ pub(super) fn map_message_row(row: &Row<'_>) -> rusqlite::Result { }) } +pub(super) fn map_message_event_row(row: &Row<'_>) -> rusqlite::Result { + let payload_json: String = row.get(6)?; + let payload = serde_json::from_str::(&payload_json).map_err(|error| { + rusqlite::Error::FromSqlConversionFailure(6, rusqlite::types::Type::Text, Box::new(error)) + })?; + Ok(MessageEvent { + id: row.get(0)?, + message_id: row.get(1)?, + action: row.get(2)?, + actor_type: actor_type_from_db(row.get::<_, String>(3)?.as_str()), + actor_id: row.get(4)?, + base_version: row.get(5)?, + payload, + created_at: row.get(7)?, + }) +} + pub(super) fn map_strand_message_row(row: &Row<'_>) -> rusqlite::Result { let content_json: String = row.get(8)?; let content = serde_json::from_str::(&content_json).map_err(|error| { diff --git a/crates/santi-core/src/store/runtime.rs b/crates/santi-core/src/store/runtime.rs index 77784af..5b3c231 100644 --- a/crates/santi-core/src/store/runtime.rs +++ b/crates/santi-core/src/store/runtime.rs @@ -270,11 +270,11 @@ impl SantiStore { if running.is_some() { return Ok(None); } - let drained_messages = drain_inbox_in_tx(&tx, strand_id)?; + let turn_id = prefixed_id("turn"); + let drained_messages = drain_inbox_in_tx(&tx, strand_id, &turn_id)?; if drained_messages.is_empty() { return Ok(None); } - let turn_id = prefixed_id("turn"); let now = timestamp_now(); tx.execute( r#" diff --git a/crates/santi-core/src/store/schema.rs b/crates/santi-core/src/store/schema.rs index b599805..03ccc57 100644 --- a/crates/santi-core/src/store/schema.rs +++ b/crates/santi-core/src/store/schema.rs @@ -165,6 +165,9 @@ CREATE TABLE IF NOT EXISTS strand_inbox ( strand_id TEXT NOT NULL, message_kind TEXT NOT NULL CHECK (message_kind IN ('text', 'santi_system')), content TEXT NOT NULL, + source_type TEXT, + source_ref TEXT, + source_metadata TEXT, created_at TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_strand_inbox_strand_created_at ON strand_inbox (strand_id, created_at); @@ -199,8 +202,10 @@ CREATE INDEX IF NOT EXISTS idx_r_strand_entries_seq ON r_strand_entries (strand_ -- ORTHOGONAL to the runtime (souls/strands/turns). These tables are the IM's own -- store — the runtime never reads them. A participant is a persistent messaging -- endpoint (a human/CLI peer with a passive inbox; a soul participant's "inbox" is --- its strand and is NOT stored here). `source` addressing lives entirely here, in --- the IM's envelope — never in the runtime primitive or `strand_inbox`. +-- its strand and is NOT stored here). Reply-routing authority lives entirely +-- here, in the IM's envelope. The runtime inbox may carry bounded diagnostic +-- source provenance, but that is not a reply capability or provider-visible +-- message content. CREATE TABLE IF NOT EXISTS im_participants ( id TEXT PRIMARY KEY, kind TEXT NOT NULL CHECK (kind IN ('human', 'soul')), diff --git a/crates/santi-core/tests/service.rs b/crates/santi-core/tests/service.rs index 8cc2b49..3b8ebe4 100644 --- a/crates/santi-core/tests/service.rs +++ b/crates/santi-core/tests/service.rs @@ -825,6 +825,115 @@ async fn request_arriving_during_running_turn_drives_one_follow_on_turn() { assert_eq!(count_messages(&runtime, "provider response 2"), 1); } +#[tokio::test] +async fn coalesced_request_drain_provenance_preserves_original_enqueue_time() { + let temp = tempfile::tempdir().expect("temp dir"); + let provider = Arc::new(GatedFirstProvider::new()); + let service = SantiService::open( + SantiServiceConfig { + database_path: temp.path().join("santi.sqlite").display().to_string(), + runtime_root: temp.path().join("runtime").display().to_string(), + execution_root: temp.path().join("execution").display().to_string(), + bind_addr: Some("127.0.0.1:0".to_string()), + }, + provider.clone(), + ) + .expect("open service"); + + let strand = service.create_strand().expect("create strand").strand; + let first = service + .send_strand( + &strand.id, + SendStrandRequest { + content: vec![MessagePart::Text { + text: "first request".to_string(), + }], + }, + ) + .await + .expect("send first request"); + let first_message_id = first + .user_message + .as_ref() + .expect("first send drains immediately") + .message + .id + .clone(); + + provider.wait_for_first_request().await; + + let second = service + .send_strand( + &strand.id, + SendStrandRequest { + content: vec![MessagePart::Text { + text: "second request".to_string(), + }], + }, + ) + .await + .expect("send second request while first runs"); + assert!(second.user_message.is_none()); + + provider.release_first_request(); + let runtime = wait_for_completed_turn_count(&service, &strand.id, 2).await; + + let second_message = runtime + .messages + .iter() + .find(|message| message.content_text == "second request") + .expect("second message drained after first turn"); + let second_event = runtime + .message_events + .iter() + .find(|event| event.message_id == second_message.message.id) + .expect("second message drain event"); + + assert_eq!(second_event.payload["kind"], "inbox_drain"); + assert_eq!( + second_event.payload["message_id"], + second_message.message.id + ); + assert_eq!( + second_event.payload["drained_at"], + second_message.message.created_at + ); + assert_eq!(second_event.created_at, second_message.message.created_at); + assert_eq!( + second_event.payload["source"]["type"], "strand_send", + "direct sends should carry caller/source shape" + ); + assert_eq!(second_event.payload["source"]["ref"], strand.id); + + let follow_on_turn = runtime + .turns + .iter() + .find(|turn| turn.id != first.turn.id) + .expect("follow-on turn"); + assert_eq!( + second_event.payload["committing_turn_id"], follow_on_turn.id, + "the drain event should name the turn that committed the pending request" + ); + + let enqueued_at = second_event.payload["enqueued_at"] + .as_str() + .expect("enqueued_at string"); + assert!( + enqueued_at <= second_message.message.created_at.as_str(), + "enqueue time should not be later than drain/message time" + ); + + let first_event = runtime + .message_events + .iter() + .find(|event| event.message_id == first_message_id) + .expect("first message drain event"); + assert_ne!( + first_event.payload["inbox_id"], second_event.payload["inbox_id"], + "each inbound request should keep its own inbox id provenance" + ); +} + #[tokio::test] async fn completed_turn_emits_turn_completed_event() { let temp = tempfile::tempdir().expect("temp dir"); diff --git a/crates/santi-core/tests/store.rs b/crates/santi-core/tests/store.rs index cdf2c4a..db33f4f 100644 --- a/crates/santi-core/tests/store.rs +++ b/crates/santi-core/tests/store.rs @@ -394,6 +394,70 @@ fn drain_commits_all_pending_inbox_entries_to_one_turn() { ); } +#[test] +fn inbox_drain_records_enqueue_and_commit_provenance() { + let temp = tempfile::tempdir().expect("temp dir"); + let db = temp.path().join("santi.sqlite"); + let store = SantiStore::open(&db).expect("open store"); + let strand = store.create_strand().expect("create strand"); + + store + .enqueue_inbox_with_source( + &strand.id, + MessageKind::Text, + MessageContent::text("needs provenance"), + Some( + santi_core::InboxSource::new("test") + .with_ref("caller-1") + .with_metadata(serde_json::json!({ "adaptor": "fake" })), + ), + ) + .expect("enqueue with source"); + + let conn = Connection::open(&db).expect("open sqlite"); + let (inbox_id, enqueued_at): (String, String) = conn + .query_row( + "SELECT id, created_at FROM strand_inbox WHERE strand_id = ?1", + [&strand.id], + |row| Ok((row.get(0)?, row.get(1)?)), + ) + .expect("inbox row"); + drop(conn); + + let started = store + .try_start_turn(&strand.id, "strand_send", None) + .expect("try") + .expect("turn started"); + assert_eq!(started.drained_messages.len(), 1); + let drained = &started.drained_messages[0]; + + let runtime = store + .runtime_snapshot(&strand.id) + .expect("runtime snapshot") + .expect("strand runtime"); + assert_eq!(runtime.message_events.len(), 1); + let event = &runtime.message_events[0]; + assert_eq!(event.action, "insert"); + assert_eq!(event.message_id, drained.message.id); + assert_eq!(event.created_at, drained.message.created_at); + + let payload = &event.payload; + assert_eq!(payload["kind"], "inbox_drain"); + assert_eq!(payload["inbox_id"], inbox_id); + assert_eq!(payload["enqueued_at"], enqueued_at); + assert_eq!(payload["drained_at"], drained.message.created_at); + assert_eq!(payload["committing_turn_id"], started.turn.id); + assert_eq!(payload["message_id"], drained.message.id); + assert_eq!(payload["strand_seq"], drained.relation.strand_seq); + assert_eq!(payload["source"]["type"], "test"); + assert_eq!(payload["source"]["ref"], "caller-1"); + assert_eq!(payload["source"]["metadata"]["adaptor"], "fake"); + + let input = store.assembly_input(&strand.id).expect("assembly input"); + assert_eq!(input.len(), 1); + assert_text(&input[0], "user", "needs provenance"); +} + #[test] fn inbox_gate_rejects_past_threshold() { let temp = tempfile::tempdir().expect("temp dir"); From 7e40c0532b80a19158cee8376f644893f9b26b98 Mon Sep 17 00:00:00 2001 From: LiberteCode <291704231+LiberteCode@users.noreply.github.com> Date: Wed, 8 Jul 2026 19:58:45 +0800 Subject: [PATCH 2/2] fix: migrate schema 21 to 22 in place --- crates/santi-core/src/store.rs | 53 ++++++++++++++++--- crates/santi-core/tests/store.rs | 90 ++++++++++++++++++++++++++++++++ 2 files changed, 136 insertions(+), 7 deletions(-) diff --git a/crates/santi-core/src/store.rs b/crates/santi-core/src/store.rs index 89f8b51..200fcff 100644 --- a/crates/santi-core/src/store.rs +++ b/crates/santi-core/src/store.rs @@ -24,8 +24,9 @@ use db::*; use rows::{actor_type_db, collect_rows, map_webhook_row, message_state_db}; use schema::SCHEMA; -/// The schema version this binary expects. On open, a DB at any other version -/// is wiped + rebuilt (beta: no back-compat migrations yet — see PHASE-07 crux #5). +/// The schema version this binary expects. On open, recognized runtime-schema +/// migrations run in place; an unrecognized mismatch still falls back to the +/// beta wipe + rebuild policy (see PHASE-07 crux #5 / PHASE-09 tier work). /// Public so ops paths (`santi doctor`) can compare a DB's `user_version` to it /// WITHOUT opening the store (which would migrate/wipe). pub const SCHEMA_VERSION: u32 = 22; @@ -92,6 +93,39 @@ pub fn soul_memory_file(runtime_root: impl AsRef, soul_id: &str) -> std::p .join(crate::workspace_uri::MEMORY_FILE) } +fn migrate_schema_21_to_22(conn: &Connection) -> Result<(), String> { + // v22 adds enqueue provenance to the inbox. These fields are nullable by + // design so existing v21 rows keep their original content/created_at + // semantics while new ingress can attach bounded source metadata. + add_column_if_missing(conn, "strand_inbox", "source_type", "TEXT")?; + add_column_if_missing(conn, "strand_inbox", "source_ref", "TEXT")?; + add_column_if_missing(conn, "strand_inbox", "source_metadata", "TEXT")?; + Ok(()) +} + +fn add_column_if_missing( + conn: &Connection, + table: &str, + column: &str, + definition: &str, +) -> Result<(), String> { + let count: i64 = conn + .query_row( + &format!("SELECT COUNT(*) FROM pragma_table_info('{table}') WHERE name = ?1"), + params![column], + |row| row.get(0), + ) + .map_err(|error| error.to_string())?; + if count == 0 { + conn.execute( + &format!("ALTER TABLE {table} ADD COLUMN {column} {definition}"), + [], + ) + .map_err(|error| error.to_string())?; + } + Ok(()) +} + impl SantiStore { pub fn open(path: impl AsRef) -> Result { if let Some(parent) = path.as_ref().parent() { @@ -115,11 +149,16 @@ impl SantiStore { let version = conn .query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0)) .map_err(|error| error.to_string())?; - if version != SCHEMA_VERSION { - // Tier boundary (PHASE-09 decision #2, thin form): every table below is - // EPHEMERAL-tier (rooms / timeline / provider replay material) — a schema - // bump may drop-recreate them. A future DURABLE tier (attention / effects - // / checkpoints ledger) must instead be MIGRATED, never listed here. + if version == 21 && SCHEMA_VERSION == 22 { + // v21 -> v22 is additive: PR #47 only adds bounded inbox-source + // provenance columns. Migrate it in place so live ingress topology + // (notably `webhooks` / the secretary subscription) cannot be + // silently severed by this schema bump. + migrate_schema_21_to_22(&conn)?; + } else if version != SCHEMA_VERSION { + // Fallback beta policy for unrecognized schema jumps: drop the + // current runtime workspace and rebuild it. This must keep shrinking + // as more runtime topology/evidence graduates into migrated tiers. conn.execute_batch( r#" DROP TABLE IF EXISTS provider_replay_material; diff --git a/crates/santi-core/tests/store.rs b/crates/santi-core/tests/store.rs index db33f4f..1459b10 100644 --- a/crates/santi-core/tests/store.rs +++ b/crates/santi-core/tests/store.rs @@ -615,6 +615,96 @@ fn read_schema_version_none_when_db_absent() { ); } +#[test] +fn schema_21_to_22_migrates_in_place_and_preserves_webhooks() { + let temp = tempfile::tempdir().expect("temp dir"); + let db = temp.path().join("santi.sqlite"); + { + let conn = Connection::open(&db).expect("open sqlite"); + conn.execute_batch( + r#" + CREATE TABLE souls ( + id TEXT PRIMARY KEY, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL + ); + CREATE TABLE webhooks ( + name TEXT PRIMARY KEY, + adaptor TEXT NOT NULL, + soul_id TEXT NOT NULL, + strand_strategy TEXT NOT NULL, + secret_env TEXT NOT NULL, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL + ); + CREATE TABLE strand_inbox ( + id TEXT PRIMARY KEY, + strand_id TEXT NOT NULL, + message_kind TEXT NOT NULL CHECK (message_kind IN ('text', 'santi_system')), + content TEXT NOT NULL, + created_at TEXT NOT NULL + ); + INSERT INTO souls (id, created_at, updated_at) + VALUES ('soul_default', '2026-07-08T00:00:00Z', '2026-07-08T00:00:00Z'); + INSERT INTO webhooks ( + name, adaptor, soul_id, strand_strategy, secret_env, created_at, updated_at + ) VALUES ( + 'secretary', 'github', 'soul_default', 'per_thread', + 'SANTI_WEBHOOK_GITHUB_SECRET', '2026-07-08T00:00:01Z', '2026-07-08T00:00:01Z' + ); + INSERT INTO strand_inbox (id, strand_id, message_kind, content, created_at) + VALUES ('inbox_old', 'ss_existing', 'text', '{"parts":[]}', '2026-07-08T00:00:02Z'); + PRAGMA user_version = 21; + "#, + ) + .expect("seed v21 db"); + } + + let store = SantiStore::open(&db).expect("open migrates v21 to v22"); + assert_eq!( + santi_core::read_schema_version(&db).expect("read version"), + Some(santi_core::SCHEMA_VERSION) + ); + let webhooks = store.list_webhooks().expect("list webhooks"); + assert_eq!(webhooks.len(), 1); + let webhook = &webhooks[0]; + assert_eq!(webhook.name, "secretary"); + assert_eq!(webhook.adaptor, "github"); + assert_eq!(webhook.soul_id, "soul_default"); + assert_eq!(webhook.strand_strategy, "per_thread"); + assert_eq!(webhook.secret_env, "SANTI_WEBHOOK_GITHUB_SECRET"); + assert!(store.soul("soul_default").expect("soul").is_some()); + drop(store); + + let conn = Connection::open(&db).expect("open sqlite"); + for column in ["source_type", "source_ref", "source_metadata"] { + let exists: i64 = conn + .query_row( + "SELECT COUNT(*) FROM pragma_table_info('strand_inbox') WHERE name = ?1", + [column], + |row| row.get(0), + ) + .expect("column lookup"); + assert_eq!(exists, 1, "missing migrated column {column}"); + } + let pending_count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM strand_inbox WHERE id = 'inbox_old'", + [], + |row| row.get(0), + ) + .expect("pending row count"); + assert_eq!(pending_count, 1, "v21 pending inbox row was wiped"); + let webhook_count: i64 = conn + .query_row( + "SELECT COUNT(*) FROM webhooks WHERE name = 'secretary'", + [], + |row| row.get(0), + ) + .expect("webhook count"); + assert_eq!(webhook_count, 1, "webhook subscription was wiped"); +} + #[test] fn read_schema_version_is_readonly_and_matches_after_open() { let temp = tempfile::tempdir().expect("temp dir");