Skip to content
Merged
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
3 changes: 3 additions & 0 deletions apps/mobile/src/features/threads/ThreadDetailScreen.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,8 @@ export interface ThreadDetailScreenProps {
readonly selectedThreadFeed: ReadonlyArray<ThreadFeedEntry>;
readonly activityRun: ThreadFeedLatestRun | null;
readonly activeWorkStartedAt: string | null;
/** The live work is a provider-native subagent's runless root turn. */
readonly runlessWorkActive?: boolean;
readonly isCompacting: boolean;
/**
* The server has not created this thread yet. "preparing" runs while the
Expand Down Expand Up @@ -1031,6 +1033,7 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread
threadTitle={props.selectedThread.title}
latestRun={props.activityRun}
activeWorkStartedAt={props.activeWorkStartedAt}
runlessWorkActive={props.runlessWorkActive ?? false}
listRef={listRef}
freeze={freeze}
anchorMessageId={anchorMessageId}
Expand Down
3 changes: 3 additions & 0 deletions apps/mobile/src/features/threads/ThreadFeed.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -277,6 +277,7 @@ export interface ThreadFeedProps {
readonly agentLabel: string;
readonly latestRun: ThreadFeedLatestRun | null;
readonly activeWorkStartedAt: string | null;
readonly runlessWorkActive?: boolean;
readonly listRef: RefObject<LegendListRef | null>;
readonly freeze: SharedValue<boolean>;
readonly anchorMessageId: MessageId | null;
Expand Down Expand Up @@ -2606,6 +2607,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) {
.map(([groupId]) => groupId),
),
props.activeWorkStartedAt,
props.runlessWorkActive ?? false,
),
props.feed,
props.queuedMessages,
Expand All @@ -2615,6 +2617,7 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) {
expandedTurnIds,
expandedWorkGroups,
props.activeWorkStartedAt,
props.runlessWorkActive,
props.feed,
props.latestRun,
],
Expand Down
1 change: 1 addition & 0 deletions apps/mobile/src/features/threads/ThreadRouteScreen.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -1005,6 +1005,7 @@ function ThreadRouteContent(
: composer.activeWorkStartedAt
}
isCompacting={composer.isCompacting}
runlessWorkActive={composer.runlessWorkActive}
creationState={creationState}
setupWorkingStartedAt={
composer.activeWorkStartedAt !== null &&
Expand Down
56 changes: 56 additions & 0 deletions apps/mobile/src/lib/threadActivity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -956,6 +956,62 @@ describe("buildThreadFeed", () => {
});
});

it("keeps a provider-native subagent's runless tool call live while it works", () => {
const startedAt = "2026-06-20T00:00:01.000Z";
const { exitCode: _exitCode, ...completedCommand } = command();
const runningCommand: OrchestrationV2TurnItem = {
...completedCommand,
runId: null,
status: "running",
completedAt: null,
output: "",
};
const feed = buildThreadFeed([
projected({ ...userMessage(), runId: null }, 0),
projected(runningCommand, 1),
]);

const presented = deriveThreadFeedPresentation(
feed,
null,
new Set(),
new Set(),
startedAt,
true,
);
expect(presented.find((entry) => entry.type === "work-toggle")).toMatchObject({
summary: "Running vp",
live: true,
shimmer: true,
});
expect(presented.some((entry) => entry.type === "thinking")).toBe(false);
});

it("keeps a runless tail settled while a normal thread waits for its sent run", () => {
// Right after a send the local clock runs before the server creates the
// run, and the latest run may still be queued: neither is runless work.
const startedAt = "2026-06-20T00:00:05.000Z";
const feed = buildThreadFeed([
projected({ ...userMessage(), runId: null }, 0),
projected({ ...command(), runId: null }, 1),
]);
for (const latestRun of [
null,
{ runId, status: "queued" as const, startedAt: null, completedAt: null },
]) {
const presented = deriveThreadFeedPresentation(
feed,
latestRun,
new Set(),
new Set(),
startedAt,
);
const toggle = presented.find((entry) => entry.type === "work-toggle");
expect(toggle).toMatchObject({ live: false, shimmer: false });
expect(presented.at(-1)?.type).toBe("thinking");
}
});

it("waits for workspace preparation before showing provider activity", () => {
const startedAt = "2026-04-01T00:00:01.000Z";
const run = { runId, status: "preparing" as const, startedAt: null, completedAt: null };
Expand Down
6 changes: 5 additions & 1 deletion apps/mobile/src/lib/threadActivity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1019,6 +1019,8 @@ export function deriveThreadFeedPresentation(
expandedRunIds: ReadonlySet<RunId>,
expandedWorkGroupIds: ReadonlySet<string> = new Set(),
activeWorkStartedAt: string | null = null,
/** The live work is a provider-native subagent's runless root turn. */
runlessWorkActive = false,
): ThreadFeedEntry[] {
const sourceFeed = feed.filter(
(entry) =>
Expand All @@ -1037,9 +1039,11 @@ export function deriveThreadFeedPresentation(
}
const result: ThreadFeedEntry[] = [];
for (const entry of sourceFeed) {
// A provider-native subagent works without a run: its null-run tail is
// live only while that runless work is active.
const isActiveTailGroup =
isWorking &&
activeRunId !== null &&
(activeRunId !== null || runlessWorkActive) &&
entry.type === "activity-group" &&
activeTailGroup?.type === "activity-group" &&
activeTailGroup.id === entry.id &&
Expand Down
22 changes: 17 additions & 5 deletions apps/mobile/src/state/use-thread-composer-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import type { ComposerTextPaste } from "../native/T3ComposerEditor.types";
import { useAtomValue } from "@effect/atom-react";
import { threadRuntimeIsActive } from "@t3tools/client-runtime/state/shell";
import {
deriveRunlessWorkStartedAt,
deriveThreadActivityRun,
deriveThreadRuntime,
threadRuntimeHasInterruptibleRun,
Expand Down Expand Up @@ -393,15 +394,25 @@ export function useThreadComposerState() {
selectedThreadVisibleTurnItems,
]);

const runlessWorkStartedAt = useMemo(
() =>
selectedThreadProjection
? deriveRunlessWorkStartedAt(selectedThreadProjection.projection)
: null,
[selectedThreadProjection],
);
const activeWorkStartedAt = useMemo(() => {
if (!selectedThreadShell) {
return null;
}
return resolveThreadWorkingStartedAt({
latestRun: selectedThreadActivityRun,
runtime: selectedThreadRuntime,
});
}, [selectedThreadActivityRun, selectedThreadRuntime, selectedThreadShell]);
return (
resolveThreadWorkingStartedAt({
latestRun: selectedThreadActivityRun,
runtime: selectedThreadRuntime,
}) ?? runlessWorkStartedAt
);
}, [selectedThreadActivityRun, runlessWorkStartedAt, selectedThreadRuntime, selectedThreadShell]);
const runlessWorkActive = runlessWorkStartedAt !== null;

// The run can start, or be cancelled from another client, while its message
// is open in the composer. Leave edit mode rather than saving into a run the
Expand Down Expand Up @@ -1024,6 +1035,7 @@ export function useThreadComposerState() {
selectedThreadQueuedMessages,
dispatchingQueuedMessageId,
activeWorkStartedAt,
runlessWorkActive,
isCompacting,
draftMessage,
draftAttachments,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import { provideDeterministicTestRuntime } from "./DeterministicRuntime.ts";
import { ORCHESTRATOR_REPLAY_FIXTURES } from "./fixtures/index.ts";
import { messageRestartInput } from "./fixtures/message_steering/input.ts";
import {
assertProviderNativeSubagentRootTurns,
materializeFixtureInput,
type OrchestratorFixtureInput,
type ProviderOrchestratorReplayVariant,
Expand Down Expand Up @@ -117,6 +118,7 @@ const runFixtureProvider = Effect.fn("runOrchestratorReplayFixture")(function* <
input.driver.runContinuationWorker === true ? { runContinuationWorker: true } : {},
).pipe(provideDeterministicTestRuntime);
input.driver.assertOutput(result, transcript);
assertProviderNativeSubagentRootTurns(result);
const expectedAbsentWorkspacePaths = input.driver.expectedAbsentWorkspacePaths;
if (expectedAbsentWorkspacePaths !== undefined) {
yield* Effect.gen(function* () {
Expand Down
61 changes: 61 additions & 0 deletions apps/server/src/orchestration-v2/testkit/fixtures/shared.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { assert } from "@effect/vitest";
import {
type ChatAttachment,
CommandId,
isOrchestrationV2WorkActive,
MessageId,
ProjectId,
ThreadId,
Expand Down Expand Up @@ -1044,6 +1045,66 @@ export function assertNoExtraAppRunsForProviderChildren(input: {
);
}

/**
* Provider-native subagent threads have no runs; clients show them working
* from the child's runless root turn. Pin that contract for every recorded
* native subagent: the child hangs off the subagent node, every root turn is
* runless, the root turn is live before the child's first item, and its
* activity mirrors the subagent's (including a resume re-opening it).
*/
export function assertProviderNativeSubagentRootTurns(result: OrchestratorV2ScenarioResult) {
const activity = (statuses: ReadonlyArray<OrchestrationV2ExecutionNode["status"]>) =>
statuses
.map((status) => (isOrchestrationV2WorkActive(status) ? "active" : status))
.filter((status, index, all) => status !== all[index - 1]);
for (const projection of result.projections.values()) {
for (const subagent of projection.subagents) {
if (subagent.origin !== "provider_native" || subagent.childThreadId === null) continue;
const childThreadId = subagent.childThreadId;
const child = result.projections.get(childThreadId);
assert.isDefined(child, `missing child thread for subagent ${subagent.id}`);
assert.equal(child.thread.creationSource, "provider");
assert.deepEqual(child.thread.forkedFrom, { type: "node", nodeId: subagent.id });
assert.lengthOf(child.runs, 0);
const roots = child.nodes.filter((node) => node.kind === "root_turn");
assert.isNotEmpty(roots, `child ${childThreadId} must have a root turn`);
for (const root of roots) assert.isNull(root.runId);

const rootEvents = result.domainEvents.flatMap((event, index) =>
event.type === "node.updated" &&
event.payload.threadId === childThreadId &&
event.payload.kind === "root_turn"
? [{ index, status: event.payload.status }]
: [],
);
const firstItemIndex = result.domainEvents.findIndex(
(event) =>
event.type === "turn-item.updated" &&
event.payload.threadId === childThreadId &&
event.payload.type !== "user_message",
);
assert.equal(rootEvents[0]?.status, "running");
if (firstItemIndex !== -1) {
assert.isBelow(
rootEvents[0]?.index ?? Infinity,
firstItemIndex,
`child ${childThreadId} must be working before its first item`,
);
}
const subagentStatuses = result.domainEvents.flatMap((event) =>
event.type === "subagent.updated" && event.payload.id === subagent.id
? [event.payload.status]
: [],
);
assert.deepEqual(
activity(rootEvents.map((event) => event.status)),
activity(subagentStatuses),
`child ${childThreadId} root turn must follow subagent ${subagent.id}`,
);
}
}
}

export function assertExecutionNodeKinds(
projection: OrchestrationV2ThreadProjection,
expectedKinds: ReadonlyArray<OrchestrationV2ExecutionNode["kind"]>,
Expand Down
21 changes: 15 additions & 6 deletions apps/web/src/components/ChatView.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ import { isPasteAsTextShortcut } from "@t3tools/client-runtime/text-paste";
import { effectiveSnoozed, threadWokeAt } from "@t3tools/client-runtime/state/thread-settled";
import { useThreadActions } from "../hooks/useThreadActions";
import {
deriveRunlessWorkStartedAt,
deriveThreadActivityRun,
deriveLatestThreadRun,
deriveThreadRuntime,
Expand Down Expand Up @@ -2003,6 +2004,10 @@ export default function ChatView(props: ChatViewProps) {
() => (serverProjection === null ? null : deriveThreadRuntime(serverProjection)),
[serverProjection],
);
const runlessWorkStartedAt = useMemo(
() => (serverProjection === null ? null : deriveRunlessWorkStartedAt(serverProjection)),
[serverProjection],
);
const supportsProviderSwitchingViaHandoff = useMemo(
() => threadSupportsProviderHandoff(serverProjection),
[serverProjection],
Expand Down Expand Up @@ -3414,7 +3419,12 @@ export default function ChatView(props: ChatViewProps) {
compactRequestIsActive &&
!compactionSettled;
const isWorking =
phase === "running" || isSendBusy || isConnecting || isRevertingCheckpoint || isCompacting;
phase === "running" ||
isSendBusy ||
isConnecting ||
isRevertingCheckpoint ||
isCompacting ||
runlessWorkStartedAt !== null;
const activeContextWindow = useMemo(
() =>
deriveLatestContextWindowSnapshot(
Expand Down Expand Up @@ -3446,11 +3456,9 @@ export default function ChatView(props: ChatViewProps) {
}),
];
}, [serverProjection]);
const activeWorkStartedAt = deriveActiveWorkStartedAt(
activeActivityRun,
activeRuntime,
localDispatchStartedAt,
);
const activeWorkStartedAt =
deriveActiveWorkStartedAt(activeActivityRun, activeRuntime, localDispatchStartedAt) ??
runlessWorkStartedAt;
// Server-side workspace preparation: unlike the local-dispatch flag this
// survives reloads and shows on remote viewers of the same thread.
const activeRunPreparing = activeActivityRun?.status === "preparing";
Expand Down Expand Up @@ -10445,6 +10453,7 @@ export default function ChatView(props: ChatViewProps) {
}
: {})}
isWorking={!paintOnlyDisplayedTimeline && isWorking}
runlessWorkActive={runlessWorkStartedAt !== null}
activeTurnInProgress={
!paintOnlyDisplayedTimeline && (isWorking || !latestRunSettled)
}
Expand Down
Loading
Loading