From b00c49b08ba1ee1e6ddacead1e85ed974f6dfe39 Mon Sep 17 00:00:00 2001 From: Jinwoo-H Date: Tue, 8 Sep 2026 23:40:44 -0400 Subject: [PATCH] fix(runtime): make files.unwatch wait for watcher release and refuse foreign connections MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../methods/files-unwatch-ownership.test.ts | 165 ++++++++++++++++++ src/main/runtime/rpc/methods/files.ts | 12 +- .../runtime-service-command-surface.ts | 3 + .../runtime/runtime-subscription-registry.ts | 38 +++- 4 files changed, 209 insertions(+), 9 deletions(-) create mode 100644 src/main/runtime/rpc/methods/files-unwatch-ownership.test.ts diff --git a/src/main/runtime/rpc/methods/files-unwatch-ownership.test.ts b/src/main/runtime/rpc/methods/files-unwatch-ownership.test.ts new file mode 100644 index 00000000000..cb0b03433a7 --- /dev/null +++ b/src/main/runtime/rpc/methods/files-unwatch-ownership.test.ts @@ -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((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((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; 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) + }) +}) diff --git a/src/main/runtime/rpc/methods/files.ts b/src/main/runtime/rpc/methods/files.ts index ef349a22f84..b022a2b7ee6 100644 --- a/src/main/runtime/rpc/methods/files.ts +++ b/src/main/runtime/rpc/methods/files.ts @@ -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 } } diff --git a/src/main/runtime/runtime-service-command-surface.ts b/src/main/runtime/runtime-service-command-surface.ts index 19545cc76e6..8acaed830a3 100644 --- a/src/main/runtime/runtime-service-command-surface.ts +++ b/src/main/runtime/runtime-service-command-surface.ts @@ -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), diff --git a/src/main/runtime/runtime-subscription-registry.ts b/src/main/runtime/runtime-subscription-registry.ts index febab3f85ff..198874bb495 100644 --- a/src/main/runtime/runtime-subscription-registry.ts +++ b/src/main/runtime/runtime-subscription-registry.ts @@ -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 { + 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 {