Skip to content

Commit 470a346

Browse files
letv1nnnjohntmyers
andauthored
fix(api): make WatchSandbox loss-aware and resumable (#3209)
* fix(api): emit warning on WatchSandbox broadcast lag instead of terminating Broadcast lag on the status, log, and platform receivers was converted to a RESOURCE_EXHAUSTED status that terminated the whole watch stream. Lag is recoverable: the receiver resumes at the oldest surviving message. Emit a SandboxStreamWarning and continue streaming instead; keep terminating on Closed. Add helpers and unit tests covering the warning payload and receiver recovery after lag. Partially addresses #3055 (cursor/resume follow up separately). Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * refactor(server): group per-sandbox log bus state and stamp sequence numbers Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * feat(proto): add resume cursor fields to sandbox watch API Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * feat(server): stamp watch cursors from a shared per-sandbox sequence Allocate cursors from a single SeqAllocator shared by the log and platform event buses, so a sandbox's merged watch stream carries unique, strictly increasing cursors. A single resume_after_cursor can then unambiguously locate a client's position across both sources. Rewrite both publish paths to allocate the sequence, stamp event.cursor, send, and append to the tail under one lock. This removes the previous get_mut().expect() TOCTOU race where a concurrent remove() between the two lock sections could panic. Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * feat(server): serve WatchSandbox resume from cursor with gap detection Add tail_after() to the log and platform event buses, returning every buffered event newer than a client's resume cursor. Each PerSandbox now tracks last_trimmed_seq (the highest seq it has evicted) so a resume is reported as an unrecoverable ResumeGap only when this bus dropped an event the client still needs. Judging gaps by evictions, not by the tail's oldest seq, is required under the shared cursor space: each bus's tail is non-contiguous in the global sequence because the other bus owns the missing seqs, so comparing against tail.front() would flag false gaps. Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * feat(server): resume WatchSandbox from cursor across log and platform buses Wire resume_after_cursor into the watch producer. On a non-zero cursor, replay events strictly after it from both the log and platform buses, merge by shared cursor, and emit in order before entering the live loop. A trimmed range on either bus is an unrecoverable gap and terminates the stream with OUT_OF_RANGE carrying the requested and earliest-available cursors, distinct from recoverable lag which warns and continues. Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * test(server): cover WatchSandbox cursor resume paths Add handler-level tests for the resumable watch stream: replay strictly after the client cursor, merge log and platform events in shared-cursor order, suppress duplicates when resuming at the latest cursor, and terminate with OUT_OF_RANGE when the requested cursor has been trimmed. Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * docs(api): document WatchSandbox loss-awareness and resume Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(server): deliver watch events once and harden cursor teardown Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * feat(sdk): add loss-aware resumable watch_logs to Rust SDK client Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(server): keep watch cursors monotonic across teardown and restart Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(server): merge live watch sources by cursor before emission Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(api): bind watch cursors to a cursor space and merge tail sources Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(server): revalidate the watch cursor space after collecting replay Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * test(server): update the public RPC schema fingerprint for the string cursor Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(server): hold watch events above the publication watermark and emit the watch lag warning before its batch Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * test(server): synchronize the watch live-order test with the end of initialization Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(sdk): use canonical sandbox name in watch_logs Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * test(sdk): guard canonical-name addressing in watch_logs Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(server): fix public rpc schema Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> * fix(api): reconcile watch resume rebase Signed-off-by: John Myers <9696606+johntmyers@users.noreply.github.com> * fix(server): bound interactive relay cleanup Signed-off-by: John Myers <9696606+johntmyers@users.noreply.github.com> --------- Signed-off-by: Artem Lytvyn <alytvyn@redhat.com> Signed-off-by: John Myers <9696606+johntmyers@users.noreply.github.com> Co-authored-by: John Myers <9696606+johntmyers@users.noreply.github.com>
1 parent 8835777 commit 470a346

21 files changed

Lines changed: 3180 additions & 158 deletions

File tree

‎Cargo.lock‎

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎architecture/gateway.md‎

Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -440,6 +440,82 @@ Domain objects use shared metadata: stable server-generated IDs, human-readable
440440
names, creation timestamps, and labels. Crate-level details live in
441441
`crates/openshell-core/README.md`.
442442

443+
### Watch streams
444+
445+
`WatchSandbox` merges three per-sandbox sources into one client stream: status
446+
snapshots, server/sandbox logs, and platform events. Logs and platform events
447+
are resumable; a shared per-sandbox allocator stamps each with a `cursor`.
448+
Cursor-ordered delivery is guaranteed for the replay phase: on resume the
449+
buffered events from both sources are sorted before emission. Live events are
450+
monotonic within each source, but the two sources are read independently, so a
451+
client should order across sources by `cursor` rather than by arrival. Status
452+
snapshots and warnings are re-read on demand and carry an empty cursor.
453+
454+
#### Cursor spaces
455+
456+
A sandbox's cursors live in a **cursor space**: a `{epoch, seq}` pair, where the
457+
epoch is a UUID minted on the first publish and `seq` counts from 1. The epoch
458+
is dropped by `TracingLogBus::remove`, so a teardown — or a gateway restart —
459+
retires the space, and the next publish mints a new one. A sequence number alone
460+
cannot distinguish a caught-up client from one holding a cursor out of a space
461+
that no longer exists, because the replacement space reuses the same numbers;
462+
the epoch answers *which counter issued this*, which is the question resume
463+
validation actually has to ask. Behind multiple replicas the same rule makes a
464+
reconnect to a different replica fail loudly rather than return the wrong
465+
events.
466+
467+
On the wire a cursor is an opaque, fixed-width token. Clients may only compare
468+
two cursors from one stream and keep the greater; the encoding zero-pads `seq`
469+
so that byte-wise comparison matches sequence order, which is what lets every
470+
SDK track a high-water mark without parsing. A stream only ever observes one
471+
epoch — a reset closes both resumable broadcast receivers, ending the stream
472+
rather than switching spaces mid-flight — so that comparison is always well
473+
defined where clients are allowed to use it. The gateway does not rely on it:
474+
server-side ordering runs on the raw `u64` seq carried alongside each event in
475+
`CursoredEvent`, never on the token.
476+
477+
The gateway holds a bounded in-memory tail per sandbox. Loss is reported with
478+
two distinct, documented behaviors:
479+
480+
- **Recoverable lag** — a broadcast receiver falls behind and the server skips
481+
ahead. The stream emits a `SandboxStreamWarning` event and continues. Since
482+
cursors are opaque, the warning is the client's only signal.
483+
- **Unrecoverable gap** — the server sends a snapshot, then terminates with
484+
`OUT_OF_RANGE`. Three cases reach it: the tail was trimmed past the requested
485+
cursor, the cursor's epoch does not match the sandbox's current space, or no
486+
space exists because nothing has been published since teardown. The status
487+
tells the client to restart with an empty cursor; retrying the same token
488+
fails identically. A token the gateway could not have issued is rejected
489+
earlier, as `INVALID_ARGUMENT` on the call itself.
490+
491+
Both resumable sources draw from one cursor space, so the server merges them by
492+
seq before emitting rather than draining each in turn: on resume it replays only
493+
events after the client's cursor, and without one it replays each bus's retained
494+
tail. Either way the batch leaves in ascending cursor order. The two tails are
495+
bounded independently (`log_tail_lines` and `event_tail`), so merging orders
496+
whatever each bus kept; it does not align their depths.
497+
498+
The broadcast receivers are subscribed before replay, so an event buffered during
499+
initialization could appear in both replay and the live receiver; the producer
500+
tracks the highest replayed seq and suppresses live events at or below it, so
501+
each event is delivered once. That mark is per source. The two tails are read at
502+
different instants and bounded independently, so one shared mark would let the
503+
deeper source censor the shallower one — with `event_tail` unset the mark rises
504+
to the newest buffered log while no platform event is replayed at all, and
505+
platform events published during initialization are discarded as duplicates of a
506+
replay that never ran. Subscribing never mints a cursor space, so a
507+
resume against a torn-down sandbox cannot create the space its stale cursor is
508+
then checked against. Clients track the highest observed `cursor` and pass it as
509+
`resume_after_cursor` on reconnect.
510+
511+
The epoch is validated twice on resume: once before reading the tails and again
512+
once both are in hand, before anything is emitted. The check and each read take
513+
their locks separately, so a teardown plus a republish can retire the validated
514+
space and install a replacement in between; the reads would then apply the old
515+
space's seq to the replacement's buffers, and a trimmed-range check that only
516+
compares numbers would report no gap while skipping the replacement's lower
517+
events. The second look ends the stream with `OUT_OF_RANGE` instead.
518+
443519
## Persistence
444520

445521
The gateway persistence layer is a protobuf object store. Domain services store

‎crates/openshell-cli/src/run.rs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -809,6 +809,7 @@ pub async fn sandbox_create(
809809
since_time: None,
810810
log_sources: vec!["gateway".to_string()],
811811
log_min_level: String::new(),
812+
resume_after_cursor: String::new(),
812813
})
813814
.await
814815
.into_diagnostic()?
@@ -3613,6 +3614,7 @@ async fn wait_for_lifecycle_phase(
36133614
since_time: None,
36143615
log_sources: Vec::new(),
36153616
log_min_level: String::new(),
3617+
resume_after_cursor: String::new(),
36163618
})
36173619
.await
36183620
.into_diagnostic()?
@@ -5847,6 +5849,7 @@ pub async fn sandbox_logs(
58475849
.into_diagnostic()?,
58485850
log_sources: source_filter,
58495851
log_min_level: level.to_uppercase(),
5852+
resume_after_cursor: String::new(),
58505853
})
58515854
.await
58525855
.into_diagnostic()?

‎crates/openshell-cli/tests/sandbox_create_lifecycle_integration.rs‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -743,6 +743,7 @@ impl OpenShell for TestOpenShell {
743743
let _ = tx
744744
.send(Ok(SandboxStreamEvent {
745745
payload: Some(sandbox_stream_event::Payload::Sandbox(provisioning)),
746+
cursor: String::new(),
746747
}))
747748
.await;
748749
if terminal_after_provisional_container_exit
@@ -753,6 +754,7 @@ impl OpenShell for TestOpenShell {
753754
payload: Some(sandbox_stream_event::Payload::Sandbox(
754755
provisional_container_exit,
755756
)),
757+
cursor: String::new(),
756758
}))
757759
.await;
758760
provisional_container_exit_sent.notify_waiters();
@@ -764,6 +766,7 @@ impl OpenShell for TestOpenShell {
764766
let _ = tx
765767
.send(Ok(SandboxStreamEvent {
766768
payload: Some(sandbox_stream_event::Payload::Sandbox(completed)),
769+
cursor: String::new(),
767770
}))
768771
.await;
769772
return;
@@ -777,11 +780,13 @@ impl OpenShell for TestOpenShell {
777780
message: "Started VM launcher".to_string(),
778781
..PlatformEvent::default()
779782
})),
783+
cursor: String::new(),
780784
}))
781785
.await;
782786
let _ = tx
783787
.send(Ok(SandboxStreamEvent {
784788
payload: Some(sandbox_stream_event::Payload::Sandbox(error)),
789+
cursor: String::new(),
785790
}))
786791
.await;
787792
tokio::time::sleep(Duration::from_secs(5)).await;
@@ -801,12 +806,14 @@ impl OpenShell for TestOpenShell {
801806
source: "gateway".to_string(),
802807
fields: HashMap::new(),
803808
})),
809+
cursor: String::new(),
804810
}))
805811
.await;
806812
}
807813
let _ = tx
808814
.send(Ok(SandboxStreamEvent {
809815
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
816+
cursor: String::new(),
810817
}))
811818
.await;
812819
return;
@@ -815,6 +822,7 @@ impl OpenShell for TestOpenShell {
815822
let _ = tx
816823
.send(Ok(SandboxStreamEvent {
817824
payload: Some(sandbox_stream_event::Payload::Sandbox(completed)),
825+
cursor: String::new(),
818826
}))
819827
.await;
820828
return;
@@ -829,6 +837,7 @@ impl OpenShell for TestOpenShell {
829837
message: "Preparing rootfs".to_string(),
830838
..PlatformEvent::default()
831839
})),
840+
cursor: String::new(),
832841
}))
833842
.await;
834843
tokio::time::sleep(Duration::from_millis(600)).await;
@@ -840,12 +849,14 @@ impl OpenShell for TestOpenShell {
840849
message: "Formatting root disk".to_string(),
841850
..PlatformEvent::default()
842851
})),
852+
cursor: String::new(),
843853
}))
844854
.await;
845855
tokio::time::sleep(Duration::from_millis(600)).await;
846856
let _ = tx
847857
.send(Ok(SandboxStreamEvent {
848858
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
859+
cursor: String::new(),
849860
}))
850861
.await;
851862
return;
@@ -857,11 +868,13 @@ impl OpenShell for TestOpenShell {
857868
message: "Sandbox scheduled".to_string(),
858869
..PlatformEvent::default()
859870
})),
871+
cursor: String::new(),
860872
}))
861873
.await;
862874
let _ = tx
863875
.send(Ok(SandboxStreamEvent {
864876
payload: Some(sandbox_stream_event::Payload::Sandbox(ready)),
877+
cursor: String::new(),
865878
}))
866879
.await;
867880
});

‎crates/openshell-sdk/Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ tokio-tungstenite = { workspace = true }
2828
tonic = { workspace = true, features = ["tls-native-roots"] }
2929
tower = { workspace = true }
3030
tracing = { workspace = true }
31+
async-stream = "0.3.6"
3132

3233
[dev-dependencies]
3334
serde_json = { workspace = true }

0 commit comments

Comments
 (0)