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
30 changes: 30 additions & 0 deletions apps/mobile/src/lib/followUpBehavior.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
import { describe, expect, it } from "vite-plus/test";

import { resolveFollowUpDispatchMode } from "./followUpBehavior";

describe("resolveFollowUpDispatchMode", () => {
const running = {
running: true,
canSteer: true,
isCompacting: false,
followUpBehavior: "steer",
} as const;

it("lets the server decide while the thread is idle", () => {
expect(resolveFollowUpDispatchMode({ ...running, running: false })).toBeNull();
});

it("sends steering as auto and queueing as queue", () => {
expect(resolveFollowUpDispatchMode(running)).toBe("auto");
expect(resolveFollowUpDispatchMode({ ...running, followUpBehavior: "queue" })).toBe("queue");
expect(resolveFollowUpDispatchMode({ ...running, followUpOverride: "queue" })).toBe("queue");
});

it("queues explicitly while compacting, matching the Queue label", () => {
// While /compact is still being dispatched the active run is the ordinary
// one, which an auto send would steer.
expect(resolveFollowUpDispatchMode({ ...running, canSteer: false, isCompacting: true })).toBe(
"queue",
);
});
});
30 changes: 29 additions & 1 deletion apps/mobile/src/lib/followUpBehavior.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,7 @@
import type { ActiveTurnComposerAction } from "@t3tools/client-runtime/state/composer-dispatch";
import {
resolveComposerDispatchMode,
type ActiveTurnComposerAction,
} from "@t3tools/client-runtime/state/composer-dispatch";

/**
* What the send button does while a turn is already running: `queue` waits for
Expand All @@ -11,3 +14,28 @@ import type { ActiveTurnComposerAction } from "@t3tools/client-runtime/state/com
export type FollowUpBehavior = Extract<ActiveTurnComposerAction, "queue" | "steer">;

export const DEFAULT_FOLLOW_UP_BEHAVIOR: FollowUpBehavior = "queue";

/**
* The outbox dispatch mode for a send, or null when the thread is idle and the
* server decides. Steering travels as "auto" so a turn that ends before the
* outbox delivers degrades to a queued run instead of failing the delivery and
* bouncing the message back into the draft.
*/
export function resolveFollowUpDispatchMode(input: {
readonly running: boolean;
readonly canSteer: boolean;
readonly isCompacting: boolean;
readonly followUpBehavior: FollowUpBehavior;
readonly followUpOverride?: ActiveTurnComposerAction;
}): "queue" | "auto" | null {
const action = resolveComposerDispatchMode({
// Compaction queues explicitly: while /compact is still being dispatched
// the active run is the ordinary one, which an auto send would steer.
running: input.running && (input.canSteer || input.isCompacting),
alternateModifier:
input.followUpOverride !== undefined && input.followUpOverride !== input.followUpBehavior,
activeTurnDefault: input.followUpBehavior,
activeTurnIsCompaction: input.isCompacting,
});
return action === "auto" ? null : action === "queue" ? "queue" : "auto";
}
27 changes: 12 additions & 15 deletions apps/mobile/src/state/use-thread-composer-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,14 +77,11 @@ import {
updateComposerDraftSettings,
useComposerDraft,
} from "./use-composer-drafts";
import {
resolveComposerDispatchMode,
type ActiveTurnComposerAction,
} from "@t3tools/client-runtime/state/composer-dispatch";
import { type ActiveTurnComposerAction } from "@t3tools/client-runtime/state/composer-dispatch";
import { Atom } from "effect/unstable/reactivity";
import { AsyncResult } from "effect/unstable/reactivity";
import { prepareTurnAttachments } from "../lib/attachmentUpload";
import { DEFAULT_FOLLOW_UP_BEHAVIOR } from "../lib/followUpBehavior";
import { DEFAULT_FOLLOW_UP_BEHAVIOR, resolveFollowUpDispatchMode } from "../lib/followUpBehavior";
import { mobilePreferencesAtom } from "./preferences";
import { environmentThreadDetails } from "./threads";
import {
Expand Down Expand Up @@ -311,7 +308,6 @@ export function useThreadComposerState() {
threadId: selectedThreadShell.id,
}),
);
const canSteerActiveTurn = queueWorkflow?.canPromoteToSteer === true;
const queuedRunEdit = useQueuedRunEdit(selectedThreadKey);
const composerDraftKey =
selectedThreadKey === null
Expand Down Expand Up @@ -394,6 +390,9 @@ export function useThreadComposerState() {
selectedThreadRuntime,
selectedThreadVisibleTurnItems,
]);
// Compaction runs cannot take a steer, so the composer labels and sends
// follow-ups as queued behind it, matching the server's dispatch policy.
const canSteerActiveTurn = queueWorkflow?.canPromoteToSteer === true && !isCompacting;
Comment thread
saphid marked this conversation as resolved.

const runlessWorkStartedAt = useMemo(
() =>
Expand Down Expand Up @@ -657,16 +656,13 @@ export function useThreadComposerState() {

// Resolved here rather than at drain time: the outbox can deliver minutes
// later, and the choice belongs to the moment the user pressed send.
// Steering travels as "auto" so a turn that ends in the meantime degrades
// to a queued run on the server instead of failing the delivery and
// bouncing the message back into the draft.
const followUpAction = resolveComposerDispatchMode({
running: activeThreadBusy && canSteerActiveTurn,
alternateModifier: followUpOverride !== undefined && followUpOverride !== followUpBehavior,
activeTurnDefault: followUpBehavior,
const followUpDispatchMode = resolveFollowUpDispatchMode({
running: activeThreadBusy,
canSteer: canSteerActiveTurn,
isCompacting,
followUpBehavior,
...(followUpOverride === undefined ? {} : { followUpOverride }),
});
const followUpDispatchMode =
followUpAction === "auto" ? null : followUpAction === "queue" ? "queue" : "auto";

const metadata = makeQueuedMessageMetadata();
const messageId = MessageId.make(metadata.messageId);
Expand Down Expand Up @@ -717,6 +713,7 @@ export function useThreadComposerState() {
activeThreadBusy,
canSteerActiveTurn,
followUpBehavior,
isCompacting,
saveQueuedRunEdit,
selectedEnvironmentRuntime?.connectionState,
selectedEnvironmentRuntime?.serverConfig,
Expand Down
91 changes: 91 additions & 0 deletions apps/server/src/orchestration-v2/CommandPolicy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import {
CommandId,
type OrchestrationV2ProviderCapabilities,
type OrchestrationV2ThreadProjection,
MessageId,
ProviderInstanceId,
ProviderSessionId,
ProviderThreadId,
Expand Down Expand Up @@ -34,6 +35,7 @@ function dispatchProjection(
const providerThreadId = ProviderThreadId.make("command-policy-provider-thread");
const providerSessionId = ProviderSessionId.make("command-policy-provider-session");
return {
messages: [],
runs:
sessionCapabilities === undefined
? []
Expand Down Expand Up @@ -141,6 +143,95 @@ it("targets the latest active run for explicit steer and restart intent", () =>
);
});

it.each(["/compact", " /COMPACT ", "/logout"])(
"queues automatic follow-ups behind %s and preserves explicit intent",
(text) => {
const messageId = MessageId.make("command-policy-compaction-message");
const base = dispatchProjection(baseCapabilities);
const projection = {
...base,
runs: base.runs.map((run) => ({ ...run, userMessageId: messageId })),
messages: [{ id: messageId, role: "user", text, attachments: [] }],
} as unknown as OrchestrationV2ThreadProjection;
const queueAfterActive = { type: "queue_after_active" };
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(projection, { type: "start_immediately" }, "auto"),
queueAfterActive,
);
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(
projection,
{ type: "start_immediately" },
"steer",
),
{ type: "steer_active", targetRunId: activeRunId },
);
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(
projection,
{ type: "start_immediately" },
"restart",
),
{ type: "restart_active", targetRunId: activeRunId },
);
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(projection, {
type: "steer_active",
targetRunId: activeRunId,
}),
{ type: "steer_active", targetRunId: activeRunId },
);
},
);

it.each(["/Logout", "/LOGOUT"])("preserves ordinary delivery for %s", (text) => {
const messageId = MessageId.make("command-policy-mixed-case-message");
const base = dispatchProjection(baseCapabilities);
const projection = {
...base,
runs: base.runs.map((run) => ({ ...run, userMessageId: messageId })),
messages: [{ id: messageId, role: "user", text, attachments: [] }],
} as unknown as OrchestrationV2ThreadProjection;
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(projection, { type: "start_immediately" }, "auto"),
{ type: "steer_active", targetRunId: activeRunId },
);
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(projection, { type: "start_immediately" }, "steer"),
{ type: "steer_active", targetRunId: activeRunId },
);
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(
projection,
{ type: "start_immediately" },
"restart",
),
{ type: "restart_active", targetRunId: activeRunId },
);
});

it("keeps steering an ordinary run alongside a compacted message in history", () => {
const messageId = MessageId.make("command-policy-ordinary-message");
const base = dispatchProjection(baseCapabilities);
const projection = {
...base,
runs: base.runs.map((run) => ({ ...run, userMessageId: messageId })),
messages: [
{ id: messageId, role: "user", text: "ship it", attachments: [] },
{
id: MessageId.make("command-policy-earlier-compaction"),
role: "user",
text: "/compact",
attachments: [],
},
],
} as unknown as OrchestrationV2ThreadProjection;
assert.deepEqual(
CommandPolicy.resolveMessageDispatchIntent(projection, { type: "start_immediately" }, "auto"),
{ type: "steer_active", targetRunId: activeRunId },
);
});

const layer = it.layer(CommandPolicy.layer);

layer("CommandPolicyV2", (it) => {
Expand Down
40 changes: 38 additions & 2 deletions apps/server/src/orchestration-v2/CommandPolicy.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import {
ChatAttachment,
CommandId,
ModelSelection,
type OrchestrationV2Command,
Expand Down Expand Up @@ -116,6 +117,35 @@ type MessageDispatchMode = Extract<
{ readonly type: "message.dispatch" }
>["dispatchMode"];

/** Native commands the provider executes as its own whole-turn task, so they
* own the run they start and cannot take steering or a restart. */
export function isNativeMaintenanceCommand(message: {
readonly text: string;
readonly attachments: ReadonlyArray<ChatAttachment>;
}): boolean {
return (
message.attachments.length === 0 &&
(message.text.trim().toLowerCase() === "/compact" || message.text.trim() === "/logout")
);
}

/** Queue automatic follow-ups behind maintenance. Explicit steer and restart
* keep their intent and receive the orchestrator's unavailable error. */
function queueBehindMaintenanceRun(
projection: OrchestrationV2ThreadProjection,
decision: MessageDispatchMode,
): MessageDispatchMode {
if (decision.type !== "steer_active" && decision.type !== "restart_active") return decision;
const targetRun = projection.runs.find((run) => run.id === decision.targetRunId);
const targetMessage =
targetRun === undefined
? undefined
: projection.messages.find((message) => message.id === targetRun.userMessageId);
return targetMessage !== undefined && isNativeMaintenanceCommand(targetMessage)
? { type: "queue_after_active" }
: decision;
}

/** Resolve client intent from the state serialized by the thread dispatch lock. */
export function resolveMessageDispatchIntent(
projection: OrchestrationV2ThreadProjection,
Expand Down Expand Up @@ -153,13 +183,19 @@ export function resolveMessageDispatchIntent(
);
const capabilities = providerSession?.capabilities.turns;
if (capabilities?.supportsActiveSteering === true) {
return { type: "steer_active", targetRunId: activeRun.id };
return queueBehindMaintenanceRun(projection, {
type: "steer_active",
targetRunId: activeRun.id,
});
}
if (capabilities?.supportsQueuedMessages === true) {
return { type: "queue_after_active" };
}
if (capabilities?.supportsSteeringByInterruptRestart === true) {
return { type: "restart_active", targetRunId: activeRun.id };
return queueBehindMaintenanceRun(projection, {
type: "restart_active",
targetRunId: activeRun.id,
});
}
return { type: "queue_after_active" };
}
Expand Down
17 changes: 5 additions & 12 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,11 @@ import {
SHARED_WORKSPACE_RESTORE_MESSAGE,
} from "./CheckpointRestoreSafety.ts";
import { CheckpointServiceV2 } from "./CheckpointService.ts";
import { CommandPolicyV2, resolveMessageDispatchIntent } from "./CommandPolicy.ts";
import {
CommandPolicyV2,
isNativeMaintenanceCommand,
resolveMessageDispatchIntent,
} from "./CommandPolicy.ts";
import { CommandReceiptStoreV2 } from "./CommandReceiptStore.ts";
import { ContextHandoffServiceV2 } from "./ContextHandoffService.ts";
import { notificationTurnItem } from "./Notification.ts";
Expand Down Expand Up @@ -336,17 +340,6 @@ function wakeWorkStartedAt(
return previous === undefined ? {} : { workStartedAt: orchestrationV2RunWorkStartedAt(previous) };
}

function isNativeMaintenanceCommand(message: {
readonly text: string;
readonly attachments: ReadonlyArray<ChatAttachment>;
readonly context?: import("@t3tools/contracts").OrchestrationMessageContext | undefined;
}): boolean {
return (
message.attachments.length === 0 &&
["/compact", "/logout"].includes(message.text.trim().toLowerCase())
);
}

const threadPullRequestLinksEqual = Schema.toEquivalence(Schema.NullOr(ThreadLinkedPullRequest));

function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId {
Expand Down
1 change: 1 addition & 0 deletions apps/web/src/components/ChatView.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -10886,6 +10886,7 @@ export default function ChatView(props: ChatViewProps) {
}
phase={phase}
canInterrupt={canInterruptRunningThread}
activeTurnIsCompaction={isCompacting}
isConnecting={isConnecting}
isSendBusy={isSendBusy || isSavingQueuedEdit || isResuming}
canResume={resumableRunId !== null || hasHeldQueuedRuns}
Expand Down
Loading
Loading