laminar_core/checkpoint/checkpoint_manifest/
sink_artifacts.rs1use sha2::{Digest, Sha256};
4
5use super::{ByteRange, CheckpointManifest, PREPARED_SINK_DESCRIPTOR_VERSION};
6
7#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
9#[serde(deny_unknown_fields)]
10pub struct PreparedSinkDescriptor {
11 pub sink_name: String,
13 pub format_version: u16,
15 pub payload: Option<ByteRange>,
17 pub sha256: String,
19}
20
21#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq, Eq)]
23#[serde(deny_unknown_fields)]
24pub struct PreparedSinkArtifactIntent {
25 pub sink_name: String,
27 pub format_version: u16,
29 pub payload: Option<ByteRange>,
31 pub sha256: String,
33}
34
35#[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)] fn 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}