Skip to content

Commit 0cc0096

Browse files
committed
fix(task): enqueue streaming input before child completion
1 parent eaef1ae commit 0cc0096

5 files changed

Lines changed: 112 additions & 3 deletions

File tree

‎src/core/task/Task.ts‎

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ import {
2929
type ContextTruncation,
3030
type ClineMessage,
3131
type ClineSay,
32+
type ClineSayTool,
3233
type ClineAsk,
3334
type ToolProgressStatus,
3435
type HistoryItem,
@@ -1150,10 +1151,10 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
11501151
return undefined
11511152
}
11521153

1153-
private drainQueuedMessageIntoAskResponse(): void {
1154+
private drainQueuedMessageIntoAskResponse(allowResolvedAskOverride = false): void {
11541155
// A synchronous auto-approval may already have resolved the ask before the
11551156
// entry queue snapshot is acted on. Never replace that resolved response.
1156-
if (this.askResponse !== undefined) {
1157+
if (this.askResponse !== undefined && !allowResolvedAskOverride) {
11571158
return
11581159
}
11591160

@@ -1328,6 +1329,14 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
13281329
// Keep queued user messages intact during command_output asks. Those asks
13291330
// are terminal flow-control, not conversational turns.
13301331
const shouldDrainQueuedMessageForAsk = type !== "command_output"
1332+
let isFinishTaskAsk = false
1333+
if (type === "tool") {
1334+
try {
1335+
isFinishTaskAsk = (JSON.parse(text || "{}") as ClineSayTool).tool === "finishTask"
1336+
} catch {
1337+
// Invalid tool payloads are handled by their caller; they are not finishTask asks.
1338+
}
1339+
}
13311340
const isStatusMutable = !partial && isBlocking && !isMessageQueued && approval.decision === "ask"
13321341

13331342
if (isStatusMutable) {
@@ -1373,7 +1382,9 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
13731382
}
13741383
} else if (isMessageQueued && shouldDrainQueuedMessageForAsk) {
13751384
// This branch acts on the queue state captured when the ask was entered.
1376-
this.drainQueuedMessageIntoAskResponse()
1385+
// A queued instruction must interrupt finishTask before the child returns
1386+
// to its parent, even when subtask completion is otherwise auto-approved.
1387+
this.drainQueuedMessageIntoAskResponse(isFinishTaskAsk)
13771388
}
13781389

13791390
// Wait for askResponse to be set

‎src/core/task/__tests__/ask-queued-message-drain.spec.ts‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,43 @@ describe("Task.ask queued message drain", () => {
9494
expect(task.messageQueueService.messages[0]?.text).toBe("change direction")
9595
})
9696

97+
it("lets queued input interrupt an auto-approved finishTask ask", async () => {
98+
const task = buildTask({
99+
autoApprovalEnabled: true,
100+
alwaysAllowSubtasks: true,
101+
})
102+
task.messageQueueService.addMessage("Use the queued instruction before completing.")
103+
104+
const result = await task.ask("tool", JSON.stringify({ tool: "finishTask" }), false)
105+
106+
expect(result).toEqual({
107+
response: "messageResponse",
108+
text: "Use the queued instruction before completing.",
109+
images: undefined,
110+
})
111+
expect(task.messageQueueService.isEmpty()).toBe(true)
112+
})
113+
114+
it("preserves queued input behind an ordinary auto-approved tool ask", async () => {
115+
const task = buildTask({ autoApprovalEnabled: true })
116+
task.messageQueueService.addMessage("change direction")
117+
118+
const result = await task.ask("tool", JSON.stringify({ tool: "updateTodoList" }), false)
119+
120+
expect(result.response).toBe("yesButtonClicked")
121+
expect(task.messageQueueService.messages).toHaveLength(1)
122+
})
123+
124+
it("treats queued input as feedback when a tool ask payload is malformed", async () => {
125+
const task = buildTask()
126+
task.messageQueueService.addMessage("change direction")
127+
128+
const result = await task.ask("tool", "{", false)
129+
130+
expect(result).toMatchObject({ response: "messageResponse", text: "change direction" })
131+
expect(task.messageQueueService.isEmpty()).toBe(true)
132+
})
133+
97134
it("does not consume messages that were queued before a command_output ask", async () => {
98135
const task = buildTask()
99136
task.messageQueueService.addMessage("1+1=?")

‎src/core/tools/__tests__/readFileTool.spec.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -687,6 +687,22 @@ describe("ReadFileTool", () => {
687687
})
688688
parseSpy.mockRestore()
689689
})
690+
691+
it("denies batch reads for a queued message response without text", async () => {
692+
const task = Object.create(Task.prototype) as Task
693+
Object.defineProperty(task, "cwd", { value: "/test/workspace", writable: true })
694+
Object.assign(task, createMockTask())
695+
task.ask = vi.fn().mockResolvedValue({ response: "messageResponse", text: undefined, images: undefined })
696+
const fileResults = [
697+
{ path: "one.ts", status: "pending" as const, entry: { path: "one.ts", mode: "slice" as const } },
698+
{ path: "two.ts", status: "pending" as const, entry: { path: "two.ts", mode: "slice" as const } },
699+
]
700+
701+
await readFileTool["requestApproval"](task, fileResults, () => {})
702+
703+
expect(task.say).not.toHaveBeenCalledWith("user_feedback", expect.anything(), expect.anything())
704+
expect(task.didRejectTool).toBe(true)
705+
})
690706
})
691707

692708
describe("output structure", () => {

‎src/extension/__tests__/api-send-message.spec.ts‎

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,42 @@ describe("API - SendMessage Command", () => {
5656
})
5757
})
5858

59+
it("should enqueue directly when the current task is streaming", async () => {
60+
const addMessage = vi.fn()
61+
const messageText = "Use this before completing"
62+
const images = ["data:image/png;base64,image1data"]
63+
const currentTask = {
64+
isStreaming: true,
65+
messageQueueService: { addMessage },
66+
}
67+
mockProvider.getCurrentTask = vi.fn().mockReturnValue(currentTask)
68+
69+
await api.sendMessage(messageText, images)
70+
71+
expect(addMessage).toHaveBeenCalledWith(messageText, images)
72+
expect(mockPostMessageToWebview).not.toHaveBeenCalled()
73+
})
74+
75+
it("should retain webview routing when the current task is not streaming", async () => {
76+
const addMessage = vi.fn()
77+
const messageText = "Answer the current ask"
78+
const currentTask = {
79+
isStreaming: false,
80+
messageQueueService: { addMessage },
81+
}
82+
mockProvider.getCurrentTask = vi.fn().mockReturnValue(currentTask)
83+
84+
await api.sendMessage(messageText)
85+
86+
expect(addMessage).not.toHaveBeenCalled()
87+
expect(mockPostMessageToWebview).toHaveBeenCalledWith({
88+
type: "invoke",
89+
invoke: "sendMessage",
90+
text: messageText,
91+
images: undefined,
92+
})
93+
})
94+
5995
it("should handle SendMessage command with text and images", async () => {
6096
// Arrange
6197
const messageText = "Analyze this image"

‎src/extension/api.ts‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -271,6 +271,15 @@ export class API extends EventEmitter<RooCodeEvents> implements RooCodeAPI {
271271
public async sendMessage(text?: string, images?: string[]) {
272272
const currentTask = this.sidebarProvider.getCurrentTask()
273273

274+
// API callers need the returned promise to mean that sequencing-critical
275+
// input has reached the active task. During a stream, the webview would
276+
// only relay this message back as queueMessage asynchronously, so enqueue
277+
// it in the extension host instead of racing task completion.
278+
if (currentTask?.isStreaming) {
279+
currentTask.messageQueueService.addMessage(text ?? "", images)
280+
return
281+
}
282+
274283
// In headless/sandbox flows the webview may not be launched, so routing
275284
// through invoke=sendMessage drops the message. Deliver directly to the
276285
// task ask-response channel instead.

0 commit comments

Comments
 (0)