Skip to main content

mkit_server/relay/
enqueue.rs

1//! Bounded source batches for relay rows produced outside an advance batch.
2
3use crate::store::{
4    Batch, BatchOutcome, MAX_BATCH_BYTES, MAX_BATCH_OPS, NamespaceStore, Partition, Precondition,
5    StoreCapabilities, StoreError, Value, Write,
6    codec::{self, RelayV1},
7    keys,
8    outbox::MAX_RELAY_PUTS,
9};
10
11/// Snapshot used to plan a chain. `deadline_ms` includes the deployment's
12/// clock-skew margin; for a D34 ref shard, `source_lease` is mandatory.
13#[derive(Debug, Clone)]
14pub struct RelayEnqueueSnapshot {
15    /// Observed outbox sequence row, if any.
16    pub sequence: Option<Value>,
17    /// Observed source epoch lease on a D34 ref shard.
18    pub source_lease: Option<Value>,
19    /// Latest time each batch may commit on the backend clock.
20    pub deadline_ms: u64,
21}
22
23fn guard(key: crate::Key, value: Option<&Value>) -> Precondition {
24    match value {
25        Some(value) => Precondition::Equals(key, value.clone()),
26        None => Precondition::Absent(key),
27    }
28}
29
30fn lease_lost() -> StoreError {
31    StoreError::unavailable(std::io::Error::other("source epoch lease lost; retry"))
32}
33
34fn batch_for(
35    prior: Option<&Value>,
36    lease: Option<&Value>,
37    seq: u64,
38    row: &RelayV1,
39    now_ms: u64,
40    deadline_ms: u64,
41) -> Result<Batch, StoreError> {
42    if row.puts.len().saturating_add(row.deletes.len()) > MAX_RELAY_PUTS {
43        return Err(StoreError::Invalid(
44            "relay row exceeds MAX_RELAY_PUTS".into(),
45        ));
46    }
47    let value = codec::encode_relay(row)?;
48    let mut batch = Batch::new()
49        .require(Precondition::NotAfter(deadline_ms))
50        .require(guard(keys::outbox_sequence(), prior));
51    if let Some(lease) = lease {
52        batch = batch.require(Precondition::Equals(keys::epoch_lease(), lease.clone()));
53    }
54    Ok(batch
55        .put(keys::outbox_sequence(), codec::encode_u64(seq))
56        .put(
57            keys::timer(now_ms, crate::timers::registry::kinds::RELAY.get(), b""),
58            Value::default(),
59        )
60        .put(keys::relay(seq), value))
61}
62
63/// Plan chained source batches. Each has a deadline, an `os` guard/put,
64/// one relay kick, optional source lease guard, and bounded row puts.
65/// The caller commits in order; after an `os` CAS loss it re-plans only
66/// the uncommitted suffix from a fresh snapshot.
67pub fn enqueue_relay_rows(
68    os: &RelayEnqueueSnapshot,
69    source: &Partition,
70    rows: &[RelayV1],
71    now_ms: u64,
72) -> Result<Vec<Batch>, StoreError> {
73    if matches!(source, Partition::Ref { .. }) && os.source_lease.is_none() {
74        return Err(StoreError::Corrupt(
75            "D34 relay batch lacks source epoch lease".into(),
76        ));
77    }
78    if rows.is_empty() {
79        return Ok(Vec::new());
80    }
81    if let Some(raw) = &os.source_lease {
82        let lease = codec::decode_epoch_lease(raw)?;
83        if now_ms >= lease.expires_at_ms || os.deadline_ms >= lease.expires_at_ms {
84            return Err(lease_lost());
85        }
86    }
87    let mut seq = os
88        .sequence
89        .as_ref()
90        .map(codec::decode_u64)
91        .transpose()?
92        .unwrap_or(0);
93    let mut prior = os.sequence.clone();
94    let mut batches = Vec::new();
95    let mut index = 0;
96    while index < rows.len() {
97        seq = seq
98            .checked_add(1)
99            .ok_or_else(|| StoreError::Corrupt("outbox sequence overflow".into()))?;
100        let mut batch = batch_for(
101            prior.as_ref(),
102            os.source_lease.as_ref(),
103            seq,
104            &rows[index],
105            now_ms,
106            os.deadline_ms,
107        )?;
108        batch.validate(&StoreCapabilities::full())?;
109        index += 1;
110        while index < rows.len() && batch.preconditions.len() + batch.writes.len() < MAX_BATCH_OPS {
111            let next = seq
112                .checked_add(1)
113                .ok_or_else(|| StoreError::Corrupt("outbox sequence overflow".into()))?;
114            if rows[index]
115                .puts
116                .len()
117                .saturating_add(rows[index].deletes.len())
118                > MAX_RELAY_PUTS
119            {
120                return Err(StoreError::Invalid(
121                    "relay row exceeds MAX_RELAY_PUTS".into(),
122                ));
123            }
124            let value = codec::encode_relay(&rows[index])?;
125            let key = keys::relay(next);
126            let mut candidate = batch.clone();
127            candidate.writes.push(Write::Put(key, value));
128            if candidate.validate(&StoreCapabilities::full()).is_err() {
129                break;
130            }
131            batch = candidate;
132            seq = next;
133            index += 1;
134        }
135        let next_os = codec::encode_u64(seq);
136        // The sequence put is the first write, after all preconditions.
137        batch.writes[0] = Write::Put(keys::outbox_sequence(), next_os.clone());
138        debug_assert!(batch.preconditions.len() + batch.writes.len() <= MAX_BATCH_OPS);
139        debug_assert!(
140            batch
141                .writes
142                .iter()
143                .map(|w| match w {
144                    Write::Put(k, v) => k.as_bytes().len() + v.as_bytes().len(),
145                    Write::Delete(k) => k.as_bytes().len(),
146                })
147                .sum::<usize>()
148                <= MAX_BATCH_BYTES
149        );
150        batches.push(batch);
151        prior = Some(next_os);
152    }
153    Ok(batches)
154}
155
156/// Commit a relay-row chain, re-reading `os` and re-planning the remaining
157/// rows when another writer wins. The caller supplies its current source
158/// lease and deadline; a stale lease fails closed.
159pub async fn commit_relay_rows<S: NamespaceStore>(
160    store: &S,
161    source: &Partition,
162    rows: &[RelayV1],
163    now_ms: u64,
164    deadline_ms: u64,
165    source_lease: Option<&Value>,
166) -> Result<(), StoreError> {
167    if matches!(source, Partition::Ref { .. }) && source_lease.is_none() {
168        return Err(StoreError::Corrupt(
169            "D34 relay batch lacks source epoch lease".into(),
170        ));
171    }
172    let mut done = 0;
173    let mut losses = 0;
174    while done < rows.len() {
175        let snapshot = RelayEnqueueSnapshot {
176            sequence: store.get(source, &keys::outbox_sequence()).await?,
177            source_lease: source_lease.cloned(),
178            deadline_ms,
179        };
180        let batches = enqueue_relay_rows(&snapshot, source, &rows[done..], now_ms)?;
181        let mut lost = false;
182        for batch in batches {
183            let count = batch.writes.iter().filter(|w| matches!(w, Write::Put(key, _) if matches!(keys::parse(key), Some(keys::ParsedKey::Relay(_))))).count();
184            match store.apply(source, batch).await? {
185                BatchOutcome::Committed => done += count,
186                BatchOutcome::PreconditionFailed { index, .. }
187                    if source_lease.is_some() && index == 2 =>
188                {
189                    return Err(lease_lost());
190                }
191                BatchOutcome::PreconditionFailed { .. } => {
192                    lost = true;
193                    break;
194                }
195                BatchOutcome::DeadlinePassed { .. } => {
196                    return Err(StoreError::unavailable(std::io::Error::other(
197                        "relay enqueue deadline passed",
198                    )));
199                }
200            }
201        }
202        if lost {
203            losses += 1;
204            if losses == 8 {
205                return Err(StoreError::unavailable(std::io::Error::other(
206                    "relay enqueue contention",
207                )));
208            }
209        }
210    }
211    Ok(())
212}
213
214/// Whether all source rows through `seq` have left its outbox. It uses the
215/// source state scan; target watermark delivery precedes source cleanup.
216pub async fn relay_delivered_through<S: NamespaceStore>(
217    store: &S,
218    source: &Partition,
219    seq: u64,
220) -> Result<bool, StoreError> {
221    let allocated = store
222        .get(source, &keys::outbox_sequence())
223        .await?
224        .as_ref()
225        .map(codec::decode_u64)
226        .transpose()?
227        .unwrap_or(0);
228    if allocated < seq {
229        return Ok(false);
230    }
231    let (start, end) = keys::class_range(keys::TAG_RELAY);
232    let page = store.scan(source, &start, &end, None, 1).await?;
233    match page.entries.first().and_then(|(key, _)| keys::parse(key)) {
234        Some(keys::ParsedKey::Relay(first)) => Ok(first > seq),
235        None if page.entries.is_empty() => Ok(true),
236        _ => Err(StoreError::Corrupt("invalid relay head".into())),
237    }
238}
239
240#[cfg(all(test, feature = "memory"))]
241mod tests {
242    use super::*;
243    use crate::memory::MemoryKv;
244    use crate::repo::{NamespaceKey, RepoName};
245    use futures_executor::block_on;
246
247    fn source() -> Partition {
248        Partition::Namespace(NamespaceKey::deployment_default())
249    }
250    fn ref_source() -> Partition {
251        Partition::Ref {
252            ns: NamespaceKey::deployment_default(),
253            repo: RepoName::new("one").expect("valid repository name"),
254            shard_ref: "refs/heads/main".into(),
255        }
256    }
257    fn row(n: u16) -> RelayV1 {
258        RelayV1 {
259            at_ms: 1_000,
260            target: source(),
261            puts: vec![(
262                crate::Key::new(n.to_be_bytes().to_vec()),
263                Value::new(vec![u8::try_from(n % 256).unwrap_or(0)]),
264            )],
265            deletes: Vec::new(),
266        }
267    }
268
269    #[test]
270    fn chains_batches_and_detects_delivery() {
271        let store = MemoryKv::default();
272        let source = source();
273        let rows: Vec<_> = (0..120).map(row).collect();
274        let snapshot = RelayEnqueueSnapshot {
275            sequence: None,
276            source_lease: None,
277            deadline_ms: u64::MAX,
278        };
279        let batches = enqueue_relay_rows(&snapshot, &source, &rows, 1_000).unwrap();
280        assert!(batches.len() >= 2);
281        for batch in batches {
282            assert_eq!(
283                block_on(store.apply(&source, batch)).unwrap(),
284                BatchOutcome::Committed
285            );
286        }
287        assert!(!block_on(relay_delivered_through(&store, &source, 120)).unwrap());
288        let deletes = (1..=120).fold(Batch::new(), |batch, seq| batch.delete(keys::relay(seq)));
289        // The store cap forbids one large delete; remove in bounded chunks.
290        for chunk in deletes.writes.chunks(80) {
291            let batch = Batch {
292                preconditions: Vec::new(),
293                writes: chunk.to_vec(),
294            };
295            assert_eq!(
296                block_on(store.apply(&source, batch)).unwrap(),
297                BatchOutcome::Committed
298            );
299        }
300        assert!(block_on(relay_delivered_through(&store, &source, 120)).unwrap());
301        assert!(!block_on(relay_delivered_through(&store, &source, 121)).unwrap());
302    }
303
304    #[test]
305    fn stale_sequence_replans_uncommitted_rows() {
306        let store = MemoryKv::default();
307        let source = source();
308        let stale = RelayEnqueueSnapshot {
309            sequence: None,
310            source_lease: None,
311            deadline_ms: u64::MAX,
312        };
313        let batch = enqueue_relay_rows(&stale, &source, &[row(1)], 1_000)
314            .unwrap()
315            .remove(0);
316        assert_eq!(
317            block_on(store.apply(
318                &source,
319                Batch::new().put(keys::outbox_sequence(), codec::encode_u64(7))
320            ))
321            .unwrap(),
322            BatchOutcome::Committed
323        );
324        assert!(matches!(
325            block_on(store.apply(&source, batch)).unwrap(),
326            BatchOutcome::PreconditionFailed { .. }
327        ));
328        block_on(commit_relay_rows(
329            &store,
330            &source,
331            &[row(1)],
332            1_000,
333            u64::MAX,
334            None,
335        ))
336        .unwrap();
337        assert!(
338            block_on(store.get(&source, &keys::relay(8)))
339                .unwrap()
340                .is_some()
341        );
342    }
343
344    #[test]
345    fn expired_source_lease_cannot_enqueue() {
346        let lease = codec::encode_epoch_lease(&codec::EpochLease {
347            authority_ready: None,
348            authority_generation: None,
349            epoch: 1,
350            expires_at_ms: 1_010,
351            config_version: 1,
352        });
353        let snapshot = RelayEnqueueSnapshot {
354            sequence: None,
355            source_lease: Some(lease),
356            deadline_ms: 1_010,
357        };
358        assert!(enqueue_relay_rows(&snapshot, &source(), &[row(1)], 1_000).is_err());
359    }
360
361    #[test]
362    fn ref_source_requires_lease_and_a_lost_lease_is_not_contention() {
363        let source = ref_source();
364        let snapshot = RelayEnqueueSnapshot {
365            sequence: None,
366            source_lease: None,
367            deadline_ms: 2_000,
368        };
369        assert!(matches!(
370            enqueue_relay_rows(&snapshot, &source, &[row(1)], 1_000),
371            Err(StoreError::Corrupt(_))
372        ));
373        let stale = codec::encode_epoch_lease(&codec::EpochLease {
374            authority_ready: None,
375            authority_generation: None,
376            epoch: 1,
377            expires_at_ms: 10_000,
378            config_version: 1,
379        });
380        let current = codec::encode_epoch_lease(&codec::EpochLease {
381            authority_ready: None,
382            authority_generation: None,
383            epoch: 2,
384            expires_at_ms: 10_000,
385            config_version: 1,
386        });
387        let store = MemoryKv::with_clock(std::sync::Arc::new(crate::ManualClock::new(1_000)));
388        block_on(store.apply(&source, Batch::new().put(keys::epoch_lease(), current))).unwrap();
389        let error = block_on(commit_relay_rows(
390            &store,
391            &source,
392            &[row(1)],
393            1_000,
394            2_000,
395            Some(&stale),
396        ))
397        .unwrap_err();
398        assert!(error.to_string().contains("source epoch lease lost"));
399    }
400}