mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 08:03:20 +00:00
fix(runtime): make files.unwatch wait for watcher release and refuse foreign connections
files.unwatch tore down whatever subscriptionId it was handed, so any paired
connection could retire another connection's file watch, and the reply landed
before @parcel/watcher was actually released — a client that rewatched on the
reply held two watchers on one path and double-delivered change events.
Route the connection-scoped path through cleanupIfOwnedByConnectionAndWait: it
reuses the ownership verdict terminal.unsubscribe already relies on, and awaits
cleanupAndWait so the reply is the release. resolveConnectionOwnership factors
that verdict out of cleanupIfOwnedByConnection, whose behaviour is unchanged —
an unregistered id is still reported as gone rather than refused.
Wire-compatible: `files.unwatch` keeps its params and its `{ unsubscribed:
boolean }` response, and no client reads the flag — the renderer and the shared
control connection only branch on the RPC envelope's `ok`. In-process callers
pass no connectionId and keep the unconditional teardown.
Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
@@ -0,0 +1,165 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import type { OrcaRuntimeService } from '../../orca-runtime'
|
||||
import { RpcDispatcher } from '../dispatcher'
|
||||
import type { RpcRequest } from '../core'
|
||||
import { RuntimeSubscriptionRegistry } from '../../runtime-subscription-registry'
|
||||
import { FILE_METHODS } from './files'
|
||||
|
||||
function unwatchRequest(subscriptionId: string): RpcRequest {
|
||||
return {
|
||||
id: 'req-1',
|
||||
authToken: 'tok',
|
||||
method: 'files.unwatch',
|
||||
params: { subscriptionId }
|
||||
}
|
||||
}
|
||||
|
||||
describe('files.unwatch ownership', () => {
|
||||
it('refuses teardown when the socket does not own the subscription', async () => {
|
||||
const cleanupSubscriptionIfOwnedByConnectionAndWait = vi.fn().mockResolvedValue(false)
|
||||
const cleanupSubscriptionAndWait = vi.fn()
|
||||
const runtime = {
|
||||
getRuntimeId: () => 'test-runtime',
|
||||
cleanupSubscriptionIfOwnedByConnectionAndWait,
|
||||
cleanupSubscriptionAndWait
|
||||
} as unknown as OrcaRuntimeService
|
||||
const dispatcher = new RpcDispatcher({ runtime, methods: FILE_METHODS })
|
||||
const replies: unknown[] = []
|
||||
|
||||
await dispatcher.dispatchStreaming(
|
||||
unwatchRequest('files-watch-conn-owner-1'),
|
||||
(reply) => replies.push(JSON.parse(reply)),
|
||||
{ connectionId: 'conn-attacker' }
|
||||
)
|
||||
|
||||
expect(cleanupSubscriptionIfOwnedByConnectionAndWait).toHaveBeenCalledWith(
|
||||
'files-watch-conn-owner-1',
|
||||
'conn-attacker'
|
||||
)
|
||||
expect(cleanupSubscriptionAndWait).not.toHaveBeenCalled()
|
||||
expect(replies).toEqual([expect.objectContaining({ result: { unsubscribed: false } })])
|
||||
})
|
||||
|
||||
// Why: returning before @parcel/watcher is released lets a rewatch hold two watchers.
|
||||
it('does not reply until the owning connection teardown settles', async () => {
|
||||
const subscriptions = new RuntimeSubscriptionRegistry()
|
||||
let releaseWatcher: (() => void) | undefined
|
||||
let released = false
|
||||
subscriptions.register(
|
||||
'files-watch-slow-1',
|
||||
() =>
|
||||
new Promise<void>((resolve) => {
|
||||
releaseWatcher = () => {
|
||||
released = true
|
||||
resolve()
|
||||
}
|
||||
}),
|
||||
'conn-owner'
|
||||
)
|
||||
const runtime = {
|
||||
getRuntimeId: () => 'test-runtime',
|
||||
cleanupSubscriptionIfOwnedByConnectionAndWait:
|
||||
subscriptions.cleanupIfOwnedByConnectionAndWait.bind(subscriptions)
|
||||
} as unknown as OrcaRuntimeService
|
||||
const dispatcher = new RpcDispatcher({ runtime, methods: FILE_METHODS })
|
||||
|
||||
let settled = false
|
||||
const pending = dispatcher
|
||||
.dispatch(unwatchRequest('files-watch-slow-1'), { connectionId: 'conn-owner' })
|
||||
.then((response) => {
|
||||
settled = true
|
||||
return response
|
||||
})
|
||||
|
||||
await vi.waitFor(() => expect(releaseWatcher).toBeDefined())
|
||||
expect(settled).toBe(false)
|
||||
expect(released).toBe(false)
|
||||
|
||||
releaseWatcher?.()
|
||||
await expect(pending).resolves.toMatchObject({ ok: true, result: { unsubscribed: true } })
|
||||
expect(released).toBe(true)
|
||||
})
|
||||
|
||||
// Why: the reply is the only signal a client waits on before rewatching, so the whole
|
||||
// point of awaiting teardown is that the second watcher never overlaps the first.
|
||||
it('holds one watcher when the owner rewatches straight after the unwatch reply', async () => {
|
||||
const subscriptions = new RuntimeSubscriptionRegistry()
|
||||
let liveWatchers = 0
|
||||
let peakLiveWatchers = 0
|
||||
const pendingReleases: (() => void)[] = []
|
||||
const watchFileExplorer = vi.fn(async () => {
|
||||
liveWatchers += 1
|
||||
peakLiveWatchers = Math.max(peakLiveWatchers, liveWatchers)
|
||||
return () =>
|
||||
new Promise<void>((resolve) => {
|
||||
pendingReleases.push(() => {
|
||||
liveWatchers -= 1
|
||||
resolve()
|
||||
})
|
||||
})
|
||||
})
|
||||
const runtime = {
|
||||
getRuntimeId: () => 'test-runtime',
|
||||
watchFileExplorer,
|
||||
registerSubscriptionCleanup: subscriptions.register.bind(subscriptions),
|
||||
cleanupSubscription: subscriptions.cleanup.bind(subscriptions),
|
||||
cleanupSubscriptionAndWait: subscriptions.cleanupAndWait.bind(subscriptions),
|
||||
cleanupSubscriptionIfOwnedByConnectionAndWait:
|
||||
subscriptions.cleanupIfOwnedByConnectionAndWait.bind(subscriptions)
|
||||
} as unknown as OrcaRuntimeService
|
||||
const dispatcher = new RpcDispatcher({ runtime, methods: FILE_METHODS })
|
||||
|
||||
const watch = (
|
||||
id: string
|
||||
): { done: Promise<unknown>; events: { type?: string; subscriptionId?: string }[] } => {
|
||||
const events: { type?: string; subscriptionId?: string }[] = []
|
||||
const done = dispatcher.dispatchStreaming(
|
||||
{ id, authToken: 'tok', method: 'files.watch', params: { worktree: 'id:wt-1' } },
|
||||
(reply) => {
|
||||
const parsed = JSON.parse(reply) as {
|
||||
result?: { type?: string; subscriptionId?: string }
|
||||
}
|
||||
if (parsed.result) {
|
||||
events.push(parsed.result)
|
||||
}
|
||||
},
|
||||
{ connectionId: 'conn-owner' }
|
||||
)
|
||||
return { done, events }
|
||||
}
|
||||
|
||||
const first = watch('watch-1')
|
||||
await vi.waitFor(() => expect(first.events.some((event) => event.type === 'ready')).toBe(true))
|
||||
const subscriptionId = first.events[0]?.subscriptionId
|
||||
expect(subscriptionId).toBeTruthy()
|
||||
expect(liveWatchers).toBe(1)
|
||||
|
||||
let unwatchSettled = false
|
||||
const unwatch = dispatcher
|
||||
.dispatch(unwatchRequest(subscriptionId!), { connectionId: 'conn-owner' })
|
||||
.then((response) => {
|
||||
unwatchSettled = true
|
||||
return response
|
||||
})
|
||||
await vi.waitFor(() => expect(pendingReleases).toHaveLength(1))
|
||||
expect(unwatchSettled).toBe(false)
|
||||
|
||||
pendingReleases[0]?.()
|
||||
await expect(unwatch).resolves.toMatchObject({ ok: true, result: { unsubscribed: true } })
|
||||
await first.done
|
||||
expect(liveWatchers).toBe(0)
|
||||
|
||||
const second = watch('watch-2')
|
||||
await vi.waitFor(() => expect(second.events.some((event) => event.type === 'ready')).toBe(true))
|
||||
expect(peakLiveWatchers).toBe(1)
|
||||
|
||||
const secondUnwatch = dispatcher.dispatch(unwatchRequest(second.events[0]!.subscriptionId!), {
|
||||
connectionId: 'conn-owner'
|
||||
})
|
||||
await vi.waitFor(() => expect(pendingReleases).toHaveLength(2))
|
||||
pendingReleases[1]?.()
|
||||
await secondUnwatch
|
||||
await second.done
|
||||
expect(liveWatchers).toBe(0)
|
||||
})
|
||||
})
|
||||
@@ -277,7 +277,17 @@ export const FILE_METHODS: RpcAnyMethod[] = [
|
||||
defineMethod({
|
||||
name: 'files.unwatch',
|
||||
params: FileUnwatch,
|
||||
handler: async (params, { runtime }) => {
|
||||
handler: async (params, { runtime, connectionId }) => {
|
||||
// Why: only the connection that owns the watch may retire it, and the reply must
|
||||
// wait for the watcher release so a rewatch cannot hold two watchers on one path.
|
||||
if (connectionId) {
|
||||
return {
|
||||
unsubscribed: await runtime.cleanupSubscriptionIfOwnedByConnectionAndWait(
|
||||
params.subscriptionId,
|
||||
connectionId
|
||||
)
|
||||
}
|
||||
}
|
||||
await runtime.cleanupSubscriptionAndWait(params.subscriptionId)
|
||||
return { unsubscribed: true }
|
||||
}
|
||||
|
||||
@@ -23,6 +23,7 @@ export type RuntimeServiceCommandSurface = {
|
||||
cleanupSubscriptionsByPrefix: RuntimeSubscriptionRegistry['cleanupByPrefix']
|
||||
cleanupSubscriptionsForConnection: RuntimeSubscriptionRegistry['cleanupForConnection']
|
||||
cleanupSubscriptionIfOwnedByConnection: RuntimeSubscriptionRegistry['cleanupIfOwnedByConnection']
|
||||
cleanupSubscriptionIfOwnedByConnectionAndWait: RuntimeSubscriptionRegistry['cleanupIfOwnedByConnectionAndWait']
|
||||
onNotificationDispatched: RuntimeMobileNotificationController['onDispatched']
|
||||
getMobileNotificationListenerCount: RuntimeMobileNotificationController['getListenerCount']
|
||||
dispatchMobileNotification: RuntimeMobileNotificationController['dispatch']
|
||||
@@ -103,6 +104,8 @@ export function installRuntimeServiceCommandSurface(
|
||||
cleanupSubscriptionsForConnection: subscriptions.cleanupForConnection.bind(subscriptions),
|
||||
cleanupSubscriptionIfOwnedByConnection:
|
||||
subscriptions.cleanupIfOwnedByConnection.bind(subscriptions),
|
||||
cleanupSubscriptionIfOwnedByConnectionAndWait:
|
||||
subscriptions.cleanupIfOwnedByConnectionAndWait.bind(subscriptions),
|
||||
onNotificationDispatched: notifications.onDispatched.bind(notifications),
|
||||
getMobileNotificationListenerCount: notifications.getListenerCount.bind(notifications),
|
||||
dispatchMobileNotification: notifications.dispatch.bind(notifications),
|
||||
|
||||
@@ -42,18 +42,40 @@ export class RuntimeSubscriptionRegistry {
|
||||
}
|
||||
|
||||
cleanupIfOwnedByConnection(subscriptionId: string, connectionId?: string): boolean {
|
||||
if (!connectionId) {
|
||||
const verdict = this.resolveConnectionOwnership(subscriptionId, connectionId)
|
||||
if (verdict === 'cleanup') {
|
||||
this.cleanup(subscriptionId)
|
||||
return true
|
||||
}
|
||||
return verdict !== 'foreign'
|
||||
}
|
||||
|
||||
// Why: an unwatch that returns before the underlying watcher is released lets a fast
|
||||
// rewatch hold two watchers on one path and double-deliver change events.
|
||||
async cleanupIfOwnedByConnectionAndWait(
|
||||
subscriptionId: string,
|
||||
connectionId?: string
|
||||
): Promise<boolean> {
|
||||
const verdict = this.resolveConnectionOwnership(subscriptionId, connectionId)
|
||||
if (verdict === 'cleanup') {
|
||||
await this.cleanupAndWait(subscriptionId)
|
||||
}
|
||||
return verdict !== 'foreign'
|
||||
}
|
||||
|
||||
private resolveConnectionOwnership(
|
||||
subscriptionId: string,
|
||||
connectionId?: string
|
||||
): 'cleanup' | 'absent' | 'foreign' {
|
||||
if (!connectionId) {
|
||||
return 'cleanup'
|
||||
}
|
||||
// An unregistered id is already gone, not refused.
|
||||
if (!this.cleanups.has(subscriptionId)) {
|
||||
return true
|
||||
return 'absent'
|
||||
}
|
||||
if (this.connectionBySubscription.get(subscriptionId) !== connectionId) {
|
||||
return false
|
||||
}
|
||||
this.cleanup(subscriptionId)
|
||||
return true
|
||||
return this.connectionBySubscription.get(subscriptionId) === connectionId
|
||||
? 'cleanup'
|
||||
: 'foreign'
|
||||
}
|
||||
|
||||
cleanup(subscriptionId: string): void {
|
||||
|
||||
Reference in New Issue
Block a user