From 1db2682d9ef4128a6215cf0ca2f8bd07c765177a Mon Sep 17 00:00:00 2001 From: Adeeb Ahmad Date: Wed, 23 Sep 2026 03:26:43 +0500 Subject: [PATCH 1/3] fix(acp): prevent assistant message fragmentation on background tool updates & filter leaked harness telemetry Fixes #13133 - Only close active assistant segments on new tool calls (isNew: previous === undefined) instead of on every in-flight tool progress or completion event. - Prevent duplicate `assistant:assistant:` message IDs in ProviderRuntimeIngestion. - Add streaming filter to strip echoed blocks and harness preambles from Antigravity session chunks. - Add negative constraint to buildRuntimeInstructions preventing models from quoting harness telemetry. --- .../Layers/ProviderRuntimeIngestion.test.ts | 19 +++ .../Layers/ProviderRuntimeIngestion.ts | 17 +- .../src/provider/RuntimeInstructions.ts | 6 +- .../src/provider/acp/AcpSessionRuntime.ts | 17 +- .../src/provider/acp/AntigravityAcpSupport.ts | 4 +- .../provider/acp/AntigravityProtocol.test.ts | 31 ++++ .../src/provider/acp/AntigravityProtocol.ts | 161 ++++++++++++++---- 7 files changed, 211 insertions(+), 44 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index d61739f72c21..7eda1a2e7096 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -61,6 +61,8 @@ import * as ThreadPlanProgress from "../ThreadPlanProgress.ts"; import { ProviderRuntimeIngestionLive, splitBufferedAssistantText, + canonicalAssistantMessageId, + assistantSegmentMessageId, } from "./ProviderRuntimeIngestion.ts"; import { DEFAULT_THREAD_TITLE } from "../threadTitles.ts"; import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; @@ -5175,3 +5177,20 @@ describe("splitBufferedAssistantText", () => { }); }); }); + +describe("assistant message id canonicalization and segmentation", () => { + it("avoids double-prefixing assistant message IDs", () => { + expect(canonicalAssistantMessageId("item-123")).toBe("assistant:item-123"); + expect(canonicalAssistantMessageId("assistant:item-123")).toBe("assistant:item-123"); + }); + + it("normalizes segments without nesting the prefix", () => { + expect(assistantSegmentMessageId("item-123", 0)).toBe("assistant:item-123"); + expect(assistantSegmentMessageId("assistant:item-123", 0)).toBe("assistant:item-123"); + expect(assistantSegmentMessageId("item-123", 1)).toBe("assistant:item-123:segment:1"); + expect(assistantSegmentMessageId("assistant:item-123", 1)).toBe("assistant:item-123:segment:1"); + expect(assistantSegmentMessageId("reasoning:thought-1", 1, "reasoning")).toBe( + "reasoning:thought-1:segment:1", + ); + }); +}); diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 0db70e491235..d55575bc35f1 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -297,14 +297,19 @@ function messageStreamRoleOf(messageId: MessageId): MessageStreamRole { return messageId.startsWith(REASONING_MESSAGE_ID_PREFIX) ? "reasoning" : "assistant"; } -function assistantSegmentMessageId( +export function canonicalAssistantMessageId(rawId: string): MessageId { + return MessageId.make(rawId.startsWith("assistant:") ? rawId : `assistant:${rawId}`); +} + +export function assistantSegmentMessageId( baseKey: string, segmentIndex: number, role: MessageStreamRole = "assistant", ): MessageId { const prefix = role === "reasoning" ? REASONING_MESSAGE_ID_PREFIX : "assistant:"; + const normalizedBaseKey = baseKey.startsWith(prefix) ? baseKey.slice(prefix.length) : baseKey; return MessageId.make( - segmentIndex === 0 ? `${prefix}${baseKey}` : `${prefix}${baseKey}:segment:${segmentIndex}`, + segmentIndex === 0 ? `${prefix}${normalizedBaseKey}` : `${prefix}${normalizedBaseKey}:segment:${segmentIndex}`, ); } @@ -2228,8 +2233,8 @@ const make = Effect.gen(function* () { const assistantCompletion = event.type === "item.completed" && event.payload.itemType === "assistant_message" ? { - messageId: MessageId.make( - `assistant:${event.itemId ?? event.turnId ?? event.eventId}`, + messageId: canonicalAssistantMessageId( + String(event.itemId ?? event.turnId ?? event.eventId), ), fallbackText: event.payload.detail, } @@ -2625,8 +2630,8 @@ const make = Effect.gen(function* () { checkpointRef: CheckpointRef.make(`provider-diff:${event.eventId}`), status: "missing", files: [], - assistantMessageId: MessageId.make( - `assistant:${event.itemId ?? event.turnId ?? event.eventId}`, + assistantMessageId: canonicalAssistantMessageId( + String(event.itemId ?? event.turnId ?? event.eventId), ), checkpointTurnCount: maxCheckpointTurnCount(checkpointContext.checkpoints) + 1, createdAt: now, diff --git a/apps/server/src/provider/RuntimeInstructions.ts b/apps/server/src/provider/RuntimeInstructions.ts index 5e72586062e5..8a95de9b1d72 100644 --- a/apps/server/src/provider/RuntimeInstructions.ts +++ b/apps/server/src/provider/RuntimeInstructions.ts @@ -1,3 +1,7 @@ +const SYSTEM_TELEMETRY_INSTRUCTIONS = ` +Internal harness notifications, background task updates, and messages enclosed in are strictly for contextual awareness. Never repeat, quote, or echo tags, background task logs, or internal harness status lines to the user. +`; + const PULL_REQUEST_LINKING_INSTRUCTIONS = ` When the t3-code MCP server exposes link_pull_request, you must use it to register every pull request you create or work on for this thread. Call link_pull_request with the full PR URL immediately after creating a PR or starting work on an existing PR. For a stack, call it for every layer, not just the current branch or the top PR. This applies when creating or updating PRs through gh, gh stack, another CLI, or the host API: those operations do not register the PRs with this thread. Linking an already-linked PR is safe. Before finishing PR work, call list_thread_pull_requests and link any PR from your work that is missing. Do not link unrelated PRs mentioned only as background. If a linking call fails, report that failure instead of claiming the PR is linked. `; @@ -13,7 +17,7 @@ export function buildRuntimeInstructions(runtime: { const effort = toSingleLine(runtime.reasoningEffort ?? ""); const modelInfo = model && model !== "auto" && model !== "default" ? `, as ${model}` : ""; const effortInfo = effort ? ` with ${effort} reasoning effort` : ""; - return `In case you're asked: you are running in T3 Code through the ${harness} harness${modelInfo}${effortInfo}. No need to mention this otherwise. You can embed images and videos in your response using Markdown with absolute file paths.\n\n${PULL_REQUEST_LINKING_INSTRUCTIONS}`; + return `In case you're asked: you are running in T3 Code through the ${harness} harness${modelInfo}${effortInfo}. No need to mention this otherwise. You can embed images and videos in your response using Markdown with absolute file paths.\n\n${PULL_REQUEST_LINKING_INSTRUCTIONS}\n\n${SYSTEM_TELEMETRY_INSTRUCTIONS}`; } function toSingleLine(value: string): string { diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index 77517c44ea9b..e5d4c253e2b0 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -1202,11 +1202,7 @@ const handleSessionUpdate = ({ } for (const event of parsed.events) { if (event._tag === "ToolCallUpdated") { - yield* closeActiveAssistantSegment({ - queue, - assistantSegmentRef, - }); - const { merged, decision } = yield* Ref.modify(toolCallsRef, (current) => { + const { merged, decision, isNew } = yield* Ref.modify(toolCallsRef, (current) => { const tracked = current.get(event.toolCall.toolCallId); const previous = tracked?.state; const nextToolCall = mergeToolCallState(previous, event.toolCall); @@ -1228,8 +1224,14 @@ const handleSessionUpdate = ({ skippedSinceEmit: decision.skippedSinceEmit, }); } - return [{ merged: nextToolCall, decision }, next] as const; + return [{ merged: nextToolCall, decision, isNew: previous === undefined }, next] as const; }); + if (isNew) { + yield* closeActiveAssistantSegment({ + queue, + assistantSegmentRef, + }); + } if (!decision.emit) { continue; } @@ -1241,6 +1243,9 @@ const handleSessionUpdate = ({ continue; } if (event._tag === "ContentDelta") { + if (event.text.length === 0) { + continue; + } if (event.text.trim().length === 0) { const assistantSegmentState = yield* Ref.get(assistantSegmentRef); if (!assistantSegmentState.activeItemId) { diff --git a/apps/server/src/provider/acp/AntigravityAcpSupport.ts b/apps/server/src/provider/acp/AntigravityAcpSupport.ts index f2f370068181..42e395a89863 100644 --- a/apps/server/src/provider/acp/AntigravityAcpSupport.ts +++ b/apps/server/src/provider/acp/AntigravityAcpSupport.ts @@ -23,7 +23,7 @@ import { makeAntigravityStdoutTransform, } from "../antigravityAuthSupport.ts"; import * as AcpSessionRuntime from "./AcpSessionRuntime.ts"; -import { normalizeAntigravitySessionUpdate } from "./AntigravityProtocol.ts"; +import { makeAntigravitySessionUpdateTransformer } from "./AntigravityProtocol.ts"; export interface AntigravityAcpRuntimeInput extends Omit< AcpSessionRuntime.AcpSessionRuntimeOptions, @@ -78,7 +78,7 @@ export const makeAntigravityAcpRuntime = Effect.fn("makeAntigravityAcpRuntime")( onStderr: makeAntigravityStderrHandler( input.onAuthorizationUrl ? { onAuthorizationUrl: input.onAuthorizationUrl } : {}, ), - transformSessionUpdate: normalizeAntigravitySessionUpdate, + transformSessionUpdate: makeAntigravitySessionUpdateTransformer(), }).pipe( Layer.provide( Layer.succeed(ChildProcessSpawner.ChildProcessSpawner, input.childProcessSpawner), diff --git a/apps/server/src/provider/acp/AntigravityProtocol.test.ts b/apps/server/src/provider/acp/AntigravityProtocol.test.ts index 4e7541ee8c10..8979c8a3ef53 100644 --- a/apps/server/src/provider/acp/AntigravityProtocol.test.ts +++ b/apps/server/src/provider/acp/AntigravityProtocol.test.ts @@ -11,6 +11,8 @@ import { classifyAntigravitySubagentToolCall, isAntigravityUserInputRequest, makeAntigravityUserInputResponse, + createAntigravityMessageFilter, + makeAntigravitySessionUpdateTransformer, normalizeAntigravitySessionUpdate, normalizeAntigravityToolCall, sanitizeAntigravityToolPayload, @@ -482,4 +484,33 @@ describe("Antigravity tool results", () => { ); expect(isAntigravityOpenCommand(completed)).toBe(false); }); + + it("sanitizes echoed SYSTEM_MESSAGE telemetry blocks from agent message chunks", () => { + const transformer = makeAntigravitySessionUpdateTransformer(); + const leaked = + "The following is a not actually sent by the user. It is provided by the system as important information to pay attention to.\n\n\n[Message] timestamp=2026-09-22T20:00:23Z sender=task-1 content=Task finished\nHere is the real answer."; + + const notification = { + sessionId: "session-1", + update: { + sessionUpdate: "agent_message_chunk" as const, + content: { type: "text" as const, text: leaked }, + }, + }; + + const normalized = transformer(notification); + expect((normalized.update as any).content.text).toBe("Here is the real answer."); + }); + + it("sanitizes streamed SYSTEM_MESSAGE blocks across chunk boundaries", () => { + const filter = createAntigravityMessageFilter(); + expect(filter("Before ")).toBe("Before "); + expect( + filter( + "The following is a not actually sent by the user.\n\n\n[Message] part 1", + ), + ).toBe(""); + expect(filter("\npart 2\nOutput: [mobile] OK\n")).toBe(""); + expect(filter(" After")).toBe(" After"); + }); }); diff --git a/apps/server/src/provider/acp/AntigravityProtocol.ts b/apps/server/src/provider/acp/AntigravityProtocol.ts index 1993c9457af5..98a25d7d74f2 100644 --- a/apps/server/src/provider/acp/AntigravityProtocol.ts +++ b/apps/server/src/provider/acp/AntigravityProtocol.ts @@ -234,39 +234,142 @@ export function sanitizeAntigravityToolPayload(payload: unknown): unknown { return sanitizeToolValue(payload, { nodes: 512, text: 64_000 }, 0); } -/** The runtime uses this before it retains tool state or dispatches raw callbacks. */ -export function normalizeAntigravitySessionUpdate( +/** + * Stateful streaming filter that discards internal Antigravity harness system messages + * (e.g. `...` and the system preamble) that the model + * may mistakenly echo back into its assistant message or thought stream. + */ +export function createAntigravityMessageFilter(): (text: string) => string { + let inSystemMessage = false; + + return (text: string): string => { + let result = ""; + let cursor = 0; + + while (cursor < text.length) { + if (inSystemMessage) { + const closeTagIndex = text.indexOf("", cursor); + if (closeTagIndex === -1) { + // Entire remainder of text is inside + break; + } + cursor = closeTagIndex + "".length; + inSystemMessage = false; + continue; + } + + const preambleText = "The following is a not actually sent by the user"; + const preambleIndex = text.indexOf(preambleText, cursor); + const openTagIndex = text.indexOf("", cursor); + + let nextIndex = -1; + let isPreamble = false; + + if (preambleIndex !== -1 && openTagIndex !== -1) { + if (preambleIndex < openTagIndex) { + nextIndex = preambleIndex; + isPreamble = true; + } else { + nextIndex = openTagIndex; + } + } else if (preambleIndex !== -1) { + nextIndex = preambleIndex; + isPreamble = true; + } else if (openTagIndex !== -1) { + nextIndex = openTagIndex; + } + + if (nextIndex === -1) { + result += text.slice(cursor); + break; + } + + result += text.slice(cursor, nextIndex); + + if (isPreamble) { + const afterPreamble = text.indexOf("", nextIndex); + if (afterPreamble !== -1) { + cursor = afterPreamble + "".length; + inSystemMessage = true; + } else { + const nextNewline = text.indexOf("\n", nextIndex); + if (nextNewline !== -1) { + cursor = nextNewline + 1; + } else { + cursor = text.length; + } + } + } else { + cursor = nextIndex + "".length; + inSystemMessage = true; + } + } + + return result; + }; +} + +export function makeAntigravitySessionUpdateTransformer(): ( notification: EffectAcpSchema.SessionNotification, -): EffectAcpSchema.SessionNotification { - const update = notification.update; - if (update.sessionUpdate !== "tool_call" && update.sessionUpdate !== "tool_call_update") { - return notification; - } - const contentBudget = { nodes: 512, text: 32_000 }; - const content = update.content?.flatMap((entry) => { - const decoded = Option.getOrUndefined( - decodeToolCallContent(sanitizeToolValue(entry, contentBudget, 0)), - ); - return decoded === undefined ? [] : [decoded]; - }); - const meta = sanitizeAntigravityToolPayload(update._meta); - return { - ...notification, - update: { - ...update, - ...(typeof update.title === "string" ? { title: boundText(update.title) } : {}), - ...(update.rawInput !== undefined - ? { rawInput: sanitizeAntigravityToolPayload(update.rawInput) } - : {}), - ...(update.rawOutput !== undefined - ? { rawOutput: sanitizeAntigravityToolPayload(update.rawOutput) } - : {}), - ...(update.content !== undefined ? { content: content ?? [] } : {}), - ...(update._meta !== undefined ? { _meta: Predicate.isObject(meta) ? meta : null } : {}), - }, +) => EffectAcpSchema.SessionNotification { + const filterMessage = createAntigravityMessageFilter(); + const filterThought = createAntigravityMessageFilter(); + + return (notification: EffectAcpSchema.SessionNotification): EffectAcpSchema.SessionNotification => { + const update = notification.update; + if (update.sessionUpdate === "agent_message_chunk" || update.sessionUpdate === "agent_thought_chunk") { + if (update.content.type === "text") { + const filter = update.sessionUpdate === "agent_message_chunk" ? filterMessage : filterThought; + const filtered = filter(update.content.text); + if (filtered !== update.content.text) { + return { + ...notification, + update: { + ...update, + content: { + ...update.content, + text: filtered, + }, + }, + }; + } + } + return notification; + } + if (update.sessionUpdate !== "tool_call" && update.sessionUpdate !== "tool_call_update") { + return notification; + } + const contentBudget = { nodes: 512, text: 32_000 }; + const content = update.content?.flatMap((entry) => { + const decoded = Option.getOrUndefined( + decodeToolCallContent(sanitizeToolValue(entry, contentBudget, 0)), + ); + return decoded === undefined ? [] : [decoded]; + }); + const meta = sanitizeAntigravityToolPayload(update._meta); + return { + ...notification, + update: { + ...update, + ...(typeof update.title === "string" ? { title: boundText(update.title) } : {}), + ...(update.rawInput !== undefined + ? { rawInput: sanitizeAntigravityToolPayload(update.rawInput) } + : {}), + ...(update.rawOutput !== undefined + ? { rawOutput: sanitizeAntigravityToolPayload(update.rawOutput) } + : {}), + ...(update.content !== undefined ? { content: content ?? [] } : {}), + ...(update._meta !== undefined ? { _meta: Predicate.isObject(meta) ? meta : null } : {}), + }, + }; }; } +/** The runtime uses this before it retains tool state or dispatches raw callbacks. */ +export const normalizeAntigravitySessionUpdate: ( + notification: EffectAcpSchema.SessionNotification, +) => EffectAcpSchema.SessionNotification = makeAntigravitySessionUpdateTransformer(); + function localImagePath(imagePath: string | undefined): string | undefined { if (!imagePath || imagePath.length > TOOL_TEXT_LIMIT) { return undefined; From 8f006a005d76a3c8d1219f19176389db581a9b2f Mon Sep 17 00:00:00 2001 From: Adeeb Ahmad Date: Wed, 23 Sep 2026 04:58:38 +0500 Subject: [PATCH 2/3] fix(acp): buffer split telemetry delimiters, fix preamble search, and reset filters at prompt boundaries --- .../Layers/ProviderRuntimeIngestion.ts | 6 +- .../src/provider/acp/AcpSessionRuntime.ts | 8 +- .../provider/acp/AntigravityProtocol.test.ts | 56 ++++++++ .../src/provider/acp/AntigravityProtocol.ts | 135 ++++++++++++++---- 4 files changed, 177 insertions(+), 28 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index d55575bc35f1..ad7dd4ef30df 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -297,10 +297,12 @@ function messageStreamRoleOf(messageId: MessageId): MessageStreamRole { return messageId.startsWith(REASONING_MESSAGE_ID_PREFIX) ? "reasoning" : "assistant"; } +/** Ensures an assistant message identifier carries exactly one `assistant:` prefix. */ export function canonicalAssistantMessageId(rawId: string): MessageId { return MessageId.make(rawId.startsWith("assistant:") ? rawId : `assistant:${rawId}`); } +/** Formats a unique segment identifier for assistant or reasoning message parts. */ export function assistantSegmentMessageId( baseKey: string, segmentIndex: number, @@ -309,7 +311,9 @@ export function assistantSegmentMessageId( const prefix = role === "reasoning" ? REASONING_MESSAGE_ID_PREFIX : "assistant:"; const normalizedBaseKey = baseKey.startsWith(prefix) ? baseKey.slice(prefix.length) : baseKey; return MessageId.make( - segmentIndex === 0 ? `${prefix}${normalizedBaseKey}` : `${prefix}${normalizedBaseKey}:segment:${segmentIndex}`, + segmentIndex === 0 + ? `${prefix}${normalizedBaseKey}` + : `${prefix}${normalizedBaseKey}:segment:${segmentIndex}`, ); } diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index e5d4c253e2b0..853e8019584c 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -100,9 +100,11 @@ export interface AcpSessionRuntimeOptions { /** Transforms provider stdout before protocol parsing and protocol logging. */ readonly transformStdout?: EffectAcpClient.AcpClientOptions["transformStdout"]; /** Normalizes provider-specific fields before notification queues or runtime state retain them. */ - readonly transformSessionUpdate?: ( + readonly transformSessionUpdate?: (( notification: EffectAcpSchema.SessionNotification, - ) => EffectAcpSchema.SessionNotification; + ) => EffectAcpSchema.SessionNotification) & { + readonly reset?: () => void; + }; /** Receives bounded stderr chunks. Redact secrets before logging. A failure closes the runtime. */ readonly onStderr?: (text: string) => Effect.Effect; readonly requestLogger?: (event: AcpSessionRequestLogEvent) => Effect.Effect; @@ -1032,6 +1034,7 @@ export const make = ( const started = yield* getStartedState; yield* closeActiveAssistantSegment({ queue: eventQueue, assistantSegmentRef }); yield* Ref.set(assistantUpdatesOpenRef, true); + options.transformSessionUpdate?.reset?.(); const requestPayload = { sessionId: started.sessionId, ...payload, @@ -1081,6 +1084,7 @@ export const make = ( yield* Fiber.interrupt(activePrompt.fiber).pipe(Effect.ignore); yield* Ref.set(activePromptRef, Option.none()); yield* Deferred.succeed(activePrompt.completed, undefined); + options.transformSessionUpdate?.reset?.(); }), ), ), diff --git a/apps/server/src/provider/acp/AntigravityProtocol.test.ts b/apps/server/src/provider/acp/AntigravityProtocol.test.ts index 8979c8a3ef53..bd7c502affe5 100644 --- a/apps/server/src/provider/acp/AntigravityProtocol.test.ts +++ b/apps/server/src/provider/acp/AntigravityProtocol.test.ts @@ -513,4 +513,60 @@ describe("Antigravity tool results", () => { expect(filter("\npart 2\nOutput: [mobile] OK\n")).toBe(""); expect(filter(" After")).toBe(" After"); }); + + it("handles closing delimiter split across chunks without dropping the assistant answer", () => { + const filter = createAntigravityMessageFilter(); + expect(filter("secret real answer")).toBe(" real answer"); + }); + + it("handles opening delimiter split across chunks without leaking telemetry", () => { + const filter = createAntigravityMessageFilter(); + expect(filter("Normal text secret more text")).toBe(" more text"); + }); + + it("handles preamble split across chunks without leaking telemetry", () => { + const filter = createAntigravityMessageFilter(); + expect(filter("Answer: The following is a not actually sent by the user.\n\nsecretDone", + ), + ).toBe("Done"); + }); + + it("discards through newline when preamble has no subsequent opening tag", () => { + const filter = createAntigravityMessageFilter(); + expect( + filter( + "The following is a not actually sent by the user\nHere is the answer.", + ), + ).toBe("Here is the answer."); + }); + + it("resets filter state at prompt boundary so incomplete prior turns do not swallow text", () => { + const transformer = makeAntigravitySessionUpdateTransformer(); + // Prompt 1 ends unexpectedly while inside + transformer({ + sessionId: "session-1", + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "unclosed secret" }, + }, + }); + + // Reset at prompt boundary via user_message_chunk or explicit reset + transformer.reset(); + + // Prompt 2 assistant text should not be swallowed + const nextResponse = transformer({ + sessionId: "session-1", + update: { + sessionUpdate: "agent_message_chunk", + content: { type: "text", text: "Legitimate prompt 2 answer." }, + }, + }); + expect((nextResponse.update as any).content.text).toBe("Legitimate prompt 2 answer."); + }); }); diff --git a/apps/server/src/provider/acp/AntigravityProtocol.ts b/apps/server/src/provider/acp/AntigravityProtocol.ts index 98a25d7d74f2..94568843ed66 100644 --- a/apps/server/src/provider/acp/AntigravityProtocol.ts +++ b/apps/server/src/provider/acp/AntigravityProtocol.ts @@ -234,33 +234,77 @@ export function sanitizeAntigravityToolPayload(payload: unknown): unknown { return sanitizeToolValue(payload, { nodes: 512, text: 64_000 }, 0); } +const OPEN_SYSTEM_MESSAGE_TAG = ""; +const CLOSE_SYSTEM_MESSAGE_TAG = ""; +const SYSTEM_MESSAGE_PREAMBLE = "The following is a not actually sent by the user"; + +/** + * Finds the length of the longest trailing suffix of `str` that matches + * a non-empty prefix of any candidate string. + */ +function findLongestCandidateSuffix(str: string, candidates: readonly string[]): number { + let longest = 0; + for (const candidate of candidates) { + const maxLen = Math.min(str.length, candidate.length - 1); + for (let len = maxLen; len > longest; len--) { + if (str.endsWith(candidate.slice(0, len))) { + longest = len; + break; + } + } + } + return longest; +} + +export interface AntigravityMessageFilter { + (text: string): string; + reset: () => void; +} + /** * Stateful streaming filter that discards internal Antigravity harness system messages * (e.g. `...` and the system preamble) that the model * may mistakenly echo back into its assistant message or thought stream. + * + * Partial opening tags, closing tags, and preamble markers are buffered across chunks + * so split delimiters do not leak telemetry or swallow assistant answers. */ -export function createAntigravityMessageFilter(): (text: string) => string { +export function createAntigravityMessageFilter(): AntigravityMessageFilter { let inSystemMessage = false; + let pending = ""; + + const reset = (): void => { + inSystemMessage = false; + pending = ""; + }; + + const filter = (text: string): string => { + const textToProcess = pending + text; + pending = ""; - return (text: string): string => { let result = ""; let cursor = 0; - while (cursor < text.length) { + while (cursor < textToProcess.length) { if (inSystemMessage) { - const closeTagIndex = text.indexOf("", cursor); + const closeTagIndex = textToProcess.indexOf(CLOSE_SYSTEM_MESSAGE_TAG, cursor); if (closeTagIndex === -1) { - // Entire remainder of text is inside + // Entire remainder of text is inside , except possibly + // a trailing partial close tag that might be completed in the next chunk. + const remainder = textToProcess.slice(cursor); + const suffixLen = findLongestCandidateSuffix(remainder, [CLOSE_SYSTEM_MESSAGE_TAG]); + if (suffixLen > 0) { + pending = remainder.slice(remainder.length - suffixLen); + } break; } - cursor = closeTagIndex + "".length; + cursor = closeTagIndex + CLOSE_SYSTEM_MESSAGE_TAG.length; inSystemMessage = false; continue; } - const preambleText = "The following is a not actually sent by the user"; - const preambleIndex = text.indexOf(preambleText, cursor); - const openTagIndex = text.indexOf("", cursor); + const preambleIndex = textToProcess.indexOf(SYSTEM_MESSAGE_PREAMBLE, cursor); + const openTagIndex = textToProcess.indexOf(OPEN_SYSTEM_MESSAGE_TAG, cursor); let nextIndex = -1; let isPreamble = false; @@ -280,46 +324,85 @@ export function createAntigravityMessageFilter(): (text: string) => string { } if (nextIndex === -1) { - result += text.slice(cursor); + const remainder = textToProcess.slice(cursor); + const suffixLen = findLongestCandidateSuffix(remainder, [ + OPEN_SYSTEM_MESSAGE_TAG, + SYSTEM_MESSAGE_PREAMBLE, + ]); + if (suffixLen > 0) { + result += remainder.slice(0, remainder.length - suffixLen); + pending = remainder.slice(remainder.length - suffixLen); + } else { + result += remainder; + } break; } - result += text.slice(cursor, nextIndex); + result += textToProcess.slice(cursor, nextIndex); if (isPreamble) { - const afterPreamble = text.indexOf("", nextIndex); + const afterPreamble = textToProcess.indexOf( + OPEN_SYSTEM_MESSAGE_TAG, + nextIndex + SYSTEM_MESSAGE_PREAMBLE.length, + ); if (afterPreamble !== -1) { - cursor = afterPreamble + "".length; + cursor = afterPreamble + OPEN_SYSTEM_MESSAGE_TAG.length; inSystemMessage = true; } else { - const nextNewline = text.indexOf("\n", nextIndex); + const nextNewline = textToProcess.indexOf("\n", nextIndex); if (nextNewline !== -1) { cursor = nextNewline + 1; } else { - cursor = text.length; + cursor = textToProcess.length; } } } else { - cursor = nextIndex + "".length; + cursor = nextIndex + OPEN_SYSTEM_MESSAGE_TAG.length; inSystemMessage = true; } } return result; }; + + filter.reset = reset; + return filter; } -export function makeAntigravitySessionUpdateTransformer(): ( - notification: EffectAcpSchema.SessionNotification, -) => EffectAcpSchema.SessionNotification { +export interface AntigravitySessionUpdateTransformer { + (notification: EffectAcpSchema.SessionNotification): EffectAcpSchema.SessionNotification; + reset: () => void; +} + +/** + * Creates a stateful session update transformer that sanitizes Antigravity tool payloads + * and filters leaked system telemetry from assistant message and thought streams. + * Resets per-channel streaming filters at prompt boundaries. + */ +export function makeAntigravitySessionUpdateTransformer(): AntigravitySessionUpdateTransformer { const filterMessage = createAntigravityMessageFilter(); const filterThought = createAntigravityMessageFilter(); - return (notification: EffectAcpSchema.SessionNotification): EffectAcpSchema.SessionNotification => { + const reset = (): void => { + filterMessage.reset(); + filterThought.reset(); + }; + + const transform = ( + notification: EffectAcpSchema.SessionNotification, + ): EffectAcpSchema.SessionNotification => { const update = notification.update; - if (update.sessionUpdate === "agent_message_chunk" || update.sessionUpdate === "agent_thought_chunk") { + if (update.sessionUpdate === "user_message_chunk") { + reset(); + return notification; + } + if ( + update.sessionUpdate === "agent_message_chunk" || + update.sessionUpdate === "agent_thought_chunk" + ) { if (update.content.type === "text") { - const filter = update.sessionUpdate === "agent_message_chunk" ? filterMessage : filterThought; + const filter = + update.sessionUpdate === "agent_message_chunk" ? filterMessage : filterThought; const filtered = filter(update.content.text); if (filtered !== update.content.text) { return { @@ -363,12 +446,14 @@ export function makeAntigravitySessionUpdateTransformer(): ( }, }; }; + + transform.reset = reset; + return transform; } /** The runtime uses this before it retains tool state or dispatches raw callbacks. */ -export const normalizeAntigravitySessionUpdate: ( - notification: EffectAcpSchema.SessionNotification, -) => EffectAcpSchema.SessionNotification = makeAntigravitySessionUpdateTransformer(); +export const normalizeAntigravitySessionUpdate: AntigravitySessionUpdateTransformer = + makeAntigravitySessionUpdateTransformer(); function localImagePath(imagePath: string | undefined): string | undefined { if (!imagePath || imagePath.length > TOOL_TEXT_LIMIT) { From 45864a0356dd3168a9da83db2153d85b07a1946f Mon Sep 17 00:00:00 2001 From: Adeeb Ahmad Date: Wed, 23 Sep 2026 05:15:01 +0500 Subject: [PATCH 3/3] fix(acp): buffer incomplete preambles containing embedded system message tag --- .../Layers/ProviderRuntimeIngestion.ts | 1 + .../src/provider/RuntimeInstructions.ts | 1 + .../src/provider/acp/AcpSessionRuntime.ts | 9 ++++++ .../src/provider/acp/AntigravityAcpSupport.ts | 2 ++ .../provider/acp/AntigravityProtocol.test.ts | 12 +++++++ .../src/provider/acp/AntigravityProtocol.ts | 31 ++++++++++++++++--- 6 files changed, 51 insertions(+), 5 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index ad7dd4ef30df..d2a7e61e15b4 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -293,6 +293,7 @@ type MessageStreamRole = "assistant" | "reasoning"; const REASONING_MESSAGE_ID_PREFIX = "reasoning:"; +/** Determines whether a message ID corresponds to reasoning or assistant output. */ function messageStreamRoleOf(messageId: MessageId): MessageStreamRole { return messageId.startsWith(REASONING_MESSAGE_ID_PREFIX) ? "reasoning" : "assistant"; } diff --git a/apps/server/src/provider/RuntimeInstructions.ts b/apps/server/src/provider/RuntimeInstructions.ts index 8a95de9b1d72..58c7864c8f96 100644 --- a/apps/server/src/provider/RuntimeInstructions.ts +++ b/apps/server/src/provider/RuntimeInstructions.ts @@ -20,6 +20,7 @@ export function buildRuntimeInstructions(runtime: { return `In case you're asked: you are running in T3 Code through the ${harness} harness${modelInfo}${effortInfo}. No need to mention this otherwise. You can embed images and videos in your response using Markdown with absolute file paths.\n\n${PULL_REQUEST_LINKING_INSTRUCTIONS}\n\n${SYSTEM_TELEMETRY_INSTRUCTIONS}`; } +/** Normalizes a string by collapsing consecutive whitespace into single spaces and trimming it. */ function toSingleLine(value: string): string { return value.replaceAll(/\s+/g, " ").trim(); } diff --git a/apps/server/src/provider/acp/AcpSessionRuntime.ts b/apps/server/src/provider/acp/AcpSessionRuntime.ts index 853e8019584c..1e918188bbad 100644 --- a/apps/server/src/provider/acp/AcpSessionRuntime.ts +++ b/apps/server/src/provider/acp/AcpSessionRuntime.ts @@ -1177,6 +1177,9 @@ function isStartupMetadataUpdate(notification: EffectAcpSchema.SessionNotificati } } +/** + * Processes an ACP session notification, manages assistant segments, and dispatches runtime events. + */ const handleSessionUpdate = ({ queue, modeStateRef, @@ -1288,6 +1291,9 @@ function updateModeState(modeState: AcpSessionModeState, nextModeId: string): Ac const assistantItemId = (sessionId: string, runtimeId: string, segmentIndex: number) => `assistant:${sessionId}:runtime:${runtimeId}:segment:${segmentIndex}`; +/** + * Ensures an active assistant segment exists, starting a new indexed segment if none is open. + */ const ensureActiveAssistantSegment = ({ queue, assistantSegmentRef, @@ -1328,6 +1334,9 @@ const ensureActiveAssistantSegment = ({ ), ); +/** + * Emits a completion event for the currently open assistant segment, if one is active. + */ const closeActiveAssistantSegment = ({ queue, assistantSegmentRef, diff --git a/apps/server/src/provider/acp/AntigravityAcpSupport.ts b/apps/server/src/provider/acp/AntigravityAcpSupport.ts index 42e395a89863..acc3f4c34d86 100644 --- a/apps/server/src/provider/acp/AntigravityAcpSupport.ts +++ b/apps/server/src/provider/acp/AntigravityAcpSupport.ts @@ -88,6 +88,7 @@ export const makeAntigravityAcpRuntime = Effect.fn("makeAntigravityAcpRuntime")( return yield* Effect.service(AcpSessionRuntime.AcpSessionRuntime).pipe(Effect.provide(context)); }); +/** Maps a user-facing runtime mode to the Antigravity ACP permission mode identifier. */ export function antigravityPermissionMode(runtimeMode: RuntimeMode): string { switch (runtimeMode) { case "full-access": @@ -100,6 +101,7 @@ export function antigravityPermissionMode(runtimeMode: RuntimeMode): string { } } +/** Extracts select model options from Antigravity session configuration options. */ export function antigravityModelOptions( configOptions: ReadonlyArray, ) { diff --git a/apps/server/src/provider/acp/AntigravityProtocol.test.ts b/apps/server/src/provider/acp/AntigravityProtocol.test.ts index bd7c502affe5..93ec3a989ca5 100644 --- a/apps/server/src/provider/acp/AntigravityProtocol.test.ts +++ b/apps/server/src/provider/acp/AntigravityProtocol.test.ts @@ -536,6 +536,18 @@ describe("Antigravity tool results", () => { ).toBe("Done"); }); + it("buffers incomplete preamble containing embedded opening tag without dropping assistant answer", () => { + const filter = createAntigravityMessageFilter(); + expect(filter("The following is a not actually")).toBe(""); + expect(filter(" sent by the user\nLegitimate answer")).toBe("Legitimate answer"); + }); + + it("buffers incomplete preamble ending immediately after opening tag", () => { + const filter = createAntigravityMessageFilter(); + expect(filter("Answer: The following is a ")).toBe("Answer: "); + expect(filter(" not actually sent by the user\nFinal text")).toBe("Final text"); + }); + it("discards through newline when preamble has no subsequent opening tag", () => { const filter = createAntigravityMessageFilter(); expect( diff --git a/apps/server/src/provider/acp/AntigravityProtocol.ts b/apps/server/src/provider/acp/AntigravityProtocol.ts index 94568843ed66..2ce86460204b 100644 --- a/apps/server/src/provider/acp/AntigravityProtocol.ts +++ b/apps/server/src/provider/acp/AntigravityProtocol.ts @@ -45,6 +45,7 @@ export function isAntigravityUserInputRequest( return request.toolCall.toolCallId.startsWith("interaction_"); } +/** Matches a provider approval decision to the corresponding Antigravity permission option identifier. */ export function selectAntigravityPermissionOptionId( request: EffectAcpSchema.RequestPermissionRequest, decision: ProviderApprovalDecision, @@ -111,10 +112,12 @@ export function antigravityApprovalOptions( return options; } +/** Creates an isolated string slice copy so V8 does not retain larger source strings in memory. */ function copyBoundedText(text: string): string { return Buffer.from(text, "utf16le").toString("utf16le"); } +/** Formats and bounds the display label for an Antigravity permission option. */ function questionLabel(option: EffectAcpSchema.PermissionOption): string { const label = option.name.trim() || option.optionId; return label.length > QUESTION_LABEL_LIMIT @@ -175,6 +178,7 @@ export function makeAntigravityUserInputResponse( return option ? { outcome: { outcome: "selected", optionId: option.optionId } } : undefined; } +/** Bounds text to a maximum length, prepending a truncation marker if the text was shortened. */ function boundText(text: string, limit = TOOL_TEXT_LIMIT): string { return text.length <= limit ? text @@ -186,6 +190,7 @@ interface ToolPayloadBudget { text: number; } +/** Recursively traverses and bounds tool payload values against node count and string length limits. */ function sanitizeToolValue(value: unknown, budget: ToolPayloadBudget, depth: number): unknown { if (depth > 12 || budget.nodes-- <= 0) { return undefined; @@ -306,21 +311,34 @@ export function createAntigravityMessageFilter(): AntigravityMessageFilter { const preambleIndex = textToProcess.indexOf(SYSTEM_MESSAGE_PREAMBLE, cursor); const openTagIndex = textToProcess.indexOf(OPEN_SYSTEM_MESSAGE_TAG, cursor); + const remainder = textToProcess.slice(cursor); + const preambleCandidateLen = findLongestCandidateSuffix(remainder, [SYSTEM_MESSAGE_PREAMBLE]); + const preambleCandidateStart = + preambleCandidateLen > 0 ? textToProcess.length - preambleCandidateLen : -1; + + // If openTagIndex points to the embedded inside an incomplete + // preamble at the end of the text, do not treat it as a standalone tag. + const effectiveOpenTagIndex = + openTagIndex !== -1 && + (preambleCandidateStart === -1 || openTagIndex < preambleCandidateStart) + ? openTagIndex + : -1; + let nextIndex = -1; let isPreamble = false; - if (preambleIndex !== -1 && openTagIndex !== -1) { - if (preambleIndex < openTagIndex) { + if (preambleIndex !== -1 && effectiveOpenTagIndex !== -1) { + if (preambleIndex < effectiveOpenTagIndex) { nextIndex = preambleIndex; isPreamble = true; } else { - nextIndex = openTagIndex; + nextIndex = effectiveOpenTagIndex; } } else if (preambleIndex !== -1) { nextIndex = preambleIndex; isPreamble = true; - } else if (openTagIndex !== -1) { - nextIndex = openTagIndex; + } else if (effectiveOpenTagIndex !== -1) { + nextIndex = effectiveOpenTagIndex; } if (nextIndex === -1) { @@ -455,6 +473,7 @@ export function makeAntigravitySessionUpdateTransformer(): AntigravitySessionUpd export const normalizeAntigravitySessionUpdate: AntigravitySessionUpdateTransformer = makeAntigravitySessionUpdateTransformer(); +/** Validates and normalizes an image file path from tool output for preview rendering. */ function localImagePath(imagePath: string | undefined): string | undefined { if (!imagePath || imagePath.length > TOOL_TEXT_LIMIT) { return undefined; @@ -476,6 +495,7 @@ function localImagePath(imagePath: string | undefined): string | undefined { return /^[a-z][a-z\d+.-]*:/i.test(path) && !/^[a-z]:[\\/]/i.test(path) ? undefined : path; } +/** Normalizes tool call payloads, populating command, cwd, exit code, and image paths. */ export function normalizeAntigravityToolCall(toolCall: AcpToolCallState): AcpToolCallState { const input = Option.getOrUndefined(decodeNativeToolFields(toolCall.data.rawInput)); const output = Option.getOrUndefined(decodeNativeToolFields(toolCall.data.rawOutput)); @@ -562,6 +582,7 @@ export function isAntigravitySubagentReplayStart(rawPayload: unknown): boolean { ); } +/** Extracts and bounds the raw output string from a completed subagent tool call. */ export function antigravitySubagentOutput(toolCall: AcpToolCallState): string | undefined { const output = toolCall.data.rawOutput; return typeof output === "string" && output.trim() ? boundText(output.trim()) : undefined;