mirror of
https://github.com/stablyai/orca.git
synced 2026-10-01 00:02:10 +00:00
fix(mobile): release streams whose ready arrives after a replayed cancel (#22945)
* fix(mobile): release streams whose ready arrives after a replayed cancel * test(mobile): cover a replayed browser stream replaced before its ready
This commit is contained in:
@@ -136,4 +136,30 @@ describe('runtime client-event stream disposal', () => {
|
||||
])
|
||||
client.close()
|
||||
})
|
||||
|
||||
it('releases the replayed registration on the new socket when disposed before its ready', async () => {
|
||||
const { client, socket: first } = connectReadyClient()
|
||||
const unsubscribe = client.subscribe('runtime.clientEvents.subscribe', null, () => {})
|
||||
const request = sentRequests(first, 'runtime.clientEvents.subscribe')[0]!
|
||||
emitReady(first, request.id, 'runtime-events:first-socket')
|
||||
|
||||
first.close()
|
||||
await vi.advanceTimersByTimeAsync(500)
|
||||
const second = sockets.at(-1)!
|
||||
expect(second).not.toBe(first)
|
||||
second.open()
|
||||
second.receive(JSON.stringify({ type: 'e2ee_ready' }))
|
||||
second.receive('encrypted:{"type":"e2ee_authenticated"}')
|
||||
expect(sentRequests(second, 'runtime.clientEvents.subscribe')).toEqual([
|
||||
expect.objectContaining({ id: request.id })
|
||||
])
|
||||
|
||||
unsubscribe()
|
||||
emitReady(second, request.id, 'runtime-events:second-socket')
|
||||
|
||||
expect(sentRequests(second, 'runtime.clientEvents.unsubscribe')).toEqual([
|
||||
expect.objectContaining({ params: { subscriptionId: 'runtime-events:second-socket' } })
|
||||
])
|
||||
client.close()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -98,4 +98,105 @@ describe('RpcClientStreamRegistry', () => {
|
||||
params: { subscriptionId: 'browser-screencast:page-1:test' }
|
||||
})
|
||||
})
|
||||
|
||||
it('releases a replayed browser stream replaced by a new one before its ready', () => {
|
||||
const { registry, sent } = createRegistry()
|
||||
registry.subscribe('browser.screencast', { page: 'page-1' }, () => {})
|
||||
const replayedId = sent[0]!.id
|
||||
registry.handleResponse(
|
||||
streamingResponse(replayedId, { type: 'ready', subscriptionId: 'page-1-old-connection' })
|
||||
)
|
||||
|
||||
registry.markForReplay()
|
||||
registry.replayAfterAuthentication()
|
||||
registry.subscribe('browser.screencast', { page: 'page-2' }, () => {})
|
||||
const replacementId = sent.at(-1)!.id
|
||||
registry.handleResponse(
|
||||
streamingResponse(replayedId, { type: 'ready', subscriptionId: 'page-1-new-connection' })
|
||||
)
|
||||
registry.handleResponse(
|
||||
streamingResponse(replacementId, { type: 'ready', subscriptionId: 'page-2' })
|
||||
)
|
||||
|
||||
expect(
|
||||
sent
|
||||
.filter((request) => request.method === 'browser.screencast.unsubscribe')
|
||||
.map((request) => request.params)
|
||||
).toEqual([{ subscriptionId: 'page-1-new-connection' }])
|
||||
expect(registry.size()).toBe(1)
|
||||
})
|
||||
|
||||
describe.each([
|
||||
['runtime.clientEvents.subscribe', 'runtime.clientEvents.unsubscribe', null],
|
||||
['browser.screencast', 'browser.screencast.unsubscribe', { page: 'page-1' }]
|
||||
])('%s ready id across a replay', (method, unsubscribeMethod, params) => {
|
||||
function unsubscribes(sent: SentRequest[]): unknown[] {
|
||||
return sent.filter((request) => request.method === unsubscribeMethod).map((r) => r.params)
|
||||
}
|
||||
|
||||
function subscribeReady() {
|
||||
const harness = createRegistry()
|
||||
const dispose = harness.registry.subscribe(method, params, () => {})
|
||||
const requestId = harness.sent[0]!.id
|
||||
harness.registry.handleResponse(
|
||||
streamingResponse(requestId, { type: 'ready', subscriptionId: 'old-connection-id' })
|
||||
)
|
||||
return { ...harness, dispose, requestId }
|
||||
}
|
||||
|
||||
it('releases the replayed registration when disposed before its new ready', () => {
|
||||
const { registry, sent, dispose, requestId } = subscribeReady()
|
||||
|
||||
registry.markForReplay()
|
||||
registry.replayAfterAuthentication()
|
||||
dispose()
|
||||
registry.handleResponse(
|
||||
streamingResponse(requestId, { type: 'ready', subscriptionId: 'new-connection-id' })
|
||||
)
|
||||
|
||||
expect(unsubscribes(sent)).toEqual([{ subscriptionId: 'new-connection-id' }])
|
||||
})
|
||||
|
||||
it('forgets the previous connection id when marked for replay', () => {
|
||||
const { registry, sent, dispose } = subscribeReady()
|
||||
|
||||
registry.markForReplay()
|
||||
dispose()
|
||||
|
||||
// A disposal while disconnected has nothing to name on the next connection.
|
||||
expect(unsubscribes(sent)).toEqual([])
|
||||
expect(registry.size()).toBe(0)
|
||||
})
|
||||
|
||||
it('still releases a stream cancelled before its first ready', () => {
|
||||
const { registry, sent } = createRegistry()
|
||||
const dispose = registry.subscribe(method, params, () => {})
|
||||
const requestId = sent[0]!.id
|
||||
|
||||
dispose()
|
||||
expect(unsubscribes(sent)).toEqual([])
|
||||
registry.handleResponse(
|
||||
streamingResponse(requestId, { type: 'ready', subscriptionId: 'first-id' })
|
||||
)
|
||||
|
||||
expect(unsubscribes(sent)).toEqual([{ subscriptionId: 'first-id' }])
|
||||
expect(registry.size()).toBe(0)
|
||||
})
|
||||
|
||||
it('sends one unsubscribe however often the stream is disposed', () => {
|
||||
const { registry, sent, dispose, requestId } = subscribeReady()
|
||||
|
||||
registry.markForReplay()
|
||||
registry.replayAfterAuthentication()
|
||||
dispose()
|
||||
dispose()
|
||||
registry.handleResponse(
|
||||
streamingResponse(requestId, { type: 'ready', subscriptionId: 'new-connection-id' })
|
||||
)
|
||||
dispose()
|
||||
|
||||
expect(unsubscribes(sent)).toEqual([{ subscriptionId: 'new-connection-id' }])
|
||||
expect(registry.size()).toBe(0)
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user