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..d2a7e61e15b4 100644
--- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
+++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
@@ -293,18 +293,28 @@ 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";
}
-function assistantSegmentMessageId(
+/** 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,
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 +2238,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 +2635,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..58c7864c8f96 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,9 +17,10 @@ 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}`;
}
+/** 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 77517c44ea9b..1e918188bbad 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?.();
}),
),
),
@@ -1173,6 +1177,9 @@ function isStartupMetadataUpdate(notification: EffectAcpSchema.SessionNotificati
}
}
+/**
+ * Processes an ACP session notification, manages assistant segments, and dispatches runtime events.
+ */
const handleSessionUpdate = ({
queue,
modeStateRef,
@@ -1202,11 +1209,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 +1231,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 +1250,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) {
@@ -1279,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,
@@ -1319,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 f2f370068181..acc3f4c34d86 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),
@@ -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 4e7541ee8c10..93ec3a989ca5 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,101 @@ 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");
+ });
+
+ 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("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(
+ 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 1993c9457af5..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;
@@ -234,39 +239,241 @@ 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(
- notification: EffectAcpSchema.SessionNotification,
-): EffectAcpSchema.SessionNotification {
- const update = notification.update;
- if (update.sessionUpdate !== "tool_call" && update.sessionUpdate !== "tool_call_update") {
- return notification;
+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;
+ }
+ }
}
- 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 } : {}),
- },
+ 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(): AntigravityMessageFilter {
+ let inSystemMessage = false;
+ let pending = "";
+
+ const reset = (): void => {
+ inSystemMessage = false;
+ pending = "";
};
+
+ const filter = (text: string): string => {
+ const textToProcess = pending + text;
+ pending = "";
+
+ let result = "";
+ let cursor = 0;
+
+ while (cursor < textToProcess.length) {
+ if (inSystemMessage) {
+ const closeTagIndex = textToProcess.indexOf(CLOSE_SYSTEM_MESSAGE_TAG, cursor);
+ if (closeTagIndex === -1) {
+ // 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 + CLOSE_SYSTEM_MESSAGE_TAG.length;
+ inSystemMessage = false;
+ continue;
+ }
+
+ 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 && effectiveOpenTagIndex !== -1) {
+ if (preambleIndex < effectiveOpenTagIndex) {
+ nextIndex = preambleIndex;
+ isPreamble = true;
+ } else {
+ nextIndex = effectiveOpenTagIndex;
+ }
+ } else if (preambleIndex !== -1) {
+ nextIndex = preambleIndex;
+ isPreamble = true;
+ } else if (effectiveOpenTagIndex !== -1) {
+ nextIndex = effectiveOpenTagIndex;
+ }
+
+ if (nextIndex === -1) {
+ 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 += textToProcess.slice(cursor, nextIndex);
+
+ if (isPreamble) {
+ const afterPreamble = textToProcess.indexOf(
+ OPEN_SYSTEM_MESSAGE_TAG,
+ nextIndex + SYSTEM_MESSAGE_PREAMBLE.length,
+ );
+ if (afterPreamble !== -1) {
+ cursor = afterPreamble + OPEN_SYSTEM_MESSAGE_TAG.length;
+ inSystemMessage = true;
+ } else {
+ const nextNewline = textToProcess.indexOf("\n", nextIndex);
+ if (nextNewline !== -1) {
+ cursor = nextNewline + 1;
+ } else {
+ cursor = textToProcess.length;
+ }
+ }
+ } else {
+ cursor = nextIndex + OPEN_SYSTEM_MESSAGE_TAG.length;
+ inSystemMessage = true;
+ }
+ }
+
+ return result;
+ };
+
+ filter.reset = reset;
+ return filter;
+}
+
+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();
+
+ const reset = (): void => {
+ filterMessage.reset();
+ filterThought.reset();
+ };
+
+ const transform = (
+ notification: EffectAcpSchema.SessionNotification,
+ ): EffectAcpSchema.SessionNotification => {
+ const update = notification.update;
+ 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 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 } : {}),
+ },
+ };
+ };
+
+ transform.reset = reset;
+ return transform;
+}
+
+/** The runtime uses this before it retains tool state or dispatches raw callbacks. */
+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;
@@ -288,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));
@@ -374,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;