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
37 changes: 34 additions & 3 deletions src/api/router/snapshot_router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,14 @@ tokio::task_local! {
#[cfg(test)]
tokio::task_local! {
static NATIVE_RESOLVE_BARRIERS: (std::sync::Arc<tokio::sync::Barrier>, std::sync::Arc<tokio::sync::Barrier>);
static REJECT_NATIVE_OBSERVATION_SOURCE: bool;
}

#[cfg(test)]
pub(crate) async fn with_rejected_native_observation_source<F: std::future::Future>(
future: F,
) -> F::Output {
REJECT_NATIVE_OBSERVATION_SOURCE.scope(true, future).await
}

#[cfg(test)]
Expand Down Expand Up @@ -460,18 +468,31 @@ 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,
)
});
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
}
Expand Down Expand Up @@ -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");
}
}
}

Expand Down
1 change: 1 addition & 0 deletions src/ceres/snapshot/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
86 changes: 85 additions & 1 deletion src/ceres/snapshot/projection_observation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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.
Expand Down Expand Up @@ -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<S: serde::Serializer>(oid: &ObjectHash, serializer: S) -> Result<S::Ok, S::Error> {
serializer.serialize_str(&oid.to_tagged_string())
}

fn invalid() -> SnapshotError {
SnapshotError::new(
SnapshotErrorCode::IntegrityError,
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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"],
Expand Down
Loading
Loading