test(sync): audit browser recovery paths
This commit is contained in:
+21
-16
@@ -842,6 +842,7 @@ async fn sync_acknowledgements(
|
||||
.get("reconnect")
|
||||
.filter(|value| !value.is_empty())
|
||||
.cloned();
|
||||
let persistent_stream = reconnect_key.is_some();
|
||||
|
||||
let mut sync = state.sync.lock().unwrap();
|
||||
if let Some(key) = reconnect_key {
|
||||
@@ -921,22 +922,26 @@ async fn sync_acknowledgements(
|
||||
};
|
||||
drop(sync);
|
||||
let heartbeat_interval = state.acknowledgement_heartbeat_interval;
|
||||
let heartbeat = stream::unfold(heartbeat_interval, |interval| async move {
|
||||
tokio::time::sleep(interval).await;
|
||||
Some((
|
||||
Ok(Event::default()
|
||||
.event("heartbeat")
|
||||
.data("{\"status\":\"ok\"}")),
|
||||
interval,
|
||||
))
|
||||
});
|
||||
Ok(
|
||||
Sse::new(stream::iter(events).chain(heartbeat).boxed()).keep_alive(
|
||||
KeepAlive::new()
|
||||
.interval(heartbeat_interval)
|
||||
.text("heartbeat"),
|
||||
),
|
||||
)
|
||||
let event_stream = stream::iter(events).boxed();
|
||||
let response_stream = if persistent_stream {
|
||||
let heartbeat = stream::unfold(heartbeat_interval, |interval| async move {
|
||||
tokio::time::sleep(interval).await;
|
||||
Some((
|
||||
Ok(Event::default()
|
||||
.event("heartbeat")
|
||||
.data("{\"status\":\"ok\"}")),
|
||||
interval,
|
||||
))
|
||||
});
|
||||
event_stream.chain(heartbeat).boxed()
|
||||
} else {
|
||||
event_stream
|
||||
};
|
||||
Ok(Sse::new(response_stream).keep_alive(
|
||||
KeepAlive::new()
|
||||
.interval(heartbeat_interval)
|
||||
.text("heartbeat"),
|
||||
))
|
||||
}
|
||||
|
||||
fn registry(state: Arc<AppState>) -> impl DispatchRegistry {
|
||||
|
||||
Reference in New Issue
Block a user