1use 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
20const MAX_FIRE_BYTES: usize = 4 * 1024 * 1024;
24const SCAN_PAGE_ROWS: u32 = 64;
25const 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#[derive(Debug)]
61pub struct RelayHandler<T, H = NoHook> {
62 pub target: T,
64 pub hook: H,
66 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 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 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 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 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 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 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 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 #[allow(clippy::too_many_lines)] 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 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 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 staged.cursor = staged.cycle_end;
631 } else if !window.corrupt && !interrupted && window.exhausted {
632 staged.cursor = staged.cycle_end;
635 }
636 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 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}