Skip to content
Draft
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
85 changes: 85 additions & 0 deletions apps/vscode-e2e/src/fixtures/subtasks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ const SUBTASK_APPROVAL_RESTORE_CHILD_MARKER = "SUBTASK_CHILD_APPROVAL_RESTORE"
const SUBTASK_XPROFILE_PARENT_MARKER = "SUBTASK_PARENT_CROSS_PROFILE"
const SUBTASK_XPROFILE_SAME_CHILD_MARKER = "SUBTASK_CHILD_SAME_PROFILE"
const SUBTASK_XPROFILE_DIFFERENT_CHILD_MARKER = "SUBTASK_CHILD_DIFFERENT_PROFILE"
export const SUBTASK_QUEUED_INPUT_PARENT_MARKER = "SUBTASK_PARENT_QUEUED_INPUT"
export const SUBTASK_QUEUED_INPUT_CHILD_MARKER = "SUBTASK_CHILD_QUEUED_INPUT"

const SUBTASK_CHILD_PROMPT = `${SUBTASK_CHILD_MARKER}: Ask the user exactly this follow-up question: What is the square root of 81? After the user answers, complete with only the answer.`
export const SUBTASK_PARENT_PROMPT = `${SUBTASK_PARENT_MARKER}: Use the new_task tool exactly once. Create an ask-mode subtask with this exact message: "${SUBTASK_CHILD_PROMPT}" Do not answer directly.`
Expand Down Expand Up @@ -59,6 +61,14 @@ export const SUBTASK_XPROFILE_SAME_CHILD_RESULT = "Same-profile child completed"
export const SUBTASK_XPROFILE_DIFFERENT_CHILD_RESULT = "Different-profile child completed"
export const SUBTASK_XPROFILE_PARENT_RESULT = "Sequential cross-profile parent resumed"

const SUBTASK_QUEUED_INPUT_INITIAL_RESULT = "Child completed before queued input"
export const SUBTASK_QUEUED_INPUT_MESSAGE = "Use the queued instruction before completing."
export const SUBTASK_QUEUED_INPUT_CHILD_RESULT = "Child processed queued input"
export const SUBTASK_QUEUED_INPUT_PARENT_RESULT = "Parent resumed after queued input"
const SUBTASK_QUEUED_INPUT_CHILD_PROMPT = `${SUBTASK_QUEUED_INPUT_CHILD_MARKER}: Complete immediately with the exact result "${SUBTASK_QUEUED_INPUT_INITIAL_RESULT}".`
export const SUBTASK_QUEUED_INPUT_PARENT_PROMPT = `${SUBTASK_QUEUED_INPUT_PARENT_MARKER}: Use the new_task tool exactly once. Create an ask-mode subtask with this exact message: "${SUBTASK_QUEUED_INPUT_CHILD_PROMPT}" Do not answer directly. When the subtask returns, complete with the exact result "${SUBTASK_QUEUED_INPUT_PARENT_RESULT}".`
export const SUBTASK_QUEUED_INPUT_RESPONSE_LATENCY_MS = 2_000

// Scheduler regression tests — exercises TaskScheduler + run() dispatch post-CodeRabbit fix.
// Separate markers to avoid collisions with the other subtask fixtures.
const SCHED_STANDALONE_MARKER = "SCHED_STANDALONE_INTERRUPT_RESUME"
Expand Down Expand Up @@ -179,6 +189,81 @@ export function addSubtaskFixtures(mock: InstanceType<typeof LLMock>) {
},
})

mock.addFixture({
match: {
userMessage: new RegExp(SUBTASK_QUEUED_INPUT_PARENT_MARKER),
sequenceIndex: 0,
},
response: {
toolCalls: [
{
name: "new_task",
arguments: JSON.stringify({
mode: "ask",
message: SUBTASK_QUEUED_INPUT_CHILD_PROMPT,
}),
id: "call_queued_input_parent_new_task_001",
},
],
},
})

mock.addFixture({
match: {
predicate: (req: ChatCompletionRequest) =>
lastUserMessageContains(req, SUBTASK_QUEUED_INPUT_CHILD_MARKER) &&
!requestContains(req, [SUBTASK_QUEUED_INPUT_PARENT_MARKER]) &&
!requestContains(req, [SUBTASK_QUEUED_INPUT_MESSAGE]),
},
streamingProfile: { ttft: SUBTASK_QUEUED_INPUT_RESPONSE_LATENCY_MS },
response: {
toolCalls: [
{
name: "attempt_completion",
arguments: JSON.stringify({ result: SUBTASK_QUEUED_INPUT_INITIAL_RESULT }),
id: "call_queued_input_child_initial_completion_002",
},
],
},
})

mock.addFixture({
match: {
predicate: (req: ChatCompletionRequest) =>
requestContains(req, [SUBTASK_QUEUED_INPUT_CHILD_MARKER, SUBTASK_QUEUED_INPUT_MESSAGE]) &&
!requestContains(req, [SUBTASK_QUEUED_INPUT_PARENT_MARKER]),
},
response: {
toolCalls: [
{
name: "attempt_completion",
arguments: JSON.stringify({ result: SUBTASK_QUEUED_INPUT_CHILD_RESULT }),
id: "call_queued_input_child_revised_completion_003",
},
],
},
})

mock.addFixture({
match: {
predicate: (req: ChatCompletionRequest) =>
requestContains(req, [
SUBTASK_QUEUED_INPUT_PARENT_MARKER,
SUBTASK_RESULT_INJECTION,
SUBTASK_QUEUED_INPUT_CHILD_RESULT,
]),
},
response: {
toolCalls: [
{
name: "attempt_completion",
arguments: JSON.stringify({ result: SUBTASK_QUEUED_INPUT_PARENT_RESULT }),
id: "call_queued_input_parent_completion_004",
},
],
},
})

mock.addFixture({
match: {
userMessage: new RegExp(SUBTASK_FAST_PARENT_MARKER),
Expand Down
73 changes: 73 additions & 0 deletions apps/vscode-e2e/src/suite/subtasks.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,12 @@ import {
SUBTASK_INTERRUPT_PARENT_PROMPT,
SUBTASK_INTERRUPT_PARENT_RESULT,
SUBTASK_PARENT_PROMPT,
SUBTASK_QUEUED_INPUT_CHILD_MARKER,
SUBTASK_QUEUED_INPUT_CHILD_RESULT,
SUBTASK_QUEUED_INPUT_MESSAGE,
SUBTASK_QUEUED_INPUT_PARENT_MARKER,
SUBTASK_QUEUED_INPUT_PARENT_PROMPT,
SUBTASK_QUEUED_INPUT_PARENT_RESULT,
SUBTASK_XPROFILE_DIFFERENT_CHILD_RESULT,
SUBTASK_XPROFILE_PARENT_PROMPT,
SUBTASK_XPROFILE_PARENT_RESULT,
Expand Down Expand Up @@ -260,6 +266,73 @@ suite("Roo Code Subtasks", function () {
}
})

test("queued input interrupts child completion before the parent resumes", async () => {
const api = globalThis.api
const says: Record<string, ClineMessage[]> = {}

const messageHandler = ({ taskId, message }: { taskId: string; message: ClineMessage }) => {
if (message.type === "say" && message.partial !== true) {
says[taskId] = says[taskId] || []
says[taskId].push(message)
}
}

api.on(RooCodeEventName.Message, messageHandler)

try {
const parentTaskId = await api.startNewTask({
configuration: {
mode: "ask",
alwaysAllowModeSwitch: true,
alwaysAllowSubtasks: true,
autoApprovalEnabled: true,
enableCheckpoints: false,
},
text: SUBTASK_QUEUED_INPUT_PARENT_PROMPT,
})

let childTaskId: string | undefined
await waitFor(() => {
const current = api.getCurrentTaskStack().at(-1)
if (current && current !== parentTaskId) {
childTaskId = current
return true
}
return false
})

await waitForAimockRequestContaining(SUBTASK_QUEUED_INPUT_CHILD_MARKER, SUBTASK_QUEUED_INPUT_PARENT_MARKER)

const completedParentTaskId = await waitUntilCompleted({
api,
start: async () => {
await api.sendMessage(SUBTASK_QUEUED_INPUT_MESSAGE)
return parentTaskId
},
})

assert.strictEqual(completedParentTaskId, parentTaskId)
assert.ok(
says[childTaskId!]?.some(
({ say, text }) =>
say === "completion_result" && text?.trim() === SUBTASK_QUEUED_INPUT_CHILD_RESULT,
),
"Child should process the queued instruction before returning to its parent",
)
assert.strictEqual(
says[parentTaskId]?.find(({ say }) => say === "completion_result")?.text?.trim(),
SUBTASK_QUEUED_INPUT_PARENT_RESULT,
"Parent should resume only after the child processes the queued instruction",
)
} finally {
api.off(RooCodeEventName.Message, messageHandler)
while (api.getCurrentTaskStack().length > 0) {
await api.clearCurrentTask()
}
await waitFor(() => api.getCurrentTaskStack().length === 0).catch(() => {})
}
})

// Smoke: child completing normally must resume the parent task.
test("child task returns to parent after normal completion", async () => {
const api = globalThis.api
Expand Down
5 changes: 4 additions & 1 deletion packages/types/src/api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,10 @@ export interface RooCodeAPI extends EventEmitter<RooCodeAPIEvents> {
*/
abandonSubtask(childTaskId: string): Promise<boolean>
/**
* Sends a message to the current task.
* Sends a message to the current task as conversational input.
* If the task is busy the message is queued and becomes the next user
* turn. Queued input never approves a pending or later tool, command, or
* MCP ask; use approveCurrentAsk() for explicit approval.
* @param message Optional message to send.
* @param images Optional array of image data URIs (e.g., "data:image/webp;base64,...").
*/
Expand Down
6 changes: 6 additions & 0 deletions packages/types/src/message.ts
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,12 @@ export const queuedMessageSchema = z.object({
id: z.string(),
text: z.string(),
images: z.array(z.string()).optional(),
/**
* Where the message was queued from. Absent means the interactive
* webview, whose queued input may answer approval asks. "api" input is
* conversational steering and never answers an approval ask.
*/
origin: z.enum(["webview", "api"]).optional(),
})

export type QueuedMessage = z.infer<typeof queuedMessageSchema>
8 changes: 7 additions & 1 deletion src/core/message-queue/MessageQueueService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,11 @@ import { v4 as uuidv4 } from "uuid"

import { QueuedMessage } from "@roo-code/types"

export interface AddMessageOptions {
/** Origin of the input. See {@link QueuedMessage.origin}. */
origin?: QueuedMessage["origin"]
}

export interface MessageQueueState {
messages: QueuedMessage[]
isProcessing: boolean
Expand Down Expand Up @@ -34,7 +39,7 @@ export class MessageQueueService extends EventEmitter<QueueEvents> {
return { index, message: this._messages[index] }
}

public addMessage(text: string, images?: string[]): QueuedMessage | undefined {
public addMessage(text: string, images?: string[], options?: AddMessageOptions): QueuedMessage | undefined {
if (!text && !images?.length) {
return undefined
}
Expand All @@ -44,6 +49,7 @@ export class MessageQueueService extends EventEmitter<QueueEvents> {
id: uuidv4(),
text,
images,
origin: options?.origin,
}

this._messages.push(message)
Expand Down
65 changes: 41 additions & 24 deletions src/core/task/Task.ts
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,16 @@ export function isBlanketDenyEngaged(
)
}

function queuedResponseForAsk(type: ClineAsk, text?: string): QueuedAskResolution | undefined {
function queuedResponseForAsk(
type: ClineAsk,
text: string | undefined,
origin: QueuedMessage["origin"],
): QueuedAskResolution | undefined {
// API-origin queued input is conversational steering, not an approval.
// Returning undefined keeps the message queued for the next turn and
// leaves the ask to auto-approval settings or an explicit decision.
const canApprove = origin !== "api"

if (type === "command_output") {
return undefined
}
Expand All @@ -223,12 +232,13 @@ function queuedResponseForAsk(type: ClineAsk, text?: string): QueuedAskResolutio
}
} catch {
// Malformed tool asks retain the existing approve-with-feedback behavior.
return canApprove ? { response: "yesButtonClicked", requiresDurableAck: false } : undefined
}

return { response: "yesButtonClicked", requiresDurableAck: false }
return canApprove ? { response: "yesButtonClicked", requiresDurableAck: false } : undefined
}
if (type === "command" || type === "use_mcp_server") {
return { response: "yesButtonClicked", requiresDurableAck: false }
return canApprove ? { response: "yesButtonClicked", requiresDurableAck: false } : undefined
}

return { response: "messageResponse", requiresDurableAck: type === "completion_result" }
Expand Down Expand Up @@ -1277,9 +1287,9 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
* blanket-denied (it answers that denial, not the current ask), and
* `hasUnclaimed()` replaces the length-only `isEmpty()`, which reports a
* queue containing nothing but claims as available for a new consumer.
* `isMessageQueued`/`isStatusMutable` keep `isEmpty()` semantics on purpose:
* flipping those would re-enable interactive prompt timers whenever a claim
* is outstanding.
* `isStatusMutable` does not read queue length: it keys on whether a claimed
* message will answer the ask, so a claim that answers suppresses the
* interactive prompt timers and a released claim leaves them armed.
*/
private mayDrainQueuedMessageForAsk(): boolean {
return !this.blanketDeniedCommandThisTurn && this.messageQueueService.hasUnclaimed()
Expand Down Expand Up @@ -1797,7 +1807,14 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
!this.mayDrainQueuedMessageForAsk()
? undefined
: this.messageQueueService.claimNextMessage()
const queuedAskResolution = queuedMessage ? queuedResponseForAsk(type, text) : undefined
const queuedAskResolution = queuedMessage ? queuedResponseForAsk(type, text, queuedMessage.origin) : undefined

if (queuedMessage && !queuedAskResolution) {
// API-origin input cannot answer this ask. Release the claim so the
// message stays queued as the next conversational turn.
this.messageQueueService.releaseMessage(queuedMessage.id)
}

// `this.cwd`, not `provider.cwd`:
// The path inside `text` was made relative to this task's workspace,
// which for a resumed or child task need not be the one the provider
Expand Down Expand Up @@ -1970,20 +1987,23 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
// The state is mutable if the message is complete and the task will
// block (via the `pWaitFor`).
const isBlocking = !(this.askResponse !== undefined || this.lastMessageTs !== askTs)
const isMessageQueued = !this.messageQueueService.isEmpty()
// Keep queued user messages intact during command_output asks. Those asks
// are terminal flow-control, not conversational turns.
const shouldDrainQueuedMessageForAsk = type !== "command_output"
const isStatusMutable = !partial && isBlocking && !isMessageQueued && approval.decision === "ask"
// The FIFO drain answers this ask only from the claimed head message.
// API-origin steering released above stays queued for the next turn, so
// queue emptiness cannot decide whether the ask waits for a response.
const isStatusMutable =
!partial && isBlocking && !(queuedMessage && queuedAskResolution) && approval.decision === "ask"

let queuedMessageId: string | undefined
// Arm the interactive/resumable/idle status timers for this ask: the
// single source of that arm, shared between the queue-free case and the
// claim-gated and queued-release paths below. A gated or released claim
// keeps the message in the queue, so `isMessageQueued` stays true and
// `isStatusMutable` — which requires an empty queue — stays false while
// the ask waits for the user; arming only from `isStatusMutable` would
// leave hands-free/API consumers seeing `Running` with no
// keeps the message in the queue, but `isStatusMutable` no longer reads
// queue length: with no claimed message answering the ask it stays true
// while the ask waits for the user, and arming only from the drain sites
// would leave hands-free/API consumers seeing `Running` with no
// `TaskInteractive`/`interactionRequired` for a prompt that is in fact
// pending. Idempotent: several arm sites can fire for one ask (e.g. the
// queue-free arm, then a drain-site release), and a second arm would
Expand Down Expand Up @@ -2121,13 +2141,6 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
} else {
queuedMessageId = this.handleQueuedAskResponse(queuedMessage, queuedAskResolution)
}
} else if (shouldDrainQueuedMessageForAsk && isMessageQueued) {
// The claim gate (per-turn latch, or blanket deny engaged for a command
// ask) left the queued message untouched. If the policy still leaves the
// prompt pending, the non-empty queue keeps `isStatusMutable` false, so
// the interactive arm must run from here — the same reason a release
// re-arms. For an auto-answered ask the arm's pending check declines.
armAskStatusTimers()
}

// At most one drain-site policy re-check is in flight per ask; the
Expand Down Expand Up @@ -2163,9 +2176,9 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
)
if (!this.abort && this.askResponse === undefined && this.lastMessageTs === askTs) {
queuedMessageId = this.applyQueuedCommandPolicyAction(action, message, resolution)
// A "release" outcome leaves the ask pending while the queue stays
// non-empty, so the arm that `isStatusMutable` gates — computed
// once, before the claim — must run here.
// A "release" outcome leaves the ask pending, and `isStatusMutable`
// — computed once, before the claim — still counts the message as
// answering the ask, so the arm must run here.
armAskStatusTimers()
}
} finally {
Expand Down Expand Up @@ -2201,7 +2214,7 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
this.mayDrainQueuedMessageForAsk()
) {
const message = this.messageQueueService.claimNextMessage()
const resolution = message ? queuedResponseForAsk(type, text) : undefined
const resolution = message ? queuedResponseForAsk(type, text, message.origin) : undefined
if (message && resolution) {
if (type === "command") {
// Claim first, then verify the policy off-predicate: a
Expand All @@ -2221,6 +2234,10 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
} else {
queuedMessageId = this.handleQueuedAskResponse(message, resolution)
}
} else if (message) {
// API-origin input cannot answer this ask. Release the
// claim so the message stays queued for the next turn.
this.messageQueueService.releaseMessage(message.id)
}
}

Expand Down
Loading
Loading