diff --git a/packages/opencode/src/index.ts b/packages/opencode/src/index.ts index 0888eba..d46e34a 100644 --- a/packages/opencode/src/index.ts +++ b/packages/opencode/src/index.ts @@ -1004,7 +1004,8 @@ export async function CodexAuthPlugin( // command.execute.before reads this; if null (auth not loaded yet), // the command is rejected with a message. let cmdCtx: CommandContext | null = null - let activeRpcServer: RpcServerHandle | null = null + const ownedCacheKeepManagers = new Map() + const ownedRpcServers = new Map() let sidebarStateFileForEvents: string | undefined // Per-loader poller: each plugin invocation owns its timer and callback, so @@ -1041,16 +1042,29 @@ export async function CodexAuthPlugin( backgroundQuotaRefresh.stop() for (const websocketFetch of websocketFetches) websocketFetch.close() websocketFetches.length = 0 - if (activeRpcServer) { - await activeRpcServer.stop().catch(() => {}) - const rpcGlobal = globalThis as { - __openaiAuthRpcServer?: RpcServerHandle + const cacheKeepGlobal = globalThis as { + __openaiAuthCacheKeepManagers?: Map + } + for (const [key, manager] of ownedCacheKeepManagers) { + if ( + cacheKeepGlobal.__openaiAuthCacheKeepManagers?.get(key) === manager + ) { + manager.stop() + cacheKeepGlobal.__openaiAuthCacheKeepManagers.delete(key) } - if (rpcGlobal.__openaiAuthRpcServer === activeRpcServer) { - rpcGlobal.__openaiAuthRpcServer = undefined + } + ownedCacheKeepManagers.clear() + + const rpcGlobal = globalThis as { + __openaiAuthRpcServers?: Map + } + for (const [key, rpcServer] of ownedRpcServers) { + if (rpcGlobal.__openaiAuthRpcServers?.get(key) === rpcServer) { + await rpcServer.stop().catch(() => {}) + rpcGlobal.__openaiAuthRpcServers.delete(key) } - activeRpcServer = null } + ownedRpcServers.clear() }, async event(input) { if (input.event.type !== 'session.deleted') return @@ -1179,6 +1193,11 @@ export async function CodexAuthPlugin( const auth = await getAuth() if (auth.type !== 'oauth') return {} + const rpcDir = input.directory + ? await resolveRpcDir(input.directory) + : undefined + const cacheKeepKey = rpcDir?.dir ?? getConfigPath() + // Migration: seed the multi-account store from the existing token (idempotent) await migrateIfNeeded( { @@ -1523,9 +1542,12 @@ export async function CodexAuthPlugin( return mainRefreshPromise } const cacheKeepGlobal = globalThis as { - __openaiAuthCacheKeepManager?: CacheKeepManager + __openaiAuthCacheKeepManagers?: Map } - cacheKeepGlobal.__openaiAuthCacheKeepManager?.stop() + const cacheKeepManagers = + cacheKeepGlobal.__openaiAuthCacheKeepManagers ?? new Map() + cacheKeepGlobal.__openaiAuthCacheKeepManagers = cacheKeepManagers + cacheKeepManagers.get(cacheKeepKey)?.stop() const cacheKeepManager = new CacheKeepManager({ fetchImpl: fetch, getMainToken: async () => { @@ -1563,7 +1585,8 @@ export async function CodexAuthPlugin( getWindow: () => cacheKeepWindow, getSustain: () => cacheKeepSustain, }) - cacheKeepGlobal.__openaiAuthCacheKeepManager = cacheKeepManager + cacheKeepManagers.set(cacheKeepKey, cacheKeepManager) + ownedCacheKeepManagers.set(cacheKeepKey, cacheKeepManager) async function pushQuota( snapshot: Record, @@ -1905,14 +1928,16 @@ export async function CodexAuthPlugin( } let rpcServer: RpcServerHandle | null = null - if (input.directory) { - const rpcDir = await resolveRpcDir(input.directory) + if (rpcDir) { const rpcGlobal = globalThis as { - __openaiAuthRpcServer?: RpcServerHandle + __openaiAuthRpcServers?: Map } - if (rpcGlobal.__openaiAuthRpcServer) { - await rpcGlobal.__openaiAuthRpcServer.stop().catch(() => {}) - rpcGlobal.__openaiAuthRpcServer = undefined + const rpcServers = rpcGlobal.__openaiAuthRpcServers ?? new Map() + rpcGlobal.__openaiAuthRpcServers = rpcServers + const existingRpcServer = rpcServers.get(rpcDir.dir) + if (existingRpcServer) { + await existingRpcServer.stop().catch(() => {}) + rpcServers.delete(rpcDir.dir) } try { rpcServer = await startRpcServer({ @@ -1934,8 +1959,8 @@ export async function CodexAuthPlugin( return { text: payload.text, knobs: payload.knobs } }, }) - rpcGlobal.__openaiAuthRpcServer = rpcServer - activeRpcServer = rpcServer + rpcServers.set(rpcDir.dir, rpcServer) + ownedRpcServers.set(rpcDir.dir, rpcServer) } catch { // RPC is best-effort; the plugin must not fail if the port file // can't be written (e.g. missing directory in test environments). @@ -3410,26 +3435,6 @@ export async function CodexAuthPlugin( ).catch(() => {}) return finalResponse }, - async dispose() { - backgroundQuotaRefresh.stop() - cacheKeepManager.stop() - if ( - cacheKeepGlobal.__openaiAuthCacheKeepManager === cacheKeepManager - ) { - cacheKeepGlobal.__openaiAuthCacheKeepManager = undefined - } - fallbackManager.stopBackgroundRefresh() - if (activeRpcServer) { - await activeRpcServer.stop().catch(() => {}) - const rpcGlobal = globalThis as { - __openaiAuthRpcServer?: RpcServerHandle - } - if (rpcGlobal.__openaiAuthRpcServer === activeRpcServer) { - rpcGlobal.__openaiAuthRpcServer = undefined - } - activeRpcServer = null - } - }, } }, methods: [ diff --git a/packages/opencode/src/rpc/notifications.ts b/packages/opencode/src/rpc/notifications.ts index 1464063..cedd89c 100644 --- a/packages/opencode/src/rpc/notifications.ts +++ b/packages/opencode/src/rpc/notifications.ts @@ -5,7 +5,6 @@ const TUI_CONNECTED_WINDOW_MS = 3_000 let queue: RpcNotification[] = [] let nextId = 1 -let lastDrainAtAny = 0 const lastDrainAtBySession = new Map() export function pushNotification( @@ -21,7 +20,6 @@ export function drainNotifications( sessionId?: string, ): RpcNotification[] { const now = Date.now() - lastDrainAtAny = now if (sessionId !== undefined) lastDrainAtBySession.set(sessionId, now) const matches = (n: RpcNotification) => sessionId === undefined || @@ -30,25 +28,21 @@ export function drainNotifications( if (lastReceivedId > 0) { queue = queue.filter((n) => { if (n.id > lastReceivedId) return true - if (sessionId === undefined) return false + if (sessionId === undefined) return true return n.sessionId !== sessionId }) } return queue.filter((n) => n.id > lastReceivedId && matches(n)) } -export function isTuiConnected(sessionId?: string): boolean { +export function isTuiConnected(sessionId: string): boolean { const now = Date.now() - if (sessionId !== undefined) { - const at = lastDrainAtBySession.get(sessionId) ?? 0 - return at > 0 && now - at < TUI_CONNECTED_WINDOW_MS - } - return lastDrainAtAny > 0 && now - lastDrainAtAny < TUI_CONNECTED_WINDOW_MS + const at = lastDrainAtBySession.get(sessionId) ?? 0 + return at > 0 && now - at < TUI_CONNECTED_WINDOW_MS } export function resetNotificationsForTest(): void { queue = [] nextId = 1 - lastDrainAtAny = 0 lastDrainAtBySession.clear() } diff --git a/packages/opencode/src/rpc/rpc-server.ts b/packages/opencode/src/rpc/rpc-server.ts index 87e946e..2c57b14 100644 --- a/packages/opencode/src/rpc/rpc-server.ts +++ b/packages/opencode/src/rpc/rpc-server.ts @@ -1,5 +1,5 @@ import { randomBytes, timingSafeEqual } from 'node:crypto' -import { unlink } from 'node:fs/promises' +import { readFile, unlink } from 'node:fs/promises' import { createServer, type IncomingMessage, @@ -78,6 +78,7 @@ export async function startRpcServer( // every endpoint holding a dead connection for 90s. const handlerTimeoutMs = options.timeoutMs ?? 90_000 const receiptTimeoutMs = options.receiptTimeoutMs ?? 2_000 + let warnedMissingNotificationSession = false const server = createServer((req, res) => { req.setTimeout(handlerTimeoutMs, () => { req.socket.destroy() @@ -107,9 +108,17 @@ export async function startRpcServer( const body = await readBody(req) const params = JSON.parse(body || '{}') as Record if (method === 'pending-notifications') { + const sessionId = + typeof params.sessionId === 'string' ? params.sessionId : undefined + if (sessionId === undefined && !warnedMissingNotificationSession) { + warnedMissingNotificationSession = true + log.warn('rpc notification drain missing session id', { + pid: process.pid, + }) + } const messages = options.drain( Number(params.lastReceivedId ?? 0), - typeof params.sessionId === 'string' ? params.sessionId : undefined, + sessionId, ) return json(200, { messages }) } @@ -164,9 +173,12 @@ export async function startRpcServer( token, async stop() { await new Promise((resolve) => server.close(() => resolve())) - await unlink(join(options.dir, `port-${process.pid}.json`)).catch( - () => {}, - ) + const portFile = join(options.dir, `port-${process.pid}.json`) + const current = await readFile(portFile, 'utf8') + .then((raw) => JSON.parse(raw) as { port?: unknown; token?: unknown }) + .catch(() => undefined) + if (current?.port === port && current.token === token) + await unlink(portFile).catch(() => {}) }, } } diff --git a/packages/opencode/src/tests/cachekeep.test.ts b/packages/opencode/src/tests/cachekeep.test.ts index f7e749f..c1fb1fb 100644 --- a/packages/opencode/src/tests/cachekeep.test.ts +++ b/packages/opencode/src/tests/cachekeep.test.ts @@ -2,6 +2,8 @@ import { afterEach, beforeEach, describe, expect, mock, test } from 'bun:test' import { mkdtemp, readdir, rm, writeFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' + +import { getConfigPath } from '../config' import type { AccountStorage } from '../core/accounts' import { buildKeepwarmBody, @@ -2466,7 +2468,9 @@ describe('CacheKeepManager token resolution', () => { if (!loaderResult?.fetch) throw new Error('No fetch override') const cacheKeepGlobal = globalThis as any - const mgr = cacheKeepGlobal.__openaiAuthCacheKeepManager + const mgr = cacheKeepGlobal.__openaiAuthCacheKeepManagers?.get( + getConfigPath(), + ) expect(mgr).toBeDefined() const mockFetch = mock(async () => new Response('{}')) @@ -2530,7 +2534,7 @@ describe('RPC server dispose', () => { await rm(tempDir, { recursive: true, force: true }) }) - test('RPC server stops and unlinks port file on loader dispose', async () => { + test('loader options do not expose an RPC lifecycle dispose hook', async () => { const originalRpcDir = process.env.OPENCODE_OPENAI_AUTH_RPC_DIR process.env.OPENCODE_OPENAI_AUTH_RPC_DIR = tempDir @@ -2564,25 +2568,19 @@ describe('RPC server dispose', () => { ) // Verify port file exists in tempDir - let files = await readdir(tempDir) + const files = await readdir(tempDir) expect( files.some((f) => f.startsWith('port-') && f.endsWith('.json')), ).toBe(true) - // Dispose the loader - await loaderResult?.dispose?.() - - // Verify port file is gone - files = await readdir(tempDir) - expect( - files.some((f) => f.startsWith('port-') && f.endsWith('.json')), - ).toBe(false) + expect(loaderResult?.dispose).toBeUndefined() + await plugin.dispose?.() } finally { process.env.OPENCODE_OPENAI_AUTH_RPC_DIR = originalRpcDir } }) - test('RPC server stops and unlinks port file on plugin dispose', async () => { + test('plugin dispose clears the RPC registry entry and unlinks the port file', async () => { const originalRpcDir = process.env.OPENCODE_OPENAI_AUTH_RPC_DIR process.env.OPENCODE_OPENAI_AUTH_RPC_DIR = tempDir @@ -2621,14 +2619,18 @@ describe('RPC server dispose', () => { files.some((f) => f.startsWith('port-') && f.endsWith('.json')), ).toBe(true) - // Dispose the plugin + const rpcGlobal = globalThis as { + __openaiAuthRpcServers?: Map + } + expect(rpcGlobal.__openaiAuthRpcServers?.size ?? 0).toBeGreaterThan(0) + await plugin.dispose?.() - // Verify port file is gone files = await readdir(tempDir) expect( files.some((f) => f.startsWith('port-') && f.endsWith('.json')), ).toBe(false) + expect(rpcGlobal.__openaiAuthRpcServers?.size ?? 0).toBe(0) } finally { process.env.OPENCODE_OPENAI_AUTH_RPC_DIR = originalRpcDir } diff --git a/packages/opencode/src/tests/command-session-isolation.test.ts b/packages/opencode/src/tests/command-session-isolation.test.ts index 15065d1..54284a2 100644 --- a/packages/opencode/src/tests/command-session-isolation.test.ts +++ b/packages/opencode/src/tests/command-session-isolation.test.ts @@ -127,7 +127,7 @@ describe('command hook session isolation', () => { experimentalWebSockets: false, }) - const loaderResult = await plugin.auth?.loader?.( + await plugin.auth?.loader?.( async () => ({ type: 'oauth', provider: 'openai', @@ -204,6 +204,6 @@ describe('command hook session isolation', () => { expect(added).toBeDefined() expect(added?.sessionId).toBe('sess-A') - await loaderResult?.dispose?.() + await plugin.dispose?.() }) }) diff --git a/packages/opencode/src/tests/integration.test.ts b/packages/opencode/src/tests/integration.test.ts index ae27ed4..56c3b7e 100644 --- a/packages/opencode/src/tests/integration.test.ts +++ b/packages/opencode/src/tests/integration.test.ts @@ -10,6 +10,7 @@ import { readFile } from 'node:fs/promises' import { tmpdir } from 'node:os' import { join } from 'node:path' import type { Hooks, PluginInput } from '@opencode-ai/plugin' +import { getConfigPath } from '../config.ts' import type { OAuthAccount } from '../core/accounts.ts' import { migrateIfNeeded } from '../core/accounts.ts' import { acquireRefreshFileLock } from '../core/refresh-file-lock.ts' @@ -4952,12 +4953,15 @@ describe('integration: active fallback routing', () => { await runCommand(hooks, 'openai-cachekeep', 'sustain on') const manager = ( globalThis as typeof globalThis & { - __openaiAuthCacheKeepManager?: { - tick(): Promise - status(): { tracked: number; sustain: boolean } - } + __openaiAuthCacheKeepManagers?: Map< + string, + { + tick(): Promise + status(): { tracked: number; sustain: boolean } + } + > } - ).__openaiAuthCacheKeepManager + ).__openaiAuthCacheKeepManagers?.get(getConfigPath()) if (!manager) throw new Error('missing cachekeep manager') await manager.tick() @@ -6900,9 +6904,12 @@ describe('integration: active fallback routing', () => { now += 30 * 60_000 const manager = ( globalThis as typeof globalThis & { - __openaiAuthCacheKeepManager?: { tick(): Promise } + __openaiAuthCacheKeepManagers?: Map< + string, + { tick(): Promise } + > } - ).__openaiAuthCacheKeepManager + ).__openaiAuthCacheKeepManagers?.get(getConfigPath()) if (!manager) throw new Error('missing cachekeep manager') await manager.tick() }, diff --git a/packages/opencode/src/tests/rpc-notifications.test.ts b/packages/opencode/src/tests/rpc-notifications.test.ts index 0df4c53..0836e03 100644 --- a/packages/opencode/src/tests/rpc-notifications.test.ts +++ b/packages/opencode/src/tests/rpc-notifications.test.ts @@ -46,6 +46,15 @@ describe('notifications', () => { expect(isTuiConnected('s1')).toBe(true) }) + test('a drain for one session does not make an unscoped probe connected', () => { + drainNotifications(0, 's2') + expect(isTuiConnected('s1')).toBe(false) + expect(isTuiConnected(undefined as never)).toBe(false) + }) + + // @ts-expect-error TUI connectivity must always be scoped to a session. + isTuiConnected() + test('queue cap evicts oldest beyond 100', () => { for (let i = 0; i < 130; i++) pushNotification(payload('openai-quota'), 's1') diff --git a/packages/opencode/src/tests/rpc-server.test.ts b/packages/opencode/src/tests/rpc-server.test.ts index 51fb081..9afb672 100644 --- a/packages/opencode/src/tests/rpc-server.test.ts +++ b/packages/opencode/src/tests/rpc-server.test.ts @@ -12,17 +12,93 @@ import { import http from 'node:http' import { tmpdir } from 'node:os' import { join } from 'node:path' +import type { PluginInput } from '@opencode-ai/plugin' +import { CodexAuthPlugin } from '../index' import { flushForTest } from '../logger' import { drainNotifications, pushNotification, resetNotificationsForTest, } from '../rpc/notifications' +import { discoverPortFile } from '../rpc/port-file' +import { resolveRpcDir } from '../rpc/rpc-dir' import { startRpcServer } from '../rpc/rpc-server' let stop: (() => Promise) | null = null let dir: string +function makePluginInput(directory: string): PluginInput { + return { + client: { + auth: { set: async () => {} }, + session: { promptAsync: async () => {} }, + } as unknown as PluginInput['client'], + project: { id: 'test', name: 'test' } as unknown as PluginInput['project'], + directory, + worktree: '/tmp/test-worktree', + experimental_workspace: { register: () => {} }, + serverUrl: new URL('http://localhost:0'), + $: {} as PluginInput['$'], + } +} + +async function loadProjectPlugin(directory: string) { + const plugin = await CodexAuthPlugin(makePluginInput(directory), { + experimentalWebSockets: false, + }) + await loadAuthPlugin(plugin) + return plugin +} + +async function loadAuthPlugin( + plugin: Awaited>, +) { + const loader = plugin.auth?.loader + if (!loader) throw new Error('missing auth loader') + const loaded = await loader( + async () => ({ + type: 'oauth', + provider: 'openai', + access: 'access-token', + refresh: 'refresh-token', + expires: Date.now() + 3600_000, + }), + { id: 'openai', label: 'OpenAI', models: [] } as never, + ) + if (!loaded) throw new Error('missing loader options') + return loaded +} + +async function writeAccountStore(path: string, accountId: string) { + const now = Date.now() + await writeFile( + path, + JSON.stringify({ + version: 1, + main: { type: 'opencode', provider: 'openai' }, + accounts: [ + { + id: accountId, + type: 'oauth', + provider: 'openai', + access: 'fallback-access', + refresh: 'fallback-refresh', + expires: now + 3600_000, + enabled: true, + addedAt: now, + lastUsed: now, + lastRefreshedAt: now, + }, + ], + }), + ) +} + +function restoreEnv(name: string, value: string | undefined) { + if (value === undefined) delete process.env[name] + else process.env[name] = value +} + afterEach(async () => { await stop?.() stop = null @@ -98,6 +174,116 @@ describe('rpc-server', () => { }) }) + test('a session-less notification drain delivers every notice but cannot prune another session', async () => { + dir = await mkdtemp(join(tmpdir(), 'oa-rpcsrv-')) + const server = await startRpcServer({ + dir, + drain: drainNotifications, + apply: async () => ({ text: 'ok', knobs: {} }), + }) + stop = server.stop + const base = `http://127.0.0.1:${server.port}` + pushNotification({ command: 'openai-quota', text: 's1', knobs: {} }, 's1') + pushNotification({ command: 'openai-account', text: 's2', knobs: {} }, 's2') + + const noSession = await fetch(`${base}/rpc/pending-notifications`, { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${server.token}`, + }, + body: JSON.stringify({ lastReceivedId: 0 }), + }) + expect(noSession.status).toBe(200) + const all = (await noSession.json()).messages as Array<{ + id: number + payload: { command: string } + }> + expect(all.map((message) => message.payload.command)).toEqual([ + 'openai-quota', + 'openai-account', + ]) + + const noSessionAck = await fetch(`${base}/rpc/pending-notifications`, { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${server.token}`, + }, + body: JSON.stringify({ lastReceivedId: all[0]?.id }), + }) + expect(noSessionAck.status).toBe(200) + expect((await noSessionAck.json()).messages).toEqual([ + expect.objectContaining({ + payload: { command: 'openai-account', text: 's2', knobs: {} }, + }), + ]) + + const s1 = await fetch(`${base}/rpc/pending-notifications`, { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${server.token}`, + }, + body: JSON.stringify({ lastReceivedId: 0, sessionId: 's1' }), + }) + expect(s1.status).toBe(200) + expect((await s1.json()).messages).toEqual([ + expect.objectContaining({ + payload: { command: 'openai-quota', text: 's1', knobs: {} }, + }), + ]) + + const s1Ack = await fetch(`${base}/rpc/pending-notifications`, { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${server.token}`, + }, + body: JSON.stringify({ lastReceivedId: all[0]?.id, sessionId: 's1' }), + }) + expect(s1Ack.status).toBe(200) + + const s2 = await fetch(`${base}/rpc/pending-notifications`, { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${server.token}`, + }, + body: JSON.stringify({ lastReceivedId: 0, sessionId: 's2' }), + }) + expect(s2.status).toBe(200) + expect((await s2.json()).messages).toEqual([ + expect.objectContaining({ + payload: { command: 'openai-account', text: 's2', knobs: {} }, + }), + ]) + }) + + test('stopping a stale server leaves its successor port file and health endpoint live', async () => { + dir = await mkdtemp(join(tmpdir(), 'oa-rpcsrv-')) + const first = await startRpcServer({ + dir, + drain: drainNotifications, + apply: async () => ({ text: 'first', knobs: {} }), + }) + const second = await startRpcServer({ + dir, + drain: drainNotifications, + apply: async () => ({ text: 'second', knobs: {} }), + }) + try { + await first.stop() + const entry = await discoverPortFile(dir, process.pid) + expect(entry?.port).toBe(second.port) + expect( + (await fetch(`http://127.0.0.1:${second.port}/health`)).status, + ).toBe(200) + } finally { + await second.stop() + } + }) + test('rejects body exceeding 1 MB byte limit', async () => { dir = await mkdtemp(join(tmpdir(), 'oa-rpcsrv-')) const server = await startRpcServer({ @@ -323,4 +509,282 @@ describe('rpc-server', () => { expect(res.status).toBe(200) expect(await res.json()).toEqual({ text: 'slow-ok', knobs: {} }) }) + + test('keeps RPC ports discoverable and applies with each project captured context', async () => { + const root = await mkdtemp(join(tmpdir(), 'oa-rpc-projects-')) + const originalFetch = globalThis.fetch + const originalStateHome = process.env.XDG_STATE_HOME + const originalConfigFile = process.env.OPENCODE_OPENAI_AUTH_FILE + const originalStateFile = process.env.OPENCODE_OPENAI_AUTH_STATE_FILE + const loaded: Array>> = [] + try { + process.env.XDG_STATE_HOME = join(root, 'state') + process.env.OPENCODE_OPENAI_AUTH_STATE_FILE = join( + root, + 'auth-state.json', + ) + globalThis.fetch = (async () => + new Response('{}')) as unknown as typeof globalThis.fetch + + const projectA = join(root, 'project-a') + const projectB = join(root, 'project-b') + await mkdir(projectA) + await mkdir(projectB) + + process.env.OPENCODE_OPENAI_AUTH_FILE = join(root, 'project-a.json') + await writeAccountStore( + process.env.OPENCODE_OPENAI_AUTH_FILE, + 'account-a', + ) + loaded.push(await loadProjectPlugin(projectA)) + + process.env.OPENCODE_OPENAI_AUTH_FILE = join(root, 'project-b.json') + await writeAccountStore( + process.env.OPENCODE_OPENAI_AUTH_FILE, + 'account-b', + ) + loaded.push(await loadProjectPlugin(projectB)) + + const rpcA = await resolveRpcDir(projectA) + const rpcB = await resolveRpcDir(projectB) + const portA = await discoverPortFile(rpcA.dir, process.pid) + const portB = await discoverPortFile(rpcB.dir, process.pid) + + expect(portA).not.toBeNull() + expect(portB).not.toBeNull() + expect(portA?.port).not.toBe(portB?.port) + + const responseA = await originalFetch( + `http://127.0.0.1:${portA?.port}/rpc/apply`, + { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${portA?.token}`, + }, + body: JSON.stringify({ + command: 'openai-account', + arguments: '', + sessionId: 'session-a', + }), + }, + ) + expect(responseA.status).toBe(200) + expect((await responseA.json()).text).toContain('account-a') + + const responseB = await originalFetch( + `http://127.0.0.1:${portB?.port}/rpc/apply`, + { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${portB?.token}`, + }, + body: JSON.stringify({ + command: 'openai-account', + arguments: '', + sessionId: 'session-b', + }), + }, + ) + expect(responseB.status).toBe(200) + expect((await responseB.json()).text).toContain('account-b') + } finally { + for (const plugin of loaded) await plugin.dispose?.() + globalThis.fetch = originalFetch + restoreEnv('XDG_STATE_HOME', originalStateHome) + restoreEnv('OPENCODE_OPENAI_AUTH_FILE', originalConfigFile) + restoreEnv('OPENCODE_OPENAI_AUTH_STATE_FILE', originalStateFile) + await rm(root, { recursive: true, force: true }) + } + }) + + test('disposal removes this test projects from the RPC and cachekeep registries', async () => { + const root = await mkdtemp(join(tmpdir(), 'oa-rpc-dispose-')) + const originalFetch = globalThis.fetch + const originalStateHome = process.env.XDG_STATE_HOME + const originalConfigFile = process.env.OPENCODE_OPENAI_AUTH_FILE + const originalStateFile = process.env.OPENCODE_OPENAI_AUTH_STATE_FILE + const loaded: Array>> = [] + try { + process.env.XDG_STATE_HOME = join(root, 'state') + process.env.OPENCODE_OPENAI_AUTH_STATE_FILE = join( + root, + 'auth-state.json', + ) + globalThis.fetch = (async () => + new Response('{}')) as unknown as typeof globalThis.fetch + + const projects: Array<{ + project: string + rpc: Awaited> + }> = [] + for (const suffix of ['a', 'b', 'c']) { + const project = join(root, `project-${suffix}`) + await mkdir(project) + process.env.OPENCODE_OPENAI_AUTH_FILE = join(root, `${suffix}.json`) + await writeAccountStore( + process.env.OPENCODE_OPENAI_AUTH_FILE, + `account-${suffix}`, + ) + loaded.push(await loadProjectPlugin(project)) + projects.push({ project, rpc: await resolveRpcDir(project) }) + } + + const registries = globalThis as typeof globalThis & { + __openaiAuthCacheKeepManagers?: Map + __openaiAuthRpcServers?: Map + } + for (const { rpc } of projects) { + expect(registries.__openaiAuthRpcServers?.get(rpc.dir)).toBeDefined() + expect( + registries.__openaiAuthCacheKeepManagers?.get(rpc.dir), + ).toBeDefined() + } + + for (const plugin of loaded) await plugin.dispose?.() + loaded.length = 0 + + for (const { rpc } of projects) { + expect(registries.__openaiAuthRpcServers?.get(rpc.dir)).toBeUndefined() + expect( + registries.__openaiAuthCacheKeepManagers?.get(rpc.dir), + ).toBeUndefined() + expect(await discoverPortFile(rpc.dir, process.pid)).toBeNull() + } + } finally { + for (const plugin of loaded) await plugin.dispose?.() + globalThis.fetch = originalFetch + restoreEnv('XDG_STATE_HOME', originalStateHome) + restoreEnv('OPENCODE_OPENAI_AUTH_FILE', originalConfigFile) + restoreEnv('OPENCODE_OPENAI_AUTH_STATE_FILE', originalStateFile) + await rm(root, { recursive: true, force: true }) + } + }) + + test('disposing a replaced plugin instance does not stop its stale RPC handle', async () => { + const root = await mkdtemp(join(tmpdir(), 'oa-rpc-replace-')) + const originalFetch = globalThis.fetch + const originalStateHome = process.env.XDG_STATE_HOME + const originalConfigFile = process.env.OPENCODE_OPENAI_AUTH_FILE + const originalStateFile = process.env.OPENCODE_OPENAI_AUTH_STATE_FILE + let first: Awaited> | undefined + let second: Awaited> | undefined + try { + process.env.XDG_STATE_HOME = join(root, 'state') + process.env.OPENCODE_OPENAI_AUTH_STATE_FILE = join( + root, + 'auth-state.json', + ) + process.env.OPENCODE_OPENAI_AUTH_FILE = join(root, 'accounts.json') + globalThis.fetch = (async () => + new Response('{}')) as unknown as typeof globalThis.fetch + + const project = join(root, 'project') + await mkdir(project) + await writeAccountStore( + process.env.OPENCODE_OPENAI_AUTH_FILE, + 'account-replace', + ) + first = await loadProjectPlugin(project) + const rpc = await resolveRpcDir(project) + const rpcServers = ( + globalThis as typeof globalThis & { + __openaiAuthRpcServers?: Map< + string, + { port: number; stop: () => Promise } + > + } + ).__openaiAuthRpcServers + const firstRpcServer = rpcServers?.get(rpc.dir) + if (!firstRpcServer) throw new Error('missing first RPC server') + second = await loadProjectPlugin(project) + + const successor = await discoverPortFile(rpc.dir, process.pid) + expect(successor).not.toBeNull() + let staleStopCalls = 0 + const stop = firstRpcServer.stop + firstRpcServer.stop = async () => { + staleStopCalls += 1 + await stop() + } + + await first.dispose?.() + expect(staleStopCalls).toBe(0) + } finally { + await second?.dispose?.() + await first?.dispose?.() + globalThis.fetch = originalFetch + restoreEnv('XDG_STATE_HOME', originalStateHome) + restoreEnv('OPENCODE_OPENAI_AUTH_FILE', originalConfigFile) + restoreEnv('OPENCODE_OPENAI_AUTH_STATE_FILE', originalStateFile) + await rm(root, { recursive: true, force: true }) + } + }) + + test('Hooks dispose clears every registry entry started by its loader runs', async () => { + const root = await mkdtemp(join(tmpdir(), 'oa-rpc-hooks-dispose-')) + const originalFetch = globalThis.fetch + const originalStateHome = process.env.XDG_STATE_HOME + const originalConfigFile = process.env.OPENCODE_OPENAI_AUTH_FILE + const originalStateFile = process.env.OPENCODE_OPENAI_AUTH_STATE_FILE + const originalRpcDir = process.env.OPENCODE_OPENAI_AUTH_RPC_DIR + let plugin: Awaited> | undefined + try { + process.env.XDG_STATE_HOME = join(root, 'state') + process.env.OPENCODE_OPENAI_AUTH_STATE_FILE = join( + root, + 'auth-state.json', + ) + process.env.OPENCODE_OPENAI_AUTH_FILE = join(root, 'accounts.json') + globalThis.fetch = (async () => + new Response('{}')) as unknown as typeof globalThis.fetch + + const project = join(root, 'project') + await mkdir(project) + await writeAccountStore( + process.env.OPENCODE_OPENAI_AUTH_FILE, + 'account-replace', + ) + plugin = await CodexAuthPlugin(makePluginInput(project), { + experimentalWebSockets: false, + }) + delete process.env.OPENCODE_OPENAI_AUTH_RPC_DIR + await loadAuthPlugin(plugin) + const firstRpc = await resolveRpcDir(project) + + process.env.OPENCODE_OPENAI_AUTH_RPC_DIR = join(root, 'alternate-rpc') + await loadAuthPlugin(plugin) + const secondRpc = await resolveRpcDir(project) + + const registries = globalThis as typeof globalThis & { + __openaiAuthRpcServers?: Map + __openaiAuthCacheKeepManagers?: Map + } + for (const rpc of [firstRpc, secondRpc]) { + expect(registries.__openaiAuthRpcServers?.get(rpc.dir)).toBeDefined() + expect( + registries.__openaiAuthCacheKeepManagers?.get(rpc.dir), + ).toBeDefined() + } + + await plugin.dispose?.() + + for (const rpc of [firstRpc, secondRpc]) { + expect(registries.__openaiAuthRpcServers?.get(rpc.dir)).toBeUndefined() + expect( + registries.__openaiAuthCacheKeepManagers?.get(rpc.dir), + ).toBeUndefined() + expect(await discoverPortFile(rpc.dir, process.pid)).toBeNull() + } + } finally { + await plugin?.dispose?.() + globalThis.fetch = originalFetch + restoreEnv('XDG_STATE_HOME', originalStateHome) + restoreEnv('OPENCODE_OPENAI_AUTH_FILE', originalConfigFile) + restoreEnv('OPENCODE_OPENAI_AUTH_STATE_FILE', originalStateFile) + restoreEnv('OPENCODE_OPENAI_AUTH_RPC_DIR', originalRpcDir) + await rm(root, { recursive: true, force: true }) + } + }) }) diff --git a/packages/opencode/src/tests/tui-packaging.test.ts b/packages/opencode/src/tests/tui-packaging.test.ts index f1e0805..4f6d5a2 100644 --- a/packages/opencode/src/tests/tui-packaging.test.ts +++ b/packages/opencode/src/tests/tui-packaging.test.ts @@ -1,5 +1,5 @@ import { describe, expect, test } from 'bun:test' -import { existsSync, readFileSync } from 'node:fs' +import { existsSync, readFileSync, statSync } from 'node:fs' import { dirname, join, relative, resolve } from 'node:path' // --------------------------------------------------------------------------- @@ -177,4 +177,34 @@ describe('tui packaging (compiled ./tui entry shim)', () => { const missing = reachable.filter((rel) => !shipped.has(rel)) expect(missing).toEqual([]) }) + + test('the built plugin bundle removes the singular RPC global', () => { + // The bundle is produced by `bun run build`, which CI runs as a separate + // step before `bun run test` (see .github/workflows/ci.yml). Reading the + // bundle directly here keeps the test dependent on the same freshness + // guarantee CI provides instead of rebuilding inside the test. + const bundle = join(PKG_DIR, 'dist', 'index.js') + if (!existsSync(bundle)) { + throw new Error( + 'Built plugin bundle is missing: dist/index.js (run `bun run build` first)', + ) + } + if (statSync(bundle).size < 1_024) { + throw new Error( + 'Built plugin bundle is unexpectedly small: dist/index.js', + ) + } + + const source = readFileSync(bundle, 'utf8') + const registryCount = source.match(/__openaiAuthRpcServers/g)?.length ?? 0 + if (registryCount < 1) { + throw new Error('Built plugin bundle is missing the RPC registry global') + } + + const singularCount = + source.match(/__openaiAuthRpcServer[^s]/g)?.length ?? 0 + if (singularCount !== 0) { + throw new Error('Built plugin bundle retains the singular RPC global') + } + }) })