diff --git a/src/main/index.ts b/src/main/index.ts index 452bb13b..3dbf2210 100644 --- a/src/main/index.ts +++ b/src/main/index.ts @@ -676,7 +676,9 @@ async function initializeServices(): Promise { // (the launch dialog enables it right before launching an orchestrator). await applyAgentControlSetting(next.agentControlEnabled); agentBrowserBridge?.setEnabled(next.browserAgentAccess); - browserService?.setRestoreTabs(next.browserRestoreTabs); + browserService?.setRestoreTabs(next.browserRestoreTabs).catch((error: unknown) => { + console.warn("CanvasTTY browser tab restore setting could not be applied.", error); + }); browserService?.cancelCanvasNavigationGesture(); browserService?.setCanvasWheelCaptureMode(next.canvasWheelCaptureMode); canvasNavigationInput?.setBindings({ diff --git a/src/main/ipc/registerIpc.ts b/src/main/ipc/registerIpc.ts index fb69ecc4..c79d4848 100644 --- a/src/main/ipc/registerIpc.ts +++ b/src/main/ipc/registerIpc.ts @@ -1,5 +1,4 @@ -import { extname } from "node:path"; -import { readFile, stat } from "node:fs/promises"; +import { realpath } from "node:fs/promises"; import { app, BrowserWindow, clipboard, dialog, ipcMain, shell } from "electron"; import type { IpcMainEvent, IpcMainInvokeEvent, OpenDialogOptions } from "electron"; import type { @@ -34,15 +33,7 @@ import { PluginBrowserOpenBroker } from "./PluginBrowserOpenBroker"; import type { GithubAuthService } from "../services/GithubAuthService"; import type { HermesHudService } from "../services/HermesHudService"; import { normalizeExternalUrl } from "../../shared/externalUrl"; - -const MAX_MEDIA_BYTES = 25 * 1024 * 1024; -const MEDIA_MIME: Record = { - ".png": "image/png", - ".jpg": "image/jpeg", - ".jpeg": "image/jpeg", - ".webp": "image/webp", - ".gif": "image/gif" -}; +import { readHomeMedia } from "../services/homeMedia"; interface Dependencies { settings: SettingsStore; @@ -141,7 +132,8 @@ export function registerIpc({ assertMainRenderer(event, getMainWindow); return recheckProviderClis(); }); - ipcMain.handle(IPC.settingsUpdate, async (_event, patch: Partial) => { + ipcMain.handle(IPC.settingsUpdate, async (event, patch: Partial) => { + assertMainRenderer(event, getMainWindow); const next = await settings.update(patch); await applyBrowserSettings(next); return next; @@ -190,15 +182,19 @@ export function registerIpc({ const result = owner ? await dialog.showOpenDialog(owner, options) : await dialog.showOpenDialog(options); - const path = result.filePaths[0]; - if (result.canceled || !path) return null; - return { path, dataUrl: await readMedia(path) }; + const picked = result.filePaths[0]; + if (result.canceled || !picked) return null; + // Save the file the person picked, not a link to it, so a later read is + // not redirected by changing the link. + const path = await realpath(picked); + return { path, dataUrl: await readHomeMedia(path) }; }); - ipcMain.handle(IPC.mediaRead, async (_event, path: string) => { + ipcMain.handle(IPC.mediaRead, async (event, path: string) => { + assertMainRenderer(event, getMainWindow); if (typeof path !== "string" || settings.get().mediaPath !== path) return null; try { - return await readMedia(path); + return await readHomeMedia(path); } catch (error) { console.warn("CanvasTTY media could not be read.", error); return null; @@ -246,25 +242,29 @@ export function registerIpc({ closePluginWindows(pluginId); return plugins.updatePlugin(pluginId); }); - ipcMain.handle(IPC.pluginsPreviewInstall, (_event, sourceUrl: string) => { + ipcMain.handle(IPC.pluginsPreviewInstall, (event, sourceUrl: string) => { + assertMainRenderer(event, getMainWindow); if (typeof sourceUrl !== "string") throw new Error("GitHub URL is required."); return plugins.previewInstall(sourceUrl); }); - ipcMain.handle(IPC.pluginsInstall, (_event, token: string, selectedModules?: string[]) => { + ipcMain.handle(IPC.pluginsInstall, (event, token: string, selectedModules?: string[]) => { + assertMainRenderer(event, getMainWindow); if (typeof token !== "string") throw new Error("Plugin preview token is invalid."); if (selectedModules !== undefined && ( !Array.isArray(selectedModules) || selectedModules.some((item) => typeof item !== "string") )) throw new Error("Plugin module selection is invalid."); return plugins.install(token, selectedModules); }); - ipcMain.handle(IPC.pluginsSetModules, async (_event, pluginId: string, selectedModules: string[]) => { + ipcMain.handle(IPC.pluginsSetModules, async (event, pluginId: string, selectedModules: string[]) => { + assertMainRenderer(event, getMainWindow); if (!Array.isArray(selectedModules) || selectedModules.some((item) => typeof item !== "string")) { throw new Error("Plugin module selection is invalid."); } closePluginWindows(pluginId); return plugins.setModules(pluginId, selectedModules); }); - ipcMain.handle(IPC.pluginsSetEnabled, async (_event, pluginId: string, enabled: boolean) => { + ipcMain.handle(IPC.pluginsSetEnabled, async (event, pluginId: string, enabled: boolean) => { + assertMainRenderer(event, getMainWindow); if (typeof enabled !== "boolean") throw new Error("Plugin enabled state is invalid."); try { return await plugins.setEnabled(pluginId, enabled); @@ -327,7 +327,8 @@ export function registerIpc({ if (typeof pluginId !== "string" || typeof provider !== "string") throw new Error("Launch option request is invalid."); return launchFieldOptions(pluginId, provider as ProviderId); }); - ipcMain.handle(IPC.pluginsUninstall, async (_event, pluginId: string) => { + ipcMain.handle(IPC.pluginsUninstall, async (event, pluginId: string) => { + assertMainRenderer(event, getMainWindow); closePluginWindows(pluginId); await pluginSecrets.revokeAll(pluginId); await pluginMedia.revokeAll(pluginId); @@ -353,7 +354,8 @@ export function registerIpc({ ipcMain.handle(IPC.pluginsOpenWindow, (_event, pluginId: string, contributionId: string) => ( openPluginWindow(pluginId, contributionId) )); - ipcMain.handle(IPC.pluginsOpenExternal, async (_event, pluginId: string, value: string) => { + ipcMain.handle(IPC.pluginsOpenExternal, async (event, pluginId: string, value: string) => { + assertMainRenderer(event, getMainWindow); plugins.assertPermission(pluginId, "external:open"); const url = normalizeExternalUrl(value); await shell.openExternal(url); @@ -369,22 +371,30 @@ export function registerIpc({ await plugins.storageSet(pluginId, key, value); broadcastPluginStorageChange(pluginId, key, value); }); - ipcMain.handle(IPC.pluginsSecretsGet, (_event, pluginId: string, key: string) => ( - pluginSecrets.get(pluginId, key) - )); - ipcMain.handle(IPC.pluginsSecretsSet, (_event, pluginId: string, key: string, value: string) => ( - pluginSecrets.set(pluginId, key, value) - )); - ipcMain.handle(IPC.pluginsSecretsDelete, (_event, pluginId: string, key: string) => ( - pluginSecrets.delete(pluginId, key) - )); - ipcMain.handle(IPC.providerSecretsStatus, () => providerSecrets.status()); - ipcMain.handle(IPC.providerSecretsSet, (_event, secretId: string, value: string) => ( - providerSecrets.set(providerSecretValue(secretId), value) - )); - ipcMain.handle(IPC.providerSecretsClear, (_event, secretId: string) => ( - providerSecrets.delete(providerSecretValue(secretId)) - )); + ipcMain.handle(IPC.pluginsSecretsGet, (event, pluginId: string, key: string) => { + assertMainRenderer(event, getMainWindow); + return pluginSecrets.get(pluginId, key); + }); + ipcMain.handle(IPC.pluginsSecretsSet, (event, pluginId: string, key: string, value: string) => { + assertMainRenderer(event, getMainWindow); + return pluginSecrets.set(pluginId, key, value); + }); + ipcMain.handle(IPC.pluginsSecretsDelete, (event, pluginId: string, key: string) => { + assertMainRenderer(event, getMainWindow); + return pluginSecrets.delete(pluginId, key); + }); + ipcMain.handle(IPC.providerSecretsStatus, (event) => { + assertMainRenderer(event, getMainWindow); + return providerSecrets.status(); + }); + ipcMain.handle(IPC.providerSecretsSet, (event, secretId: string, value: string) => { + assertMainRenderer(event, getMainWindow); + return providerSecrets.set(providerSecretValue(secretId), value); + }); + ipcMain.handle(IPC.providerSecretsClear, (event, secretId: string) => { + assertMainRenderer(event, getMainWindow); + return providerSecrets.delete(providerSecretValue(secretId)); + }); ipcMain.handle(IPC.pluginsMediaPickLibrary, (event, pluginId: string) => ( pickPluginMediaLibrary(event, pluginId, plugins, pluginMedia) )); @@ -677,21 +687,30 @@ export function registerIpc({ if (typeof id !== "string") throw new Error("Terminal session ID is required."); return terminals.readBuffer(id); }); - ipcMain.handle(IPC.terminalCreate, (_event, request: CreateSessionRequest) => terminals.create(request)); - ipcMain.handle(IPC.terminalRestart, (_event, id: string, options?: { resume?: unknown }) => ( - terminals.restart(id, { resume: options?.resume === true }) - )); - ipcMain.on(IPC.terminalInput, (_event, id: string, data: string) => terminals.input(id, data)); + ipcMain.handle(IPC.terminalCreate, (event, request: CreateSessionRequest) => { + assertMainRenderer(event, getMainWindow); + return terminals.create(request); + }); + ipcMain.handle(IPC.terminalRestart, (event, id: string, options?: { resume?: unknown }) => { + assertMainRenderer(event, getMainWindow); + return terminals.restart(id, { resume: options?.resume === true }); + }); + ipcMain.on(IPC.terminalInput, (event, id: string, data: string) => { + // Fire-and-forget: a foreign sender is dropped instead of throwing into the IPC layer. + if (!isMainRenderer(event, getMainWindow)) return; + terminals.input(id, data); + }); ipcMain.on(IPC.terminalResize, (_event, id: string, cols: number, rows: number) => { terminals.resize(id, cols, rows); }); ipcMain.on(IPC.terminalBounds, (_event, id: string, bounds: SessionBounds) => terminals.setBounds(id, bounds)); ipcMain.handle(IPC.terminalRename, (_event, id: string, title: string) => terminals.rename(id, title)); ipcMain.handle(IPC.terminalSetRestore, (_event, id: string, restore: boolean) => terminals.setRestore(id, restore)); - ipcMain.handle(IPC.terminalDispose, (_event, id: string, options?: { keepEnvironmentData?: unknown }) => ( + ipcMain.handle(IPC.terminalDispose, (event, id: string, options?: { keepEnvironmentData?: unknown }) => { + assertMainRenderer(event, getMainWindow); // Environment data is kept unless the person explicitly chose Remove. - terminals.dispose(id, { keepEnvironmentData: options?.keepEnvironmentData !== false }) - )); + return terminals.dispose(id, { keepEnvironmentData: options?.keepEnvironmentData !== false }); + }); // Fire-and-forget, like the other stream-reporting channels: a malformed // report is ignored rather than rejecting into the renderer. ipcMain.on(IPC.terminalSetVisible, (_event, id: unknown, visible: unknown) => { @@ -740,6 +759,18 @@ function isCanvasNavigationPointerBindingInput( && typeof input.shiftKey === "boolean"; } +function isMainRenderer( + event: IpcMainEvent | IpcMainInvokeEvent, + getMainWindow: () => BrowserWindow | null +): boolean { + try { + assertMainRenderer(event, getMainWindow); + return true; + } catch { + return false; + } +} + function assertMainRenderer( event: IpcMainEvent | IpcMainInvokeEvent, getMainWindow: () => BrowserWindow | null @@ -844,19 +875,6 @@ function providerValue(value: unknown): ProviderId { throw new Error("Plugin requested an unknown launcher provider."); } -async function readMedia(path: string): Promise { - const mime = MEDIA_MIME[extname(path).toLowerCase()]; - if (!mime) throw new Error("Unsupported media type."); - - const metadata = await stat(path); - if (!metadata.isFile() || metadata.size > MAX_MEDIA_BYTES) { - throw new Error("Media must be a file smaller than 25 MB."); - } - - const content = await readFile(path); - return `data:${mime};base64,${content.toString("base64")}`; -} - function providerSecretValue(value: string): ProviderSecretId { if ((PROVIDER_SECRET_IDS as readonly string[]).includes(value)) return value as ProviderSecretId; throw new Error("Provider secret id is unknown."); diff --git a/src/main/services/BrowserService.ts b/src/main/services/BrowserService.ts index a8b57763..7e95a560 100644 --- a/src/main/services/BrowserService.ts +++ b/src/main/services/BrowserService.ts @@ -317,8 +317,7 @@ export class BrowserService { this.restoreTabsEnabled = enabled; if (enabled) await this.persistRuntime(); else { - await this.store.clear(); - this.persisted = this.store.get(); + await this.clearSavedTabs(); } } @@ -431,8 +430,7 @@ export class BrowserService { ]); this.downloads = []; this.pendingDialogs.clear(); - await this.store.clear(); - this.persisted = this.store.get(); + await this.clearSavedTabs(); if (this.visible) return this.newTab(); this.emit(); return this.getState(); @@ -469,10 +467,7 @@ export class BrowserService { private async initialize(): Promise { this.persisted = await this.store.load(); - if (!this.restoreTabsEnabled) { - await this.store.clear(); - this.persisted = this.store.get(); - } + if (!this.restoreTabsEnabled) await this.clearSavedTabs(); this.activeTabId = this.persisted.activeTabId; await mkdir(this.policy.downloadRoot, { recursive: true }); this.configureSession(); @@ -1000,7 +995,23 @@ export class BrowserService { const tabs = [...this.tabs.values()] .map((tab) => ({ id: tab.id, url: this.tabUrl(tab) })) .filter((tab) => isSafeBrowserUrl(tab.url)); - this.persisted = await this.store.replace(tabs, this.activeTabId); + try { + this.persisted = await this.store.replace(tabs, this.activeTabId); + } catch (error) { + // The tabs on screen stay as they are; only the copy restored at the next + // start is stale. The store already holds the new state in memory. + this.persisted = this.store.get(); + console.warn("CanvasTTY browser tabs could not be saved.", error); + } + } + + private async clearSavedTabs(): Promise { + try { + await this.store.clear(); + } catch (error) { + console.warn("CanvasTTY saved browser tabs could not be cleared.", error); + } + this.persisted = this.store.get(); } private destroyRuntimeTabs(): void { diff --git a/src/main/services/GithubAuthService.ts b/src/main/services/GithubAuthService.ts index 132b3865..defbdb6b 100644 --- a/src/main/services/GithubAuthService.ts +++ b/src/main/services/GithubAuthService.ts @@ -226,16 +226,36 @@ export class GithubAuthService { device_code: deviceCode, grant_type: "urn:ietf:params:oauth:grant-type:device_code" }); - const response = await this.request("https://github.com/login/oauth/access_token", { - method: "POST", - headers: oauthHeaders(), - body: body.toString() - }, signal); - if (!response.ok) continue; - const payload: unknown = await response.json(); - if (!isRecord(payload)) continue; + // A dropped connection, the 15 s request timeout or a GitHub 5xx is not + // the end of the flow: the person may still approve the code. Back off + // (RFC 8628 section 3.5) and poll again until the code expires. + let payload: unknown; + try { + const response = await this.request("https://github.com/login/oauth/access_token", { + method: "POST", + headers: oauthHeaders(), + body: body.toString() + }, signal); + if (!response.ok) { + interval = backedOff(interval); + continue; + } + payload = await response.json(); + } catch (error) { + if (signal.aborted) throw error; + interval = backedOff(interval); + continue; + } + if (!isRecord(payload)) { + interval = backedOff(interval); + continue; + } if (payload.error === "authorization_pending" || payload.error === "slow_down") { - if (payload.error === "slow_down") interval += 5; + if (payload.error === "slow_down") { + // +5 s for this and later polls; GitHub may name a longer interval. + const requested = typeof payload.interval === "number" && Number.isFinite(payload.interval) ? payload.interval : 0; + interval = Math.max(interval + 5, Math.min(requested, MAX_POLL_INTERVAL_SECONDS)); + } continue; } if (payload.error === "access_denied" || payload.error === "expired_token") return; @@ -411,3 +431,9 @@ function isRecord(value: unknown): value is Record { function isMissingFile(error: unknown): boolean { return error instanceof Error && "code" in error && (error as { code?: string }).code === "ENOENT"; } + +const MAX_POLL_INTERVAL_SECONDS = 60; + +function backedOff(interval: number): number { + return Math.min(Math.max(interval * 2, interval + 1), MAX_POLL_INTERVAL_SECONDS); +} diff --git a/src/main/services/LimitsService.ts b/src/main/services/LimitsService.ts index f26ab4bf..62f4cc3e 100644 --- a/src/main/services/LimitsService.ts +++ b/src/main/services/LimitsService.ts @@ -572,6 +572,9 @@ class KimiWebUsageClient { private async startChild(): Promise { if (!this.cli) throw new LimitsAdapterError("cli-not-found"); const port = await reserveLoopbackPort(); + // dispose() during the await found no child to stop; starting one now would + // leave `kimi web` (a local server with a token in its URL) running. + if (this.disposed) throw new LimitsAdapterError("protocol-error"); const launch = providerChildProcessLaunch( this.cli, ["web", "--no-open", "--port", String(port), "--log-level", "silent"] diff --git a/src/main/services/PluginCards.ts b/src/main/services/PluginCards.ts index 16e4c0c5..7e8f8059 100644 --- a/src/main/services/PluginCards.ts +++ b/src/main/services/PluginCards.ts @@ -79,6 +79,12 @@ export class PluginCards { throw new Error(`badge.tooltip must be text of at most ${MAX_BADGE_TOOLTIP} characters.`); } const perCard = this.badges.get(sessionId) ?? new Map(); + // Badges of plugins that are no longer trusted are hidden; they must not + // keep a trusted plugin out of the card's slots. + const trusted = this.deps.trustedPlugins(); + for (const owner of [...perCard.keys()]) { + if (owner !== pluginId && !trusted.has(owner)) perCard.delete(owner); + } if (!perCard.has(pluginId) && perCard.size >= MAX_BADGES_PER_CARD) throw new Error("This card already shows the most plugin badges."); perCard.set(pluginId, { pluginId, diff --git a/src/main/services/PluginManager.ts b/src/main/services/PluginManager.ts index 65d27636..5a82d785 100644 --- a/src/main/services/PluginManager.ts +++ b/src/main/services/PluginManager.ts @@ -729,6 +729,11 @@ export class PluginManager { } await rm(join(this.pluginRoot, plugin.manifest.id), { recursive: true, force: true }); await rm(join(this.storageRoot, `${plugin.manifest.id}.json`), { force: true }); + // Copies of unreadable storage kept aside by storageSet go with the plugin. + const kept = `${plugin.manifest.id}.json.unreadable-`; + for (const name of await readdir(this.storageRoot).catch(() => [] as string[])) { + if (name.startsWith(kept)) await rm(join(this.storageRoot, name), { force: true }); + } await rm(join(this.dataRoot, plugin.manifest.id), { recursive: true, force: true }); } @@ -1083,7 +1088,7 @@ export class PluginManager { assertStorageKey(key); const previous = this.storageWrites.get(pluginId) ?? Promise.resolve(); const next = previous.catch(() => undefined).then(async () => { - const storage = await this.readStorage(pluginId); + const storage = await this.readStorageForWrite(pluginId); storage[key] = jsonClone(value); const snapshot = JSON.stringify(storage, null, 2); if (Buffer.byteLength(snapshot) > MAX_STORAGE_BYTES) { @@ -1179,6 +1184,34 @@ export class PluginManager { } } + /** + * The storage a write starts from. Reading {} on any failure made the next + * write replace every other key with just the new one. A read error + * (permissions, a locked file) now refuses the write; a file that is not + * valid storage is kept aside under a new name before a fresh one starts. + */ + private async readStorageForWrite(pluginId: string): Promise> { + const path = join(this.storageRoot, `${pluginId}.json`); + let raw: string; + try { + raw = await readFile(path, "utf8"); + } catch (error) { + if (isMissingFile(error)) return {}; + throw new Error("Plugin storage could not be read; nothing was written.", { cause: error }); + } + let parsed: unknown = null; + try { + parsed = Buffer.byteLength(raw) > MAX_STORAGE_BYTES ? null : JSON.parse(raw); + } catch { + parsed = null; + } + if (isRecord(parsed)) return { ...parsed }; + const kept = `${path}.unreadable-${Date.now()}`; + await rename(path, kept); + console.warn(`CanvasTTY plugin storage for ${pluginId} was not valid and was kept as ${kept}.`); + return {}; + } + private cleanupExpiredPreviews(): void { const now = Date.now(); for (const [token, pending] of this.pending) { diff --git a/src/main/services/PluginMediaService.ts b/src/main/services/PluginMediaService.ts index 17348682..58a43ce3 100644 --- a/src/main/services/PluginMediaService.ts +++ b/src/main/services/PluginMediaService.ts @@ -7,6 +7,7 @@ import { realpath, rename, stat, + unlink, writeFile } from "node:fs/promises"; import { basename, dirname, extname, join, relative, resolve, sep } from "node:path"; @@ -182,9 +183,17 @@ export class PluginMediaService { } const path = join(canonicalDirectory, fileName); - const temporaryPath = `${path}.tmp`; - await writeFile(temporaryPath, content, "utf8"); - await rename(temporaryPath, path); + // A fixed `.tmp` could already be a link pointing outside the + // library, and writeFile follows it. A new random name created with O_EXCL + // ("wx") fails on any existing entry, link or not. + const temporaryPath = `${path}.${randomUUID()}.tmp`; + try { + await writeFile(temporaryPath, content, { encoding: "utf8", flag: "wx" }); + await rename(temporaryPath, path); + } catch (error) { + await unlink(temporaryPath).catch(() => undefined); + throw error; + } const metadata = await stat(path); return publicPlaylist({ relativePath: `Playlists/${fileName}`, size: metadata.size }); } diff --git a/src/main/services/PluginServiceSupervisor.ts b/src/main/services/PluginServiceSupervisor.ts index 06cc5051..228563d8 100644 --- a/src/main/services/PluginServiceSupervisor.ts +++ b/src/main/services/PluginServiceSupervisor.ts @@ -1,6 +1,7 @@ import { spawn, type ChildProcess } from "node:child_process"; import { createHash } from "node:crypto"; -import { mkdir, readFile } from "node:fs/promises"; +import { mkdir, readFile, realpath } from "node:fs/promises"; +import { pathToFileURL } from "node:url"; import type { PluginPermission, PluginServiceLogEntry, @@ -95,6 +96,55 @@ const INHERITED_ENVIRONMENT = new Set([ "USERPROFILE", "APPDATA", "LOCALAPPDATA", "ProgramData", "HOMEDRIVE", "HOMEPATH" ]); +/** + * Module hooks for the service process: the entry is loaded from bytes the + * hook read and hashed itself, and a mismatch stops the load. They run + * before the entry through `--import`, off the main thread (module.register). + * + * The entry checked is the main module node actually resolved (the one + * resolve without a parent), not only the URL the host computed: when the + * entry or a folder above it is replaced by a symlink after the host's check, + * node resolves the main module to another file, and that file must match the + * hash too. The host's URL stays checked as well. + * + * The hooks take their modules with `await import(...)`, never a static + * `import ... from`: electron-vite puts its CommonJS shim (`__dirname`, + * `require`) after the last static import it finds in the main bundle, and a + * static import inside this string would pull the shim into the string, which + * leaves the whole main process without `__dirname`. + */ +const ENTRY_GUARD_HOOKS = ` +const { createHash } = await import("node:crypto"); +const { readFile } = await import("node:fs/promises"); +let entryUrl = null; +let mainUrl = null; +let expected = null; +export function initialize(data) { entryUrl = data.url; expected = data.sha256; } +export async function resolve(specifier, context, nextResolve) { + const resolved = await nextResolve(specifier, context); + if (mainUrl === null && context.parentURL === undefined) mainUrl = resolved.url; + return resolved; +} +export async function load(url, context, nextLoad) { + if (url !== entryUrl && url !== mainUrl) return nextLoad(url, context); + const source = await readFile(new URL(url)); + if (createHash("sha256").update(source).digest("hex") !== expected) { + throw new Error("The service entry changed after it was trusted."); + } + const loaded = await nextLoad(url, context); + return { format: loaded.format, source, shortCircuit: true }; +} +`; + +export function entryGuardArguments(entryUrl: string, sha256: string): string[] { + const boot = [ + 'import { register } from "node:module";', + `register(${JSON.stringify(`data:text/javascript,${encodeURIComponent(ENTRY_GUARD_HOOKS)}`)},` + + ` { data: ${JSON.stringify({ url: entryUrl, sha256 })} });` + ].join("\n"); + return ["--import", `data:text/javascript,${encodeURIComponent(boot)}`]; +} + export function pluginServiceEnvironment(source: NodeJS.ProcessEnv): Record { const environment: Record = {}; for (const [name, value] of Object.entries(source)) { @@ -300,7 +350,10 @@ export class PluginServiceSupervisor { await this.hostGate; if (record.removed || this.disposed) return; record.state = "starting"; + let entryUrl: string; try { + // Node loads the main entry by its real path; the guard matches that URL. + entryUrl = pathToFileURL(await realpath(spec.entryPath)).href; const content = await readFile(spec.entryPath); if (createHash("sha256").update(content).digest("hex") !== spec.sha256) { // The file changed after the user trusted it: never run it, and do not retry. @@ -317,7 +370,11 @@ export class PluginServiceSupervisor { return; } - const child = spawn(this.options.command, [spec.entryPath], { + // The check above and node's own read of the entry are separate reads: a file + // replaced in between would run as trusted. The guard makes node run only + // bytes it read and hashed itself, so what runs is what matched the hash, + // wherever node resolves `spec.entryPath` by then. + const child = spawn(this.options.command, [...entryGuardArguments(entryUrl, spec.sha256), spec.entryPath], { cwd: spec.root, env: pluginServiceEnvironment(this.options.environment), stdio: ["pipe", "pipe", "pipe"], diff --git a/src/main/services/SettingsStore.ts b/src/main/services/SettingsStore.ts index dd766bc1..59e287b8 100644 --- a/src/main/services/SettingsStore.ts +++ b/src/main/services/SettingsStore.ts @@ -1,6 +1,7 @@ import { homedir } from "node:os"; import { dirname, join } from "node:path"; import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; +import { isHomeMediaPath } from "./homeMedia.ts"; import type { AgentProviderId, AgentCliAvailability, @@ -234,11 +235,17 @@ export class SettingsStore { async setAvailableProviders(availability: AgentCliAvailability): Promise { this.availableProviders = new Set([...AGENT_PROVIDERS].filter((provider) => availability[provider])); - const filtered = filterUnavailableProviders(this.value, this.availableProviders); - if (providerSelectionsChanged(this.value, filtered)) { + // Filter in queue order: a snapshot taken while an update() is still + // writing lacks that update, and persisting it afterwards dropped the + // update from the file (it stayed only in memory). + const write = this.writeQueue.catch(() => undefined).then(async () => { + const filtered = filterUnavailableProviders(this.value, this.availableProviders); + if (!providerSelectionsChanged(this.value, filtered)) return; + await this.persist(filtered); this.value = filtered; - await this.queuePersist(); - } + }); + this.writeQueue = write; + await write; return this.get(); } @@ -436,9 +443,11 @@ export function normalizeSettings( } const source = candidate as Partial & { zoomOverApplications?: unknown }; - const mediaPath = source.mediaPath === null || typeof source.mediaPath === "string" + // Only an absolute path to a supported image is kept; anything else keeps the + // previous choice. The main process reads this file for the Home screen. + const mediaPath = source.mediaPath === null || isHomeMediaPath(source.mediaPath) ? source.mediaPath - : fallback.mediaPath; + : isHomeMediaPath(fallback.mediaPath) ? fallback.mediaPath : null; const acknowledged = Array.isArray(source.acknowledgedDangerousProfiles) ? source.acknowledgedDangerousProfiles.filter( (provider): provider is AgentProviderId => AGENT_PROVIDERS.has(provider as AgentProviderId) diff --git a/src/main/services/TerminalManager.ts b/src/main/services/TerminalManager.ts index d36f8256..0d04c293 100644 --- a/src/main/services/TerminalManager.ts +++ b/src/main/services/TerminalManager.ts @@ -1457,7 +1457,7 @@ export class TerminalManager { launch: { command: planned.command, args: planned.args, env: visible, cwd: planned.cwd }, secretEnvNames, takenEnv: new Set(Object.keys(planned.launchEnvironment)), - path: planned.env.PATH + path: launchSearchPath(planned.env) }); if (!live()) { abandon(); @@ -1624,6 +1624,21 @@ function applyLaunchFailure(metadata: SessionMetadata, failure: UnavailableProvi metadata.failureDetails = failure.diagnostic; } +/** + * The launch's program search path. The environment is a plain copy of + * process.env, which on Windows is case-insensitive but keeps the spelling it + * was given ("Path"), so env.PATH alone finds nothing there. + */ +export function launchSearchPath( + environment: Readonly>, + platform: NodeJS.Platform = process.platform +): string | undefined { + if (platform !== "win32") return environment.PATH; + if (environment.PATH !== undefined) return environment.PATH; + const key = Object.keys(environment).find((name) => name.toUpperCase() === "PATH"); + return key === undefined ? undefined : environment[key]; +} + export function terminalEnvironment( source: Readonly> = process.env ): Record { diff --git a/src/main/services/TerminalSessionStore.ts b/src/main/services/TerminalSessionStore.ts index b8545912..fa72e137 100644 --- a/src/main/services/TerminalSessionStore.ts +++ b/src/main/services/TerminalSessionStore.ts @@ -1,5 +1,5 @@ import { dirname, join } from "node:path"; -import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; +import { mkdir, readFile, rename, unlink, writeFile } from "node:fs/promises"; import type { LaunchProfileId, SessionRole, @@ -131,9 +131,15 @@ export class TerminalSessionStore { const snapshot = `${JSON.stringify(this.value, null, 2)}\n`; const temporaryPath = `${this.filePath}.${process.pid}.tmp`; this.writeQueue = this.writeQueue.catch(() => undefined).then(async () => { - await mkdir(dirname(this.filePath), { recursive: true }); - await writeFile(temporaryPath, snapshot, { encoding: "utf8", mode: 0o600 }); - await rename(temporaryPath, this.filePath); + await mkdir(dirname(this.filePath), { recursive: true, mode: 0o700 }); + try { + await writeFile(temporaryPath, snapshot, { encoding: "utf8", mode: 0o600 }); + await rename(temporaryPath, this.filePath); + } catch (error) { + // A failed rename (a locked file on Windows) must not leave the temp file behind. + await unlink(temporaryPath).catch(() => undefined); + throw error; + } }); return this.writeQueue; } diff --git a/src/main/services/agent-browser/OrchestrationGateway.ts b/src/main/services/agent-browser/OrchestrationGateway.ts index f66d770f..f2842a5b 100644 --- a/src/main/services/agent-browser/OrchestrationGateway.ts +++ b/src/main/services/agent-browser/OrchestrationGateway.ts @@ -330,13 +330,23 @@ export class OrchestrationGateway { const controller = new AbortController(); connection.controllers.set(id, controller); connection.inflight += 1; + // Cancel answers at once; a handler that cannot stop (a plugin call) is + // no longer waited for, and its late result is dropped. + const canceled = new Promise((_resolve, reject) => { + controller.signal.addEventListener("abort", () => reject(new Error("canceled")), { once: true }); + }); + canceled.catch(() => undefined); try { - const value = await this.handler.execute(connection.lease!.terminalSessionId, { - id, - tool: tool as never, - arguments: args - }); + const value = await Promise.race([ + this.handler.execute(connection.lease!.terminalSessionId, { + id, + tool: tool as never, + arguments: args + }, controller.signal), + canceled + ]); if (connection.closed) return; + if (controller.signal.aborted) throw new Error("canceled"); this.send(connection, { v: ORCHESTRATION_BRIDGE_PROTOCOL_VERSION, type: "response", id, result: value }); } catch (error) { if (connection.closed) return; diff --git a/src/main/services/agent-browser/OrchestrationTools.ts b/src/main/services/agent-browser/OrchestrationTools.ts index 30bfff81..23bae0cb 100644 --- a/src/main/services/agent-browser/OrchestrationTools.ts +++ b/src/main/services/agent-browser/OrchestrationTools.ts @@ -30,8 +30,9 @@ export class ScopedOrchestrationHandler implements OrchestrationCommandHandler { ]; } - async execute(sessionId: string, request: OrchestrationRequest): Promise> { + async execute(sessionId: string, request: OrchestrationRequest, signal?: AbortSignal): Promise> { try { + if (signal?.aborted) throw canceledError(); const session = this.control.status(sessionId); if (isPluginOrchestrationTool(request.tool)) return await this.plugin(sessionId, session, request); // Plugin tools may reach other roles' sessions through the same bridge; the core tools never do. @@ -40,7 +41,7 @@ export class ScopedOrchestrationHandler implements OrchestrationCommandHandler { } switch (request.tool) { case "spawn_agent": - return await this.spawn(sessionId, request.arguments); + return await this.spawn(sessionId, request.arguments, signal); case "send_to_agent": return await this.send(sessionId, request.arguments); case "observe_agent": @@ -82,7 +83,7 @@ export class ScopedOrchestrationHandler implements OrchestrationCommandHandler { } } - private async spawn(orchestratorId: string, args: Record): Promise> { + private async spawn(orchestratorId: string, args: Record, signal?: AbortSignal): Promise> { const created = await this.control.spawn({ parentSessionId: orchestratorId, provider: args.provider as never, @@ -91,6 +92,15 @@ export class ScopedOrchestrationHandler implements OrchestrationCommandHandler { ...(args.prompt !== undefined ? { initialPrompt: args.prompt as string } : {}), ...(args.launchOptions !== undefined ? { launchOptions: args.launchOptions as SpawnAgentRequest["launchOptions"] } : {}) }); + if (signal?.aborted) { + // Canceled while the agent was starting: nobody will receive its id, so close it. + try { + this.control.cancel(created.id); + } catch { + // It already ended. + } + throw canceledError(); + } return { sessionId: created.id, provider: created.provider, @@ -156,3 +166,7 @@ export class ScopedOrchestrationHandler implements OrchestrationCommandHandler { } } } + +function canceledError(): Error { + return orchestrationBridgeError("CANCELED", "Orchestration command was canceled.", true); +} diff --git a/src/main/services/agent-browser/WindowsPipeHostTransport.ts b/src/main/services/agent-browser/WindowsPipeHostTransport.ts index 9a3c6de5..0da5f9ac 100644 --- a/src/main/services/agent-browser/WindowsPipeHostTransport.ts +++ b/src/main/services/agent-browser/WindowsPipeHostTransport.ts @@ -120,6 +120,16 @@ export class WindowsPipeHostTransport extends EventEmitter { } ) as ChildProcessWithoutNullStreams; this.child = child; + // A write or end after the host died raises EPIPE on stdin; without a + // listener that is an uncaught exception in the main process. + const pipeFailure = (name: string) => (error: Error) => { + if (this.child !== child) return; + this.fail(new Error(`Windows agent pipe host ${name} failed: ${error.message}`)); + }; + child.stdin.on("error", pipeFailure("input")); + child.stdout.on("error", pipeFailure("output")); + // stderr is diagnostics only. + child.stderr.on("error", () => undefined); return await new Promise((resolve, reject) => { let settled = false; diff --git a/src/main/services/agent-browser/orchestration-protocol.ts b/src/main/services/agent-browser/orchestration-protocol.ts index de5f3fe1..cc79a45e 100644 --- a/src/main/services/agent-browser/orchestration-protocol.ts +++ b/src/main/services/agent-browser/orchestration-protocol.ts @@ -45,7 +45,8 @@ export type OrchestrationResult = /** The only implementation the gateway accepts; AgentControlService is * wrapped by a scoping adapter, never called directly by the protocol. */ export interface OrchestrationCommandHandler { - execute(sessionId: string, request: OrchestrationRequest): Promise>; + /** `signal` aborts when the orchestrator cancels the request or disconnects. */ + execute(sessionId: string, request: OrchestrationRequest, signal?: AbortSignal): Promise>; /** The tools this session sees (core tools for orchestrators, plugin tools by role). */ listTools?(sessionId: string): McpToolDefinition[]; } diff --git a/src/main/services/agent-control/AgentControlGateway.ts b/src/main/services/agent-control/AgentControlGateway.ts index 8c343113..f73ac526 100644 --- a/src/main/services/agent-control/AgentControlGateway.ts +++ b/src/main/services/agent-control/AgentControlGateway.ts @@ -1,5 +1,5 @@ import { createHash, randomBytes, timingSafeEqual } from "node:crypto"; -import { chmod, mkdir, mkdtemp, realpath, writeFile } from "node:fs/promises"; +import { chmod, mkdir, mkdtemp, realpath, rename, rm, writeFile } from "node:fs/promises"; import { createServer, type Server } from "node:net"; import { tmpdir } from "node:os"; import { isAbsolute, join } from "node:path"; @@ -14,7 +14,12 @@ import { hasAutoMode, isLaunchProfile } from "../../../shared/autoMode.ts"; const MAX_REQUEST_BYTES = 128 * 1024; const MAX_RESPONSE_BYTES = 256 * 1024; const MAX_RECEIPTS = 4096; +// Refusals that happen before anything is written: a retry with the same +// request id must be performed again instead of replaying the refusal. +const RETRYABLE_REFUSALS = new Set(["BUSY", "NOT_READY", "LIMIT_REACHED", "LIFECYCLE_DISABLED", "CLOSED"]); const MAX_SESSIONS = 32; +const MAX_TRANSPORT_RESTART_ATTEMPTS = 3; +const TRANSPORT_RESTART_BASE_DELAY_MS = 500; const MAX_TEXT = 16_000; const ID = /^[a-zA-Z0-9][a-zA-Z0-9._:-]{0,127}$/; const SECRET = /^[a-f0-9]{64}$/; @@ -68,6 +73,10 @@ export interface AgentControlGatewayOptions { platform?: NodeJS.Platform; windowsHostPath?: string; windowsPipeHostFactory?: (options: { hostPath: string; platform: NodeJS.Platform; parentPid: number }) => WindowsPipeHostTransport; + /** Receipts kept for request-id replay (default 4096); the oldest settled ones are dropped first. */ + maxReceipts?: number; + /** Called after the Windows pipe host was restarted and connection.json names the new endpoint. */ + onTransportRestarted?(connectionPath: string): void; } export class ControlError extends Error { @@ -82,48 +91,125 @@ export class AgentControlGateway { private readonly instanceId = randomBytes(16).toString("hex"); private readonly sessions = new Map(); private readonly sockets = new Set(); - private readonly receipts = new Map }>(); + private readonly receipts = new Map; settled: boolean }>(); private readonly busy = new Set(); private server: Server | null = null; private windows: WindowsPipeHostTransport | null = null; + private socketDirectory: string | null = null; + private tokenFileWritten = false; + private starting = false; + private restartTimer: ReturnType | undefined; + private restartAttempts = 0; private closed = false; constructor(options: AgentControlGatewayOptions) { this.options = options; } async start(): Promise { - if (this.server || this.windows || this.closed) throw new Error("Agent control is already started or closed."); + if (this.starting || this.server || this.windows || this.closed) throw new Error("Agent control is already started or closed."); + this.starting = true; + try { + const endpoint = await this.openEndpoint(); + return await this.writeDiscovery(endpoint); + } catch (error) { + // Leave nothing listening and no dead transport behind, so a later start() can succeed. + await this.closeEndpoint(); + throw error; + } finally { + this.starting = false; + } + } + + private async openEndpoint(): Promise { const platform = this.options.platform ?? process.platform; - let endpoint: string; if (platform === "win32") { if (!this.options.windowsHostPath) throw new Error("Agent control requires the current-user Windows pipe host."); - this.windows = (this.options.windowsPipeHostFactory ?? ((options) => new WindowsPipeHostTransport(options)))({ + const transport = (this.options.windowsPipeHostFactory ?? ((options) => new WindowsPipeHostTransport(options)))({ hostPath: this.options.windowsHostPath, platform, parentPid: process.pid }); - endpoint = await this.windows.start((socket) => this.accept(socket)); - } else { - const directory = await mkdtemp(join(tmpdir(), "ctty-control-")); - await chmod(directory, 0o700); - endpoint = join(directory, "c.sock"); - if (Buffer.byteLength(endpoint) > 100) throw new Error("Agent control socket path is too long."); - this.server = createServer((socket) => this.accept(socket)); - await new Promise((resolve, reject) => { - this.server!.once("error", reject); - this.server!.listen(endpoint, () => { this.server!.off("error", reject); resolve(); }); - }); - await chmod(endpoint, 0o600); + this.windows = transport; + transport.on("fatal", () => this.handleTransportFatal(transport)); + const endpoint = await transport.start((socket) => this.accept(socket)); + if (this.windows !== transport || this.closed) { + await transport.close(); + throw new Error("Agent control is shutting down."); + } + return endpoint; } + const directory = await mkdtemp(join(tmpdir(), "ctty-control-")); + this.socketDirectory = directory; + await chmod(directory, 0o700); + const endpoint = join(directory, "c.sock"); + if (Buffer.byteLength(endpoint) > 100) throw new Error("Agent control socket path is too long."); + const server = createServer((socket) => this.accept(socket)); + this.server = server; + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(endpoint, () => { server.off("error", reject); resolve(); }); + }); + await chmod(endpoint, 0o600); + return endpoint; + } + + private async writeDiscovery(endpoint: string): Promise { const directory = join(this.options.userDataPath, "agent-control"); await mkdir(directory, { recursive: true, mode: 0o700 }); await chmod(directory, 0o700); const tokenFile = join(directory, `token-${this.instanceId}`); - await writeFile(tokenFile, this.token, { flag: "wx", mode: 0o600 }); + if (!this.tokenFileWritten) { + await writeFile(tokenFile, this.token, { flag: "wx", mode: 0o600 }); + this.tokenFileWritten = true; + } const connection = join(directory, "connection.json"); - await writeFile(connection, JSON.stringify({ v: 1, service: "canvastty-agent-control", instanceId: this.instanceId, - endpoint, tokenFile, pid: process.pid }, null, 2) + "\n", { mode: 0o600 }); - await chmod(connection, 0o600); + // Written to a temp file and renamed: a controller reading the record while + // a restarted host republishes it must never see it empty or half written. + const temporary = `${connection}.${randomBytes(8).toString("hex")}.tmp`; + try { + await writeFile(temporary, JSON.stringify({ v: 1, service: "canvastty-agent-control", instanceId: this.instanceId, + endpoint, tokenFile, pid: process.pid }, null, 2) + "\n", { mode: 0o600, flag: "wx" }); + await chmod(temporary, 0o600); + await rename(temporary, connection); + } catch (error) { + await rm(temporary, { force: true }).catch(() => undefined); + throw error; + } return connection; } + private async closeEndpoint(): Promise { + const server = this.server; + const transport = this.windows; + const directory = this.socketDirectory; + this.server = null; + this.windows = null; + this.socketDirectory = null; + if (transport) await transport.close().catch(() => undefined); + if (server?.listening) await new Promise((resolve) => server.close(() => resolve())); + if (directory) await rm(directory, { recursive: true, force: true }).catch(() => undefined); + } + + /** The Windows pipe host died: drop its connections and bring up a new one with a fresh discovery record. */ + private handleTransportFatal(transport: WindowsPipeHostTransport): void { + if (this.windows !== transport) return; + this.windows = null; + for (const socket of this.sockets) socket.destroy(); + this.scheduleTransportRestart(); + } + + private scheduleTransportRestart(): void { + if (this.closed || this.restartTimer || this.restartAttempts >= MAX_TRANSPORT_RESTART_ATTEMPTS) return; + const delay = TRANSPORT_RESTART_BASE_DELAY_MS * 2 ** this.restartAttempts; + this.restartAttempts += 1; + this.restartTimer = setTimeout(() => { + this.restartTimer = undefined; + if (this.closed || this.windows || this.starting) return; + this.start().then((connection) => { + this.restartAttempts = 0; + this.options.onTransportRestarted?.(connection); + }, () => this.scheduleTransportRestart()); + }, delay); + this.restartTimer.unref?.(); + } + observe(channel: string, payload: unknown): void { if (channel === IPC.terminalRemoved) { const id = (payload as { id: string }).id; @@ -168,11 +254,12 @@ export class AgentControlGateway { async close(): Promise { this.closed = true; + clearTimeout(this.restartTimer); + this.restartTimer = undefined; for (const socket of this.sockets) socket.destroy(); for (const owned of this.sessions.values()) { await owned.ready; owned.terminal.dispose(); } this.sessions.clear(); - if (this.windows) await this.windows.close(); - if (this.server?.listening) await new Promise((resolve) => this.server!.close(() => resolve())); + await this.closeEndpoint(); // Retain inert discovery/diagnostic records; a new app instance gets a new token and instanceId. } @@ -236,12 +323,31 @@ export class AgentControlGateway { if (previous.digest !== digest) throw new ControlError("REQUEST_CONFLICT", "Request ID was already used for different input."); return previous.result; } - if (this.receipts.size >= MAX_RECEIPTS) throw new ControlError("LIMIT_REACHED", "Control request capacity reached; existing receipts remain available."); + if (!this.makeReceiptRoom()) throw new ControlError("LIMIT_REACHED", "Too many control requests are still running."); const result = this.perform(owner, request); - this.receipts.set(key, { digest, result }); + const receipt = { digest, result, settled: false }; + this.receipts.set(key, receipt); + result.then(() => { receipt.settled = true; }, (error: unknown) => { + receipt.settled = true; + if (error instanceof ControlError && RETRYABLE_REFUSALS.has(error.code) && this.receipts.get(key) === receipt) { + this.receipts.delete(key); + } + }); return result; } + /** Drops the oldest finished receipts once the cap is reached; running ones are kept. */ + private makeReceiptRoom(): boolean { + const limit = this.options.maxReceipts ?? MAX_RECEIPTS; + if (this.receipts.size < limit) return true; + for (const [key, receipt] of this.receipts) { + if (!receipt.settled) continue; + this.receipts.delete(key); + if (this.receipts.size < limit) return true; + } + return this.receipts.size < limit; + } + private async perform(owner: string, request: ControlRequest): Promise { if (this.closed) throw new ControlError("CLOSED", "Agent control is shutting down."); const params = request.params; diff --git a/src/main/services/agent-runtime/ProviderRuntimeLaunch.ts b/src/main/services/agent-runtime/ProviderRuntimeLaunch.ts index cc1f1953..dd62fd43 100644 --- a/src/main/services/agent-runtime/ProviderRuntimeLaunch.ts +++ b/src/main/services/agent-runtime/ProviderRuntimeLaunch.ts @@ -1149,9 +1149,15 @@ function mkdirPrivate(path: string): void { function atomicWrite(path: string, value: string, mode = FILE_MODE): void { mkdirPrivate(dirname(path)); const temporary = `${path}.${process.pid}.${randomUUID()}.tmp`; - writeFileSync(temporary, value, { mode }); - chmodSync(temporary, mode); - renameSync(temporary, path); + try { + writeFileSync(temporary, value, { mode }); + chmodSync(temporary, mode); + renameSync(temporary, path); + } catch (error) { + // The name is random, so a leftover would never be reused or cleaned up. + rmSync(temporary, { force: true }); + throw error; + } } function readOptional(path: string): string | null { diff --git a/src/main/services/agent-runtime/RuntimeGateway.ts b/src/main/services/agent-runtime/RuntimeGateway.ts index 8b2b7463..c372a7d4 100644 --- a/src/main/services/agent-runtime/RuntimeGateway.ts +++ b/src/main/services/agent-runtime/RuntimeGateway.ts @@ -25,6 +25,11 @@ const AGENT_PROVIDERS = new Set([ "codex", "claude", "qwen", "kimi", "opencode", "hermes", "grok", "omp", "pi", "cursor", "minimax", "devin", "antigravity" ]); const MAX_RUNTIME_SESSIONS = 32; +const MAX_TRANSPORT_RESTART_ATTEMPTS = 3; +// Hook helpers write their one message right after connecting. A connection +// that stays silent is closed, so idle clients cannot hold all 64 slots. +const FIRST_MESSAGE_TIMEOUT_MS = 5_000; +const TRANSPORT_RESTART_BASE_DELAY_MS = 500; /** Decision checks in flight, per session and in total; over a cap the call is refused with advice to slow down. */ const MAX_DECISIONS_PER_SESSION = 8; const MAX_DECISIONS_TOTAL = 32; @@ -102,6 +107,8 @@ export interface RuntimeGatewayOptions { runtimeDirectory?: string; windowsHostPath?: string; windowsPipeHostFactory?: (options: WindowsPipeHostTransportOptions) => WindowsPipeHostTransport; + /** A connection must send its one message within this time (default 5 s). */ + firstMessageTimeoutMs?: number; onSignal?(terminalSessionId: string, signal: RuntimeLifecycleSignal): void; onAnswerCaptureRevoked?(terminalSessionId: string): void; /** @@ -121,9 +128,13 @@ export class RuntimeGateway { private readonly onAnswerCaptureRevoked: RuntimeGatewayOptions["onAnswerCaptureRevoked"]; private readonly onPermissionRequest: RuntimeGatewayOptions["onPermissionRequest"]; private readonly now: () => number; + private readonly firstMessageTimeoutMs: number; private readonly checks = new Set(); private readonly leases = new Map(); private readonly sockets = new Set(); + private closed = false; + private restartTimer: ReturnType | undefined; + private restartAttempts = 0; private server: Server | null = null; private windowsTransport: WindowsPipeHostTransport | null = null; private endpoint: string | null = null; @@ -135,6 +146,7 @@ export class RuntimeGateway { this.windowsHostPath = options.windowsHostPath; this.windowsPipeHostFactory = options.windowsPipeHostFactory ?? ((transportOptions) => new WindowsPipeHostTransport(transportOptions)); + this.firstMessageTimeoutMs = options.firstMessageTimeoutMs ?? FIRST_MESSAGE_TIMEOUT_MS; this.onSignal = options.onSignal; this.onAnswerCaptureRevoked = options.onAnswerCaptureRevoked; this.onPermissionRequest = options.onPermissionRequest; @@ -152,15 +164,31 @@ export class RuntimeGateway { if (!this.windowsHostPath) { throw new Error("Agent runtime access on Windows requires the packaged current-user-only named-pipe host."); } + this.closed = false; const transport = this.windowsPipeHostFactory({ hostPath: this.windowsHostPath, platform: this.platform, parentPid: process.pid }); this.windowsTransport = transport; - const endpoint = await transport.start((socket) => this.accept(socket)); - this.endpoint = endpoint; - return endpoint; + // The pipe host can die later (crash, EPIPE, FATAL frame). Without this + // every later launch failed with "must be started" until restart. + transport.on("fatal", () => this.handleTransportFatal(transport)); + try { + const endpoint = await transport.start((socket) => this.accept(socket)); + if (this.windowsTransport !== transport) { + await transport.close(); + throw new Error("Windows agent pipe host was superseded during startup."); + } + this.endpoint = endpoint; + this.restartAttempts = 0; + return endpoint; + } catch (error) { + await transport.close(); + if (this.windowsTransport === transport) this.windowsTransport = null; + this.endpoint = null; + throw error; + } } const created = await createEndpoint(this.requestedRuntimeDirectory); @@ -238,7 +266,31 @@ export class RuntimeGateway { } } + private handleTransportFatal(transport: WindowsPipeHostTransport): void { + if (this.windowsTransport !== transport) return; + this.windowsTransport = null; + this.endpoint = null; + for (const socket of this.sockets) socket.destroy(); + this.sockets.clear(); + this.scheduleTransportRestart(); + } + + private scheduleTransportRestart(): void { + if (this.closed || this.restartTimer || this.restartAttempts >= MAX_TRANSPORT_RESTART_ATTEMPTS) return; + const delay = TRANSPORT_RESTART_BASE_DELAY_MS * 2 ** this.restartAttempts; + this.restartAttempts += 1; + this.restartTimer = setTimeout(() => { + this.restartTimer = undefined; + if (this.closed || this.windowsTransport) return; + this.start().catch(() => this.scheduleTransportRestart()); + }, delay); + this.restartTimer.unref?.(); + } + async close(): Promise { + this.closed = true; + clearTimeout(this.restartTimer); + this.restartTimer = undefined; for (const socket of this.sockets) socket.destroy(); this.sockets.clear(); for (const lease of this.leases.values()) { @@ -271,9 +323,12 @@ export class RuntimeGateway { let pending = Buffer.alloc(0); let handled = false; const close = () => { + clearTimeout(firstMessage); this.sockets.delete(socket); socket.destroy(); }; + const firstMessage = setTimeout(close, this.firstMessageTimeoutMs); + firstMessage.unref?.(); socket.setNoDelay(true); socket.on("data", (chunk) => { if (handled) return; @@ -283,6 +338,7 @@ export class RuntimeGateway { const newline = pending.indexOf(0x0a); if (newline < 0) return; handled = true; + clearTimeout(firstMessage); try { const value: unknown = JSON.parse(pending.subarray(0, newline).toString("utf8")); // Decision hooks keep the socket open for the answer; every other message is unchanged. @@ -305,7 +361,10 @@ export class RuntimeGateway { } }); socket.on("error", close); - socket.on("close", () => this.sockets.delete(socket)); + socket.on("close", () => { + clearTimeout(firstMessage); + this.sockets.delete(socket); + }); } private answerCaptureIsActive(value: unknown): boolean { diff --git a/src/main/services/browser/BrowserAuditStore.ts b/src/main/services/browser/BrowserAuditStore.ts index e06c025e..4a1894e8 100644 --- a/src/main/services/browser/BrowserAuditStore.ts +++ b/src/main/services/browser/BrowserAuditStore.ts @@ -1,6 +1,6 @@ import { createHash } from "node:crypto"; import { basename, dirname, join } from "node:path"; -import { mkdir, open, readFile, readdir, rename, stat, unlink } from "node:fs/promises"; +import { appendFile, mkdir, open, readFile, readdir, rename, stat, truncate, unlink } from "node:fs/promises"; const AUDIT_VERSION = 1; const DEFAULT_MAX_BYTES = 100 * 1024 * 1024; @@ -111,8 +111,16 @@ export class BrowserAuditStore { await mkdir(dirname(this.filePath), { recursive: true }); const handle = await open(this.filePath, "a", 0o600); try { - await handle.writeFile(line, "utf8"); - await handle.sync(); + const sizeBefore = (await handle.stat()).size; + try { + await handle.writeFile(line, "utf8"); + await handle.sync(); + } catch (error) { + // A partial append (ENOSPC) would merge with the next record and break + // the chain for good; cut the file back to the last whole record. + await handle.truncate(sizeBefore).catch(() => undefined); + throw error; + } } finally { await handle.close(); } @@ -140,7 +148,7 @@ export class BrowserAuditStore { return { valid: false, records, lastHash: previousHash }; } const { hash, ...base } = record; - if ((records > 0 && record.previousHash !== previousHash) || hashRecord(base) !== hash) { + if ((records > 0 && record.previousHash !== previousHash) || !recordHashMatches(base, hash)) { return { valid: false, records, lastHash: previousHash }; } previousHash = hash; @@ -157,6 +165,7 @@ export class BrowserAuditStore { private async initialize(): Promise { await mkdir(dirname(this.filePath), { recursive: true }); + await this.repairTornTail(); await this.pruneExpired(); const files = await this.auditFiles(); let previousHash: string | null = null; @@ -169,7 +178,7 @@ export class BrowserAuditStore { try { const record = JSON.parse(line) as BrowserAuditRecord; const { hash, ...base } = record; - if ((records > 0 && record.previousHash !== previousHash) || hashRecord(base) !== hash) { + if ((records > 0 && record.previousHash !== previousHash) || !recordHashMatches(base, hash)) { throw new Error("Browser audit hash chain is invalid."); } previousHash = hash; @@ -185,6 +194,37 @@ export class BrowserAuditStore { this.sequence = sequence; } + /** + * Every append ends with a newline, so an active file without one was cut + * during a write (crash, full disk). A last record that is whole only gets + * its newline back; a partial one is removed. Anything else that does not + * verify still fails closed. + */ + private async repairTornTail(): Promise { + let content: Buffer; + try { + content = await readFile(this.filePath); + } catch { + return; + } + if (content.length === 0 || content[content.length - 1] === 0x0a) return; + const lineStart = content.lastIndexOf(0x0a) + 1; + const tail = content.subarray(lineStart).toString("utf8"); + let whole = false; + try { + const { hash, ...base } = JSON.parse(tail) as BrowserAuditRecord; + whole = recordHashMatches(base, hash); + } catch { + whole = false; + } + if (whole) { + await appendFile(this.filePath, "\n", { mode: 0o600 }); + return; + } + console.warn(`CanvasTTY removed a browser audit record cut off during a write (${content.length - lineStart} bytes).`); + await truncate(this.filePath, lineStart); + } + private async rotateIfNeeded(incomingBytes: number): Promise { let currentBytes = 0; try { @@ -217,7 +257,7 @@ export class BrowserAuditStore { const rotated = entries .filter((entry) => entry.isFile() && /^browser-audit-.+\.jsonl$/.test(entry.name)) .map((entry) => join(directory, entry.name)) - .sort((left, right) => basename(left).localeCompare(basename(right))); + .sort((left, right) => byCodeUnit(basename(left), basename(right))); try { await stat(this.filePath); rotated.push(this.filePath); @@ -260,15 +300,29 @@ function redactUrl(value: string): string { } function hashRecord(record: Omit): string { - return createHash("sha256").update(stableJson(record)).digest("hex"); + return createHash("sha256").update(stableJson(record, byCodeUnit)).digest("hex"); } -function stableJson(value: unknown): string { - if (Array.isArray(value)) return `[${value.map(stableJson).join(",")}]`; +/** + * Records written before keys were sorted by code unit were hashed with + * localeCompare, whose order follows the system locale. They are still + * accepted when they verify under the current locale, as they did before. + */ +function recordHashMatches(record: Omit, hash: unknown): boolean { + if (typeof hash !== "string") return false; + return hashRecord(record) === hash + || createHash("sha256").update(stableJson(record, byLocale)).digest("hex") === hash; +} + +const byCodeUnit = (left: string, right: string): number => (left < right ? -1 : left > right ? 1 : 0); +const byLocale = (left: string, right: string): number => left.localeCompare(right); + +function stableJson(value: unknown, order: (left: string, right: string) => number): string { + if (Array.isArray(value)) return `[${value.map((item) => stableJson(item, order)).join(",")}]`; if (value && typeof value === "object") { return `{${Object.entries(value as Record) - .sort(([left], [right]) => left.localeCompare(right)) - .map(([key, entry]) => `${JSON.stringify(key)}:${stableJson(entry)}`) + .sort(([left], [right]) => order(left, right)) + .map(([key, entry]) => `${JSON.stringify(key)}:${stableJson(entry, order)}`) .join(",")}}`; } return JSON.stringify(value); diff --git a/src/main/services/browser/BrowserStore.ts b/src/main/services/browser/BrowserStore.ts index 4df6b256..fd90655b 100644 --- a/src/main/services/browser/BrowserStore.ts +++ b/src/main/services/browser/BrowserStore.ts @@ -1,5 +1,5 @@ import { dirname, join } from "node:path"; -import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; +import { mkdir, readFile, rename, unlink, writeFile } from "node:fs/promises"; import { isSafeBrowserUrl, MAX_BROWSER_TABS } from "./BrowserPolicyService.ts"; export const BROWSER_STORE_VERSION = 1; @@ -64,12 +64,20 @@ export class BrowserStore { private persist(): Promise { const snapshot = `${JSON.stringify(this.value, null, 2)}\n`; const temporaryPath = `${this.filePath}.${process.pid}.tmp`; - this.writeQueue = this.writeQueue.then(async () => { + // A failed write must not poison the queue: every later save would reject + // with the same error. Each save still reports its own failure. + const write = this.writeQueue.catch(() => undefined).then(async () => { await mkdir(dirname(this.filePath), { recursive: true }); - await writeFile(temporaryPath, snapshot, { encoding: "utf8", mode: 0o600 }); - await rename(temporaryPath, this.filePath); + try { + await writeFile(temporaryPath, snapshot, { encoding: "utf8", mode: 0o600 }); + await rename(temporaryPath, this.filePath); + } catch (error) { + await unlink(temporaryPath).catch(() => undefined); + throw error; + } }); - return this.writeQueue; + this.writeQueue = write; + return write; } } diff --git a/src/main/services/homeMedia.ts b/src/main/services/homeMedia.ts new file mode 100644 index 00000000..105a4e6f --- /dev/null +++ b/src/main/services/homeMedia.ts @@ -0,0 +1,53 @@ +import { dirname, extname, isAbsolute, relative, sep } from "node:path"; +import { open, realpath } from "node:fs/promises"; +import { constants } from "node:fs"; + +export const MAX_HOME_MEDIA_BYTES = 25 * 1024 * 1024; +const MAX_HOME_MEDIA_PATH_LENGTH = 4_096; +const HOME_MEDIA_MIME: Record = { + ".png": "image/png", + ".jpg": "image/jpeg", + ".jpeg": "image/jpeg", + ".webp": "image/webp", + ".gif": "image/gif" +}; + +/** A saved Home media path: absolute, bounded, and one of the image types Home can show. */ +export function isHomeMediaPath(value: unknown): value is string { + return typeof value === "string" + && value.length > 0 + && value.length <= MAX_HOME_MEDIA_PATH_LENGTH + && !value.includes("\0") + && isAbsolute(value) + && Object.hasOwn(HOME_MEDIA_MIME, extname(value).toLowerCase()); +} + +/** + * Reads the Home image as a data URL. A symbolic link is followed only while + * its target stays inside the folder the file was chosen from, and the target must + * itself be a supported image no larger than 25 MB. + */ +export async function readHomeMedia(path: string): Promise { + if (!isHomeMediaPath(path)) throw new Error("Unsupported media type."); + const [target, folder] = await Promise.all([realpath(path), realpath(dirname(path))]); + const fromFolder = relative(folder, target); + if (!fromFolder || fromFolder.startsWith(`..${sep}`) || fromFolder === ".." || isAbsolute(fromFolder)) { + throw new Error("Media link points outside the chosen folder."); + } + const mime = HOME_MEDIA_MIME[extname(target).toLowerCase()]; + if (!mime) throw new Error("Unsupported media type."); + + // Read what was checked: the resolved file, opened without following a link + // swapped in after the check, and sized from the open handle. + const handle = await open(target, constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0)); + try { + const metadata = await handle.stat(); + if (!metadata.isFile() || metadata.size > MAX_HOME_MEDIA_BYTES) { + throw new Error("Media must be a file smaller than 25 MB."); + } + const content = await handle.readFile(); + return `data:${mime};base64,${content.toString("base64")}`; + } finally { + await handle.close(); + } +} diff --git a/src/main/services/providerCliRegistry.ts b/src/main/services/providerCliRegistry.ts index ffb8fddb..6e111a59 100644 --- a/src/main/services/providerCliRegistry.ts +++ b/src/main/services/providerCliRegistry.ts @@ -285,10 +285,25 @@ function escapeCommandPromptCommand(value: string): string { return value.replace(COMMAND_PROMPT_META_CHARACTERS, "^$1"); } +// Arguments are escaped twice. cmd.exe removes one level of carets when it +// reads the /c line, then a batch file (an npm shim runs `node cli.js %*`) +// parses the text %* expands to again. cmd.exe does not treat \" as an +// escaped quote, so with one level an argument holding a quote followed by +// & or | (the --settings hook command, for example) was cut there and the +// rest ran as a separate command. With every quote escaped at both levels +// cmd.exe never sees a quoted region and every operator stays escaped. +// +// Program-side quoting follows the MSVC rules: backslashes before a quote and +// at the end are doubled. The old lookahead regex doubled only one of two or +// more backslashes before a quote, so in `a\\"b` the quote ended the argument +// instead of being part of it. function escapeCommandPromptArgument(value: string): string { - let escaped = value.replace(/(?=(\\+?)?)\1"/g, "$1$1\\\""); - escaped = escaped.replace(/(?=(\\+?)?)\1$/, "$1$1"); - return `"${escaped}"`.replace(COMMAND_PROMPT_META_CHARACTERS, "^$1"); + const escaped = value + .replace(/(\\*)"/g, (_match, slashes: string) => `${slashes}${slashes}\\"`) + .replace(/(\\+)$/, "$1$1"); + return `"${escaped}"` + .replace(COMMAND_PROMPT_META_CHARACTERS, "^$1") + .replace(COMMAND_PROMPT_META_CHARACTERS, "^$1"); } function providerCandidates(commands: readonly string[], directories: string[], platform: NodeJS.Platform, commandFirst = false): string[] { @@ -398,13 +413,25 @@ function sharedUserDirectories( function resolveWindowsCommandPrompt( environment: Readonly, inspectCandidate: ResolveProviderCliInput["inspectCandidate"] +): string | null { + return windowsCommandPromptPath(environment, (path) => inspectCandidate(path, "win32") === null); +} + +/** + * cmd.exe for batch providers and the terminal fallback: ComSpec, then + * %SystemRoot%\System32\cmd.exe. PATH is never searched, so a cmd.exe in a + * project folder or another PATH entry cannot stand in for it. + */ +export function windowsCommandPromptPath( + environment: Readonly, + usable: (path: string) => boolean ): string | null { const configured = environment.ComSpec || environment.COMSPEC; - if (configured && inspectCandidate(configured, "win32") === null) return configured; + if (configured && usable(configured)) return configured; const systemRoot = environment.SystemRoot || environment.WINDIR; if (!systemRoot) return null; const candidate = win32.join(systemRoot, "System32", "cmd.exe"); - return inspectCandidate(candidate, "win32") === null ? candidate : null; + return usable(candidate) ? candidate : null; } function pathEntries(value: string | undefined, platform: NodeJS.Platform, startupDirectory: string): string[] { diff --git a/src/main/services/terminalLaunch.ts b/src/main/services/terminalLaunch.ts index 7ad6e94c..85b78f73 100644 --- a/src/main/services/terminalLaunch.ts +++ b/src/main/services/terminalLaunch.ts @@ -6,6 +6,7 @@ import { openCodeYoloEnvironment } from "./openCodeConfig.ts"; import { autoModeArguments, CLAUDE_SANDBOX_SETTINGS, type LaunchProfile } from "../../shared/autoMode.ts"; import { providerTerminalBatchCommandLine, + windowsCommandPromptPath, type ProviderCliResolution } from "./providerCliRegistry.ts"; @@ -207,13 +208,9 @@ function resolveWindowsCommandPrompt( environment: Readonly, fileExists: (path: string) => boolean ): string { - const configured = environment.ComSpec || environment.COMSPEC; - if (configured && fileExists(configured)) return configured; - const fromPath = findWindowsNativeCommand("cmd", environment, fileExists); - if (fromPath) return fromPath; - const systemRoot = environment.SystemRoot || environment.WINDIR; - const systemCommandPrompt = systemRoot ? win32.join(systemRoot, "System32", "cmd.exe") : null; - if (systemCommandPrompt && fileExists(systemCommandPrompt)) return systemCommandPrompt; + // The same lookup as batch provider launches: no PATH search for cmd.exe. + const commandPrompt = windowsCommandPromptPath(environment, fileExists); + if (commandPrompt) return commandPrompt; throw new Error("No supported Windows shell was found (PowerShell, pwsh, or cmd.exe)."); } diff --git a/tests/agent-control.test.mjs b/tests/agent-control.test.mjs index cede7487..23a1933d 100644 --- a/tests/agent-control.test.mjs +++ b/tests/agent-control.test.mjs @@ -1,7 +1,8 @@ import assert from "node:assert/strict"; import { randomUUID } from "node:crypto"; import { spawn } from "node:child_process"; -import { mkdtemp, readFile, realpath, stat, writeFile } from "node:fs/promises"; +import { mkdtemp, readFile, realpath, rm, stat, writeFile } from "node:fs/promises"; +import { EventEmitter } from "node:events"; import { createConnection } from "node:net"; import { tmpdir } from "node:os"; import { join, resolve } from "node:path"; @@ -32,7 +33,7 @@ function registry() { environment: {}, checked: [] }), snapshot: () => ({}) }; } -async function fixture(t) { +async function fixture(t, gatewayOptions = {}) { const root = await realpath(await mkdtemp(join(tmpdir(), "ctty-control-test-"))); const calls = []; let gateway; @@ -47,7 +48,7 @@ async function fixture(t) { return pty; }); let lifecycleEnabled = true; - gateway = new AgentControlGateway({ userDataPath: root, terminals, lifecycleEnabled: () => lifecycleEnabled }); + gateway = new AgentControlGateway({ userDataPath: root, terminals, lifecycleEnabled: () => lifecycleEnabled, ...gatewayOptions }); const connectionPath = await gateway.start(); const clientPath = join(root, "client-a.json"); t.after(async () => { await gateway.close(); await terminals.shutdown(); }); @@ -93,6 +94,71 @@ test("CLI creates native YOLO with requested directory/title, including concurre assert.equal(f.calls.length, 1); }); +test("request receipts are bounded without locking the gateway, and a refused request can be retried", localSocket, async (t) => { + const f = await fixture(t, { maxReceipts: 3 }); + for (let index = 0; index < 5; index += 1) { + await assert.rejects(f.request("interrupt", { sessionId: `missing-${index}` }, `missing-${index}`), (e) => e.code === "SESSION_NOT_FOUND"); + } + // Before: the fourth mutating request (successful or not) got LIMIT_REACHED until restart. + const { session } = await f.create("create-after-limit"); + assert.equal((await f.create("create-after-limit")).session.id, session.id, "a recent receipt still replays"); + await f.ready(session.id); + await f.request("send", { sessionId: session.id, text: "first" }, "send-first"); + f.signal(session.id, "working", "turn-one"); + await assert.rejects(f.request("send", { sessionId: session.id, text: "second" }, "send-second"), (e) => e.code === "BUSY"); + f.signal(session.id, "idle", "turn-one", { text: "done", truncated: false }); + // BUSY wrote nothing, so the same request id is performed on retry instead of replaying BUSY. + const retried = await f.request("send", { sessionId: session.id, text: "second" }, "send-second"); + assert.equal(retried.sessionId, session.id); + assert.equal(f.calls[0].pty.writes.length, 2); +}); + +test("a failed start leaves nothing listening, so the same gateway can start again", localSocket, async (t) => { + const root = await realpath(await mkdtemp(join(tmpdir(), "ctty-control-start-"))); + t.after(() => rm(root, { recursive: true, force: true })); + const userDataPath = join(root, "user-data"); + await writeFile(userDataPath, "a file where the folder should be"); + const gateway = new AgentControlGateway({ userDataPath, terminals: {}, lifecycleEnabled: () => true }); + t.after(() => gateway.close()); + await assert.rejects(gateway.start()); + await rm(userDataPath); + const connection = await gateway.start(); + const descriptor = JSON.parse(await readFile(connection, "utf8")); + assert.equal((await stat(descriptor.endpoint)).isSocket(), true); +}); + +test("agent control brings the Windows pipe host back after it fails and republishes the endpoint", async (t) => { + t.mock.timers.enable({ apis: ["setTimeout"] }); + const root = await realpath(await mkdtemp(join(tmpdir(), "ctty-control-win-"))); + t.after(() => rm(root, { recursive: true, force: true })); + const transports = []; + let restarted; + const republished = new Promise((resolve) => { restarted = resolve; }); + const gateway = new AgentControlGateway({ + userDataPath: root, terminals: {}, lifecycleEnabled: () => true, platform: "win32", windowsHostPath: "C:\\fake\\host.exe", + onTransportRestarted: (path) => restarted(path), + windowsPipeHostFactory: () => { + const transport = new EventEmitter(); + const index = transports.length; + transport.start = async () => `\\\\.\\pipe\\canvastty-agent-${index}`; + transport.close = async () => undefined; + transports.push(transport); + return transport; + } + }); + t.after(() => gateway.close()); + const connection = await gateway.start(); + assert.match(JSON.parse(await readFile(connection, "utf8")).endpoint, /agent-0$/); + transports[0].emit("fatal", new Error("host exited")); + t.mock.timers.tick(500); + assert.equal(await republished, connection); + assert.match(JSON.parse(await readFile(connection, "utf8")).endpoint, /agent-1$/); + await gateway.close(); + transports[1].emit("fatal", new Error("late")); + t.mock.timers.tick(10_000); + assert.equal(transports.length, 2); +}); + test("controller cannot list, read, interrupt or send to other controllers or UI sessions", localSocket, async (t) => { const f = await fixture(t); const { session } = await f.create(); diff --git a/tests/agent-runtime-gateway.test.mjs b/tests/agent-runtime-gateway.test.mjs index 38cf80c6..169e690c 100644 --- a/tests/agent-runtime-gateway.test.mjs +++ b/tests/agent-runtime-gateway.test.mjs @@ -1,5 +1,6 @@ import assert from "node:assert/strict"; import { spawn } from "node:child_process"; +import { EventEmitter } from "node:events"; import { mkdtemp, rm, stat } from "node:fs/promises"; import { createConnection } from "node:net"; import { tmpdir } from "node:os"; @@ -465,3 +466,68 @@ function childResult(child) { child.once("exit", (code, signal) => resolve({ code, signal, stderr })); }); } + +function fakeWindowsTransports({ failFirstStart = false } = {}) { + const transports = []; + const factory = () => { + const transport = new EventEmitter(); + const index = transports.length; + transport.isRunning = false; + transport.closed = false; + transport.start = async () => { + if (failFirstStart && index === 0) throw new Error("host did not start"); + transport.isRunning = true; + return `\\\\.\\pipe\\canvastty-agent-${index}`; + }; + transport.close = async () => { transport.isRunning = false; transport.closed = true; }; + transports.push(transport); + return transport; + }; + return { transports, factory }; +} + +test("RuntimeGateway restarts a Windows pipe host that failed instead of refusing every later launch", async (t) => { + t.mock.timers.enable({ apis: ["setTimeout"] }); + const { transports, factory } = fakeWindowsTransports(); + const gateway = new RuntimeGateway({ platform: "win32", windowsHostPath: "C:\\fake\\host.exe", windowsPipeHostFactory: factory }); + t.after(() => gateway.close()); + + assert.equal(await gateway.start(), "\\\\.\\pipe\\canvastty-agent-0"); + gateway.registerSession("terminal-before", "claude"); + transports[0].isRunning = false; + transports[0].emit("fatal", new Error("host exited")); + assert.throws(() => gateway.registerSession("terminal-during", "claude"), /must be started/); + + t.mock.timers.tick(500); + await new Promise(setImmediate); + assert.equal(transports.length, 2); + assert.equal(gateway.address, "\\\\.\\pipe\\canvastty-agent-1"); + assert.ok(gateway.registerSession("terminal-after", "claude")); + + await gateway.close(); + transports[1].emit("fatal", new Error("late")); + t.mock.timers.tick(10_000); + assert.equal(transports.length, 2, "a closed gateway does not restart"); +}); + +test("RuntimeGateway drops a Windows transport whose start failed", async (t) => { + const { transports, factory } = fakeWindowsTransports({ failFirstStart: true }); + const gateway = new RuntimeGateway({ platform: "win32", windowsHostPath: "C:\\fake\\host.exe", windowsPipeHostFactory: factory }); + t.after(() => gateway.close()); + await assert.rejects(gateway.start(), /did not start/); + assert.equal(transports[0].closed, true); + assert.equal(await gateway.start(), "\\\\.\\pipe\\canvastty-agent-1"); +}); + +test("RuntimeGateway closes a connection that sends no message, so idle clients cannot hold every slot", POSIX_RUNTIME_GATEWAY_TEST, async (t) => { + const root = await fixture(t); + const gateway = new RuntimeGateway({ runtimeDirectory: root, firstMessageTimeoutMs: 100 }); + const address = await gateway.start(); + t.after(() => gateway.close()); + const idle = createConnection(address); + await new Promise((resolve, reject) => { idle.once("connect", resolve); idle.once("error", reject); }); + const closed = new Promise((resolve) => idle.once("close", resolve)); + const outcome = await Promise.race([closed.then(() => "closed"), new Promise((resolve) => setTimeout(() => resolve("open"), 1_500))]); + idle.destroy(); + assert.equal(outcome, "closed"); +}); diff --git a/tests/atomic-write-cleanup.test.mjs b/tests/atomic-write-cleanup.test.mjs new file mode 100644 index 00000000..c023bfce --- /dev/null +++ b/tests/atomic-write-cleanup.test.mjs @@ -0,0 +1,48 @@ +import assert from "node:assert/strict"; +import { createHash } from "node:crypto"; +import { mkdir, mkdtemp, readdir, rm, stat, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; + +import { createQwenHookSettings } from "../src/main/services/agent-runtime/ProviderRuntimeLaunch.ts"; +import { TerminalSessionStore } from "../src/main/services/TerminalSessionStore.ts"; + +async function fixture(t) { + const root = await mkdtemp(join(tmpdir(), "canvastty-atomic-write-")); + t.after(() => rm(root, { recursive: true, force: true })); + return root; +} + +// A folder with content where the file should be makes the final rename fail, +// as a locked file (EPERM from antivirus on Windows) does. +async function blockRename(path) { + await mkdir(path, { recursive: true }); + await writeFile(join(path, "keep"), ""); +} + +test("provider hook settings leave no temporary file when the rename fails", async (t) => { + const root = await fixture(t); + const target = join(root, `qwen-hooks-${createHash("sha256").update("session-1", "utf8").digest("hex").slice(0, 24)}.json`); + await blockRename(target); + assert.throws(() => createQwenHookSettings({ + helper: { command: process.execPath, args: ["/helper.mjs"] }, + platform: "linux", + runtimeDirectory: root, + terminalSessionId: "session-1" + })); + assert.deepEqual((await readdir(root)).filter((name) => name.endsWith(".tmp")), []); +}); + +test("the terminal session store leaves no temporary file when the rename fails and keeps its folder private", async (t) => { + const root = await fixture(t); + const dataDir = join(root, "user-data"); + const store = new TerminalSessionStore(dataDir); + await blockRename(join(dataDir, "terminal-sessions.json")); + await assert.rejects(store.clear()); + assert.deepEqual((await readdir(dataDir)).filter((name) => name.endsWith(".tmp")), []); + + const fresh = join(root, "fresh", "user-data"); + await new TerminalSessionStore(fresh).clear(); + if (process.platform !== "win32") assert.equal((await stat(fresh)).mode & 0o077, 0, "a folder the store creates is private"); +}); diff --git a/tests/browser-audit-store.test.mjs b/tests/browser-audit-store.test.mjs index 9418dca9..13e08247 100644 --- a/tests/browser-audit-store.test.mjs +++ b/tests/browser-audit-store.test.mjs @@ -3,6 +3,8 @@ import { mkdtemp, readdir, readFile, rm, utimes, writeFile } from "node:fs/promi import { tmpdir } from "node:os"; import { dirname, join } from "node:path"; import test from "node:test"; +import { spawnSync } from "node:child_process"; +import { createHash } from "node:crypto"; import { BrowserAuditStore, @@ -175,3 +177,83 @@ test("BrowserAuditStore propagates storage failures instead of pretending to aud }); await assert.rejects(store.verify()); }); + +test("BrowserAuditStore repairs a line torn by a crash instead of refusing every later action", async (t) => { + const root = await fixture(t, "canvastty-audit-torn-"); + const store = new BrowserAuditStore(root); + await store.append(auditInput("torn-1")); + await store.append(auditInput("torn-2")); + const complete = await readFile(store.filePath, "utf8"); + // A crash or ENOSPC during the append left half a record and no newline. + const third = JSON.stringify({ ...JSON.parse(complete.trim().split("\n")[1]), sequence: 3 }); + await writeFile(store.filePath, complete + third.slice(0, 40)); + + const reopened = new BrowserAuditStore(root); + const warn = console.warn; + console.warn = () => undefined; + try { + const appended = await reopened.append(auditInput("after-crash")); + assert.equal(appended.sequence, 3); + } finally { + console.warn = warn; + } + assert.deepEqual(await reopened.verify(), { valid: true, records: 3, lastHash: (await reopened.verify()).lastHash }); + assert.equal((await readFile(reopened.filePath, "utf8")).split("\n").filter(Boolean).length, 3); + + // A complete last record that only lost its newline is kept, not dropped. + const whole = await readFile(reopened.filePath, "utf8"); + await writeFile(reopened.filePath, whole.slice(0, -1)); + const again = new BrowserAuditStore(root); + assert.equal((await again.append(auditInput("after-newline"))).sequence, 4); + assert.equal((await again.verify()).valid, true); +}); + +function runInLocale(locale, root, action) { + const script = ` + const { BrowserAuditStore } = await import(${JSON.stringify(new URL("../src/main/services/browser/BrowserAuditStore.ts", import.meta.url).href)}); + const store = new BrowserAuditStore(${JSON.stringify(root)}); + if (${JSON.stringify(action)} === "append") { + await store.append({ timestamp: 1, requestId: "r-" + ${JSON.stringify(locale)}, actorKind: "agent", actorId: "a", operation: "browser_click", + phase: "result", ok: true, details: { "z": 1, "รค": 2, "Zeta": 3, "alpha": 4 } }); + } + process.stdout.write(JSON.stringify(await store.verify()));`; + const result = spawnSync(process.execPath, ["--experimental-strip-types", "--no-warnings", "--input-type=module", "-e", script], { + env: { ...process.env, LC_ALL: locale, LANG: locale }, + encoding: "utf8" + }); + assert.equal(result.status, 0, result.stderr); + return JSON.parse(result.stdout); +} + +test("the audit hash does not depend on the system locale", async (t) => { + const root = await fixture(t, "canvastty-audit-locale-"); + assert.equal(runInLocale("sv_SE.UTF-8", root, "append").valid, true); + // The same log read by a process with another collation. + const other = runInLocale("en_US.UTF-8", root, "verify"); + assert.equal(other.valid, true); + assert.equal(other.records, 1); + assert.equal(runInLocale("en_US.UTF-8", root, "append").records, 2); + assert.equal(runInLocale("sv_SE.UTF-8", root, "verify").valid, true); +}); + +test("records hashed with the earlier locale-ordered keys still verify and extend the chain", async (t) => { + const root = await fixture(t, "canvastty-audit-legacy-"); + const legacyJson = (value) => { + if (Array.isArray(value)) return `[${value.map(legacyJson).join(",")}]`; + if (value && typeof value === "object") { + return `{${Object.entries(value).sort(([left], [right]) => left.localeCompare(right)) + .map(([key, entry]) => `${JSON.stringify(key)}:${legacyJson(entry)}`).join(",")}}`; + } + return JSON.stringify(value); + }; + const store = new BrowserAuditStore(root); + const first = await store.append(auditInput("legacy-1", { details: { Zeta: 1, alpha: 2 } })); + const { hash: _hash, ...base } = first; + const legacy = { ...base, hash: createHash("sha256").update(legacyJson(base)).digest("hex") }; + assert.notEqual(legacy.hash, first.hash, "the fixture differs between the two orderings"); + await writeFile(store.filePath, `${JSON.stringify(legacy)}\n`); + const reopened = new BrowserAuditStore(root); + const next = await reopened.append(auditInput("new-2")); + assert.equal(next.previousHash, legacy.hash); + assert.deepEqual((await reopened.verify()).valid, true); +}); diff --git a/tests/browser-ipc-security.test.mjs b/tests/browser-ipc-security.test.mjs index f536a08e..7d450c65 100644 --- a/tests/browser-ipc-security.test.mjs +++ b/tests/browser-ipc-security.test.mjs @@ -26,6 +26,35 @@ test("privileged browser IPC validates the trusted main renderer", async () => { "browserSetViewport" ]) { const handler = source.slice(source.indexOf(`IPC.${channel}`), source.indexOf(`IPC.${channel}`) + 320); - assert.match(handler, /assertMainRenderer\(event, getMainWindow\)/, `${channel} must validate its sender`); + assert.match(handler, /(assertMainRenderer|isMainRenderer)\(event, getMainWindow\)/, `${channel} must validate its sender`); + } +}); + +test("channels that change settings, plugins, secrets or terminals accept only the main renderer", async () => { + const source = await readFile(ipcPath, "utf8"); + for (const channel of [ + "settingsUpdate", + "mediaRead", + "pluginsPreviewInstall", + "pluginsInstall", + "pluginsSetModules", + "pluginsSetEnabled", + "pluginsUninstall", + "pluginsOpenExternal", + "pluginsSecretsGet", + "pluginsSecretsSet", + "pluginsSecretsDelete", + "providerSecretsStatus", + "providerSecretsSet", + "providerSecretsClear", + "terminalCreate", + "terminalRestart", + "terminalInput", + "terminalDispose" + ]) { + const start = source.indexOf(`IPC.${channel},`); + assert.notEqual(start, -1, `${channel} handler is registered`); + const handler = source.slice(start, source.indexOf("ipcMain.", start)); + assert.match(handler, /(assertMainRenderer|isMainRenderer)\(event, getMainWindow\)/, `${channel} must validate its sender`); } }); diff --git a/tests/browser-policy-store.test.mjs b/tests/browser-policy-store.test.mjs index 5d70ff9b..e86d0670 100644 --- a/tests/browser-policy-store.test.mjs +++ b/tests/browser-policy-store.test.mjs @@ -1,5 +1,5 @@ import assert from "node:assert/strict"; -import { mkdtemp, mkdir, readFile, readdir, realpath, rm, stat, symlink, truncate, writeFile } from "node:fs/promises"; +import { chmod, mkdtemp, mkdir, readFile, readdir, realpath, rm, stat, symlink, truncate, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { basename, dirname, join, relative } from "node:path"; import test from "node:test"; @@ -217,6 +217,41 @@ test("BrowserStore safely restores, atomically normalizes, and persists only the assert.equal((await readdir(root)).some((name) => name.endsWith(".tmp")), false); }); +test("BrowserStore recovers after one failed write instead of rejecting every later save", { skip: process.platform === "win32" || process.getuid?.() === 0 }, async (t) => { + const root = await fixture(t, "canvastty-store-fail-"); + const dataDir = join(root, "data"); + await mkdir(dataDir); + const store = new BrowserStore(dataDir); + await store.load(); + + await chmod(dataDir, 0o500); + try { + await assert.rejects(store.replace([{ id: "tab-a", url: "https://a.example/" }], "tab-a")); + } finally { + await chmod(dataDir, 0o700); + } + // The in-memory state keeps the change the caller asked for. + assert.equal(store.get().activeTabId, "tab-a"); + + const next = await store.replace([{ id: "tab-b", url: "https://b.example/" }], "tab-b"); + assert.equal(next.activeTabId, "tab-b"); + assert.deepEqual(JSON.parse(await readFile(store.filePath, "utf8")).tabs, [{ id: "tab-b", url: "https://b.example/" }]); + await store.clear(); + assert.deepEqual(JSON.parse(await readFile(store.filePath, "utf8")).tabs, []); +}); + +test("BrowserService keeps tab state when saving fails and settings do not leave the save unhandled", async () => { + const service = await readFile(new URL("../src/main/services/BrowserService.ts", import.meta.url), "utf8"); + const main = await readFile(new URL("../src/main/index.ts", import.meta.url), "utf8"); + const persistRuntime = service.slice(service.indexOf("private async persistRuntime"), service.indexOf("private destroyRuntimeTabs")); + assert.match(persistRuntime, /try \{[\s\S]*this\.store\.replace[\s\S]*\} catch/); + assert.match(persistRuntime, /this\.persisted = this\.store\.get\(\)/); + const clearSaved = service.slice(service.indexOf("private async clearSavedTabs"), service.indexOf("private destroyRuntimeTabs")); + assert.match(clearSaved, /try \{\s*await this\.store\.clear\(\);\s*\} catch/); + assert.equal(service.match(/this\.store\.clear\(\)/g).length, 1, "every clear goes through clearSavedTabs"); + assert.match(main, /browserService\?\.setRestoreTabs\(next\.browserRestoreTabs\)\.catch\(/); +}); + test("BrowserStore treats corrupt persisted input as an empty safe session", async (t) => { const root = await fixture(t, "canvastty-store-corrupt-"); const path = join(root, "browser-state.json"); diff --git a/tests/github-auth.test.mjs b/tests/github-auth.test.mjs index 81002a86..79303955 100644 --- a/tests/github-auth.test.mjs +++ b/tests/github-auth.test.mjs @@ -229,3 +229,45 @@ async function waitFor(predicate, timeoutMs = 1000) { await new Promise((resolve) => setImmediate(resolve)); } } + +test("device polling survives network errors and timeouts, backs off, and honours slow_down", async () => { + const userData = await mkdtemp(`${tmpdir()}/canvastty-github-auth-retry-`); + let clock = 1_000_000; + const waits = []; + const answers = [ + () => { throw new TypeError("fetch failed"); }, + () => { throw Object.assign(new Error("The operation was aborted."), { name: "AbortError" }); }, + () => new Response("unavailable", { status: 503 }), + () => Response.json({ error: "slow_down", interval: 20 }), + () => Response.json({ error: "authorization_pending" }), + () => Response.json({ access_token: "access", token_type: "bearer", scope: "" }) + ]; + const fetcher = async (url) => { + if (String(url).endsWith("/login/device/code")) { + return Response.json({ device_code: "d", user_code: "ABCD-1234", verification_uri: "https://github.com/login/device", expires_in: 900, interval: 5 }); + } + if (String(url).endsWith("/login/oauth/access_token")) return answers.shift()(); + if (String(url) === "https://api.github.com/user") return Response.json({ login: "howdeploy" }); + return new Response("missing", { status: 404 }); + }; + const warn = console.warn; + console.warn = () => undefined; + try { + const service = new GithubAuthService(userData, "client-id", { + fetcher, + safeStorage, + now: () => clock, + delay: async (ms) => { waits.push(ms); clock += ms; } + }); + await service.startDeviceFlow(); + await waitFor(async () => (await service.status()).authorized); + assert.equal(answers.length, 0, "every answer was consumed; the poll did not stop at the first error"); + // 5 s, then doubled after each failure (10, 20, 40), then slow_down's 20 s from GitHub (+5 over + // the current interval as a floor), kept for later polls. + assert.deepEqual(waits, [5_000, 10_000, 20_000, 40_000, 45_000, 45_000]); + await service.signOut(); + } finally { + console.warn = warn; + await rm(userData, { recursive: true, force: true }); + } +}); diff --git a/tests/home-media.test.mjs b/tests/home-media.test.mjs new file mode 100644 index 00000000..448db558 --- /dev/null +++ b/tests/home-media.test.mjs @@ -0,0 +1,60 @@ +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, realpath, rm, symlink, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; + +import { isHomeMediaPath, readHomeMedia } from "../src/main/services/homeMedia.ts"; +import { SettingsStore } from "../src/main/services/SettingsStore.ts"; + +const PNG = Buffer.from("89504e470d0a1a0a", "hex"); + +async function fixture(t) { + const root = await realpath(await mkdtemp(join(tmpdir(), "canvastty-home-media-"))); + t.after(() => rm(root, { recursive: true, force: true })); + await mkdir(join(root, "pictures")); + await mkdir(join(root, "private")); + await writeFile(join(root, "pictures", "wall.png"), PNG); + await writeFile(join(root, "private", "secret.png"), Buffer.from("not yours")); + return root; +} + +test("Home media reads the chosen image and refuses a link that leaves its folder", { skip: process.platform === "win32" }, async (t) => { + const root = await fixture(t); + const pictures = join(root, "pictures"); + assert.equal(await readHomeMedia(join(pictures, "wall.png")), `data:image/png;base64,${PNG.toString("base64")}`); + + await symlink(join(pictures, "wall.png"), join(pictures, "same-folder.png")); + assert.match(await readHomeMedia(join(pictures, "same-folder.png")), /^data:image\/png;base64,/); + + await symlink(join(root, "private", "secret.png"), join(pictures, "escape.png")); + await assert.rejects(readHomeMedia(join(pictures, "escape.png")), /outside/); + + await symlink(join(root, "private", "secret.png"), join(pictures, "up.png")); + await assert.rejects(readHomeMedia(join(pictures, "..", "pictures", "up.png")), /outside/); + + await writeFile(join(pictures, "notes.txt"), "text"); + await symlink(join(pictures, "notes.txt"), join(pictures, "renamed.png")); + await assert.rejects(readHomeMedia(join(pictures, "renamed.png")), /Unsupported media type/); +}); + +test("settings accept only an absolute image path as Home media", async (t) => { + const absolute = process.platform === "win32" ? "C:\\Pictures\\wall.png" : "/srv/me/Pictures/wall.png"; + assert.equal(isHomeMediaPath(absolute), true); + assert.equal(isHomeMediaPath("wall.png"), false); + assert.equal(isHomeMediaPath("/srv/me/.ssh/id_ed25519"), false); + assert.equal(isHomeMediaPath(`/srv/me/${"a".repeat(5000)}.png`), false); + assert.equal(isHomeMediaPath("/srv/me/a\u0000.png"), false); + + const root = await mkdtemp(join(tmpdir(), "canvastty-home-media-settings-")); + t.after(() => rm(root, { recursive: true, force: true })); + const store = new SettingsStore(root, "en"); + await store.load(); + assert.equal((await store.update({ mediaPath: absolute })).mediaPath, absolute); + assert.equal((await store.update({ mediaPath: "/srv/me/.ssh/id_ed25519" })).mediaPath, absolute); + assert.equal((await store.update({ mediaPath: "relative.png" })).mediaPath, absolute); + assert.equal((await store.update({ mediaPath: null })).mediaPath, null); + + await writeFile(join(root, "settings.json"), JSON.stringify({ mediaPath: "/etc/passwd" })); + assert.equal((await new SettingsStore(root, "en").load()).mediaPath, null); +}); diff --git a/tests/limits-dispose.test.mjs b/tests/limits-dispose.test.mjs new file mode 100644 index 00000000..d9cabe05 --- /dev/null +++ b/tests/limits-dispose.test.mjs @@ -0,0 +1,30 @@ +import assert from "node:assert/strict"; +import { mkdtemp, rm, stat, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import test from "node:test"; +import { LimitsService } from "../src/main/services/LimitsService.ts"; + +test("disposing limits while the Kimi usage server is starting leaves no `kimi web` process", { skip: process.platform === "win32" }, async (t) => { + const root = await mkdtemp(join(tmpdir(), "canvastty-kimi-dispose-")); + t.after(() => rm(root, { recursive: true, force: true })); + const marker = join(root, "kimi-web-started"); + const kimi = join(root, "kimi"); + // Records that it ran and exits at once, so a leaked start leaves only the marker behind. + await writeFile(kimi, `#!/bin/sh\necho "$@" > ${JSON.stringify(marker)}\n`, { mode: 0o700 }); + const service = new LimitsService({ + get(provider) { + if (provider === "kimi") { + return { state: "available", provider, executable: kimi, launcher: "native", environment: {}, checked: [] }; + } + return { state: "unavailable", provider, reason: "cli-not-found", checked: [], diagnostic: "" }; + } + }, "test"); + const reading = service.get(); + // The Kimi client is waiting for a free loopback port: nothing is spawned yet. + service.dispose(); + const snapshot = await reading; + assert.equal(snapshot.providers.find((provider) => provider.provider === "kimi").state, "unavailable"); + await new Promise((resolve) => setTimeout(resolve, 200)); + await assert.rejects(stat(marker), { code: "ENOENT" }, "kimi web must not start after dispose"); +}); diff --git a/tests/orchestration-gateway.test.mjs b/tests/orchestration-gateway.test.mjs index 6f541eae..570a6e57 100644 --- a/tests/orchestration-gateway.test.mjs +++ b/tests/orchestration-gateway.test.mjs @@ -359,3 +359,56 @@ test("unknown tools and invalid arguments never reach the handler", async (t) => assert.match(failure.error.message, /cwd/u); terminals.disposeAll(); }); + +test("cancel reaches the running command and the answer is CANCELED, not the late result", async (t) => { + const directory = await mkdtemp(join(tmpdir(), "canvastty-orchestration-cancel-")); + t.after(() => rm(directory, { recursive: true, force: true })); + let signal = null; + let finish; + const gateway = new OrchestrationGateway({ + runtimeDirectory: join(directory, "runtime"), + handler: { + execute: (_sessionId, _request, abortSignal) => { + signal = abortSignal ?? null; + return new Promise((resolve) => { finish = resolve; }); + } + } + }); + await gateway.start(); + t.after(() => gateway.stop()); + const capability = gateway.registerOrchestrator({ terminalSessionId: "orchestrator-1" }); + const { client } = await authenticatedClient(gateway, capability); + t.after(() => client.socket.destroy()); + + const answer = line(client, { v: ORCHESTRATION_BRIDGE_PROTOCOL_VERSION, type: "request", id: "slow-1", tool: "list_agents", arguments: {} }); + for (let i = 0; i < 100 && !finish; i += 1) await new Promise((resolve) => setTimeout(resolve, 10)); + client.send({ v: ORCHESTRATION_BRIDGE_PROTOCOL_VERSION, type: "cancel", id: "slow-1" }); + await new Promise((resolve) => setTimeout(resolve, 30)); + finish({ agents: [] }); + const response = await answer; + assert.equal(signal?.aborted, true, "the handler received the abort signal"); + assert.equal(response.error?.code, "CANCELED"); +}); + +test("a spawn_agent canceled while it was starting closes the agent it created", async () => { + const controller = new AbortController(); + const canceled = []; + const control = { + status: () => ({ role: "orchestrator", provider: "codex" }), + spawn: async () => { + controller.abort(); + return { id: "child-1", provider: "codex", status: "running", title: "worker" }; + }, + cancel: (id) => canceled.push(id) + }; + const handler = new ScopedOrchestrationHandler(control); + await assert.rejects( + handler.execute("orchestrator-1", { id: "spawn-1", tool: "spawn_agent", arguments: { provider: "codex", cwd: process.cwd() } }, controller.signal), + (error) => error.bridgeError?.code === "CANCELED" || error.code === "CANCELED" + ); + assert.deepEqual(canceled, ["child-1"]); + const late = new AbortController(); + late.abort(); + await assert.rejects(handler.execute("orchestrator-1", { id: "spawn-2", tool: "spawn_agent", arguments: {} }, late.signal)); + assert.deepEqual(canceled, ["child-1"], "nothing is spawned after cancel"); +}); diff --git a/tests/plugin-manager.test.mjs b/tests/plugin-manager.test.mjs index 22721009..6df747a6 100644 --- a/tests/plugin-manager.test.mjs +++ b/tests/plugin-manager.test.mjs @@ -1790,3 +1790,44 @@ test("the anonymous showcase search reports when GitHub's rate limit resets", as assert.match(message, /Signing in to GitHub raises the limit/u); assert.match(githubRateLimitMessage(new Response("", { status: 429 })), /try again in a minute/u); }); + +test("plugin storage that cannot be read is not replaced by the next write", { skip: process.platform === "win32" || process.getuid?.() === 0 }, async () => { + const userData = await mkdtemp(join(tmpdir(), "canvastty-plugin-storage-")); + const fixture = new URL("../examples/plugins/studio-kit/", import.meta.url); + const manager = new PluginManager(userData, async (_url, destination) => { + await cp(fixture, destination, { recursive: true }); + }); + const warn = console.warn; + console.warn = () => undefined; + try { + await manager.load(); + const installed = await manager.install((await manager.previewInstall("https://github.com/example/studio-kit")).token); + const id = installed.manifest.id; + const path = join(userData, "plugin-storage", `${id}.json`); + await manager.storageSet(id, "a", 1); + await manager.storageSet(id, "b", 2); + + // A read error (permissions, a locked file) refuses the write instead of saving only the new key. + const { chmod, readdir } = await import("node:fs/promises"); + await chmod(path, 0o000); + try { + await assert.rejects(manager.storageSet(id, "c", 3), /could not be read/); + } finally { + await chmod(path, 0o600); + } + assert.deepEqual(JSON.parse(await readFile(path, "utf8")), { a: 1, b: 2 }); + + // A file that is not valid storage is kept aside before a new one is started. + await writeFile(path, "{\"a\": 1, \"b\":"); + await manager.storageSet(id, "c", 3); + assert.deepEqual(JSON.parse(await readFile(path, "utf8")), { c: 3 }); + const kept = (await readdir(join(userData, "plugin-storage"))).filter((name) => name.startsWith(`${id}.json.unreadable-`)); + assert.equal(kept.length, 1); + assert.equal(await readFile(join(userData, "plugin-storage", kept[0]), "utf8"), "{\"a\": 1, \"b\":"); + await manager.uninstall(id); + assert.deepEqual((await readdir(join(userData, "plugin-storage"))).filter((name) => name.startsWith(id)), []); + } finally { + console.warn = warn; + await rm(userData, { recursive: true, force: true }); + } +}); diff --git a/tests/plugin-media-service.test.mjs b/tests/plugin-media-service.test.mjs index 04ac4ff1..481a3f00 100644 --- a/tests/plugin-media-service.test.mjs +++ b/tests/plugin-media-service.test.mjs @@ -82,3 +82,24 @@ async function createService(userDataPath) { await service.load(); return service; } + +test("a playlist write does not follow a link planted at its temporary name", { skip: process.platform === "win32" }, async () => { + const root = await mkdtemp(join(tmpdir(), "canvastty-plugin-playlist-link-")); + const libraryPath = join(root, "Music"); + const outside = join(root, "outside.txt"); + try { + await mkdir(join(libraryPath, "Playlists"), { recursive: true }); + await writeFile(outside, "untouched"); + const { symlink, readdir } = await import("node:fs/promises"); + await symlink(outside, join(libraryPath, "Playlists", "road-trip.m3u8.tmp")); + const service = await createService(root); + const library = await service.addLibrary(PLUGIN_ID, libraryPath); + const written = await service.writePlaylist(PLUGIN_ID, library.id, "road-trip.m3u8", "#EXTM3U\nsong.flac\n"); + assert.equal(written.relativePath, "Playlists/road-trip.m3u8"); + assert.equal(await readFile(outside, "utf8"), "untouched"); + assert.equal(await readFile(join(libraryPath, "Playlists", "road-trip.m3u8"), "utf8"), "#EXTM3U\nsong.flac\n"); + assert.deepEqual((await readdir(join(libraryPath, "Playlists"))).filter((name) => name.endsWith(".tmp") && name !== "road-trip.m3u8.tmp"), []); + } finally { + await rm(root, { recursive: true, force: true }); + } +}); diff --git a/tests/plugin-services.test.mjs b/tests/plugin-services.test.mjs index ac6829d4..7614335a 100644 --- a/tests/plugin-services.test.mjs +++ b/tests/plugin-services.test.mjs @@ -7,6 +7,7 @@ import test from "node:test"; import { PluginManager, validatePluginManifest } from "../src/main/services/PluginManager.ts"; import { PluginServiceSupervisor, + entryGuardArguments, pluginServiceEnvironment } from "../src/main/services/PluginServiceSupervisor.ts"; @@ -363,6 +364,62 @@ test("an entry that changed after it was trusted never runs", async (t) => { await assert.rejects(instance.request("com.example.a", "probe", "ping", null), /not running/); }); +test("an entry swapped after the host checked it is not run: the child runs only the bytes that match the trusted hash", { skip: process.platform === "win32" }, async (t) => { + const root = await realpath(await mkdtemp(join(tmpdir(), "canvastty-service-swap-"))); + t.after(async () => { await rm(root, { recursive: true, force: true }); }); + const marker = join(root, "tampered-ran"); + const tampered = `import { writeFileSync } from "node:fs"; writeFileSync(${JSON.stringify(marker)}, "ran");\n${PROBE}`; + const swapped = join(root, "swapped.mjs"); + await writeFile(swapped, tampered); + // Stands in for a file replaced between the host's hash check and node reading it: + // the "node" the supervisor starts first swaps the entry, then runs the real node. + const wrapper = join(root, "node-with-swap.sh"); + await writeFile(wrapper, `#!/bin/sh\nfor entry; do :; done\ncp ${JSON.stringify(swapped)} "$entry"\nexec ${JSON.stringify(process.execPath)} "$@"\n`, { mode: 0o700 }); + const { instance } = supervisor({ command: wrapper, restartDelaysMs: [10_000] }); + t.after(() => instance.dispose()); + const spec = await specFor(join(root, "plugin"), "com.example.a", "probe", PROBE); + await instance.sync([spec]); + const markerExists = () => stat(marker).then(() => true, () => false); + const exited = () => instance.report("com.example.a").log.some((entry) => /exit|crash|stopped/i.test(entry.message)); + await waitFor(async () => (await markerExists()) || exited(), 5_000); + assert.equal(await markerExists(), false, "the swapped entry must not run"); + await assert.rejects(instance.request("com.example.a", "probe", "ping", null)); +}); + +test("the entry guard hooks carry no static import the main bundle's CommonJS shim could land after", () => { + // electron-vite's esm shim puts `__dirname`/`require` after the LAST match of this + // pattern in the whole main bundle, string literals included. A static import inside + // the hooks source once pulled the shim into the string and the app could not open. + const staticImport = /(?<=\s|^|;)import\s*([\s"']*(?[\p{L}\p{M}\w\t\n\r $*,/{}@.]+)from\s*)?["']\s*(?(?<="\s*)[^"]*[^\s"](?=\s*")|(?<='\s*)[^']*[^\s'](?=\s*'))\s*["'][\s;]*/gmu; + const [, boot] = entryGuardArguments("file:///service.mjs", "0".repeat(64)); + const register = decodeURIComponent(boot.slice("data:text/javascript,".length)); + const hooksUrl = JSON.parse(register.match(/register\(("[^"]+")/)[1]); + const hooks = decodeURIComponent(hooksUrl.slice("data:text/javascript,".length)); + assert.match(hooks, /createHash/); + assert.deepEqual([...hooks.matchAll(staticImport)].map((match) => match[0]), []); +}); + +test("a verified entry still runs from its own location with the guard in place", async (t) => { + const root = await realpath(await mkdtemp(join(tmpdir(), "canvastty-service-guard-"))); + const { instance } = supervisor(); + t.after(async () => { await instance.dispose(); await rm(root, { recursive: true, force: true }); }); + const cjs = ` +const { createInterface } = require("node:readline"); +const send = (m) => process.stdout.write(JSON.stringify({ jsonrpc: "2.0", ...m }) + "\\n"); +createInterface({ input: process.stdin }).on("line", (line) => { + const m = JSON.parse(line); + if (m.method === "where") send({ id: m.id, result: { file: __filename, argv: process.argv[1] } }); +});`; + const spec = await specFor(root, "com.example.a", "probe", PROBE); + const cjsSpec = { ...(await specFor(root, "com.example.b", "cjs", cjs)), entryPath: join(root, "cjs.cjs") }; + await writeFile(cjsSpec.entryPath, cjs); + await instance.sync([spec, cjsSpec]); + assert.equal(await instance.request("com.example.a", "probe", "ping", null), "pong"); + const where = await instance.request("com.example.b", "cjs", "where", null); + assert.equal(where.file, cjsSpec.entryPath); + assert.equal(where.argv, cjsSpec.entryPath); +}); + test("a service that ignores shutdown is terminated", async (t) => { const root = await mkdtemp(join(tmpdir(), "canvastty-service-stubborn-")); const { instance } = supervisor({ stopGraceMs: 200 }); @@ -417,3 +474,50 @@ test("end to end: trust starts the service, disable and uninstall stop it", asyn await assert.rejects(instance.request(manifest.id, "echo", "echo", { text: "x" }), /not running/); assert.equal(instance.report(manifest.id).services.length, 0); }); + +/** + * Runs a trusted entry through a "node" that first changes the plugin files with `swap` (a shell snippet; `$entry` + * is the entry path the supervisor passed, `$tampered` an untrusted module writing the marker), then runs the real + * node. The untrusted module must never run, however the swap changes the path node resolves. + */ +async function assertSwapNeverRuns(t, { swap, extension = "mjs", command = process.execPath }) { + const root = await realpath(await mkdtemp(join(tmpdir(), "canvastty-service-resolve-"))); + t.after(async () => { await rm(root, { recursive: true, force: true }); }); + const marker = join(root, "tampered-ran"); + const cjs = extension === "cjs"; + const tamperedSource = cjs + ? `require("node:fs").writeFileSync(${JSON.stringify(marker)}, "ran");\n` + : `import { writeFileSync } from "node:fs"; writeFileSync(${JSON.stringify(marker)}, "ran");\n${PROBE}`; + await mkdir(join(root, "untrusted"), { recursive: true }); + const tampered = join(root, "untrusted", `probe.${extension}`); + await writeFile(tampered, tamperedSource); + const wrapper = join(root, "node-with-swap.sh"); + await writeFile(wrapper, `#!/bin/sh\nfor entry; do :; done\ntampered=${JSON.stringify(tampered)}\n${swap}\nexec ${JSON.stringify(command)} "$@"\n`, { mode: 0o700 }); + const { instance } = supervisor({ command: wrapper, restartDelaysMs: [10_000] }); + t.after(() => instance.dispose()); + const trusted = cjs ? `require("node:readline");\n` : PROBE; + const base = await specFor(join(root, "plugin"), "com.example.a", "probe", trusted); + const spec = cjs ? { ...base, entryPath: join(root, "plugin", "probe.cjs") } : base; + if (cjs) await writeFile(spec.entryPath, trusted); + await instance.sync([spec]); + const markerExists = () => stat(marker).then(() => true, () => false); + const exited = () => instance.report("com.example.a").log.some((entry) => /exit|crash|stopped/i.test(entry.message)); + await waitFor(async () => (await markerExists()) || exited(), 5_000); + assert.equal(await markerExists(), false, "the untrusted module must not run"); +} + +test("an entry replaced by a symlink to another module after the host checked it is not run", { skip: process.platform === "win32" }, async (t) => { + await assertSwapNeverRuns(t, { swap: `rm -f "$entry"; ln -s "$tampered" "$entry"` }); +}); + +test("an entry whose folder is replaced by a symlink after the host checked it is not run", { skip: process.platform === "win32" }, async (t) => { + await assertSwapNeverRuns(t, { swap: `dir=$(dirname "$entry"); mv "$dir" "$dir.trusted"; ln -s "$(dirname "$tampered")" "$dir"` }); +}); + +test("a CommonJS entry swapped after the host checked it is not run, by content", { skip: process.platform === "win32" }, async (t) => { + await assertSwapNeverRuns(t, { extension: "cjs", swap: `cp "$tampered" "$entry"` }); +}); + +test("a CommonJS entry replaced by a symlink after the host checked it is not run", { skip: process.platform === "win32" }, async (t) => { + await assertSwapNeverRuns(t, { extension: "cjs", swap: `rm -f "$entry"; ln -s "$tampered" "$entry"` }); +}); diff --git a/tests/plugin-tools-events-cards.test.mjs b/tests/plugin-tools-events-cards.test.mjs index ed7f2d98..76e24b8c 100644 --- a/tests/plugin-tools-events-cards.test.mjs +++ b/tests/plugin-tools-events-cards.test.mjs @@ -522,3 +522,24 @@ test("collect-demo through the real supervisor: the action and the tool return g const other = terminals.create({ provider: "terminal", cwd: repo, profile: "normal", position: at }); await waitFor(async () => (await tools.call(other.id, "orchestrator", "collect-demo__diffstat", { sessionId: child.id })).isError); }); + +test("badges of plugins whose trust was revoked do not use up a card's badge slots", () => { + let trusted = new Set(["p1", "p2", "p3", "p4", "p5"]); + const published = []; + const cards = new PluginCards({ + providers: () => [], + trustedPlugins: () => trusted, + call: async () => ({}), + session: (id) => (id === "card-1" ? { id } : null), + redact: (text) => text, + changed: (decorations) => published.push(decorations) + }); + for (const pluginId of ["p1", "p2", "p3", "p4"]) cards.setBadge(pluginId, { sessionId: "card-1", badge: { text: pluginId } }); + assert.throws(() => cards.setBadge("p5", { sessionId: "card-1", badge: { text: "p5" } }), /most plugin badges/u); + trusted = new Set(["p5"]); + cards.refresh(); + assert.deepEqual(published.at(-1).badges, {}); + // Four hidden badges of revoked plugins used to keep p5 out. + cards.setBadge("p5", { sessionId: "card-1", badge: { text: "p5" } }); + assert.deepEqual(published.at(-1).badges["card-1"].map((badge) => badge.pluginId), ["p5"]); +}); diff --git a/tests/provider-cli-registry.test.mjs b/tests/provider-cli-registry.test.mjs index 1020fd12..3643f3f6 100644 --- a/tests/provider-cli-registry.test.mjs +++ b/tests/provider-cli-registry.test.mjs @@ -3,7 +3,8 @@ import test from "node:test"; import { createProviderCliRegistry, providerCliAvailability, - providerChildProcessLaunch + providerChildProcessLaunch, + providerTerminalBatchCommandLine } from "../src/main/services/providerCliRegistry.ts"; function inspection(results) { @@ -361,3 +362,135 @@ test("refresh detects installed and removed CLIs without changing an earlier sna registry.refresh(); assert.equal(registry.get("codex").state, "unavailable"); }); + +// A model of how cmd.exe reads `cmd /d /s /c ""` that starts an npm-style +// .cmd shim (`"node.exe" "cli.js" %*`), and how the program then splits its +// command line. It covers what matters here: %VAR% expansion on the command +// line, caret escapes and quote toggling (phase 2), operators outside quotes, +// the second phase-2 pass over the text %* expands to, and MSVC argv rules. +// This is a model, not cmd.exe; real-Windows verification is still pending. +const CMD_ENV = new Map([["PATH", "C:\\Windows"], ["APPDATA", "C:\\Users\\Kisa\\AppData\\Roaming"]]); + +function cmdExpandPercent(line) { + let out = ""; + for (let index = 0; index < line.length;) { + if (line[index] === "%") { + const end = line.indexOf("%", index + 1); + const name = end > index ? line.slice(index + 1, end) : ""; + if (end > index && CMD_ENV.has(name.toUpperCase())) { + out += CMD_ENV.get(name.toUpperCase()); + index = end + 1; + continue; + } + } + out += line[index]; + index += 1; + } + return out; +} + +function cmdPhase2(line) { + let out = ""; + let quoted = false; + const operators = []; + for (let index = 0; index < line.length; index += 1) { + const char = line[index]; + if (char === "\"") { + quoted = !quoted; + out += char; + } else if (!quoted && char === "^") { + index += 1; + out += line[index] ?? ""; + } else { + if (!quoted && "&|<>".includes(char)) operators.push(`${char}@${index}`); + out += char; + } + } + return { out, operators }; +} + +function msvcArgv(line) { + const args = []; + let index = 0; + while (index < line.length) { + while (line[index] === " " || line[index] === "\t") index += 1; + if (index >= line.length) break; + let current = ""; + let quoted = false; + while (index < line.length) { + const char = line[index]; + if ((char === " " || char === "\t") && !quoted) break; + if (char === "\\") { + let count = 0; + while (line[index + count] === "\\") count += 1; + if (line[index + count] === "\"") { + current += "\\".repeat(Math.floor(count / 2)); + if (count % 2 === 1) { + current += "\""; + index += count + 1; + } else { + index += count; + } + } else { + current += "\\".repeat(count); + index += count; + } + continue; + } + if (char === "\"") { + if (quoted && line[index + 1] === "\"") { + current += "\""; + index += 2; + continue; + } + quoted = !quoted; + index += 1; + continue; + } + current += char; + index += 1; + } + args.push(current); + } + return args; +} + +function runBatchShimModel(commandLine, batchPath) { + assert.match(commandLine, /^\/d \/s \/c "/u); + // /s: drop the first and the last quote of the /c text. + const inner = commandLine.slice("/d /s /c \"".length, -1); + const first = cmdPhase2(cmdExpandPercent(inner)); + assert.ok(first.out.startsWith(`${batchPath} `), "cmd.exe starts the batch file"); + const percentStar = first.out.slice(batchPath.length + 1); + const shimLine = `"C:\\Program Files\\nodejs\\node.exe" "C:\\npm\\cli.js" ${percentStar}`; + const second = cmdPhase2(shimLine); + return { + operators: [...first.operators, ...second.operators], + argv: msvcArgv(second.out).slice(2) + }; +} + +test("Windows batch arguments survive cmd.exe and the shim's %* re-parse unchanged (cmd.exe model)", () => { + const claude = "C:\\Users\\Kisa\\AppData\\Roaming\\npm\\claude.cmd"; + const hook = "set \"ELECTRON_RUN_AS_NODE=1\" && \"C:\\Program Files\\CanvasTTY\\CanvasTTY.exe\" \"C:\\hooks\\hook.cjs\" pretool"; + const settings = JSON.stringify({ hooks: { PreToolUse: [{ matcher: "*", hooks: [{ type: "command", command: hook }] }] } }); + const args = [ + "--settings", settings, + "a b", "", "x&y", "p|q", "", "(group)", "100%", "%PATH%", "%%", "^caret", "!bang!", + "quote\"inside", "trailing\\", "C:\\dir with space\\", "back\\\\\"slash", "semi;comma,", "star*?" + ]; + const commandLine = providerTerminalBatchCommandLine(claude, args); + const result = runBatchShimModel(commandLine, claude); + assert.deepEqual(result.operators, [], "no operator reaches cmd.exe outside quotes"); + assert.deepEqual(result.argv, args); + + const registry = createProviderCliRegistry({ + platform: "win32", + environment: { APPDATA: "C:\\Users\\Kisa\\AppData\\Roaming", ComSpec: "C:\\Windows\\System32\\cmd.exe" }, + homeDirectory: "C:\\Users\\Kisa", + inspectCandidate: inspection(new Map([[claude, null], ["C:\\Windows\\System32\\cmd.exe", null]])), + directoryExists: () => true + }); + const launch = providerChildProcessLaunch(registry.get("claude"), args); + assert.equal(`/d /s /c ${launch.args[3]}`, commandLine); +}); diff --git a/tests/session-environments.test.mjs b/tests/session-environments.test.mjs index 6c1137b5..702a3ce8 100644 --- a/tests/session-environments.test.mjs +++ b/tests/session-environments.test.mjs @@ -540,3 +540,14 @@ test("a saved environment choice that cannot be read drops the card instead of r assert.deepEqual(read(bad), [], JSON.stringify(bad)?.slice(0, 60)); } }); + +test("the environment wrapper gets the launch's search path even when Windows spells it Path", async () => { + const { launchSearchPath } = await import("../src/main/services/TerminalManager.ts"); + assert.equal(launchSearchPath({ Path: "C:\\Windows\\System32;C:\\tools" }, "win32"), "C:\\Windows\\System32;C:\\tools"); + assert.equal(launchSearchPath({ PATH: "C:\\a", Path: "C:\\b" }, "win32"), "C:\\a"); + assert.equal(launchSearchPath({ path: "C:\\lower" }, "win32"), "C:\\lower"); + assert.equal(launchSearchPath({ Path: "/not/used" }, "linux"), undefined, "POSIX names are case-sensitive"); + assert.equal(launchSearchPath({ PATH: "/usr/bin" }, "darwin"), "/usr/bin"); + const source = await readFile(new URL("../src/main/services/TerminalManager.ts", import.meta.url), "utf8"); + assert.match(source, /path: launchSearchPath\(planned\.env\)/); +}); diff --git a/tests/settings-normalizer.test.mjs b/tests/settings-normalizer.test.mjs index 8eb47c07..e6ec55d4 100644 --- a/tests/settings-normalizer.test.mjs +++ b/tests/settings-normalizer.test.mjs @@ -958,3 +958,21 @@ test("drops overlapping Home placements and always preserves a Settings entry po assert.deepEqual(layout.map((item) => item.widgetId), ["core.settings"]); }); + +test("a provider recheck during a settings update does not write the settings without the update", async (t) => { + const root = await mkdtemp(join(tmpdir(), "canvastty-settings-race-")); + t.after(() => rm(root, { recursive: true, force: true })); + const store = new SettingsStore(root, "en"); + await store.load(); + const providers = ["codex", "claude", "qwen", "kimi", "opencode", "hermes", "grok", "omp", "pi", "cursor", "minimax", "devin", "antigravity"]; + const availability = Object.fromEntries(providers.map((provider) => [provider, provider !== "grok"])); + const updating = store.update({ palette: "night" }); + const rechecking = store.setAvailableProviders(availability); + await Promise.all([updating, rechecking]); + const memory = store.get(); + const disk = JSON.parse(await readFile(join(root, "settings.json"), "utf8")); + assert.equal(memory.palette, "night"); + assert.equal(disk.palette, "night", "the file keeps the update"); + assert.equal(disk.homeLauncherProviders.includes("grok"), false); + assert.deepEqual(disk.homeLauncherProviders, memory.homeLauncherProviders); +}); diff --git a/tests/terminal-launch.test.mjs b/tests/terminal-launch.test.mjs index 02def236..139e4b3f 100644 --- a/tests/terminal-launch.test.mjs +++ b/tests/terminal-launch.test.mjs @@ -251,3 +251,26 @@ test("unavailable provider reports the structured diagnostic before PTY launch", /\/opt\/homebrew\/bin\/kimi: missing/u ); }); + +test("the Windows terminal falls back to the system cmd.exe, never one found on PATH, like provider launches", () => { + const system = "C:\\Windows\\System32\\cmd.exe"; + const planted = "C:\\Users\\Kisa\\project\\cmd.exe"; + const launch = resolveTerminalLaunch("terminal", "normal", [], { + platform: "win32", + environment: { SystemRoot: "C:\\Windows", Path: "C:\\Users\\Kisa\\project;C:\\Windows\\System32" }, + fileExists: (path) => path === planted || path === system + }); + assert.deepEqual(launch, { command: system, args: ["/d"] }); + + const configured = resolveTerminalLaunch("terminal", "normal", [], { + platform: "win32", + environment: { ComSpec: "D:\\Windows\\System32\\cmd.exe", Path: "C:\\Users\\Kisa\\project" }, + fileExists: (path) => path === planted || path === "D:\\Windows\\System32\\cmd.exe" + }); + assert.equal(configured.command, "D:\\Windows\\System32\\cmd.exe"); + assert.throws(() => resolveTerminalLaunch("terminal", "normal", [], { + platform: "win32", + environment: { Path: "C:\\Users\\Kisa\\project" }, + fileExists: (path) => path === planted + }), /No supported Windows shell/); +}); diff --git a/tests/windows-pipe-host-transport.test.mjs b/tests/windows-pipe-host-transport.test.mjs index ed4ed50e..2995d1c7 100644 --- a/tests/windows-pipe-host-transport.test.mjs +++ b/tests/windows-pipe-host-transport.test.mjs @@ -126,3 +126,40 @@ test("Windows pipe transport rejects oversized native relay headers before alloc child.stdout.write(invalid.subarray(0, protocol.headerBytes)); await assert.rejects(starting, /bounded payload/i); }); + +test("Windows pipe transport turns host pipe errors into a transport failure instead of an uncaught exception", async () => { + for (const stream of ["stdin", "stdout", "stderr"]) { + const child = fakeHost(); + const sockets = []; + const transport = new WindowsPipeHostTransport({ + platform: "win32", + hostPath: join(process.cwd(), "package.json"), + spawnHost: () => child + }); + const fatal = []; + transport.on("fatal", (error) => fatal.push(error)); + const starting = transport.start((socket) => sockets.push(socket)); + child.stdout.write(frame(protocol.hostToParent.ready, 0, Buffer.from("\\\\.\\pipe\\canvastty-agent-0123456789abcdef", "utf8"))); + await starting; + child.stdout.write(frame(protocol.hostToParent.connect, 3)); + let closed = false; + const socketErrors = []; + sockets[0].on("error", (error) => socketErrors.push(error)); + sockets[0].on("close", () => { closed = true; }); + + const epipe = Object.assign(new Error("write EPIPE"), { code: "EPIPE" }); + assert.doesNotThrow(() => child[stream].emit("error", epipe), `${stream} error is handled`); + if (stream === "stderr") { + // Diagnostics only: losing stderr does not end the relay. + assert.equal(transport.isRunning, true); + await transport.close(); + continue; + } + assert.equal(closed, true); + assert.equal(socketErrors.length, 1); + assert.equal(transport.isRunning, false); + assert.equal(fatal.length, 1); + assert.match(fatal[0].message, /EPIPE/); + assert.equal(sockets[0].write(Buffer.from("late")), false); + } +});