1use std::collections::BTreeMap;
3
4use mkit_core::hash::{hash, to_hex, to_hex_bytes};
5use serde::{Deserialize, Serialize};
6use serde_json::json;
7
8use super::{
9 Response, auth,
10 ledger::{Head, audit_entry, decode_head, encode, guarded, head_key},
11};
12use crate::{
13 Batch, BoxFuture, Key, NamespaceStore, Partition, Precondition, StoreError, Value, Write,
14 purge,
15 relay::{RelayEnqueueSnapshot, RelayHook, enqueue_relay_rows},
16 store::codec::RelayV1,
17};
18
19#[derive(Serialize, Deserialize, PartialEq, Eq)]
20#[serde(rename_all = "camelCase", deny_unknown_fields)]
21struct Event {
22 source_identity: String,
23 operation_id: String,
24 recorded_at_ms: u64,
25 request: purge::Request,
26}
27fn invalid(message: &'static str) -> StoreError {
28 StoreError::Invalid(message.into())
29}
30fn event_key(source: &str, purge_id: &str) -> Key {
31 Key::new(
32 format!(
33 "ai\0{}",
34 to_hex(&hash(format!("{source}\0{purge_id}").as_bytes()))
35 )
36 .into_bytes(),
37 )
38}
39fn events(writes: &[Write]) -> Vec<(Key, Value)> {
40 writes
41 .iter()
42 .filter_map(|w| match w {
43 Write::Put(k, v) if k.as_bytes().starts_with(b"ai\0") => Some((k.clone(), v.clone())),
44 _ => None,
45 })
46 .collect()
47}
48
49#[derive(Clone, Debug)]
51pub struct SystemAudit<S> {
52 _store: S,
53 root: Partition,
54}
55impl<S> SystemAudit<S> {
56 pub fn new(store: S, root: Partition) -> Self {
58 Self {
59 _store: store,
60 root,
61 }
62 }
63 pub fn relay_row(
69 &self,
70 source: &Partition,
71 request: &purge::Request,
72 operation_id: &str,
73 now_ms: u64,
74 ) -> Result<RelayV1, StoreError> {
75 request.validate()?;
76 if !auth::identifier(operation_id, 128, true) || request.trigger == purge::Trigger::Manual {
77 return Err(invalid("invalid automatic audit identity"));
78 }
79 let source_identity = to_hex_bytes(&source.encode()?);
80 let event = Event {
81 source_identity: source_identity.clone(),
82 operation_id: operation_id.into(),
83 recorded_at_ms: now_ms,
84 request: request.clone(),
85 };
86 let value =
87 encode(&event).map_err(|_| invalid("automatic audit event exceeds storage bound"))?;
88 Ok(RelayV1 {
89 at_ms: now_ms,
90 target: self.root.clone(),
91 puts: vec![(event_key(&source_identity, &request.purge_id), value)],
92 deletes: Vec::new(),
93 })
94 }
95}
96impl<S: NamespaceStore> purge::AutomaticAudit for SystemAudit<S> {
97 fn plan<'a>(
98 &'a self,
99 partition: &'a Partition,
100 request: &'a purge::Request,
101 operation_id: &'a str,
102 now_ms: u64,
103 snapshot: RelayEnqueueSnapshot,
104 ) -> BoxFuture<'a, Result<Batch, StoreError>> {
105 Box::pin(async move {
106 let row = self.relay_row(partition, request, operation_id, now_ms)?;
107 enqueue_relay_rows(&snapshot, partition, &[row], now_ms)?
108 .pop()
109 .ok_or_else(|| invalid("missing automatic audit relay"))
110 })
111 }
112}
113
114pub fn extend_audit_batch(
119 partition: &Partition,
120 batch: &mut Batch,
121 mut get: impl FnMut(&Key) -> Result<Option<Value>, StoreError>,
122) -> Result<(), StoreError> {
123 let events = events(&batch.writes);
124 if events.is_empty() {
125 return Ok(());
126 }
127 let root = crate::NamespaceKey::deployment_default();
128 if partition != &Partition::Namespace(root.clone())
129 && partition != &Partition::Coordinator(root)
130 {
131 return Err(invalid("automatic audit target is not deployment root"));
132 }
133 let old = get(&head_key())?;
134 let mut head = decode_head(old.as_ref()).map_err(|_| invalid("corrupt audit head"))?;
135 let mut additions = guarded(Batch::new(), head_key(), old);
136 let mut receipts = BTreeMap::new();
137 for (key, value) in events {
138 let event: Event = serde_json::from_slice(value.as_bytes())
139 .map_err(|_| invalid("invalid automatic audit event"))?;
140 event.request.validate()?;
141 if !auth::identifier(&event.operation_id, 128, true)
142 || event.request.trigger == purge::Trigger::Manual
143 || key != event_key(&event.source_identity, &event.request.purge_id)
144 {
145 return Err(invalid("invalid automatic audit identity"));
146 }
147 if event.source_identity.len() > crate::MAX_KEY_BYTES * 2
148 || !event.source_identity.len().is_multiple_of(2)
149 {
150 return Err(invalid("invalid audit source identity"));
151 }
152 let source = event
153 .source_identity
154 .as_bytes()
155 .chunks_exact(2)
156 .map(|pair| {
157 std::str::from_utf8(pair)
158 .ok()
159 .and_then(|text| u8::from_str_radix(text, 16).ok())
160 .ok_or_else(|| invalid("invalid audit source identity"))
161 })
162 .collect::<Result<Vec<_>, _>>()?;
163 let source = Partition::decode(&source)?;
165 if to_hex_bytes(&source.encode()?) != event.source_identity {
166 return Err(invalid("invalid audit source identity"));
167 }
168 if let Some(existing) = receipts.get(&key) {
169 if existing != &value {
170 return Err(invalid("automatic purge identity reused"));
171 }
172 continue;
173 }
174 let observed = get(&key)?;
175 if let Some(existing) = &observed {
176 if existing != &value {
177 return Err(invalid("automatic purge identity reused"));
178 }
179 } else {
180 let actor = match event.request.trigger {
181 purge::Trigger::Takedown | purge::Trigger::Suspension => "system:inspector",
182 purge::Trigger::LeaseDeletion => "system:timer",
183 _ => "system:relay",
184 };
185 let details = json!({"purgeId":event.request.purge_id,
186 "sourcePartitionHash":to_hex(&hash(&source.encode()?)),"trigger":event.request.trigger})
187 .to_string();
188 let (entry, next) = audit_entry(
189 &head,
190 actor,
191 &format!("{actor}/cache-purge"),
192 "",
193 "",
194 &event.operation_id,
195 "",
196 &[event.request.scope().to_owned()],
197 &Response::json(&json!({})),
198 &details,
199 event.recorded_at_ms,
200 )
201 .map_err(|_| invalid("invalid automatic audit entry"))?;
202 additions = additions.put(
203 Key::new([b"ae\0".as_slice(), &next.seq.to_be_bytes()].concat()),
204 encode(&entry).map_err(|_| invalid("automatic audit entry too large"))?,
205 );
206 head = next;
207 }
208 additions = guarded(additions, key.clone(), observed);
209 receipts.insert(key, value);
210 }
211 additions = additions.put(
212 head_key(),
213 encode(&head).map_err(|_| invalid("invalid audit head"))?,
214 );
215 batch.preconditions.extend(additions.preconditions);
216 batch.writes.extend(additions.writes);
217 Ok(())
218}
219
220#[derive(Clone, Debug)]
222pub struct AuditRelayHook<S> {
223 store: S,
224 root: Partition,
225}
226impl<S> AuditRelayHook<S> {
227 pub fn new(store: S, root: Partition) -> Self {
229 Self { store, root }
230 }
231}
232impl<S: NamespaceStore> RelayHook for AuditRelayHook<S> {
233 fn before_apply<'a>(
234 &'a self,
235 target: &'a Partition,
236 rows: &'a [(u64, RelayV1)],
237 pre: &'a mut Vec<Precondition>,
238 writes: &'a mut Vec<Write>,
239 ) -> BoxFuture<'a, Result<(), StoreError>> {
240 Box::pin(async move {
241 if target != &self.root
242 || !rows.iter().any(|(_, r)| {
243 r.puts
244 .iter()
245 .any(|(k, _)| k.as_bytes().starts_with(b"ai\0"))
246 })
247 {
248 return Ok(());
249 }
250 let mut keys = vec![head_key()];
251 keys.extend(events(writes).into_iter().map(|(k, _)| k));
252 let values = self.store.get_many(target, &keys).await?;
253 let snapshot: BTreeMap<_, _> = keys.into_iter().zip(values).collect();
254 let mut batch = Batch {
255 preconditions: pre.clone(),
256 writes: writes.clone(),
257 };
258 extend_audit_batch(target, &mut batch, |k| {
259 snapshot
260 .get(k)
261 .cloned()
262 .ok_or_else(|| invalid("missing audit snapshot"))
263 })?;
264 *pre = batch.preconditions;
265 *writes = batch.writes;
266 Ok(())
267 })
268 }
269}
270
271#[derive(Clone, Debug)]
273pub struct AuditReserveHook {
274 root: Partition,
275}
276impl AuditReserveHook {
277 #[must_use]
279 pub fn new(root: Partition) -> Self {
280 Self { root }
281 }
282}
283impl RelayHook for AuditReserveHook {
284 fn reserved_ops(&self, target: &Partition, rows: &[(u64, RelayV1)]) -> usize {
285 let n = rows
286 .iter()
287 .flat_map(|(_, r)| &r.puts)
288 .filter(|(k, _)| k.as_bytes().starts_with(b"ai\0"))
289 .count();
290 if target == &self.root && n > 0 {
291 n.saturating_mul(2).saturating_add(2)
292 } else {
293 0
294 }
295 }
296 fn before_apply<'a>(
297 &'a self,
298 target: &'a Partition,
299 rows: &'a [(u64, RelayV1)],
300 pre: &'a mut Vec<Precondition>,
301 writes: &'a mut Vec<Write>,
302 ) -> BoxFuture<'a, Result<(), StoreError>> {
303 Box::pin(async move {
304 if target != &self.root {
305 return Ok(());
306 }
307 if writes
308 .iter()
309 .filter(|write| {
310 matches!(write,
311 Write::Put(key, _) if key.as_bytes().starts_with(b"ai\0"))
312 })
313 .take(crate::MAX_BATCH_OPS + 1)
314 .count()
315 > crate::MAX_BATCH_OPS
316 {
317 return Err(invalid("invalid automatic audit entry count"));
318 }
319 let receipts: BTreeMap<_, _> = events(writes).into_iter().collect();
320 if receipts.is_empty() {
321 return Ok(());
322 }
323 let count = u64::try_from(receipts.len())
324 .map_err(|_| invalid("invalid automatic audit entry count"))?;
325 let synthetic = encode(&Head {
326 seq: u64::MAX
327 .checked_sub(count)
328 .ok_or_else(|| invalid("invalid automatic audit entry count"))?,
329 hash: "0".repeat(64),
330 })
331 .map_err(|_| invalid("invalid audit head"))?;
332 let mut estimate = Batch {
333 preconditions: pre.clone(),
334 writes: writes.clone(),
335 };
336 extend_audit_batch(target, &mut estimate, |key| {
339 Ok((key == &head_key()).then(|| synthetic.clone()))
340 })?;
341 let bytes = estimate
342 .preconditions
343 .iter()
344 .map(|pre| match pre {
345 Precondition::Absent(k) | Precondition::Present(k) => k.as_bytes().len(),
346 Precondition::Equals(k, v) => k.as_bytes().len() + v.as_bytes().len(),
347 Precondition::NotAfter(_) => 0,
348 })
349 .chain(estimate.writes.iter().map(|write| match write {
350 Write::Put(k, v) => k.as_bytes().len() + v.as_bytes().len(),
351 Write::Delete(k) => k.as_bytes().len(),
352 }))
353 .sum::<usize>()
354 .saturating_add(crate::MAX_VALUE_BYTES - synthetic.as_bytes().len())
357 .saturating_add(receipts.values().map(|v| v.as_bytes().len()).sum::<usize>());
358 let validation = estimate.validate(&crate::StoreCapabilities::full());
359 let capacity = match validation {
360 Ok(()) => bytes > crate::MAX_BATCH_BYTES,
361 Err(StoreError::Invalid(message))
362 if matches!(
363 message.as_ref(),
364 "batch exceeds MAX_BATCH_OPS" | "batch exceeds MAX_BATCH_BYTES"
365 ) =>
366 {
367 true
368 }
369 Err(error) => return Err(error),
370 };
371 if capacity && rows.len() > 1 {
374 return Err(invalid(crate::relay::AUDIT_CAPACITY));
375 }
376 Ok(())
377 })
378 }
379}