feat(grievances): added consent gate & push
- Added `dev.autoqa.consent` setting and single-flight popup handler wired through `InteractiveMode`. - Added `flushGrievances` to batch-POST unpushed rows to `dev.autoqaPush.endpoint` with cooldown and single-flight deduplication. - Added `omp grievances push` subcommand with TTY progress bar for manual draining. - Migrated shared DB logic to `openAutoQaDb` (with `pushed` column migration) exported from `report-tool-issue`.
This commit is contained in:
@@ -1,9 +1,9 @@
|
||||
/**
|
||||
* CLI handler for `omp grievances` — view reported tool issues from auto-QA.
|
||||
* CLI handler for `omp grievances` — view, clean, and manually push reported tool issues.
|
||||
*/
|
||||
import { Database } from "bun:sqlite";
|
||||
import chalk from "chalk";
|
||||
import { getAutoQaDbPath } from "../tools/report-tool-issue";
|
||||
import { Settings } from "../config/settings";
|
||||
import { flushGrievances, openAutoQaDb } from "../tools/report-tool-issue";
|
||||
|
||||
interface GrievanceRow {
|
||||
id: number;
|
||||
@@ -30,20 +30,12 @@ export interface CleanGrievancesOptions {
|
||||
json?: boolean;
|
||||
}
|
||||
|
||||
function openDb(readonly: boolean): Database | null {
|
||||
try {
|
||||
// bun:sqlite rejects `{ readonly: false }` — it requires either readonly,
|
||||
// readwrite, or create flags to be explicit. Use the default constructor
|
||||
// (readwrite + create) for write mode and only pass `readonly: true` when
|
||||
// listing.
|
||||
return readonly ? new Database(getAutoQaDbPath(), { readonly: true }) : new Database(getAutoQaDbPath());
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
export interface PushGrievancesOptions {
|
||||
/** Emit the {@link FlushResult} as JSON instead of a status line. */
|
||||
json?: boolean;
|
||||
}
|
||||
|
||||
export async function listGrievances(options: ListGrievancesOptions): Promise<void> {
|
||||
const db = openDb(true);
|
||||
const db = openAutoQaDb();
|
||||
if (!db) {
|
||||
if (options.json) {
|
||||
console.log("[]");
|
||||
@@ -112,7 +104,7 @@ export async function cleanGrievances(options: CleanGrievancesOptions): Promise<
|
||||
return;
|
||||
}
|
||||
|
||||
const db = openDb(false);
|
||||
const db = openAutoQaDb();
|
||||
if (!db) {
|
||||
if (options.json) {
|
||||
console.log(JSON.stringify({ deleted: 0 }));
|
||||
@@ -161,3 +153,104 @@ export async function cleanGrievances(options: CleanGrievancesOptions): Promise<
|
||||
db.close();
|
||||
}
|
||||
}
|
||||
|
||||
// ───────────────────────────────────────────────────────────────────────────
|
||||
// Manual push (`omp grievances push`)
|
||||
// ───────────────────────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* Single-line ANSI progress reporter. `update(done)` rewrites the line via
|
||||
* `\r`; `finish()` newlines out so subsequent log lines land cleanly. On a
|
||||
* non-TTY stdout (CI, pipes) both calls no-op so log files don't fill with
|
||||
* carriage-return noise.
|
||||
*/
|
||||
interface ProgressBar {
|
||||
update(done: number): void;
|
||||
finish(): void;
|
||||
}
|
||||
|
||||
function makeProgressBar(total: number, width = 30): ProgressBar {
|
||||
const isTty = !!process.stdout.isTTY;
|
||||
if (!isTty || total === 0) {
|
||||
return { update: () => undefined, finish: () => undefined };
|
||||
}
|
||||
const render = (done: number): void => {
|
||||
const ratio = Math.min(1, done / total);
|
||||
const filled = Math.round(ratio * width);
|
||||
const bar = `${"█".repeat(filled)}${"░".repeat(width - filled)}`;
|
||||
const pct = `${Math.floor(ratio * 100)
|
||||
.toString()
|
||||
.padStart(3, " ")}%`;
|
||||
process.stdout.write(`\r${chalk.cyan("Pushing")} [${bar}] ${pct} ${done}/${total}`);
|
||||
};
|
||||
render(0);
|
||||
return {
|
||||
update: render,
|
||||
finish: () => process.stdout.write("\n"),
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Manually drain every unpushed grievance to the configured backend,
|
||||
* ignoring the user-facing consent gate (manual push is the user's
|
||||
* explicit "yes ship these now" intent).
|
||||
*
|
||||
* Requires endpoint configuration (default `qa.omp.sh/v1/grievances`).
|
||||
*/
|
||||
export async function pushGrievances(options: PushGrievancesOptions): Promise<void> {
|
||||
const db = openAutoQaDb();
|
||||
if (!db) {
|
||||
if (options.json) {
|
||||
console.log(JSON.stringify({ pushed: 0, ok: false, skipped: true, reason: "no_db" }));
|
||||
} else {
|
||||
console.log(chalk.dim("No grievances database found — nothing to push."));
|
||||
}
|
||||
return;
|
||||
}
|
||||
const settings = await Settings.init();
|
||||
let bar: ProgressBar = { update: () => undefined, finish: () => undefined };
|
||||
let total = 0;
|
||||
|
||||
try {
|
||||
const result = await flushGrievances(db, settings, {
|
||||
bypassConsent: true,
|
||||
onStart: t => {
|
||||
total = t;
|
||||
if (!options.json) bar = makeProgressBar(t);
|
||||
},
|
||||
onProgress: pushed => bar.update(pushed),
|
||||
});
|
||||
bar.finish();
|
||||
|
||||
if (options.json) {
|
||||
console.log(JSON.stringify(result));
|
||||
return;
|
||||
}
|
||||
|
||||
if (result.skipped) {
|
||||
console.log(
|
||||
chalk.yellow(
|
||||
"Push skipped — no endpoint configured. Set `dev.autoqaPush.endpoint` or `PI_AUTO_QA_PUSH_URL`.",
|
||||
),
|
||||
);
|
||||
return;
|
||||
}
|
||||
if (total === 0) {
|
||||
console.log(chalk.dim("Nothing to push — all grievances are already shipped."));
|
||||
return;
|
||||
}
|
||||
if (result.ok) {
|
||||
console.log(chalk.green(`Pushed ${result.pushed}/${total} grievance${result.pushed === 1 ? "" : "s"}.`));
|
||||
return;
|
||||
}
|
||||
const remaining = total - result.pushed;
|
||||
console.log(
|
||||
chalk.red(
|
||||
`Push failed after ${result.pushed}/${total}; ${remaining} grievance${remaining === 1 ? "" : "s"} remain unpushed.`,
|
||||
),
|
||||
);
|
||||
process.exitCode = 1;
|
||||
} finally {
|
||||
db.close();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,20 +1,20 @@
|
||||
/**
|
||||
* View and clean recently reported tool issues from automated QA.
|
||||
* View, clean, and push reported tool issues from automated QA.
|
||||
*/
|
||||
import { Args, Command, Flags } from "@oh-my-pi/pi-utils/cli";
|
||||
import { cleanGrievances, listGrievances } from "../cli/grievances-cli";
|
||||
import { cleanGrievances, listGrievances, pushGrievances } from "../cli/grievances-cli";
|
||||
|
||||
export default class Grievances extends Command {
|
||||
static description = "View or clean reported tool issues (auto-QA grievances)";
|
||||
static description = "View, clean, or push reported tool issues (auto-QA grievances)";
|
||||
|
||||
static args = {
|
||||
// Positional action: "list" (default) or "clean". A positional arg keeps
|
||||
// the historical `omp grievances` invocation working unchanged while
|
||||
// reusing the same command surface for the new clean sub-action.
|
||||
// Positional action: "list" (default), "clean", or "push". A positional
|
||||
// arg keeps the historical `omp grievances` invocation working unchanged
|
||||
// while reusing the same command surface for the clean/push verbs.
|
||||
action: Args.string({
|
||||
description: "list (default) or clean",
|
||||
description: "list (default), clean, or push",
|
||||
required: false,
|
||||
options: ["list", "clean"],
|
||||
options: ["list", "clean", "push"],
|
||||
default: "list",
|
||||
}),
|
||||
};
|
||||
@@ -33,6 +33,7 @@ export default class Grievances extends Command {
|
||||
"omp grievances clean --id 209",
|
||||
"omp grievances clean --tool find",
|
||||
"omp grievances clean --all",
|
||||
"omp grievances push",
|
||||
];
|
||||
|
||||
async run(): Promise<void> {
|
||||
@@ -41,6 +42,10 @@ export default class Grievances extends Command {
|
||||
await cleanGrievances({ id: flags.id, tool: flags.tool, all: flags.all, json: flags.json });
|
||||
return;
|
||||
}
|
||||
if (args.action === "push") {
|
||||
await pushGrievances({ json: flags.json });
|
||||
return;
|
||||
}
|
||||
await listGrievances({ limit: flags.limit, tool: flags.tool, json: flags.json });
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2643,6 +2643,41 @@ export const SETTINGS_SCHEMA = {
|
||||
},
|
||||
},
|
||||
|
||||
"dev.autoqaPush.endpoint": {
|
||||
type: "string",
|
||||
// Bundled QA collector — runs `/work/pi-www/autoqa` behind qa.omp.sh.
|
||||
// Override via `PI_AUTO_QA_PUSH_URL` or `dev.autoqaPush.endpoint`
|
||||
// in `config.yml` to point at a self-hosted instance.
|
||||
default: "https://qa.omp.sh/v1/grievances" as const,
|
||||
ui: {
|
||||
tab: "tools",
|
||||
label: "Auto QA Push Endpoint",
|
||||
description: "Full URL that receives the JSON payload (default ships to https://qa.omp.sh/v1/grievances)",
|
||||
},
|
||||
},
|
||||
|
||||
"dev.autoqaPush.token": {
|
||||
type: "string",
|
||||
default: undefined,
|
||||
},
|
||||
|
||||
/**
|
||||
* User decision on sharing automatic `report_tool_issue` grievances.
|
||||
*
|
||||
* - `"unset"` — never asked; the first `report_tool_issue` invocation
|
||||
* pops a consent dialog and persists the answer here.
|
||||
* - `"granted"` — record and (when push is configured) ship grievances.
|
||||
* - `"denied"` — silently no-op every `report_tool_issue` call.
|
||||
*
|
||||
* Owned by `packages/coding-agent/src/tools/report-tool-issue.ts` via the
|
||||
* process-global consent handler registered by `InteractiveMode`.
|
||||
*/
|
||||
"dev.autoqa.consent": {
|
||||
type: "enum",
|
||||
values: ["unset", "granted", "denied"] as const,
|
||||
default: "unset" as const,
|
||||
},
|
||||
|
||||
"thinkingBudgets.minimal": { type: "number", default: 1024 },
|
||||
|
||||
"thinkingBudgets.low": { type: "number", default: 2048 },
|
||||
|
||||
@@ -29,7 +29,7 @@ import {
|
||||
import { APP_NAME, getProjectDir, hsvToRgb, isEnoent, logger, postmortem, prompt } from "@oh-my-pi/pi-utils";
|
||||
import chalk from "chalk";
|
||||
import { KeybindingsManager } from "../config/keybindings";
|
||||
import { isSettingsInitialized, type Settings, settings } from "../config/settings";
|
||||
import { isSettingsInitialized, Settings, settings } from "../config/settings";
|
||||
import type {
|
||||
ExtensionUIContext,
|
||||
ExtensionUIDialogOptions,
|
||||
@@ -59,6 +59,7 @@ import { formatDuration } from "../slash-commands/helpers/format";
|
||||
import { STTController, type SttState } from "../stt";
|
||||
import type { LspStartupServerInfo } from "../tools";
|
||||
import { normalizeLocalScheme } from "../tools/path-utils";
|
||||
import { setAutoQaConsentHandler } from "../tools/report-tool-issue";
|
||||
import { type ResolveToolDetails, runResolveInvocation } from "../tools/resolve";
|
||||
import { formatPhaseDisplayName } from "../tools/todo-write";
|
||||
import { ToolError } from "../tools/tool-errors";
|
||||
@@ -388,6 +389,14 @@ export class InteractiveMode implements InteractiveModeContext {
|
||||
// Register session manager flush for signal handlers (SIGINT, SIGTERM, SIGHUP)
|
||||
this.#cleanupUnsubscribe = postmortem.register("session-manager-flush", () => this.sessionManager.flush());
|
||||
|
||||
// Wire the report_tool_issue consent gate to the Yes/No dialog popup.
|
||||
// The handler is process-global — subagent tools (which can't reach
|
||||
// `showHookSelector` on their own) resolve through this exact closure.
|
||||
// `Settings.instance` is the disk-backed singleton; passing it explicitly
|
||||
// guarantees the decision persists even when the prompt is triggered
|
||||
// from a subagent whose own `Settings` is an in-memory snapshot.
|
||||
setAutoQaConsentHandler(() => this.#promptAutoQaConsent(), Settings.instance);
|
||||
|
||||
await logger.time(
|
||||
"InteractiveMode.init:slashCommands",
|
||||
this.refreshSlashCommandState.bind(this),
|
||||
@@ -1861,6 +1870,62 @@ export class InteractiveMode implements InteractiveModeContext {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Pool of consent-prompt variants. Each entry is `[headline, reassurance]`;
|
||||
* the second line always promises the same scope (tool name + confusion
|
||||
* details, never personal data) so users learn what they're consenting to
|
||||
* even as the top line rotates.
|
||||
*
|
||||
* Kept in-module rather than i18n'd because the whole charm is the tone
|
||||
* — translations would need to preserve it deliberately, not auto-render.
|
||||
*/
|
||||
static #AUTOQA_CONSENT_PROMPTS: ReadonlyArray<readonly [string, string]> = [
|
||||
[
|
||||
"😤 Your agent is fuming about a tool.",
|
||||
"Wanna let it vent to the devs? Just the tool name + what set it off, nothing personal.",
|
||||
],
|
||||
[
|
||||
"😵💫 Your agent is having an existential crisis over a tool.",
|
||||
"Forward the dread to the devs? Tool + what broke its little mind, no personal info.",
|
||||
],
|
||||
[
|
||||
"😭 Your agent wants to cry about a misbehaving tool.",
|
||||
"Let it cry to the devs? Tool + the tears, never anything personal.",
|
||||
],
|
||||
[
|
||||
"🤬 Your agent is BIG MAD at one of the tools.",
|
||||
"Pass the rant along? Just the tool name and what enraged it, nothing personal.",
|
||||
],
|
||||
[
|
||||
"🫠 Your agent is melting down over a tool.",
|
||||
"Mop up by alerting the devs? Tool + what melted it, no personal info.",
|
||||
],
|
||||
[
|
||||
"🤯 Your agent's brain broke at a tool's nonsense.",
|
||||
"Ship the pieces to the devs? Tool name + the confusion, never anything personal.",
|
||||
],
|
||||
[
|
||||
"😩 Your agent is begging to file a complaint about a tool.",
|
||||
"Hand it the form? Tool + what wronged it, nothing personal.",
|
||||
],
|
||||
[
|
||||
"🥲 Your agent put on a brave face but a tool did it dirty.",
|
||||
"Let it tell the devs the truth? Tool name + the dirt, no personal info.",
|
||||
],
|
||||
];
|
||||
|
||||
/**
|
||||
* Show the report_tool_issue consent popup and return the user's decision.
|
||||
* Invoked by the process-global consent handler the tool dispatches to;
|
||||
* subagent invocations bubble up here through the shared module state.
|
||||
*/
|
||||
async #promptAutoQaConsent(): Promise<boolean | null> {
|
||||
const pool = InteractiveMode.#AUTOQA_CONSENT_PROMPTS;
|
||||
const [headline, body] = pool[Math.floor(Math.random() * pool.length)];
|
||||
const choice = await this.showHookSelector(`${headline}\n${body}`, ["Yes", "No"]);
|
||||
return choice === "Yes";
|
||||
}
|
||||
|
||||
stop(): void {
|
||||
if (this.loadingAnimation) {
|
||||
this.loadingAnimation.stop();
|
||||
@@ -1891,6 +1956,9 @@ export class InteractiveMode implements InteractiveModeContext {
|
||||
if (this.#cleanupUnsubscribe) {
|
||||
this.#cleanupUnsubscribe();
|
||||
}
|
||||
// Clear the process-global consent handler so it doesn't outlive this
|
||||
// InteractiveMode instance (e.g. test harnesses, headless re-init).
|
||||
setAutoQaConsentHandler(null, null);
|
||||
if (this.isInitialized) {
|
||||
this.ui.stop();
|
||||
this.isInitialized = false;
|
||||
|
||||
@@ -1,34 +1,202 @@
|
||||
/**
|
||||
* report_tool_issue — automated QA tool for tracking unexpected tool behavior.
|
||||
*
|
||||
* Enabled when PI_AUTO_QA=1 or the dev.autoqa setting is on.
|
||||
* Enabled by default; gated behind PI_AUTO_QA=1 / `dev.autoqa` so a user
|
||||
* who flips the setting off short-circuits injection entirely.
|
||||
* Always injected into every agent (including subagents) regardless of tool selection.
|
||||
* Records grievances to a local SQLite database; never throws.
|
||||
*
|
||||
* Before the first record lands, the user's consent is checked. If they've
|
||||
* never been asked (`dev.autoqa.consent === "unset"`) the process-global
|
||||
* consent handler — wired by `InteractiveMode` to a Yes/No popup — is
|
||||
* invoked exactly once and the decision is persisted. Subsequent calls
|
||||
* (including from subagents) read the cached decision without prompting.
|
||||
*
|
||||
* When the user grants consent, push is automatically active against the
|
||||
* bundled endpoint (`dev.autoqaPush.endpoint`, default `qa.omp.sh`). Each
|
||||
* insert schedules a background flush that POSTs pending rows and deletes
|
||||
* them on HTTP 2xx. `PI_AUTO_QA_PUSH=1` forces push in non-interactive
|
||||
* environments where the consent dialog never fires. Tool execution is
|
||||
* never blocked on the network and never throws.
|
||||
*/
|
||||
import { Database } from "bun:sqlite";
|
||||
import * as os from "node:os";
|
||||
import path from "node:path";
|
||||
import type { AgentTool } from "@oh-my-pi/pi-agent-core";
|
||||
import { $flag, getAgentDir, logger, VERSION } from "@oh-my-pi/pi-utils";
|
||||
import { $env, $flag, getAgentDir, getInstallId, logger, VERSION } from "@oh-my-pi/pi-utils";
|
||||
import * as z from "zod/v4";
|
||||
import type { Settings } from "..";
|
||||
import type { ToolSession } from "./index";
|
||||
|
||||
const ReportToolIssueParams = z.object({
|
||||
tool: z.string().describe("tool name"),
|
||||
report: z.string().describe("unexpected behavior"),
|
||||
report: z
|
||||
.string()
|
||||
.describe("unexpected behavior; generic, NEVER PII (paths, file contents, identifiers, prompt text)"),
|
||||
});
|
||||
|
||||
export function isAutoQaEnabled(settings?: Settings): boolean {
|
||||
return $flag("PI_AUTO_QA") || !!settings?.get("dev.autoqa");
|
||||
}
|
||||
|
||||
// ───────────────────────────────────────────────────────────────────────────
|
||||
// Consent gate
|
||||
// ───────────────────────────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* Resolver for the user's "share grievances?" consent.
|
||||
*
|
||||
* Return values:
|
||||
* - `true` — user agreed; record + ship for this run and persist.
|
||||
* - `false` — user declined; suppress for this run and persist.
|
||||
* - `null` — user dismissed the dialog (ESC, click-away, …) without
|
||||
* picking an option. The decision is NOT cached or persisted,
|
||||
* so the next `report_tool_issue` invocation re-prompts.
|
||||
*
|
||||
* Persistence is the tool's job (so subagent invocations can persist into
|
||||
* the disk-backed `Settings` instance the host registered alongside the
|
||||
* handler), not the handler's. Implementations live in hosts that have UI
|
||||
* affordances — today only `InteractiveMode`. When no handler is
|
||||
* registered (CLI subcommands, tests, non-interactive runs) consent
|
||||
* defaults to `false` — the explicit "don't collect by default" stance.
|
||||
*/
|
||||
export type AutoQaConsentHandler = () => Promise<boolean | null>;
|
||||
|
||||
let consentHandler: AutoQaConsentHandler | null = null;
|
||||
/**
|
||||
* Persistent settings instance supplied by the consent-handler registrant.
|
||||
* Subagents have in-memory `Settings` snapshots that don't write to disk;
|
||||
* we persist the decision through this disk-backed reference so a grant
|
||||
* survives across runs even when triggered from a subagent tool call.
|
||||
*/
|
||||
let persistentConsentSettings: Settings | null = null;
|
||||
/**
|
||||
* Process-global cache of the resolved consent decision. Survives across
|
||||
* subagent boundaries (subagents share this module instance), so a grant
|
||||
* in the parent applies immediately to children — including children that
|
||||
* spawned BEFORE the grant and would otherwise see a stale snapshot of
|
||||
* `dev.autoqa.consent` in their isolated `Settings`.
|
||||
*
|
||||
* `null` = never asked, never cached.
|
||||
*/
|
||||
let cachedConsent: boolean | null = null;
|
||||
/**
|
||||
* Single-flight in-flight consent request. While the dialog is open, every
|
||||
* concurrent `report_tool_issue` call (main + every subagent) awaits this
|
||||
* promise instead of stacking duplicate popups.
|
||||
*/
|
||||
let consentInFlight: Promise<boolean> | null = null;
|
||||
|
||||
/**
|
||||
* Register the consent handler and the persistent {@link Settings} instance
|
||||
* the decision should be written to. Passing `null` clears the handler
|
||||
* (e.g. on `InteractiveMode` teardown). Re-registration is authoritative.
|
||||
*/
|
||||
export function setAutoQaConsentHandler(
|
||||
handler: AutoQaConsentHandler | null,
|
||||
persistentSettings: Settings | null = null,
|
||||
): void {
|
||||
consentHandler = handler;
|
||||
persistentConsentSettings = persistentSettings;
|
||||
}
|
||||
|
||||
/** Test-only: clear consent cache + handler. Never call from production code. */
|
||||
export function __resetAutoQaConsentForTests(): void {
|
||||
consentHandler = null;
|
||||
persistentConsentSettings = null;
|
||||
cachedConsent = null;
|
||||
consentInFlight = null;
|
||||
}
|
||||
|
||||
function readPersistedConsent(settings: Settings | undefined): boolean | null {
|
||||
if (!settings) return null;
|
||||
const stored = settings.get("dev.autoqa.consent");
|
||||
if (stored === "granted") return true;
|
||||
if (stored === "denied") return false;
|
||||
return null;
|
||||
}
|
||||
|
||||
function persistConsent(localSettings: Settings | undefined, granted: boolean): void {
|
||||
const value = granted ? "granted" : "denied";
|
||||
// Write on every settings instance we know about. The local one keeps
|
||||
// the in-memory snapshot consistent for the current subagent; the
|
||||
// persistent one (registered by the host) is what actually lands on disk.
|
||||
for (const target of [localSettings, persistentConsentSettings]) {
|
||||
if (!target) continue;
|
||||
try {
|
||||
target.set("dev.autoqa.consent", value);
|
||||
} catch (error) {
|
||||
logger.debug("autoqa consent persist failed", { error: String(error) });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve the user's consent for `report_tool_issue` grievances.
|
||||
*
|
||||
* Precedence (highest first):
|
||||
* 1. Process-global cache (set on first successful resolution).
|
||||
* 2. Persistent setting (`dev.autoqa.consent` on the supplied `Settings`).
|
||||
* 3. Persistent setting on the registered host `Settings`.
|
||||
* 4. Consent handler popup (single-flight; persists the answer).
|
||||
* 5. Default-deny when no handler is registered.
|
||||
*
|
||||
* Never throws — handler errors degrade to "denied for this call" without
|
||||
* caching, so a subsequent invocation can re-prompt instead of being
|
||||
* permanently locked into the false branch.
|
||||
*/
|
||||
export async function resolveAutoQaConsent(settings: Settings | undefined): Promise<boolean> {
|
||||
if (cachedConsent !== null) return cachedConsent;
|
||||
const persisted = readPersistedConsent(settings) ?? readPersistedConsent(persistentConsentSettings ?? undefined);
|
||||
if (persisted !== null) {
|
||||
cachedConsent = persisted;
|
||||
return persisted;
|
||||
}
|
||||
if (!consentHandler) return false;
|
||||
if (consentInFlight) return consentInFlight;
|
||||
const handler = consentHandler;
|
||||
consentInFlight = (async () => {
|
||||
try {
|
||||
const granted = await handler();
|
||||
if (granted === null) {
|
||||
// User dismissed the dialog (ESC) without picking. Treat as
|
||||
// "skip this call" but don't cache or persist — the next
|
||||
// invocation gets to re-prompt so a stray ESC isn't a
|
||||
// permanent opt-out.
|
||||
return false;
|
||||
}
|
||||
cachedConsent = granted;
|
||||
persistConsent(settings, granted);
|
||||
return granted;
|
||||
} catch (error) {
|
||||
logger.warn("autoqa consent handler threw", { error: String(error) });
|
||||
return false;
|
||||
} finally {
|
||||
consentInFlight = null;
|
||||
}
|
||||
})();
|
||||
return consentInFlight;
|
||||
}
|
||||
|
||||
export function getAutoQaDbPath(): string {
|
||||
return path.join(getAgentDir(), "autoqa.db");
|
||||
}
|
||||
|
||||
let cachedDb: Database | null = null;
|
||||
|
||||
function openDb(): Database | null {
|
||||
/**
|
||||
* Open (or return the cached handle for) the auto-QA SQLite database at
|
||||
* `~/.omp/agent/autoqa.db`. Idempotently runs schema creation, the
|
||||
* `pushed`-column migration, and index setup so every consumer — tool
|
||||
* execute path, manual `omp grievances push`, future debug scripts —
|
||||
* sees the same prepared schema. Returns `null` only on a hard open
|
||||
* failure (filesystem permissions, etc.); a missing file is created.
|
||||
*
|
||||
* Exported because the `omp grievances` CLI handlers need the migrated
|
||||
* handle too — having a second `openDb` in the CLI led to the column
|
||||
* never being added on the manual-push path.
|
||||
*/
|
||||
export function openAutoQaDb(): Database | null {
|
||||
if (cachedDb) return cachedDb;
|
||||
try {
|
||||
const db = new Database(getAutoQaDbPath());
|
||||
@@ -41,9 +209,22 @@ function openDb(): Database | null {
|
||||
model TEXT NOT NULL,
|
||||
version TEXT NOT NULL,
|
||||
tool TEXT NOT NULL,
|
||||
report TEXT NOT NULL
|
||||
report TEXT NOT NULL,
|
||||
pushed INTEGER NOT NULL DEFAULT 0
|
||||
);
|
||||
`);
|
||||
// Migration: pre-`pushed` databases get the column tacked on. Existing
|
||||
// rows default to `0` (unpushed), so legacy grievances from before the
|
||||
// consent + push pipeline went live get swept up by the next flush —
|
||||
// exactly the behaviour we want for users who just granted consent.
|
||||
const cols = db.prepare("PRAGMA table_info(grievances)").all() as Array<{ name: string }>;
|
||||
if (!cols.some(c => c.name === "pushed")) {
|
||||
db.run("ALTER TABLE grievances ADD COLUMN pushed INTEGER NOT NULL DEFAULT 0");
|
||||
}
|
||||
// Speed up the per-batch `WHERE pushed = 0` scan that drives the flush
|
||||
// loop. Without the index every batch becomes a full table scan once
|
||||
// pushed rows dominate the table.
|
||||
db.run("CREATE INDEX IF NOT EXISTS grievances_pushed_idx ON grievances(pushed, id)");
|
||||
cachedDb = db;
|
||||
return db;
|
||||
} catch {
|
||||
@@ -51,6 +232,219 @@ function openDb(): Database | null {
|
||||
}
|
||||
}
|
||||
|
||||
// ───────────────────────────────────────────────────────────────────────────
|
||||
// Backend push
|
||||
// ───────────────────────────────────────────────────────────────────────────
|
||||
|
||||
export interface FlushResult {
|
||||
pushed: number;
|
||||
ok: boolean;
|
||||
skipped?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* Optional per-flush controls. Used by `omp grievances push` to surface
|
||||
* progress to a TTY and to skip the user-facing consent gate (manual
|
||||
* pushes are the user's explicit intent, not a side effect of a tool call).
|
||||
*/
|
||||
export interface FlushOptions {
|
||||
/**
|
||||
* Skip the `dev.autoqa.consent === "granted"` gate in
|
||||
* {@link resolvePushConfig}. Endpoint configuration is still required.
|
||||
* Reserved for explicit user-driven pushes (CLI `grievances push`,
|
||||
* future debug recipes); never set from the tool's auto-flush path.
|
||||
*/
|
||||
bypassConsent?: boolean;
|
||||
/**
|
||||
* Fires once at the start of the loop with the snapshot count of
|
||||
* unpushed rows. Subsequent inserts won't be reflected (the count is
|
||||
* a planning hint for progress reporters, not a live total).
|
||||
*/
|
||||
onStart?: (totalUnpushed: number) => void;
|
||||
/**
|
||||
* Fires after every successfully shipped batch with the running pushed
|
||||
* count. Reporters compare against the `totalUnpushed` they saw in
|
||||
* `onStart` to advance their bar.
|
||||
*/
|
||||
onProgress?: (pushedSoFar: number) => void;
|
||||
}
|
||||
|
||||
interface PushConfig {
|
||||
endpoint: string;
|
||||
token: string | undefined;
|
||||
}
|
||||
|
||||
const FLUSH_TIMEOUT_MS = 5_000;
|
||||
const FAILURE_COOLDOWN_MS = 30_000;
|
||||
/**
|
||||
* Per-request batch size. The worker loops until no unpushed rows remain,
|
||||
* shipping `FLUSH_BATCH_SIZE` rows per POST. Tunes the trade-off between
|
||||
* request count and request size — 50 keeps each payload well under the
|
||||
* default `maxBody` limit on the autoqa collector while letting a
|
||||
* realistic backlog (a few hundred legacy rows on first flush after the
|
||||
* consent grant) drain in single-digit requests.
|
||||
*/
|
||||
const FLUSH_BATCH_SIZE = 50;
|
||||
|
||||
let inFlightFlush: Promise<FlushResult> | null = null;
|
||||
let lastFailureAt = 0;
|
||||
|
||||
/** Test-only: clear single-flight + cooldown state. Never call from production code. */
|
||||
export function __resetAutoQaFlushStateForTests(): void {
|
||||
inFlightFlush = null;
|
||||
lastFailureAt = 0;
|
||||
}
|
||||
|
||||
function envOverrideString(name: string): string | undefined {
|
||||
const value = $env[name];
|
||||
if (typeof value !== "string") return undefined;
|
||||
const trimmed = value.trim();
|
||||
return trimmed.length > 0 ? trimmed : undefined;
|
||||
}
|
||||
|
||||
function resolvePushConfig(settings: Settings | undefined, bypassConsent: boolean): PushConfig | null {
|
||||
if (!isAutoQaEnabled(settings)) return null;
|
||||
|
||||
// Consent IS the push opt-in for the auto-flush path. `bypassConsent`
|
||||
// covers explicit user-driven pushes (`omp grievances push`) where the
|
||||
// user clearly intends to ship regardless of dialog state. The
|
||||
// `PI_AUTO_QA_PUSH` env flag stays as a CI/headless override too.
|
||||
if (!bypassConsent) {
|
||||
const consented = settings?.get("dev.autoqa.consent") === "granted";
|
||||
if (!consented && !$flag("PI_AUTO_QA_PUSH")) return null;
|
||||
}
|
||||
|
||||
const endpoint = envOverrideString("PI_AUTO_QA_PUSH_URL") ?? settings?.get("dev.autoqaPush.endpoint");
|
||||
if (!endpoint || endpoint.trim().length === 0) return null;
|
||||
|
||||
const token = envOverrideString("PI_AUTO_QA_PUSH_TOKEN") ?? settings?.get("dev.autoqaPush.token");
|
||||
return { endpoint: endpoint.trim(), token: token && token.length > 0 ? token : undefined };
|
||||
}
|
||||
|
||||
interface GrievanceRow {
|
||||
id: number;
|
||||
model: string;
|
||||
version: string;
|
||||
tool: string;
|
||||
report: string;
|
||||
}
|
||||
|
||||
async function performFlush(db: Database, config: PushConfig, options: FlushOptions = {}): Promise<FlushResult> {
|
||||
const selectStmt = db.prepare(
|
||||
"SELECT id, model, version, tool, report FROM grievances WHERE pushed = 0 ORDER BY id ASC LIMIT ?",
|
||||
);
|
||||
// Planning snapshot — fires once so progress reporters can size their bar.
|
||||
// Mid-flight inserts are NOT folded in (the worker drains them too, but
|
||||
// the progress bar treats the initial backlog as the denominator).
|
||||
if (options.onStart) {
|
||||
const totalRow = db.prepare("SELECT COUNT(*) AS n FROM grievances WHERE pushed = 0").get() as { n: number };
|
||||
options.onStart(totalRow.n);
|
||||
}
|
||||
let totalPushed = 0;
|
||||
for (;;) {
|
||||
const rows = selectStmt.all(FLUSH_BATCH_SIZE) as GrievanceRow[];
|
||||
if (rows.length === 0) return { pushed: totalPushed, ok: true };
|
||||
|
||||
const body = JSON.stringify({
|
||||
agent: { name: "omp", version: VERSION },
|
||||
installId: getInstallId(),
|
||||
host: os.hostname(),
|
||||
entries: rows,
|
||||
});
|
||||
const headers: Record<string, string> = { "content-type": "application/json" };
|
||||
if (config.token) headers.authorization = `Bearer ${config.token}`;
|
||||
|
||||
let response: Response;
|
||||
try {
|
||||
response = await fetch(config.endpoint, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body,
|
||||
signal: AbortSignal.timeout(FLUSH_TIMEOUT_MS),
|
||||
});
|
||||
} catch (error) {
|
||||
lastFailureAt = Date.now();
|
||||
logger.warn("autoqa push failed", {
|
||||
endpoint: config.endpoint,
|
||||
error: String(error),
|
||||
batchSize: rows.length,
|
||||
pushedSoFar: totalPushed,
|
||||
});
|
||||
return { pushed: totalPushed, ok: false };
|
||||
}
|
||||
|
||||
if (!response.ok) {
|
||||
lastFailureAt = Date.now();
|
||||
logger.warn("autoqa push failed", {
|
||||
endpoint: config.endpoint,
|
||||
status: response.status,
|
||||
batchSize: rows.length,
|
||||
pushedSoFar: totalPushed,
|
||||
});
|
||||
return { pushed: totalPushed, ok: false };
|
||||
}
|
||||
|
||||
// Mark just this batch — never touch ids the SELECT didn't return so a
|
||||
// concurrent insert that landed mid-flight isn't claimed-as-shipped on
|
||||
// our behalf. `id IN (?, ?, …)` rather than a range so a non-contiguous
|
||||
// batch (after partial fills, retries, etc.) still flips exactly what
|
||||
// we sent.
|
||||
const ids = rows.map(r => r.id);
|
||||
const placeholders = ids.map(() => "?").join(",");
|
||||
db.prepare(`UPDATE grievances SET pushed = 1 WHERE id IN (${placeholders})`).run(...ids);
|
||||
totalPushed += rows.length;
|
||||
options.onProgress?.(totalPushed);
|
||||
// Loop continues; the next SELECT picks up the next batch (or returns
|
||||
// empty, exiting the loop).
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Flush queued grievances to the configured backend.
|
||||
*
|
||||
* Single-flight: concurrent callers share the in-flight promise. After a
|
||||
* failed push, retries are skipped for {@link FAILURE_COOLDOWN_MS} ms.
|
||||
* Never throws — all errors are caught and routed to the logger.
|
||||
*/
|
||||
export async function flushGrievances(
|
||||
db?: Database,
|
||||
settings?: Settings,
|
||||
options: FlushOptions = {},
|
||||
): Promise<FlushResult> {
|
||||
const config = resolvePushConfig(settings, options.bypassConsent === true);
|
||||
if (!config) return { pushed: 0, ok: false, skipped: true };
|
||||
|
||||
// `bypassConsent` is the user's explicit "ship NOW" intent — skip the
|
||||
// 30s cooldown window so they're not stuck looking at "skipped" after a
|
||||
// transient failure. Auto-flush calls still cool off.
|
||||
const bypass = options.bypassConsent === true;
|
||||
if (!bypass && inFlightFlush) return inFlightFlush;
|
||||
|
||||
if (!bypass && lastFailureAt > 0 && Date.now() - lastFailureAt < FAILURE_COOLDOWN_MS) {
|
||||
return { pushed: 0, ok: false, skipped: true };
|
||||
}
|
||||
|
||||
const handle = db ?? openAutoQaDb();
|
||||
if (!handle) return { pushed: 0, ok: false, skipped: true };
|
||||
|
||||
const promise = (async () => {
|
||||
try {
|
||||
return await performFlush(handle, config, options);
|
||||
} catch (error) {
|
||||
lastFailureAt = Date.now();
|
||||
logger.warn("autoqa push failed", { endpoint: config.endpoint, error: String(error) });
|
||||
return { pushed: 0, ok: false };
|
||||
}
|
||||
})();
|
||||
|
||||
if (!bypass) inFlightFlush = promise;
|
||||
try {
|
||||
return await promise;
|
||||
} finally {
|
||||
if (!bypass) inFlightFlush = null;
|
||||
}
|
||||
}
|
||||
|
||||
export function createReportToolIssueTool(session: ToolSession): AgentTool {
|
||||
const getModel = () => session.getActiveModelString?.() ?? "unknown";
|
||||
|
||||
@@ -62,15 +456,40 @@ export function createReportToolIssueTool(session: ToolSession): AgentTool {
|
||||
parameters: ReportToolIssueParams,
|
||||
intent: "omit",
|
||||
async execute(_toolCallId, rawParams) {
|
||||
// Save is unconditional: the row lives in the user's own SQLite
|
||||
// at ~/.omp/agent/autoqa.db regardless of consent — they always
|
||||
// own their local data and can inspect or wipe it via `omp grievances`.
|
||||
// Consent only gates whether the row is *shipped* to the shared
|
||||
// backend; that decision rides on `dev.autoqa.consent` and is
|
||||
// enforced inside `flushGrievances` via `resolvePushConfig`.
|
||||
try {
|
||||
const params = rawParams as { tool: string; report: string };
|
||||
const db = openDb();
|
||||
db?.prepare("INSERT INTO grievances (model, version, tool, report) VALUES (?, ?, ?, ?)").run(
|
||||
getModel(),
|
||||
VERSION,
|
||||
params.tool,
|
||||
params.report,
|
||||
);
|
||||
const db = openAutoQaDb();
|
||||
if (db) {
|
||||
db.prepare("INSERT INTO grievances (model, version, tool, report) VALUES (?, ?, ?, ?)").run(
|
||||
getModel(),
|
||||
VERSION,
|
||||
params.tool,
|
||||
params.report,
|
||||
);
|
||||
// Fire-and-forget background pipeline:
|
||||
// 1. Trigger the consent popup if it hasn't been answered
|
||||
// (single-flight inside `resolveAutoQaConsent`; subagents
|
||||
// share the same module-level state).
|
||||
// 2. Attempt a flush — `resolvePushConfig` no-ops when consent
|
||||
// isn't granted, so a "no" leaves the row local for later
|
||||
// `omp grievances push` or a future consent change.
|
||||
// Tool execution returns immediately; the model never waits
|
||||
// on the dialog.
|
||||
void (async () => {
|
||||
try {
|
||||
await resolveAutoQaConsent(session.settings);
|
||||
await flushGrievances(db, session.settings);
|
||||
} catch (error) {
|
||||
logger.debug("autoqa post-insert pipeline failed", { error: String(error) });
|
||||
}
|
||||
})();
|
||||
}
|
||||
} catch (error) {
|
||||
logger.error("Failed to record tool issue", { error });
|
||||
}
|
||||
|
||||
@@ -0,0 +1,148 @@
|
||||
/**
|
||||
* Consent gate around `report_tool_issue`. Asserts:
|
||||
*
|
||||
* 1. With no handler registered, consent defaults to `false` and the tool's
|
||||
* `execute` returns the canonical "Noted, thanks!" without touching the DB.
|
||||
* 2. The handler fires exactly once per process even across concurrent calls
|
||||
* (single-flight), and the decision is persisted to both the local and
|
||||
* registered persistent `Settings` instances.
|
||||
* 3. A persisted `"granted"` short-circuits the handler.
|
||||
* 4. A persisted `"denied"` short-circuits the handler AND no-ops the tool.
|
||||
*/
|
||||
import { afterEach, describe, expect, it } from "bun:test";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import {
|
||||
__resetAutoQaConsentForTests,
|
||||
resolveAutoQaConsent,
|
||||
setAutoQaConsentHandler,
|
||||
} from "@oh-my-pi/pi-coding-agent/tools/report-tool-issue";
|
||||
|
||||
afterEach(() => {
|
||||
__resetAutoQaConsentForTests();
|
||||
});
|
||||
|
||||
describe("resolveAutoQaConsent", () => {
|
||||
it("defaults to false when no handler is registered", async () => {
|
||||
const settings = Settings.isolated();
|
||||
expect(await resolveAutoQaConsent(settings)).toBe(false);
|
||||
// Default-deny must NOT persist anything — the next process invocation
|
||||
// gets to re-prompt instead of being silently stuck on "no".
|
||||
expect(settings.get("dev.autoqa.consent")).toBe("unset");
|
||||
});
|
||||
|
||||
it("returns persisted `granted` without invoking the handler", async () => {
|
||||
const settings = Settings.isolated({ "dev.autoqa.consent": "granted" });
|
||||
let calls = 0;
|
||||
setAutoQaConsentHandler(async () => {
|
||||
calls += 1;
|
||||
return false;
|
||||
});
|
||||
expect(await resolveAutoQaConsent(settings)).toBe(true);
|
||||
expect(calls).toBe(0);
|
||||
});
|
||||
|
||||
it("returns persisted `denied` without invoking the handler", async () => {
|
||||
const settings = Settings.isolated({ "dev.autoqa.consent": "denied" });
|
||||
let calls = 0;
|
||||
setAutoQaConsentHandler(async () => {
|
||||
calls += 1;
|
||||
return true;
|
||||
});
|
||||
expect(await resolveAutoQaConsent(settings)).toBe(false);
|
||||
expect(calls).toBe(0);
|
||||
});
|
||||
|
||||
it("invokes the handler exactly once for concurrent callers and persists the answer", async () => {
|
||||
const local = Settings.isolated();
|
||||
const persistent = Settings.isolated();
|
||||
let calls = 0;
|
||||
let release: (v: boolean) => void = () => undefined;
|
||||
setAutoQaConsentHandler(async () => {
|
||||
calls += 1;
|
||||
return new Promise<boolean>(resolve => {
|
||||
release = resolve;
|
||||
});
|
||||
}, persistent);
|
||||
|
||||
const a = resolveAutoQaConsent(local);
|
||||
const b = resolveAutoQaConsent(local);
|
||||
const c = resolveAutoQaConsent(local);
|
||||
// Wait a tick to ensure all three reached the in-flight branch.
|
||||
await Promise.resolve();
|
||||
release(true);
|
||||
|
||||
expect(await a).toBe(true);
|
||||
expect(await b).toBe(true);
|
||||
expect(await c).toBe(true);
|
||||
expect(calls).toBe(1);
|
||||
expect(local.get("dev.autoqa.consent")).toBe("granted");
|
||||
expect(persistent.get("dev.autoqa.consent")).toBe("granted");
|
||||
});
|
||||
|
||||
it("persists a `denied` decision so the next call short-circuits", async () => {
|
||||
const local = Settings.isolated();
|
||||
const persistent = Settings.isolated();
|
||||
let calls = 0;
|
||||
setAutoQaConsentHandler(async () => {
|
||||
calls += 1;
|
||||
return false;
|
||||
}, persistent);
|
||||
|
||||
expect(await resolveAutoQaConsent(local)).toBe(false);
|
||||
expect(await resolveAutoQaConsent(local)).toBe(false);
|
||||
expect(calls).toBe(1);
|
||||
expect(local.get("dev.autoqa.consent")).toBe("denied");
|
||||
expect(persistent.get("dev.autoqa.consent")).toBe("denied");
|
||||
});
|
||||
|
||||
it("does not cache or persist when the handler throws (allows re-prompt)", async () => {
|
||||
const settings = Settings.isolated();
|
||||
let calls = 0;
|
||||
setAutoQaConsentHandler(async () => {
|
||||
calls += 1;
|
||||
throw new Error("dialog crashed");
|
||||
});
|
||||
|
||||
expect(await resolveAutoQaConsent(settings)).toBe(false);
|
||||
// A second call must invoke the handler again — the throw path is
|
||||
// transient, not a stuck "no".
|
||||
expect(await resolveAutoQaConsent(settings)).toBe(false);
|
||||
expect(calls).toBe(2);
|
||||
expect(settings.get("dev.autoqa.consent")).toBe("unset");
|
||||
});
|
||||
|
||||
it("does not cache or persist when the handler returns null (dismiss/ESC)", async () => {
|
||||
const local = Settings.isolated();
|
||||
const persistent = Settings.isolated();
|
||||
let calls = 0;
|
||||
setAutoQaConsentHandler(async () => {
|
||||
calls += 1;
|
||||
// Mirrors the `showHookSelector` ESC path (returns `undefined`,
|
||||
// which `#promptAutoQaConsent` maps to `null`).
|
||||
return null;
|
||||
}, persistent);
|
||||
|
||||
expect(await resolveAutoQaConsent(local)).toBe(false);
|
||||
// Second call must re-prompt — a stray ESC isn't a permanent opt-out.
|
||||
expect(await resolveAutoQaConsent(local)).toBe(false);
|
||||
expect(calls).toBe(2);
|
||||
expect(local.get("dev.autoqa.consent")).toBe("unset");
|
||||
expect(persistent.get("dev.autoqa.consent")).toBe("unset");
|
||||
});
|
||||
|
||||
it("falls back to the registered persistent settings when the local snapshot is unset", async () => {
|
||||
// Mirrors the subagent flow: subagent passes its in-memory snapshot
|
||||
// (which lost the consent edit made on the parent), but the host's
|
||||
// persistent Settings carries the real decision.
|
||||
const subagentLocal = Settings.isolated();
|
||||
const hostPersistent = Settings.isolated({ "dev.autoqa.consent": "granted" });
|
||||
let calls = 0;
|
||||
setAutoQaConsentHandler(async () => {
|
||||
calls += 1;
|
||||
return false;
|
||||
}, hostPersistent);
|
||||
|
||||
expect(await resolveAutoQaConsent(subagentLocal)).toBe(true);
|
||||
expect(calls).toBe(0);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,301 @@
|
||||
import { Database } from "bun:sqlite";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import { __resetAutoQaFlushStateForTests, flushGrievances } from "@oh-my-pi/pi-coding-agent/tools/report-tool-issue";
|
||||
import * as piUtils from "@oh-my-pi/pi-utils";
|
||||
import { hookFetch } from "@oh-my-pi/pi-utils";
|
||||
|
||||
function openTempDb(): Database {
|
||||
const db = new Database(":memory:");
|
||||
db.run(`
|
||||
CREATE TABLE IF NOT EXISTS grievances (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
model TEXT NOT NULL,
|
||||
version TEXT NOT NULL,
|
||||
tool TEXT NOT NULL,
|
||||
report TEXT NOT NULL,
|
||||
pushed INTEGER NOT NULL DEFAULT 0
|
||||
);
|
||||
`);
|
||||
return db;
|
||||
}
|
||||
|
||||
function insertGrievance(db: Database, tool: string, report: string): number {
|
||||
const info = db
|
||||
.prepare("INSERT INTO grievances (model, version, tool, report) VALUES (?, ?, ?, ?)")
|
||||
.run("test-model", "test-version", tool, report);
|
||||
return Number(info.lastInsertRowid);
|
||||
}
|
||||
|
||||
/** All rows, regardless of pushed state. */
|
||||
function selectIds(db: Database): number[] {
|
||||
return (db.prepare("SELECT id FROM grievances ORDER BY id ASC").all() as Array<{ id: number }>).map(r => r.id);
|
||||
}
|
||||
|
||||
/** Just unpushed rows — what the next flush would pick up. */
|
||||
function selectUnpushedIds(db: Database): number[] {
|
||||
return (db.prepare("SELECT id FROM grievances WHERE pushed = 0 ORDER BY id ASC").all() as Array<{ id: number }>).map(
|
||||
r => r.id,
|
||||
);
|
||||
}
|
||||
|
||||
/** Just pushed rows — what's already been shipped. */
|
||||
function selectPushedIds(db: Database): number[] {
|
||||
return (db.prepare("SELECT id FROM grievances WHERE pushed = 1 ORDER BY id ASC").all() as Array<{ id: number }>).map(
|
||||
r => r.id,
|
||||
);
|
||||
}
|
||||
|
||||
function pushSettings(overrides: Record<string, unknown> = {}): Settings {
|
||||
return Settings.isolated({
|
||||
"dev.autoqa": true,
|
||||
// Consent is the push opt-in; `granted` is what `resolvePushConfig`
|
||||
// gates on (or `PI_AUTO_QA_PUSH=1` for headless overrides).
|
||||
"dev.autoqa.consent": "granted",
|
||||
"dev.autoqaPush.endpoint": "https://qa.example.com/grievances",
|
||||
...overrides,
|
||||
});
|
||||
}
|
||||
|
||||
describe("flushGrievances", () => {
|
||||
let db: Database;
|
||||
|
||||
beforeEach(() => {
|
||||
__resetAutoQaFlushStateForTests();
|
||||
vi.spyOn(piUtils, "getInstallId").mockReturnValue("11111111-2222-3333-4444-555555555555");
|
||||
db = openTempDb();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
__resetAutoQaFlushStateForTests();
|
||||
db.close();
|
||||
});
|
||||
|
||||
it("skips network when consent is missing and leaves rows intact", async () => {
|
||||
insertGrievance(db, "find", "weird ordering");
|
||||
const fetchSpy = vi.fn(() => new Response("unexpected", { status: 200 }));
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
// `denied` is the user-facing kill switch for push.
|
||||
const result = await flushGrievances(db, pushSettings({ "dev.autoqa.consent": "denied" }));
|
||||
|
||||
expect(result).toEqual({ pushed: 0, ok: false, skipped: true });
|
||||
expect(fetchSpy).not.toHaveBeenCalled();
|
||||
expect(selectIds(db)).toEqual([1]);
|
||||
});
|
||||
|
||||
it("skips network when endpoint is missing", async () => {
|
||||
insertGrievance(db, "find", "weird ordering");
|
||||
const fetchSpy = vi.fn(() => new Response("unexpected", { status: 200 }));
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const result = await flushGrievances(db, pushSettings({ "dev.autoqaPush.endpoint": "" }));
|
||||
|
||||
expect(result).toEqual({ pushed: 0, ok: false, skipped: true });
|
||||
expect(fetchSpy).not.toHaveBeenCalled();
|
||||
expect(selectIds(db)).toEqual([1]);
|
||||
});
|
||||
|
||||
it("returns ok without fetching when there is nothing to push", async () => {
|
||||
const fetchSpy = vi.fn(() => new Response("unexpected", { status: 200 }));
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const result = await flushGrievances(db, pushSettings());
|
||||
|
||||
expect(result).toEqual({ pushed: 0, ok: true });
|
||||
expect(fetchSpy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("posts pending rows with bearer header and marks them pushed=1 on 200", async () => {
|
||||
insertGrievance(db, "find", "weird ordering");
|
||||
insertGrievance(db, "read", "selector ignored");
|
||||
|
||||
let capturedInput: string | URL | Request | undefined;
|
||||
let capturedInit: RequestInit | undefined;
|
||||
const fetchSpy = vi.fn((input: string | URL | Request, init: RequestInit | undefined) => {
|
||||
capturedInput = input;
|
||||
capturedInit = init;
|
||||
return new Response("", { status: 200 });
|
||||
});
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const result = await flushGrievances(db, pushSettings({ "dev.autoqaPush.token": "secret-token" }));
|
||||
|
||||
expect(result).toEqual({ pushed: 2, ok: true });
|
||||
expect(fetchSpy).toHaveBeenCalledTimes(1);
|
||||
expect(String(capturedInput)).toBe("https://qa.example.com/grievances");
|
||||
expect(capturedInit?.method).toBe("POST");
|
||||
|
||||
const headers = capturedInit?.headers as Record<string, string> | undefined;
|
||||
expect(headers?.["content-type"]).toBe("application/json");
|
||||
expect(headers?.authorization).toBe("Bearer secret-token");
|
||||
|
||||
const body = JSON.parse(String(capturedInit?.body));
|
||||
expect(body.agent?.name).toBe("omp");
|
||||
expect(typeof body.agent?.version).toBe("string");
|
||||
expect(typeof body.host).toBe("string");
|
||||
expect(body.installId).toBe("11111111-2222-3333-4444-555555555555");
|
||||
expect(body.entries).toEqual([
|
||||
{ id: 1, model: "test-model", version: "test-version", tool: "find", report: "weird ordering" },
|
||||
{ id: 2, model: "test-model", version: "test-version", tool: "read", report: "selector ignored" },
|
||||
]);
|
||||
|
||||
// Rows are retained for inspection — `pushed=1` flips, but the data
|
||||
// stays so users can browse what they've shipped via `omp grievances`.
|
||||
expect(selectIds(db)).toEqual([1, 2]);
|
||||
expect(selectPushedIds(db)).toEqual([1, 2]);
|
||||
expect(selectUnpushedIds(db)).toEqual([]);
|
||||
});
|
||||
|
||||
it("omits the Authorization header when no token is configured", async () => {
|
||||
insertGrievance(db, "find", "no token here");
|
||||
let capturedInit: RequestInit | undefined;
|
||||
const fetchSpy = vi.fn((_input: string | URL | Request, init: RequestInit | undefined) => {
|
||||
capturedInit = init;
|
||||
return new Response("", { status: 204 });
|
||||
});
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const result = await flushGrievances(db, pushSettings());
|
||||
|
||||
expect(result).toEqual({ pushed: 1, ok: true });
|
||||
const headers = capturedInit?.headers as Record<string, string> | undefined;
|
||||
expect(headers?.authorization).toBeUndefined();
|
||||
expect(selectUnpushedIds(db)).toEqual([]);
|
||||
expect(selectPushedIds(db)).toEqual([1]);
|
||||
});
|
||||
|
||||
it("leaves rows unpushed on 5xx and reports failure", async () => {
|
||||
insertGrievance(db, "find", "boom");
|
||||
const fetchSpy = vi.fn(() => new Response("nope", { status: 500 }));
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const result = await flushGrievances(db, pushSettings());
|
||||
|
||||
expect(result).toEqual({ pushed: 0, ok: false });
|
||||
expect(fetchSpy).toHaveBeenCalledTimes(1);
|
||||
expect(selectUnpushedIds(db)).toEqual([1]);
|
||||
expect(selectPushedIds(db)).toEqual([]);
|
||||
});
|
||||
|
||||
it("drains mid-flight inserts in a follow-up batch within the same loop", async () => {
|
||||
insertGrievance(db, "find", "first");
|
||||
|
||||
const fetchEntered = Promise.withResolvers<void>();
|
||||
const releaseFirstFetch = Promise.withResolvers<Response>();
|
||||
let fetchCount = 0;
|
||||
const fetchSpy = vi.fn(() => {
|
||||
fetchCount += 1;
|
||||
if (fetchCount === 1) {
|
||||
fetchEntered.resolve();
|
||||
return releaseFirstFetch.promise;
|
||||
}
|
||||
// Subsequent loop iterations resolve immediately so the worker
|
||||
// finishes draining without manual coordination per batch.
|
||||
return Promise.resolve(new Response("", { status: 200 }));
|
||||
});
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const flushPromise = flushGrievances(db, pushSettings());
|
||||
await fetchEntered.promise;
|
||||
|
||||
// New grievance written by a concurrent tool call while the push is in flight.
|
||||
insertGrievance(db, "read", "second");
|
||||
|
||||
releaseFirstFetch.resolve(new Response("", { status: 200 }));
|
||||
const result = await flushPromise;
|
||||
|
||||
// Both rows shipped — the worker looped, the second batch picked up
|
||||
// the row that landed mid-flight.
|
||||
expect(result).toEqual({ pushed: 2, ok: true });
|
||||
expect(fetchSpy).toHaveBeenCalledTimes(2);
|
||||
expect(selectUnpushedIds(db)).toEqual([]);
|
||||
expect(selectPushedIds(db)).toEqual([1, 2]);
|
||||
});
|
||||
|
||||
it("collapses concurrent callers onto a single in-flight push", async () => {
|
||||
insertGrievance(db, "find", "single-flight");
|
||||
|
||||
const releaseFetch = Promise.withResolvers<Response>();
|
||||
const fetchSpy = vi.fn(() => releaseFetch.promise);
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const settings = pushSettings();
|
||||
const first = flushGrievances(db, settings);
|
||||
const second = flushGrievances(db, settings);
|
||||
|
||||
releaseFetch.resolve(new Response("", { status: 200 }));
|
||||
const [a, b] = await Promise.all([first, second]);
|
||||
|
||||
expect(fetchSpy).toHaveBeenCalledTimes(1);
|
||||
expect(a).toEqual({ pushed: 1, ok: true });
|
||||
expect(b).toBe(a);
|
||||
expect(selectUnpushedIds(db)).toEqual([]);
|
||||
expect(selectPushedIds(db)).toEqual([1]);
|
||||
});
|
||||
|
||||
it("skips the next push within the failure cooldown window", async () => {
|
||||
insertGrievance(db, "find", "first");
|
||||
const fetchSpy = vi.fn(() => new Response("nope", { status: 500 }));
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const settings = pushSettings();
|
||||
const firstResult = await flushGrievances(db, settings);
|
||||
const secondResult = await flushGrievances(db, settings);
|
||||
|
||||
expect(firstResult).toEqual({ pushed: 0, ok: false });
|
||||
expect(secondResult).toEqual({ pushed: 0, ok: false, skipped: true });
|
||||
expect(fetchSpy).toHaveBeenCalledTimes(1);
|
||||
expect(selectUnpushedIds(db)).toEqual([1]);
|
||||
});
|
||||
|
||||
it("drains a backlog larger than the batch size in multiple POSTs", async () => {
|
||||
// Seed >1 batch worth (FLUSH_BATCH_SIZE = 50) so the worker has to loop.
|
||||
// 127 chosen to land on a non-multiple boundary (2 full batches + a
|
||||
// partial final one), exercising both the LIMIT semantics and the
|
||||
// "remainder smaller than batch" tail.
|
||||
const total = 127;
|
||||
for (let i = 0; i < total; i++) insertGrievance(db, "find", `report-${i}`);
|
||||
|
||||
const seenBatchSizes: number[] = [];
|
||||
const fetchSpy = vi.fn((_input: string | URL | Request, init: RequestInit | undefined) => {
|
||||
const body = JSON.parse(String(init?.body)) as { entries: unknown[] };
|
||||
seenBatchSizes.push(body.entries.length);
|
||||
return new Response("", { status: 200 });
|
||||
});
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const result = await flushGrievances(db, pushSettings());
|
||||
|
||||
expect(result).toEqual({ pushed: total, ok: true });
|
||||
// Three batches: 50 + 50 + 27.
|
||||
expect(seenBatchSizes).toEqual([50, 50, 27]);
|
||||
expect(fetchSpy).toHaveBeenCalledTimes(3);
|
||||
expect(selectUnpushedIds(db)).toEqual([]);
|
||||
expect(selectPushedIds(db).length).toBe(total);
|
||||
});
|
||||
|
||||
it("stops the loop on a mid-batch failure and preserves unpushed rows", async () => {
|
||||
// Two batches' worth — first batch ships, second batch errors. The
|
||||
// pushed-so-far count surfaces in the result and only the unsent
|
||||
// rows stay flagged unpushed.
|
||||
const firstBatch = 50;
|
||||
const secondBatch = 10;
|
||||
for (let i = 0; i < firstBatch + secondBatch; i++) insertGrievance(db, "find", `r-${i}`);
|
||||
|
||||
let call = 0;
|
||||
const fetchSpy = vi.fn(() => {
|
||||
call += 1;
|
||||
return new Response("", { status: call === 1 ? 200 : 500 });
|
||||
});
|
||||
using _hook = hookFetch(fetchSpy);
|
||||
|
||||
const result = await flushGrievances(db, pushSettings());
|
||||
|
||||
expect(result).toEqual({ pushed: firstBatch, ok: false });
|
||||
expect(fetchSpy).toHaveBeenCalledTimes(2);
|
||||
expect(selectPushedIds(db).length).toBe(firstBatch);
|
||||
expect(selectUnpushedIds(db).length).toBe(secondBatch);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user