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
1 change: 1 addition & 0 deletions apps/mobile/src/features/threads/NewTaskDraftScreen.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -454,6 +454,7 @@ export function NewTaskDraftScreen(props: {
projectCwd: composerWorkspaceCwd,
selectedProviderStatus: flow.selectedProviderStatus,
hasThread: false,
hasSession: false,
hasCompactableConversation: false,
offersUsageLimits: offersUsageLimits,
enabled: isComposerFocused && !isComposerInteractionLocked,
Expand Down
1 change: 1 addition & 0 deletions apps/mobile/src/features/threads/ThreadComposer.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -383,6 +383,7 @@ export const ThreadComposer = memo(function ThreadComposer(props: ThreadComposer
pullRequestRepository: project?.repositoryIdentity?.displayName ?? null,
selectedProviderStatus,
hasThread: true,
hasSession: props.selectedThread.session !== null,
hasCompactableConversation: props.hasCompactableConversation,
onChangeDraftMessage: props.onChangeDraftMessage,
onUpdateInteractionMode:
Expand Down
40 changes: 40 additions & 0 deletions apps/mobile/src/features/threads/use-composer-command-menu.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ describe("mobile slash commands", () => {
query: "pl",
atMessageStart: true,
hasThread: true,
hasSession: true,
allowInteractionMode,
selectedProviderStatus: antigravity,
});
Expand All @@ -62,6 +63,7 @@ describe("mobile slash commands", () => {
query: "plan",
atMessageStart: false,
hasThread: false,
hasSession: false,
allowInteractionMode: true,
selectedProviderStatus: antigravity,
}),
Expand All @@ -73,6 +75,7 @@ describe("mobile slash commands", () => {
query: "plan",
atMessageStart: true,
hasThread: true,
hasSession: true,
allowInteractionMode: true,
selectedProviderStatus: {
driver: ProviderDriverKind.make("codex"),
Expand Down Expand Up @@ -100,4 +103,41 @@ describe("mobile slash commands", () => {
}),
).toEqual({ text: "/plan ", cursor: 6, interactionMode: null });
});

it("offers MCP reconnect only for threads with a session", () => {
const claude = {
driver: ProviderDriverKind.make("claudeAgent"),
slashCommands: [],
};
expect(
buildComposerSlashCommandItems({
query: "reconnect",
atMessageStart: true,
hasThread: true,
hasSession: true,
allowInteractionMode: false,
selectedProviderStatus: claude,
}).map((item) => item.id),
).toEqual(["cmd:reconnect-mcp"]);
expect(
buildComposerSlashCommandItems({
query: "reconnect",
atMessageStart: true,
hasThread: true,
hasSession: false,
allowInteractionMode: false,
selectedProviderStatus: claude,
}),
).toEqual([]);
expect(
buildComposerSlashCommandItems({
query: "reconnect",
atMessageStart: true,
hasThread: false,
hasSession: false,
allowInteractionMode: false,
selectedProviderStatus: claude,
}),
).toEqual([]);
});
});
21 changes: 20 additions & 1 deletion apps/mobile/src/features/threads/use-composer-command-menu.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,8 @@ export function buildComposerSlashCommandItems(input: {
readonly query: string;
readonly atMessageStart: boolean;
readonly hasThread: boolean;
/** Whether the thread has a provider session (false for threads with no message yet). */
readonly hasSession: boolean;
readonly hasCompactableConversation?: boolean;
/** Whether T3 itself offers /usage-limits for the selected provider. */
readonly offersUsageLimits?: boolean;
Expand Down Expand Up @@ -87,9 +89,22 @@ export function buildComposerSlashCommandItems(input: {
label: "/default",
description: "Switch to default mode",
},
...(input.hasThread && input.hasSession
? [
{
id: "cmd:reconnect-mcp",
type: "slash-command" as const,
command: "reconnect-mcp",
label: "/reconnect-mcp",
description: "Reconnect MCP servers before the next turn",
},
]
: []),
] satisfies ComposerCommandItem[];
const items: ComposerCommandItem[] = builtIn.filter(
(item) => item.command.includes(query) && (item.command === "model" || allowInteractionMode),
(item) =>
item.command.includes(query) &&
(item.command === "model" || item.command === "reconnect-mcp" || allowInteractionMode),
);

// Providers expand commands only at the start of a message. T3 commands
Expand Down Expand Up @@ -169,6 +184,7 @@ export function useComposerCommandMenu({
pullRequestRepository = null,
selectedProviderStatus,
hasThread,
hasSession,
hasCompactableConversation,
offersUsageLimits = false,
enabled = true,
Expand All @@ -184,6 +200,7 @@ export function useComposerCommandMenu({
readonly pullRequestRepository?: string | null;
readonly selectedProviderStatus: ServerProvider | null;
readonly hasThread: boolean;
readonly hasSession: boolean;
readonly hasCompactableConversation: boolean;
/** Whether T3 itself offers /usage-limits for the selected provider. */
readonly offersUsageLimits?: boolean;
Expand Down Expand Up @@ -336,6 +353,7 @@ export function useComposerCommandMenu({
query: q,
atMessageStart: trigger.rangeStart === 0,
hasThread,
hasSession,
hasCompactableConversation,
offersUsageLimits,
allowInteractionMode: onUpdateInteractionMode !== undefined,
Expand Down Expand Up @@ -463,6 +481,7 @@ export function useComposerCommandMenu({
return [];
}, [
hasThread,
hasSession,
hasCompactableConversation,
onUpdateInteractionMode,
pathSearch.entries,
Expand Down
52 changes: 52 additions & 0 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 { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { Alert } from "react-native";
import * as Cause from "effect/Cause";

import {
CommandId,
Expand All @@ -16,9 +17,11 @@ import {
type ThreadId,
} from "@t3tools/contracts";
import { safeErrorLogAttributes } from "@t3tools/client-runtime/errors";
import { isAtomCommandInterrupted } from "@t3tools/client-runtime/state/runtime";
import { clampFileAttachmentUploadBytes } from "@t3tools/client-runtime/state/attachments";
import { nextPastedTextFileName, pastedTextDisposition } from "@t3tools/client-runtime/text-paste";
import {
isReconnectMcpCommand,
parseCodexFeedbackCommand,
submitCodexFeedback,
type CodexFeedbackSubmission,
Expand Down Expand Up @@ -140,6 +143,9 @@ export function useThreadComposerState() {
const uploadThreadFeedback = useAtomCommand(threadEnvironment.uploadFeedback, {
reportFailure: false,
});
const stopThreadSession = useAtomCommand(threadEnvironment.stopSession, {
reportFailure: false,
});
const pastedTextFileNamesRef = useRef<{ threadKey: string | null; names: Set<string> }>({
threadKey: null,
names: new Set(),
Expand Down Expand Up @@ -387,6 +393,51 @@ export function useThreadComposerState() {
const provider = serverConfig?.providers.find(
(entry) => entry.instanceId === modelSelection.instanceId,
);
const reconnectMcpCommand =
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
attachments.length === 0 &&
(draft.context?.records.length ?? 0) === 0 &&
isReconnectMcpCommand(text);
if (reconnectMcpCommand) {
if (thread.session === null) {
Alert.alert("Start a thread first", "Send a message before reconnecting its MCP servers.");
return null;
}
if (thread.session.status === "starting" || thread.session.status === "running") {
Alert.alert(
"Agent is still working",
"Wait for the current turn to finish, then reconnect MCP servers.",
);
return null;
}
clearComposerDraftContent(threadKey);
const result = await stopThreadSession({
environmentId: selectedThreadShell.environmentId,
input: { threadId: selectedThreadShell.id, onlyIfIdle: true },
});
if (result._tag === "Failure") {
const currentDraft = getComposerDraftSnapshot(threadKey);
if (
currentDraft.text.length === 0 &&
currentDraft.attachments.length === 0 &&
(currentDraft.context?.records.length ?? 0) === 0
) {
setComposerDraftText(threadKey, text);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if (!isAtomCommandInterrupted(result)) {
const error = Cause.squash(result.cause);
Alert.alert(
"Could not reconnect MCP servers",
error instanceof Error ? error.message : "An error occurred.",
);
}
return null;
}
Alert.alert(
"MCP reconnect requested",
"Wait for the session to stop, then send a message to reconnect with fresh tool schemas. Any stop failure appears in the thread.",
);
return null;
}
const feedbackCommand =
attachments.length === 0 &&
(provider?.driver === "codex" || thread.session?.providerName === "codex")
Expand Down Expand Up @@ -483,6 +534,7 @@ export function useThreadComposerState() {
selectedThreadCreation,
selectedThreadDetail,
selectedThreadShell,
stopThreadSession,
uploadThreadFeedback,
]);

Expand Down
96 changes: 96 additions & 0 deletions apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -567,6 +567,102 @@ describe("OrchestrationEngine", () => {
}).pipe(Effect.provide(makeOrchestrationLayer())),
);

effectIt.effect("idle-only stops preserve background work until it finishes", () =>
Effect.gen(function* () {
yield* TestClock.setTime(Date.parse(now()));
const engine = yield* OrchestrationEngineService;
const backgroundLiveness = yield* ThreadBackgroundLiveness.ThreadBackgroundLivenessService;
const projectId = ProjectId.make("project-reconnect");
const threadId = ThreadId.make("thread-reconnect");
const unrelatedThreadId = ThreadId.make("thread-unrelated");
yield* engine.dispatch({
type: "project.create",
commandId: CommandId.make("create-reconnect-project"),
projectId,
title: "Project",
workspaceRoot: "/tmp/project-reconnect",
createdAt: now(),
});
for (const id of [threadId, unrelatedThreadId]) {
yield* engine.dispatch({
type: "thread.create",
commandId: CommandId.make(`create-${id}`),
threadId: id,
projectId,
title: "Thread",
modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "gpt-5-codex" },
interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
runtimeMode: "full-access",
branch: null,
worktreePath: null,
createdAt: now(),
});
}
yield* engine.dispatch({
type: "thread.session.set",
commandId: CommandId.make("ready-reconnect-session"),
threadId,
session: {
threadId,
status: "ready",
providerName: "codex",
runtimeMode: "full-access",
activeTurnId: null,
lastError: null,
updatedAt: now(),
},
createdAt: now(),
});
for (const taskType of ["subagent", "local_bash"]) {
backgroundLiveness.recordTaskLiveness({
threadId,
taskId: taskType,
taskType,
kind: "started",
status: "running",
...(taskType === "local_bash" ? { agentId: "subagent-owner" } : {}),
});
const sequence = yield* engine.latestSequence;
const error = yield* engine
.dispatch({
type: "thread.session.stop",
commandId: CommandId.make(`reject-reconnect-${taskType}`),
threadId,
onlyIfIdle: true,
createdAt: now(),
})
.pipe(Effect.flip);
expect(error).toMatchObject({
_tag: "OrchestrationCommandInvariantError",
detail: `thread ${threadId} has live background work`,
});
expect(yield* engine.latestSequence).toBe(sequence);
yield* engine.dispatch({
type: "thread.session.stop",
commandId: CommandId.make(`stop-unrelated-${taskType}`),
threadId: unrelatedThreadId,
onlyIfIdle: true,
createdAt: now(),
});
backgroundLiveness.recordTaskLiveness({
threadId,
taskId: taskType,
taskType,
kind: "completed",
status: "completed",
...(taskType === "local_bash" ? { agentId: "subagent-owner" } : {}),
});
}
yield* engine.dispatch({
type: "thread.session.stop",
commandId: CommandId.make("reconnect-after-background-finished"),
threadId,
onlyIfIdle: true,
createdAt: now(),
});
}).pipe(Effect.provide(makeOrchestrationLayer())),
);

effectIt.effect(
"rejects persisted changes and live background work without blocking unrelated threads",
() =>
Expand Down
4 changes: 3 additions & 1 deletion apps/server/src/orchestration/Layers/OrchestrationEngine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -203,7 +203,9 @@ const makeOrchestrationEngine = Effect.gen(function* () {
}

if (
envelope.command.type === "thread.auto-settle" &&
(envelope.command.type === "thread.auto-settle" ||
(envelope.command.type === "thread.session.stop" &&
envelope.command.onlyIfIdle === true)) &&
threadBackgroundLiveness.getThreadBackgroundLiveness(envelope.command.threadId) !== null
) {
return yield* new OrchestrationCommandInvariantError({
Expand Down
16 changes: 13 additions & 3 deletions apps/server/src/orchestration/ThreadBackgroundLiveness.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,7 +139,7 @@ describe("ThreadBackgroundLiveness", () => {
expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull();
});

it("untyped rows count as agents; idle is not live; agent-owned tasks are ignored", () => {
it("untyped rows count as agents; idle is not live; agent-owned tasks remain live", () => {
const liveness = ThreadBackgroundLiveness.make();
const threadId = "t-live-3";
liveness.recordTaskLiveness({
Expand All @@ -166,6 +166,15 @@ describe("ThreadBackgroundLiveness", () => {
kind: "started",
agentId: "owner",
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("monitoring");
liveness.recordTaskLiveness({
threadId,
taskId: "sh:1",
taskType: "local_bash",
status: "completed",
kind: "completed",
agentId: "owner",
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull();
});

Expand All @@ -191,7 +200,8 @@ describe("ThreadBackgroundLiveness", () => {
kind: "progress",
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("monitoring");
// Turning out to be inert or agent-owned drops the prior entry too.
// Turning out to be inert drops the prior entry too; an agent-owned shell
// remains independently live until its own terminal transition.
liveness.recordTaskLiveness({
threadId,
taskId: "x1",
Expand All @@ -200,7 +210,7 @@ describe("ThreadBackgroundLiveness", () => {
kind: "progress",
agentId: "owner",
});
expect(liveness.getThreadBackgroundLiveness(threadId)).toBeNull();
expect(liveness.getThreadBackgroundLiveness(threadId)).toBe("monitoring");
});

it("plan tasks are inert; clear removes everything; instances are isolated", () => {
Expand Down
Loading
Loading