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
63 changes: 63 additions & 0 deletions crates/santi-api/src/ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,29 @@ pub fn inbox_seed_at(
if store.strand(strand_id)?.is_none() {
return Err(format!("unknown strand: {strand_id}"));
}
inbox_seed_existing_strand(&store, strand_id, text)
}

/// Enqueue into the strand anchored by a stable label, creating it if it is
/// missing. This is the offline twin of webhook per-thread ingest's
/// label→strand materialization: the label is the durable routing anchor, while
/// the concrete strand id is a replaceable room.
pub(crate) fn inbox_seed_by_label_at(
paths: &RuntimePaths,
soul_id: &str,
label: &str,
text: &str,
) -> Result<SeedReport, String> {
let store = santi_core::SantiStore::open(&paths.database_path)?;
let strand = store.find_or_create_strand_by_label(soul_id, label)?;
inbox_seed_existing_strand(&store, &strand.id, text)
}

fn inbox_seed_existing_strand(
store: &santi_core::SantiStore,
strand_id: &str,
text: &str,
) -> Result<SeedReport, String> {
let outcome = store.enqueue_inbox(
strand_id,
santi_core::MessageKind::SantiSystem,
Expand Down Expand Up @@ -276,6 +299,46 @@ mod tests {
);
}

#[test]
fn seed_by_label_creates_labeled_strand_and_is_drainable_on_boot() {
let temp = tempfile::tempdir().expect("temp dir");
let paths = paths_under(temp.path());
santi_core::SantiStore::open(&paths.database_path).expect("open");

let label = "soul:soul_default:ops";
let report = inbox_seed_by_label_at(
&paths,
santi_core::DEFAULT_SOUL_ID,
label,
"upgrade finished — come look",
)
.unwrap();
assert!(report.accepted);

let store = santi_core::SantiStore::open(&paths.database_path).expect("reopen");
let strand = store
.strand(&report.strand_id)
.unwrap()
.expect("labeled strand exists");
assert_eq!(strand.external_label.as_deref(), Some(label));
assert!(
store
.strands_with_pending_requests()
.unwrap()
.contains(&report.strand_id),
"boot recovery would re-drive this strand"
);
let started = store
.try_start_turn(&report.strand_id, "strand_send", None)
.unwrap()
.expect("a turn starts by draining the labeled seed");
assert_eq!(started.drained_messages.len(), 1);
assert_eq!(
started.drained_messages[0].content_text,
"upgrade finished — come look"
);
}

#[test]
fn im_reply_delivers_into_the_conversations_participant_inbox() {
let temp = tempfile::tempdir().expect("temp dir");
Expand Down
199 changes: 176 additions & 23 deletions crates/santi-api/src/upgrade.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,9 +68,21 @@ pub struct UpgradeReport {
pub outcome: Outcome,
/// The text seeded into the soul's inbox ("come look" / failure feedback).
pub record: String,
/// Whether the record was durably enqueued (false ⟺ no self-strand
/// configured or the seed was rejected; the upgrade still completes).
/// Whether the record was durably enqueued (false ⟺ the seed failed; the
/// upgrade still completes, but `warnings` makes the failure loud).
pub seeded: bool,
/// The concrete strand that received the record. It is a materialized room;
/// the durable addressing anchor is a stable label.
pub seeded_strand_id: Option<String>,
/// Explicit warnings that must not be hidden behind `seeded: false`.
pub warnings: Vec<String>,
}

/// The successful result of seeding a come-look record.
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct SeedOutcome {
pub strand_id: String,
pub warnings: Vec<String>,
}

/// The launcher's fast return: "started listening, max timeout, log location".
Expand Down Expand Up @@ -98,9 +110,10 @@ pub trait UpgradeHost {
fn trial_probe(&mut self) -> Result<bool, String>;
/// Restore the snapshot + reinstall the previous version (the final = OLD).
fn rollback(&mut self) -> Result<(), String>;
/// Seed one durable "come look" record into the soul's self-strand, offline,
/// using the FINAL version's binary/schema. Best-effort at the call site.
fn seed(&mut self, text: &str) -> Result<(), String>;
/// Seed one durable "come look" record into the soul's stable ops strand,
/// offline, using the FINAL version's binary/schema. Best-effort at the call
/// site, but failures are reported loudly in `UpgradeReport::warnings`.
fn seed(&mut self, text: &str) -> Result<SeedOutcome, String>;
/// Start the FINAL version for real (boot recovery then drains the seed).
fn start(&mut self) -> Result<(), String>;
}
Expand Down Expand Up @@ -138,9 +151,22 @@ pub fn run_upgrade<H: UpgradeHost>(

// Seed the truthful record (offline, into the FINAL DB) BEFORE the final
// start, so boot recovery is guaranteed to wake the soul into it. Seeding is
// best-effort: a missing self-strand must not abort an otherwise-done upgrade.
// best-effort, but failures must be loud: never hide them behind `seeded:false`.
let record = compose_record(deb, &outcome);
let seeded = host.seed(&record).is_ok();
let mut warnings = Vec::new();
let (seeded, seeded_strand_id) = match host.seed(&record) {
Ok(seed) => {
warnings.extend(seed.warnings);
(true, Some(seed.strand_id))
}
Err(error) => {
warnings.push(format!("come-look seed failed: {error}"));
(false, None)
}
};
for warning in &warnings {
eprintln!("santi: upgrade warning: {warning}");
}

// Start the final version for real.
host.start()?;
Expand All @@ -149,6 +175,8 @@ pub fn run_upgrade<H: UpgradeHost>(
outcome,
record,
seeded,
seeded_strand_id,
warnings,
})
}

Expand Down Expand Up @@ -182,6 +210,73 @@ pub fn upgrade_timeout() -> Duration {
Duration::from_secs(secs)
}

fn self_ops_label(soul_id: &str) -> String {
format!("soul:{soul_id}:ops")
}

fn seed_come_look_at(
paths: &RuntimePaths,
soul_id: &str,
configured_strand: Option<&str>,
text: &str,
) -> Result<SeedOutcome, String> {
let label = self_ops_label(soul_id);
let configured_strand = configured_strand
.map(str::trim)
.filter(|value| !value.is_empty());

match crate::ops::inbox_seed_by_label_at(paths, soul_id, &label, text) {
Ok(report) if report.accepted => Ok(SeedOutcome {
strand_id: report.strand_id,
warnings: Vec::new(),
}),
Ok(report) => seed_come_look_via_configured_strand(
paths,
configured_strand,
text,
format!(
"stable self-strand label {label} rejected the come-look seed: {}",
report.reason.unwrap_or_else(|| "seed rejected".to_string())
),
),
Err(error) => seed_come_look_via_configured_strand(
paths,
configured_strand,
text,
format!(
"stable self-strand label {label} could not receive the come-look seed: {error}"
),
),
}
}

fn seed_come_look_via_configured_strand(
paths: &RuntimePaths,
configured_strand: Option<&str>,
text: &str,
stable_label_error: String,
) -> Result<SeedOutcome, String> {
let Some(strand_id) = configured_strand else {
return Err(stable_label_error);
};

match crate::ops::inbox_seed_at(paths, strand_id, text) {
Ok(report) if report.accepted => Ok(SeedOutcome {
strand_id: report.strand_id,
warnings: vec![format!(
"{stable_label_error}; fell back to configured SANTI_STRAND_ID {strand_id}"
)],
}),
Ok(report) => Err(format!(
"{stable_label_error}; configured SANTI_STRAND_ID {strand_id} also rejected the come-look seed: {}",
report.reason.unwrap_or_else(|| "seed rejected".to_string())
)),
Err(error) => Err(format!(
"{stable_label_error}; configured SANTI_STRAND_ID {strand_id} also could not receive the come-look seed: {error}"
)),
}
}

fn request_path(paths: &RuntimePaths) -> PathBuf {
paths.runtime_root.join("upgrade.request")
}
Expand Down Expand Up @@ -349,18 +444,14 @@ impl UpgradeHost for SystemHost {
Ok(())
}

fn seed(&mut self, text: &str) -> Result<(), String> {
let strand = env::var("SANTI_STRAND_ID")
fn seed(&mut self, text: &str) -> Result<SeedOutcome, String> {
let soul_id = env::var("SANTI_SOUL_ID")
.ok()
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
.ok_or("no SANTI_STRAND_ID configured for the come-look record")?;
let report = crate::ops::inbox_seed_at(&self.paths, &strand, text)?;
if report.accepted {
Ok(())
} else {
Err(report.reason.unwrap_or_else(|| "seed rejected".to_string()))
}
.unwrap_or_else(|| santi_core::DEFAULT_SOUL_ID.to_string());
let configured_strand = env::var("SANTI_STRAND_ID").ok();
seed_come_look_at(&self.paths, &soul_id, configured_strand.as_deref(), text)
}

fn start(&mut self) -> Result<(), String> {
Expand All @@ -376,18 +467,33 @@ mod tests {
calls: Vec<String>,
install_result: Result<(), String>,
probe_result: Result<bool, String>,
seed_result: Result<(), String>,
seed_result: Result<SeedOutcome, String>,
seeded_text: Option<String>,
}

fn fake_seed(strand_id: &str) -> SeedOutcome {
SeedOutcome {
strand_id: strand_id.to_string(),
warnings: Vec::new(),
}
}

fn paths_under(root: &std::path::Path) -> RuntimePaths {
RuntimePaths {
database_path: root.join("runtime").join("db"),
runtime_root: root.join("runtime"),
execution_root: root.join("execution"),
}
}

impl Default for FakeHost {
fn default() -> Self {
// Healthy defaults; each test overrides the field it exercises.
Self {
calls: Vec::new(),
install_result: Ok(()),
probe_result: Ok(true),
seed_result: Ok(()),
seed_result: Ok(fake_seed("ss_seeded")),
seeded_text: None,
}
}
Expand All @@ -414,7 +520,7 @@ mod tests {
self.calls.push("rollback".into());
Ok(())
}
fn seed(&mut self, text: &str) -> Result<(), String> {
fn seed(&mut self, text: &str) -> Result<SeedOutcome, String> {
self.calls.push("seed".into());
self.seeded_text = Some(text.to_string());
self.seed_result.clone()
Expand All @@ -434,7 +540,7 @@ mod tests {
let mut host = FakeHost {
install_result: Ok(()),
probe_result: Ok(true),
seed_result: Ok(()),
seed_result: Ok(fake_seed("ss_seeded")),
..Default::default()
};
let report = run(&mut host);
Expand All @@ -451,6 +557,8 @@ mod tests {
);
assert_eq!(report.outcome, Outcome::Upgraded);
assert!(report.seeded);
assert_eq!(report.seeded_strand_id.as_deref(), Some("ss_seeded"));
assert!(report.warnings.is_empty());
// Seed happens BEFORE the final start (guaranteed drain on boot).
let seed_i = host.calls.iter().position(|c| c == "seed").unwrap();
let start_i = host.calls.iter().position(|c| c == "start").unwrap();
Expand All @@ -463,7 +571,7 @@ mod tests {
let mut host = FakeHost {
install_result: Err("bad package signature".into()),
probe_result: Ok(true), // never reached
seed_result: Ok(()),
seed_result: Ok(fake_seed("ss_seeded")),
..Default::default()
};
let report = run(&mut host);
Expand Down Expand Up @@ -493,7 +601,7 @@ mod tests {
let mut host = FakeHost {
install_result: Ok(()),
probe_result: Ok(false),
seed_result: Ok(()),
seed_result: Ok(fake_seed("ss_seeded")),
..Default::default()
};
let report = run(&mut host);
Expand Down Expand Up @@ -521,7 +629,7 @@ mod tests {
let mut host = FakeHost {
install_result: Ok(()),
probe_result: Err("probe timed out".into()),
seed_result: Ok(()),
seed_result: Ok(fake_seed("ss_seeded")),
..Default::default()
};
let report = run(&mut host);
Expand All @@ -543,7 +651,52 @@ mod tests {
let report = run(&mut host);
assert_eq!(report.outcome, Outcome::Upgraded);
assert!(!report.seeded, "seed reported not durably enqueued");
assert_eq!(report.seeded_strand_id, None);
assert_eq!(
report.warnings,
vec!["come-look seed failed: no self-strand configured".to_string()]
);
// The final start still happens.
assert_eq!(host.calls.last().unwrap(), "start");
}

#[test]
fn come_look_seed_uses_stable_label_even_when_configured_strand_is_stale() {
let temp = tempfile::tempdir().expect("temp dir");
let paths = paths_under(temp.path());
santi_core::SantiStore::open(&paths.database_path).expect("open");

let outcome = seed_come_look_at(
&paths,
santi_core::DEFAULT_SOUL_ID,
Some("ss_deleted_self_strand"),
"you were upgrading — come look",
)
.expect("seed via stable label");
assert!(outcome.warnings.is_empty());

let store = santi_core::SantiStore::open(&paths.database_path).expect("reopen");
let strand = store
.strand(&outcome.strand_id)
.unwrap()
.expect("ops strand exists");
let label = self_ops_label(santi_core::DEFAULT_SOUL_ID);
assert_eq!(strand.external_label.as_deref(), Some(label.as_str()));
assert!(
store
.strands_with_pending_requests()
.unwrap()
.contains(&outcome.strand_id),
"boot recovery would re-drive the stable ops strand"
);
let started = store
.try_start_turn(&outcome.strand_id, "strand_send", None)
.unwrap()
.expect("a turn starts by draining the stable ops seed");
assert_eq!(started.drained_messages.len(), 1);
assert_eq!(
started.drained_messages[0].content_text,
"you were upgrading — come look"
);
}
}