diff --git a/src/api/router/snapshot_router.rs b/src/api/router/snapshot_router.rs index b6fe5813..6e996d4b 100644 --- a/src/api/router/snapshot_router.rs +++ b/src/api/router/snapshot_router.rs @@ -126,6 +126,14 @@ tokio::task_local! { #[cfg(test)] tokio::task_local! { static NATIVE_RESOLVE_BARRIERS: (std::sync::Arc, std::sync::Arc); + static REJECT_NATIVE_OBSERVATION_SOURCE: bool; +} + +#[cfg(test)] +pub(crate) async fn with_rejected_native_observation_source( + future: F, +) -> F::Output { + REJECT_NATIVE_OBSERVATION_SOURCE.scope(true, future).await } #[cfg(test)] @@ -460,11 +468,21 @@ async fn resolve( ) }) .and_then(|(commit, tree)| { + let certificate = head.token.certificate; + #[cfg(test)] + let certificate = if REJECT_NATIVE_OBSERVATION_SOURCE + .try_with(|reject| *reject) + .unwrap_or(false) + { + None + } else { + certificate + }; NativeResolveSource::capture( &head.instance_id, commit, tree, - head.token.certificate, + certificate, head.token.epoch, head.token.sequence, ) @@ -472,6 +490,9 @@ async fn resolve( native_source = match observation_source { Ok(source) => Some(source), Err(_) => { + if let Some(sink) = &state.storage.projection_observation_sink { + sink.reject_binding(); + } tracing::warn!("native resolve observation source rejected"); None } @@ -561,8 +582,18 @@ async fn resolve( projection_work, projection_elapsed, ) { - Ok(observation) => observation.emit(), - Err(_) => tracing::warn!("native resolve observation context rejected"), + Ok(observation) => { + if let Some(sink) = &state.storage.projection_observation_sink { + let _ = sink.enqueue(&observation); + } + observation.emit(); + } + Err(_) => { + if let Some(sink) = &state.storage.projection_observation_sink { + sink.reject_binding(); + } + tracing::warn!("native resolve observation context rejected"); + } } } diff --git a/src/ceres/snapshot/mod.rs b/src/ceres/snapshot/mod.rs index 3c408ace..b6218db8 100644 --- a/src/ceres/snapshot/mod.rs +++ b/src/ceres/snapshot/mod.rs @@ -9,6 +9,7 @@ pub mod error; pub mod frame_stream; pub mod pages; pub(crate) mod projection_observation; +pub(crate) mod projection_writer; pub mod publication; pub mod resolver; pub mod retention; diff --git a/src/ceres/snapshot/projection_observation.rs b/src/ceres/snapshot/projection_observation.rs index 3e0dc083..819e0e86 100644 --- a/src/ceres/snapshot/projection_observation.rs +++ b/src/ceres/snapshot/projection_observation.rs @@ -7,6 +7,7 @@ use mst2_codec::descriptor::{ ACCESS_PROJECTION_EXACT_FULL, FS_SEMANTICS_LINUX_CODE_V1, MATERIALIZATION_POLICY_GIT_RAW_V1, METADATA_CODEC, SCHEMA_VERSION, ServingDescriptor, }; +use serde::Serialize; use sha2::{Digest, Sha256}; use uuid::Uuid; @@ -19,7 +20,7 @@ pub(crate) const NATIVE_PROJECTION_REVISION: u16 = 1; /// Work in one directory projection. Page counters cover returned directory /// roots, excluding codec-internal radix encoding and later route traversal. -#[derive(Debug, Clone, Default, PartialEq, Eq)] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize)] pub struct ProjectionWork { pub directories_rebuilt: u64, /// Cache-hit boundary visits, not all descendants or unique source OIDs. @@ -75,6 +76,44 @@ pub(crate) struct NativeProjectionObservation { work: ProjectionWork, } +/// Closed wire fields, borrowed only from the already validated observation. +#[derive(Serialize)] +pub(crate) struct ProjectionWireRecord<'a> { + observation_revision: u16, + phase: &'static str, + source_domain: &'static str, + request_id: &'a str, + instance_id: String, + #[serde(serialize_with = "tagged_oid")] + root_commit_oid: ObjectHash, + #[serde(serialize_with = "tagged_oid")] + root_tree_oid: ObjectHash, + native_certificate_receipt_id: u64, + native_writer_epoch: u64, + native_publication_sequence: u64, + scope: &'a str, + schema_version: u16, + metadata_codec: u16, + materialization_policy: u16, + fs_semantics: u16, + access_projection: u16, + verification_revision: i32, + projection_revision: u16, + namespace_view_id: &'a str, + snapshot_id: &'a str, + metadata_root: &'a str, + projection_elapsed_micros: u64, + page_counter_scope: &'static str, + codec_radix_work: &'static str, + #[serde(flatten)] + work: &'a ProjectionWork, + message: &'static str, +} + +fn tagged_oid(oid: &ObjectHash, serializer: S) -> Result { + serializer.serialize_str(&oid.to_tagged_string()) +} + fn invalid() -> SnapshotError { SnapshotError::new( SnapshotErrorCode::IntegrityError, @@ -156,6 +195,37 @@ impl NativeResolveSource { } impl NativeProjectionObservation { + pub(crate) fn wire_record(&self) -> ProjectionWireRecord<'_> { + ProjectionWireRecord { + observation_revision: 1, + phase: "resolve_directory_projection", + source_domain: "native-git", + request_id: &self.request_id, + instance_id: self.source.instance.to_string(), + root_commit_oid: self.source.commit, + root_tree_oid: self.source.tree, + native_certificate_receipt_id: self.source.certificate_receipt_id, + native_writer_epoch: self.source.writer_epoch, + native_publication_sequence: self.source.publication_sequence, + scope: &self.scope, + schema_version: SCHEMA_VERSION, + metadata_codec: METADATA_CODEC, + materialization_policy: MATERIALIZATION_POLICY_GIT_RAW_V1, + fs_semantics: FS_SEMANTICS_LINUX_CODE_V1, + access_projection: ACCESS_PROJECTION_EXACT_FULL, + verification_revision: MST2_VERIFICATION_VERSION, + projection_revision: NATIVE_PROJECTION_REVISION, + namespace_view_id: &self.namespace_view_id, + snapshot_id: &self.snapshot_id, + metadata_root: &self.metadata_root, + projection_elapsed_micros: self.projection_elapsed_micros, + page_counter_scope: "returned-directory-root-pages", + codec_radix_work: "NOT_EXPOSED", + work: &self.work, + message: "native resolve directory projection succeeded", + } + } + pub(crate) fn emit(self) { tracing::debug!( target: "mst2::native_projection_observation", @@ -459,6 +529,7 @@ mod tests { Duration::from_micros(4), ) .unwrap(); + let wire = serde_json::to_value(observation.wire_record()).unwrap(); let events = Arc::new(Mutex::new(Vec::new())); tracing::subscriber::with_default(TraceCapture(events.clone()), || observation.emit()); let events = events.lock().unwrap(); @@ -506,6 +577,19 @@ mod tests { ]; assert_eq!(fields.len(), expected.len()); assert!(expected.iter().all(|field| fields.contains_key(*field))); + let wire = wire.as_object().unwrap(); + assert_eq!(wire.len(), expected.len()); + for key in expected { + let typed = &wire[key]; + let expected = if key == "scope" { + typed.to_string() + } else if let Some(text) = typed.as_str() { + text.into() + } else { + typed.to_string() + }; + assert_eq!(fields[key], expected, "typed observation diverged at {key}"); + } assert_eq!(fields["phase"], "resolve_directory_projection"); assert_eq!( fields["page_counter_scope"], diff --git a/src/ceres/snapshot/projection_writer.rs b/src/ceres/snapshot/projection_writer.rs new file mode 100644 index 00000000..a3d51852 --- /dev/null +++ b/src/ceres/snapshot/projection_writer.rs @@ -0,0 +1,506 @@ +//! Default-off typed projection observations. Buffer capacity is reserved +//! before serialization; logger filters and HTTP response bodies are unchanged. + +use std::{ + fs::{self, File, OpenOptions}, + io::{self, Write}, + path::{Path, PathBuf}, + sync::{ + Arc, Mutex, + atomic::{AtomicU8, AtomicU64, Ordering}, + mpsc::{self, Receiver, SyncSender, TrySendError}, + }, + thread::{self, JoinHandle}, + time::{Duration, Instant}, +}; + +use serde::Serialize; +use sha2::{Digest, Sha256}; +use uuid::Uuid; + +use super::projection_observation::NativeProjectionObservation; + +pub(crate) const RECORD_BYTES: usize = 32 * 1024; +pub(crate) const RECORD_LIMIT: usize = 64; +const STATUS_BYTES: usize = 4096; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +#[repr(u8)] +pub(crate) enum WriterFailure { + QueueSaturated = 1, + WriterUnavailable, + RecordTooLarge, + RecordLimitExceeded, + ObservationBindingRejected, + IoFailure, + WorkerPanic, + DrainTimeout, +} + +impl WriterFailure { + fn label(self) -> &'static str { + match self { + Self::QueueSaturated => "QUEUE_SATURATED", + Self::WriterUnavailable => "WRITER_UNAVAILABLE", + Self::RecordTooLarge => "RECORD_TOO_LARGE", + Self::RecordLimitExceeded => "RECORD_LIMIT_EXCEEDED", + Self::ObservationBindingRejected => "OBSERVATION_BINDING_REJECTED", + Self::IoFailure => "IO_FAILURE", + Self::WorkerPanic => "WORKER_PANIC", + Self::DrainTimeout => "DRAIN_TIMEOUT", + } + } +} + +struct Health { + first_error: AtomicU8, + accepted: AtomicU64, +} + +impl Health { + fn fail(&self, error: WriterFailure) { + if self + .first_error + .compare_exchange(0, error as u8, Ordering::SeqCst, Ordering::SeqCst) + .is_ok() + { + tracing::warn!(target: "mst2::projection_writer", failure_code = error.label(), "projection observation delivery failed"); + } + } +} + +// Exactly 64 preallocated payload slots exist, including the slot currently +// being serialized and the worker's slot. No serialize-to-Vec-then-check path. +struct RecordBuffer { + bytes: Box<[u8; CAPACITY]>, + len: usize, +} + +impl Write for RecordBuffer { + fn write(&mut self, bytes: &[u8]) -> io::Result { + let end = self + .len + .checked_add(bytes.len()) + .filter(|end| *end <= CAPACITY) + .ok_or_else(|| io::Error::other("projection record capacity exhausted"))?; + self.bytes[self.len..end].copy_from_slice(bytes); + self.len = end; + Ok(bytes.len()) + } + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } +} + +struct QueuedRecord { + sequence: u64, + buffer: RecordBuffer, +} +type Pool = Arc>>; +struct Producer { + sender: Option>, +} + +pub(crate) struct ProjectionObservationSink { + id: Uuid, + producer: Mutex, + pool: Pool, + health: Arc, + worker: Mutex>>, +} + +impl ProjectionObservationSink { + pub(crate) fn start(cache: &Path) -> io::Result> { + let id = Uuid::new_v4(); + let root = prepare_directory(cache, id)?; + let records = create_private(&root.join("records.jsonl"))?; + let health = Arc::new(Health { + first_error: AtomicU8::new(0), + accepted: AtomicU64::new(0), + }); + let pool = Arc::new(Mutex::new( + (0..RECORD_LIMIT) + .map(|_| RecordBuffer { + bytes: Box::new([0; RECORD_BYTES]), + len: 0, + }) + .collect::>(), + )); + let mut writer = FileWriter::new(root, records, id, health.clone()); + writer.status(false)?; + let (sender, receiver) = mpsc::sync_channel(RECORD_LIMIT); + let worker_pool = pool.clone(); + let worker_health = health.clone(); + let worker = thread::Builder::new() + .name("mst2-projection-writer".into()) + .spawn(move || { + let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + writer.run(receiver, worker_pool); + })); + if outcome.is_err() { + worker_health.fail(WriterFailure::WorkerPanic); + let _ = writer.status(true); + } + })?; + tracing::info!(target: "mst2::projection_writer", writer_revision = 1u16, sink_instance = %id, + payload_slots = RECORD_LIMIT, payload_slot_bytes = RECORD_BYTES, + "typed projection writer started"); + Ok(Arc::new(Self { + id, + producer: Mutex::new(Producer { + sender: Some(sender), + }), + pool, + health, + worker: Mutex::new(Some(worker)), + })) + } + + pub(crate) fn reject_binding(&self) { + self.health.fail(WriterFailure::ObservationBindingRejected); + } + + pub(crate) fn enqueue( + &self, + observation: &NativeProjectionObservation, + ) -> Result<(), WriterFailure> { + let result = self.enqueue_inner(observation); + if let Err(error) = result { + self.health.fail(error); + } + result + } + + fn enqueue_inner( + &self, + observation: &NativeProjectionObservation, + ) -> Result<(), WriterFailure> { + if self.health.first_error.load(Ordering::SeqCst) != 0 { + return Err(WriterFailure::WriterUnavailable); + } + let producer = self + .producer + .try_lock() + .map_err(|_| WriterFailure::QueueSaturated)?; + let sender = producer + .sender + .as_ref() + .ok_or(WriterFailure::WriterUnavailable)?; + let sequence = self.health.accepted.load(Ordering::SeqCst) + 1; + if sequence > RECORD_LIMIT as u64 { + return Err(WriterFailure::RecordLimitExceeded); + } + let mut buffer = self + .pool + .try_lock() + .map_err(|_| WriterFailure::QueueSaturated)? + .pop() + .ok_or(WriterFailure::QueueSaturated)?; + buffer.len = 0; + let serialized = (|| -> io::Result<()> { + write!( + buffer, + "{{\"writer_revision\":1,\"sink_instance\":\"{}\",\"record_sequence\":{},\"payload\":", + self.id, sequence + )?; + let begin = buffer.len; + serde_json::to_writer(&mut buffer, &observation.wire_record()) + .map_err(io::Error::other)?; + let digest = Sha256::digest(&buffer.bytes[begin..buffer.len]); + writeln!( + buffer, + ",\"payload_sha256\":\"sha256:{}\"}}", + hex::encode(digest) + ) + })(); + if serialized.is_err() { + self.return_slot(buffer); + return Err(WriterFailure::RecordTooLarge); + } + self.health.accepted.store(sequence, Ordering::SeqCst); + match sender.try_send(QueuedRecord { sequence, buffer }) { + Ok(()) => Ok(()), + Err(error) => { + self.health.accepted.store(sequence - 1, Ordering::SeqCst); + let (error, record) = match error { + TrySendError::Full(record) => (WriterFailure::QueueSaturated, record), + TrySendError::Disconnected(record) => { + (WriterFailure::WriterUnavailable, record) + } + }; + self.return_slot(record.buffer); + Err(error) + } + } + } + + fn return_slot(&self, mut buffer: RecordBuffer) { + buffer.len = 0; + if let Ok(mut pool) = self.pool.lock() { + pool.push(buffer); + } + } + + pub(crate) async fn shutdown(&self, deadline: Instant) -> Result<(), WriterFailure> { + self.producer + .lock() + .map_err(|_| WriterFailure::WriterUnavailable)? + .sender + .take(); + loop { + let finished = self + .worker + .lock() + .map_err(|_| WriterFailure::WriterUnavailable)? + .as_ref() + .is_none_or(JoinHandle::is_finished); + if finished { + break; + } + if Instant::now() >= deadline { + self.health.fail(WriterFailure::DrainTimeout); + return Err(WriterFailure::DrainTimeout); + } + tokio::time::sleep( + Duration::from_millis(10).min(deadline.saturating_duration_since(Instant::now())), + ) + .await; + } + if let Some(worker) = self + .worker + .lock() + .map_err(|_| WriterFailure::WriterUnavailable)? + .take() + && worker.join().is_err() + { + self.health.fail(WriterFailure::WorkerPanic); + } + if self.health.first_error.load(Ordering::SeqCst) != 0 { + return Err(WriterFailure::WriterUnavailable); + } + if Instant::now() >= deadline { + self.health.fail(WriterFailure::DrainTimeout); + return Err(WriterFailure::DrainTimeout); + } + Ok(()) + } + + #[cfg(test)] + pub(crate) fn test_directory(&self, cache: &Path) -> PathBuf { + cache + .join("logs/mst2-native-projection") + .join(self.id.to_string()) + } +} + +#[derive(Serialize)] +struct WriterStatus { + writer_revision: u16, + sink_instance: String, + accepted_records: u64, + written_sequence: u64, + written_records: u64, + written_bytes: u64, + rolling_sha256: String, + first_error_code: u8, + closed: bool, +} + +struct FileWriter { + root: PathBuf, + records: File, + id: Uuid, + health: Arc, + sequence: u64, + bytes: u64, + digest: Sha256, +} +impl FileWriter { + fn new(root: PathBuf, records: File, id: Uuid, health: Arc) -> Self { + Self { + root, + records, + id, + health, + sequence: 0, + bytes: 0, + digest: Sha256::new(), + } + } + fn status(&self, closed: bool) -> io::Result<()> { + let status = WriterStatus { + writer_revision: 1, + sink_instance: self.id.to_string(), + accepted_records: self.health.accepted.load(Ordering::SeqCst), + written_sequence: self.sequence, + written_records: self.sequence, + written_bytes: self.bytes, + rolling_sha256: format!("sha256:{}", hex::encode(self.digest.clone().finalize())), + first_error_code: self.health.first_error.load(Ordering::SeqCst), + closed, + }; + let mut buffer = RecordBuffer:: { + bytes: Box::new([0; STATUS_BYTES]), + len: 0, + }; + serde_json::to_writer(&mut buffer, &status).map_err(io::Error::other)?; + let temporary = self.root.join("status.tmp"); + let mut file = create_private(&temporary)?; + file.write_all(&buffer.bytes[..buffer.len])?; + file.sync_all()?; + fs::rename(temporary, self.root.join("status.json"))?; + sync_directory(&self.root) + } + fn run(&mut self, receiver: Receiver, pool: Pool) { + let mut reported_error = 0; + loop { + match receiver.recv_timeout(Duration::from_millis(100)) { + Ok(mut record) => { + let result = (|| -> io::Result<()> { + if record.sequence != self.sequence + 1 + || record.sequence > RECORD_LIMIT as u64 + { + return Err(io::Error::other("projection writer sequence mismatch")); + } + self.records + .write_all(&record.buffer.bytes[..record.buffer.len])?; + self.records.flush()?; + self.records.sync_all()?; + self.digest + .update(&record.buffer.bytes[..record.buffer.len]); + self.bytes += record.buffer.len as u64; + self.sequence = record.sequence; + self.status(false) + })(); + record.buffer.len = 0; + if let Ok(mut slots) = pool.lock() { + slots.push(record.buffer); + } + if result.is_err() { + self.health.fail(WriterFailure::IoFailure); + break; + } + } + Err(mpsc::RecvTimeoutError::Timeout) => {} + Err(mpsc::RecvTimeoutError::Disconnected) => break, + } + let error = self.health.first_error.load(Ordering::SeqCst); + if error != 0 && error != reported_error { + if self.status(false).is_err() { + self.health.fail(WriterFailure::IoFailure); + break; + } + reported_error = error; + } + } + if self.status(true).is_err() { + self.health.fail(WriterFailure::IoFailure); + } + } +} + +fn create_private(path: &Path) -> io::Result { + let mut options = OpenOptions::new(); + options.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + options.mode(0o600); + } + options.open(path) +} + +fn prepare_directory(cache: &Path, id: Uuid) -> io::Result { + prepare_directory_with(cache, id, sync_directory) +} + +fn real_directory(path: &Path) -> io::Result<()> { + if !fs::symlink_metadata(path)?.is_dir() { + return Err(io::Error::other( + "projection storage must be a real directory", + )); + } + Ok(()) +} + +fn create_directory(path: &Path, private: bool) -> io::Result<()> { + #[cfg(unix)] + let mut directory = fs::DirBuilder::new(); + #[cfg(not(unix))] + let directory = fs::DirBuilder::new(); + #[cfg(unix)] + { + use std::os::unix::fs::DirBuilderExt; + directory.mode(if private { 0o700 } else { 0o755 }); + } + #[cfg(not(unix))] + let _ = private; + directory.create(path) +} + +fn private_directory(path: &Path) -> io::Result<()> { + real_directory(path)?; + #[cfg(unix)] + { + use std::os::unix::fs::{MetadataExt, PermissionsExt}; + let metadata = fs::symlink_metadata(path)?; + // SAFETY: geteuid reads the process identity and has no pointer arguments. + let owner = unsafe { libc::geteuid() }; + if metadata.uid() != owner || metadata.permissions().mode() & 0o7777 != 0o700 { + return Err(io::Error::other( + "projection directory must be privately owned", + )); + } + } + Ok(()) +} + +fn prepare_directory_with( + cache: &Path, + id: Uuid, + sync: impl Fn(&Path) -> io::Result<()>, +) -> io::Result { + // The actual configured cache is an existing anchor. No arbitrary output + // path or recursively created ancestor is accepted by this writer. + real_directory(cache)?; + let logs = cache.join("logs"); + match create_directory(&logs, false) { + Ok(()) => {} + Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {} + Err(error) => return Err(error), + } + real_directory(&logs)?; + sync(cache)?; + let parent = logs.join("mst2-native-projection"); + match create_directory(&parent, true) { + Ok(()) => {} + Err(error) if error.kind() == io::ErrorKind::AlreadyExists => {} + Err(error) => return Err(error), + } + private_directory(&parent)?; + sync(&logs)?; + let root = parent.join(id.to_string()); + create_directory(&root, true)?; + private_directory(&root)?; + sync(&parent)?; + Ok(root) +} + +fn sync_directory(path: &Path) -> io::Result<()> { + #[cfg(unix)] + { + File::open(path)?.sync_all() + } + #[cfg(not(unix))] + { + let _ = path; + Err(io::Error::new( + io::ErrorKind::Unsupported, + "durable projection writer requires directory fsync", + )) + } +} + +#[cfg(test)] +#[path = "projection_writer_tests.rs"] +mod tests; diff --git a/src/ceres/snapshot/projection_writer_tests.rs b/src/ceres/snapshot/projection_writer_tests.rs new file mode 100644 index 00000000..2de737a2 --- /dev/null +++ b/src/ceres/snapshot/projection_writer_tests.rs @@ -0,0 +1,358 @@ +use git_internal::hash::{HashKind, ObjectHash}; +use mst2_codec::descriptor::ServingDescriptor; + +use super::*; +use crate::ceres::snapshot::projection_observation::{ + NativeResolveSource, ProjectionWork, ResolvedProjection, +}; + +fn observation(scope: &str) -> NativeProjectionObservation { + let commit = ObjectHash::from_hex_for_kind(HashKind::Sha256, &"a".repeat(64)).unwrap(); + let tree = ObjectHash::from_hex_for_kind(HashKind::Sha256, &"b".repeat(64)).unwrap(); + let instance = Uuid::new_v4(); + let mut namespace = Sha256::new(); + namespace.update(b"mega.mst2.namespaceview\0"); + namespace.update(commit.to_string().as_bytes()); + let descriptor = ServingDescriptor { + instance_uuid: *instance.as_bytes(), + namespace_view_id: namespace.finalize().into(), + scope: scope.into(), + metadata_root: [3; 32], + }; + let snapshot = format!("sha256:{}", hex::encode(descriptor.snapshot_id().unwrap())); + let root = format!("sha256:{}", hex::encode(descriptor.metadata_root)); + NativeResolveSource::capture(&instance.to_string(), commit, tree, Some(1), 2, 3) + .unwrap() + .observe( + ResolvedProjection { + descriptor: &descriptor, + snapshot_id: &snapshot, + metadata_root: &root, + context_commit: &commit.to_string(), + context_root_tree: &tree.to_string(), + fixed_root_tree: tree, + requested_scope: scope, + request_id: "writer:test:a1", + }, + ProjectionWork::default(), + Duration::from_micros(4), + ) + .unwrap() +} + +fn idle_sink() -> (ProjectionObservationSink, Receiver) { + let (sender, receiver) = mpsc::sync_channel(RECORD_LIMIT); + let sink = ProjectionObservationSink { + id: Uuid::new_v4(), + producer: Mutex::new(Producer { + sender: Some(sender), + }), + pool: Arc::new(Mutex::new( + (0..RECORD_LIMIT) + .map(|_| RecordBuffer { + bytes: Box::new([0; RECORD_BYTES]), + len: 0, + }) + .collect(), + )), + health: Arc::new(Health { + first_error: AtomicU8::new(0), + accepted: AtomicU64::new(0), + }), + worker: Mutex::new(None), + }; + (sink, receiver) +} + +#[test] +fn bounded_slots_precede_serialization_and_saturation_is_sticky_with_no_extra_allocation() { + let (sink, receiver) = idle_sink(); + let observation = observation("/project"); + for _ in 0..RECORD_LIMIT { + sink.enqueue(&observation).unwrap(); + } + assert_eq!(sink.pool.lock().unwrap().len(), 0); + assert_eq!( + sink.health.accepted.load(Ordering::SeqCst), + RECORD_LIMIT as u64 + ); + assert_eq!( + sink.enqueue(&observation), + Err(WriterFailure::RecordLimitExceeded) + ); + assert_eq!( + sink.health.first_error.load(Ordering::SeqCst), + WriterFailure::RecordLimitExceeded as u8 + ); + let record = receiver.try_recv().unwrap(); + sink.return_slot(record.buffer); + assert_eq!( + sink.enqueue(&observation), + Err(WriterFailure::WriterUnavailable) + ); + assert_eq!(sink.pool.lock().unwrap().len(), 1); +} + +#[test] +fn pool_saturation_and_disconnection_have_typed_failure_and_return_the_actual_slot() { + let (sink, receiver) = idle_sink(); + let observation = observation("/project"); + let held = sink.pool.lock().unwrap().drain(..).collect::>(); + assert_eq!( + sink.enqueue(&observation), + Err(WriterFailure::QueueSaturated) + ); + assert_eq!(held.len(), RECORD_LIMIT); + assert_eq!(sink.health.accepted.load(Ordering::SeqCst), 0); + let (sink, receiver2) = idle_sink(); + drop(receiver2); + assert_eq!( + sink.enqueue(&observation), + Err(WriterFailure::WriterUnavailable) + ); + assert_eq!(sink.pool.lock().unwrap().len(), RECORD_LIMIT); + assert_eq!(sink.health.accepted.load(Ordering::SeqCst), 0); + drop(receiver); +} + +#[test] +fn typed_wire_contract_has_exact_38_fields_and_rejects_oversize_without_growing_buffer() { + let (sink, receiver) = idle_sink(); + let captured = observation("/project\"quoted"); + sink.enqueue(&captured).unwrap(); + let record = receiver.try_recv().unwrap(); + assert_eq!(record.buffer.bytes.len(), RECORD_BYTES); + let envelope: serde_json::Value = + serde_json::from_slice(&record.buffer.bytes[..record.buffer.len]).unwrap(); + let payload = &envelope["payload"]; + assert_eq!(payload.as_object().unwrap().len(), 38); + assert_eq!(payload["scope"], "/project\"quoted"); + assert_eq!(payload["request_id"], "writer:test:a1"); + assert_eq!(payload["projection_elapsed_micros"], 4); + assert_eq!(payload["codec_radix_work"], "NOT_EXPOSED"); + let mut small = RecordBuffer::<64> { + bytes: Box::new([0; 64]), + len: 0, + }; + // Serialize a genuine validated observation into an insufficient reserved + // slot, rather than bypassing the descriptor's component length checks. + assert!(serde_json::to_writer(&mut small, &captured.wire_record()).is_err()); + assert!(small.len <= 64); + assert_eq!(small.bytes.len(), 64); + let len = small.len; + assert!(small.write_all(&[0; 65]).is_err()); + assert_eq!(small.len, len); + let wire = &record.buffer.bytes[..record.buffer.len]; + let prefix = b"\"payload\":"; + let begin = wire + .windows(prefix.len()) + .position(|part| part == prefix) + .unwrap() + + prefix.len(); + let suffix = b",\"payload_sha256\":"; + let end = wire + .windows(suffix.len()) + .position(|part| part == suffix) + .unwrap(); + assert_eq!( + envelope["payload_sha256"], + format!("sha256:{}", hex::encode(Sha256::digest(&wire[begin..end]))) + ); +} + +#[cfg(unix)] +#[tokio::test] +async fn actual_writer_fsync_ack_tracks_exact_bytes_and_closed_drain() { + let temp = tempfile::tempdir().unwrap(); + let sink = ProjectionObservationSink::start(temp.path()).unwrap(); + sink.enqueue(&observation("/project")).unwrap(); + sink.shutdown(Instant::now() + Duration::from_secs(5)) + .await + .unwrap(); + let root = temp + .path() + .join("logs/mst2-native-projection") + .join(sink.id.to_string()); + let wire = fs::read(root.join("records.jsonl")).unwrap(); + let status: serde_json::Value = + serde_json::from_slice(&fs::read(root.join("status.json")).unwrap()).unwrap(); + assert_eq!(status["written_bytes"], wire.len()); + assert_eq!(status["written_sequence"], 1); + assert_eq!(status["accepted_records"], 1); + assert_eq!(status["first_error_code"], 0); + assert_eq!(status["closed"], true); + assert_eq!( + status["rolling_sha256"], + format!("sha256:{}", hex::encode(Sha256::digest(&wire))) + ); + let record: serde_json::Value = serde_json::from_slice(&wire).unwrap(); + assert_eq!(record["payload"].as_object().unwrap().len(), 38); + assert_eq!(sink.pool.lock().unwrap().len(), RECORD_LIMIT); + use std::os::unix::fs::PermissionsExt; + assert_eq!( + fs::metadata(&root).unwrap().permissions().mode() & 0o777, + 0o700 + ); + assert_eq!( + fs::metadata(root.parent().unwrap()) + .unwrap() + .permissions() + .mode() + & 0o7777, + 0o700 + ); + for path in [root.join("records.jsonl"), root.join("status.json")] { + assert_eq!( + fs::metadata(path).unwrap().permissions().mode() & 0o777, + 0o600 + ); + } +} + +#[tokio::test] +async fn expired_drain_deadline_cannot_report_a_successful_capture() { + let (sink, _) = idle_sink(); + assert_eq!( + sink.shutdown(Instant::now()).await, + Err(WriterFailure::DrainTimeout) + ); + assert_eq!( + sink.health.first_error.load(Ordering::SeqCst), + WriterFailure::DrainTimeout as u8 + ); + assert_eq!( + sink.enqueue(&observation("/project")), + Err(WriterFailure::WriterUnavailable) + ); +} + +#[cfg(unix)] +#[tokio::test] +async fn actual_status_write_failure_and_binding_rejection_cannot_drain_successfully() { + let temp = tempfile::tempdir().unwrap(); + let sink = ProjectionObservationSink::start(temp.path()).unwrap(); + let root = temp + .path() + .join("logs/mst2-native-projection") + .join(sink.id.to_string()); + fs::write(root.join("status.tmp"), b"blocked exclusive temporary").unwrap(); + sink.enqueue(&observation("/project")).unwrap(); + assert!( + sink.shutdown(Instant::now() + Duration::from_secs(5)) + .await + .is_err() + ); + let status: serde_json::Value = + serde_json::from_slice(&fs::read(root.join("status.json")).unwrap()).unwrap(); + assert_eq!( + status["written_sequence"], 0, + "no durable sidecar acknowledgement for this record" + ); + assert_eq!( + sink.health.first_error.load(Ordering::SeqCst), + WriterFailure::IoFailure as u8 + ); + let temp = tempfile::tempdir().unwrap(); + let sink = ProjectionObservationSink::start(temp.path()).unwrap(); + sink.reject_binding(); + assert!( + sink.shutdown(Instant::now() + Duration::from_secs(5)) + .await + .is_err() + ); + let root = temp + .path() + .join("logs/mst2-native-projection") + .join(sink.id.to_string()); + let status: serde_json::Value = + serde_json::from_slice(&fs::read(root.join("status.json")).unwrap()).unwrap(); + assert_eq!( + status["first_error_code"], + WriterFailure::ObservationBindingRejected as u8 + ); +} + +#[cfg(unix)] +#[test] +fn actual_writer_rejects_a_symlinked_output_directory_before_creating_records() { + use std::os::unix::fs::symlink; + let temp = tempfile::tempdir().unwrap(); + let target = tempfile::tempdir().unwrap(); + fs::create_dir(temp.path().join("logs")).unwrap(); + symlink( + target.path(), + temp.path().join("logs/mst2-native-projection"), + ) + .unwrap(); + assert!(ProjectionObservationSink::start(temp.path()).is_err()); + assert_eq!(target.path().read_dir().unwrap().count(), 0); + let temp = tempfile::tempdir().unwrap(); + symlink(target.path(), temp.path().join("logs")).unwrap(); + assert!(ProjectionObservationSink::start(temp.path()).is_err()); + assert_eq!(target.path().read_dir().unwrap().count(), 0); +} + +#[cfg(unix)] +#[test] +fn directory_chain_is_synced_from_existing_anchor_before_any_record_can_be_created() { + use std::cell::RefCell; + let temp = tempfile::tempdir().unwrap(); + let calls = RefCell::new(Vec::new()); + let root = prepare_directory_with(temp.path(), Uuid::new_v4(), |path| { + sync_directory(path)?; + calls.borrow_mut().push(path.to_path_buf()); + Ok(()) + }) + .unwrap(); + assert_eq!( + *calls.borrow(), + [ + temp.path().to_path_buf(), + temp.path().join("logs"), + temp.path().join("logs/mst2-native-projection") + ] + ); + assert_eq!(root.read_dir().unwrap().count(), 0); + for fail_at in 0..3 { + let temp = tempfile::tempdir().unwrap(); + let id = Uuid::new_v4(); + let calls = RefCell::new(0); + assert!( + prepare_directory_with(temp.path(), id, |path| { + let current = *calls.borrow(); + *calls.borrow_mut() += 1; + if current == fail_at { + Err(io::Error::other("injected directory durability failure")) + } else { + sync_directory(path) + } + }) + .is_err() + ); + assert_eq!(*calls.borrow(), fail_at + 1); + let root = temp + .path() + .join("logs/mst2-native-projection") + .join(id.to_string()); + assert!(!root.join("records.jsonl").exists()); + assert!(!root.join("status.json").exists()); + } +} + +#[cfg(unix)] +#[test] +fn missing_or_symlink_anchor_and_nonprivate_existing_parent_fail_before_records() { + use std::os::unix::fs::{PermissionsExt, symlink}; + let temp = tempfile::tempdir().unwrap(); + assert!(ProjectionObservationSink::start(&temp.path().join("missing")).is_err()); + assert!(!temp.path().join("missing").exists()); + let alias = temp.path().join("alias"); + symlink(temp.path(), &alias).unwrap(); + assert!(ProjectionObservationSink::start(&alias).is_err()); + assert!(!temp.path().join("logs").exists()); + let parent = temp.path().join("logs/mst2-native-projection"); + fs::create_dir_all(&parent).unwrap(); + fs::set_permissions(&parent, fs::Permissions::from_mode(0o755)).unwrap(); + assert!(ProjectionObservationSink::start(temp.path()).is_err()); + assert_eq!(parent.read_dir().unwrap().count(), 0); +} diff --git a/src/commands/service/mod.rs b/src/commands/service/mod.rs index 666dfb7c..939674c3 100644 --- a/src/commands/service/mod.rs +++ b/src/commands/service/mod.rs @@ -96,7 +96,7 @@ pub(crate) async fn exec(ctx: CommandContext, args: &ArgMatches) -> MegaResult { .await }; - match setup.await { + let result = match setup.await { Err(error) => cleanup_tail(Err(error), emitter, None, signal_forwarder).await, Ok(reload_watcher) => { let service = context.clone(); @@ -117,7 +117,16 @@ pub(crate) async fn exec(ctx: CommandContext, args: &ArgMatches) -> MegaResult { ) .await } + }; + if let Some(sink) = &context.storage.projection_observation_sink { + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + if sink.shutdown(deadline).await.is_err() && result.is_ok() { + return Err(MegaError::Other( + "typed projection writer did not drain successfully".into(), + )); + } } + result } /// Runs the service body under the unified cleanup tail (WH-13). diff --git a/src/config/model.rs b/src/config/model.rs index 65ca6756..2075eca5 100644 --- a/src/config/model.rs +++ b/src/config/model.rs @@ -878,6 +878,9 @@ impl Default for LFSSshConfig { /// MST/2 snapshot feature flags (default off; spec 00 ยง6 rollout). #[derive(Serialize, Deserialize, Debug, Clone, Default)] pub struct Mst2Config { + /// Startup-only typed projection diagnostics, independent of log.level. + #[serde(default)] + pub projection_observation_enabled: bool, /// Master switch for the `/api/v2/snapshots` surface. #[serde(default)] pub enabled: bool, diff --git a/src/config/reload.rs b/src/config/reload.rs index f117460b..b097de47 100644 --- a/src/config/reload.rs +++ b/src/config/reload.rs @@ -617,6 +617,12 @@ fn collect_static_restart_fields( candidate: &Config, report: &mut ConfigReloadReport, ) { + if current.mst2.projection_observation_enabled != candidate.mst2.projection_observation_enabled + { + report + .restart_required_fields + .push("mst2.projection_observation_enabled"); + } if current.base_dir != candidate.base_dir { report.restart_required_fields.push("base_dir"); } @@ -981,6 +987,33 @@ mod tests { .expect("test runtime") } + #[test] + fn projection_writer_toggle_requires_restart_and_preserves_live_snapshot() { + let temp = tempfile::tempdir().unwrap(); + let mut current = isolated_config(temp.path().join("config")); + current.mst2.enabled = true; + current.mst2.publication_enabled = true; + current.mst2.instance_uuid = Some("12345678-1234-4234-9234-123456789abc".into()); + let handle = ConfigHandle::new(current); + let before = handle.snapshot().unwrap(); + let mut candidate = before.as_ref().clone(); + candidate.mst2.projection_observation_enabled = true; + let report = handle.reload(candidate).unwrap(); + assert_eq!( + report.restart_required_fields, + vec!["mst2.projection_observation_enabled"] + ); + assert!(!report.applied()); + assert!(Arc::ptr_eq(&before, &handle.snapshot().unwrap())); + assert!( + !handle + .snapshot() + .unwrap() + .mst2 + .projection_observation_enabled + ); + } + #[test] fn reload_applies_log_fields_and_preserves_restart_required_database_fields() { let temp_dir = tempfile::tempdir().expect("temp dir"); diff --git a/src/config/validate.rs b/src/config/validate.rs index 7804ff4d..11fd9350 100644 --- a/src/config/validate.rs +++ b/src/config/validate.rs @@ -150,6 +150,11 @@ impl Config { } fn validate_mst2_config(config: &Mst2Config) -> Result<(), MegaError> { + if config.projection_observation_enabled && (!config.enabled || !config.publication_enabled) { + return Err(MegaError::Other( + "mst2.projection_observation_enabled requires mst2.enabled and mst2.publication_enabled".into(), + )); + } if !config.enabled { return Ok(()); } @@ -2129,6 +2134,7 @@ pub(crate) fn known_fields(path: &str) -> Option<&'static [&'static str]> { "enabled", "instance_uuid", "publication_enabled", + "projection_observation_enabled", "auth_token", ]), "buck" => Some(&[ @@ -2235,6 +2241,31 @@ mod tests { isolated_config(std::env::temp_dir().join("mega2-config-validate-tests")) } + #[test] + fn typed_projection_writer_is_default_off_and_requires_native_publication() { + assert!(!Mst2Config::default().projection_observation_enabled); + let loaded: Mst2Config = toml::from_str("enabled = false").unwrap(); + assert!(!loaded.projection_observation_enabled); + let mut config = valid_config(); + config.mst2.projection_observation_enabled = true; + config.mst2.instance_uuid = Some("12345678-1234-4234-9234-123456789abc".into()); + for (enabled, publication) in [(false, false), (false, true), (true, false)] { + config.mst2.enabled = enabled; + config.mst2.publication_enabled = publication; + assert!( + config + .validate() + .unwrap_err() + .to_string() + .contains("projection_observation_enabled") + ); + } + config.mst2.enabled = true; + config.mst2.publication_enabled = true; + config.validate().unwrap(); + assert!(is_known_field_path("mst2.projection_observation_enabled")); + } + #[test] fn config_validate_mst2_checks_enabled_identity_and_configured_token() { let mut config = valid_config(); diff --git a/src/context/mod.rs b/src/context/mod.rs index dd97ed20..a4c05d22 100644 --- a/src/context/mod.rs +++ b/src/context/mod.rs @@ -295,6 +295,22 @@ impl AppContext { } }; + let mut storage = storage; + if config.mst2.projection_observation_enabled { + match crate::ceres::snapshot::projection_writer::ProjectionObservationSink::start( + &crate::config::mega_cache(), + ) { + Ok(sink) => storage.projection_observation_sink = Some(sink), + Err(_) => { + notification_shutdown.cancel(); + storage.storage_event_emitter.shutdown().await; + return Err(MegaError::Other( + "typed projection writer failed startup".into(), + )); + } + } + } + Ok(Self { storage, vault, diff --git a/src/jupiter/service/native_publication_push_tests.rs b/src/jupiter/service/native_publication_push_tests.rs index 2bd0f211..89290152 100644 --- a/src/jupiter/service/native_publication_push_tests.rs +++ b/src/jupiter/service/native_publication_push_tests.rs @@ -225,6 +225,10 @@ async fn http_resolve(state: crate::api::MonoApiServiceState,scope:&str) -> (u16 #[tokio::test] async fn http_resolve_captured_before_commit_never_mixes_new_sequence_with_old_descriptor() { let (_temp,storage,tip,path)=native_fixture().await; + #[cfg(unix)] + let writer=crate::ceres::snapshot::projection_writer::ProjectionObservationSink::start(_temp.path()).unwrap(); + #[cfg(unix)] + let storage={ let mut storage=storage; storage.projection_observation_sink=Some(writer.clone()); storage }; let (first,payload)=save_same_tree_commit(&storage,&tip).await; let id=wh03_enqueue_push(&storage,&path,&tip.id.to_string(),&first,&payload).await; assert!(matches!(wh03_exec(&storage,id).await,ExecuteOutcome::Done{..})); @@ -283,4 +287,53 @@ async fn http_resolve_captured_before_commit_never_mixes_new_sequence_with_old_d let ((status,_),failed_observations)=crate::ceres::snapshot::projection_observation::with_observations(http_resolve(state,"/absent-observation-scope")).await; assert_eq!(status,404); assert!(failed_observations.is_empty(),"failed projection emitted success observation"); + #[cfg(unix)] + { + writer.shutdown(std::time::Instant::now()+Duration::from_secs(5)).await.unwrap(); + let directory=writer.test_directory(_temp.path()); + let records=std::fs::read_to_string(directory.join("records.jsonl")).unwrap(); + let records:Vec=records.lines().map(|line|serde_json::from_str(line).unwrap()).collect(); + assert_eq!(records.len(),4,"failed resolve must not be a successful observation"); + for (record,observation) in [(&records[0],&first_observations[0]),(&records[2],&delayed_observations[0]),(&records[3],&next_observations[0])] { + let typed=serde_json::to_value(observation.wire_record()).unwrap(); + assert_eq!(record["payload"],typed,"writer must preserve the full validated operation tuple"); + assert_eq!(record["payload"].as_object().unwrap().len(),38); + } + assert_ne!(records[0]["payload"]["request_id"],records[2]["payload"]["request_id"]); + let status:serde_json::Value=serde_json::from_slice(&std::fs::read(directory.join("status.json")).unwrap()).unwrap(); + assert_eq!(status["closed"],true); + assert_eq!(status["accepted_records"],4); + assert_eq!(status["written_records"],4); + assert_eq!(status["first_error_code"],0); + } +} + +#[cfg(unix)] +#[tokio::test] +async fn rejected_native_observation_source_does_not_change_ready_resolve_and_is_a_sticky_writer_failure() { + let (temp,mut storage,tip,path)=native_fixture().await; + let (first,payload)=save_same_tree_commit(&storage,&tip).await; + let id=wh03_enqueue_push(&storage,&path,&tip.id.to_string(),&first,&payload).await; + assert!(matches!(wh03_exec(&storage,id).await,ExecuteOutcome::Done{..})); + let head=storage.mono_storage().read_native_publication_head(NATIVE_INSTANCE).await.unwrap(); + assert_eq!(head.token.sequence,1); + assert!(head.token.certificate.unwrap()>0); + let writer=crate::ceres::snapshot::projection_writer::ProjectionObservationSink::start(temp.path()).unwrap(); + storage.projection_observation_sink=Some(writer.clone()); + let state=api_state(storage).await; + // Corrupt only the observation's captured certificate. The actual native + // head is READY; INITIALIZING heads remain correctly rejected with 503. + let ((status,response),observations)=crate::ceres::snapshot::projection_observation::with_observations( + crate::api::router::snapshot_router::with_rejected_native_observation_source(http_resolve(state,"/")) + ).await; + assert_eq!(status,200,"{response}"); + assert_eq!(response["publication_sequence"],"1"); + assert!(observations.is_empty()); + assert!(writer.shutdown(std::time::Instant::now()+Duration::from_secs(5)).await.is_err()); + let directory=writer.test_directory(temp.path()); + assert!(std::fs::read(directory.join("records.jsonl")).unwrap().is_empty()); + let status:serde_json::Value=serde_json::from_slice(&std::fs::read(directory.join("status.json")).unwrap()).unwrap(); + assert_eq!(status["first_error_code"],crate::ceres::snapshot::projection_writer::WriterFailure::ObservationBindingRejected as u8); + assert_eq!(status["accepted_records"],0); + assert_eq!(status["closed"],true); } diff --git a/src/jupiter/storage/mod.rs b/src/jupiter/storage/mod.rs index 8ce75f10..c1ab0504 100644 --- a/src/jupiter/storage/mod.rs +++ b/src/jupiter/storage/mod.rs @@ -130,6 +130,8 @@ pub struct Storage { /// Derived native projection memoization, scoped to this storage assembly. /// Clones share it; independent databases/backends never share entries. pub(crate) native_projection_cache: Arc, + pub(crate) projection_observation_sink: + Option>, pub cl_service: CLService, pub push_queue_service: PushQueueService, pub artifact_service: ArtifactService, @@ -288,6 +290,7 @@ impl Storage { Ok(Storage { app_service: app_service.into(), native_projection_cache: Arc::default(), + projection_observation_sink: None, config_handle, config, cl_service: CLService::new(base.clone()), @@ -571,6 +574,7 @@ impl Storage { Storage { app_service, native_projection_cache: Arc::default(), + projection_observation_sink: None, // app_service: AppService::mock(), cl_service: CLService::mock(), push_queue_service: PushQueueService::new( diff --git a/src/jupiter/tests.rs b/src/jupiter/tests.rs index 07546cfc..f7d92d59 100644 --- a/src/jupiter/tests.rs +++ b/src/jupiter/tests.rs @@ -375,6 +375,7 @@ pub async fn test_storage_with_config(temp_dir: impl AsRef, config: Config Storage { app_service: Arc::new(svc), native_projection_cache: Arc::default(), + projection_observation_sink: None, cl_service: CLService::mock(), push_queue_service: PushQueueService::new( base.clone(),