Skip to content
Open
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
22 changes: 11 additions & 11 deletions apps/mobile/src/features/threads/ThreadFeed.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,7 @@ import {
} from "@t3tools/mobile-markdown-text/links";
import {
failedFeedRunIds,
terminalFeedAssistantMessageIds,
deriveThreadFeedPresentation,
threadFeedRunIsUnsettled,
isContextCompactionActivityGroup,
Expand Down Expand Up @@ -321,6 +322,8 @@ function AssistantForkButton(props: {
readonly environmentId: EnvironmentId;
readonly iconColor: ColorValue;
readonly projectedItem: OrchestrationV2ProjectedTurnItem;
/** A later response in the run follows, so the fork cuts at this item. */
readonly midRun: boolean;
readonly sourceTitle: string;
}) {
const support = useV2ItemSupport({
Expand Down Expand Up @@ -354,6 +357,7 @@ function AssistantForkButton(props: {
sourceThreadId: props.projectedItem.sourceThreadId,
targetThreadId,
runId,
...(props.midRun ? { turnItemId: props.projectedItem.item.id } : {}),
title: `${props.sourceTitle} fork`,
creationSource: "mobile",
},
Expand Down Expand Up @@ -1525,7 +1529,7 @@ function renderFeedEntry(
readonly expandedWorkRows: Record<string, boolean>;
readonly workRowSizing: ReturnType<typeof deriveThreadWorkLogSizing>;
readonly workGroupScrollPositions: Map<string, ThreadWorkGroupScrollPosition>;
readonly terminalAssistantMessageIds: ReadonlySet<string>;
readonly terminalAssistantMessageIds: ReturnType<typeof terminalFeedAssistantMessageIds>;
readonly unsettledTurnId: RunId | null;
readonly failedRunIds: ReadonlySet<RunId>;
readonly onCopyWorkRow: (rowId: string, value: string) => void;
Expand Down Expand Up @@ -1694,7 +1698,7 @@ function renderFeedEntry(
message.runId === props.unsettledTurnId;
const showAssistantMeta =
message.role === "assistant" &&
props.terminalAssistantMessageIds.has(message.id) &&
props.terminalAssistantMessageIds.terminalIds.has(message.id) &&
!assistantTurnStillInProgress &&
!message.streaming;

Expand Down Expand Up @@ -1966,6 +1970,7 @@ function renderFeedEntry(
environmentId={props.environmentId}
iconColor={iconSubtleColor}
projectedItem={message.projectedItem}
midRun={props.terminalAssistantMessageIds.midRunIds.has(message.id)}
sourceTitle={props.threadTitle}
/>
) : null}
Expand Down Expand Up @@ -2761,15 +2766,10 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) {
() => failedFeedRunIds(props.feed, props.latestRun),
[props.feed, props.latestRun],
);
const terminalAssistantMessageIds = useMemo(() => {
const terminalIdsByTurn = new Map<RunId, string>();
for (const entry of props.feed) {
if (entry.type === "message" && entry.message.role === "assistant" && entry.message.runId) {
terminalIdsByTurn.set(entry.message.runId, entry.message.id);
}
}
return new Set(terminalIdsByTurn.values());
}, [props.feed]);
const terminalAssistantMessageIds = useMemo(
() => terminalFeedAssistantMessageIds(props.feed),
[props.feed],
);
useEffect(() => {
const previous = previousLatestTurnRef.current;
previousLatestTurnRef.current = props.latestRun;
Expand Down
62 changes: 62 additions & 0 deletions apps/mobile/src/lib/threadActivity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ import {
deriveThreadFeedPresentation,
threadFeedActivityIsVisible,
threadFeedRunIsUnsettled,
terminalFeedAssistantMessageIds,
type ThreadFeedActivity,
type ThreadFeedEntry,
togglePendingUserInputOptionSelection,
Expand Down Expand Up @@ -2028,6 +2029,67 @@ const multiSelectQuestion = {
multiSelect: true,
} as const;

describe("terminalFeedAssistantMessageIds", () => {
const runId = RunId.make("steered-run");
const message = (
id: string,
role: "user" | "assistant",
inputIntent?: "turn_start" | "steer",
): ThreadFeedEntry => ({
type: "message",
id,
createdAt: "2026-04-01T00:00:00.000Z",
message: {
id: MessageId.make(id),
role,
text: id,
attachments: [],
runId,
streaming: false,
...(inputIntent === undefined ? {} : { inputIntent }),
visibility: "local",
sourceThreadId: ThreadId.make("thread"),
createdAt: "2026-04-01T00:00:00.000Z",
updatedAt: "2026-04-01T00:00:00.000Z",
},
});

it("ends a response at each steer and marks only the cut-off one as mid-run", () => {
const ids = terminalFeedAssistantMessageIds([
message("prompt", "user", "turn_start"),
message("commentary", "assistant"),
message("cut-off", "assistant"),
message("steer", "user", "steer"),
message("final", "assistant"),
]);

expect([...ids.terminalIds]).toEqual(["cut-off", "final"]);
expect([...ids.midRunIds]).toEqual(["cut-off"]);
});

it("marks a response as mid-run when a steer follows it with no reply yet", () => {
const ids = terminalFeedAssistantMessageIds([
message("prompt", "user", "turn_start"),
message("cut-off", "assistant"),
message("steer", "user", "steer"),
]);

expect([...ids.terminalIds]).toEqual(["cut-off"]);
expect([...ids.midRunIds]).toEqual(["cut-off"]);
});

it("keeps one terminal message for an unsteered run", () => {
const ids = terminalFeedAssistantMessageIds([
message("prompt", "user", "turn_start"),
message("commentary", "assistant"),
message("final", "assistant"),
]);

expect([...ids.terminalIds]).toEqual(["final"]);
expect(ids.midRunIds.size).toBe(0);
});
});

describe("pending user input answers", () => {
it("preserves exact editor text, including a deliberately cleared answer", () => {
const question = { ...singleSelectQuestion, initialAnswer: " Proposed message\n" };
Expand Down
32 changes: 32 additions & 0 deletions apps/mobile/src/lib/threadActivity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1048,6 +1048,38 @@ export function failedFeedRunIds(
return failed;
}

/**
* The assistant messages that end a stretch of output and so carry metadata:
* each run's last one, plus the last one before a later user message in the
* same run (a steer). `midRunIds` holds those a user message of the same run
* follows, replied to or not; forking from one must cut inside its run.
*/
export function terminalFeedAssistantMessageIds(feed: ReadonlyArray<ThreadFeedEntry>) {
const lastBySegment = new Map<
string,
{ readonly runId: RunId; readonly segment: number; readonly id: string }
>();
const segmentByRun = new Map<RunId, number>();
for (const entry of feed) {
if (entry.type !== "message" || !entry.message.runId) continue;
const runId = entry.message.runId;
const segment = segmentByRun.get(runId) ?? 0;
if (entry.message.role === "user") {
segmentByRun.set(runId, segment + 1);
continue;
}
lastBySegment.set(`${runId}:${segment}`, { runId, segment, id: entry.message.id });
}
const terminals = [...lastBySegment.values()];
const terminalIds = new Set(terminals.map((terminal) => terminal.id));
const midRunIds = new Set(
terminals
.filter((terminal) => terminal.segment < (segmentByRun.get(terminal.runId) ?? 0))
.map((terminal) => terminal.id),
);
return { terminalIds, midRunIds };
}

/**
* A prompt without a run (a provider-native subagent, or a turn imported from
* V1) folds its response like a run. `runlessWorkActive` keeps the latest
Expand Down
2 changes: 1 addition & 1 deletion apps/server/src/mcp/toolkits/thread/tools.ts
Original file line number Diff line number Diff line change
Expand Up @@ -203,7 +203,7 @@ const transferResult = Schema.Struct({ sequence: NonNegativeInt, targetThreadId:
const ThreadForkTool = Tool.make("t3_thread_fork", {
...commandTool,
description:
"Fork a thread from a stable run or checkpoint using the existing fork command. Omit threadId to fork this thread. The fork inherits the source configuration. Acceptance does not mean a provider turn has completed.",
"Fork a thread from a stable run, a checkpoint, or one finished assistant response (turn_item, which leaves out the rest of its run such as a steer) using the existing fork command. Omit threadId to fork this thread. The fork inherits the source configuration. Acceptance does not mean a provider turn has completed.",
parameters: Schema.Struct({
threadId: Schema.optional(ThreadId),
sourcePoint: OrchestrationV2ThreadForkSourcePoint,
Expand Down
38 changes: 38 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ import {
RunAttemptId,
RunId,
ThreadId,
TurnItemId,
} from "@t3tools/contracts";
import { assert, describe, it } from "@effect/vitest";
import * as Context from "effect/Context";
Expand Down Expand Up @@ -1846,6 +1847,43 @@ describe("ClaudeAdapterV2 native fork", () => {
assert.equal(forkedProviderThread.forkedFrom?.providerThreadId, sourceProviderThread.id);
assert.equal(forkedProviderThread.forkedFrom?.providerTurnId, providerTurnId);

// A response a steer cut off is mid-turn: its assistant item carries the
// SDK message uuid, which forkSession takes as an inclusive cut.
yield* runtime.forkThread({
sourceProviderThread,
sourceProviderTurns: [],
providerTurnId,
throughTurnItem: {
id: TurnItemId.make("turn-item-claude-cut-off"),
threadId: sourceThreadId,
runId: RunId.make("run-claude-steered"),
nodeId: null,
providerThreadId: sourceProviderThread.id,
providerTurnId,
nativeItemRef: {
driver: ClaudeAdapterV2.CLAUDE_PROVIDER,
nativeId: "cut-off-assistant-uuid",
strength: "strong",
},
parentItemId: null,
ordinal: 2,
status: "completed",
title: null,
startedAt: now,
completedAt: now,
updatedAt: now,
type: "assistant_message",
messageId: MessageId.make("message-claude-cut-off"),
text: "cut off by a steer",
streaming: false,
},
targetThreadId: ThreadId.make("thread-claude-fork-steered-target"),
});
assert.deepEqual(forkCalls[1]?.options, {
dir: "/workspace",
upToMessageId: "cut-off-assistant-uuid",
});

yield* runtime.startTurn({
appThread: {
createdBy: "user",
Expand Down
20 changes: 20 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -223,6 +223,7 @@ export const ClaudeProviderCapabilitiesV2 = {
canRollbackThread: true,
canForkThread: true,
canForkFromTurn: true,
canForkFromItem: true,
canForkFromSubagentThread: false,
exposesNativeThreadId: true,
},
Expand Down Expand Up @@ -1301,6 +1302,25 @@ const getNativeConversationHeadId = Effect.fnUntraced(function* (

const resolveClaudeForkUpToMessageId = Effect.fn("ClaudeAdapterV2.resolveForkUpToMessageId")(
function* (input: ProviderAdapter.ProviderAdapterV2ForkThreadInput) {
// Assistant items carry the SDK message uuid, which forkSession accepts as
// an inclusive cut, so a fork can end before a steer in the same turn.
const throughItem = input.throughTurnItem;
if (throughItem !== undefined) {
const nativeId = throughItem.nativeItemRef?.nativeId;
if (
throughItem.type !== "assistant_message" ||
throughItem.nativeItemRef?.driver !== CLAUDE_PROVIDER ||
nativeId === undefined ||
nativeId === null
Comment thread
juliusmarminge marked this conversation as resolved.
) {
return yield* new ProviderAdapter.ProviderAdapterForkThreadError({
driver: CLAUDE_PROVIDER,
providerThreadId: input.sourceProviderThread.id,
cause: `Cannot fork Claude thread at item ${throughItem.id}: it has no SDK assistant message id.`,
});
}
return nativeId;
}
if (input.providerTurnId === undefined || input.sourceProviderTurns === undefined) {
return undefined;
}
Expand Down
10 changes: 10 additions & 0 deletions apps/server/src/orchestration-v2/CommandPolicy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,7 @@ export interface CommandPolicyV2Shape {
readonly ensureNativeFork: (
input: CapabilityCheckInput & {
readonly fromSpecificTurn: boolean;
readonly fromSpecificItem?: boolean;
},
) => Effect.Effect<void, CommandPolicyV2Error>;
readonly decideForkExecution: (
Expand All @@ -205,6 +206,8 @@ export interface CommandPolicyV2Shape {
readonly hasStrongNativeSource: boolean;
readonly sourceRunStatus: OrchestrationV2Run["status"];
readonly fromSpecificTurn: boolean;
/** The fork ends at a response inside the run, e.g. one a steer cut off. */
readonly fromSpecificItem?: boolean;
},
) => Effect.Effect<ForkExecutionPolicyV2, CommandPolicyV2Error>;
readonly ensureRollback: (
Expand Down Expand Up @@ -291,6 +294,11 @@ const ensureNativeFork: CommandPolicyV2Shape["ensureNativeFork"] = (input) => {
),
);
}
if (input.fromSpecificItem === true && input.capabilities.threads.canForkFromItem !== true) {
return Effect.fail(
unsupported(input, "fork_from_turn", "providerInstanceId cannot fork inside a turn"),
);
}
if (input.capabilities.identity.nativeThreadIds !== "strong") {
return Effect.fail(
unsupported(
Expand Down Expand Up @@ -371,6 +379,8 @@ const decideForkExecution: CommandPolicyV2Shape["decideForkExecution"] = (input)
input.hasStrongNativeSource &&
input.capabilities.threads.canForkThread &&
(!input.fromSpecificTurn || input.capabilities.threads.canForkFromTurn) &&
// A native fork that can only cut at turn ends would keep a steer made after the response.
(input.fromSpecificItem !== true || input.capabilities.threads.canForkFromItem === true) &&
input.capabilities.identity.nativeThreadIds === "strong";

if (canForkNatively) {
Expand Down
Loading
Loading