diff --git a/AGENTS.md b/AGENTS.md index 6a56f979..9b50199b 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -55,23 +55,17 @@ deps/ codex/ # acking-you/codex fork, branch `pocket-codex` # (git submodule) = upstream openai/codex main + # our adaptations; see §8 - pb-mapper/ # upstream pb-mapper (git submodule) - kanal/ # fork pinned to a known-good commit; transitively - # required by pb-mapper, redirected via [patch] - uni-stream/ # ditto; transitively required by pb-mapper docs/ # design notes, protocol references, CLI verification scripts/ # install scripts, local CI, CI affected-surface gate ``` `Cargo.toml` is a workspace root; every crate under `crates/` is a workspace member (see the `members` list for the canonical set). -Submodules under `deps/` are kept **out** of the workspace via the -`exclude` list — the pinned -upstream crates use their own lints/profiles and we depend on them -through explicit path or git deps where needed. The root manifest's -`[patch]` table redirects `acking-you/kanal` and `acking-you/uni-stream` -to the local submodules so the build stays reproducible across -contributor checkouts and CI even after the upstream forks evolve. +The Codex submodule under `deps/` is kept **out** of the workspace via the +`exclude` list and retains its upstream lints/profiles. Cargo fetches pb-mapper +from its dedicated `pocket-codex` Git branch; `Cargo.lock` pins the exact commit +and its registry dependencies. No pb-mapper, kanal, or uni-stream submodules +or local dependency patches are needed. ## 3. Crate responsibilities @@ -81,7 +75,7 @@ Shared / host side: | --------------------------- | ---------------------------------------------------------------------------------------------- | | `pocket-codex-core` | configuration schema, on-disk `state.toml`, well-known paths, error types, `service::{ServiceId, ServiceKind, sanitize_component, default_device_id}` for `pcx:::` relay keys — small, dependency-light | | `pocket-codex-codex` | spawning / supervising / inspecting the `codex app-server` child process (out-of-process *and* the in-process `embedded-codex` path), JSON-RPC envelope types | -| `pocket-codex-pb` | async wrappers around the published `pb-mapper` client SDK: `RelaySession` (address + credential), register / subscribe / status, `publish` (and the one name-conflict failure a caller must not retry), admin credential issuance, and credential keep-alive | +| `pocket-codex-pb` | async wrappers around the Git-pinned `pb-mapper` client SDK: `RelaySession` (address + credential), register / subscribe / status, `publish` (and the one name-conflict failure a caller must not retry), admin credential issuance, and credential keep-alive | | `pocket-codex-api-proxy` | local Responses API proxy: forwards `/v1/responses` (HTTP + WS) to ChatGPT's Codex backend, reusing the host's `codex login`; shared by the CLI worker and the in-app host | | `pocket-codex-host-svc` | host-side meta service — remote-viewable codex sessions, per-thread config, attachment upload — published on the relay as a third `meta:` service | | `pocket-codex-cli` | user-facing `pocket-codex` binary; account (`login` / `logout` / `account`), setup (`init`), high-level `serve` / `connect` / `api {serve,connect}` / `services {list,default set}` / `status` / `stop`, low-level `codex {start,stop,status}`, `pb {register,subscribe,status}`, `remote-hint`, `version` | @@ -233,9 +227,12 @@ git checkout -- apps/flutter/pubspec.yaml # restore before committing `deps/codex` is a git submodule pinned to a specific commit — the only one left. `deps/pb-mapper` (plus the `deps/kanal` and `deps/uni-stream` forks it -pulled in transitively) is gone: pb-mapper is a registry dependency now, so its -own transitive pins come from the lockfile rather than a mirrored `[patch]` -table. After pulling this repo, materialise the submodule with: +pulled in transitively) is gone. pb-mapper is a Cargo Git dependency on +`https://github.com/acking-you/pb-mapper`, branch `pocket-codex`; `Cargo.lock` +pins its exact commit. Push SDK fixes to that branch, then run +`cargo update -p pb-mapper` here and commit the resulting lockfile after +verification. Normal builds use `--locked` and do not automatically follow +branch updates. After pulling this repo, materialise the Codex submodule with: ```bash git submodule update --init --recursive diff --git a/Cargo.lock b/Cargo.lock index a95a6efd..109d3538 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -387,7 +387,7 @@ version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -398,7 +398,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" dependencies = [ "anstyle", "once_cell_polyfill", - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -4710,7 +4710,7 @@ dependencies = [ "libc", "option-ext", "redox_users 0.5.2", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -5022,7 +5022,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -7022,7 +7022,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.6.3", + "socket2 0.5.10", "system-configuration", "tokio", "tower-service", @@ -7547,7 +7547,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" dependencies = [ "hermit-abi", "libc", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -8430,7 +8430,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -9112,8 +9112,7 @@ checksum = "df94ce210e5bc13cb6651479fa48d14f601d9858cfe0467f43ae157023b938d3" [[package]] name = "pb-mapper" version = "0.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a54a66e600806660cdaf43ef091c31b5bceb1f25eefad309ae1d0420b23567ea" +source = "git+https://github.com/acking-you/pb-mapper?branch=pocket-codex#f0ed4271de962b8759d11171bcf901aa6df458cc" dependencies = [ "pb-mapper-client", ] @@ -9121,8 +9120,7 @@ dependencies = [ [[package]] name = "pb-mapper-auth" version = "0.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "82fee20467d2c38ba1a3a3d5736f7a43c0516aa899f673c4c507abfe4b7ed6d3" +source = "git+https://github.com/acking-you/pb-mapper?branch=pocket-codex#f0ed4271de962b8759d11171bcf901aa6df458cc" dependencies = [ "parking_lot", "pb-mapper-core", @@ -9139,8 +9137,7 @@ dependencies = [ [[package]] name = "pb-mapper-client" version = "0.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "aff91f22c290b1ddcf9c7252c2ff21927951ab2a47334a253f9e4fcb46118359" +source = "git+https://github.com/acking-you/pb-mapper?branch=pocket-codex#f0ed4271de962b8759d11171bcf901aa6df458cc" dependencies = [ "pb-mapper-auth", "pb-mapper-core", @@ -9156,8 +9153,7 @@ dependencies = [ [[package]] name = "pb-mapper-core" version = "0.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9337e59b81313428c219ea20b3f639d2bd457c18d493ee8549eb13015f753554" +source = "git+https://github.com/acking-you/pb-mapper?branch=pocket-codex#f0ed4271de962b8759d11171bcf901aa6df458cc" dependencies = [ "base64 0.23.1", "clap", @@ -9175,8 +9171,7 @@ dependencies = [ [[package]] name = "pb-mapper-protocol" version = "0.5.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3e0c9f49be8e717f80a270f9f398a75aac711e192fcca3d9572d1f5fb3533f40" +source = "git+https://github.com/acking-you/pb-mapper?branch=pocket-codex#f0ed4271de962b8759d11171bcf901aa6df458cc" dependencies = [ "bytes", "parking_lot", @@ -10064,7 +10059,7 @@ dependencies = [ "quinn-udp", "rustc-hash 2.1.2", "rustls", - "socket2 0.6.3", + "socket2 0.5.10", "thiserror 2.0.18", "tokio", "tracing", @@ -10101,7 +10096,7 @@ dependencies = [ "cfg_aliases 0.2.1", "libc", "once_cell", - "socket2 0.6.3", + "socket2 0.5.10", "tracing", "windows-sys 0.60.2", ] @@ -10938,7 +10933,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys 0.12.1", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -11846,7 +11841,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3a766e1110788c36f4fa1c2b71b387a7815aa65f88ce0229841826633d93723e" dependencies = [ "libc", - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -12558,7 +12553,7 @@ dependencies = [ "getrandom 0.4.2", "once_cell", "rustix 1.1.4", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -12577,7 +12572,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "230a1b821ccbd75b185820a1f1ff7b14d21da1e442e22c0863ea5f08771a8874" dependencies = [ "rustix 1.1.4", - "windows-sys 0.61.2", + "windows-sys 0.59.0", ] [[package]] @@ -13407,7 +13402,7 @@ checksum = "f2f6fb2847f6742cd76af783a2a2c49e9375d0a111c7bef6f71cd9e738c72d6e" dependencies = [ "memoffset", "tempfile", - "windows-sys 0.61.2", + "windows-sys 0.60.2", ] [[package]] @@ -13991,7 +13986,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.61.2", + "windows-sys 0.48.0", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index e47d6695..84aeae1b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,8 +20,8 @@ members = [ # # `deps/pb-mapper` used to sit here too, alongside the `deps/kanal` and # `deps/uni-stream` forks it pulled in transitively. All three are gone: -# pb-mapper is a registry dependency now, so its own transitive pins are -# its business rather than something this workspace has to mirror. +# pb-mapper is fetched by Cargo from its dedicated Git branch; its transitive +# dependencies are locked here without vendoring additional submodules. exclude = [ "deps/codex", "apps/flutter", @@ -41,12 +41,9 @@ readme = "README.md" # from here via `dep = { workspace = true }` so we have a single source # of truth for upgrades. [workspace.dependencies] -# The pb-mapper client SDK, from the registry rather than the git submodule it -# used to be. Pinned here so every consumer moves together: a client older than -# the deployed relay cannot decode the relay's structured error frames, which is -# how the 0.2.14-client/0.5.0-relay skew turned an over-quota refusal into an -# unbacked-off reconnect storm. -pb-mapper = "0.5.0" +# The dedicated branch carries the SDK fixes used by Pocket-Codex. Cargo.lock +# pins its exact commit; update it deliberately with `cargo update -p pb-mapper`. +pb-mapper = { git = "https://github.com/acking-you/pb-mapper", branch = "pocket-codex" } # Internal crates pocket-codex-core = { path = "crates/pocket-codex-core" } @@ -197,8 +194,8 @@ opt-level = 2 # pb-mapper submodule. They existed because that submodule declared both as git # deps on `acking-you/*` branch HEADs, and one of those HEADs renamed a crate out # from under us — so both forks were vendored and re-routed to keep the build -# reproducible. A registry dependency resolves its own transitive pins from the -# lockfile, so none of that is our problem any more. +# reproducible. The current SDK uses registry dependencies for both crates, +# whose versions are recorded in Cargo.lock without local patches. # codex (deps/codex) requires forked tokio-tungstenite / tungstenite (they add a # `proxy` feature on the 0.28 line). The optional `embedded-codex` feature pulls diff --git a/README.md b/README.md index 1ecbbf89..44c8b729 100644 --- a/README.md +++ b/README.md @@ -68,7 +68,7 @@ or a per-account GitHub login (hosted). | ------------------------------ | -------------------------------------- | | Workspace / lints / CI | bootstrapped | | `pocket-codex` CLI | `login`, `logout`, `account`, `init`, `serve`, `connect`, `api {serve,connect}`, `services {list,default set}`, top-level `status`/`stop`, `codex {start,stop,status}`, `pb {register,subscribe,status}`, `remote-hint`, `version` | -| `pb-mapper` register/subscribe | the published `pb-mapper` client SDK | +| `pb-mapper` register/subscribe | Git SDK from the `pocket-codex` branch, pinned by `Cargo.lock`; includes connection timeout and interactive TCP latency fixes | | `codex app-server` supervision | spawn/stop/status via PID + state.toml | | App-server protocol | synced to upstream main `db0568dbb` (2026-09-07); acknowledged initialization, v2 account reads, current thread model/effort, and asynchronous question replies | | Embedded codex (desktop) | desktop builds compile codex **in-process** behind the `embedded-codex` feature, so a machine can host without a separate `codex` install (Windows/macOS) | diff --git a/apps/flutter/lib/src/screens/app_session_screen.dart b/apps/flutter/lib/src/screens/app_session_screen.dart index 3d127e2f..5a0c873c 100644 --- a/apps/flutter/lib/src/screens/app_session_screen.dart +++ b/apps/flutter/lib/src/screens/app_session_screen.dart @@ -422,6 +422,7 @@ class _AppSessionState extends ConsumerState bool _settlingToEnd = false; String? _error; VoidCallback? _retry; // action for the error banner's retry button + int _threadLoadGeneration = 0; bool _connectionLost = false; // True while an automatic reconnect is in progress (drives the status bar's // "reconnecting" state). Auto-reconnect is attempted on stream close, on a @@ -967,6 +968,7 @@ class _AppSessionState extends ConsumerState ref.read(uiPrefsProvider.notifier).setLastThread(widget.serviceKey, tid); } _cancelExternalWriterSubscription(); + _threadLoadGeneration++; setState(() { _threadId = tid; _cwd = cwd; @@ -1369,9 +1371,7 @@ class _AppSessionState extends ConsumerState ]; } - /// Open an existing thread: resume it into the session (so reads and turns - /// resolve — otherwise the server returns "thread not found"), then load - /// its history. + /// Attach to an existing thread for live events and turns, then load history. Future _resumeAndLoad() async { // Guard: a stale event (e.g. thread/compacted from a prior thread) can // arrive after switching to a new, unsaved conversation — don't `_threadId!` @@ -1383,16 +1383,23 @@ class _AppSessionState extends ConsumerState _retry = null; }); final startTid = _threadId!; + final generation = ++_threadLoadGeneration; + bool current() => + mounted && _threadId == startTid && generation == _threadLoadGeneration; try { final api = ref.read(bridgeApiProvider); await api.appThreadResume(widget.serviceKey, startTid); + // An obsolete resume must not fan out into more history/config requests. + if (!current()) return; // Read the thread history and its persisted config concurrently. The // config is best-effort (an unreachable host meta tunnel yields an // all-unset config and we fall back to the server / in-memory restore). final historyFuture = api.appThreadRead(widget.serviceKey, startTid); final persistedFuture = _loadPersistedConfig(startTid); final history = await historyFuture; + if (!current()) return; final persisted = await persistedFuture; + if (!current()) return; // Restore the model from the server's own report first (the resume // response says what the thread actually runs with); fall back to the // persisted pick for older servers that don't report one. Resolve the id @@ -1410,7 +1417,7 @@ class _AppSessionState extends ConsumerState } } // The user may have switched threads during the awaits above. - if (!mounted || _threadId != startTid) return; + if (!current()) return; setState(() { _loading = false; _replaceTranscriptItems(history.items); @@ -1533,7 +1540,7 @@ class _AppSessionState extends ConsumerState // events that would normally flush it were missed during the drop). _maybeFlushQueue(); } catch (e) { - if (!mounted || _threadId != startTid) return; + if (!current()) return; if (_isActiveWriterError(e)) { _enterExternalWriterMode(startTid); return; diff --git a/apps/flutter/test/fake_bridge_api.dart b/apps/flutter/test/fake_bridge_api.dart index 5b85ab08..63b174d0 100644 --- a/apps/flutter/test/fake_bridge_api.dart +++ b/apps/flutter/test/fake_bridge_api.dart @@ -661,6 +661,9 @@ class FakeBridgeApi implements BridgeApi { /// Records the last resumed thread id for assertions. String? lastResumed; + final Map>> pendingResumes = {}; + final Map>> pendingReads = {}; + final List threadReads = []; /// Optional failure thrown by [appThreadResume]. Object? appThreadResumeError; @@ -669,6 +672,8 @@ class FakeBridgeApi implements BridgeApi { Future appThreadResume(String serviceKey, String threadId) async { if (appThreadResumeError != null) throw appThreadResumeError!; lastResumed = threadId; + final pending = pendingResumes[threadId]; + if (pending != null && pending.isNotEmpty) await pending.removeAt(0); } /// Seedable history for resume tests. @@ -678,7 +683,12 @@ class FakeBridgeApi implements BridgeApi { Future appThreadRead( String serviceKey, String threadId, - ) async => readResult; + ) async { + threadReads.add(threadId); + final pending = pendingReads[threadId]; + if (pending != null && pending.isNotEmpty) return await pending.removeAt(0); + return readResult; + } /// Older pages a paginated thread hands back, oldest batch LAST — each call /// to [appThreadOlderPage] pops the last one, so seeding diff --git a/apps/flutter/test/screens/app_session_test.dart b/apps/flutter/test/screens/app_session_test.dart index f80bd80f..3db07852 100644 --- a/apps/flutter/test/screens/app_session_test.dart +++ b/apps/flutter/test/screens/app_session_test.dart @@ -4768,6 +4768,80 @@ void main() { expect(find.text('past chat'), findsNothing); }); + for (final stage in ['resume', 'history']) { + testWidgets('rapid A B A switching ignores obsolete $stage results', ( + t, + ) async { + final api = FakeBridgeApi( + config: const ConfigInfo(relay: 'lb7666.top:7666', hasKey: true), + ); + await api.appConnect('pcx:lb7666:app:default', 28080); + api.appThreads.addAll([ + const ThreadMeta(id: 'a', preview: 'chat A', cwd: '', updatedAt: 0), + const ThreadMeta(id: 'b', preview: 'chat B', cwd: '', updatedAt: 0), + ]); + final resume = Completer(); + final history = Completer(); + if (stage == 'resume') api.pendingResumes['a'] = [resume.future]; + if (stage == 'history') api.pendingReads['a'] = [history.future]; + api.readResult = const ThreadHistory( + items: [ + ThreadItem( + id: 'fresh', + itemType: 'agentMessage', + title: '', + text: 'fresh transcript', + ), + ], + running: false, + ); + t.view.devicePixelRatio = 1; + t.view.physicalSize = const Size(1200, 900); + addTearDown(t.view.reset); + await t.pumpWidget( + host(const AppSessionScreen(serviceKey: 'pcx:lb7666:app:default'), api), + ); + await t.pumpAndSettle(); + await t.tap(find.text('chat A')); + await t.pump(const Duration(milliseconds: 50)); + await t.tap(find.text('chat B')); + await t.pumpAndSettle(); + await t.tap(find.text('chat A')); + await t.pumpAndSettle(); + if (stage == 'resume') { + resume.complete(); + } else { + history.complete( + const ThreadHistory( + items: [ + ThreadItem( + id: 'stale', + itemType: 'agentMessage', + title: '', + text: 'obsolete transcript', + ), + ], + running: false, + ), + ); + } + await t.pumpAndSettle(); + expect( + find.textContaining('fresh transcript', findRichText: true), + findsOneWidget, + ); + expect( + find.textContaining('obsolete transcript', findRichText: true), + findsNothing, + ); + expect( + api.threadReads.where((id) => id == 'a').length, + stage == 'resume' ? 1 : 2, + ); + expect(t.takeException(), isNull); + }); + } + testWidgets('Sessions pane buttons switch threads without crashing', ( t, ) async { diff --git a/crates/pocket-codex-backend/src/credentials.rs b/crates/pocket-codex-backend/src/credentials.rs index b38cac65..7c6a48f4 100644 --- a/crates/pocket-codex-backend/src/credentials.rs +++ b/crates/pocket-codex-backend/src/credentials.rs @@ -50,15 +50,15 @@ struct Cached { expires_at: u64, } +type AccountCache = Arc>>; + /// Mints and caches one relay credential per account. #[derive(Clone)] pub struct Credentials { relay: RelaySession, - /// Keyed by internal user id. A mutex rather than a lock-free map because - /// the contended path is a relay round trip, and holding the lock - /// across it is what stops a thundering herd of first-time requests - /// from minting one credential each. - cache: Arc>>, + /// Serialize credential changes per account. A slow relay operation must + /// not block unrelated accounts from reading their cached credentials. + cache: Arc>>, } impl Credentials { @@ -83,17 +83,29 @@ impl Credentials { /// admin-side listing or retire to one account. `None` is not an error: an /// account with no credential has no services either. pub async fn namespace_of(&self, user_id: &str) -> anyhow::Result> { - if let Some(cached) = self.cache.lock().await.get(user_id) { - return Ok(Some(cached.key_id)); + let account = self.account(user_id).await; + let mut cached = account.lock().await; + if cached.is_none() { + *cached = self.adopt_from_relay(user_id).await?; } - Ok(self.adopt_from_relay(user_id).await?.map(|c| c.key_id)) + Ok(cached.as_ref().map(|entry| entry.key_id)) + } + + async fn account(&self, user_id: &str) -> AccountCache { + self.cache + .lock() + .await + .entry(user_id.to_string()) + .or_default() + .clone() } /// This account's credential, minting or renewing it if needed. /// /// Returns `(credential, expires_at)`. pub async fn for_account(&self, user_id: &str) -> anyhow::Result<(String, u64)> { - let mut cache = self.cache.lock().await; + let account = self.account(user_id).await; + let mut cache = account.lock().await; let now = now_secs(); // Nothing in memory does NOT mean the account has no credential — this // process may just have restarted. Adopting the live one matters because @@ -101,12 +113,12 @@ impl Credentials { // account, so devices still holding the old credential and devices // fetching after the restart would land in different namespaces and stop // seeing each other. - if !cache.contains_key(user_id) { + if cache.is_none() { if let Some(adopted) = self.adopt_from_relay(user_id).await? { - cache.insert(user_id.to_string(), adopted); + *cache = Some(adopted); } } - if let Some(cached) = cache.get(user_id) { + if let Some(cached) = cache.as_ref() { if cached.expires_at > now + RENEW_MARGIN.as_secs() && !cached.credential.is_empty() { return Ok((cached.credential.clone(), cached.expires_at)); } @@ -123,19 +135,22 @@ impl Credentials { key_id: renewed.key_id, expires_at: renewed.expires_at, }; - cache.insert(user_id.to_string(), entry.clone()); + *cache = Some(entry.clone()); return Ok((entry.credential, entry.expires_at)); }, Err(err) => { - // Expired past renewal, revoked, or lost to a state reset. - // Falling through to a fresh mint is the only way back, and - // it is worth a log line because the account's namespace - // changes with it. + // A failed renewal can be a lost response or a temporary + // relay failure. Only a successful listing that confirms + // this key is gone justifies changing the namespace. + let live = pocket_codex_pb::live_credentials(&self.relay).await?; + if live.iter().any(|key| key.key_id == cached.key_id) { + return Err(err); + } tracing::warn!( user = %user_id, key_id = cached.key_id, error = %format!("{err:#}"), - "renewing the relay credential failed; minting a new one" + "relay confirmed the old credential is gone; minting a new one" ); }, } @@ -151,7 +166,7 @@ impl Credentials { key_id: issued.key_id, expires_at: issued.expires_at, }; - cache.insert(user_id.to_string(), entry.clone()); + *cache = Some(entry.clone()); Ok((entry.credential, entry.expires_at)) } @@ -162,10 +177,10 @@ impl Credentials { /// forgotten satisfies it. Dropping the cache entry is the part that must /// not be skipped, so it happens regardless. pub async fn revoke_account(&self, user_id: &str) { - let cached = match self.cache.lock().await.remove(user_id) { + let account = self.account(user_id).await; + let mut entry = account.lock().await; + let cached = match entry.take() { Some(cached) => Some(cached), - // Not in memory: it may still be live on the relay from before a - // restart, and "revoked" has to mean revoked. None => self.adopt_from_relay(user_id).await.ok().flatten(), }; let Some(cached) = cached else { return }; @@ -234,3 +249,7 @@ fn now_secs() -> u64 { .map(|d| d.as_secs()) .unwrap_or_default() } + +#[cfg(test)] +#[path = "credentials_tests.rs"] +mod tests; diff --git a/crates/pocket-codex-backend/src/credentials_tests.rs b/crates/pocket-codex-backend/src/credentials_tests.rs new file mode 100644 index 00000000..60778bae --- /dev/null +++ b/crates/pocket-codex-backend/src/credentials_tests.rs @@ -0,0 +1,84 @@ +use super::*; + +fn cached() -> Cached { + Cached { + credential: "cached-test-credential".into(), + key_id: 7, + expires_at: now_secs() + 7200, + } +} + +#[tokio::test] +async fn a_busy_account_does_not_block_other_cached_accounts() { + let vendor = Credentials::new(RelaySession::for_test("127.0.0.1:1")); + let busy = vendor.account("busy").await; + let _busy = busy.lock().await; + *vendor.account("ready").await.lock().await = Some(cached()); + let (credential, _) = + tokio::time::timeout(Duration::from_millis(100), vendor.for_account("ready")) + .await + .expect("unrelated account must not wait") + .expect("cached credential"); + assert_eq!(credential, "cached-test-credential"); +} + +#[tokio::test] +async fn renewal_transport_failure_preserves_the_cached_namespace() { + let vendor = Credentials::new(RelaySession::for_test("127.0.0.1:1")); + let mut entry = cached(); + entry.expires_at = now_secs() + 60; + *vendor.account("user").await.lock().await = Some(entry); + assert!(vendor.for_account("user").await.is_err()); + assert_eq!(vendor.namespace_of("user").await.expect("cached namespace"), Some(7)); +} + +#[tokio::test] +async fn a_lost_renewal_reply_never_mints_a_second_namespace() { + use tokio::net::{TcpListener, TcpStream}; + let (Ok(addr), Ok(key)) = + (std::env::var("PCX_TEST_RELAY"), std::env::var("PCX_TEST_RELAY_KEY")) + else { + eprintln!("skipping lost-reply test: PCX_TEST_RELAY / PCX_TEST_RELAY_KEY unset"); + return; + }; + let relay = RelaySession::new(addr.clone(), key.clone()); + let issued = pocket_codex_pb::issue_credential( + &relay, + Duration::from_secs(7200), + Some("lost-renewal-test".into()), + ) + .await + .expect("issue fixture"); + let listener = TcpListener::bind("127.0.0.1:0").await.expect("proxy"); + let vendor = Credentials::new(RelaySession::new( + listener.local_addr().expect("address").to_string(), + key, + )); + *vendor.account("user").await.lock().await = Some(Cached { + credential: issued.credential, + key_id: issued.key_id, + expires_at: now_secs() + 60, + }); + let proxy = tokio::spawn(async move { + // Lose the renewal connection, then allow the authoritative key listing. + drop(listener.accept().await.expect("renewal").0); + let (mut downstream, _) = listener.accept().await.expect("key listing"); + let mut upstream = TcpStream::connect(addr).await.expect("relay"); + tokio::io::copy_bidirectional(&mut downstream, &mut upstream) + .await + .expect("forward listing"); + listener + }); + assert!(vendor.for_account("user").await.is_err()); + let listener = proxy.await.expect("proxy task"); + assert!( + tokio::time::timeout(Duration::from_millis(100), listener.accept()) + .await + .is_err(), + "must not mint after an ambiguous failure" + ); + assert_eq!(vendor.namespace_of("user").await.expect("namespace"), Some(issued.key_id)); + pocket_codex_pb::revoke_credential(&relay, issued.key_id) + .await + .expect("cleanup"); +} diff --git a/crates/pocket-codex-bridge/src/engine/app_session.rs b/crates/pocket-codex-bridge/src/engine/app_session.rs index f9f9a236..ea302e0b 100644 --- a/crates/pocket-codex-bridge/src/engine/app_session.rs +++ b/crates/pocket-codex-bridge/src/engine/app_session.rs @@ -221,16 +221,36 @@ fn sessions() -> &'static Mutex> { SESSIONS.get_or_init(|| Mutex::new(HashMap::new())) } +// Weak entries serialize connect/disconnect for one service without keeping +// obsolete services alive or holding the session registry over network I/O. +fn lifecycle(service_key: &str) -> Arc> { + type Locks = Mutex>>>; + static LOCKS: OnceCell = OnceCell::new(); + let mut locks = LOCKS + .get_or_init(Mutex::default) + .lock() + .unwrap_or_else(|e| e.into_inner()); + if let Some(lock) = locks.get(service_key).and_then(std::sync::Weak::upgrade) { + return lock; + } + locks.retain(|_, lock| lock.strong_count() > 0); + let lock = Arc::new(Mutex::new(())); + locks.insert(service_key.to_string(), Arc::downgrade(&lock)); + lock +} + /// Subscribe to `service_key` (materialising the local ws endpoint), open a /// JSON-RPC client over it and run the `initialize` handshake. Idempotent: a /// live session for the same key is reused. pub fn connect(service_key: String, local_port: u16, transport: &Transport) -> Result<()> { + let lifecycle = lifecycle(&service_key); + let _guard = lifecycle.lock().unwrap_or_else(|e| e.into_inner()); if reuse_live(&service_key) { return Ok(()); } // No live session: drop any stale one (and its subscription) so we reconnect // cleanly rather than reusing a closed socket. - disconnect(&service_key); + disconnect_inner(&service_key); // Materialise the local ws endpoint via pb-mapper (kind-agnostic subscribe). let sub = runtime::subscribe_service(service_key.clone(), local_port, transport)?; establish(service_key, &sub.local_addr) @@ -725,6 +745,12 @@ pub fn is_connected(service_key: &str) -> bool { /// Drop the session for `service_key` and its pb-mapper subscription. pub fn disconnect(service_key: &str) { + let lifecycle = lifecycle(service_key); + let _guard = lifecycle.lock().unwrap_or_else(|e| e.into_inner()); + disconnect_inner(service_key); +} + +fn disconnect_inner(service_key: &str) { sessions() .lock() .expect("sessions poisoned") @@ -1225,8 +1251,9 @@ fn track_pending_approval(pending: &Mutex>, inbound: &Inb /// view's gists are what break first. pub fn thread_resume(service_key: &str, thread_id: &str) -> Result<()> { let client = client_for(service_key)?; - let res = runtime::runtime() - .block_on(client.request("thread/resume", json!({ "threadId": thread_id })))?; + let res = runtime::runtime().block_on( + client.request("thread/resume", json!({ "threadId": thread_id, "excludeTurns": true })), + )?; // The resume response carries the thread's effective runtime config — // model, modelProvider, reasoningEffort, approvalPolicy, sandbox — none of // which `thread/read` exposes. Refresh the cache unconditionally: the diff --git a/crates/pocket-codex-bridge/src/engine/app_session_pagination_tests.rs b/crates/pocket-codex-bridge/src/engine/app_session_pagination_tests.rs index b03e89cb..0769a7bb 100644 --- a/crates/pocket-codex-bridge/src/engine/app_session_pagination_tests.rs +++ b/crates/pocket-codex-bridge/src/engine/app_session_pagination_tests.rs @@ -59,6 +59,12 @@ fn mock_client( let request: Value = serde_json::from_str(&frame).expect("pagination test operation"); assert_eq!(request["method"], method); + if method == "thread/resume" { + assert_eq!( + request["params"]["excludeTurns"], true, + "resume must not hydrate full history" + ); + } socket .send(Message::text(json!({"id": request["id"], "result": result}).to_string())) .await @@ -228,3 +234,58 @@ fn pending_summary_yields_on_a_single_async_worker() { }); runtime::runtime().block_on(peer).expect("test peer"); } + +#[test] +fn resume_requests_metadata_and_retains_the_runtime_configuration() { + runtime::init(std::env::temp_dir()).expect("init runtime"); + let (client, peer) = mock_client(vec![( + "thread/resume", + json!({ + "thread": {"id": "thread"}, "model": "test-model", "reasoningEffort": "high" + }), + )]); + let session = TestSession::new(client); + thread_resume(&session.0, "thread").expect("resume"); + assert_eq!( + thread_runtime_config(&session.0, "thread") + .expect("config") + .model + .as_deref(), + Some("test-model") + ); + runtime::runtime().block_on(peer).expect("peer"); +} + +#[test] +#[ignore = "manual: PCX_SOAK_WS and PCX_SOAK_THREAD_IDS select existing idle threads"] +fn real_session_switch_soak() { + use std::time::Instant; + runtime::init(std::env::temp_dir()).expect("init runtime"); + let addr = std::env::var("PCX_SOAK_WS").expect("PCX_SOAK_WS host:port"); + let threads = std::env::var("PCX_SOAK_THREAD_IDS").expect("comma-separated idle thread ids"); + let threads: Vec<&str> = threads.split(',').collect(); + let rounds = std::env::var("PCX_SOAK_ROUNDS") + .ok() + .and_then(|value| value.parse().ok()) + .unwrap_or(100); + let session = TestSession("session-switch-soak".into()); + establish(session.0.clone(), &addr).expect("connect"); + let mut times = Vec::new(); + for i in 0..rounds { + let thread = threads[i % threads.len()]; + let start = Instant::now(); + thread_resume(&session.0, thread).expect("resume"); + let history = thread_read(&session.0, thread).expect("history"); + assert!(!history.items.is_empty(), "fixture should include a transcript"); + assert!(is_connected(&session.0), "same connection survives every switch"); + times.push(start.elapsed()); + } + times.sort(); + eprintln!( + "{rounds} switches, p50={:?}, p95={:?}, p99={:?}, max={:?}", + times[rounds / 2], + times[rounds * 95 / 100], + times[rounds * 99 / 100], + times.last().expect("samples") + ); +} diff --git a/crates/pocket-codex-codex/src/client.rs b/crates/pocket-codex-codex/src/client.rs index a8a755a3..4bca2f29 100644 --- a/crates/pocket-codex-codex/src/client.rs +++ b/crates/pocket-codex-codex/src/client.rs @@ -35,7 +35,9 @@ use tokio::{ task::JoinHandle, }; use tokio_tungstenite::{ - connect_async, tungstenite::Message as WsMessage, MaybeTlsStream, WebSocketStream, + connect_async_with_config, + tungstenite::{protocol::WebSocketConfig, Message as WsMessage}, + MaybeTlsStream, WebSocketStream, }; use tokio_util::sync::CancellationToken; @@ -55,6 +57,9 @@ pub struct Inbound { pub request_id: Option, } +/// Largest accepted JSON-RPC message and frame. +const MAX_MESSAGE_BYTES: usize = 64 << 20; + /// Default per-request timeout. A model turn streams via notifications, so /// individual request/response round-trips (initialize, thread/start, …) are /// quick; 60s is generous headroom for a slow relay hop. @@ -87,6 +92,7 @@ pub struct AppClient { reader: JoinHandle<()>, keepalive: JoinHandle<()>, closed: CancellationToken, + close_reason: Arc>>, } impl Drop for AppClient { @@ -100,7 +106,12 @@ impl AppClient { /// Connect to `ws_url` (e.g. `ws://127.0.0.1:28080`) and start the reader /// task. Returns the client plus the receiver of inbound notifications. pub async fn connect(ws_url: &str) -> Result<(Self, mpsc::UnboundedReceiver)> { - let (stream, _resp) = connect_async(ws_url) + // A JSON-RPC response is one frame. Match the message limit so valid + // legacy histories between 16 and 64 MiB do not close the connection. + let config = WebSocketConfig::default() + .max_frame_size(Some(MAX_MESSAGE_BYTES)) + .max_message_size(Some(MAX_MESSAGE_BYTES)); + let (stream, _resp) = connect_async_with_config(ws_url, Some(config), true) .await .with_context(|| format!("connecting app-server websocket {ws_url}"))?; let (sink, mut read) = stream.split(); @@ -115,18 +126,23 @@ impl AppClient { // live-but-quiet socket from a dead half-open one. let activity = Arc::new(AtomicU64::new(0)); let closed = CancellationToken::new(); + let close_reason = Arc::new(StdMutex::new(None)); let reader_pending = Arc::clone(&pending); let reader_server_reqs = Arc::clone(&server_reqs); let reader_activity = Arc::clone(&activity); let reader_closed = closed.clone(); + let reader_reason = close_reason.clone(); let reader = tokio::spawn(async move { loop { let frame = tokio::select! { _ = reader_closed.cancelled() => break, frame = read.next() => match frame { Some(frame) => frame, - None => break, + None => { + record_close_reason(&reader_reason, "peer ended the websocket stream".into()); + break; + }, }, }; // Any frame — including the Pong answering our keepalive Ping — @@ -135,7 +151,17 @@ impl AppClient { let text = match frame { Ok(WsMessage::Text(t)) => t.to_string(), Ok(WsMessage::Binary(b)) => String::from_utf8_lossy(&b).into_owned(), - Ok(WsMessage::Close(_)) | Err(_) => break, + Ok(WsMessage::Close(frame)) => { + record_close_reason(&reader_reason, format!("peer sent close: {frame:?}")); + break; + }, + Err(error) => { + record_close_reason( + &reader_reason, + format!("websocket read failed: {error}"), + ); + break; + }, // Ping/Pong/Frame: nothing to dispatch. Ok(_) => continue, }; @@ -181,7 +207,7 @@ impl AppClient { } // Connection closed: fail every in-flight request so callers don't // hang on a oneshot that will never resolve. - close_connection(&reader_closed, &reader_pending); + close_connection(&reader_closed, &reader_pending, &reader_reason); }); // Keepalive + liveness watchdog. Each tick pings (keeping the relay @@ -196,6 +222,7 @@ impl AppClient { let keepalive_activity = Arc::clone(&activity); let keepalive_closed = closed.clone(); let keepalive_pending = Arc::clone(&pending); + let keepalive_reason = close_reason.clone(); let keepalive = tokio::spawn(async move { let mut tick = tokio::time::interval(KEEPALIVE_INTERVAL); tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); @@ -213,6 +240,10 @@ impl AppClient { }) => sent, }; if !matches!(sent, Ok(Ok(()))) { + record_close_reason( + &keepalive_reason, + format!("websocket ping write failed: {sent:?}"), + ); break; } // Give the Pong (or any traffic) a bounded window to arrive. @@ -221,11 +252,14 @@ impl AppClient { _ = tokio::time::sleep(LIVENESS_DEADLINE) => {}, } if keepalive_activity.load(Ordering::Relaxed) == before { - // No return frame within the deadline → half-open/dead. + record_close_reason( + &keepalive_reason, + "websocket liveness deadline exceeded".into(), + ); break; } } - close_connection(&keepalive_closed, &keepalive_pending); + close_connection(&keepalive_closed, &keepalive_pending, &keepalive_reason); let _ = tokio::time::timeout(LIVENESS_DEADLINE, async { keepalive_sink.lock().await.close().await }) @@ -241,13 +275,14 @@ impl AppClient { reader, keepalive, closed, + close_reason, }, notify_rx, )) } /// Whether the socket is still considered live. Goes `false` once the - /// reader closes, a write fails, a request times out, or the keepalive + /// reader closes, a write fails, or the keepalive /// watchdog sees silence past [`LIVENESS_DEADLINE`]. Higher layers poll /// this to reconnect instead of reusing an unresponsive connection. pub fn is_alive(&self) -> bool { @@ -325,7 +360,7 @@ impl AppClient { let in_flight = { let mut pending = self.pending.lock().unwrap_or_else(|e| e.into_inner()); if !self.is_alive() { - return Err(anyhow!("app-server connection closed")); + return Err(connection_error(&self.close_reason)); } pending.insert(id.clone(), tx); pending.len() @@ -345,12 +380,13 @@ impl AppClient { // Bound the write and sink lock too: a half-open peer can stop reading // before the request has even reached its response wait. + let mut sent = false; let exchange = async { self.send_frame(frame) .await .with_context(|| format!("sending request `{method}`"))?; - rx.await - .map_err(|_| anyhow!("app-server connection closed"))? + sent = true; + rx.await.map_err(|_| connection_error(&self.close_reason))? }; match tokio::time::timeout(timeout, exchange).await { Ok(result) => { @@ -372,12 +408,25 @@ impl AppClient { result }, Err(_) => { - close_connection(&self.closed, &self.pending); + // A slow handler is not a dead transport. Late responses are + // discarded by id; other requests and streaming keep working. + // A cancelled partial write, however, cannot safely be reused. + if !sent { + record_close_reason( + &self.close_reason, + format!("request `{method}` write timed out"), + ); + close_connection(&self.closed, &self.pending, &self.close_reason); + } tracing::error!( target: "pocket_codex_codex::rpc", "<- {method} id={id} TIMED OUT after {:?}", started.elapsed() ); - Err(anyhow!("request `{method}` timed out; app-server connection closed")) + if sent { + Err(anyhow!("request `{method}` timed out")) + } else { + Err(connection_error(&self.close_reason)) + } }, } } @@ -397,11 +446,11 @@ impl AppClient { async fn send_frame(&self, frame: String) -> Result<()> { if !self.is_alive() { - return Err(anyhow!("app-server connection closed")); + return Err(connection_error(&self.close_reason)); } let result = tokio::select! { biased; - _ = self.closed.cancelled() => return Err(anyhow!("app-server connection closed")), + _ = self.closed.cancelled() => return Err(connection_error(&self.close_reason)), result = tokio::time::timeout(REQUEST_TIMEOUT, async { self.sink.lock().await.send(WsMessage::text(frame)).await }) => result, @@ -409,8 +458,12 @@ impl AppClient { match result { Ok(Ok(())) => Ok(()), result => { - close_connection(&self.closed, &self.pending); - Err(anyhow!("app-server connection closed while sending: {result:?}")) + record_close_reason( + &self.close_reason, + format!("websocket write failed: {result:?}"), + ); + close_connection(&self.closed, &self.pending, &self.close_reason); + Err(connection_error(&self.close_reason)) }, } } @@ -427,10 +480,27 @@ fn take_pending(pending: &Pending, id: &RequestId) -> Option>, message: String) { + let mut reason = reason.lock().unwrap_or_else(|e| e.into_inner()); + if reason.is_none() { + tracing::warn!(reason = %message, "app-server connection closed"); + *reason = Some(message); + } +} + +fn connection_error(reason: &StdMutex>) -> anyhow::Error { + let reason = reason.lock().unwrap_or_else(|e| e.into_inner()); + anyhow!("app-server connection closed: {}", reason.as_deref().unwrap_or("connection ended")) +} + +fn close_connection( + closed: &CancellationToken, + pending: &Pending, + reason: &StdMutex>, +) { closed.cancel(); for (_, tx) in pending.lock().unwrap_or_else(|e| e.into_inner()).drain() { - let _ = tx.send(Err(anyhow!("app-server connection closed"))); + let _ = tx.send(Err(connection_error(reason))); } } diff --git a/crates/pocket-codex-codex/src/client_tests.rs b/crates/pocket-codex-codex/src/client_tests.rs index 094bf49e..5cde2279 100644 --- a/crates/pocket-codex-codex/src/client_tests.rs +++ b/crates/pocket-codex-codex/src/client_tests.rs @@ -55,43 +55,118 @@ async fn cancelling_a_request_removes_its_pending_entry() { } #[tokio::test] -async fn rpc_timeout_closes_even_a_peer_that_keeps_sending_pongs() { - let (client, mut inbound, mut server) = connection().await; +async fn slow_rpc_does_not_disconnect_other_requests_or_late_replies() { + let (client, _inbound, mut server) = connection().await; let peer = tokio::spawn(async move { - let mut tick = tokio::time::interval(Duration::from_millis(10)); - loop { - tokio::select! { - frame = server.next() => match frame { - Some(Ok(WsMessage::Text(_))) | Some(Ok(WsMessage::Ping(_))) => {}, - _ => break, - }, - _ = tick.tick() => { - if server.send(WsMessage::Pong(Vec::new().into())).await.is_err() { - break; - } - } - } + let mut requests = Vec::new(); + for _ in 0..2 { + let frame = server + .next() + .await + .expect("frame") + .expect("read") + .into_text() + .expect("text"); + requests.push(serde_json::from_str::(&frame).expect("json")); + } + tokio::time::sleep(Duration::from_millis(100)).await; + for request in requests.into_iter().rev() { + server + .send(WsMessage::text( + json!({"id": request["id"], "result": {"ok": true}}).to_string(), + )) + .await + .expect("reply"); } + let frame = server + .next() + .await + .expect("frame") + .expect("read") + .into_text() + .expect("text"); + let request: Value = serde_json::from_str(&frame).expect("json"); + server + .send(WsMessage::text(json!({"id": request["id"], "result": {"ok": true}}).to_string())) + .await + .expect("reply"); + server }); - let (timed_out, other) = tokio::join!( - client.request_inner("thread/list", Some(json!({})), Duration::from_millis(100)), + let (slow, other) = tokio::join!( + client.request_inner("thread/resume", Some(json!({})), Duration::from_millis(30)), client.request("thread/read", json!({"threadId": "other"})), ); - assert!(timed_out - .expect_err("operation must fail") - .to_string() - .contains("timed out")); - assert!(other - .expect_err("operation must fail") - .to_string() - .contains("connection closed")); - assert!(!client.is_alive()); - assert!(tokio::time::timeout(Duration::from_secs(1), inbound.recv()) + let error = slow.expect_err("slow handler").to_string(); + assert!(error.contains("timed out"), "{error}"); + assert!(!error.contains("connection closed"), "{error}"); + assert_eq!(other.expect("other request survives")["ok"], true); + assert!(client.is_alive()); + assert_eq!( + client + .request("thread/list", json!({})) + .await + .expect("still usable")["ok"], + true + ); + assert!(client.pending.lock().expect("pending").is_empty()); + let _server = peer.await.expect("peer"); +} + +#[tokio::test] +async fn large_history_frame_and_repeated_switches_keep_the_same_connection() { + let (client, _inbound, mut server) = connection().await; + let peer = tokio::spawn(async move { + for i in 0..100 { + let frame = server + .next() + .await + .expect("frame") + .expect("read") + .into_text() + .expect("text"); + let request: Value = serde_json::from_str(&frame).expect("json"); + let text = if i == 0 { "x".repeat(17 << 20) } else { i.to_string() }; + server + .send(WsMessage::text( + json!({"id": request["id"], "result": {"text": text}}).to_string(), + )) + .await + .expect("reply"); + } + server + }); + for i in 0..100 { + let reply = client + .request("thread/read", json!({"threadId": format!("t{}", i % 5)})) + .await + .expect("history"); + assert_eq!( + reply["text"].as_str().expect("text").len(), + if i == 0 { 17 << 20 } else { i.to_string().len() } + ); + } + let _server = peer.await.expect("peer"); + assert!(client.is_alive()); +} + +#[tokio::test] +async fn websocket_protocol_failure_preserves_the_underlying_reason() { + use tokio::io::AsyncWriteExt; + let (client, mut inbound, mut server) = connection().await; + // An unmasked, final reserved opcode is invalid regardless of payload. + server + .get_mut() + .write_all(&[0x83, 0x00]) .await - .expect("test operation") - .is_none()); - assert!(client.pending.lock().expect("test operation").is_empty()); - peer.abort(); + .expect("invalid frame"); + assert!(inbound.recv().await.is_none()); + let error = client + .request("thread/list", json!({})) + .await + .expect_err("invalid protocol") + .to_string(); + assert!(error.contains("websocket read failed"), "{error}"); + assert!(error.contains("invalid opcode"), "{error}"); } #[tokio::test] diff --git a/crates/pocket-codex-host-svc/src/resume.rs b/crates/pocket-codex-host-svc/src/resume.rs index 98198cf7..511e5ce4 100644 --- a/crates/pocket-codex-host-svc/src/resume.rs +++ b/crates/pocket-codex-host-svc/src/resume.rs @@ -100,7 +100,7 @@ async fn resume_into(app_ws_addr: SocketAddr, thread_id: &str) -> Result<()> { .await .context("app-server initialize")?; client - .request("thread/resume", json!({ "threadId": thread_id })) + .request("thread/resume", json!({ "threadId": thread_id, "excludeTurns": true })) .await .context("thread/resume")?; Ok(()) diff --git a/crates/pocket-codex-pb/Cargo.toml b/crates/pocket-codex-pb/Cargo.toml index 56235cd4..26513271 100644 --- a/crates/pocket-codex-pb/Cargo.toml +++ b/crates/pocket-codex-pb/Cargo.toml @@ -25,13 +25,7 @@ tokio = { workspace = true } tracing = { workspace = true } url = { workspace = true } -# The published pb-mapper client SDK. Was a `path` dep on the -# `deps/pb-mapper` submodule pinned at 0.2.14, which is what let the client -# and the (already-0.5) production relay drift apart: the relay answers an -# over-quota register with a structured error frame that a 0.2.x client cannot -# decode, so it read a permanent refusal as a transport fault and retried -# without backoff. Depending on the release removes the skew by construction. -# +# The workspace selects the SDK's dedicated Git branch and locks its commit. # `pb-mapper` is a one-line re-export of `pb_mapper_client::sdk`, which is the # whole surface we need — no `uni-stream` type parameters to thread any more, # because the SDK picks the TCP providers itself. diff --git a/crates/pocket-codex-pb/src/session.rs b/crates/pocket-codex-pb/src/session.rs index 42e2e88e..2384208e 100644 --- a/crates/pocket-codex-pb/src/session.rs +++ b/crates/pocket-codex-pb/src/session.rs @@ -1,4 +1,4 @@ -//! Async wrappers around the published `pb-mapper` client SDK. +//! Async wrappers around the Git-pinned `pb-mapper` client SDK. //! //! ```text //! Pocket-Codex helper pb-mapper SDK entrypoint diff --git a/docs/session-switching-stability.md b/docs/session-switching-stability.md new file mode 100644 index 00000000..31087e27 --- /dev/null +++ b/docs/session-switching-stability.md @@ -0,0 +1,109 @@ +# Session 切换断线排查与验证(2026-09-08) + +## 已确认的主因 + +UI 打开 session 时先调用 `thread/resume`,再通过 `thread/read`、 +`thread/turns/list`、`thread/items/list` 读取历史。resume 没有设置 +`excludeTurns`,所以服务器还会额外恢复并返回一次完整历史。 + +真实桌面宿主的两个历史样本分别返回 25,820,325 和 36,948,729 字节。 +客户端原先使用 tungstenite 默认配置,单帧上限为 16 MiB,消息上限却为 +64 MiB。服务器以一个 WebSocket 帧发送 JSON-RPC 响应,因此这些合法历史 +会触发客户端协议层关闭连接。读循环原先丢弃底层错误,最终只剩 +`app-server connection closed`。 + +用旧上限再次复现时,实际错误为 WebSocket 1009:37,976,717 字节的帧 +超过 16,777,216 字节上限。同一历史以 `excludeTurns: true` 返回 2,251 字节。 +样本会随实时会话变化,以上大小来自各次独立采样。 + +## 修复 + +- Bridge 和 host meta service 的 resume 请求都设置 `excludeTurns: true`。 + 历史仍通过现有分页接口读取,保留完整历史导航和运行时配置。 +- AppClient 的帧上限与原来的 64 MiB 消息上限对齐,兼容较大的旧版历史。 + 保留底层读错误、关闭帧、写失败和存活检测失败的原因。 +- 已发送 RPC 的单次响应超时只结束该请求。迟到响应按 id 丢弃,不再关闭 + 其他 session 共用的连接;写失败和真实连接故障仍关闭连接。 +- UI 每次加载使用递增代次。切走后旧 resume 不再追加 history/config 请求, + A→B→A 时第一次 A 的结果也不能覆盖第二次 A。 +- 每个 service 的连接建立和断开串行化,避免并发重连互相移除刚建立的连接。 +- Backend 续期失败后先查询 relay:只有确认旧 key 已不存在才重新签发, + 避免一次网络失败把同一账号的设备分到不同命名空间。缓存锁改为每账号一把, + 一个账号的网络操作不再阻塞其他账号读取缓存。 +- `../pb-mapper`:订阅端连接 relay、注册端连接本地服务都具有超时; + 转发的两个本地 TCP 端也关闭 Nagle。补充真实大流量与小消息连续传输测试, + 并从订阅端 tracing span 中排除 credential。 + +不改变 CLI、持久化格式或 relay 协议。`excludeTurns` 是当前 Codex 已有字段; +没有改动 `deps/codex`。 + +## 验证结果 + +- Pocket-Codex 完整 workspace:309 passed、7 ignored;ignored 项为手工测试。 +- Flutter:457 passed、3 skipped;analyze 无问题;Rust/Dart 格式检查通过。 +- 两个仓库的 Clippy 通过。 +- pb-mapper 完整 workspace:207 passed、3 ignored;随后新增的本地 TCP + `TCP_NODELAY` 属性测试也通过。 +- Backend 使用隔离的真实 relay 额外运行 13 项测试,覆盖签发、复用和失效。 + 故障注入主动丢弃续期连接,确认不会再签发第二个命名空间;跨账号锁测试通过。 +- WebSocket 回归覆盖 17 MiB 单帧、同连接 100 次请求、单次超时后其他请求继续、 + 迟到响应、协议错误原因保留、peer close 和写超时。 +- Flutter 回归覆盖 A→B→A 的迟到 resume 和迟到 history。 +- pb-mapper 对明文/codec 两种模式分别传输 20 MiB 后,在同一连接继续进行 + 200 次小消息往返,并验证半关闭后的尾部完整性。小消息 P95 均约 2.1 ms。 + +### 实际 Rust bridge 的完整历史加载 + +使用真实历史数据库的隔离副本,包含上述大历史及其 fork 祖先;不启动模型 turn。 +每次切换执行实际 `thread_resume` + `thread_read`,包括历史分页和摘要读取, +在同一 WebSocket 上轮换三个 session。数据为本机开发构建,测试时还在编译其他目标。 + +| 路径 | 连续切换 | 中位数 | P95 | 最大值 | 断线 | +| --- | ---: | ---: | ---: | ---: | ---: | +| 直连独立 app-server | 100 | 49.5 ms | 155.1 ms | 636.3 ms | 0 | +| 本地真实 pb-mapper relay | 100 | 45.0 ms | 129.1 ms | 166.5 ms | 0 | +| 修复后的 bridge/pb-mapper 联编,经本地 relay | 300 | 13.8 ms | 45.5 ms | 63.3 ms | 0 | + +联编组单独运行,缓存也已预热,不能把两组差值全归因于 TCP_NODELAY。 + +单独比较 resume:两个真实样本从 1.44 s / 5.94 s 降到 9 ms / 25 ms, +响应缩小到约 2.2 KiB。该比较包括冷热状态差异,不能直接当作端到端 UI 加速倍数。 + +### 复跑 + +先准备独立 app-server(推荐复制历史数据库和相关 rollout,避免与正在使用的 +session 争抢 writer),然后运行: + +```sh +PCX_SOAK_WS=127.0.0.1:18880 \ +PCX_SOAK_THREAD_IDS=',,' \ +PCX_SOAK_ROUNDS=100 \ +cargo test -p pocket_codex_bridge real_session_switch_soak -- --ignored --nocapture +``` + +要覆盖 relay,将 `PCX_SOAK_WS` 改为本地订阅端口。测试要求样本有历史内容。 +它会 resume 指定 session,因此应使用独立副本或确认没有其他 writer 的会话。 + +## 接入与验证边界 + +pb-mapper 修复已通过 [PR #10](https://github.com/acking-you/pb-mapper/pull/10) +合入远端,并包含在供 Pocket-Codex 使用的 `pocket-codex` 分支中。 +Pocket-Codex 直接使用该 Git 分支,不再等待 crates.io 版本发布: + +```sh +cargo test --workspace --locked +``` + +`Cargo.toml` 使用 `git = "https://github.com/acking-you/pb-mapper"` 和 +`branch = "pocket-codex"`;`Cargo.lock` 固定到包含本次修复的 +`f0ed4271de962b8759d11171bcf901aa6df458cc`。普通构建不会自动追随分支变化。 +以后将兼容修复推送到该分支,再在本仓库运行 `cargo update -p pb-mapper`, +检查锁定提交并完成回归验证后提交 Cargo.lock。 + +上述 300 次切换在该 Git 接入之前,使用相同修复源码的临时本地联编完成; +现在正常依赖解析即可包含这些修复,无需绝对路径依赖或本地 patch。 + +上述结果验证本地链路和隔离历史上的稳定性,不是公网延迟或长期在线 SLA。 +当前运行的旧桌面进程在调查中也出现过 `/readyz` 无响应;线程采样已保存, +本轮没有把该进程的停顿归因于未经确认的 App Nap 或 backend。 +源码修复需要重新构建/启动应用才会用于现有桌面会话;本轮未部署生产 backend/relay。