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
59 changes: 59 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/MuseAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -809,6 +809,65 @@ describe("MuseAdapterV2", () => {
}).pipe(Effect.scoped, Effect.provide(testLayer)),
);

it.effect("approves a workflow child's approval in full access after its turn ended", () =>
Effect.gen(function* () {
const fake = yield* makeFakeMuse();
const harness = yield* makeHarness(fake);
const { nativeId } = yield* startConversation(harness, fake);
yield* fake.emit("turn/completed", { turnId: nativeId, terminal: "completed" });
yield* harness.takeEvent("turn.terminal");
// Muse runs workflow children without the session's allowAll mode.
const childApproval = { ...approval("child-run-1"), subagentOrigin: { subagentId: "a7" } };
yield* fake.emit("approval/requested", childApproval);
const asked = yield* harness.takeEvent("runtime_request.updated");
assert.strictEqual(asked.runtimeRequest.status, "pending");
assert.isNull(asked.runtimeRequest.providerTurnId);
const decision = yield* fake.takeCall("approval/decide");
assert.strictEqual(decision.params.choiceId, "once");
yield* fake.emit("approval/resolved", childApproval);
const resolved = yield* harness.takeEvent("runtime_request.updated");
assert.strictEqual(resolved.runtimeRequest.status, "resolved");
assert.strictEqual(resolved.runtimeRequest.decision, "accept");
}).pipe(Effect.scoped, Effect.provide(testLayer)),
);

it.effect("keeps a workflow child's approval pending for the user past the turn's end", () =>
Effect.gen(function* () {
const fake = yield* makeFakeMuse();
const policy = ProviderAdapterV2RuntimePolicy.make({
...runtimePolicy,
runtimeMode: "approval-required",
});
const harness = yield* makeHarness(
fake,
INSTANCE_ID,
undefined,
undefined,
undefined,
policy,
);
const { nativeId } = yield* startConversation(harness, fake);
const childApproval = { ...approval("child-run-1"), subagentOrigin: { subagentId: "a7" } };
yield* fake.emit("approval/requested", childApproval);
const asked = yield* harness.takeEvent("runtime_request.updated");
assert.isNull(asked.runtimeRequest.providerTurnId);
yield* fake.emit("turn/completed", { turnId: nativeId, terminal: "completed" });
yield* harness.takeEvent("turn.terminal");
assert.isFalse(
harness.allEvents.some(
(event) =>
event.type === "runtime_request.updated" && event.runtimeRequest.status !== "pending",
),
);
yield* harness.runtime.respondToRuntimeRequest({
requestId: asked.runtimeRequest.id,
decision: "decline",
});
const decision = yield* fake.takeCall("approval/decide");
assert.strictEqual(decision.params.choiceId, "deny");
}).pipe(Effect.scoped, Effect.provide(testLayer)),
);

it.effect("validates structured answers and preserves selected labels with a custom note", () =>
Effect.gen(function* () {
const fake = yield* makeFakeMuse();
Expand Down
192 changes: 126 additions & 66 deletions apps/server/src/orchestration-v2/Adapters/MuseAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import {
type ProviderInstanceId,
type PlanId,
type OrchestrationV2PlanStep,
type ProviderApprovalDecision,
type RuntimeRequestId,
type ServerProviderModel,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -157,7 +158,7 @@ const MuseProviderCapabilitiesV2 = {
supportsApplyPatchApproval: false,
approvalsHaveNativeRequestIds: true,
approvalCallbacksAreLiveOnly: true,
approvalsCanOriginateFromSubagents: false,
approvalsCanOriginateFromSubagents: true,
},
planning: {
emitsPlanUpdated: true,
Expand Down Expand Up @@ -664,6 +665,8 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
}
observedChildren.clear();
if (thread?.pendingBackgroundTasks?.length) yield* syncBackground();
// Only workflow children's requests outlive a turn; they end with the workflow.
for (const entry of pending.values()) yield* resolvePending(entry, "cancelled");
});
const resolvePending = Effect.fnUntraced(function* (
entry: PendingRequest,
Expand Down Expand Up @@ -718,7 +721,9 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
yield* publishItem(turn, item, item.status === "inProgress" ? status : undefined);
}
turn.dirty.clear();
for (const entry of pending.values()) yield* resolvePending(entry, "cancelled");
// A workflow child's request is runless and keeps waiting after the turn ends.
for (const entry of pending.values())
if (entry.request.providerTurnId !== null) yield* resolvePending(entry, "cancelled");
turn.providerTurn = { ...turn.providerTurn, status, completedAt };
lastProviderTurn = turn.providerTurn;
yield* emit({
Expand Down Expand Up @@ -801,28 +806,54 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
});
const owns = (turn: ActiveTurn, turnId: string) =>
turnId === turn.nativeId || turn.joined.has(turnId);
/** Sends T3's decision for a pending approval to Muse. */
const decideApproval = Effect.fnUntraced(function* (
entry: PendingRequest,
decision: ProviderApprovalDecision | undefined,
) {
if (entry.native.type !== "approval") return;
const approval = entry.native.value;
const choice = decision && museApprovalChoices(approval).get(decision);
if (!choice) return yield* protocolError("Muse did not offer this approval decision");
// Recorded first: Muse can settle the approval before it acknowledges the command.
entry.response = { decision };
yield* request("approval/decide", {
approvalId: approval.approvalId,
requirementId: approval.currentRequirementId,
choiceId: choice.choiceId,
}).pipe(Effect.tapError(() => Effect.sync(() => delete entry.response)));
});
/**
* Shows a Muse approval or question as a T3 request. A workflow child's approval
* is runless: the turn that started the workflow may end long before the child
* asks, and T3 keeps a runless request on the thread until Muse settles it.
*/
const publishRequest = Effect.fnUntraced(function* (native: PendingRequest["native"]) {
const turn = active;
const child = native.type === "approval" && native.value.subagentOrigin !== undefined;
const turn = child ? undefined : active;
if (
!turn ||
!thread ||
(!child && !turn) ||
native.value.sessionId !== nativeSessionId ||
(native.value.turnId && !owns(turn, native.value.turnId))
(turn && native.value.turnId && !owns(turn, native.value.turnId))
)
return;
const nativeId =
native.type === "approval" ? native.value.approvalId : native.value.userInputId;
if (turn.settledRequests.has(`${native.type}:${nativeId}`)) return;
if (turn?.settledRequests.has(`${native.type}:${nativeId}`)) return;
const previous = [...pending.values()].find(
(entry) =>
entry.native.type === native.type &&
entry.request.nativeRequestRef?.nativeId === nativeId,
);
const providerTurnId = turn?.providerTurn.id ?? null;
const threadId = turn?.input.threadId ?? thread.appThreadId ?? input.threadId;
const requestId =
previous?.request.id ??
(yield* idAllocator.allocate
.runtimeRequest({
driver: MUSE_PROVIDER,
providerTurnId: turn.providerTurn.id,
...(providerTurnId ? { providerTurnId } : {}),
nativeRequestId: nativeId,
})
.pipe(
Expand All @@ -843,7 +874,7 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
const runtimeRequest: OrchestrationV2RuntimeRequest = {
id: requestId,
nodeId,
providerTurnId: turn.providerTurn.id,
providerTurnId,
nativeRequestRef: nativeRef(nativeId),
kind,
status: "pending",
Expand All @@ -853,23 +884,37 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
};
const node: OrchestrationV2ExecutionNode = {
id: nodeId,
threadId: turn.input.threadId,
runId: turn.input.runId,
parentNodeId: turn.input.rootNodeId,
rootNodeId: turn.input.rootNodeId,
threadId,
runId: turn?.input.runId ?? null,
parentNodeId: turn?.input.rootNodeId ?? null,
rootNodeId: turn?.input.rootNodeId ?? nodeId,
kind: native.type === "approval" ? "approval_request" : "user_input_request",
status: "waiting",
countsForRun: false,
providerThreadId: turn.input.providerThread.id,
providerTurnId: turn.providerTurn.id,
providerThreadId: thread.id,
providerTurnId,
nativeItemRef: nativeRef(nativeId),
runtimeRequestId: requestId,
checkpointScopeId: null,
startedAt: runtimeRequest.createdAt,
completedAt: null,
};
const latestTurn = active?.providerTurn ?? lastProviderTurn;
const base = {
...baseItem(turn, `request:${nativeId}`, time),
...(turn
? baseItem(turn, `request:${nativeId}`, time)
: {
threadId,
runId: null,
providerThreadId: thread.id,
providerTurnId: null,
nativeItemRef: nativeRef(nativeId),
parentItemId: null,
// Sorts after the latest turn's items.
ordinal: previous?.item.ordinal ?? ((latestTurn?.ordinal ?? 0) + 1) * 100,
startedAt: runtimeRequest.createdAt,
updatedAt: time,
}),
id: idAllocator.derive.approvalTurnItem({ requestId }),
nodeId,
status: "waiting" as const,
Expand Down Expand Up @@ -909,16 +954,71 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
required: true,
})),
};
pending.set(requestId, { native, request: runtimeRequest, node, item });
const entry: PendingRequest = { native, request: runtimeRequest, node, item };
pending.set(requestId, entry);
yield* emit({
type: "runtime_request.updated",
driver: MUSE_PROVIDER,
threadId: turn.input.threadId,
threadId,
runtimeRequest,
});
yield* emit({ type: "node.updated", driver: MUSE_PROVIDER, node });
yield* emit({ type: "turn_item.updated", driver: MUSE_PROVIDER, turnItem: item });
yield* updateSession("waiting");
// Muse does not give workflow children the session's allowAll mode, so in full
// access T3 approves for them. If Muse rejects that, the request waits for the user.
if (child && input.runtimePolicy.runtimeMode === "full-access")
yield* decideApproval(entry, "accept").pipe(Effect.ignore, Effect.forkIn(scope));
});
/** Approval and question events; `turn` is undefined for a workflow child's approval. */
const handleRequestEvent = Effect.fnUntraced(function* (
method: string,
params: Record<string, unknown>,
turn: ActiveTurn | undefined,
) {
switch (method) {
case "approval/requested":
case "approval/updated":
if (
method === "approval/updated" &&
![...pending.values()].some(
(entry) =>
entry.native.type === "approval" &&
entry.native.value.approvalId === params.approvalId,
)
)
return;
return yield* publishRequest({
type: "approval",
value: yield* decode(MuseApproval, params),
});
case "userInput/requested":
return yield* publishRequest({
type: "question",
value: yield* decode(MuseUserInput, params),
});
case "approval/resolved":
case "userInput/settled": {
const id = method === "approval/resolved" ? params.approvalId : params.userInputId;
if (typeof id === "string")
turn?.settledRequests.add(
`${method === "approval/resolved" ? "approval" : "question"}:${id}`,
);
const entry = [...pending.values()].find(
(candidate) => candidate.request.nativeRequestRef?.nativeId === id,
);
if (entry) {
// An approval settled outside T3 still shows Muse's own decision.
const native =
method === "approval/resolved" && typeof params.decision === "string"
? museApprovalDecision({ decision: params.decision, scope: "once" })
: undefined;
if (!entry.response && native) entry.response = { decision: native };
yield* resolvePending(entry, "resolved");
}
if (!pending.size) yield* updateSession(active ? "running" : "ready");
}
}
});
const handleNotification = Effect.fnUntraced(function* (method: string, data: unknown) {
const params = yield* decode(recordSchema, data);
Expand Down Expand Up @@ -962,6 +1062,10 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
return;
}
}
// A workflow child's approval comes under the child's own run id, during a turn
// or after it, so it skips turn routing.
if (params.subagentOrigin !== undefined && method.startsWith("approval/"))
return yield* handleRequestEvent(method, params, undefined);
// Muse starts a turn on its own when a workflow finishes, to report its result.
// During a run, that run shows it. Otherwise hold it and ask the orchestrator for
// a run; that run takes the held events, unless a user turn takes them first.
Expand Down Expand Up @@ -1097,45 +1201,11 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
}
case "approval/requested":
case "approval/updated":
if (
method === "approval/updated" &&
![...pending.values()].some(
(entry) =>
entry.native.type === "approval" &&
entry.native.value.approvalId === params.approvalId,
)
)
break;
yield* publishRequest({ type: "approval", value: yield* decode(MuseApproval, params) });
break;
case "userInput/requested":
yield* publishRequest({
type: "question",
value: yield* decode(MuseUserInput, params),
});
break;
case "approval/resolved":
case "userInput/settled": {
const id = method === "approval/resolved" ? params.approvalId : params.userInputId;
if (typeof id === "string")
turn.settledRequests.add(
`${method === "approval/resolved" ? "approval" : "question"}:${id}`,
);
const entry = [...pending.values()].find(
(candidate) => candidate.request.nativeRequestRef?.nativeId === id,
);
if (entry) {
// An approval settled outside T3 still shows Muse's own decision.
const native =
method === "approval/resolved" && typeof params.decision === "string"
? museApprovalDecision({ decision: params.decision, scope: "once" })
: undefined;
if (!entry.response && native) entry.response = { decision: native };
yield* resolvePending(entry, "resolved");
}
if (!pending.size) yield* updateSession("running");
case "userInput/requested":
case "userInput/settled":
yield* handleRequestEvent(method, params, turn);
break;
}
case "session/tokenUsage": {
const usage = yield* decode(MuseTokenUsageEvent, params);
const previous = turn.providerTurn.turnTokenUsage;
Expand Down Expand Up @@ -1892,18 +1962,8 @@ export function makeMuseAdapterV2(options: MuseAdapterV2Options): ProviderAdapte
Effect.gen(function* () {
const entry = pending.get(args.requestId);
if (!entry) return yield* protocolError("This Muse request is no longer pending");
if (entry.native.type === "approval") {
const approval = entry.native.value;
const choice = args.decision && museApprovalChoices(approval).get(args.decision);
if (!choice) return yield* protocolError("Muse did not offer this approval decision");
// Recorded first: Muse can settle the approval before it acknowledges the command.
entry.response = { decision: args.decision };
yield* request("approval/decide", {
approvalId: approval.approvalId,
requirementId: approval.currentRequirementId,
choiceId: choice.choiceId,
}).pipe(Effect.tapError(() => Effect.sync(() => delete entry.response)));
} else {
if (entry.native.type === "approval") yield* decideApproval(entry, args.decision);
else {
const questionRequest = entry.native.value;
const answers = [];
for (const question of questionRequest.questions) {
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/provider/museProtocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,8 @@ export const MuseApproval = Schema.Struct({
itemId: Schema.optional(Schema.String),
toolName: Schema.optional(Schema.String),
rawArgs: Schema.optional(Schema.String),
/** Set when a workflow child asks; `turnId` is then the child's run, not a session turn. */
subagentOrigin: Schema.optional(Schema.Struct({ subagentId: NonEmptyString })),
currentRequirementId: Schema.Struct({ approvalId: NonEmptyString, sourceIndex: Schema.Int }),
availableChoices: Schema.Array(
Schema.Struct({
Expand Down
Loading