fix: salvage truncated worker envelopes and log worker lifecycle
Models occasionally drop the closing tag of the result envelope while the JSON itself is complete; the strict parser rejected these and one repair attempt could not always recover, cascading into worker_error and dropping valid attachments. The parser now falls back to brace-balanced salvage when the closing tag is missing, and worker lifecycle events (start, resume, invalid envelope, repair, settle, worker_error) are logged.
This commit is contained in:
@@ -12,6 +12,8 @@ import type { AcpPromptContent } from "./client.js";
|
||||
import { AcpSessionRestoreError, AcpWorker } from "./worker.js";
|
||||
|
||||
const ASSISTANT_POLICY: PermissionPolicy = { mode: "deny", allowedTools: [], allowedCommandPatterns: [] };
|
||||
const ASSISTANT_START_TAG = "<GORI_ASSISTANT_ACTION_V2>";
|
||||
const WORKER_START_TAG = "<GORI_WORKER_RESULT_V2>";
|
||||
const ASSISTANT_ENVELOPE = /<GORI_ASSISTANT_ACTION_V2>(?<json>[\s\S]*?)<\/GORI_ASSISTANT_ACTION_V2>\s*$/;
|
||||
const WORKER_ENVELOPE = /<GORI_WORKER_RESULT_V2>(?<json>[\s\S]*?)<\/GORI_WORKER_RESULT_V2>\s*$/;
|
||||
|
||||
@@ -355,6 +357,7 @@ export class AssistantManager implements ConversationRuntime {
|
||||
if (!proposal) return false;
|
||||
if (this.active || this.proposals.list({ status: "working" }).length > 0) return false;
|
||||
await this.proposals.update(proposal.id, { status: "working", pending: undefined, startedAt: Date.now() });
|
||||
console.log(`Worker follow-up for proposal ${proposalTag(proposal.id)}: resuming session ${proposal.workerNativeSessionId || "none"}`);
|
||||
this.launchWorker(proposal.id, {
|
||||
resumeSessionId: proposal.workerNativeSessionId,
|
||||
firstTurn: workerFollowUpPrompt(instruction, proposal.pending?.question),
|
||||
@@ -469,11 +472,14 @@ export class AssistantManager implements ConversationRuntime {
|
||||
}
|
||||
|
||||
private async startWorkerRun(active: ActiveWorker, options: WorkerLaunchOptions): Promise<void> {
|
||||
let resumed = false;
|
||||
if (options.resumeSessionId) {
|
||||
try {
|
||||
await active.worker.start(options.resumeSessionId);
|
||||
resumed = true;
|
||||
} catch (error) {
|
||||
if (!(error instanceof AcpSessionRestoreError) || this.active !== active || active.settled) throw error;
|
||||
console.log(`Worker resume failed for proposal ${proposalTag(active.proposalId)}; starting a fresh session`);
|
||||
active.worker = this.spawnProposalWorker(active.proposalId);
|
||||
await active.worker.start();
|
||||
options = {
|
||||
@@ -486,6 +492,7 @@ export class AssistantManager implements ConversationRuntime {
|
||||
await active.worker.start();
|
||||
}
|
||||
if (this.active !== active || active.settled) throw new Error("Worker was superseded before its session started");
|
||||
console.log(`Worker started for proposal ${proposalTag(active.proposalId)} (session ${active.worker.nativeSessionId}, ${resumed ? "resumed" : "new"})`);
|
||||
await this.proposals.update(active.proposalId, { workerNativeSessionId: active.worker.nativeSessionId });
|
||||
await active.worker.prompt(this.bot.workerBootstrap, "initializing");
|
||||
if (this.active !== active || active.settled) throw new Error("Worker was superseded during bootstrap");
|
||||
@@ -500,11 +507,16 @@ export class AssistantManager implements ConversationRuntime {
|
||||
if (this.active !== active || active.settled) throw new Error("Worker was superseded during its turn");
|
||||
let result = parseWorkerResult(reply);
|
||||
if (!result) {
|
||||
console.error(`Worker turn for proposal ${proposalTag(active.proposalId)} returned no valid envelope (strict or salvaged); requesting repair`);
|
||||
const repaired = await active.worker.prompt(workerRepairPrompt());
|
||||
if (this.active !== active || active.settled) throw new Error("Worker was superseded during result repair");
|
||||
result = parseWorkerResult(repaired);
|
||||
if (result) console.log(`Worker result repair succeeded for proposal ${proposalTag(active.proposalId)}`);
|
||||
}
|
||||
if (!result) {
|
||||
console.error(`Worker result repair failed for proposal ${proposalTag(active.proposalId)}`);
|
||||
throw new Error("Worker did not return a valid GORI_WORKER_RESULT_V2 envelope after one repair attempt");
|
||||
}
|
||||
if (!result) throw new Error("Worker did not return a valid GORI_WORKER_RESULT_V2 envelope after one repair attempt");
|
||||
await this.settleWorkerResult(active, result);
|
||||
}
|
||||
|
||||
@@ -548,6 +560,7 @@ export class AssistantManager implements ConversationRuntime {
|
||||
pending,
|
||||
lastWorkerSummary: result.summary
|
||||
});
|
||||
console.log(`Worker settled proposal ${proposalTag(active.proposalId)} as pending (attachments=${attachments.length} dropped=${droppedAttachments.length})`);
|
||||
await active.worker.terminate().catch(() => undefined);
|
||||
await this.proposals.update(active.proposalId, { workerProcessGroup: undefined });
|
||||
return true;
|
||||
@@ -566,6 +579,7 @@ export class AssistantManager implements ConversationRuntime {
|
||||
active.settled = true;
|
||||
this.active = undefined;
|
||||
const summary = `worker_error: ${error instanceof Error ? error.message : String(error)}`;
|
||||
console.error(`Worker error for proposal ${proposalTag(active.proposalId)}: ${error instanceof Error ? error.message : String(error)}`);
|
||||
await this.proposals.update(active.proposalId, {
|
||||
status: "pending",
|
||||
pending: { summary, workspaceDirty: true, receivedAt: Date.now() },
|
||||
@@ -816,10 +830,22 @@ export function assistantKeyFor(chatKey: string, userId: string): string {
|
||||
return `${chatKey}#${userId}`;
|
||||
}
|
||||
|
||||
export function parseAssistantActions(text: string): ParsedAssistantEnvelope | undefined { const match = ASSISTANT_ENVELOPE.exec(text);
|
||||
if (!match?.groups) return undefined;
|
||||
function proposalTag(proposalId: string): string {
|
||||
return proposalId.slice(0, 8);
|
||||
}
|
||||
|
||||
export function parseAssistantActions(text: string): ParsedAssistantEnvelope | undefined {
|
||||
const direct = parseAssistantEnvelopeValue(strictEnvelopeJson(text, ASSISTANT_ENVELOPE));
|
||||
if (direct) return direct;
|
||||
const salvaged = parseAssistantEnvelopeValue(salvageEnvelopeJson(text, ASSISTANT_START_TAG));
|
||||
if (salvaged) console.log("Assistant action envelope salvaged (missing closing tag)");
|
||||
return salvaged;
|
||||
}
|
||||
|
||||
function parseAssistantEnvelopeValue(json: string | undefined): ParsedAssistantEnvelope | undefined {
|
||||
if (!json) return undefined;
|
||||
let value: unknown;
|
||||
try { value = JSON.parse(match.groups.json); } catch { return undefined; }
|
||||
try { value = JSON.parse(json); } catch { return undefined; }
|
||||
if (!isRecord(value) || typeof value.reply !== "string" || !Array.isArray(value.actions)) return undefined;
|
||||
const actions: AssistantAction[] = [];
|
||||
for (const item of value.actions) {
|
||||
@@ -878,10 +904,17 @@ function parseAssistantAction(value: unknown): AssistantAction | undefined {
|
||||
}
|
||||
|
||||
export function parseWorkerResult(text: string): ParsedWorkerResult | undefined {
|
||||
const match = WORKER_ENVELOPE.exec(text);
|
||||
if (!match?.groups) return undefined;
|
||||
const direct = parseWorkerResultValue(strictEnvelopeJson(text, WORKER_ENVELOPE));
|
||||
if (direct) return direct;
|
||||
const salvaged = parseWorkerResultValue(salvageEnvelopeJson(text, WORKER_START_TAG));
|
||||
if (salvaged) console.log("Worker result envelope salvaged (missing closing tag)");
|
||||
return salvaged;
|
||||
}
|
||||
|
||||
function parseWorkerResultValue(json: string | undefined): ParsedWorkerResult | undefined {
|
||||
if (!json) return undefined;
|
||||
let value: unknown;
|
||||
try { value = JSON.parse(match.groups.json); } catch { return undefined; }
|
||||
try { value = JSON.parse(json); } catch { return undefined; }
|
||||
if (!isRecord(value) || typeof value.summary !== "string") return undefined;
|
||||
if (value.status !== "PENDING") return undefined;
|
||||
if ((value.question !== undefined && typeof value.question !== "string")
|
||||
@@ -905,6 +938,47 @@ export function parseWorkerResult(text: string): ParsedWorkerResult | undefined
|
||||
};
|
||||
}
|
||||
|
||||
function strictEnvelopeJson(text: string, envelope: RegExp): string | undefined {
|
||||
return envelope.exec(text)?.groups?.json;
|
||||
}
|
||||
|
||||
// Fallback for model responses truncated at the envelope's closing tag: brace-balance the JSON
|
||||
// after the start tag (string- and escape-aware) and accept it only when the balanced object
|
||||
// ends exactly at the end of the text. Anything unbalanced or with trailing content stays invalid.
|
||||
function salvageEnvelopeJson(text: string, startTag: string): string | undefined {
|
||||
const start = text.lastIndexOf(startTag);
|
||||
if (start < 0) return undefined;
|
||||
const after = text.slice(start + startTag.length);
|
||||
if (!after.startsWith("{")) return undefined;
|
||||
const end = balancedJsonObjectEnd(after);
|
||||
if (end === undefined) return undefined;
|
||||
if (after.slice(end).trim().length > 0) return undefined;
|
||||
return after.slice(0, end);
|
||||
}
|
||||
|
||||
function balancedJsonObjectEnd(text: string): number | undefined {
|
||||
let depth = 0;
|
||||
let inString = false;
|
||||
let escaped = false;
|
||||
for (let index = 0; index < text.length; index++) {
|
||||
const char = text[index]!;
|
||||
if (inString) {
|
||||
if (escaped) escaped = false;
|
||||
else if (char === "\\") escaped = true;
|
||||
else if (char === '"') inString = false;
|
||||
continue;
|
||||
}
|
||||
if (char === '"') inString = true;
|
||||
else if (char === "{") depth++;
|
||||
else if (char === "}") {
|
||||
depth--;
|
||||
if (depth === 0) return index + 1;
|
||||
if (depth < 0) return undefined;
|
||||
}
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function withWorkerResultProtocol(input: string | AcpPromptContent[]): string | AcpPromptContent[] {
|
||||
const instruction = "When this turn is finished, end your response with exactly one hidden worker result envelope: <GORI_WORKER_RESULT_V2>{\"status\":\"PENDING\",\"summary\":\"...\"}</GORI_WORKER_RESULT_V2>. The summary is a short user-readable description of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing, and set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. When you produced image files the user should see (png/jpg only), save them inside the workspace (prefer .gori-outbox/) and report up to 3 of them as \"attachments\": [{\"path\":\"relative/or/absolute/path\"}]; never report paths outside the workspace. PENDING is the only status; do not emit any other status value or any text after the envelope.";
|
||||
if (typeof input === "string") return `${input}\n\n${instruction}`;
|
||||
|
||||
Reference in New Issue
Block a user