From 60c18256b6f857bfbbf9c2db8944295de16c14b8 Mon Sep 17 00:00:00 2001 From: jinye Date: Tue, 11 Aug 2026 15:09:28 +0800 Subject: [PATCH] feat(cli): clean up OpenAI logs in non-interactive sessions (#8893) * feat(cli): clean up OpenAI logs in non-interactive sessions Co-authored-by: Qwen-Coder * codex: address PR review feedback (#8893) Co-authored-by: Qwen-Coder * codex: address PR review feedback (#8893) Co-authored-by: Qwen-Coder --------- Co-authored-by: Qwen-Coder --- docs/design/openai-log-retention.md | 12 +- docs/users/configuration/settings.md | 2 +- .../cli/src/acp-integration/acpAgent.test.ts | 11 + packages/cli/src/acp-integration/acpAgent.ts | 2 + packages/cli/src/config/settingsSchema.ts | 2 +- packages/cli/src/gemini.test.tsx | 145 ++++++- packages/cli/src/gemini.tsx | 72 ++-- .../src/utils/housekeeping/cleanup.test.ts | 102 ++++- .../cli/src/utils/housekeeping/cleanup.ts | 23 +- .../non-interactive-scheduler.test.ts | 368 ++++++++++++++++++ .../cli/src/utils/housekeeping/scheduler.ts | 240 +++++++++++- .../utils/housekeeping/throttledOnce.test.ts | 87 ++++- .../src/utils/housekeeping/throttledOnce.ts | 46 ++- .../schemas/settings.schema.json | 2 +- 14 files changed, 1028 insertions(+), 86 deletions(-) create mode 100644 packages/cli/src/utils/housekeeping/non-interactive-scheduler.test.ts diff --git a/docs/design/openai-log-retention.md b/docs/design/openai-log-retention.md index 7a32696c35..01efe9e82e 100644 --- a/docs/design/openai-log-retention.md +++ b/docs/design/openai-log-retention.md @@ -8,7 +8,7 @@ Historical files have one in-tree reader, which searches only recent logs for th ## Design -Interactive sessions register an OpenAI log cleaner in the existing background housekeeping pipeline. The cleaner runs at most once per resolved log directory per day and uses a dedicated retention setting with a seven-day default. A zero value uses the housekeeping minimum of approximately one hour. +Interactive sessions register an OpenAI log cleaner in the existing background housekeeping pipeline. Non-interactive CLI, stream-json SDK, and ACP sessions register the same cleaner through a process-local queue. A completed cleaner runs at most once per resolved log directory per day and uses a dedicated retention setting with a seven-day default. A zero value uses the housekeeping minimum of approximately one hour. Deletion is restricted to the exact filename shape emitted by `OpenAILogger`: a UTC timestamp, an eight-character hexadecimal ID, and an optional sanitized diagnostic suffix. Prefix-only lookalikes are never deleted. The cleaner streams directory entries and processes at most 20 files concurrently so the first sweep of a very large directory does not retain the full listing in memory. @@ -16,6 +16,10 @@ The UTC date in a valid filename avoids one `stat` call per file. Files on the c The first-pass scheduler checks both the global file-history marker and the marker for the resolved OpenAI log directory. A missing or stale OpenAI marker selects the one-minute catch-up delay even when file-history cleanup ran recently. +One-shot and stream-json cleanup starts immediately before model-capable execution, while ACP registers a target after each workspace session config initializes successfully. Cleanup never blocks startup or session creation. One directory is scanned at a time per process, while the existing cross-process marker and lock prevent duplicate work between interactive, headless, SDK, and ACP processes. A long-lived process retries completed work daily, a fresh marker when its remaining interval expires, lock contention after one minute, and failures after ten minutes. + +On exit, the non-interactive queue rejects new targets, discards queued work, and aborts the active scan. Cancellation is checked for every directory entry and between bounded file batches. The process waits up to 250 milliseconds for the current filesystem operation to settle and release its lock; abrupt termination falls back to the existing one-hour stale-lock recovery. + ## Configuration ownership The default log directory is relative to the workspace, so its merged workspace retention setting and its directory have the same owner. @@ -24,7 +28,9 @@ A custom log directory can be shared by multiple workspaces, but its flat files ## Scope -The housekeeping scheduler starts only for interactive sessions. Headless CLI and SDK-only processes can still write logs without starting a sweep; the setting documentation states this limitation. Moving cleanup onto the write path remains a separate follow-up because it changes core logging behavior and process coordination. +The CLI lifecycle covers interactive sessions, one-shot headless invocations, the TypeScript, Python, and Java stream-json SDK transports, and ACP children used by IDE and daemon clients. ACP resolves a cleanup target for each initialized workspace session rather than its bootstrap working directory. Direct core embeddings and in-memory channels that bypass the CLI lifecycle remain outside this design. + +Short processes provide best-effort progress: the flat directory layout has no portable persistent iteration cursor, so an interrupted process may revisit the same directory prefix next time. Long-lived stream-json and ACP processes keep their iterator alive through EOF. Moving cleanup onto the write path or changing to a sharded/indexed log layout remains out of scope because either would change core logging behavior and compatibility. ## Alternatives considered @@ -35,4 +41,4 @@ The housekeeping scheduler starts only for interactive sessions. Headless CLI an ## Verification -Focused tests cover writer-format recognition, preservation of lookalikes, cutoff boundaries, custom-directory policy ownership, default and zero retention, oversized retention values, per-directory throttling, settings fallback, root scan failures, and catch-up delay selection. The streaming loop is inspected directly to confirm that only one bounded batch is retained. +Focused tests cover writer-format recognition, preservation of lookalikes, cutoff boundaries, cancellation, custom-directory policy ownership, default and zero retention, oversized retention values, per-directory throttling, settings fallback, root scan failures, and catch-up delay selection. Non-interactive scheduler tests cover directory deduplication, FIFO serialization, retry timing, graceful stop, and the CLI/ACP lifecycle hooks. The streaming loop is inspected directly to confirm that only one bounded batch is retained. diff --git a/docs/users/configuration/settings.md b/docs/users/configuration/settings.md index 4a8ac57ece..8ef92c17b7 100644 --- a/docs/users/configuration/settings.md +++ b/docs/users/configuration/settings.md @@ -181,7 +181,7 @@ Settings are organized into categories. Most settings should be placed within th | `model.skipStartupContext` | boolean | Skips sending the startup workspace context (environment summary and acknowledgement) at the beginning of each session. Enable this if you prefer to provide context manually or want to save tokens on startup. | `false` | | `model.enableOpenAILogging` | boolean | Enables logging of OpenAI API calls for debugging and analysis. When enabled, API requests and responses are logged to JSON files. | `false` | | `model.openAILoggingDir` | string | Custom directory path for OpenAI API logs. If not specified, defaults to `logs/openai` in the current working directory. Supports absolute paths, relative paths (resolved from current working directory), and `~` expansion (home directory). | `undefined` | -| `model.openAILogRetentionDays` | number | Days to retain OpenAI API log files written when `model.enableOpenAILogging` is on. Log files older than this are removed by an interactive-session background housekeeping pass that runs at most once per day; headless and SDK-only use does not start the pass. `0` = minimum retention (~1 hour). For a custom `model.openAILoggingDir`, configure retention at user or system scope; workspace-scoped retention is skipped because one custom directory can be shared by multiple workspaces. Changes take effect after restart. | `7` | +| `model.openAILogRetentionDays` | number | Days to retain OpenAI API log files written when `model.enableOpenAILogging` is on. Completed background housekeeping passes run at most once per day in interactive, headless, stream-json SDK, and ACP sessions. Short-lived non-interactive processes make best-effort progress, while persistent processes scan to completion. `0` = minimum retention (~1 hour). For a custom `model.openAILoggingDir`, configure retention at user or system scope; workspace-scoped retention is skipped because one custom directory can be shared by multiple workspaces. Changes take effect after restart. | `7` | **Example model.generationConfig:** diff --git a/packages/cli/src/acp-integration/acpAgent.test.ts b/packages/cli/src/acp-integration/acpAgent.test.ts index 9ad32f4bb9..7d41d9da52 100644 --- a/packages/cli/src/acp-integration/acpAgent.test.ts +++ b/packages/cli/src/acp-integration/acpAgent.test.ts @@ -24,6 +24,9 @@ import { ACP_EVENT_LOOP_STALL_RESTART_MS } from '@qwen-code/channel-base'; const { mockRunExitCleanup } = vi.hoisted(() => ({ mockRunExitCleanup: vi.fn().mockResolvedValue(undefined), })); +const { mockStartNonInteractiveOpenAILogHousekeeping } = vi.hoisted(() => ({ + mockStartNonInteractiveOpenAILogHousekeeping: vi.fn(), +})); const { mockMcpPoolDrainAll } = vi.hoisted(() => ({ mockMcpPoolDrainAll: vi .fn() @@ -32,6 +35,10 @@ const { mockMcpPoolDrainAll } = vi.hoisted(() => ({ vi.mock('../utils/cleanup.js', () => ({ runExitCleanup: mockRunExitCleanup, })); +vi.mock('../utils/housekeeping/scheduler.js', () => ({ + startNonInteractiveOpenAILogHousekeeping: + mockStartNonInteractiveOpenAILogHousekeeping, +})); // Mock the ACP SDK const { mockConnectionState } = vi.hoisted(() => { @@ -2963,6 +2970,10 @@ describe('QwenAgent MCP SSE/HTTP support', () => { ); } expect(attributes['session.id']).toBe('test-session-id'); + expect(mockStartNonInteractiveOpenAILogHousekeeping).toHaveBeenCalledWith( + innerConfig, + expect.any(Object), + ); mockConnectionState.resolve(); await agentPromise; diff --git a/packages/cli/src/acp-integration/acpAgent.ts b/packages/cli/src/acp-integration/acpAgent.ts index 0401b7d5b1..f11d8f9341 100644 --- a/packages/cli/src/acp-integration/acpAgent.ts +++ b/packages/cli/src/acp-integration/acpAgent.ts @@ -248,6 +248,7 @@ import { import { runWithAcpRuntimeOutputDir } from './runtimeOutputDirContext.js'; import { ACP_ERROR_CODES } from './errorCodes.js'; import { runExitCleanup } from '../utils/cleanup.js'; +import { startNonInteractiveOpenAILogHousekeeping } from '../utils/housekeeping/scheduler.js'; import { appEvents, AppEvent } from '../utils/events.js'; import { setLanguageAsync, @@ -12002,6 +12003,7 @@ class QwenAgent implements Agent { this.cleanupUnstoredConfig(config), ); } + startNonInteractiveOpenAILogHousekeeping(config, settings); // ACP sessions served to WebUI clients are interactive: MCP tools can // arrive progressively, but session creation/loading must not wait for a // slow or wedged server discovery. diff --git a/packages/cli/src/config/settingsSchema.ts b/packages/cli/src/config/settingsSchema.ts index 73888e84be..a805c6f83a 100644 --- a/packages/cli/src/config/settingsSchema.ts +++ b/packages/cli/src/config/settingsSchema.ts @@ -1602,7 +1602,7 @@ const SETTINGS_SCHEMA = { default: DEFAULT_OPENAI_LOG_RETENTION_DAYS, minimum: 0, description: - 'Number of days to retain OpenAI API log files written when enableOpenAILogging is on. Log files older than this are removed by an interactive-session background housekeeping pass that runs at most once per day. Set to 0 for minimum retention (~1 hour). For a custom openAILoggingDir, configure this at user or system scope; workspace-scoped retention is skipped because one directory can be shared by multiple workspaces.', + 'Number of days to retain OpenAI API log files written when enableOpenAILogging is on. Completed background housekeeping passes run at most once per day in interactive, headless, stream-json SDK, and ACP sessions. Short-lived non-interactive processes make best-effort progress, while persistent processes scan to completion. Set to 0 for minimum retention (~1 hour). For a custom openAILoggingDir, configure this at user or system scope; workspace-scoped retention is skipped because one directory can be shared by multiple workspaces.', showInDialog: false, }, generationConfig: { diff --git a/packages/cli/src/gemini.test.tsx b/packages/cli/src/gemini.test.tsx index 52191c2634..09f7820792 100644 --- a/packages/cli/src/gemini.test.tsx +++ b/packages/cli/src/gemini.test.tsx @@ -45,6 +45,10 @@ const mockHandleListExtensions = vi.hoisted(() => vi.fn()); const mockStartEarlyStartupPrefetches = vi.hoisted(() => vi.fn()); const mockStartPostRenderPrefetches = vi.hoisted(() => vi.fn()); const mockRunAcpAgent = vi.hoisted(() => vi.fn()); +const mockStartNonInteractiveOpenAILogHousekeeping = vi.hoisted(() => vi.fn()); +const mockStopNonInteractiveOpenAILogHousekeeping = vi.hoisted(() => + vi.fn(async () => {}), +); const mockUpdateBeforeRelaunch = vi.hoisted(() => vi.fn()); const mockGetInstallationInfo = vi.hoisted(() => vi.fn()); const lspConfigWatcherMock = vi.hoisted(() => ({ @@ -193,6 +197,18 @@ vi.mock('./acp-integration/acpAgent.js', () => ({ runAcpAgent: (...args: unknown[]) => mockRunAcpAgent(...args), })); +vi.mock('./utils/housekeeping/scheduler.js', async (importOriginal) => { + const actual = + await importOriginal(); + return { + ...actual, + startNonInteractiveOpenAILogHousekeeping: (...args: unknown[]) => + mockStartNonInteractiveOpenAILogHousekeeping(...args), + stopNonInteractiveOpenAILogHousekeeping: () => + mockStopNonInteractiveOpenAILogHousekeeping(), + }; +}); + vi.mock('./commands/extensions/list.js', () => ({ handleList: mockHandleListExtensions, })); @@ -261,6 +277,8 @@ describe('gemini.tsx main function', () => { []; beforeEach(() => { + mockStartNonInteractiveOpenAILogHousekeeping.mockClear(); + mockStopNonInteractiveOpenAILogHousekeeping.mockClear(); lspConfigWatcherMock.instances.length = 0; mockUpdateBeforeRelaunch.mockResolvedValue(true); mockGetInstallationInfo.mockReturnValue({ @@ -308,12 +326,10 @@ describe('gemini.tsx main function', () => { } const currentListeners = process.listeners('unhandledRejection'); - const addedListener = currentListeners.find( - (listener) => !initialUnhandledRejectionListeners.includes(listener), - ); - - if (addedListener) { - process.removeListener('unhandledRejection', addedListener); + for (const listener of currentListeners) { + if (!initialUnhandledRejectionListeners.includes(listener)) { + process.removeListener('unhandledRejection', listener); + } } vi.restoreAllMocks(); }); @@ -1030,7 +1046,10 @@ describe('gemini.tsx main function', () => { ); mockWriteStderrLine.mockClear(); - vi.mocked(cleanupModule.runExitCleanup).mockResolvedValue(undefined); + const runExitCleanupMock = vi.mocked(cleanupModule.runExitCleanup); + runExitCleanupMock.mockResolvedValue(undefined); + const cleanupRegistrationStart = vi.mocked(cleanupModule.registerCleanup) + .mock.calls.length; vi.spyOn(initializerModule, 'initializeApp').mockResolvedValue({ authError: null, themeError: null, @@ -1042,7 +1061,9 @@ describe('gemini.tsx main function', () => { userStartupWarningsModule, 'getUserStartupWarnings', ).mockResolvedValue([]); - vi.spyOn(nonInteractiveModule, 'runNonInteractive').mockResolvedValue(0); + const runNonInteractiveSpy = vi + .spyOn(nonInteractiveModule, 'runNonInteractive') + .mockResolvedValue(0); let initialized = false; const configStub = { @@ -1121,6 +1142,48 @@ describe('gemini.tsx main function', () => { expect.any(Object), { deferIdeConnection: false }, ); + expect(mockStartNonInteractiveOpenAILogHousekeeping).toHaveBeenCalledWith( + configStub, + expect.any(Object), + ); + const housekeepingCleanup = vi.mocked(cleanupModule.registerCleanup).mock + .calls[cleanupRegistrationStart]?.[0]; + expect(housekeepingCleanup).toBeTypeOf('function'); + await housekeepingCleanup?.(); + expect(mockStopNonInteractiveOpenAILogHousekeeping).toHaveBeenCalledOnce(); + + const runFailure = new Error('headless run failed'); + runNonInteractiveSpy.mockRejectedValueOnce(runFailure); + process.env['QWEN_CODE_NO_RELAUNCH'] = 'true'; + Object.defineProperty(process.stdin, 'isTTY', { + value: true, + configurable: true, + }); + const secondProcessExitSpy = vi + .spyOn(process, 'exit') + .mockImplementation((code) => { + throw new MockProcessExitError(code); + }); + try { + await expect(main()).rejects.toBe(runFailure); + } finally { + secondProcessExitSpy.mockRestore(); + if (originalIsTTY) { + Object.defineProperty(process.stdin, 'isTTY', originalIsTTY); + } else { + delete (process.stdin as { isTTY?: unknown }).isTTY; + } + if (originalNoRelaunch !== undefined) { + process.env['QWEN_CODE_NO_RELAUNCH'] = originalNoRelaunch; + } else { + delete process.env['QWEN_CODE_NO_RELAUNCH']; + } + } + + expect(runExitCleanupMock).toHaveBeenCalledTimes(2); + expect(mockStartNonInteractiveOpenAILogHousekeeping).toHaveBeenCalledTimes( + 2, + ); }); it('creates non-interactive prompt ids that preserve session correlation', () => { @@ -1343,7 +1406,7 @@ describe('gemini.tsx main function', () => { processExitSpy.mockRestore(); }); - it('invokes runNonInteractiveStreamJson and performs cleanup in stream-json mode', async () => { + it('starts housekeeping and cleans up stream-json success and failure', async () => { const originalIsTTY = Object.getOwnPropertyDescriptor( process.stdin, 'isTTY', @@ -1382,6 +1445,8 @@ describe('gemini.tsx main function', () => { vi.mocked(cleanupModule.cleanupCheckpoints).mockResolvedValue(undefined); vi.mocked(cleanupModule.registerCleanup).mockImplementation(() => () => {}); + const cleanupRegistrationStart = vi.mocked(cleanupModule.registerCleanup) + .mock.calls.length; const runExitCleanupMock = vi.mocked(cleanupModule.runExitCleanup); runExitCleanupMock.mockResolvedValue(undefined); vi.spyOn(initializerModule, 'initializeApp').mockResolvedValue({ @@ -1494,7 +1559,53 @@ describe('gemini.tsx main function', () => { expect.any(Object), { deferIdeConnection: false }, ); - expect(runExitCleanupMock).toHaveBeenCalledTimes(1); + expect(mockStartNonInteractiveOpenAILogHousekeeping).toHaveBeenCalledWith( + validatedConfig, + settingsArg, + ); + const housekeepingCleanup = vi.mocked(cleanupModule.registerCleanup).mock + .calls[cleanupRegistrationStart]?.[0]; + expect(housekeepingCleanup).toBeTypeOf('function'); + await housekeepingCleanup?.(); + expect(mockStopNonInteractiveOpenAILogHousekeeping).toHaveBeenCalledOnce(); + + const streamFailure = new Error('stream failed'); + runStreamJsonSpy.mockRejectedValueOnce(streamFailure); + Object.defineProperty(process.stdin, 'isTTY', { + value: false, + configurable: true, + }); + Object.defineProperty(process.stdin, 'isRaw', { + value: false, + configurable: true, + }); + const secondProcessExitSpy = vi + .spyOn(process, 'exit') + .mockImplementation((code) => { + throw new MockProcessExitError(code); + }); + process.env['SANDBOX'] = '1'; + try { + await expect(main()).rejects.toBe(streamFailure); + } finally { + secondProcessExitSpy.mockRestore(); + if (originalIsTTY) { + Object.defineProperty(process.stdin, 'isTTY', originalIsTTY); + } else { + delete (process.stdin as { isTTY?: unknown }).isTTY; + } + if (originalIsRaw) { + Object.defineProperty(process.stdin, 'isRaw', originalIsRaw); + } else { + delete (process.stdin as { isRaw?: unknown }).isRaw; + } + delete process.env['SANDBOX']; + } + + expect(runExitCleanupMock).toHaveBeenCalledTimes(2); + expect(mockStartNonInteractiveOpenAILogHousekeeping).toHaveBeenCalledTimes( + 2, + ); }); }); @@ -1932,6 +2043,7 @@ describe('gemini.tsx main function kitty protocol', () => { './config/config.js' ); const { loadSettings } = await import('./config/settings.js'); + const cleanupModule = await import('./utils/cleanup.js'); const initializerModule = await import('./core/initializer.js'); const initializeAppSpy = vi .spyOn(initializerModule, 'initializeApp') @@ -2059,6 +2171,19 @@ describe('gemini.tsx main function kitty protocol', () => { expect(mockStartEarlyStartupPrefetches).toHaveBeenCalledWith( expect.any(Object), ); + + const acpFailure = new Error('ACP failed'); + mockRunAcpAgent.mockRejectedValueOnce(acpFailure); + process.exit = ((code?: string | number | null | undefined) => { + throw new MockProcessExitError(code); + }) as unknown as typeof process.exit; + try { + await expect(main()).rejects.toBe(acpFailure); + } finally { + process.exit = originalExit; + } + + expect(cleanupModule.runExitCleanup).toHaveBeenCalledTimes(2); }); // Shared config/settings mocks for the interactive signal-handler tests. diff --git a/packages/cli/src/gemini.tsx b/packages/cli/src/gemini.tsx index d98dd2e250..5b5cc17f21 100644 --- a/packages/cli/src/gemini.tsx +++ b/packages/cli/src/gemini.tsx @@ -857,6 +857,16 @@ export async function main() { markAcpStartup('configConstructionEnd'); profileCheckpoint('after_load_cli_config'); + const nonInteractiveHousekeeping = + !config.isInteractive() || config.getExperimentalZedIntegration() + ? await import('./utils/housekeeping/scheduler.js') + : undefined; + if (nonInteractiveHousekeeping) { + registerCleanup(() => + nonInteractiveHousekeeping.stopNonInteractiveOpenAILogHousekeeping(), + ); + } + // Subscribe the running Config to settings changes so MCP servers // reconnect / disconnect / restart without a session restart (#3696, // sub-task 3). Skipped in bare mode (no watcher). @@ -1053,17 +1063,20 @@ export async function main() { markAcpStartup('acpImportStart'); const { runAcpAgent } = await import('./acp-integration/acpAgent.js'); markAcpStartup('acpImportEnd'); - await runAcpAgent(config, settings, argv, { - privateParentCapability: isAcpMode - ? privateAcpParentCapability - : undefined, - externalToolGuardRequired: - isAcpMode && - privateAcpParentCapability !== undefined && - privateExternalToolGuard === EXTERNAL_TOOL_GUARD_REQUIRED_VALUE, - }); - // Clean up child processes and force exit, matching other non-interactive modes - await runExitCleanup(); + try { + await runAcpAgent(config, settings, argv, { + privateParentCapability: isAcpMode + ? privateAcpParentCapability + : undefined, + externalToolGuardRequired: + isAcpMode && + privateAcpParentCapability !== undefined && + privateExternalToolGuard === EXTERNAL_TOOL_GUARD_REQUIRED_VALUE, + }); + } finally { + // Clean up child processes even when ACP setup or shutdown fails. + await runExitCleanup(); + } process.exit(0); } @@ -1277,12 +1290,19 @@ export async function main() { './nonInteractive/session.js' ); - await runNonInteractiveStreamJson( + nonInteractiveHousekeeping?.startNonInteractiveOpenAILogHousekeeping( nonInteractiveConfig, - trimmedInput.length > 0 ? trimmedInput : '', settings, ); - await runExitCleanup(); + try { + await runNonInteractiveStreamJson( + nonInteractiveConfig, + trimmedInput.length > 0 ? trimmedInput : '', + settings, + ); + } finally { + await runExitCleanup(); + } // `runNonInteractiveStreamJson` doesn't return an explicit exit // code yet, so a cleanup task that mutates `process.exitCode` // could clobber a non-zero failure signal. This is currently safe @@ -1314,18 +1334,24 @@ export async function main() { debugLogger.debug(`Session ID: ${config.getSessionId()}`); const { runNonInteractive } = await import('./nonInteractiveCli.js'); - const exitCode = await runNonInteractive( + nonInteractiveHousekeeping?.startNonInteractiveOpenAILogHousekeeping( nonInteractiveConfig, settings, - input, - prompt_id, ); - // Call cleanup before process.exit, which causes cleanup to not run. - // Capture the exit code BEFORE cleanup so any cleanup task that - // mutates process.exitCode can't silently turn a structured-output - // failure (or other explicit non-zero return from runNonInteractive) - // into a zero exit. - await runExitCleanup(); + let exitCode: number; + try { + exitCode = await runNonInteractive( + nonInteractiveConfig, + settings, + input, + prompt_id, + ); + } finally { + // Call cleanup before process.exit, which causes cleanup to not run. + await runExitCleanup(); + } + // Capture the exit code BEFORE cleanup so any cleanup task that mutates + // process.exitCode can't silently turn an explicit failure into success. process.exit(exitCode); } } diff --git a/packages/cli/src/utils/housekeeping/cleanup.test.ts b/packages/cli/src/utils/housekeeping/cleanup.test.ts index 4f80dcc846..5f7e7a2987 100644 --- a/packages/cli/src/utils/housekeeping/cleanup.test.ts +++ b/packages/cli/src/utils/housekeeping/cleanup.test.ts @@ -6,6 +6,7 @@ import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; import * as fs from 'node:fs'; +import * as fsPromises from 'node:fs/promises'; import * as path from 'node:path'; import * as os from 'node:os'; import { OpenAILogger } from '@qwen-code/qwen-code-core'; @@ -16,6 +17,8 @@ import { getCutoffDate, } from './cleanup.js'; +vi.mock('node:fs/promises', { spy: true }); + const MS_PER_HOUR = 60 * 60 * 1000; const MS_PER_DAY = 24 * 60 * 60 * 1000; const FILE_HISTORY_DIR = 'file-history'; @@ -272,6 +275,7 @@ describe('cleanupOldOpenAILogs', () => { }); afterEach(() => { + vi.restoreAllMocks(); fs.rmSync(logDir, { recursive: true, force: true }); }); @@ -295,7 +299,7 @@ describe('cleanupOldOpenAILogs', () => { logDir: path.join(logDir, 'nope'), cutoffDate: cutoff, }); - expect(r).toEqual({ removed: 0, errors: 0 }); + expect(r).toEqual({ removed: 0, errors: 0, completed: true }); }); it('removes logs whose filename date is older than the cutoff, even with a fresh mtime', async () => { @@ -306,7 +310,7 @@ describe('cleanupOldOpenAILogs', () => { new Date(), ); const r = await cleanupOldOpenAILogs({ logDir, cutoffDate: cutoff }); - expect(r).toEqual({ removed: 1, errors: 0 }); + expect(r).toEqual({ removed: 1, errors: 0, completed: true }); expect(fs.existsSync(old)).toBe(false); // The project-local log dir itself is never removed. expect(fs.existsSync(logDir)).toBe(true); @@ -317,7 +321,7 @@ describe('cleanupOldOpenAILogs', () => { const name = openAILogName(recent, 'b2c3d4e5', 'side-query-session-title'); const fresh = mkLog(name, recent); const r = await cleanupOldOpenAILogs({ logDir, cutoffDate: cutoff }); - expect(r).toEqual({ removed: 0, errors: 0 }); + expect(r).toEqual({ removed: 0, errors: 0, completed: true }); expect(fs.existsSync(fresh)).toBe(true); }); @@ -331,7 +335,7 @@ describe('cleanupOldOpenAILogs', () => { new Date(Date.now() - 30 * MS_PER_DAY), ); const r = await cleanupOldOpenAILogs({ logDir, cutoffDate: cutoff }); - expect(r).toEqual({ removed: 0, errors: 0 }); + expect(r).toEqual({ removed: 0, errors: 0, completed: true }); expect(fs.existsSync(oldPrefixedFile)).toBe(true); expect(fs.existsSync(missingId)).toBe(true); }); @@ -349,7 +353,7 @@ describe('cleanupOldOpenAILogs', () => { logDir, cutoffDate: new Date(Date.now() + MS_PER_DAY), }); - expect(r).toEqual({ removed: 1, errors: 0 }); + expect(r).toEqual({ removed: 1, errors: 0, completed: true }); expect(fs.existsSync(generated)).toBe(false); }); @@ -364,7 +368,7 @@ describe('cleanupOldOpenAILogs', () => { new Date(cutoff.getTime() + MS_PER_HOUR), ); const r = await cleanupOldOpenAILogs({ logDir, cutoffDate: cutoff }); - expect(r).toEqual({ removed: 1, errors: 0 }); + expect(r).toEqual({ removed: 1, errors: 0, completed: true }); expect(fs.existsSync(olderThanCutoff)).toBe(false); expect(fs.existsSync(newerThanCutoff)).toBe(true); }); @@ -382,12 +386,96 @@ describe('cleanupOldOpenAILogs', () => { fs.mkdirSync(dirWithMatchingName); const r = await cleanupOldOpenAILogs({ logDir, cutoffDate: cutoff }); - expect(r).toEqual({ removed: 0, errors: 0 }); + expect(r).toEqual({ removed: 0, errors: 0, completed: true }); expect(fs.existsSync(note)).toBe(true); expect(fs.existsSync(otherLog)).toBe(true); expect(fs.existsSync(dirWithMatchingName)).toBe(true); }); + it('returns incomplete without opening the directory when already aborted', async () => { + const controller = new AbortController(); + controller.abort(); + + const r = await cleanupOldOpenAILogs({ + logDir, + cutoffDate: cutoff, + signal: controller.signal, + }); + + expect(r).toEqual({ removed: 0, errors: 0, completed: false }); + }); + + it('checks cancellation for every directory entry', async () => { + const old = new Date(Date.now() - 30 * MS_PER_DAY); + for (let i = 0; i < 10; i++) { + mkLog(`notes-${i}.txt`, old); + } + for (let i = 0; i < 10; i++) { + mkLog(openAILogName(old, i.toString(16).padStart(8, '0')), old); + } + + let checks = 0; + const signal = { + get aborted() { + checks++; + return checks > 5; + }, + } as AbortSignal; + + const r = await cleanupOldOpenAILogs({ + logDir, + cutoffDate: cutoff, + signal, + }); + + expect(r.completed).toBe(false); + expect(checks).toBeGreaterThan(5); + expect(r.removed).toBeLessThan(10); + }); + + it('settles a partially populated deletion batch before returning', async () => { + const old = new Date(Date.now() - 30 * MS_PER_DAY); + for (let i = 0; i < 10; i++) { + mkLog(openAILogName(old, i.toString(16).padStart(8, '0')), old); + } + + let unlinkStarted = false; + let releaseUnlink: (() => void) | undefined; + vi.mocked(fsPromises.unlink).mockImplementation(async (filePath) => { + unlinkStarted = true; + await new Promise((resolve) => { + releaseUnlink = resolve; + }); + fs.unlinkSync(filePath); + }); + const signal = { + get aborted() { + return unlinkStarted; + }, + } as AbortSignal; + + let settled = false; + const cleanup = cleanupOldOpenAILogs({ + logDir, + cutoffDate: cutoff, + signal, + }).then((result) => { + settled = true; + return result; + }); + + await vi.waitFor(() => expect(releaseUnlink).toBeTypeOf('function')); + await Promise.resolve(); + try { + expect(settled).toBe(false); + } finally { + releaseUnlink?.(); + } + const r = await cleanup; + + expect(r).toEqual({ removed: 1, errors: 0, completed: false }); + }); + it.skipIf(process.platform === 'win32')( 'rejects when the log directory cannot be scanned', async () => { diff --git a/packages/cli/src/utils/housekeeping/cleanup.ts b/packages/cli/src/utils/housekeeping/cleanup.ts index 22ce6ababd..ce8a5166b8 100644 --- a/packages/cli/src/utils/housekeeping/cleanup.ts +++ b/packages/cli/src/utils/housekeeping/cleanup.ts @@ -42,6 +42,11 @@ export interface OpenAILogCleanupOptions { cutoffDate: Date; /** Resolved OpenAI log directory (see resolveOpenAILogDir in core). */ logDir: string; + signal?: AbortSignal; +} + +export interface OpenAILogCleanupResult extends CleanupResult { + completed: boolean; } // cleanupPeriodDays = 0 means "minimum retention", not "delete everything @@ -143,8 +148,17 @@ const OPENAI_LOG_FILE_PATTERN = // default location lives inside the user's project checkout. export async function cleanupOldOpenAILogs( opts: OpenAILogCleanupOptions, -): Promise { - const result: CleanupResult = { removed: 0, errors: 0 }; +): Promise { + const result: OpenAILogCleanupResult = { + removed: 0, + errors: 0, + completed: true, + }; + + if (opts.signal?.aborted) { + result.completed = false; + return result; + } let dir; try { @@ -159,6 +173,11 @@ export async function cleanupOldOpenAILogs( let batch: Array> = []; for await (const entry of dir) { + if (opts.signal?.aborted) { + result.completed = false; + break; + } + const filenameDate = entry.isFile() ? OPENAI_LOG_FILE_PATTERN.exec(entry.name)?.[1] : undefined; diff --git a/packages/cli/src/utils/housekeeping/non-interactive-scheduler.test.ts b/packages/cli/src/utils/housekeeping/non-interactive-scheduler.test.ts new file mode 100644 index 0000000000..57ede7a79b --- /dev/null +++ b/packages/cli/src/utils/housekeeping/non-interactive-scheduler.test.ts @@ -0,0 +1,368 @@ +/** + * @license + * Copyright 2025 Google LLC + * SPDX-License-Identifier: Apache-2.0 + */ + +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import * as fs from 'node:fs'; +import * as os from 'node:os'; +import * as path from 'node:path'; +import type { Config } from '@qwen-code/qwen-code-core'; +import type { LoadedSettings } from '../../config/settings.js'; + +const mocks = vi.hoisted(() => ({ + cleanupOldOpenAILogs: vi.fn(), + runThrottledOnce: vi.fn(), +})); + +vi.mock('./cleanup.js', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + cleanupOldOpenAILogs: mocks.cleanupOldOpenAILogs, + }; +}); + +vi.mock('./throttledOnce.js', () => ({ + runThrottledOnce: mocks.runThrottledOnce, +})); + +import { + _resetNonInteractiveForTesting, + startNonInteractiveOpenAILogHousekeeping, + stopNonInteractiveOpenAILogHousekeeping, +} from './scheduler.js'; + +const MS_PER_MINUTE = 60 * 1000; +const MS_PER_HOUR = 60 * MS_PER_MINUTE; + +function makeConfig(logDir: string): Config { + return { + getContentGeneratorConfig: () => ({ openAILoggingDir: logDir }), + getModelsConfig: () => ({ getGenerationConfig: () => ({}) }), + getWorkingDir: () => process.cwd(), + } as unknown as Config; +} + +function makeSettings(): LoadedSettings { + return { + merged: { model: { openAILogRetentionDays: 7 } }, + isTrusted: true, + system: { settings: {} }, + systemDefaults: { settings: {} }, + user: { settings: { model: { openAILogRetentionDays: 7 } } }, + workspace: { settings: {} }, + } as unknown as LoadedSettings; +} + +describe('non-interactive OpenAI log housekeeping', () => { + let qwenHome: string; + + beforeEach(async () => { + await _resetNonInteractiveForTesting(); + qwenHome = fs.mkdtempSync( + path.join(os.tmpdir(), 'qwen-noninteractive-housekeeping-'), + ); + vi.stubEnv('QWEN_HOME', qwenHome); + mocks.cleanupOldOpenAILogs.mockReset(); + mocks.runThrottledOnce.mockReset(); + mocks.cleanupOldOpenAILogs.mockResolvedValue({ + removed: 0, + errors: 0, + completed: true, + }); + mocks.runThrottledOnce.mockImplementation( + async (_options, task: () => Promise) => + (await task()) === false + ? { status: 'incomplete' } + : { status: 'completed' }, + ); + }); + + afterEach(async () => { + vi.useRealTimers(); + await _resetNonInteractiveForTesting(); + vi.unstubAllEnvs(); + fs.rmSync(qwenHome, { recursive: true, force: true }); + }); + + it('deduplicates the same resolved log directory', async () => { + const logDir = path.join(qwenHome, 'logs'); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(logDir), + makeSettings(), + ); + startNonInteractiveOpenAILogHousekeeping( + makeConfig(logDir), + makeSettings(), + ); + + await vi.waitFor(() => + expect(mocks.cleanupOldOpenAILogs).toHaveBeenCalledOnce(), + ); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(mocks.cleanupOldOpenAILogs).toHaveBeenCalledOnce(); + }); + + it('serializes different directories in FIFO order', async () => { + const firstDir = path.join(qwenHome, 'first'); + const secondDir = path.join(qwenHome, 'second'); + let releaseFirst: ((value: unknown) => void) | undefined; + mocks.cleanupOldOpenAILogs.mockImplementation(({ logDir, signal }) => { + if (logDir === firstDir) { + return new Promise((resolve) => { + releaseFirst = resolve; + signal?.addEventListener('abort', () => { + resolve({ removed: 0, errors: 0, completed: false }); + }); + }); + } + return Promise.resolve({ removed: 0, errors: 0, completed: true }); + }); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(firstDir), + makeSettings(), + ); + startNonInteractiveOpenAILogHousekeeping( + makeConfig(secondDir), + makeSettings(), + ); + + await vi.waitFor(() => expect(releaseFirst).toBeDefined()); + expect(mocks.cleanupOldOpenAILogs).toHaveBeenCalledTimes(1); + + releaseFirst?.({ removed: 0, errors: 0, completed: true }); + await vi.waitFor(() => + expect(mocks.cleanupOldOpenAILogs).toHaveBeenCalledTimes(2), + ); + expect(mocks.cleanupOldOpenAILogs.mock.calls[1]?.[0].logDir).toBe( + secondDir, + ); + }); + + it('uses ModelsConfig as the fallback for the CLI logging directory', async () => { + const modelLogDir = path.join(qwenHome, 'from-models-config'); + const settingsLogDir = path.join(qwenHome, 'from-settings'); + const config = { + getContentGeneratorConfig: () => undefined, + getModelsConfig: () => ({ + getGenerationConfig: () => ({ openAILoggingDir: modelLogDir }), + }), + getWorkingDir: () => process.cwd(), + } as unknown as Config; + const settings = makeSettings(); + settings.merged.model = { + ...settings.merged.model, + openAILoggingDir: settingsLogDir, + }; + + startNonInteractiveOpenAILogHousekeeping(config, settings); + + await vi.waitFor(() => + expect(mocks.cleanupOldOpenAILogs).toHaveBeenCalledOnce(), + ); + expect(mocks.cleanupOldOpenAILogs.mock.calls[0]?.[0].logDir).toBe( + modelLogDir, + ); + }); + + it('prefers the initialized content-generator logging directory', async () => { + const contentGeneratorLogDir = path.join( + qwenHome, + 'from-content-generator-config', + ); + const modelLogDir = path.join(qwenHome, 'from-models-config'); + const config = { + getContentGeneratorConfig: () => ({ + openAILoggingDir: contentGeneratorLogDir, + }), + getModelsConfig: () => ({ + getGenerationConfig: () => ({ openAILoggingDir: modelLogDir }), + }), + getWorkingDir: () => process.cwd(), + } as unknown as Config; + + startNonInteractiveOpenAILogHousekeeping(config, makeSettings()); + + await vi.waitFor(() => + expect(mocks.cleanupOldOpenAILogs).toHaveBeenCalledOnce(), + ); + expect(mocks.cleanupOldOpenAILogs.mock.calls[0]?.[0].logDir).toBe( + contentGeneratorLogDir, + ); + }); + + it('retries a held lock after one minute', async () => { + vi.useFakeTimers(); + mocks.runThrottledOnce.mockResolvedValue({ status: 'locked' }); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'logs')), + makeSettings(), + ); + await vi.waitFor(() => + expect(mocks.runThrottledOnce).toHaveBeenCalledOnce(), + ); + + await vi.advanceTimersByTimeAsync(MS_PER_MINUTE - 1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(1); + await vi.advanceTimersByTimeAsync(1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(2); + }); + + it('retries a fresh marker after its remaining interval', async () => { + vi.useFakeTimers(); + mocks.runThrottledOnce + .mockResolvedValueOnce({ + status: 'fresh', + retryAfterMs: 30 * MS_PER_MINUTE, + }) + .mockResolvedValue({ status: 'completed' }); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'logs')), + makeSettings(), + ); + await vi.waitFor(() => + expect(mocks.runThrottledOnce).toHaveBeenCalledOnce(), + ); + + await vi.advanceTimersByTimeAsync(30 * MS_PER_MINUTE - 1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(1); + await vi.advanceTimersByTimeAsync(1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(2); + }); + + it.each([ + ['minimum', 5_000, MS_PER_MINUTE], + ['maximum', 25 * MS_PER_HOUR, 24 * MS_PER_HOUR], + ])( + 'clamps a fresh marker retry to the %s delay', + async (_name, retryAfterMs, expectedDelayMs) => { + vi.useFakeTimers(); + mocks.runThrottledOnce + .mockResolvedValueOnce({ status: 'fresh', retryAfterMs }) + .mockResolvedValue({ status: 'completed' }); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'logs')), + makeSettings(), + ); + await vi.waitFor(() => + expect(mocks.runThrottledOnce).toHaveBeenCalledOnce(), + ); + + await vi.advanceTimersByTimeAsync(expectedDelayMs - 1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(1); + await vi.advanceTimersByTimeAsync(1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(2); + }, + ); + + it('retries a failed cleanup after ten minutes', async () => { + vi.useFakeTimers(); + mocks.runThrottledOnce + .mockRejectedValueOnce(new Error('disk unavailable')) + .mockResolvedValue({ status: 'completed' }); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'logs')), + makeSettings(), + ); + await vi.waitFor(() => + expect(mocks.runThrottledOnce).toHaveBeenCalledOnce(), + ); + + await vi.advanceTimersByTimeAsync(10 * MS_PER_MINUTE - 1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(1); + await vi.advanceTimersByTimeAsync(1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(2); + }); + + it('runs a completed cleanup again after twenty-four hours', async () => { + vi.useFakeTimers(); + mocks.runThrottledOnce.mockResolvedValue({ status: 'completed' }); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'logs')), + makeSettings(), + ); + await vi.waitFor(() => + expect(mocks.runThrottledOnce).toHaveBeenCalledOnce(), + ); + + await vi.advanceTimersByTimeAsync(24 * 60 * MS_PER_MINUTE - 1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(1); + await vi.advanceTimersByTimeAsync(1); + expect(mocks.runThrottledOnce).toHaveBeenCalledTimes(2); + }); + + it('aborts the active scan during stop', async () => { + let observedSignal: AbortSignal | undefined; + mocks.cleanupOldOpenAILogs.mockImplementation( + ({ signal }: { signal?: AbortSignal }) => + new Promise((resolve) => { + observedSignal = signal; + signal?.addEventListener('abort', () => { + resolve({ removed: 0, errors: 0, completed: false }); + }); + }), + ); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'logs')), + makeSettings(), + ); + await vi.waitFor(() => expect(observedSignal).toBeDefined()); + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'queued')), + makeSettings(), + ); + + await stopNonInteractiveOpenAILogHousekeeping(); + + expect(observedSignal?.aborted).toBe(true); + expect(mocks.cleanupOldOpenAILogs).toHaveBeenCalledOnce(); + mocks.runThrottledOnce.mockClear(); + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'after-stop')), + makeSettings(), + ); + await Promise.resolve(); + expect(mocks.runThrottledOnce).not.toHaveBeenCalled(); + }); + + it('caps stop waiting at 250ms when a filesystem task does not settle', async () => { + vi.useFakeTimers(); + let releaseWorker: ((value: unknown) => void) | undefined; + mocks.runThrottledOnce.mockImplementation( + () => + new Promise((resolve) => { + releaseWorker = resolve; + }), + ); + + startNonInteractiveOpenAILogHousekeeping( + makeConfig(path.join(qwenHome, 'logs')), + makeSettings(), + ); + await vi.waitFor(() => expect(releaseWorker).toBeDefined()); + + let stopped = false; + const firstStop = stopNonInteractiveOpenAILogHousekeeping(); + expect(stopNonInteractiveOpenAILogHousekeeping()).toBe(firstStop); + const stopPromise = firstStop.then(() => { + stopped = true; + }); + await vi.advanceTimersByTimeAsync(249); + expect(stopped).toBe(false); + await vi.advanceTimersByTimeAsync(1); + await stopPromise; + expect(stopped).toBe(true); + + releaseWorker?.({ status: 'incomplete' }); + }); +}); diff --git a/packages/cli/src/utils/housekeeping/scheduler.ts b/packages/cli/src/utils/housekeeping/scheduler.ts index d86d1e35dd..f9f8c9d4b5 100644 --- a/packages/cli/src/utils/housekeeping/scheduler.ts +++ b/packages/cli/src/utils/housekeeping/scheduler.ts @@ -40,6 +40,9 @@ const RECENT_INTERACTION_MS = 60 * 1000; // the typical sporadic user still gets periodic cleanup". const CATCHUP_THRESHOLD_MS = 7 * 24 * 60 * 60 * 1000; const STARTUP_DELAY_CATCHUP_MS = 60 * 1000; +const NON_INTERACTIVE_LOCK_RETRY_MS = 60 * 1000; +const NON_INTERACTIVE_FAILURE_RETRY_MS = 10 * 60 * 1000; +const NON_INTERACTIVE_STOP_GRACE_MS = 250; const FILE_HISTORY_MARKER = '.file-history-cleanup'; const SUBAGENT_MARKER = '.subagent-cleanup'; @@ -47,6 +50,21 @@ const OPENAI_LOGS_MARKER = '.openai-logs-cleanup'; let started = false; +interface NonInteractiveOpenAILogJob { + target: OpenAILogCleanupTarget; + markerPath: string; + queued: boolean; + timer?: NodeJS.Timeout; +} + +const nonInteractiveJobs = new Map(); +const nonInteractiveQueue: NonInteractiveOpenAILogJob[] = []; +let activeNonInteractiveJob: NonInteractiveOpenAILogJob | undefined; +let activeNonInteractiveAbortController: AbortController | undefined; +let nonInteractiveWorker: Promise | undefined; +let nonInteractiveStopping = false; +let nonInteractiveStopPromise: Promise | undefined; + export function startBackgroundHousekeeping( config: Config, settings: LoadedSettings, @@ -132,6 +150,7 @@ function getOpenAILogCleanupTarget( try { const customLogDir = config.getContentGeneratorConfig?.()?.openAILoggingDir ?? + config.getModelsConfig?.()?.getGenerationConfig?.().openAILoggingDir ?? settings.merged.model?.openAILoggingDir; const systemRetention = settings.system?.settings.model?.openAILogRetentionDays; @@ -168,6 +187,191 @@ function getOpenAILogCleanupTarget( } } +export function startNonInteractiveOpenAILogHousekeeping( + config: Config, + settings: LoadedSettings, +): void { + if (nonInteractiveStopping) return; + + try { + const target = getOpenAILogCleanupTarget(config, settings); + if (!target || nonInteractiveJobs.has(target.logDir)) return; + + const markerPath = getOpenAILogsMarkerPath( + Storage.getGlobalQwenDir(), + target.logDir, + ); + const job: NonInteractiveOpenAILogJob = { + target, + markerPath, + queued: false, + }; + nonInteractiveJobs.set(target.logDir, job); + enqueueNonInteractiveJob(job); + } catch (err) { + debugLogger.error( + 'failed to start non-interactive OpenAI log cleanup; skipping', + err, + ); + } +} + +export function stopNonInteractiveOpenAILogHousekeeping(): Promise { + if (nonInteractiveStopPromise) return nonInteractiveStopPromise; + + nonInteractiveStopping = true; + for (const job of nonInteractiveJobs.values()) { + if (job.timer) { + clearTimeout(job.timer); + job.timer = undefined; + } + job.queued = false; + } + nonInteractiveQueue.length = 0; + activeNonInteractiveAbortController?.abort(); + + nonInteractiveStopPromise = waitForNonInteractiveWorkerToStop(); + return nonInteractiveStopPromise; +} + +function enqueueNonInteractiveJob(job: NonInteractiveOpenAILogJob): void { + if (nonInteractiveStopping || job.queued || activeNonInteractiveJob === job) { + return; + } + + job.queued = true; + nonInteractiveQueue.push(job); + startNonInteractiveWorker(); +} + +function startNonInteractiveWorker(): void { + if (nonInteractiveWorker || nonInteractiveStopping) return; + + nonInteractiveWorker = drainNonInteractiveQueue() + .catch((err) => { + debugLogger.error('non-interactive OpenAI log worker failed', err); + }) + .finally(() => { + nonInteractiveWorker = undefined; + if (nonInteractiveQueue.length > 0 && !nonInteractiveStopping) { + startNonInteractiveWorker(); + } + }); +} + +async function drainNonInteractiveQueue(): Promise { + while (!nonInteractiveStopping) { + const job = nonInteractiveQueue.shift(); + if (!job) return; + + job.queued = false; + activeNonInteractiveJob = job; + const abortController = new AbortController(); + activeNonInteractiveAbortController = abortController; + + try { + const result = await runOpenAILogCleanup( + job.target, + job.markerPath, + abortController.signal, + ); + if (nonInteractiveStopping) continue; + + switch (result.status) { + case 'completed': + scheduleNonInteractiveJob(job, RECURRING_INTERVAL_MS); + break; + case 'fresh': + scheduleNonInteractiveJob( + job, + Math.min( + RECURRING_INTERVAL_MS, + Math.max(NON_INTERACTIVE_LOCK_RETRY_MS, result.retryAfterMs), + ), + ); + break; + case 'locked': + scheduleNonInteractiveJob(job, NON_INTERACTIVE_LOCK_RETRY_MS); + break; + case 'incomplete': + break; + default: + break; + } + } catch (err) { + debugLogger.error( + `non-interactive OpenAI log cleanup failed for ${job.target.logDir}`, + err, + ); + if (!nonInteractiveStopping) { + scheduleNonInteractiveJob(job, NON_INTERACTIVE_FAILURE_RETRY_MS); + } + } finally { + activeNonInteractiveJob = undefined; + activeNonInteractiveAbortController = undefined; + } + } +} + +async function runOpenAILogCleanup( + target: OpenAILogCleanupTarget, + markerPath: string, + signal?: AbortSignal, +) { + return runThrottledOnce( + { + name: 'openai-logs-cleanup', + markerPath, + lockPath: markerPath + '.lock', + }, + async () => { + const r = await cleanupOldOpenAILogs({ + logDir: target.logDir, + cutoffDate: getCutoffDate(target.retentionDays), + signal, + }); + debugLogger.debug( + `openai-logs: removed=${r.removed} errors=${r.errors} completed=${r.completed}`, + ); + return r.completed ? undefined : false; + }, + ); +} + +function scheduleNonInteractiveJob( + job: NonInteractiveOpenAILogJob, + delayMs: number, +): void { + if (nonInteractiveStopping) return; + + job.timer = setTimeout(() => { + job.timer = undefined; + enqueueNonInteractiveJob(job); + }, delayMs); + job.timer.unref(); +} + +async function waitForNonInteractiveWorkerToStop(): Promise { + const worker = nonInteractiveWorker; + if (!worker) { + nonInteractiveJobs.clear(); + return; + } + + let timeout: NodeJS.Timeout | undefined; + try { + await Promise.race([ + worker, + new Promise((resolve) => { + timeout = setTimeout(resolve, NON_INTERACTIVE_STOP_GRACE_MS); + }), + ]); + } finally { + if (timeout) clearTimeout(timeout); + nonInteractiveJobs.clear(); + } +} + async function runPass( config: Config, settings: LoadedSettings, @@ -279,23 +483,9 @@ async function runHousekeeping( // 30-day default would retain tens of GB for heavy users. const openaiTarget = getOpenAILogCleanupTarget(config, settings); if (openaiTarget) { - const { logDir, retentionDays } = openaiTarget; - const markerPath = getOpenAILogsMarkerPath(qwenDir, logDir); - await runThrottledOnce( - { - name: 'openai-logs-cleanup', - markerPath, - lockPath: markerPath + '.lock', - }, - async () => { - const r = await cleanupOldOpenAILogs({ - logDir, - cutoffDate: getCutoffDate(retentionDays), - }); - debugLogger.debug( - `openai-logs: removed=${r.removed} errors=${r.errors}`, - ); - }, + await runOpenAILogCleanup( + openaiTarget, + getOpenAILogsMarkerPath(qwenDir, openaiTarget.logDir), ); } } @@ -307,6 +497,22 @@ async function runHousekeeping( export function _resetForTesting(): void { started = false; } +export async function _resetNonInteractiveForTesting(): Promise { + nonInteractiveStopping = true; + for (const job of nonInteractiveJobs.values()) { + if (job.timer) clearTimeout(job.timer); + } + nonInteractiveQueue.length = 0; + activeNonInteractiveAbortController?.abort(); + await nonInteractiveWorker; + + nonInteractiveJobs.clear(); + activeNonInteractiveJob = undefined; + activeNonInteractiveAbortController = undefined; + nonInteractiveWorker = undefined; + nonInteractiveStopping = false; + nonInteractiveStopPromise = undefined; +} export const _needsCatchUpForTesting = needsCatchUp; export const _getFirstPassDelayForTesting = getFirstPassDelay; export const _runHousekeepingForTesting = runHousekeeping; diff --git a/packages/cli/src/utils/housekeeping/throttledOnce.test.ts b/packages/cli/src/utils/housekeeping/throttledOnce.test.ts index 224c2f58e5..84fab0d248 100644 --- a/packages/cli/src/utils/housekeeping/throttledOnce.test.ts +++ b/packages/cli/src/utils/housekeeping/throttledOnce.test.ts @@ -6,10 +6,14 @@ import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'; import * as fs from 'node:fs'; +import * as fsPromises from 'node:fs/promises'; import * as path from 'node:path'; import * as os from 'node:os'; +import type { FileHandle } from 'node:fs/promises'; import { runThrottledOnce } from './throttledOnce.js'; +vi.mock('node:fs/promises', { spy: true }); + const MS_PER_HOUR = 60 * 60 * 1000; describe('runThrottledOnce', () => { @@ -24,6 +28,7 @@ describe('runThrottledOnce', () => { }); afterEach(() => { + vi.restoreAllMocks(); fs.rmSync(tempDir, { recursive: true, force: true }); }); @@ -33,7 +38,7 @@ describe('runThrottledOnce', () => { { name: 'test', markerPath, lockPath }, task, ); - expect(ran).toBe(true); + expect(ran).toEqual({ status: 'completed' }); expect(task).toHaveBeenCalledOnce(); expect(fs.existsSync(markerPath)).toBe(true); expect(fs.existsSync(lockPath)).toBe(false); @@ -48,10 +53,33 @@ describe('runThrottledOnce', () => { task2, ); expect(task1).toHaveBeenCalledOnce(); - expect(ran2).toBe(false); + expect(ran2.status).toBe('fresh'); + if (ran2.status === 'fresh') { + expect(ran2.retryAfterMs).toBeGreaterThan(23 * MS_PER_HOUR); + expect(ran2.retryAfterMs).toBeLessThanOrEqual(24 * MS_PER_HOUR + 1000); + } expect(task2).not.toHaveBeenCalled(); }); + it('reports only the remaining freshness interval for an old marker', async () => { + fs.writeFileSync(markerPath, ''); + const past = new Date(Date.now() - 23 * MS_PER_HOUR); + fs.utimesSync(markerPath, past, past); + const task = vi.fn(async () => {}); + + const result = await runThrottledOnce( + { name: 'test', markerPath, lockPath }, + task, + ); + + expect(result.status).toBe('fresh'); + if (result.status === 'fresh') { + expect(result.retryAfterMs).toBeGreaterThan(59 * 60 * 1000); + expect(result.retryAfterMs).toBeLessThanOrEqual(60 * 60 * 1000); + } + expect(task).not.toHaveBeenCalled(); + }); + it('runs again after marker mtime is older than interval', async () => { const task1 = vi.fn(async () => {}); await runThrottledOnce({ name: 'test', markerPath, lockPath }, task1); @@ -64,7 +92,7 @@ describe('runThrottledOnce', () => { { name: 'test', markerPath, lockPath }, task2, ); - expect(ran2).toBe(true); + expect(ran2).toEqual({ status: 'completed' }); expect(task2).toHaveBeenCalledOnce(); }); @@ -77,7 +105,26 @@ describe('runThrottledOnce', () => { runThrottledOnce({ name: 'test', markerPath, lockPath }, task), ]); expect(task).toHaveBeenCalledOnce(); - expect([a, b].filter(Boolean)).toHaveLength(1); + expect([a.status, b.status].sort()).toEqual(['completed', 'locked']); + }); + + it('rechecks marker freshness after acquiring the lock', async () => { + vi.mocked(fsPromises.open).mockImplementationOnce(async () => { + fs.writeFileSync(markerPath, 'completed elsewhere'); + return { + close: vi.fn(async () => {}), + } as unknown as FileHandle; + }); + const task = vi.fn(async () => {}); + + const result = await runThrottledOnce( + { name: 'test', markerPath, lockPath }, + task, + ); + + expect(result.status).toBe('fresh'); + expect(task).not.toHaveBeenCalled(); + expect(fs.existsSync(lockPath)).toBe(false); }); it('skips when a fresh lock exists (lock held by another process)', async () => { @@ -88,7 +135,7 @@ describe('runThrottledOnce', () => { { name: 'test', markerPath, lockPath, staleLockMs: MS_PER_HOUR }, task, ); - expect(ran).toBe(false); + expect(ran).toEqual({ status: 'locked' }); expect(task).not.toHaveBeenCalled(); // We did not own the lock, so we must not remove it. expect(fs.existsSync(lockPath)).toBe(true); @@ -105,12 +152,24 @@ describe('runThrottledOnce', () => { { name: 'test', markerPath, lockPath, staleLockMs: MS_PER_HOUR }, task, ); - expect(ran).toBe(true); + expect(ran).toEqual({ status: 'completed' }); expect(task).toHaveBeenCalledOnce(); expect(fs.existsSync(markerPath)).toBe(true); expect(fs.existsSync(lockPath)).toBe(false); }); + it('does not write marker when task reports incomplete, but releases lock', async () => { + const task = vi.fn(async () => false as const); + const result = await runThrottledOnce( + { name: 'test', markerPath, lockPath }, + task, + ); + + expect(result).toEqual({ status: 'incomplete' }); + expect(fs.existsSync(markerPath)).toBe(false); + expect(fs.existsSync(lockPath)).toBe(false); + }); + it('does not write marker when task throws, but releases lock', async () => { const task = vi.fn(async () => { throw new Error('boom'); @@ -121,4 +180,20 @@ describe('runThrottledOnce', () => { expect(fs.existsSync(markerPath)).toBe(false); expect(fs.existsSync(lockPath)).toBe(false); }); + + it('treats marker write failure as benign and still releases lock', async () => { + const task = vi.fn(async () => { + fs.mkdirSync(markerPath); + }); + + const result = await runThrottledOnce( + { name: 'test', markerPath, lockPath }, + task, + ); + + expect(result).toEqual({ status: 'completed' }); + expect(task).toHaveBeenCalledOnce(); + expect(fs.statSync(markerPath).isDirectory()).toBe(true); + expect(fs.existsSync(lockPath)).toBe(false); + }); }); diff --git a/packages/cli/src/utils/housekeeping/throttledOnce.ts b/packages/cli/src/utils/housekeeping/throttledOnce.ts index 7fb9fe2b82..66a6cbc185 100644 --- a/packages/cli/src/utils/housekeeping/throttledOnce.ts +++ b/packages/cli/src/utils/housekeeping/throttledOnce.ts @@ -26,13 +26,20 @@ export interface ThrottledOnceOptions { name: string; } +export type ThrottledOnceResult = + | { status: 'completed' } + | { status: 'fresh'; retryAfterMs: number } + | { status: 'locked' } + | { status: 'incomplete' }; + // Run task at most once per minIntervalMs per machine across concurrent // processes. Cooperative: no waiting, no retries — losers return immediately. -// Returns true if task ran, false if skipped (recently done or lock held). +// Returns why the task completed or skipped so persistent schedulers can +// retry at the right cadence without guessing from one boolean. export async function runThrottledOnce( opts: ThrottledOnceOptions, - task: () => Promise, -): Promise { + task: () => Promise, +): Promise { const minIntervalMs = opts.minIntervalMs ?? ONE_DAY_MS; const staleLockMs = opts.staleLockMs ?? STALE_LOCK_MS; @@ -44,8 +51,13 @@ export async function runThrottledOnce( () => {}, ); - if (await markerIsFresh(opts.markerPath, minIntervalMs, opts.name)) { - return false; + const firstFreshForMs = await markerFreshForMs( + opts.markerPath, + minIntervalMs, + opts.name, + ); + if (firstFreshForMs !== undefined) { + return { status: 'fresh', retryAfterMs: firstFreshForMs }; } let acquired = await tryAcquire(opts.lockPath); @@ -66,7 +78,7 @@ export async function runThrottledOnce( } if (!acquired) { debugLogger.debug(`${opts.name}: skipping, lock held`); - return false; + return { status: 'locked' }; } } @@ -74,14 +86,18 @@ export async function runThrottledOnce( // Re-check marker AFTER acquiring the lock. Closes the TOCTOU window // where another process completed the work between our initial mtime // check and our lock acquisition. One extra `stat` per run; cheap. - if (await markerIsFresh(opts.markerPath, minIntervalMs, opts.name)) { - return false; + const secondFreshForMs = await markerFreshForMs( + opts.markerPath, + minIntervalMs, + opts.name, + ); + if (secondFreshForMs !== undefined) { + return { status: 'fresh', retryAfterMs: secondFreshForMs }; } let taskCompleted = false; try { - await task(); - taskCompleted = true; + taskCompleted = (await task()) !== false; } finally { // Persist the marker only after successful task completion. Marker // write failure is treated as benign: cleanup already ran, and the @@ -99,7 +115,7 @@ export async function runThrottledOnce( } } } - return taskCompleted; + return { status: taskCompleted ? 'completed' : 'incomplete' }; } finally { await unlink(opts.lockPath).catch(() => { debugLogger.debug(`${opts.name}: lock unlink failed (harmless)`); @@ -107,22 +123,22 @@ export async function runThrottledOnce( } } -async function markerIsFresh( +async function markerFreshForMs( markerPath: string, minIntervalMs: number, name: string, -): Promise { +): Promise { try { const s = await stat(markerPath); const age = Date.now() - s.mtimeMs; if (age < minIntervalMs) { debugLogger.debug(`${name}: skipping, ran ${age}ms ago`); - return true; + return minIntervalMs - age; } } catch { // marker missing — treat as not fresh. } - return false; + return undefined; } async function tryAcquire(lockPath: string): Promise { diff --git a/packages/vscode-ide-companion/schemas/settings.schema.json b/packages/vscode-ide-companion/schemas/settings.schema.json index b30c393bf2..3ddb64a6b0 100644 --- a/packages/vscode-ide-companion/schemas/settings.schema.json +++ b/packages/vscode-ide-companion/schemas/settings.schema.json @@ -705,7 +705,7 @@ "type": "string" }, "openAILogRetentionDays": { - "description": "Number of days to retain OpenAI API log files written when enableOpenAILogging is on. Log files older than this are removed by an interactive-session background housekeeping pass that runs at most once per day. Set to 0 for minimum retention (~1 hour). For a custom openAILoggingDir, configure this at user or system scope; workspace-scoped retention is skipped because one directory can be shared by multiple workspaces.", + "description": "Number of days to retain OpenAI API log files written when enableOpenAILogging is on. Completed background housekeeping passes run at most once per day in interactive, headless, stream-json SDK, and ACP sessions. Short-lived non-interactive processes make best-effort progress, while persistent processes scan to completion. Set to 0 for minimum retention (~1 hour). For a custom openAILoggingDir, configure this at user or system scope; workspace-scoped retention is skipped because one directory can be shared by multiple workspaces.", "type": "number", "default": 7, "minimum": 0