Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions src/browser/daemon-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,36 @@ export async function releaseSiteSessionLease(params: {
}
}

/**
* Best-effort pacing outcome report on adapter command completion. Feeds the
* daemon's per-site security-block circuit breaker (see site-pacing.ts); a
* lost report only costs the breaker one sample, so this must never block or
* fail the caller.
*/
export async function reportPacingOutcome(params: {
session: string;
outcome: 'ok' | 'security_block';
contextId?: string;
}): Promise<void> {
try {
await requestDaemon('/command', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
id: generateId(),
action: 'pacing-report',
session: params.session,
surface: 'adapter',
outcome: params.outcome,
...(params.contextId ? { contextId: params.contextId } : {}),
}),
timeout: 2000,
});
} catch {
// Best-effort: the breaker just misses one sample.
}
}

/**
* Transport-level deadlines share one source of truth: `body.timeout` (seconds).
* The daemon arms its per-command timer from it, the extension derives its CDP
Expand Down
64 changes: 64 additions & 0 deletions src/daemon.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,12 @@ import {
getSessionLeaseKey,
isSessionLeaseCommand,
} from './session-lease.js';
import {
SitePacer,
buildSecurityCooldownFailure,
classifyPacedNavigation,
parseSiteFromSession,
} from './site-pacing.js';

const PORT = DEFAULT_DAEMON_PORT;
if (!isIgnorableDaemonPortEnv(process.env.OPENCLI_DAEMON_PORT)) {
Expand Down Expand Up @@ -92,6 +98,11 @@ const pending = new Map<string, PendingEntry>();
// session-lease.ts).
const sessionLeases = new SessionLeaseRegistry();

// Per-site navigation pacing + security-block circuit breaker (site-pacing.ts).
// Kill switch is read once at daemon startup: OPENCLI_PACING=off disables it.
const sitePacer = new SitePacer();
const sitePacingEnabled = process.env.OPENCLI_PACING !== 'off';

/** A TTL-stale lease holder with a command still in flight is alive, not dead. */
function runHasPendingWork(runId: string): boolean {
for (const entry of pending.values()) {
Expand Down Expand Up @@ -365,6 +376,27 @@ async function handleRequest(req: IncomingMessage, res: ServerResponse): Promise
return;
}

// ─── Site pacing: outcome report ─────────────────────────────────
// Daemon-local, never dispatched to the extension. Feeds the per-site
// security-block circuit breaker; unpaced sites are ignored inside
// reportOutcome. Best-effort contextId resolution — a report for a
// profile that just disconnected is simply dropped.
if (body.action === 'pacing-report') {
const site = parseSiteFromSession(body.session);
const outcome = body.outcome === 'security_block' || body.outcome === 'ok' ? body.outcome : null;
if (site && outcome) {
const reportRoute = resolveExtensionConnection(
typeof body.contextId === 'string' ? body.contextId : undefined,
undefined,
);
const ctx = reportRoute.connection?.contextId
?? (typeof body.contextId === 'string' && body.contextId ? body.contextId : null);
if (ctx) sitePacer.reportOutcome(ctx, site, outcome, Date.now());
}
jsonResponse(res, 200, { id: body.id, ok: true });
return;
}

const route = resolveExtensionConnection(
typeof body.contextId === 'string' ? body.contextId : undefined,
typeof body.preferredContextId === 'string' ? body.preferredContextId : undefined,
Expand Down Expand Up @@ -419,6 +451,38 @@ async function handleRequest(req: IncomingMessage, res: ServerResponse): Promise
leaseRunId = body.runId;
}

// ─── Site pacing: navigation slots + security cooldown ───────────
// Space adapter navigations for velocity-sensitive sites and fail fast
// while a security cooldown is open (site-pacing.ts). Applies to
// `navigate` only — evaluates against a warm tab load no pages.
if (sitePacingEnabled) {
const pacedSite = classifyPacedNavigation(body);
if (pacedSite) {
const pacing = sitePacer.acquireNavigationSlot(route.connection.contextId, pacedSite, Date.now());
if (!pacing.granted) {
const failure = buildSecurityCooldownFailure(pacedSite, pacing.retryAfterMs);
// The daemon's own stdio is discarded (spawned with stdio: 'ignore');
// the /logs ring buffer is the only operator-visible channel.
pushLog({ level: 'warn', msg: `[pacing] ${pacedSite} navigate refused — security cooldown, retry in ${Math.ceil(pacing.retryAfterMs / 1000)}s`, ts: Date.now() });
log.warn(`[daemon] ${pacedSite} navigate refused — security cooldown, retry in ${Math.ceil(pacing.retryAfterMs / 1000)}s`);
jsonResponse(res, failure.status, {
id: body.id,
ok: false,
errorCode: failure.errorCode,
error: failure.message,
errorHint: failure.errorHint,
retryAfterMs: failure.retryAfterMs,
});
return;
}
if (pacing.delayMs > 0) {
pushLog({ level: 'info', msg: `[pacing] spacing ${pacedSite} navigate by ${Math.round(pacing.delayMs)}ms`, ts: Date.now() });
log.info(`[daemon] pacing ${pacedSite} navigate by ${Math.round(pacing.delayMs)}ms`);
await new Promise((resolve) => setTimeout(resolve, pacing.delayMs));
}
}
}

// Absolute deadline wins over the legacy duration field: all hops share
// one wall clock, so remaining budget absorbs queueing/transit time.
const timeoutMs = typeof body.deadlineAt === 'number' && body.deadlineAt > 0
Expand Down
76 changes: 75 additions & 1 deletion src/execution.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import * as os from 'node:os';
import * as path from 'node:path';
import type { CliCommand } from './registry.js';
import { coerceAndValidateArgs, executeCommand, prepareCommandArgs } from './execution.js';
import { ArgumentError, TimeoutError, toEnvelope } from './errors.js';
import { ArgumentError, CliError, TimeoutError, toEnvelope } from './errors.js';
import { cli, Strategy } from './registry.js';
import { withTimeoutMs } from './runtime.js';
import * as runtime from './runtime.js';
Expand Down Expand Up @@ -966,3 +966,77 @@ describe('executeCommand — persistent write lease release', () => {
vi.restoreAllMocks();
});
});

describe('executeCommand — pacing outcome report', () => {
function pacedReadCmd(name: string, func: () => Promise<unknown>): CliCommand {
return cli({
site: 'test-pacing-report',
name,
access: 'read',
description: 'test pacing outcome report',
browser: true,
strategy: Strategy.PUBLIC,
func,
});
}

it('reports ok after a successful adapter browser command', async () => {
vi.spyOn(capRouting, 'shouldUseBrowserSession').mockReturnValue(true);
vi.spyOn(runtime, 'browserSession').mockImplementation(async (_Factory, fn) => fn({} as any));
const reportSpy = vi.spyOn(daemonClient, 'reportPacingOutcome').mockResolvedValue(undefined);

await executeCommand(pacedReadCmd('pace-ok', async () => [{ ok: true }]), {});

expect(reportSpy).toHaveBeenCalledTimes(1);
expect(reportSpy).toHaveBeenCalledWith(expect.objectContaining({
session: expect.stringMatching(/^site:test-pacing-report/),
outcome: 'ok',
}));
vi.restoreAllMocks();
});

it('reports security_block when the command fails with a SECURITY_BLOCK error', async () => {
vi.spyOn(capRouting, 'shouldUseBrowserSession').mockReturnValue(true);
vi.spyOn(runtime, 'browserSession').mockImplementation(async (_Factory, fn) => fn({} as any));
const reportSpy = vi.spyOn(daemonClient, 'reportPacingOutcome').mockResolvedValue(undefined);

const cmd = pacedReadCmd('pace-block', async () => {
throw new CliError('SECURITY_BLOCK', 'risk control blocked the page');
});
await expect(executeCommand(cmd, {})).rejects.toThrow('risk control');

expect(reportSpy).toHaveBeenCalledTimes(1);
expect(reportSpy).toHaveBeenCalledWith(expect.objectContaining({ outcome: 'security_block' }));
vi.restoreAllMocks();
});

it('reports nothing for other failures — unrelated errors must not reset the breaker', async () => {
vi.spyOn(capRouting, 'shouldUseBrowserSession').mockReturnValue(true);
vi.spyOn(runtime, 'browserSession').mockImplementation(async (_Factory, fn) => fn({} as any));
const reportSpy = vi.spyOn(daemonClient, 'reportPacingOutcome').mockResolvedValue(undefined);

const cmd = pacedReadCmd('pace-other-failure', async () => {
throw new BrowserCommandError('boom', 'attach_failed');
});
await expect(executeCommand(cmd, {})).rejects.toThrow('boom');

expect(reportSpy).not.toHaveBeenCalled();
vi.restoreAllMocks();
});

it('reports nothing for non-browser commands', async () => {
const reportSpy = vi.spyOn(daemonClient, 'reportPacingOutcome').mockResolvedValue(undefined);
const cmd = cli({
site: 'test-pacing-report',
name: 'pace-non-browser',
access: 'read',
description: 'non-browser command',
browser: false,
strategy: Strategy.PUBLIC,
func: async () => [{ ok: true }],
});
await executeCommand(cmd, {});
expect(reportSpy).not.toHaveBeenCalled();
vi.restoreAllMocks();
});
});
13 changes: 11 additions & 2 deletions src/execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,11 +26,11 @@ import * as crypto from 'node:crypto';
import * as fs from 'node:fs';
import * as os from 'node:os';
import { executePipeline } from './pipeline/index.js';
import { adapterLoadError, ArgumentError, CommandExecutionError, SessionBusyError, attachTraceReceipt, getErrorMessage } from './errors.js';
import { adapterLoadError, ArgumentError, CliError, CommandExecutionError, SessionBusyError, attachTraceReceipt, getErrorMessage } from './errors.js';
import { shouldUseBrowserSession } from './capabilityRouting.js';
import { getBrowserFactory, browserSession, runWithTimeout, DEFAULT_BROWSER_COMMAND_TIMEOUT, type BrowserWindowMode } from './runtime.js';
import { profileRouteParams, resolveProfileSelection } from './browser/profile.js';
import { clearDaemonRunContext, generateRunId, isUnknownOutcomeError, releaseSiteSessionLease, setDaemonCommandTimeoutSeconds, setDaemonRunContext } from './browser/daemon-client.js';
import { clearDaemonRunContext, generateRunId, isUnknownOutcomeError, releaseSiteSessionLease, reportPacingOutcome, setDaemonCommandTimeoutSeconds, setDaemonRunContext } from './browser/daemon-client.js';
import { emitHook, type HookContext } from './hooks.js';
import { log } from './logger.js';
import { isElectronApp } from './electron-apps.js';
Expand Down Expand Up @@ -443,6 +443,15 @@ export async function executeCommand(
// result-evicted, anywhere in the cause chain) means the browser-side
// command may STILL be running against the persistent tab; there is
// nothing to await client-side, so the TTL is the quiet period.
// Pacing outcome report (best-effort, never awaited): the daemon's
// per-site circuit breaker counts SECURITY_BLOCK endings and resets on
// success. Other failures say nothing about the site's risk state, and
// an unknown-outcome ending reports nothing either.
if (browserRunError === undefined) {
void reportPacingOutcome({ session, outcome: 'ok', ...(contextId ? { contextId } : {}) });
} else if (browserRunError instanceof CliError && browserRunError.code === 'SECURITY_BLOCK') {
void reportPacingOutcome({ session, outcome: 'security_block', ...(contextId ? { contextId } : {}) });
}
if (leaseRun) {
if (adapterStillRunning && adapterRun) {
const runId = leaseRun.runId;
Expand Down
Loading
Loading