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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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",
);
});
});
22 changes: 16 additions & 6 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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}`,
);
}

Expand Down Expand Up @@ -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,
}
Expand Down Expand Up @@ -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,
Expand Down
7 changes: 6 additions & 1 deletion apps/server/src/provider/RuntimeInstructions.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
const SYSTEM_TELEMETRY_INSTRUCTIONS = `<system_messages>
Internal harness notifications, background task updates, and messages enclosed in <SYSTEM_MESSAGE> are strictly for contextual awareness. Never repeat, quote, or echo <SYSTEM_MESSAGE> tags, background task logs, or internal harness status lines to the user.
</system_messages>`;

const PULL_REQUEST_LINKING_INSTRUCTIONS = `<pull_request_linking>
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.
</pull_request_linking>`;
Expand All @@ -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 `<runtime_info>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.</runtime_info>\n\n${PULL_REQUEST_LINKING_INSTRUCTIONS}`;
return `<runtime_info>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.</runtime_info>\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();
}
34 changes: 26 additions & 8 deletions apps/server/src/provider/acp/AcpSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void, EffectAcpErrors.AcpError>;
readonly requestLogger?: (event: AcpSessionRequestLogEvent) => Effect.Effect<void, never>;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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?.();
}),
),
),
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
Expand All @@ -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;
}
Expand All @@ -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) {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -1319,6 +1334,9 @@ const ensureActiveAssistantSegment = ({
),
);

/**
* Emits a completion event for the currently open assistant segment, if one is active.
*/
const closeActiveAssistantSegment = ({
queue,
assistantSegmentRef,
Expand Down
6 changes: 4 additions & 2 deletions apps/server/src/provider/acp/AntigravityAcpSupport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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),
Expand All @@ -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":
Expand All @@ -100,6 +101,7 @@ export function antigravityPermissionMode(runtimeMode: RuntimeMode): string {
}
}

/** Extracts select model options from Antigravity session configuration options. */
export function antigravityModelOptions(
configOptions: ReadonlyArray<EffectAcpSchema.SessionConfigOption>,
) {
Expand Down
99 changes: 99 additions & 0 deletions apps/server/src/provider/acp/AntigravityProtocol.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ import {
classifyAntigravitySubagentToolCall,
isAntigravityUserInputRequest,
makeAntigravityUserInputResponse,
createAntigravityMessageFilter,
makeAntigravitySessionUpdateTransformer,
normalizeAntigravitySessionUpdate,
normalizeAntigravityToolCall,
sanitizeAntigravityToolPayload,
Expand Down Expand Up @@ -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 <SYSTEM_MESSAGE> not actually sent by the user. It is provided by the system as important information to pay attention to.\n\n<SYSTEM_MESSAGE>\n[Message] timestamp=2026-09-22T20:00:23Z sender=task-1 content=Task finished\n</SYSTEM_MESSAGE>Here 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 <SYSTEM_MESSAGE> not actually sent by the user.\n\n<SYSTEM_MESSAGE>\n[Message] part 1",
),
).toBe("");
expect(filter("\npart 2\nOutput: [mobile] OK\n")).toBe("");
expect(filter("</SYSTEM_MESSAGE> After")).toBe(" After");
});

it("handles closing delimiter split across chunks without dropping the assistant answer", () => {
const filter = createAntigravityMessageFilter();
expect(filter("<SYSTEM_MESSAGE>secret</SYSTEM_MESS")).toBe("");
expect(filter("AGE> real answer")).toBe(" real answer");
});

it("handles opening delimiter split across chunks without leaking telemetry", () => {
const filter = createAntigravityMessageFilter();
expect(filter("Normal text <SYSTEM_")).toBe("Normal text ");
expect(filter("MESSAGE>secret</SYSTEM_MESSAGE> more text")).toBe(" more text");
});

it("handles preamble split across chunks without leaking telemetry", () => {
const filter = createAntigravityMessageFilter();
expect(filter("Answer: The following is a <SYSTEM_")).toBe("Answer: ");
expect(
filter(
"MESSAGE> not actually sent by the user.\n\n<SYSTEM_MESSAGE>secret</SYSTEM_MESSAGE>Done",
),
).toBe("Done");
});

it("buffers incomplete preamble containing embedded opening tag without dropping assistant answer", () => {
const filter = createAntigravityMessageFilter();
expect(filter("The following is a <SYSTEM_MESSAGE> 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 <SYSTEM_MESSAGE>")).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 <SYSTEM_MESSAGE> 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 <SYSTEM_MESSAGE>
transformer({
sessionId: "session-1",
update: {
sessionUpdate: "agent_message_chunk",
content: { type: "text", text: "<SYSTEM_MESSAGE>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.");
});
});
Loading
Loading