From 4f8c6bce09dab3e6e95fdab28d54f6b38cf5be7d Mon Sep 17 00:00:00 2001 From: wgqqqqq Date: Tue, 22 Sep 2026 20:51:23 +0800 Subject: [PATCH] perf(remote): load session history pages on demand --- .../tools/tests/host-stream.test.cjs | 31 + .../mobile/core/transport/HostStreamTest.kt | 18 + src/crates/assembly/core/AGENTS.md | 9 + .../src/agentic/coordination/coordinator.rs | 67 ++- .../core/src/agentic/persistence/manager.rs | 230 +++++++ .../core/src/service/remote_connect/mod.rs | 17 +- .../core/src/service_agent_runtime.rs | 33 ++ .../services/services-integrations/AGENTS.md | 8 + .../src/remote_connect/host_stream.rs | 560 +++++++++++++++++- 9 files changed, 942 insertions(+), 31 deletions(-) diff --git a/src/apps/mobile/harmonyos/tools/tests/host-stream.test.cjs b/src/apps/mobile/harmonyos/tools/tests/host-stream.test.cjs index 8d7efa613c..a72c7789da 100644 --- a/src/apps/mobile/harmonyos/tools/tests/host-stream.test.cjs +++ b/src/apps/mobile/harmonyos/tools/tests/host-stream.test.cjs @@ -271,3 +271,34 @@ test('wire parsers reject foreign pages, invalid shapes, and non-hint device eve assert.throws(() => checkHostStreamPage('s', { resp: 'error', message: 'invalid RPC command' }), HostStreamUnsupportedError); assert.throws(() => checkHostStreamPage('s', { resp: 'ok' }), /Unexpected stream response/); }); + + +test('lazy disk history keeps JS-safe backward cursors separate from live updates', async () => { + const ceiling = 2 ** 52; + const c = callbacks(); const reads = []; let hint; + const record = (seq, id, turn) => ({ seq, event: 'session-record', payload: { id, revision: seq, turn: { turnId: turn } } }); + const source = { + target: 'desktop', unsubscribe: async () => {}, + onHint: listener => { hint = listener; return () => {}; }, onReconnect: () => () => {}, + read: async request => { + reads.push({ ...request }); + const events = request.before !== undefined ? [record(ceiling - 3, 'old', 'old-turn')] + : request.after !== undefined ? [record(ceiling, 'new', 'new-turn')] + : [record(ceiling - 2, 'a', 'turn'), record(ceiling - 1, 'b', 'turn')]; + return { resp: 'stream_page', stream_id: 'session', epoch: 8, events, + cursor: request.after !== undefined ? ceiling : ceiling - 1, + oldest_seq: events[0].seq, has_more: request.before === undefined && request.after === undefined, truncated: false }; + } + }; + const stream = new HostSessionStream('session', source, c.hooks); + try { + await settle(() => c.caught === 1); + stream.loadOlder(); await settle(() => c.applied.length === 3); + assert.deepEqual(reads[1], { before: ceiling - 2, epoch: 8 }); + hint({ sourceDeviceId: 'desktop', streamId: 'session', epoch: 8, cursor: ceiling }); + await settle(() => c.applied.length === 4); + assert.deepEqual(reads[2], { after: ceiling - 1, epoch: 8 }); + assert.deepEqual(c.applied.map(e => e.payload.id), ['a', 'b', 'old', 'new']); + assert.equal(c.errors.length, 0); + } finally { stream.close(); } +}); diff --git a/src/apps/mobile/shared/core-transport/src/commonTest/kotlin/com/openbitfun/mobile/core/transport/HostStreamTest.kt b/src/apps/mobile/shared/core-transport/src/commonTest/kotlin/com/openbitfun/mobile/core/transport/HostStreamTest.kt index 5f7901583e..398932003a 100644 --- a/src/apps/mobile/shared/core-transport/src/commonTest/kotlin/com/openbitfun/mobile/core/transport/HostStreamTest.kt +++ b/src/apps/mobile/shared/core-transport/src/commonTest/kotlin/com/openbitfun/mobile/core/transport/HostStreamTest.kt @@ -77,6 +77,24 @@ class HostStreamTest { return received to job } + @Test fun jsSafeHistoryCursorsDoNotMoveTheForwardCursor() = runTest { + val ceiling = 1L shl 52 + val host = FakeHost("s1", pageSize = 2) + host.nextSeq = ceiling - 3 + repeat(3) { host.append("session-record", turnRecord("t$it", it)) } + val hints = MutableSharedFlow() + val older = Channel>() + val (received, _) = open(host, hints = hints, older = older) + runCurrent() + val page = CompletableDeferred(); older.send(page); runCurrent() + assertTrue(page.isCompleted) + assertEquals(Triple(null, ceiling - 2, 1L), host.reads.last()) + host.append("session-record", turnRecord("live", 3)) + hints.emit(StreamHint("desktop-1", "s1",host.epoch,host.cursor)); runCurrent() + assertEquals(Triple(ceiling - 1, null, 1L),host.reads.last()) + assertEquals(listOf("t1","t2","t0","live"), emittedTurnIds(received)) + } + @Test fun completionWaitsUntilTheCollectorConsumesTheWholePage() = runTest { val host = FakeHost("s1", pageSize = 2) repeat(4) { host.append("session-record", turnRecord("turn-$it", it)) } diff --git a/src/crates/assembly/core/AGENTS.md b/src/crates/assembly/core/AGENTS.md index f11ecf0e93..a97ce7b6b8 100644 --- a/src/crates/assembly/core/AGENTS.md +++ b/src/crates/assembly/core/AGENTS.md @@ -212,6 +212,15 @@ or test-target layout. Workspace checks and product-wide tests are CI-backed and are not the default Core precheck. For documentation-only changes, run `git diff --check`. +For disk-backed history paging and legacy sessions without a catalog: +`cargo test --locked -p openbitfun-core --no-default-features --features remote-connect,git --lib history_page_`. +Also run the `staged_revert_catalog_projection` and `load_relay_session_turns_` +filters for the same target when changing visibility. Paging must not parse +unrelated turn bodies or rewrite history. To compare real-file first-page work +against full materialization locally, use the same target with +`history_page_benchmark -- --ignored --nocapture`; it checks content equivalence +and reports timings without asserting a machine-dependent latency in CI. + For built-in provider overlay, trusted endpoint validation, and reasoning catalog changes: ```bash diff --git a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs index 8b68be6038..03654702f8 100644 --- a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs +++ b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs @@ -8588,17 +8588,48 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet session_id: &str, turn_id: Option<&str>, ) -> OpenBitFunResult> { + self.load_relay_session_selection(storage, session_id, turn_id, None) + .await + .map(|page| page.0) + } + + pub async fn load_relay_history_turn( + &self, + storage: &Path, + session_id: &str, + before: Option, + ) -> OpenBitFunResult<(Vec, Option)> { + self.load_relay_session_selection(storage, session_id, None, Some(before)) + .await + } + + async fn load_relay_session_selection( + &self, + storage: &Path, + session_id: &str, + turn_id: Option<&str>, + history: Option>, + ) -> OpenBitFunResult<(Vec, Option)> { let _mutation = self .session_manager .acquire_session_mutation(session_id) .await?; self.prepare_persisted_session_read_locked(storage, session_id) .await?; - let mut turns = self - .session_manager - .persistence_manager() - .load_visible_session_turns(storage, session_id) - .await?; + let (mut turns, next) = if let Some(before) = history { + self.session_manager + .persistence_manager() + .load_visible_history_turn(storage, session_id, before) + .await? + } else { + ( + self.session_manager + .persistence_manager() + .load_visible_session_turns(storage, session_id) + .await?, + None, + ) + }; if let Some(turn_id) = turn_id { turns.retain(|turn| turn.turn_id == turn_id); if turns.is_empty() { @@ -8607,10 +8638,16 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet ))); } } - let context = self - .session_manager - .get_context_messages(session_id) - .await?; + let context = if turns + .iter() + .any(|turn| turn.status == TurnStatus::InProgress) + { + self.session_manager + .get_context_messages(session_id) + .await? + } else { + Vec::new() + }; for turn in &mut turns { if turn.status == TurnStatus::InProgress { let messages: Vec<_> = context @@ -8625,7 +8662,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet SessionManager::append_generation_rounds(turn, &id, &messages, timestamp); } } - Ok(turns) + Ok((turns, next)) } /// Export a transcript while retaining the same Session history boundary @@ -17193,6 +17230,16 @@ mod tests { assert_eq!(one.len(), 1); assert_eq!(one[0].turn_id, turn_id); + let (page, next) = coordinator + .load_relay_history_turn(&storage, &session.session_id, None) + .await + .expect("paged history must not require an in-memory writer"); + assert_eq!( + serde_json::to_value(&page).unwrap(), + serde_json::to_value(&one).unwrap() + ); + assert_eq!(next, None); + session_manager .restore_session(workspace.path(), &session.session_id) .await diff --git a/src/crates/assembly/core/src/agentic/persistence/manager.rs b/src/crates/assembly/core/src/agentic/persistence/manager.rs index e49668e004..b899294ec4 100644 --- a/src/crates/assembly/core/src/agentic/persistence/manager.rs +++ b/src/crates/assembly/core/src/agentic/persistence/manager.rs @@ -3785,6 +3785,47 @@ impl PersistenceManager { .await } + /// Read one historical turn by storage position without parsing unrelated bodies. + /// Filename indices remain authoritative even when a legacy sidecar is absent + /// or stale. `next` is exclusive, and respects the current undo boundary. + pub async fn load_visible_history_turn( + &self, + workspace_path: &Path, + session_id: &str, + before: Option, + ) -> OpenBitFunResult<(Vec, Option)> { + Self::validate_session_id(session_id)?; + let _writer = match self.lock_session_write_operation(workspace_path, session_id) { + Ok(lock) => Some(lock), + Err(OpenBitFunError::SessionInUse { .. }) => None, + Err(error) => return Err(error), + }; + let boundary = self + .load_session_revert_state(workspace_path, session_id) + .await? + .map(|state| state.boundary_turn) + .unwrap_or(usize::MAX); + let ceiling = before.unwrap_or(usize::MAX).min(boundary); + let mut paths = self + .list_indexed_turn_paths(workspace_path, session_id) + .await?; + paths.retain(|(index, _)| *index < ceiling); + paths.sort_by_key(|(index, _)| *index); + let Some((index, path)) = paths.pop() else { + return Ok((Vec::new(), None)); + }; + let file = self + .read_json_optional::(&path) + .await? + .ok_or_else(|| OpenBitFunError::NotFound("History changed during page read".into()))?; + if file.turn.session_id != session_id || file.turn.turn_index != index { + return Err(OpenBitFunError::Validation( + "History turn identity mismatch".into(), + )); + } + Ok((vec![file.turn], (!paths.is_empty()).then_some(index))) + } + async fn project_visible_session_turns( &self, workspace_path: &Path, @@ -6242,6 +6283,183 @@ mod tests { assert_eq!(rebuilt.entries[2].preview.as_deref(), Some("prompt 2")); } + /// Opt-in measurement: real persisted files, identical first-page content, + /// and fresh stream hubs on each sample. No timing assertion in CI. + #[cfg(feature = "remote-connect")] + #[tokio::test] + #[ignore = "local first-page performance comparison"] + async fn history_page_benchmark_against_full_materialization() { + use openbitfun_services_integrations::remote_connect::{ + host_stream::{HistoryBatch, HostStreamHub, HostStreamNotifier, StreamReadRequest}, + session_records::records_from_turns, + }; + use std::sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, + }; + struct Quiet; + impl HostStreamNotifier for Quiet { + fn notify(&self, _: &str, _: serde_json::Value) {} + } + let workspace = TestWorkspace::new(); + let manager = PersistenceManager::new(workspace.path_manager()).unwrap(); + let id = Uuid::new_v4().to_string(); + let session = Session::new_with_id( + id.clone(), + "History benchmark".into(), + "Standard".into(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().into()), + ..Default::default() + }, + ); + manager + .save_session(workspace.path(), &session) + .await + .unwrap(); + for index in 0..128 { + manager + .save_dialog_turn( + workspace.path(), + &DialogTurnData::new( + format!("turn-{index}"), + index, + id.clone(), + user_message(&"x".repeat(128 * 1024)), + ), + ) + .await + .unwrap(); + } + let request = StreamReadRequest { + stream_id: id.clone(), + subscribe: true, + ..Default::default() + }; + let mut old_times = Vec::new(); + let mut new_times = Vec::new(); + for _ in 0..5 { + let legacy = HostStreamHub::start(Arc::new(Quiet)); + legacy.activate(&id); + let start = std::time::Instant::now(); + legacy + .synchronize_records(id.clone(), true, || async { + let turns = manager + .load_visible_session_turns(workspace.path(), &id) + .await?; + records_from_turns(&turns, &|_| None) + }) + .await + .unwrap(); + let old_page = legacy.read("phone", &request).unwrap(); + old_times.push(start.elapsed().as_micros()); + let paged = HostStreamHub::start(Arc::new(Quiet)); + let reads = AtomicUsize::new(0); + let start = std::time::Instant::now(); + let new_page = paged + .read_history("phone", &request, |before| { + reads.fetch_add(1, Ordering::SeqCst); + let manager = &manager; + let workspace = &workspace; + let id = &id; + async move { + let (turns, before) = manager + .load_visible_history_turn(workspace.path(), id, before) + .await?; + Ok(HistoryBatch { + records: records_from_turns(&turns, &|_| None)?, + before, + }) + } + }) + .await + .unwrap(); + new_times.push(start.elapsed().as_micros()); + let content = |events: Vec< + openbitfun_services_integrations::remote_connect::host_stream::StreamEvent, + >| { + events + .into_iter() + .map(|event| { + let mut payload = event.payload; + payload.as_object_mut().unwrap().remove("revision"); + payload + }) + .collect::>() + }; + assert_eq!(content(old_page.events), content(new_page.events)); + assert_eq!(old_page.has_more, new_page.has_more); + assert!(reads.load(Ordering::SeqCst) < 128); + eprintln!( + "history comparison: full_us={} paged_us={} full_turns=128 paged_turns={}", + old_times.last().unwrap(), + new_times.last().unwrap(), + reads.load(Ordering::SeqCst) + ); + } + old_times.sort_unstable(); + new_times.sort_unstable(); + eprintln!( + "history median: full_us={} paged_us={}", + old_times[2], new_times[2] + ); + } + + #[tokio::test] + async fn history_page_reads_only_selected_body_without_a_catalog() { + let workspace = TestWorkspace::new(); + let manager = PersistenceManager::new(workspace.path_manager()).unwrap(); + let id = Uuid::new_v4().to_string(); + let session = Session::new_with_id( + id.clone(), + "Paged".into(), + "Standard".into(), + SessionConfig { + workspace_path: Some(workspace.path().to_string_lossy().into()), + ..Default::default() + }, + ); + manager + .save_session(workspace.path(), &session) + .await + .unwrap(); + for index in [0, 3, 9] { + manager + .save_dialog_turn( + workspace.path(), + &DialogTurnData::new( + format!("turn-{index}"), + index, + id.clone(), + user_message("page"), + ), + ) + .await + .unwrap(); + } + std::fs::remove_file(manager.turn_catalog_path(workspace.path(), &id)).unwrap(); + std::fs::write(manager.turn_path(workspace.path(), &id, 0), "invalid json").unwrap(); + let (latest, next) = manager + .load_visible_history_turn(workspace.path(), &id, None) + .await + .unwrap(); + assert_eq!(latest[0].turn_index, 9); + assert_eq!(next, Some(9)); + let (older, next) = manager + .load_visible_history_turn(workspace.path(), &id, next) + .await + .unwrap(); + assert_eq!(older[0].turn_index, 3); + assert_eq!(next, Some(3)); + assert!( + manager + .load_visible_history_turn(workspace.path(), &id, next) + .await + .is_err(), + "corrupt selected history must not silently disappear" + ); + } + #[tokio::test] async fn staged_revert_catalog_projection_hides_the_physical_suffix() { let workspace = TestWorkspace::new(); @@ -6296,6 +6514,18 @@ mod tests { .await .expect("staged revert should save"); + std::fs::write( + manager.turn_path(workspace.path(), &session_id, 1), + "invalid json", + ) + .unwrap(); + let (page, next) = manager + .load_visible_history_turn(workspace.path(), &session_id, None) + .await + .unwrap(); + assert_eq!(page[0].turn_id, "turn-0"); + assert_eq!(next, None); + let projected = manager .load_session_turn_catalog( workspace.path(), diff --git a/src/crates/assembly/core/src/service/remote_connect/mod.rs b/src/crates/assembly/core/src/service/remote_connect/mod.rs index e0b57ce429..6be7c97742 100644 --- a/src/crates/assembly/core/src/service/remote_connect/mod.rs +++ b/src/crates/assembly/core/src/service/remote_connect/mod.rs @@ -1116,13 +1116,16 @@ pub async fn handle_host_stream_command( .to_string(), }); }; - if is_session_stream(&request.stream_id) && hub.needs_full_synchronization(request) { - // A fresh subscriber, or one whose epoch the host no longer holds, - // needs the runtime's stable records before its first page. - hub.activate(&request.stream_id); - if let Err(error) = synchronize_session_records(hub, &request.stream_id).await { - return Some(RemoteResponse::Error { message: error }); - } + if is_session_stream(&request.stream_id) { + let page = hub.read_history(source_device_id, request, |before| { + crate::service_agent_runtime::CoreServiceAgentRuntime::load_relay_history_batch(&request.stream_id, before) + }).await; + return Some(match page { + Ok(page) => RemoteResponse::StreamPage { page }, + Err(error) => RemoteResponse::Error { + message: error.to_string(), + }, + }); } Some(match hub.read(source_device_id, request) { Ok(page) => RemoteResponse::StreamPage { page }, diff --git a/src/crates/assembly/core/src/service_agent_runtime.rs b/src/crates/assembly/core/src/service_agent_runtime.rs index 14379b960d..541686d430 100644 --- a/src/crates/assembly/core/src/service_agent_runtime.rs +++ b/src/crates/assembly/core/src/service_agent_runtime.rs @@ -2023,6 +2023,36 @@ impl CoreServiceAgentRuntime { image_context_from_remote_image_context(context) } + /// Read just one persisted historical turn for the host's backward cursor. + #[cfg(feature = "remote-connect")] + pub(crate) async fn load_relay_history_batch( + session_id: &str, + before: Option, + ) -> anyhow::Result + { + let directory = Self::resolve_session_storage_dir(session_id) + .await + .ok_or_else(|| anyhow::anyhow!("Session storage is unavailable on this host"))?; + let coordinator = + get_global_coordinator().ok_or_else(|| anyhow::anyhow!("Runtime is unavailable"))?; + let (turns, before) = coordinator + .load_relay_history_turn(&directory, session_id, before) + .await?; + let records = tokio::task::spawn_blocking(move || { + openbitfun_services_integrations::remote_connect::session_records::records_from_turns( + &turns, + &read_remote_chat_image_pixels, + ) + }) + .await??; + Ok( + openbitfun_services_integrations::remote_connect::host_stream::HistoryBatch { + records, + before, + }, + ) + } + /// One source read/commit owner for both migration and live block updates. #[cfg(feature = "remote-connect")] pub(crate) async fn synchronize_relay_session( @@ -2030,6 +2060,9 @@ impl CoreServiceAgentRuntime { session_id: &str, turn_id: Option<&str>, ) -> Result<(), String> { + if turn_id.is_none() && hub.invalidate_paged_history(session_id).await { + return Ok(()); + } hub.synchronize_records(session_id.to_owned(),turn_id.is_none(),||async { let directory=Self::resolve_session_storage_dir(session_id).await .ok_or_else(||anyhow::anyhow!("Session storage is unavailable on this host"))?; diff --git a/src/crates/services/services-integrations/AGENTS.md b/src/crates/services/services-integrations/AGENTS.md index 25e84ae291..f57f00b4e4 100644 --- a/src/crates/services/services-integrations/AGENTS.md +++ b/src/crates/services/services-integrations/AGENTS.md @@ -34,6 +34,14 @@ slices that are outside pure product logic but still platform-neutral. hint leases. `remote_connect::host_stream_subscriber` is the Rust controller reader. Neither the relay nor any client persists stream content; do not add relay-stored history, `get_session_key`, or a durable stream cache here. +- Session history reads use `HostStreamHub::read_history` with a Core-owned + source loader. Backfill allocates decreasing JS-safe sequences below the + live sequence range; it never advances the forward cursor. One source gate + orders backfill, live publication, and history invalidation. Undo/import and + cache misses across evicted pages require an epoch fence, never silent skips. + Only the requested page's bodies enter the bounded log; the pending source + batch is at most one persisted turn, which can itself be large. No wire or + persisted-format migration is required; all controllers retain `read_stream`. - The `remote-persistence` feature is the lightweight persisted-shape owner shared by Remote Connect, remote SSH, and offline migration. Keep it free of network, SSH transport, and runtime orchestration dependencies so owner readers and diff --git a/src/crates/services/services-integrations/src/remote_connect/host_stream.rs b/src/crates/services/services-integrations/src/remote_connect/host_stream.rs index 57d088bbf7..e2f500078b 100644 --- a/src/crates/services/services-integrations/src/remote_connect/host_stream.rs +++ b/src/crates/services/services-integrations/src/remote_connect/host_stream.rs @@ -90,7 +90,26 @@ struct RecordIndexEntry { turn: String, } +// Two disjoint, JS-safe sequence ranges keep historical backfill below every +// live update. They are opaque cursors on the existing wire, not timestamps or +// array indices. Backfill never advances a subscriber's forward cursor. +const HISTORY_SEQUENCE_CEILING: u64 = 1 << 52; + +pub struct HistoryBatch { + pub records: Vec, + pub before: Option, +} + +struct HistoryReader { + before: Option, + exhausted: bool, + pending: Vec, + next_seq: u64, + newer_evicted: bool, +} + struct StreamLog { + history: Option, epoch: u64, next_seq: u64, events: BTreeMap, @@ -105,6 +124,7 @@ struct StreamLog { impl StreamLog { fn new() -> Self { Self { + history: None, epoch: fresh_epoch(), next_seq: 1, events: BTreeMap::new(), @@ -164,7 +184,17 @@ impl StreamLog { } } while self.bytes > STREAM_BYTES_BUDGET && self.events.len() > 1 { - let Some((&seq, _)) = self.events.iter().next() else { + // A controller scrolling backward already has newer history. Evict + // those bodies first, retaining this requested page. Reopening or a + // slower reader crossing the evicted range gets a fresh epoch and + // reloads from disk; no unread suffix is silently skipped. + let historical = self.history.as_ref().and_then(|_| { + self.events + .range(..HISTORY_SEQUENCE_CEILING) + .next_back() + .map(|(&seq, _)| seq) + }); + let Some(seq) = historical.or_else(|| self.events.keys().next().copied()) else { break; }; if let Some((event, size)) = self.events.remove(&seq) { @@ -174,7 +204,11 @@ impl StreamLog { self.record_seq.remove(id); } } - self.truncated = true; + if historical.is_some() { + self.history.as_mut().unwrap().newer_evicted = true; + } else { + self.truncated = true; + } } } } @@ -226,7 +260,12 @@ impl StreamLog { stream_id: stream_id.to_owned(), epoch: self.epoch, events, - has_more, + has_more: has_more + || (request.after.is_none() + && self + .history + .as_ref() + .is_some_and(|history| !history.exhausted || !history.pending.is_empty())), cursor: self.latest(), oldest_seq: self.oldest(), truncated: self.truncated, @@ -406,6 +445,219 @@ impl HostStreamHub { .last_activity = Instant::now(); } + async fn source_gate(&self, session_id: &str) -> Arc> { + self.source_gates + .lock() + .await + .entry(session_id.to_owned()) + .or_default() + .clone() + } + + /// Undo/import invalidate the chronological source. Reusing old backfill + /// positions would resurrect deleted turns; an epoch fence makes all readers + /// discard their derived view and request the latest bounded page again. + pub async fn invalidate_paged_history(&self, session_id: &str) -> bool { + let gate = self.source_gate(session_id).await; + let _source = gate.lock().await; + let mut state = self.lock(); + let Some(log) = state.streams.get_mut(session_id) else { + return false; + }; + if log.history.is_none() { + return false; + } + let subscribers = std::mem::take(&mut log.subscribers); + let previous = log.epoch; + *log = Self::paged_log(); + if log.epoch == previous { + log.epoch = (previous + 1) & JS_MAX_SAFE_INTEGER; + } + log.subscribers = subscribers; + state.dirty.insert(session_id.to_owned()); + drop(state); + self.wake.notify_one(); + true + } + + fn paged_log() -> StreamLog { + let mut log = StreamLog::new(); + log.next_seq = HISTORY_SEQUENCE_CEILING; + log.history = Some(HistoryReader { + before: None, + exhausted: false, + pending: Vec::new(), + next_seq: HISTORY_SEQUENCE_CEILING, + newer_evicted: false, + }); + log + } + + /// Materialize only enough source records to answer this backward page. + /// The same source gate covers reads and live publication, so stale disk + /// records cannot replace newer updates or resurrect a tombstone. + pub async fn read_history( + &self, + device: &str, + request: &StreamReadRequest, + mut load: F, + ) -> Result + where + F: FnMut(Option) -> Fut, + Fut: std::future::Future>, + { + if request.after.is_some() && request.before.is_some() { + bail!("read_stream accepts either after or before, not both"); + } + let gate = self.source_gate(&request.stream_id).await; + let _source = gate.lock().await; + { + let mut state = self.lock(); + state + .streams + .entry(request.stream_id.clone()) + .or_insert_with(Self::paged_log); + let log = state.streams.get_mut(&request.stream_id).unwrap(); + let missing_history = request.after.is_none() + && log.history.as_ref().is_some_and(|history| { + history.newer_evicted + && request.before.is_none_or(|before| { + log.events + .range(..HISTORY_SEQUENCE_CEILING) + .next_back() + .is_none_or(|(&last, _)| before > last.saturating_add(1)) + }) + }); + if missing_history { + let subscribers = std::mem::take(&mut log.subscribers); + let previous = log.epoch; + *log = Self::paged_log(); + if log.epoch == previous { + log.epoch = (previous + 1) & JS_MAX_SAFE_INTEGER; + } + log.subscribers = subscribers; + state.dirty.insert(request.stream_id.clone()); + self.wake.notify_one(); + } + } + // A forward read never loads historical bodies. On epoch mismatch the + // reader observes the new fence and reopens through the latest page. + if request.after.is_some() + || request.epoch.is_some_and(|epoch| { + self.lock() + .streams + .get(&request.stream_id) + .is_some_and(|log| log.epoch != epoch) + }) + { + return self.read(device, request); + } + let started = Instant::now(); + let mut source_reads = 0; + let limit = request + .limit + .unwrap_or(DEFAULT_PAGE_EVENTS) + .clamp(1, MAX_PAGE_EVENTS); + loop { + let next = { + let mut state = self.lock(); + let log = state + .streams + .get_mut(&request.stream_id) + .context("History stream expired")?; + let before = request.before.unwrap_or(u64::MAX); + let mut count = 0; + let mut bytes = 0; + let mut page_full = false; + // Inspect sizes by reference. Cloning the growing page once per + // inserted record turns a bounded read into quadratic copying. + for (_, (_, size)) in log.events.range(..before).rev() { + if count > 0 && (count >= limit || bytes + size > PAGE_BYTES) { + page_full = true; + break; + } + count += 1; + bytes += size; + } + if count >= limit || bytes >= PAGE_BYTES || page_full { + break; + } + if log.truncated && count == 0 { + bail!("Older history exceeded the host stream memory budget; reopen the session to reload its latest page"); + } + let Some(history) = log.history.as_mut() else { + break; + }; + if let Some(mut record) = history.pending.pop() { + let id = record["id"] + .as_str() + .context("record identity missing")? + .to_owned(); + // A live value (including a tombstone) wins over backfill. + if log.record_index.contains_key(&id) || log.record_seq.contains_key(&id) { + continue; + } + let turn = record["turn"]["turnId"] + .as_str() + .context("record turn missing")? + .to_owned(); + if record["sessionId"].as_str() != Some(request.stream_id.as_str()) { + bail!("record session mismatch"); + } + let hash = record_hash(&record)?; + history.next_seq = history + .next_seq + .checked_sub(1) + .filter(|seq| *seq > 0) + .context("History sequence space exhausted")?; + let seq = history.next_seq; + record["revision"] = Value::from(seq); + let size = estimate_bytes(&record) + "session-record".len() + 32; + log.bytes += size; + log.record_seq.insert(id.clone(), seq); + log.record_index.insert(id, RecordIndexEntry { hash, turn }); + log.events.insert( + seq, + ( + StreamEvent { + seq, + event: "session-record".into(), + payload: record, + }, + size, + ), + ); + log.evict(); + continue; + } + if history.exhausted { + break; + } + history.before + }; + let batch = load(next).await?; + source_reads += 1; + if let (Some(previous), Some(next)) = (next, batch.before) { + if next >= previous { + bail!("History source did not advance"); + } + } + let mut state = self.lock(); + let history = state + .streams + .get_mut(&request.stream_id) + .and_then(|log| log.history.as_mut()) + .context("History stream expired")?; + history.before = batch.before; + history.exhausted = batch.before.is_none(); + history.pending = batch.records; + } + let page = self.read(device, request)?; + log::debug!("Read paged session history: stream_id={} source_turns={} events={} has_more={} elapsed_ms={}", + request.stream_id, source_reads, page.events.len(), page.has_more, started.elapsed().as_millis()); + Ok(page) + } + pub fn read(&self, source_device_id: &str, request: &StreamReadRequest) -> Result { if request.stream_id.is_empty() { bail!("stream id is required"); @@ -490,13 +742,13 @@ impl HostStreamHub { .into_iter() .filter(|id| id != HOST_CATALOG_ID && !id.starts_with("terminal-")) .collect(); - let events = sessions - .into_iter() - .map(|session| { + let mut events = Vec::new(); + for session in sessions { + if !self.invalidate_paged_history(&session).await { let payload = serde_json::json!({"sessionId":session,"reason":reason}); - (session, "relay://session-gap".to_string(), payload) - }) - .collect(); + events.push((session, "relay://session-gap".to_string(), payload)); + } + } self.append_batch(events).await } @@ -546,6 +798,14 @@ impl Drop for HostStreamHub { } } +fn record_hash(record: &Value) -> Result { + let body = record + .get("item") + .or_else(|| record.get("round")) + .unwrap_or(&record["turn"]); + Ok(format!("{:x}", Sha256::digest(serde_json::to_vec(body)?))) +} + /// Parent status changes get their own small header record; they must not /// re-send every unchanged large tool body. Removed records become tombstones, /// scoped to the loaded turns unless the load was a full snapshot. @@ -555,6 +815,17 @@ fn diff_records( records: Vec, full: bool, ) -> Result { + if let Some(history) = log.history.as_mut() { + let updated_turns: HashSet<_> = records + .iter() + .filter_map(|record| record["turn"]["turnId"].as_str()) + .collect(); + // The pending batch was captured before this update. All current records + // of these turns are published below; never replay removed old items. + history.pending.retain(|record| { + !updated_turns.contains(record["turn"]["turnId"].as_str().unwrap_or("")) + }); + } let mut changed = Vec::new(); let mut present = HashSet::new(); let mut scope = HashSet::new(); @@ -572,11 +843,7 @@ fn diff_records( if record["sessionId"].as_str() != Some(session_id) { bail!("record session identity mismatch"); } - let body = record - .get("item") - .or_else(|| record.get("round")) - .unwrap_or(&record["turn"]); - let hash = format!("{:x}", Sha256::digest(serde_json::to_vec(body)?)); + let hash = record_hash(&record)?; let old_hash = log.record_index.get(&id).map(|entry| entry.hash.as_str()); if old_hash != Some(hash.as_str()) { log.record_index.insert(id, RecordIndexEntry { hash, turn }); @@ -625,6 +892,271 @@ mod tests { } } + fn history_record(turn: usize, item: usize) -> Value { + serde_json::json!({"sessionId":"s", "id":format!("item/{turn}-{item}"), + "turn":{"turnId":format!("t{turn}"), "turnIndex":turn}, + "item":{"type":"text","data":{"id":format!("{turn}-{item}"),"content":"original"}}}) + } + + async fn history_fixture(before: Option) -> Result { + let turn = before.unwrap_or(100).saturating_sub(1); + Ok(HistoryBatch { + records: (0..3).map(|item| history_record(turn, item)).collect(), + before: (turn > 0).then_some(turn), + }) + } + + #[tokio::test] + async fn paged_history_reads_only_needed_turns_and_preserves_legacy_wire() { + let hub = HostStreamHub::start(Arc::new(Recorder(Default::default()))); + let reads = std::sync::atomic::AtomicUsize::new(0); + let load = |before| { + reads.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + history_fixture(before) + }; + let request = StreamReadRequest { + limit: Some(2), + ..request("s") + }; + let first = hub.read_history("phone", &request, load).await.unwrap(); + assert_eq!(reads.load(std::sync::atomic::Ordering::SeqCst), 1); + assert_eq!(first.events.len(), 2); + assert!(first.has_more); + assert_eq!(first.events[0].payload["id"], "item/99-1"); + assert!(first.cursor < JS_MAX_SAFE_INTEGER); + let wire = serde_json::to_string(&first).unwrap(); + assert_eq!(serde_json::from_str::(&wire).unwrap(), first); + let older = hub + .read_history( + "phone", + &StreamReadRequest { + before: Some(first.events[0].seq), + epoch: Some(first.epoch), + ..request.clone() + }, + load, + ) + .await + .unwrap(); + assert_eq!(reads.load(std::sync::atomic::Ordering::SeqCst), 2); + assert_eq!( + older.cursor, first.cursor, + "backfill must never advance live cursor" + ); + assert!(older + .events + .iter() + .all(|event| event.seq < first.events[0].seq)); + let forward = hub + .read_history( + "phone", + &StreamReadRequest { + after: Some(first.cursor), + epoch: Some(first.epoch), + ..request + }, + load, + ) + .await + .unwrap(); + assert!(forward.events.is_empty()); + assert!(!forward.has_more, "disk history is not forward catch-up"); + assert_eq!(reads.load(std::sync::atomic::Ordering::SeqCst), 2); + } + + #[tokio::test] + async fn paged_history_live_updates_remove_stale_pending_records() { + let hub = HostStreamHub::start(Arc::new(Recorder(Default::default()))); + let request = StreamReadRequest { + limit: Some(1), + ..request("s") + }; + let first = hub + .read_history("phone", &request, history_fixture) + .await + .unwrap(); + // Two old items still wait in the pending batch. The live turn deletes + // them, and changes the item already delivered in the first page. + let mut updated = history_record(99, 2); + updated["item"]["data"]["content"] = Value::from("updated"); + hub.synchronize_records("s".into(), false, || async { Ok(vec![updated]) }) + .await + .unwrap(); + let forward = hub + .read_history( + "phone", + &StreamReadRequest { + after: Some(first.cursor), + ..request.clone() + }, + history_fixture, + ) + .await + .unwrap(); + assert_eq!(forward.events.len(), 1); + assert!(forward.events[0].seq > first.cursor); + assert_eq!( + forward.events[0].payload["item"]["data"]["content"], + "updated" + ); + let older = hub + .read_history( + "phone", + &StreamReadRequest { + before: Some(first.events[0].seq), + ..request + }, + history_fixture, + ) + .await + .unwrap(); + assert_eq!(older.events[0].payload["turn"]["turnId"], "t98"); + } + + #[tokio::test] + async fn paged_history_undo_and_restart_fence_old_cursors() { + let hub = HostStreamHub::start(Arc::new(Recorder(Default::default()))); + let first = hub + .read_history("phone", &request("s"), history_fixture) + .await + .unwrap(); + assert!(hub.invalidate_paged_history("s").await); + let reset = hub + .read_history( + "phone", + &StreamReadRequest { + after: Some(first.cursor), + epoch: Some(first.epoch), + ..request("s") + }, + |_| async { panic!("forward read must not materialize history") }, + ) + .await + .unwrap(); + assert_ne!(reset.epoch, first.epoch); + assert_eq!(hub.subscriber_count("s"), 1); + let empty = hub + .read_history("phone", &request("s"), |_| async { + Ok(HistoryBatch { + records: vec![], + before: None, + }) + }) + .await + .unwrap(); + assert!(empty.events.is_empty()); + assert!(!empty.has_more); + hub.report_source_gap("runtime journal gap").await.unwrap(); + let gap = hub + .read_history( + "phone", + &StreamReadRequest { + after: Some(empty.cursor), + epoch: Some(empty.epoch), + ..request("s") + }, + |_| async { panic!("gap must fence before source replay") }, + ) + .await + .unwrap(); + assert_ne!(gap.epoch, empty.epoch); + } + + #[tokio::test] + async fn paged_history_eviction_preserves_backscroll_and_fences_slow_readers() { + let hub = HostStreamHub::start(Arc::new(Recorder(Default::default()))); + let load = |_| async { + Ok(HistoryBatch { + records: (0..40) + .map(|i| { + let mut record = history_record(0, i); + record["item"]["data"]["content"] = Value::from("x".repeat(1024 * 1024)); + record + }) + .collect(), + before: None, + }) + }; + let mut request = request("s"); + let mut seen = HashSet::new(); + let first = hub.read_history("phone", &request, load).await.unwrap(); + seen.insert(first.events[0].payload["id"].as_str().unwrap().to_owned()); + request.before = Some(first.events[0].seq); + request.epoch = Some(first.epoch); + loop { + let page = hub.read_history("phone", &request, load).await.unwrap(); + assert_eq!( + page.epoch, first.epoch, + "active backward reader must not loop on eviction" + ); + assert!(!page.events.is_empty()); + for event in &page.events { + assert!(seen.insert(event.payload["id"].as_str().unwrap().to_owned())); + } + if !page.has_more { + break; + } + request.before = Some(page.events[0].seq); + } + assert_eq!(seen.len(), 40); + assert!(hub.lock().streams["s"].bytes <= STREAM_BYTES_BUDGET); + let slow = hub + .read_history( + "other", + &StreamReadRequest { + before: Some(first.events[0].seq), + epoch: Some(first.epoch), + ..request.clone() + }, + load, + ) + .await + .unwrap(); + assert_ne!( + slow.epoch, first.epoch, + "never silently skip evicted records" + ); + let reopened = hub + .read_history("other", &super::tests::request("s"), load) + .await + .unwrap(); + assert_eq!( + reopened.events[0].payload["id"], + first.events[0].payload["id"] + ); + } + + #[tokio::test] + async fn paged_history_large_turn_resumes_without_reloading_it() { + let hub = HostStreamHub::start(Arc::new(Recorder(Default::default()))); + let reads = std::sync::atomic::AtomicUsize::new(0); + let load = |_| { + reads.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + async { + Ok(HistoryBatch { + records: (0..1000).map(|i| history_record(0, i)).collect(), + before: None, + }) + } + }; + let mut request = request("s"); + let mut seen = std::collections::HashSet::new(); + loop { + let page = hub.read_history("phone", &request, load).await.unwrap(); + assert!(!page.events.is_empty()); + for event in &page.events { + assert!(seen.insert(event.payload["id"].as_str().unwrap().to_owned())); + } + if !page.has_more { + break; + } + request.before = Some(page.events[0].seq); + request.epoch = Some(page.epoch); + } + assert_eq!(seen.len(), 1000); + assert_eq!(reads.load(std::sync::atomic::Ordering::SeqCst), 1); + } + #[tokio::test] async fn streams_materialize_on_first_read_and_notify_subscribers() { let recorder = Arc::new(Recorder(Default::default()));