diff --git a/core/Cargo.lock b/core/Cargo.lock index 278bd3c86..d28d739cd 100644 --- a/core/Cargo.lock +++ b/core/Cargo.lock @@ -1410,6 +1410,15 @@ dependencies = [ "tokio-tungstenite", ] +[[package]] +name = "solstone-core-callosum" +version = "1.0.22" +dependencies = [ + "serde", + "serde_json", + "solstone-core-journal-io", +] + [[package]] name = "solstone-core-cli" version = "1.0.22" @@ -1473,6 +1482,7 @@ dependencies = [ "serde", "serde_json", "sha2", + "solstone-core-callosum", "solstone-core-convey-http", "solstone-core-ingest-resolve", "solstone-core-segment", diff --git a/core/Cargo.toml b/core/Cargo.toml index eaf99335e..7445e2895 100644 --- a/core/Cargo.toml +++ b/core/Cargo.toml @@ -8,6 +8,7 @@ members = [ "crates/solstone-core-entity", "crates/solstone-core-journal", "crates/solstone-core-journal-io", + "crates/solstone-core-callosum", "crates/solstone-core-segment", "crates/solstone-core-sol", "crates/solstone-core-sol-client", @@ -40,6 +41,7 @@ solstone-core-indexer-store = { path = "crates/solstone-core-indexer-store" } solstone-core-ingest-resolve = { path = "crates/solstone-core-ingest-resolve" } solstone-core-journal = { path = "crates/solstone-core-journal" } solstone-core-journal-io = { path = "crates/solstone-core-journal-io" } +solstone-core-callosum = { path = "crates/solstone-core-callosum" } solstone-core-segment = { path = "crates/solstone-core-segment" } solstone-core-sol = { path = "crates/solstone-core-sol" } solstone-core-sol-client = { path = "crates/solstone-core-sol-client" } diff --git a/core/crates/solstone-core-callosum/Cargo.toml b/core/crates/solstone-core-callosum/Cargo.toml new file mode 100644 index 000000000..c639159f2 --- /dev/null +++ b/core/crates/solstone-core-callosum/Cargo.toml @@ -0,0 +1,18 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +[package] +name = "solstone-core-callosum" +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true +publish = false + +[dependencies] +serde = { workspace = true } +serde_json = { workspace = true } +solstone-core-journal-io = { workspace = true } + +[lints] +workspace = true diff --git a/core/crates/solstone-core-callosum/src/lib.rs b/core/crates/solstone-core-callosum/src/lib.rs new file mode 100644 index 000000000..78428c331 --- /dev/null +++ b/core/crates/solstone-core-callosum/src/lib.rs @@ -0,0 +1,19 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +//! Callosum wire envelopes and their durable per-segment event log. + +#![deny(clippy::disallowed_methods, clippy::disallowed_types)] + +mod model; +mod reader; +mod registry; +mod writer; + +pub use model::{CallosumEnvelope, DeviceIngestEvent, DurableEvent, FileDescriptor}; +pub use reader::{ + CallosumReadError, DeviceIngestReport, DurableEventsReport, read_device_ingest_events, + read_durable_events, +}; +pub use registry::callosum_registry; +pub use writer::{CallosumWriteError, append_durable_event}; diff --git a/core/crates/solstone-core-callosum/src/model.rs b/core/crates/solstone-core-callosum/src/model.rs new file mode 100644 index 000000000..006594d8d --- /dev/null +++ b/core/crates/solstone-core-callosum/src/model.rs @@ -0,0 +1,53 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use serde::{Deserialize, Serialize}; +use serde_json::{Map, Value}; + +/// A Callosum wire message with all extension keys preserved. +#[derive(Clone, Debug, Deserialize, Serialize)] +pub struct CallosumEnvelope { + pub tract: String, + pub event: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub ts: Option, + #[serde(flatten)] + pub extra: Map, +} + +/// A file attributed to a device-ingest record. +#[derive(Clone, Debug, Deserialize, Serialize)] +pub struct FileDescriptor { + pub submitted: String, + pub written: String, + pub size: u64, + pub sha256: String, + #[serde(flatten)] + pub extra: Map, +} + +/// Durable attribution for a linked-device ingest. +#[derive(Clone, Debug, Deserialize, Serialize)] +pub struct DeviceIngestEvent { + pub record_type: String, + pub record_version: u8, + pub outcome: String, + pub protocol_version: u8, + pub did: String, + pub source: String, + pub stream: String, + pub day: String, + pub segment: String, + pub files: Vec, + pub meta: Map, + #[serde(flatten)] + pub extra: Map, +} + +/// One recognized durable event-log row. +#[derive(Clone, Debug)] +#[allow(clippy::large_enum_variant)] +pub enum DurableEvent { + Callosum(CallosumEnvelope), + DeviceIngest(DeviceIngestEvent), +} diff --git a/core/crates/solstone-core-callosum/src/reader.rs b/core/crates/solstone-core-callosum/src/reader.rs new file mode 100644 index 000000000..40e8869fc --- /dev/null +++ b/core/crates/solstone-core-callosum/src/reader.rs @@ -0,0 +1,279 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::error::Error; +use std::fmt; +use std::fs; +use std::io; +use std::path::{Path, PathBuf}; + +use serde_json::Value; + +use crate::{CallosumEnvelope, DeviceIngestEvent, DurableEvent}; + +const EVENTS_FILE: &str = "events.jsonl"; + +/// Parsed durable rows and row-level recovery counts. +#[derive(Clone, Debug, Default)] +pub struct DurableEventsReport { + pub records: Vec, + pub unparseable: usize, + pub unrecognized: usize, +} + +/// Device-ingest rows from a mixed durable event log. +#[derive(Clone, Debug, Default)] +pub struct DeviceIngestReport { + pub records: Vec, + pub wrong_family: usize, + pub unparseable: usize, + pub unrecognized: usize, +} + +/// A durable event-log filesystem read failure. +#[derive(Debug)] +pub enum CallosumReadError { + Io { path: PathBuf, source: io::Error }, +} + +impl fmt::Display for CallosumReadError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Io { path, source } => write!(formatter, "{}: {source}", path.display()), + } + } +} + +impl Error for CallosumReadError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::Io { source, .. } => Some(source), + } + } +} + +/// Read a segment's mixed durable event log with per-row recovery. +/// +/// The file is read as bytes so an invalid UTF-8 tail from an interrupted write +/// only skips that line rather than aborting valid surrounding records. +pub fn read_durable_events(segment_path: &Path) -> Result { + let path = segment_path.join(EVENTS_FILE); + let contents = match fs::read(&path) { + Ok(contents) => contents, + Err(error) if error.kind() == io::ErrorKind::NotFound => { + return Ok(DurableEventsReport::default()); + } + Err(source) => return Err(CallosumReadError::Io { path, source }), + }; + + let mut report = DurableEventsReport::default(); + for line in contents.split(|byte| *byte == b'\n') { + if line.iter().all(u8::is_ascii_whitespace) { + continue; + } + match classify(line) { + Ok(Some(record)) => report.records.push(record), + Ok(None) => report.unrecognized += 1, + Err(()) => report.unparseable += 1, + } + } + Ok(report) +} + +/// Read only device-ingest rows while reporting skipped bus-envelope rows. +pub fn read_device_ingest_events( + segment_path: &Path, +) -> Result { + let durable = read_durable_events(segment_path)?; + let mut report = DeviceIngestReport { + unparseable: durable.unparseable, + unrecognized: durable.unrecognized, + ..DeviceIngestReport::default() + }; + for record in durable.records { + match record { + DurableEvent::DeviceIngest(record) => report.records.push(record), + DurableEvent::Callosum(_) => report.wrong_family += 1, + } + } + Ok(report) +} + +fn classify(line: &[u8]) -> Result, ()> { + let value: Value = match serde_json::from_slice(line) { + Ok(value) => value, + Err(_) => return Err(()), + }; + let Some(object) = value.as_object() else { + return Ok(None); + }; + if object.get("record_type").and_then(Value::as_str) == Some("device_ingest") { + return serde_json::from_value(value) + .map(DurableEvent::DeviceIngest) + .map(Some) + .map_err(|_| ()); + } + if object.get("tract").and_then(Value::as_str).is_some() + && object.get("event").and_then(Value::as_str).is_some() + { + return serde_json::from_value::(value) + .map(DurableEvent::Callosum) + .map(Some) + .map_err(|_| ()); + } + Ok(None) +} + +#[cfg(test)] +#[allow(clippy::disallowed_methods, clippy::disallowed_types)] +mod tests { + use std::fs; + use std::path::{Path, PathBuf}; + use std::sync::atomic::{AtomicUsize, Ordering}; + + use serde_json::{Value, json}; + + use super::*; + + static NEXT_PATH: AtomicUsize = AtomicUsize::new(0); + + fn segment_path(name: &str) -> PathBuf { + let suffix = NEXT_PATH.fetch_add(1, Ordering::Relaxed); + std::env::temp_dir().join(format!("solstone-core-callosum-{name}-{suffix}")) + } + + fn write_events(segment: &Path, contents: &[u8]) { + fs::create_dir_all(segment).unwrap(); + fs::write(segment.join(EVENTS_FILE), contents).unwrap(); + } + + fn device_event(extra: Value) -> Value { + let mut event = json!({ + "record_type":"device_ingest", + "record_version":1, + "outcome":"accepted", + "protocol_version":3, + "did":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + "source":"", + "stream":"device", + "day":"20260804", + "segment":"120000_1", + "files":[], + "meta":{} + }); + event + .as_object_mut() + .unwrap() + .extend(extra.as_object().unwrap().clone()); + event + } + + #[test] + fn torn_multibyte_tail_does_not_abort_valid_records() { + let segment = segment_path("torn-multibyte"); + write_events( + &segment, + b"{\"tract\":\"observe\",\"event\":\"status\"}\n{\"record_type\":\"other\"}\n{\"tail\":\"\xe2", + ); + + let report = read_durable_events(&segment).unwrap(); + assert_eq!(report.records.len(), 1); + assert_eq!(report.unrecognized, 1); + assert_eq!(report.unparseable, 1); + let _ = fs::remove_dir_all(segment); + } + + #[test] + fn invalid_utf8_and_bad_json_preserve_valid_rows_and_counts() { + let segment = segment_path("invalid-rows"); + write_events( + &segment, + b"{\"tract\":\"observe\",\"event\":\"status\"}\n{bad json}\n{\"record_type\":\"device_ingest\",\"record_version\":1,\"outcome\":\"accepted\",\"protocol_version\":3,\"did\":\"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\",\"source\":\"\",\"stream\":\"device\",\"day\":\"20260804\",\"segment\":\"120000_1\",\"files\":[],\"meta\":{}}\n{\"tail\":\"\xe2", + ); + + let report = read_durable_events(&segment).unwrap(); + assert_eq!(report.records.len(), 2); + assert_eq!(report.unparseable, 2); + assert_eq!(report.unrecognized, 0); + let _ = fs::remove_dir_all(segment); + } + + #[test] + fn envelope_round_trip_preserves_unknown_values_and_number_kinds() { + let original = json!({ + "tract":"unregistered", + "event":"future", + "ts":123, + "integer":3, + "float":3.0, + "nested":{"future":true} + }); + let envelope: CallosumEnvelope = serde_json::from_value(original.clone()).unwrap(); + let round_trip = serde_json::to_value(envelope).unwrap(); + + assert_eq!(round_trip, original); + assert!(round_trip["integer"].as_number().unwrap().is_i64()); + assert!(round_trip["float"].as_number().unwrap().is_f64()); + } + + #[test] + fn matching_envelopes_keep_independent_extra_keys() { + let segment = segment_path("independent-extra"); + write_events( + &segment, + b"{\"tract\":\"future\",\"event\":\"same\",\"one\":1}\n{\"tract\":\"future\",\"event\":\"same\",\"two\":2}\n", + ); + + let report = read_durable_events(&segment).unwrap(); + let DurableEvent::Callosum(first) = &report.records[0] else { + panic!("expected envelope"); + }; + let DurableEvent::Callosum(second) = &report.records[1] else { + panic!("expected envelope"); + }; + assert_eq!(first.extra["one"], json!(1)); + assert!(first.extra.get("two").is_none()); + assert_eq!(second.extra["two"], json!(2)); + assert!(second.extra.get("one").is_none()); + let _ = fs::remove_dir_all(segment); + } + + #[test] + fn mixed_families_are_attributed_and_device_reader_counts_bus_rows() { + let segment = segment_path("mixed-families"); + let device = device_event(json!({"future_device_key":true})); + write_events( + &segment, + format!("{{\"tract\":\"future\",\"event\":\"open\",\"unknown\":true}}\n{device}\n") + .as_bytes(), + ); + + let durable = read_durable_events(&segment).unwrap(); + assert!(matches!(durable.records[0], DurableEvent::Callosum(_))); + let DurableEvent::DeviceIngest(event) = &durable.records[1] else { + panic!("expected device ingest"); + }; + assert_eq!(event.extra["future_device_key"], json!(true)); + + let device_only = read_device_ingest_events(&segment).unwrap(); + assert_eq!(device_only.records.len(), 1); + assert_eq!(device_only.wrong_family, 1); + let _ = fs::remove_dir_all(segment); + } + + #[test] + fn device_reader_reports_disjoint_skip_reasons() { + let segment = segment_path("skip-reasons"); + write_events( + &segment, + b"{\"tract\":\"observe\",\"event\":\"status\"}\n{\"record_type\":\"device_ingest\"}\n{\"record_type\":\"other\"}\n", + ); + + let report = read_device_ingest_events(&segment).unwrap(); + assert!(report.records.is_empty()); + assert_eq!(report.wrong_family, 1); + assert_eq!(report.unparseable, 1); + assert_eq!(report.unrecognized, 1); + let _ = fs::remove_dir_all(segment); + } +} diff --git a/core/crates/solstone-core-callosum/src/registry.rs b/core/crates/solstone-core-callosum/src/registry.rs new file mode 100644 index 000000000..c79cfcca6 --- /dev/null +++ b/core/crates/solstone-core-callosum/src/registry.rs @@ -0,0 +1,32 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::sync::OnceLock; + +use serde_json::Value; + +const REGISTRY_FIXTURE: &str = include_str!("../../../fixtures/callosum_registry.json"); +static REGISTRY: OnceLock = OnceLock::new(); + +/// Return the generated Callosum vocabulary fixture for inspection only. +pub fn callosum_registry() -> &'static Value { + REGISTRY.get_or_init(|| { + serde_json::from_str(REGISTRY_FIXTURE).expect("callosum registry fixture is valid JSON") + }) +} + +#[cfg(test)] +mod tests { + use serde_json::Value; + + use super::callosum_registry; + + #[test] + fn registry_matches_independently_parsed_fixture() { + let expected: Value = + serde_json::from_str(include_str!("../../../fixtures/callosum_registry.json")) + .expect("callosum registry fixture is valid JSON"); + + assert_eq!(callosum_registry(), &expected); + } +} diff --git a/core/crates/solstone-core-callosum/src/writer.rs b/core/crates/solstone-core-callosum/src/writer.rs new file mode 100644 index 000000000..ac7c9bc2b --- /dev/null +++ b/core/crates/solstone-core-callosum/src/writer.rs @@ -0,0 +1,148 @@ +// SPDX-License-Identifier: AGPL-3.0-only +// Copyright (c) 2026 sol pbc + +use std::error::Error; +use std::fmt; +use std::path::{Path, PathBuf}; + +use solstone_core_journal_io::{AppendError, append_jsonl}; + +use crate::DurableEvent; + +const EVENTS_FILE: &str = "events.jsonl"; + +/// A durable Callosum event-log append failure. +#[derive(Debug)] +pub enum CallosumWriteError { + SegmentDirectoryMissing(PathBuf), + Append(AppendError), +} + +impl fmt::Display for CallosumWriteError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::SegmentDirectoryMissing(path) => { + write!( + formatter, + "segment directory is unavailable: {}", + path.display() + ) + } + Self::Append(error) => error.fmt(formatter), + } + } +} + +impl Error for CallosumWriteError { + fn source(&self) -> Option<&(dyn Error + 'static)> { + match self { + Self::SegmentDirectoryMissing(_) => None, + Self::Append(error) => Some(error), + } + } +} + +/// Append one recognized event row to an existing segment's `events.jsonl`. +/// +/// This checks the segment directory before delegating to journal-io because +/// journal-io's append primitive creates missing parents. A concurrent delete +/// after this check remains an accepted race for this write path. +pub fn append_durable_event( + segment_path: &Path, + event: &DurableEvent, +) -> Result<(), CallosumWriteError> { + if !segment_path.is_dir() { + return Err(CallosumWriteError::SegmentDirectoryMissing( + segment_path.to_path_buf(), + )); + } + let path = segment_path.join(EVENTS_FILE); + match event { + DurableEvent::Callosum(event) => append_jsonl(path, event), + DurableEvent::DeviceIngest(event) => append_jsonl(path, event), + } + .map_err(CallosumWriteError::Append) +} + +#[cfg(test)] +#[allow(clippy::disallowed_methods, clippy::disallowed_types)] +mod tests { + use std::fs; + use std::path::PathBuf; + use std::sync::atomic::{AtomicUsize, Ordering}; + + use serde_json::Map; + + use super::*; + use crate::CallosumEnvelope; + + static NEXT_PATH: AtomicUsize = AtomicUsize::new(0); + + fn path(name: &str) -> PathBuf { + let suffix = NEXT_PATH.fetch_add(1, Ordering::Relaxed); + std::env::temp_dir().join(format!("solstone-core-callosum-writer-{name}-{suffix}")) + } + + fn event() -> DurableEvent { + DurableEvent::Callosum(CallosumEnvelope { + tract: "observe".to_owned(), + event: "status".to_owned(), + ts: None, + extra: Map::new(), + }) + } + + #[test] + fn appends_to_an_existing_segment() { + let segment = path("success"); + fs::create_dir_all(&segment).unwrap(); + + append_durable_event(&segment, &event()).unwrap(); + + assert_eq!( + fs::read_to_string(segment.join(EVENTS_FILE)).unwrap(), + "{\"tract\":\"observe\",\"event\":\"status\"}\n" + ); + let _ = fs::remove_dir_all(segment); + } + + #[test] + fn missing_segment_directory_is_not_materialized() { + let segment = path("missing"); + + assert!(matches!( + append_durable_event(&segment, &event()), + Err(CallosumWriteError::SegmentDirectoryMissing(_)) + )); + assert!(!segment.exists()); + } + + #[test] + fn blocked_parent_is_not_materialized() { + let root = path("blocked-parent"); + let blocked = root.join("chronicle/20260804/workstation"); + fs::create_dir_all(blocked.parent().unwrap()).unwrap(); + fs::write(&blocked, b"not a directory").unwrap(); + let segment = blocked.join("120000_60"); + + assert!(matches!( + append_durable_event(&segment, &event()), + Err(CallosumWriteError::SegmentDirectoryMissing(_)) + )); + assert!(blocked.is_file()); + assert!(!segment.exists()); + let _ = fs::remove_dir_all(root); + } + + #[test] + fn journal_io_append_failure_is_propagated() { + let segment = path("append-failure"); + fs::create_dir_all(segment.join(EVENTS_FILE)).unwrap(); + + assert!(matches!( + append_durable_event(&segment, &event()), + Err(CallosumWriteError::Append(_)) + )); + let _ = fs::remove_dir_all(segment); + } +} diff --git a/core/crates/solstone-core-ingest/Cargo.toml b/core/crates/solstone-core-ingest/Cargo.toml index 11fa83beb..2e142b4db 100644 --- a/core/crates/solstone-core-ingest/Cargo.toml +++ b/core/crates/solstone-core-ingest/Cargo.toml @@ -15,6 +15,7 @@ log = { workspace = true } serde = { workspace = true } serde_json = { workspace = true } sha2.workspace = true +solstone-core-callosum = { workspace = true } solstone-core-convey-http = { workspace = true } solstone-core-ingest-resolve = { workspace = true } solstone-core-segment = { workspace = true } diff --git a/core/crates/solstone-core-ingest/src/events.rs b/core/crates/solstone-core-ingest/src/events.rs deleted file mode 100644 index 1906bb1bd..000000000 --- a/core/crates/solstone-core-ingest/src/events.rs +++ /dev/null @@ -1,60 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-only -// Copyright (c) 2026 sol pbc - -use std::fs; -use std::io; -use std::path::Path; - -use crate::model::{DeviceIngestEvent, ReasonCode}; - -pub fn read_events(segment: &Path) -> Result, ReasonCode> { - let path = segment.join("events.jsonl"); - let text = match fs::read_to_string(&path) { - Ok(text) => text, - Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()), - Err(_) => return Err(ReasonCode::JournalReadFailed), - }; - let mut records = Vec::new(); - for line in text.lines().filter(|line| !line.is_empty()) { - let value: serde_json::Value = - serde_json::from_str(line).map_err(|_| ReasonCode::IngestEventLogMalformed)?; - if value.get("record_type").and_then(serde_json::Value::as_str) != Some("device_ingest") { - continue; - } - records - .push(serde_json::from_value(value).map_err(|_| ReasonCode::IngestEventLogMalformed)?); - } - Ok(records) -} - -#[cfg(test)] -#[allow(clippy::disallowed_methods, clippy::disallowed_types)] -mod tests { - use std::fs; - use std::time::{SystemTime, UNIX_EPOCH}; - - use serde_json::json; - - use super::read_events; - - #[test] - fn ignores_unrelated_event_records() { - let suffix = SystemTime::now() - .duration_since(UNIX_EPOCH) - .unwrap() - .as_nanos(); - let root = std::env::temp_dir().join(format!("solstone-core-ingest-events-{suffix}")); - fs::create_dir_all(&root).unwrap(); - let ingest = json!({"record_type":"device_ingest","record_version":1,"outcome":"accepted","protocol_version":3,"did":"sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","source":"","stream":"_default","day":"20260804","segment":"120000_1","files":[],"meta":{}}); - fs::write( - root.join("events.jsonl"), - format!("{}\n{{\"record_type\":\"other\"}}\n", ingest), - ) - .unwrap(); - - let records = read_events(&root).unwrap(); - assert_eq!(records.len(), 1); - assert_eq!(records[0].record_type, "device_ingest"); - let _ = fs::remove_dir_all(root); - } -} diff --git a/core/crates/solstone-core-ingest/src/lib.rs b/core/crates/solstone-core-ingest/src/lib.rs index 71afcda5b..a23564b62 100644 --- a/core/crates/solstone-core-ingest/src/lib.rs +++ b/core/crates/solstone-core-ingest/src/lib.rs @@ -23,7 +23,6 @@ #![deny(clippy::disallowed_methods, clippy::disallowed_types)] -mod events; mod model; mod read_routes; mod router; @@ -37,7 +36,6 @@ mod architecture_tests { // Unlike solstone-core-segment's older scanner, `public_signatures` also // recognizes `pub async fn`; handlers must not open a raw byte-write door. const SOURCES: &[&str] = &[ - include_str!("events.rs"), include_str!("model.rs"), include_str!("read_routes.rs"), include_str!("router.rs"), diff --git a/core/crates/solstone-core-ingest/src/model.rs b/core/crates/solstone-core-ingest/src/model.rs index 402c38f82..ea80a5b51 100644 --- a/core/crates/solstone-core-ingest/src/model.rs +++ b/core/crates/solstone-core-ingest/src/model.rs @@ -1,7 +1,6 @@ // SPDX-License-Identifier: AGPL-3.0-only // Copyright (c) 2026 sol pbc -use serde::{Deserialize, Serialize}; use serde_json::{Map, Value}; /// Closed vocabulary for every ingest refusal. @@ -41,12 +40,11 @@ pub enum ReasonCode { EventAppendFailed, StreamAdvanceFailed, JournalReadFailed, - IngestEventLogMalformed, } impl ReasonCode { #[cfg(test)] - const ALL: [Self; 35] = [ + const ALL: [Self; 34] = [ Self::ProtocolVersionRequired, Self::ProtocolVersionMalformed, Self::ProtocolVersionLegacy, @@ -81,7 +79,6 @@ impl ReasonCode { Self::EventAppendFailed, Self::StreamAdvanceFailed, Self::JournalReadFailed, - Self::IngestEventLogMalformed, ]; pub const fn as_str(self) -> &'static str { @@ -120,7 +117,6 @@ impl ReasonCode { Self::EventAppendFailed => "event_append_failed", Self::StreamAdvanceFailed => "stream_advance_failed", Self::JournalReadFailed => "journal_read_failed", - Self::IngestEventLogMalformed => "ingest_event_log_malformed", } } } @@ -141,31 +137,6 @@ mod tests { } } -#[derive(Clone, Debug, Deserialize, Serialize)] -pub struct FileDescriptor { - pub submitted: String, - pub written: String, - pub size: u64, - pub sha256: String, - #[serde(flatten)] - pub extra: Map, -} - -#[derive(Clone, Debug, Deserialize, Serialize)] -pub struct DeviceIngestEvent { - pub record_type: String, - pub record_version: u8, - pub outcome: String, - pub protocol_version: u8, - pub did: String, - pub source: String, - pub stream: String, - pub day: String, - pub segment: String, - pub files: Vec, - pub meta: Map, -} - #[derive(Clone, Debug)] pub struct IncomingFile { pub submitted: String, diff --git a/core/crates/solstone-core-ingest/src/read_routes.rs b/core/crates/solstone-core-ingest/src/read_routes.rs index 7cc3106d1..17bdb1563 100644 --- a/core/crates/solstone-core-ingest/src/read_routes.rs +++ b/core/crates/solstone-core-ingest/src/read_routes.rs @@ -9,12 +9,12 @@ use axum::extract::{Extension, Path, Query, State}; use axum::http::{HeaderMap, StatusCode}; use axum::response::{IntoResponse, Response}; use serde_json::{Map, Value, json}; +use solstone_core_callosum::{DeviceIngestEvent, read_device_ingest_events}; use solstone_core_convey_http::identity::AccessBasis; use solstone_core_segment::lookup_stream; use solstone_core_segment::{list_days, list_segments, list_segments_in}; -use crate::events::read_events; -use crate::model::{DeviceIngestEvent, ReasonCode}; +use crate::model::ReasonCode; use crate::router::{IngestState, refusal}; use crate::validation::{validate_access, validate_day, validate_protocol, validate_source}; @@ -52,11 +52,11 @@ pub async fn ingest_manifest( .into_iter() .filter(|segment| segment.stream == stream) { - let events = match read_events(&segment.path) { - Ok(events) => events, - Err(code) => { + let events = match read_device_ingest_events(&segment.path) { + Ok(report) => report.records, + Err(_) => { return refusal( - code, + ReasonCode::JournalReadFailed, StatusCode::INTERNAL_SERVER_ERROR, "cannot read journal", ); @@ -242,7 +242,9 @@ fn stream_events( .into_iter() .filter(|segment| segment.stream == stream) { - for event in read_events(&segment.path)? { + let report = + read_device_ingest_events(&segment.path).map_err(|_| ReasonCode::JournalReadFailed)?; + for event in report.records { if event.did == did { events.push(event); } diff --git a/core/crates/solstone-core-ingest/src/router.rs b/core/crates/solstone-core-ingest/src/router.rs index 76dae3419..6dcf25282 100644 --- a/core/crates/solstone-core-ingest/src/router.rs +++ b/core/crates/solstone-core-ingest/src/router.rs @@ -11,6 +11,9 @@ use axum::response::{IntoResponse, Response}; use axum::routing::{get, post}; use axum::{Json, Router}; use serde_json::{Map, Value, json}; +use solstone_core_callosum::{ + DeviceIngestEvent, DurableEvent, FileDescriptor, append_durable_event, +}; use solstone_core_convey_http::envelope::{error_envelope, not_found_fallback}; use solstone_core_convey_http::identity::AccessBasis; use solstone_core_ingest_resolve::{ @@ -18,12 +21,10 @@ use solstone_core_ingest_resolve::{ IngestNotice, IngestNotifier, LoggingIngestNotifier, Resolution, apply_plan, quarantine_failed, resolve_ingest, }; -use solstone_core_segment::{ - ContentName, Kind, StreamHints, advance_bound_stream, append_event, bind_stream, -}; +use solstone_core_segment::{ContentName, Kind, StreamHints, advance_bound_stream, bind_stream}; use tower_http::limit::RequestBodyLimitLayer; -use crate::model::{DeviceIngestEvent, FileDescriptor, IncomingFile, ReasonCode}; +use crate::model::{IncomingFile, ReasonCode}; use crate::read_routes::{ingest_manifest, ingest_manifest_day, ingest_segments}; use crate::validation::{ validate_access, validate_day, validate_protocol, validate_segment, validate_source, @@ -507,8 +508,10 @@ fn write_envelope(state: &IngestState, did: &str, envelope: Envelope) -> Respons segment: applied.landed_segment.clone(), files: descriptors.clone(), meta: envelope.meta.clone(), + extra: Map::new(), }; - if append_event(&applied.segment, &event).is_err() { + let durable_event = DurableEvent::DeviceIngest(event); + if append_durable_event(applied.segment.path(), &durable_event).is_err() { return outcome_error( "failed", ReasonCode::EventAppendFailed, @@ -1654,7 +1657,7 @@ mod tests { } #[tokio::test] - async fn manifest_surfaces_malformed_device_ingest_events() { + async fn manifest_skips_malformed_device_ingest_events() { let root = root(); let app = router(&root); let request = envelope("20260804", "120000_1", json!([{"submitted":"audio.flac"}])); @@ -1680,8 +1683,8 @@ mod tests { &[], ) .await; - assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR); - assert_eq!(body["reason_code"], "ingest_event_log_malformed"); + assert_eq!(status, StatusCode::OK); + assert_eq!(body["days"], json!({})); let _ = fs::remove_dir_all(root); } diff --git a/core/crates/solstone-core-segment/src/lib.rs b/core/crates/solstone-core-segment/src/lib.rs index fec0a90fe..042464dbe 100644 --- a/core/crates/solstone-core-segment/src/lib.rs +++ b/core/crates/solstone-core-segment/src/lib.rs @@ -23,7 +23,6 @@ mod identity; mod manifest; mod projection; mod segment_dir; -mod sidecars; mod stream_record; #[cfg(test)] pub(crate) mod test_support; @@ -43,7 +42,6 @@ pub use identity::{ }; pub use projection::project_stream_name; pub use segment_dir::{SegmentDir, list_days, list_segments, list_segments_in}; -pub use sidecars::append_event; pub use solstone_core_journal_io::Segment; pub use stream_record::{ BoundStream, ResolvedStream, StreamAdvance, StreamHints, StreamRecord, advance_bound_stream, @@ -65,7 +63,6 @@ mod architecture_tests { Manifest, Projection, SegmentDir, - Sidecars, StreamRecord, Write, } @@ -78,7 +75,6 @@ mod architecture_tests { (Source::Manifest, include_str!("manifest.rs")), (Source::Projection, include_str!("projection.rs")), (Source::SegmentDir, include_str!("segment_dir.rs")), - (Source::Sidecars, include_str!("sidecars.rs")), (Source::StreamRecord, include_str!("stream_record.rs")), (Source::Write, include_str!("write.rs")), ]; @@ -121,9 +117,6 @@ mod architecture_tests { assert!(source.contains("hold_lock")); assert!(source.contains("write_json")); } - Source::Sidecars => { - assert!(source.contains("append_jsonl")); - } Source::ContentName | Source::Error | Source::Identity diff --git a/core/crates/solstone-core-segment/src/sidecars.rs b/core/crates/solstone-core-segment/src/sidecars.rs deleted file mode 100644 index 7a1e8c1c2..000000000 --- a/core/crates/solstone-core-segment/src/sidecars.rs +++ /dev/null @@ -1,34 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-only -// Copyright (c) 2026 sol pbc - -use serde::Serialize; -use solstone_core_journal_io::append_jsonl; - -use crate::{SegmentDir, SegmentError}; - -/// Append one durable event record to a segment's journal-authored event log. -pub fn append_event(segment: &SegmentDir, record: &T) -> Result<(), SegmentError> { - append_jsonl(segment.path.join("events.jsonl"), record)?; - Ok(()) -} - -#[cfg(test)] -#[allow(clippy::disallowed_methods, clippy::disallowed_types)] -mod tests { - use std::fs; - - use crate::test_support::TempDir; - - use super::*; - - #[test] - fn append_failure_propagates() { - let temporary = TempDir::new(); - let segment = - SegmentDir::resolve(temporary.path(), "20260804", "120000_60", "workstation").unwrap(); - let blocked = temporary.path().join("chronicle/20260804/workstation"); - fs::create_dir_all(blocked.parent().unwrap()).unwrap(); - fs::write(&blocked, b"not a directory").unwrap(); - assert!(append_event(&segment, &serde_json::json!({"event": "x"})).is_err()); - } -}