Skip to main content

laminar_core/checkpoint/checkpoint_manifest/
sink_artifacts.rs

1//! Checkpoint-owned sink intents and phase-one descriptor validation.
2
3use sha2::{Digest, Sha256};
4
5use super::{ByteRange, CheckpointManifest, PREPARED_SINK_DESCRIPTOR_VERSION};
6
7/// Opaque phase-one sink descriptor stored in the current node data object.
8#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
9#[serde(deny_unknown_fields)]
10pub struct PreparedSinkDescriptor {
11    /// Stable registered sink name.
12    pub sink_name: String,
13    /// Runtime descriptor-envelope version.
14    pub format_version: u16,
15    /// `None` and an explicit empty range have different connector semantics.
16    pub payload: Option<ByteRange>,
17    /// Presence-domain-separated SHA-256 of the optional payload.
18    pub sha256: String,
19}
20
21/// Recovery evidence persisted before a checkpoint-committable sink may write.
22#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
23#[serde(deny_unknown_fields)]
24pub struct PreparedSinkArtifactIntent {
25    /// Stable registered sink name.
26    pub sink_name: String,
27    /// Runtime intent-envelope version.
28    pub format_version: u16,
29    /// Connector intent bytes in the current node data object.
30    pub payload: Option<ByteRange>,
31    /// Presence-domain-separated SHA-256 of the optional payload.
32    pub sha256: String,
33}
34
35/// SHA-256 for one optional sink artifact intent.
36#[must_use]
37pub fn checkpoint_artifact_intent_sha256(payload: Option<&[u8]>) -> String {
38    let mut digest = Sha256::new();
39    digest.update(b"laminardb/checkpoint-sink-artifact-intent/v1\0");
40    match payload {
41        Some(payload) => {
42            digest.update([1]);
43            digest.update(payload);
44        }
45        None => digest.update([0]),
46    }
47    format!("{:x}", digest.finalize())
48}
49
50pub(super) fn validate_sink_artifacts(
51    manifest: &CheckpointManifest,
52    current_ranges: &mut Vec<(ByteRange, String)>,
53    error: &mut impl FnMut(String),
54) {
55    if !manifest.sink_artifact_intents.is_empty() {
56        validate_intent_roster(manifest, error);
57        for intent in &manifest.sink_artifact_intents {
58            validate_entry(
59                "sink artifact intent",
60                &intent.sink_name,
61                intent.format_version,
62                intent.payload,
63                &intent.sha256,
64                checkpoint_artifact_intent_sha256,
65                manifest,
66                current_ranges,
67                error,
68            );
69        }
70    }
71    if !manifest
72        .prepared_sinks
73        .windows(2)
74        .all(|pair| pair[0].sink_name < pair[1].sink_name)
75    {
76        error("prepared_sinks must be strictly ordered by sink_name".into());
77    }
78    for sink in &manifest.prepared_sinks {
79        validate_entry(
80            "prepared sink",
81            &sink.sink_name,
82            sink.format_version,
83            sink.payload,
84            &sink.sha256,
85            super::checkpoint_descriptor_sha256,
86            manifest,
87            current_ranges,
88            error,
89        );
90    }
91}
92
93fn validate_intent_roster(manifest: &CheckpointManifest, error: &mut impl FnMut(String)) {
94    if !manifest
95        .sink_artifact_intents
96        .windows(2)
97        .all(|pair| pair[0].sink_name < pair[1].sink_name)
98    {
99        error("sink_artifact_intents must be strictly ordered by sink_name".into());
100    }
101    let intent_names = manifest
102        .sink_artifact_intents
103        .iter()
104        .map(|intent| intent.sink_name.as_str());
105    let prepared_names = manifest
106        .prepared_sinks
107        .iter()
108        .map(|sink| sink.sink_name.as_str());
109    if !intent_names.eq(prepared_names) {
110        error("sink artifact intents must name every prepared sink exactly once".into());
111    }
112}
113
114#[allow(clippy::too_many_arguments)] // One generic envelope check keeps both domains identical.
115fn validate_entry(
116    label: &str,
117    sink_name: &str,
118    format_version: u16,
119    payload: Option<ByteRange>,
120    sha256: &str,
121    absent_digest: fn(Option<&[u8]>) -> String,
122    manifest: &CheckpointManifest,
123    current_ranges: &mut Vec<(ByteRange, String)>,
124    error: &mut impl FnMut(String),
125) {
126    if sink_name.is_empty() {
127        error(format!("{label} name must not be empty"));
128    }
129    if manifest
130        .sink_names
131        .binary_search_by(|candidate| candidate.as_str().cmp(sink_name))
132        .is_err()
133    {
134        error(format!("{label} '{sink_name}' is not in sink_names"));
135    }
136    if format_version != PREPARED_SINK_DESCRIPTOR_VERSION {
137        error(format!(
138            "{label} '{sink_name}' format_version must be {PREPARED_SINK_DESCRIPTOR_VERSION}"
139        ));
140    }
141    if !super::is_sha256(sha256) {
142        error(format!(
143            "{label} '{sink_name}' digest must be lowercase SHA-256"
144        ));
145    }
146    match payload {
147        Some(range) => {
148            super::validate_range(range, Some(manifest.node_data.object_length), label, error);
149            current_ranges.push((range, format!("{label} '{sink_name}'")));
150        }
151        None if sha256 != absent_digest(None) => error(format!(
152            "{label} '{sink_name}' without a payload has the wrong domain-separated digest"
153        )),
154        None => {}
155    }
156}