fix(mobile): release notification and account streams from the transport

This commit is contained in:
Brennan Benson
2026-09-25 15:01:40 -07:00
parent c220d92c03
commit ea0c4458f0
11 changed files with 295 additions and 104 deletions
@@ -1,25 +0,0 @@
import { bindDeferredRpcOperation, defineRpcOperation } from '../transport/rpc-operation'
import { rpcResultVariant } from '../transport/rpc-operation-result-reader'
import { notificationUnreadReplySchema } from './notification-reply-schema'
/**
* Closing the desktop notification stream on the host.
*
* Its own module rather than a line in `mobile-push-registration-operations.ts`: that module is the
* push route this device holds with a gateway, and this is the socket subscription the paired
* connection holds. They are two different deliveries of the same alert and neither implies the
* other.
*
* A skip rather than a throw, and the reply is unread either way: the disposer sends this on its
* way out with nothing left to show a host message on, and main's `.catch(() => {})` already made a
* refusal and a dropped connection the same non-event.
*/
export const desktopNotificationStreamUnsubscribe = bindDeferredRpcOperation(
defineRpcOperation({
name: 'notifications.unsubscribe-or-skip',
method: 'notifications.unsubscribe',
acceptance: 'success-result-or-skip',
barrier: 'after-caller-barrier',
read: rpcResultVariant('notification-stream-closed', notificationUnreadReplySchema)
})
)
@@ -3,6 +3,7 @@ import { subscribeToDesktopNotifications } from './mobile-notifications'
import { dismissHostPushNotification } from './push-socket-dismissal'
import { requestNotificationCatchup } from './push-dismissal-reconciliation'
import { RpcClientStreamRegistry } from '../transport/rpc-client-stream-registry'
import { MobileRelayRpcStreams } from '../transport/mobile-relay-rpc-streams'
import type { RpcClient } from '../transport/rpc-client'
import type { RpcResponse } from '../transport/types'
@@ -37,10 +38,28 @@ function readSentFrame(request: unknown): SentFrame {
}
}
function transportClient(subscribe: RpcClient['subscribe']) {
const requests: { method: string; params: unknown }[] = []
const client: RpcClient = {
sendRequest: async (method, params) => {
requests.push({ method, params })
return { id: 'reply-1', ok: true, result: {}, _meta: { runtimeId: 'runtime-1' } }
},
subscribe,
updateTerminalSubscriptionViewport: () => {},
getState: () => 'connected',
getReconnectAttempt: () => 0,
getLastConnectedAt: () => null,
onStateChange: () => () => {},
notifyForeground: () => {},
close: () => {}
}
return { requests, client }
}
/** The real stream registry, so dispose-before-ready is answered by the transport, not by a fake. */
function registryClient() {
const sent: SentFrame[] = []
const requests: { method: string; params: unknown }[] = []
let id = 0
const registry = new RpcClientStreamRegistry({
nextId: () => `rpc-${++id}`,
@@ -51,22 +70,44 @@ function registryClient() {
return true
}
})
const client: RpcClient = {
sendRequest: async (method, params) => {
requests.push({ method, params })
return { id: 'reply-1', ok: true, result: {}, _meta: { runtimeId: 'runtime-1' } }
},
subscribe: (method, params, onData, options) =>
registry.subscribe(method, params, onData, options),
updateTerminalSubscriptionViewport: () => {},
getState: () => 'connected',
getReconnectAttempt: () => 0,
getLastConnectedAt: () => null,
onStateChange: () => () => {},
notifyForeground: () => {},
close: () => {}
return {
registry,
sent,
...transportClient((method, params, onData, options) =>
registry.subscribe(method, params, onData, options)
)
}
return { registry, sent, requests, client }
}
/** The real relay stream manager, the other transport a paired phone reaches a host through. */
function relayClient() {
const sent: SentFrame[] = []
let id = 0
const streams = new MobileRelayRpcStreams({
nextId: () => `relay-${++id}`,
sendFrame: (frame) => {
sent.push(readSentFrame(frame))
return true
},
waitForConnected: async () => {}
})
return {
streams,
sent,
...transportClient((method, params, onData, options) =>
streams.subscribe(method, params, onData, options)
)
}
}
/** Every `notifications.unsubscribe` the phone put on the wire, by either route. */
function notificationReleases(rpc: {
sent: SentFrame[]
requests: { method: string; params: unknown }[]
}): unknown[] {
return [...rpc.sent, ...rpc.requests]
.filter((frame) => frame.method === 'notifications.unsubscribe')
.map((frame) => frame.params)
}
function readyReply(id: string, subscriptionId: string): RpcResponse {
@@ -123,7 +164,7 @@ describe('subscribeToDesktopNotifications', () => {
expect(dismissHostPushNotification).toHaveBeenCalledWith(dismissal, 'host-1')
})
it('never runs the ready arm when the disposer ran before the reply landed', () => {
it('releases the host stream once a ready lands after the disposer ran (direct)', () => {
const rpc = registryClient()
const stop = subscribeToDesktopNotifications(rpc.client, 'host-1')
const subscribeFrame = rpc.sent[0]!
@@ -133,12 +174,24 @@ describe('subscribeToDesktopNotifications', () => {
rpc.registry.handleResponse(readyReply(subscribeFrame.id, 'sub-1'))
expect(requestNotificationCatchup).not.toHaveBeenCalled()
// The subscription id never reaches this module, so nothing closes the host's stream.
expect(rpc.requests).toEqual([])
expect(rpc.sent).toHaveLength(1)
expect(notificationReleases(rpc)).toEqual([{ subscriptionId: 'sub-1' }])
})
it('closes the host stream when the disposer runs after the ready reply', async () => {
it('releases the host stream once a ready lands after the disposer ran (relay)', async () => {
const rpc = relayClient()
const stop = subscribeToDesktopNotifications(rpc.client, 'host-1')
await Promise.resolve()
const subscribeFrame = rpc.sent[0]!
expect(subscribeFrame.method).toBe('notifications.subscribe')
stop()
rpc.streams.handleResponse(readyReply(subscribeFrame.id, 'sub-1'))
expect(requestNotificationCatchup).not.toHaveBeenCalled()
expect(notificationReleases(rpc)).toEqual([{ subscriptionId: 'sub-1' }])
})
it('closes the host stream once when the disposer runs after the ready reply (direct)', async () => {
const rpc = registryClient()
const stop = subscribeToDesktopNotifications(rpc.client, 'host-1')
rpc.registry.handleResponse(readyReply(rpc.sent[0]!.id, 'sub-1'))
@@ -146,8 +199,47 @@ describe('subscribeToDesktopNotifications', () => {
stop()
await Promise.resolve()
expect(rpc.requests).toEqual([
{ method: 'notifications.unsubscribe', params: { subscriptionId: 'sub-1' } }
])
expect(notificationReleases(rpc)).toEqual([{ subscriptionId: 'sub-1' }])
})
it('closes the host stream once when the disposer runs after the ready reply (relay)', async () => {
const rpc = relayClient()
const stop = subscribeToDesktopNotifications(rpc.client, 'host-1')
await Promise.resolve()
rpc.streams.handleResponse(readyReply(rpc.sent[0]!.id, 'sub-1'))
stop()
await Promise.resolve()
expect(notificationReleases(rpc)).toEqual([{ subscriptionId: 'sub-1' }])
})
it('releases the replayed stream by its new id, never the one the closed socket assigned', () => {
const rpc = registryClient()
const stop = subscribeToDesktopNotifications(rpc.client, 'host-1')
const subscribeFrame = rpc.sent[0]!
rpc.registry.handleResponse(readyReply(subscribeFrame.id, 'sub-1'))
rpc.registry.markForReplay()
rpc.registry.replayAfterAuthentication()
expect(rpc.sent[1]).toMatchObject({ id: subscribeFrame.id, method: 'notifications.subscribe' })
stop()
rpc.registry.handleResponse(readyReply(subscribeFrame.id, 'sub-2'))
expect(notificationReleases(rpc)).toEqual([{ subscriptionId: 'sub-2' }])
})
it('catches up again when a replayed subscribe is ready', () => {
const rpc = registryClient()
subscribeToDesktopNotifications(rpc.client, 'host-1')
const subscribeFrame = rpc.sent[0]!
rpc.registry.handleResponse(readyReply(subscribeFrame.id, 'sub-1'))
rpc.registry.markForReplay()
rpc.registry.replayAfterAuthentication()
rpc.registry.handleResponse(readyReply(subscribeFrame.id, 'sub-2'))
expect(requestNotificationCatchup).toHaveBeenCalledTimes(2)
expect(notificationReleases(rpc)).toEqual([])
})
})
@@ -1,5 +1,4 @@
import { requestNotificationCatchup } from './push-dismissal-reconciliation'
import { desktopNotificationStreamUnsubscribe } from './desktop-notification-stream-operations'
import { dismissHostPushNotification } from './push-socket-dismissal'
import type { DismissNotificationEvent } from './desktop-notification-events'
import type { RpcClient } from '../transport/rpc-client'
@@ -10,29 +9,16 @@ export {
type NotificationPermissionState
} from './notification-permissions'
type SubscribeResult = {
type: 'ready'
subscriptionId: string
}
export function subscribeToDesktopNotifications(client: RpcClient, hostId: string): () => void {
let subscriptionId: string | null = null
let disposed = false
function unsubscribeServer(id: string) {
if (client.getState() === 'connected') {
// The reply is never read: the stream is already gone locally either way.
desktopNotificationStreamUnsubscribe.request(client, { subscriptionId: id }).catch(() => {})
}
}
const params = { includeDesktopSuppressed: true }
// The transport releases the host registration with the id from the current `ready`.
const unsubscribeStream = client.subscribe('notifications.subscribe', params, (data: unknown) => {
const event = data as DismissNotificationEvent | SubscribeResult | { type: string }
const event = data as DismissNotificationEvent | { type: string }
// No dispose-before-ready arm: every transport detaches this listener inside
// `unsubscribeStream()`, so a callback that runs at all runs before disposal.
if (event.type === 'ready') {
subscriptionId = (event as SubscribeResult).subscriptionId
// A max watermark asks only which delivered pushes are stale; socket history
// never becomes a second OS-notification delivery route.
void requestNotificationCatchup(client, hostId, () => disposed).catch(() => {})
@@ -46,8 +32,5 @@ export function subscribeToDesktopNotifications(client: RpcClient, hostId: strin
return () => {
disposed = true
unsubscribeStream()
if (subscriptionId) {
unsubscribeServer(subscriptionId)
}
}
}
@@ -5,7 +5,7 @@ const HOST = 'host-1'
/**
* The desktop notification socket: one subscribe, the catch-up read its `ready` arms, the tray
* dismissals its events drive, and the server unsubscribe the disposer sends.
* dismissals its events drive, and the server unsubscribe the transport sends once it is disposed.
*
* The disposer is the whole output — it is what a host connection calls when the client goes away —
* so the recording drives `start` and `stop` and observes what each put on the wire. Everything the
@@ -53,7 +53,7 @@ export async function runRecording(
// byte-identical.
await scheduler.flush()
// A stream the product forgot to close is only visible on the wire when its method has an
// unsubscribe builder; `notifications.subscribe` has none, so closing it writes nothing and the
// unsubscribe builder; `agentSession.subscribe` has none, so closing it writes nothing and the
// leak stays a live registry record until some later cutover replays it. Observed here, after
// the product's own cleanup and before the transport tears the registries down, so a
// builder-less subscription is pinned without a scenario that cuts over to expose it.
@@ -164,9 +164,11 @@ function readyReply(id: string): RpcSuccess {
describe('MobileRelayRpcStreams cancel fencing', () => {
function subscribed() {
const listener = vi.fn()
let id = 0
const sendFrame = vi.fn((_frame: { id: string; method: string; params?: unknown }) => true)
const streams = new MobileRelayRpcStreams({
nextId: () => 'stream-1',
sendFrame: vi.fn(() => true),
nextId: () => `stream-${++id}`,
sendFrame,
waitForConnected: async () => {}
})
const cancel = streams.subscribe(
@@ -174,7 +176,7 @@ describe('MobileRelayRpcStreams cancel fencing', () => {
{ includeDesktopSuppressed: true },
listener
)
return { listener, streams, cancel }
return { listener, streams, cancel, sendFrame }
}
it('delivers a ready reply to a live subscription', async () => {
@@ -185,13 +187,16 @@ describe('MobileRelayRpcStreams cancel fencing', () => {
expect(listener).toHaveBeenCalledExactlyOnceWith({ type: 'ready', subscriptionId: 'sub-1' })
})
it('drops a ready reply that lands after the caller cancelled', async () => {
const { listener, streams, cancel } = subscribed()
it('releases, without delivering, a ready reply that lands after the caller cancelled', async () => {
const { listener, streams, cancel, sendFrame } = subscribed()
await Promise.resolve()
cancel()
expect(streams.handleResponse(readyReply('stream-1'))).toBe(false)
expect(streams.handleResponse(readyReply('stream-1'))).toBe(true)
expect(listener).not.toHaveBeenCalled()
expect(sendFrame.mock.calls.map(([frame]) => frame).slice(1)).toEqual([
{ id: 'stream-2', method: 'notifications.unsubscribe', params: { subscriptionId: 'sub-1' } }
])
})
})
@@ -8,7 +8,7 @@ import {
buildTerminalUnsubscribeParams,
updateTerminalSubscriptionViewport
} from './rpc-client-terminal-subscription'
import { buildReadyStreamUnsubscribe } from './rpc-client-server-subscription'
import { buildReadyStreamUnsubscribe, isReadyIdStream } from './rpc-client-server-subscription'
import { isStreamingOpenerReply } from './rpc-acceptance-policies'
import type { RpcClient } from './rpc-client'
import type { RpcResponse, RpcSuccess } from './types'
@@ -205,17 +205,9 @@ export class MobileRelayRpcStreams {
this.cancelledSubscriptions.set(id, { method: stream.method, unsubscribe: byParams })
} else if (unsubscribe || byParams) {
this.sendUnsubscribe((unsubscribe ?? byParams)!)
} else if (
stream.method === 'browser.screencast' ||
stream.method === 'runtime.clientEvents.subscribe'
) {
} else if (isReadyIdStream(stream.method)) {
// Keep only the cleanup route while the server assigns its subscription ID.
this.cancelledSubscriptions.set(id, { method: stream.method })
} else if (stream.subscriptionId) {
this.sendUnsubscribe({
method: stream.method.replace(/\.subscribe$/, '.unsubscribe'),
params: { subscriptionId: stream.subscriptionId }
})
}
}
}
@@ -0,0 +1,138 @@
import { describe, expect, it } from 'vitest'
import { MobileRelayRpcStreams } from './mobile-relay-rpc-streams'
import { RpcClientStreamRegistry } from './rpc-client-stream-registry'
import { READY_STREAM_RELEASE_METHODS } from './rpc-client-server-subscription'
import type { RpcSuccess } from './types'
type SentFrame = { id: string; method: string; params?: unknown }
/** The registry sends through an `unknown` port, so name the shape the assertions read. */
function readSentFrame(request: unknown): SentFrame {
if (
typeof request !== 'object' ||
request === null ||
!('id' in request) ||
typeof request.id !== 'string' ||
!('method' in request) ||
typeof request.method !== 'string'
) {
throw new Error('The stream registry sent a frame without a string id and method')
}
return {
id: request.id,
method: request.method,
params: 'params' in request ? request.params : undefined
}
}
function readyReply(id: string, subscriptionId: string): RpcSuccess {
return {
id,
ok: true,
streaming: true,
result: { type: 'ready', subscriptionId, snapshot: { accounts: [] } },
_meta: { runtimeId: 'runtime-1' }
}
}
function directTransport() {
const sent: SentFrame[] = []
let id = 0
const registry = new RpcClientStreamRegistry({
nextId: () => `rpc-${++id}`,
deviceToken: 'device-token',
getState: () => 'connected',
sendEncrypted: (request) => {
sent.push(readSentFrame(request))
return true
}
})
return {
sent,
subscribe: (method: string) => registry.subscribe(method, null, () => {}),
reply: (response: RpcSuccess) => registry.handleResponse(response)
}
}
function relayTransport() {
const sent: SentFrame[] = []
let id = 0
const streams = new MobileRelayRpcStreams({
nextId: () => `relay-${++id}`,
sendFrame: (frame) => {
sent.push(frame)
return true
},
waitForConnected: async () => {}
})
return {
sent,
subscribe: (method: string) => streams.subscribe(method, null, () => {}),
reply: (response: RpcSuccess) => streams.handleResponse(response)
}
}
function releases(sent: SentFrame[], method: string): unknown[] {
return sent.filter((frame) => frame.method === method).map((frame) => frame.params)
}
describe.each([
['direct', directTransport],
['relay', relayTransport]
])('%s transport releases the accounts stream', (_name, transport) => {
it('sends accounts.unsubscribe with the ready id on dispose', async () => {
const wire = transport()
const dispose = wire.subscribe('accounts.subscribe')
await Promise.resolve()
wire.reply(readyReply(wire.sent[0]!.id, 'accounts-conn-1'))
dispose()
expect(releases(wire.sent, 'accounts.unsubscribe')).toEqual([
{ subscriptionId: 'accounts-conn-1' }
])
})
it('holds a dispose that beat the ready and releases once the ready lands', async () => {
const wire = transport()
const dispose = wire.subscribe('accounts.subscribe')
await Promise.resolve()
dispose()
expect(releases(wire.sent, 'accounts.unsubscribe')).toEqual([])
wire.reply(readyReply(wire.sent[0]!.id, 'accounts-conn-1'))
expect(releases(wire.sent, 'accounts.unsubscribe')).toEqual([
{ subscriptionId: 'accounts-conn-1' }
])
})
})
// Iterates the mapping itself, so a method added to it is held to the same release on both routes.
describe.each(
[...READY_STREAM_RELEASE_METHODS].flatMap(([method, release]) => [
['direct', method, release, directTransport] as const,
['relay', method, release, relayTransport] as const
])
)('%s transport releases %s through the ready id', (_name, method, release, transport) => {
it('releases once after ready, and once when the dispose beat the ready', async () => {
const afterReady = transport()
const disposeAfterReady = afterReady.subscribe(method)
await Promise.resolve()
afterReady.reply(readyReply(afterReady.sent[0]!.id, 'host-id-1'))
disposeAfterReady()
const beforeReady = transport()
const disposeBeforeReady = beforeReady.subscribe(method)
await Promise.resolve()
disposeBeforeReady()
beforeReady.reply(readyReply(beforeReady.sent[0]!.id, 'host-id-2'))
expect(afterReady.sent.slice(1)).toEqual([
expect.objectContaining({ method: release, params: { subscriptionId: 'host-id-1' } })
])
expect(beforeReady.sent.slice(1)).toEqual([
expect.objectContaining({ method: release, params: { subscriptionId: 'host-id-2' } })
])
})
})
@@ -1,12 +1,20 @@
// Streams whose host names its registration only in the `ready` frame. The transport that saw that
// frame is the only holder of a current id, so it alone sends the release.
export const READY_STREAM_RELEASE_METHODS: ReadonlyMap<string, string> = new Map([
['browser.screencast', 'browser.screencast.unsubscribe'],
['runtime.clientEvents.subscribe', 'runtime.clientEvents.unsubscribe'],
['notifications.subscribe', 'notifications.unsubscribe'],
['accounts.subscribe', 'accounts.unsubscribe']
])
export function isReadyIdStream(method: string | undefined): boolean {
return method !== undefined && READY_STREAM_RELEASE_METHODS.has(method)
}
export function buildReadyStreamUnsubscribe(
method: string,
subscriptionId: string
): { method: string; params: { subscriptionId: string } } | null {
if (method === 'browser.screencast') {
return { method: 'browser.screencast.unsubscribe', params: { subscriptionId } }
}
if (method === 'runtime.clientEvents.subscribe') {
return { method: 'runtime.clientEvents.unsubscribe', params: { subscriptionId } }
}
return null
const release = READY_STREAM_RELEASE_METHODS.get(method)
return release ? { method: release, params: { subscriptionId } } : null
}
@@ -7,7 +7,7 @@ import {
buildTerminalUnsubscribeParams,
updateTerminalSubscriptionViewport
} from './rpc-client-terminal-subscription'
import { buildReadyStreamUnsubscribe } from './rpc-client-server-subscription'
import { buildReadyStreamUnsubscribe, isReadyIdStream } from './rpc-client-server-subscription'
import { isStreamingOpenerReply } from './rpc-acceptance-policies'
import {
isStreamingSubscriptionReadyResult,
@@ -108,6 +108,8 @@ export class RpcClientStreamRegistry {
this.pendingBrowserRequestId = null
for (const [id, stream] of this.streams) {
stream.sent = false
// The id named a registration on the closed socket; the replay's ready brings the new one.
stream.subscriptionId = undefined
this.resetTerminalRouting(id)
}
}
@@ -197,13 +199,10 @@ export class RpcClientStreamRegistry {
private dispose(id: string): void {
const stream = this.streams.get(id)
if (stream?.method === 'browser.screencast') {
stream.cancelled = true
this.clearBrowserRequest(id)
this.disposeServerSubscription(id, stream)
return
}
if (stream?.method === 'runtime.clientEvents.subscribe') {
if (stream && isReadyIdStream(stream.method)) {
if (stream.method === 'browser.screencast') {
this.clearBrowserRequest(id)
}
this.disposeServerSubscription(id, stream)
return
}
@@ -112,9 +112,8 @@ export const UNVALIDATED_RPC_REQUEST_PORT_PENDING: readonly UnvalidatedRpcReques
// src/notifications/ — push registration and delivery. Nothing is left here. Registration and
// unregistration migrated in step 4; see mobile-push-registration-operations.ts. Tray
// reconciliation followed once a scenario could declare the notification tray and the stored host
// list it resolves against; see push-dismissal-operations.ts. The stream unsubscribe inside the
// `notifications.subscribe` callback migrated in step 6 once the recorder could script the
// `ready` frame that hands it a subscription id; see desktop-notification-stream-operations.ts.
// list it resolves against; see push-dismissal-operations.ts. The stream's `notifications.unsubscribe`
// is no longer a request here: the stream transport sends it with the id from the current `ready`.
// src/session/ — session screen: chat, diff review, PR actions, tabs. The github.* PR surface,
// the diff-review loaders and the rest of the screen migrated in step 4; see