Skip to main content

mkit_server/relay/
hook.rs

1//! Target-batch extension point for later index and content consumers.
2
3use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
4use crate::store::{Key, Partition, Precondition, StoreError, Value, Write, codec::RelayV1};
5
6pub(crate) const AUDIT_CAPACITY: &str = "automatic audit combined batch capacity";
7
8/// Runs before each target apply attempt, including retries on contention or
9/// shrinking a combined group to fit hook additions within the store limits.
10/// Added effects must fit the batch limits and tolerate repeated delivery.
11/// An error leaves this target's rows queued and does not block other targets.
12pub trait RelayHook: MaybeSend + MaybeSync {
13    /// Additional raw observations, batched with the target watermark read.
14    /// The delivery engine deduplicates and bounds this declaration before IO.
15    fn read_keys(
16        &self,
17        _target: &Partition,
18        _rows: &[(u64, RelayV1)],
19    ) -> Result<Vec<Key>, StoreError> {
20        Ok(Vec::new())
21    }
22
23    /// Extend using the declared snapshot, without hidden IO. Older hooks
24    /// retain their original callback; production content hooks override this.
25    fn before_apply_observed<'a>(
26        &'a self,
27        target: &'a Partition,
28        rows: &'a [(u64, RelayV1)],
29        _observed: &'a [(Key, Option<Value>)],
30        pre: &'a mut Vec<Precondition>,
31        writes: &'a mut Vec<Write>,
32    ) -> BoxFuture<'a, Result<(), StoreError>> {
33        self.before_apply(target, rows, pre, writes)
34    }
35
36    /// Target-local extensions may reserve operations before a remote apply.
37    /// The relay shrinks groups until their base effects and this reserve fit.
38    fn reserved_ops(&self, _target: &Partition, _rows: &[(u64, RelayV1)]) -> usize {
39        0
40    }
41
42    /// Extend the atomic target batch, or fail delivery for this target.
43    fn before_apply<'a>(
44        &'a self,
45        target: &'a Partition,
46        rows: &'a [(u64, RelayV1)],
47        pre: &'a mut Vec<Precondition>,
48        writes: &'a mut Vec<Write>,
49    ) -> BoxFuture<'a, Result<(), StoreError>>;
50}
51
52/// Default hook: adds nothing.
53#[derive(Debug, Default, Clone, Copy)]
54pub struct NoHook;
55impl RelayHook for NoHook {
56    fn before_apply<'a>(
57        &'a self,
58        _target: &'a Partition,
59        _rows: &'a [(u64, RelayV1)],
60        _pre: &'a mut Vec<Precondition>,
61        _writes: &'a mut Vec<Write>,
62    ) -> BoxFuture<'a, Result<(), StoreError>> {
63        Box::pin(async { Ok(()) })
64    }
65}