Skip to main content

mkit_server/admin/
automatic.rs

1//! Automatic audit events ride the existing source relay and target transaction.
2use 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/// Source-local audit enqueue planner. No target effect happens before commit.
50#[derive(Clone, Debug)]
51pub struct SystemAudit<S> {
52    _store: S,
53    root: Partition,
54}
55impl<S> SystemAudit<S> {
56    /// Use the same root as the deployment admin engine and relay audit hook.
57    pub fn new(store: S, root: Partition) -> Self {
58        Self {
59            _store: store,
60            root,
61        }
62    }
63    /// Produce an event for callers already allocating source relay sequences.
64    /// Add this row to the same `OutboxBuilder`/`enqueue_relay_rows` allocation as
65    /// other effects; do not merge independently allocated `os` batches.
66    /// # Errors
67    /// Invalid automatic operation, selectors, source identity, or storage bounds.
68    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
114/// Extend an existing relay target transaction using reads from that transaction.
115/// Dedup receipts, gapless chain entries, and relay watermarks commit atomically.
116/// # Errors
117/// Malformed events, changed receipts, or corrupt/unavailable audit metadata.
118pub 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        // Re-encoding rejects noncanonical partition bytes.
164        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/// Native relay hook: extend the actual atomic target apply after bounded reads.
221#[derive(Clone, Debug)]
222pub struct AuditRelayHook<S> {
223    store: S,
224    root: Partition,
225}
226impl<S> AuditRelayHook<S> {
227    /// The target metadata store and canonical admin root.
228    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/// Worker relay hook reserves target-local SQL extension space without DO reads.
272#[derive(Clone, Debug)]
273pub struct AuditReserveHook {
274    root: Partition,
275}
276impl AuditReserveHook {
277    /// The canonical admin root extended inside the target SQL transaction.
278    #[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            // Validate identities with the existing extension, using only owned local data.
337            // Missing receipts maximize appends; existing receipts add Equals values instead.
338            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                // Accepted heads need not have canonical JSON spelling. Reserve their raw
355                // storage bound numerically, without allocating a maximum-sized dummy.
356                .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            // Conservative estimates only shrink groups. A single row reaches unchanged
372            // target SQL validation, which knows the actual head and receipt state.
373            if capacity && rows.len() > 1 {
374                return Err(invalid(crate::relay::AUDIT_CAPACITY));
375            }
376            Ok(())
377        })
378    }
379}