Skip to main content

mkit_server/relay/
deliver.rs

1//! Ordered target batches, guarded watermarks, and bounded source cleanup.
2
3use std::collections::BTreeSet;
4
5use super::{NoHook, RELAY_LAG_BOUND_MS, RelayBudget, RelayHook};
6use crate::rt::BoxFuture;
7use crate::store::{
8    Batch, BatchOutcome, Key, NamespaceStore, Partition, Precondition, StoreCapabilities,
9    StoreError, Value, Write,
10    codec::{self, MAX_BLOCKED_TARGETS, RelayScanV1, RelayV1},
11    keys,
12};
13use crate::telemetry::{
14    METRIC_RELAY_BACKLOG_ROWS, METRIC_RELAY_LAG_EXCEEDED, Metrics, NoopMetrics,
15};
16use crate::timers::{
17    DueTimer, Fired, RETRY_BACKOFF_MS, TimerCtx, TimerHandler, TimerKind, registry::kinds,
18};
19
20// Keep encoded pages and decoded groups comfortably below a Worker isolate's
21// 128 MiB limit, even when rows approach MAX_VALUE_BYTES. The row/target budget
22// remains an upper bound; leftover rows schedule another tick.
23const MAX_FIRE_BYTES: usize = 4 * 1024 * 1024;
24const SCAN_PAGE_ROWS: u32 = 64;
25// At most 4.5 MiB raw snapshot values even for corrupt maximum-sized rows.
26// Together with source pages and JSON/JS copies this leaves hook headroom.
27const MAX_HOOK_READ_KEYS: usize = 9;
28
29type QueuedRow = (u64, RelayV1, Key, Value);
30type TargetRows = (Partition, Vec<QueuedRow>);
31type InspectedRow = (u64, Partition, Key, Value);
32
33struct ScanWindow {
34    rows: Vec<InspectedRow>,
35    groups: Vec<TargetRows>,
36    corrupt: bool,
37    exhausted: bool,
38    lag_exceeded: bool,
39}
40
41struct Dispatch {
42    delivered: BTreeSet<u64>,
43    block: BTreeSet<Partition>,
44    pause_at: Option<u64>,
45    overflow_at: Option<u64>,
46}
47
48struct TargetResult {
49    completed: usize,
50    failed: bool,
51    selected: bool,
52    calls: u32,
53}
54
55/// Pushes a source's queued rows to a separately supplied target store.
56/// Distinct producers may upsert the same key when its value is identical.
57/// A key that is ever relay-deleted MUST have exactly one producer, so its
58/// source sequence orders every update and delete.
59/// Target rh rows never expire.
60#[derive(Debug)]
61pub struct RelayHandler<T, H = NoHook> {
62    /// Target store (may be a clone of the source store on native).
63    pub target: T,
64    /// Atomic target-batch extension.
65    pub hook: H,
66    /// Per-fire work cap.
67    pub budget: RelayBudget,
68}
69
70impl<S: NamespaceStore, T: NamespaceStore, H: RelayHook> TimerHandler<S> for RelayHandler<T, H> {
71    fn kind(&self) -> TimerKind {
72        kinds::RELAY
73    }
74    fn fire<'a>(
75        &'a self,
76        ctx: &'a TimerCtx<'a, S>,
77        timer: &'a DueTimer,
78    ) -> BoxFuture<'a, Result<Fired, StoreError>> {
79        Box::pin(self.deliver_with_metrics(ctx, timer, &NoopMetrics))
80    }
81}
82
83#[cfg(feature = "test-faults")]
84async fn apply_relay_delay<S: NamespaceStore>(
85    ctx: &TimerCtx<'_, S>,
86    timer: &DueTimer,
87) -> Result<Option<Fired>, StoreError> {
88    let marker = crate::pipeline::faults::relay_delay_key();
89    let Some(value) = ctx.store.get(ctx.partition, &marker).await? else {
90        return Ok(None);
91    };
92    let until = codec::decode_u64(&value)?;
93    if ctx.now_ms < until {
94        return Ok(Some(Fired::Reschedule {
95            due_at_ms: until.max(timer.due_at_ms.saturating_add(1)),
96            value: timer.value.clone(),
97            batch: Batch::new(),
98        }));
99    }
100    Ok(
101        match ctx
102            .store
103            .apply(
104                ctx.partition,
105                Batch::new()
106                    .require(Precondition::Equals(marker.clone(), value))
107                    .delete(marker),
108            )
109            .await?
110        {
111            BatchOutcome::Committed => None,
112            BatchOutcome::PreconditionFailed { .. } | BatchOutcome::DeadlinePassed { .. } => {
113                Some(Fired::Retry)
114            }
115        },
116    )
117}
118
119impl<T: NamespaceStore, H: RelayHook> RelayHandler<T, H> {
120    /// Deliver one relay fire while recording source backlog and lag.
121    pub async fn deliver_with_metrics<S: NamespaceStore>(
122        &self,
123        ctx: &TimerCtx<'_, S>,
124        timer: &DueTimer,
125        metrics: &dyn Metrics,
126    ) -> Result<Fired, StoreError> {
127        #[cfg(feature = "test-faults")]
128        if let Some(fired) = apply_relay_delay(ctx, timer).await? {
129            return Ok(fired);
130        }
131        let os_key = keys::outbox_sequence();
132        let sequence_value = ctx.store.get(ctx.partition, &os_key).await?;
133        let os = sequence_value
134            .as_ref()
135            .map(codec::decode_u64)
136            .transpose()?
137            .unwrap_or(0);
138        let rs_key = keys::relay_scan();
139        let mut scan_value = ctx.store.get(ctx.partition, &rs_key).await?;
140        let mut scan = scan_value
141            .as_ref()
142            .and_then(|value| match codec::decode_relay_scan(value) {
143                Ok(state) => Some(state),
144                Err(error) => {
145                    tracing::warn!(
146                        source = ?ctx.partition,
147                        %error,
148                        "corrupt relay scan state; restarting from source head"
149                    );
150                    None
151                }
152            })
153            .filter(|state| state.cursor < state.cycle_end)
154            .unwrap_or(RelayScanV1 {
155                cycle_end: os,
156                cursor: 0,
157                blocked: Vec::new(),
158            });
159        let max_rows = self.budget.max_rows.max(1);
160        let (rows, exhausted) = read_rows(ctx, &scan, max_rows.saturating_mul(4)).await?;
161        let window = decode_window(ctx.partition, ctx.now_ms, rows, exhausted);
162        // `os - first_queued_seq + 1` is an upper bound: out-of-order
163        // delivered rows may leave holes, but a bounded scan window must
164        // never hide a large queued tail from the backlog gauge.
165        let (start, end) = keys::class_range(keys::TAG_RELAY);
166        let head = ctx.store.scan(ctx.partition, &start, &end, None, 1).await?;
167        let backlog = match head.entries.first().and_then(|(key, _)| keys::parse(key)) {
168            Some(keys::ParsedKey::Relay(first)) => os.saturating_sub(first).saturating_add(1),
169            _ => 0,
170        };
171        #[allow(clippy::cast_precision_loss)]
172        metrics.gauge(
173            METRIC_RELAY_BACKLOG_ROWS,
174            &[("source_kind", source_kind(ctx.partition))],
175            backlog as f64,
176        );
177        if window.lag_exceeded {
178            metrics.incr(
179                METRIC_RELAY_LAG_EXCEEDED,
180                &[("source_kind", source_kind(ctx.partition))],
181                1,
182            );
183        }
184        let rh = keys::relay_high_water(ctx.partition)?;
185        let dispatch = self.dispatch(&window.groups, &scan.blocked, &rh).await;
186        if !checkpoint_window(ctx, &rs_key, &mut scan_value, &mut scan, &window, &dispatch).await? {
187            return Ok(Fired::Retry);
188        }
189        if window.corrupt {
190            // The valid prefix is durable; the malformed row stays at the
191            // front of the next scan and is never bypassed.
192            return Ok(Fired::Retry);
193        }
194        let has_remaining = if !scan.blocked.is_empty()
195            || scan.cursor < scan.cycle_end
196            || window.rows.len() > dispatch.delivered.len()
197        {
198            true
199        } else {
200            let (start, end) = keys::class_range(keys::TAG_RELAY);
201            let remaining = ctx.store.scan(ctx.partition, &start, &end, None, 1).await?;
202            !remaining.entries.is_empty() || remaining.next.is_some()
203        };
204        if has_remaining {
205            let due = ctx
206                .now_ms
207                .saturating_add(if dispatch.delivered.is_empty() {
208                    RETRY_BACKOFF_MS
209                } else {
210                    1
211                })
212                .max(timer.due_at_ms.saturating_add(1));
213            Ok(Fired::Reschedule {
214                due_at_ms: due,
215                value: Value::default(),
216                batch: Batch::new(),
217            })
218        } else {
219            // Guard `os` either way: a first-ever relay row committed during
220            // this fire moves `os` from absent, so `Done` races and the timer
221            // stays.
222            Ok(Fired::Done(Batch::new().require(match sequence_value {
223                Some(value) => Precondition::Equals(os_key, value),
224                None => Precondition::Absent(os_key),
225            })))
226        }
227    }
228
229    async fn dispatch(&self, groups: &[TargetRows], blocked: &[Partition], rh: &Key) -> Dispatch {
230        let mut progress = Dispatch {
231            delivered: BTreeSet::new(),
232            block: BTreeSet::new(),
233            pause_at: None,
234            overflow_at: None,
235        };
236        let mut selected = 0u32;
237        let mut calls = 0u32;
238        let max_targets = self.budget.max_targets.max(1);
239        // A duplicate-only group consumes a watermark read but no target
240        // slot. Keep those reads within the Worker's alarm subrequest budget.
241        let per_target_calls = self.budget.max_target_calls.map(|cap| cap.max(2));
242        let total_call_cap = per_target_calls.map(|cap| cap.saturating_mul(max_targets));
243        // Groups are ordered by their first sequence. Reaching the target
244        // budget pauses at the next new target without marking it blocked.
245        for (target, rows) in groups {
246            if progress.pause_at.is_some_and(|at| rows[0].0 >= at) {
247                break;
248            }
249            if blocked.binary_search(target).is_ok() {
250                continue;
251            }
252            let call_limit = match (per_target_calls, total_call_cap) {
253                (Some(per_target), Some(total)) => {
254                    Some(per_target.min(total.saturating_sub(calls)))
255                }
256                _ => None,
257            };
258            if call_limit == Some(0) {
259                progress.pause_at = Some(rows[0].0);
260                break;
261            }
262            let target_rows = rows
263                .iter()
264                .map(|(seq, row, _, _)| (*seq, row.clone()))
265                .collect::<Vec<_>>();
266            let result = self
267                .deliver_target(target, rh, &target_rows, call_limit, selected < max_targets)
268                .await;
269            calls = calls.saturating_add(result.calls);
270            selected += u32::from(result.selected);
271            progress.delivered.extend(
272                rows.iter()
273                    .take(result.completed)
274                    .map(|(seq, _, _, _)| *seq),
275            );
276            if result.failed {
277                if blocked.len() + progress.block.len() == MAX_BLOCKED_TARGETS {
278                    progress.overflow_at = Some(rows[result.completed].0);
279                    break;
280                }
281                progress.block.insert(target.clone());
282            } else if result.completed < rows.len() {
283                // A healthy target cut by a row or Worker-call budget resumes
284                // at its own first unattempted row on the next fire.
285                progress.pause_at =
286                    Some(progress.pause_at.map_or(rows[result.completed].0, |at| {
287                        at.min(rows[result.completed].0)
288                    }));
289            }
290        }
291        progress
292    }
293
294    /// Returns the committed/duplicate prefix, retaining progress if a later chunk fails.
295    #[allow(clippy::too_many_lines)] // One loop owns prefix shrinking, retry accounting and atomic delivery.
296    async fn deliver_target(
297        &self,
298        target: &Partition,
299        rh: &Key,
300        rows: &[(u64, RelayV1)],
301        call_limit: Option<u32>,
302        allow_apply: bool,
303    ) -> TargetResult {
304        let mut result = TargetResult {
305            completed: 0,
306            failed: false,
307            selected: false,
308            calls: 0,
309        };
310        while result.completed < rows.len() {
311            let mut committed = false;
312            for _ in 0..2 {
313                if call_limit.is_some_and(|cap| result.calls >= cap) {
314                    return result;
315                }
316                let mut prefix_end = result.completed
317                    + fitting_prefix(
318                        &Batch::new().require(Precondition::Absent(rh.clone())),
319                        rh,
320                        &rows[result.completed..],
321                    );
322                if prefix_end == result.completed {
323                    result.failed = true;
324                    return result;
325                }
326                let declared = loop {
327                    let Ok(keys) = self
328                        .hook
329                        .read_keys(target, &rows[result.completed..prefix_end])
330                    else {
331                        result.failed = true;
332                        return result;
333                    };
334                    let keys = keys.into_iter().collect::<BTreeSet<_>>();
335                    if keys.len() + usize::from(!keys.contains(rh)) <= MAX_HOOK_READ_KEYS {
336                        break keys;
337                    }
338                    if prefix_end == result.completed + 1 {
339                        result.failed = true;
340                        return result;
341                    }
342                    prefix_end = result.completed + (prefix_end - result.completed) / 2;
343                };
344                result.calls = result.calls.saturating_add(1);
345                let snapshot = async {
346                    let mut observations = Vec::new();
347                    let observed = if declared.is_empty() {
348                        self.target.get(target, rh).await?
349                    } else {
350                        let mut keys = vec![rh.clone()];
351                        keys.extend(declared.into_iter().filter(|key| key != rh));
352                        let values = self.target.get_many(target, &keys).await?;
353                        if values.len() != keys.len() {
354                            return Err(StoreError::Corrupt("short relay snapshot".into()));
355                        }
356                        observations = keys.into_iter().zip(values).collect();
357                        observations[0].1.clone()
358                    };
359                    let hw = observed
360                        .as_ref()
361                        .map(codec::decode_u64)
362                        .transpose()?
363                        .unwrap_or(0);
364                    Ok::<_, StoreError>((observed, hw, observations))
365                }
366                .await;
367                let (observed, hw, observations) = match snapshot {
368                    Ok(snapshot) => snapshot,
369                    Err(error) => {
370                        tracing::warn!(?target, %error, "relay watermark read failed");
371                        result.failed = true;
372                        result.selected = true;
373                        return result;
374                    }
375                };
376                while result.completed < rows.len() && rows[result.completed].0 <= hw {
377                    result.completed += 1;
378                }
379                if result.completed == rows.len() || !allow_apply {
380                    return result;
381                }
382                if result.completed >= prefix_end {
383                    continue;
384                }
385                result.selected = true;
386                let Some((batch, end)) = self
387                    .prepare_target_batch(
388                        target,
389                        rh,
390                        &rows[..prefix_end],
391                        observed.as_ref(),
392                        result.completed,
393                        &observations,
394                    )
395                    .await
396                else {
397                    result.failed = true;
398                    return result;
399                };
400                // Hook additions share the apply's atomicity and size boundary.
401                if call_limit.is_some_and(|cap| result.calls >= cap) {
402                    return result;
403                }
404                result.calls = result.calls.saturating_add(1);
405                match self.target.apply(target, batch).await {
406                    Ok(BatchOutcome::Committed) => {
407                        result.completed = end;
408                        committed = true;
409                        break;
410                    }
411                    Ok(BatchOutcome::PreconditionFailed { .. }) => {}
412                    other => {
413                        tracing::warn!(?target, ?other, "relay target apply failed");
414                        result.failed = true;
415                        return result;
416                    }
417                }
418            }
419            if !committed {
420                result.failed = true;
421                return result;
422            }
423        }
424        result
425    }
426
427    async fn prepare_target_batch(
428        &self,
429        target: &Partition,
430        rh: &Key,
431        rows: &[(u64, RelayV1)],
432        observed: Option<&Value>,
433        start: usize,
434        observations: &[(Key, Option<Value>)],
435    ) -> Option<(Batch, usize)> {
436        let base = Batch::new().require(match observed {
437            Some(value) => Precondition::Equals(rh.clone(), value.clone()),
438            None => Precondition::Absent(rh.clone()),
439        });
440        let mut end = start + fitting_prefix(&base, rh, &rows[start..]);
441        if end == start {
442            tracing::warn!(
443                seq = rows[start].0,
444                "relay row cannot fit one target batch; delivery to this target is stalled"
445            );
446            return None;
447        }
448        // A hook can use more space than the remaining headroom. Shrink a
449        // combined group rather than stalling rows that fit individually.
450        loop {
451            let mut batch = target_batch(rh, observed, &rows[start..end]);
452            if let Err(error) = self
453                .hook
454                .before_apply_observed(
455                    target,
456                    &rows[start..end],
457                    observations,
458                    &mut batch.preconditions,
459                    &mut batch.writes,
460                )
461                .await
462            {
463                if matches!(&error, StoreError::Invalid(message) if message.as_ref() == super::AUDIT_CAPACITY)
464                    && end > start + 1
465                {
466                    end = start + (end - start) / 2;
467                    continue;
468                }
469                tracing::warn!(?target, %error, "relay hook failed");
470                return None;
471            }
472            let mut caps = self.target.capabilities();
473            if !matches!(target, Partition::RefIndex { .. }) {
474                caps.reserved_batch_ops = 0;
475            }
476            caps.reserved_batch_ops = caps
477                .reserved_batch_ops
478                .saturating_add(self.hook.reserved_ops(target, &rows[start..end]));
479            if let Err(error) = batch.validate(&caps) {
480                if matches!(error, StoreError::Invalid(_)) && end > start + 1 {
481                    end = start + (end - start) / 2;
482                    continue;
483                }
484                tracing::warn!(?target, %error, "relay target batch invalid");
485                return None;
486            }
487            return Some((batch, end));
488        }
489    }
490}
491
492fn decode_window(
493    source: &Partition,
494    now_ms: u64,
495    rows: Vec<(Key, Value)>,
496    exhausted: bool,
497) -> ScanWindow {
498    let mut window = ScanWindow {
499        rows: Vec::new(),
500        groups: Vec::new(),
501        corrupt: false,
502        exhausted,
503        lag_exceeded: false,
504    };
505    let mut warned_lag = false;
506    for (key, value) in rows {
507        let decoded = (|| {
508            let Some(keys::ParsedKey::Relay(seq)) = keys::parse(&key) else {
509                return Err(StoreError::Corrupt("bad relay queue key".into()));
510            };
511            if seq == 0 {
512                return Err(StoreError::Corrupt("relay sequence is zero".into()));
513            }
514            Ok((seq, codec::decode_relay(&value)?))
515        })();
516        let (seq, row) = match decoded {
517            Ok(row) => row,
518            Err(error) => {
519                tracing::warn!(source = ?source, ?key, %error, "corrupt relay row; tick stopped");
520                window.corrupt = true;
521                window.exhausted = false;
522                break;
523            }
524        };
525        if !warned_lag {
526            let age_ms = now_ms.saturating_sub(row.at_ms);
527            if age_ms > RELAY_LAG_BOUND_MS {
528                tracing::warn!(source = ?source, age_ms, "outbox relay lag bound exceeded");
529                warned_lag = true;
530                window.lag_exceeded = true;
531            }
532        }
533        let target = row.target.clone();
534        window
535            .rows
536            .push((seq, target.clone(), key.clone(), value.clone()));
537        let i = window
538            .groups
539            .iter()
540            .position(|(partition, _)| partition == &target)
541            .unwrap_or_else(|| {
542                window.groups.push((target, Vec::new()));
543                window.groups.len() - 1
544            });
545        window.groups[i].1.push((seq, row, key, value));
546    }
547    window
548}
549
550fn source_kind(source: &Partition) -> &'static str {
551    match source {
552        Partition::Ref { .. } => "ref",
553        Partition::Namespace(_) => "namespace",
554        Partition::Coordinator(_) => "coordinator",
555        Partition::RepoIndex { .. } => "repo_index",
556        Partition::RefIndex { .. } => "ref_index",
557        Partition::ContentShard(_) => "content",
558    }
559}
560
561async fn checkpoint_window<S: NamespaceStore>(
562    ctx: &TimerCtx<'_, S>,
563    rs_key: &Key,
564    observed: &mut Option<Value>,
565    scan: &mut RelayScanV1,
566    window: &ScanWindow,
567    dispatch: &Dispatch,
568) -> Result<bool, StoreError> {
569    let mut staged = scan.clone();
570    let mut deletions = Vec::new();
571    let mut interrupted = false;
572    let mut overflow = false;
573    for (seq, target, key, value) in &window.rows {
574        if dispatch.overflow_at.is_some_and(|at| *seq >= at) {
575            overflow = true;
576            break;
577        }
578        if dispatch.pause_at.is_some_and(|at| *seq >= at) {
579            interrupted = true;
580            break;
581        }
582        let mut next = staged.clone();
583        next.cursor = *seq;
584        if !dispatch.delivered.contains(seq)
585            && let Err(at) = next.blocked.binary_search(target)
586        {
587            if !dispatch.block.contains(target) {
588                interrupted = true;
589                break;
590            }
591            if next.blocked.len() == MAX_BLOCKED_TARGETS {
592                overflow = true;
593                break;
594            }
595            next.blocked.insert(at, target.clone());
596        }
597        let mut candidate = deletions.clone();
598        if dispatch.delivered.contains(seq) {
599            candidate.push((key.clone(), value.clone()));
600        }
601        if checkpoint_batch(rs_key, observed.as_ref(), &next, &candidate)?
602            .validate(&ctx.store.capabilities())
603            .is_err()
604        {
605            if deletions.is_empty()
606                || !apply_checkpoint(ctx, rs_key, observed, &staged, &deletions).await?
607            {
608                return Ok(false);
609            }
610            deletions.clear();
611            checkpoint_batch(
612                rs_key,
613                observed.as_ref(),
614                &next,
615                &candidate[candidate.len() - usize::from(dispatch.delivered.contains(seq))..],
616            )?
617            .validate(&ctx.store.capabilities())?;
618            if dispatch.delivered.contains(seq) {
619                deletions.push((key.clone(), value.clone()));
620            }
621        } else {
622            deletions = candidate;
623        }
624        staged = next;
625    }
626    let stopped_after = staged.cursor;
627    if overflow {
628        // Terminal marker: the next fire starts from the head with an empty
629        // blocked set. The row that would exceed the cap is retained.
630        staged.cursor = staged.cycle_end;
631    } else if !window.corrupt && !interrupted && window.exhausted {
632        // No more queued rows exist in this cycle's bounded sequence range;
633        // missing sequence numbers are holes left by prior cleanup.
634        staged.cursor = staged.cycle_end;
635    }
636    // Target batches can contain rows beyond a later target's pause or
637    // overflow point. Their target watermark already covers those rows, so
638    // delete them now instead of spending another fire on duplicate groups.
639    // The cursor still stops before the uninspected row.
640    for (seq, _, key, value) in &window.rows {
641        if *seq <= stopped_after || !dispatch.delivered.contains(seq) {
642            continue;
643        }
644        let mut candidate = deletions.clone();
645        candidate.push((key.clone(), value.clone()));
646        if checkpoint_batch(rs_key, observed.as_ref(), &staged, &candidate)?
647            .validate(&ctx.store.capabilities())
648            .is_err()
649        {
650            if deletions.is_empty()
651                || !apply_checkpoint(ctx, rs_key, observed, &staged, &deletions).await?
652            {
653                return Ok(false);
654            }
655            deletions.clear();
656            checkpoint_batch(
657                rs_key,
658                observed.as_ref(),
659                &staged,
660                &[(key.clone(), value.clone())],
661            )?
662            .validate(&ctx.store.capabilities())?;
663            deletions.push((key.clone(), value.clone()));
664        } else {
665            deletions = candidate;
666        }
667    }
668    if !apply_checkpoint(ctx, rs_key, observed, &staged, &deletions).await? {
669        return Ok(false);
670    }
671    *scan = staged;
672    Ok(true)
673}
674
675fn checkpoint_batch(
676    rs_key: &Key,
677    observed: Option<&Value>,
678    state: &RelayScanV1,
679    deletions: &[(Key, Value)],
680) -> Result<Batch, StoreError> {
681    let encoded = codec::encode_relay_scan(state)?;
682    let guard = match observed {
683        Some(value) => Precondition::Equals(rs_key.clone(), value.clone()),
684        None => Precondition::Absent(rs_key.clone()),
685    };
686    let mut batch = Batch::new().require(guard).put(rs_key.clone(), encoded);
687    for (key, value) in deletions {
688        batch = batch
689            .require(Precondition::Equals(key.clone(), value.clone()))
690            .delete(key.clone());
691    }
692    Ok(batch)
693}
694
695async fn apply_checkpoint<S: NamespaceStore>(
696    ctx: &TimerCtx<'_, S>,
697    rs_key: &Key,
698    observed: &mut Option<Value>,
699    state: &RelayScanV1,
700    deletions: &[(Key, Value)],
701) -> Result<bool, StoreError> {
702    let batch = checkpoint_batch(rs_key, observed.as_ref(), state, deletions)?;
703    batch.validate(&ctx.store.capabilities())?;
704    let encoded = codec::encode_relay_scan(state)?;
705    match ctx.store.apply(ctx.partition, batch).await? {
706        BatchOutcome::Committed => {
707            *observed = Some(encoded);
708            Ok(true)
709        }
710        BatchOutcome::PreconditionFailed { .. } | BatchOutcome::DeadlinePassed { .. } => Ok(false),
711    }
712}
713
714fn fitting_prefix(base: &Batch, rh: &Key, rows: &[(u64, RelayV1)]) -> usize {
715    let mut batch = base.clone();
716    for (end, (seq, row)) in rows.iter().enumerate() {
717        let mut candidate = batch.clone();
718        candidate.writes.extend(
719            row.puts
720                .iter()
721                .cloned()
722                .map(|(key, value)| Write::Put(key, value)),
723        );
724        candidate
725            .writes
726            .extend(row.deletes.iter().cloned().map(Write::Delete));
727        let sized = candidate.clone().put(rh.clone(), codec::encode_u64(*seq));
728        if candidate.writes.len() > crate::store::outbox::MAX_RELAY_PUTS
729            || sized.validate(&StoreCapabilities::full()).is_err()
730        {
731            return end;
732        }
733        batch = candidate;
734    }
735    rows.len()
736}
737
738fn target_batch(rh: &Key, observed: Option<&Value>, rows: &[(u64, RelayV1)]) -> Batch {
739    let mut batch = Batch::new().require(match observed {
740        Some(value) => Precondition::Equals(rh.clone(), value.clone()),
741        None => Precondition::Absent(rh.clone()),
742    });
743    for (_, row) in rows {
744        batch.writes.extend(
745            row.puts
746                .iter()
747                .cloned()
748                .map(|(key, value)| Write::Put(key, value)),
749        );
750        batch
751            .writes
752            .extend(row.deletes.iter().cloned().map(Write::Delete));
753    }
754    batch.put(
755        rh.clone(),
756        codec::encode_u64(rows.last().expect("non-empty relay batch").0),
757    )
758}
759
760async fn read_rows<S: NamespaceStore>(
761    ctx: &TimerCtx<'_, S>,
762    state: &RelayScanV1,
763    limit: u32,
764) -> Result<(Vec<(Key, Value)>, bool), StoreError> {
765    if state.cursor >= state.cycle_end || limit == 0 {
766        return Ok((Vec::new(), true));
767    }
768    // Include the exact cursor key (which may remain for a blocked target)
769    // so malformed keys between it and the next sequence cannot be skipped.
770    // At cursor zero the class head also catches malformed keys before seq 1.
771    let anchor = (state.cursor != 0).then(|| keys::relay(state.cursor));
772    let start = anchor
773        .clone()
774        .unwrap_or_else(|| keys::class_range(keys::TAG_RELAY).0);
775    let end = state
776        .cycle_end
777        .checked_add(1)
778        .map_or_else(|| keys::class_range(keys::TAG_RELAY).1, keys::relay);
779    let mut rows = Vec::new();
780    let mut encoded_bytes = 0;
781    let mut remaining = limit;
782    let mut cursor = None;
783    loop {
784        let page = ctx
785            .store
786            .scan(
787                ctx.partition,
788                &start,
789                &end,
790                cursor.as_ref(),
791                SCAN_PAGE_ROWS.min(remaining),
792            )
793            .await?;
794        for (key, value) in page.entries {
795            remaining -= 1;
796            if anchor.as_ref() == Some(&key) {
797                continue;
798            }
799            let size = key.as_bytes().len() + value.as_bytes().len();
800            if !rows.is_empty() && encoded_bytes + size > MAX_FIRE_BYTES {
801                return Ok((rows, false));
802            }
803            encoded_bytes += size;
804            rows.push((key, value));
805        }
806        if remaining == 0 || page.next.is_none() {
807            return Ok((rows, page.next.is_none()));
808        }
809        cursor = page.next;
810    }
811}