Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion crates/santi-api/src/ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,10 +128,11 @@ fn inbox_seed_existing_strand(
strand_id: &str,
text: &str,
) -> Result<SeedReport, String> {
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 {
Expand Down
28 changes: 22 additions & 6 deletions crates/santi-api/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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 { .. } => {}
Expand Down
31 changes: 31 additions & 0 deletions crates/santi-api/src/webhook.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Value>,
/// 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,
Expand Down Expand Up @@ -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,
}));
Expand Down Expand Up @@ -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,
}))
Expand Down Expand Up @@ -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,
}));
Expand Down Expand Up @@ -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,
}))
Expand Down
31 changes: 31 additions & 0 deletions crates/santi-core/src/model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String>,
pub metadata: Option<Value>,
}

impl InboxSource {
pub fn new(source_type: impl Into<String>) -> Self {
Self {
source_type: source_type.into(),
source_ref: None,
metadata: None,
}
}

pub fn with_ref(mut self, source_ref: impl Into<String>) -> 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:<participant_id>` labels; the reply-routing correlation
/// strips this back to the participant. Shared by the IM store, the API send
Expand Down Expand Up @@ -483,6 +513,7 @@ pub enum SantiStreamPayload {
pub struct StrandRuntimeSnapshot {
pub strand: Strand,
pub messages: Vec<StrandMessage>,
pub message_events: Vec<MessageEvent>,
pub turns: Vec<Turn>,
pub thinking_spans: Vec<ThinkingSpan>,
pub tool_calls: Vec<ToolCall>,
Expand Down
40 changes: 33 additions & 7 deletions crates/santi-core/src/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -232,9 +232,20 @@ impl SantiService {
content: MessageContent,
kind: MessageKind,
trigger_type: &str,
) -> Result<IngestOutcome, String> {
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<InboxSource>,
) -> Result<IngestOutcome, String> {
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)
}

Expand All @@ -251,8 +262,11 @@ impl SantiService {
content: MessageContent,
kind: MessageKind,
trigger_type: &str,
source: Option<InboxSource>,
) -> 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,
Expand All @@ -270,6 +284,16 @@ impl SantiService {
soul_id: &str,
label: &str,
system_text: String,
) -> Result<IngestOutcome, String> {
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<InboxSource>,
) -> Result<IngestOutcome, String> {
let strand = self
.store
Expand All @@ -282,6 +306,7 @@ impl SantiService {
MessageContent::text(system_text),
MessageKind::SantiSystem,
"system",
source,
)?;
Ok(outcome)
}
Expand Down Expand Up @@ -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);
Expand Down
11 changes: 7 additions & 4 deletions crates/santi-core/src/service/im.rs
Original file line number Diff line number Diff line change
@@ -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:<participant>` conversation label; the participant address lives only
//! in the IM envelope (the label + the IM store), never in the runtime.
//! by an `im:<participant>` 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;
Expand All @@ -23,14 +25,15 @@ impl SantiService {
) -> Result<IngestOutcome, String> {
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,
},
MessageContent::text(content.to_string()),
MessageKind::Text,
"strand_send",
Some(InboxSource::new("im").with_ref(participant_id.to_string())),
)
}

Expand Down
Loading