Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
53 commits
Select commit Hold shift + click to select a range
993cf7a
fix(tools): process queued messages when terminal command finishes
Sep 20, 2026
2d08001
fix(tools): gate queued-message drain on actual command submission
Sep 20, 2026
59d20fc
test(tools): cover execa retry working-directory validation gate
Sep 20, 2026
4765b92
fix(task): make queued-message drain awaitable and failure-aware
Sep 22, 2026
8ca6ec8
test(tools): cover interrupted and timeout queue drains
Sep 22, 2026
33aaa1e
test(tools): resolve processQueuedMessages mocks so success paths exe…
myk1yt Sep 27, 2026
fd9d422
fix(tools): separate command drain failures from shell-integration ha…
myk1yt Sep 27, 2026
6903902
fix(tools): gate background-completion drain on tool result publication
myk1yt Sep 27, 2026
cefd735
test(tools): cover queued-message drain rejection paths for patch cov…
myk1yt Sep 27, 2026
5183a1e
fix(task): serialize queued-message drains per task
myk1yt Sep 27, 2026
f2981f6
fix(task): retain queued messages until an ask consumes them
myk1yt Sep 27, 2026
6bc900e
fix(task): consume queued messages by identity with durable acks
myk1yt Sep 28, 2026
14125ee
test(tools): assert drain runs once after persisted output results
myk1yt Sep 28, 2026
4860a9c
fix(task): ack intercepted queued messages at every consuming ask
myk1yt Sep 28, 2026
8764212
fix(task): make queued-feedback persistence idempotent and drains can…
myk1yt Sep 28, 2026
fbb06bc
test(tools,assistant-message): cover queued-ack wrapper branches
myk1yt Sep 28, 2026
59f0e3f
fix(task): keep queued conversational messages out of approval asks
myk1yt Sep 28, 2026
844e2c2
fix(tools): settle the command publication signal on terminal errors
myk1yt Sep 28, 2026
31eb232
test(tools): strengthen the publication-signal settle test
myk1yt Sep 28, 2026
b6112c5
fix(task): keep unclaimed ask resolutions inert so auto-approval stil…
myk1yt Sep 28, 2026
43f1cd0
fix(tools): carry the publication outcome and never drain after a fai…
myk1yt Sep 28, 2026
b25c04a
fix(task): await the reconciled feedback row and make backoff cancell…
myk1yt Sep 28, 2026
3a5ae4d
test(task): pin the awaited reconciled-row update failure path
myk1yt Sep 28, 2026
7e8528d
Merge remote-tracking branch 'refs/remotes/upstream/main' into fix/93…
myk1yt Oct 3, 2026
af09e70
fix(task): gate queued-message drains against in-flight approval asks
myk1yt Oct 3, 2026
e7fbd33
fix(task): clear the pending drain tracker when the durable ack settles
myk1yt Oct 3, 2026
9cb2686
fix(tools): settle toolResultPublished on execa fallback and warning …
myk1yt Oct 3, 2026
e7613b1
fix(task): discard queued messages consumed by the truncation retry gate
myk1yt Oct 3, 2026
b38f196
fix(task): associate queued feedback rows before the history append
myk1yt Oct 3, 2026
8eba814
fix(task,assistant-message): ack queued feedback before it enters the…
myk1yt Oct 3, 2026
cf554c7
fix(webview): validate queued-message edits and handle condense rejec…
myk1yt Oct 3, 2026
306626b
test(task): widen the drain-spec emit test-double signature
myk1yt Oct 3, 2026
c14114d
fix(task): harden queued-drain guards (tracker release, gate arming, …
myk1yt Oct 3, 2026
3984e8b
fix(task): reconcile registered feedback rows through non-durable claims
myk1yt Oct 3, 2026
2fcf21d
fix(api): reject headless sendMessage when the task refuses delivery
myk1yt Oct 3, 2026
c42dd9c
fix(webview): surface refused edit resubmissions and condense failures
myk1yt Oct 3, 2026
f1d1a13
fix(task): keep queued messages out of failure-gate retry asks
myk1yt Oct 3, 2026
8be15bc
refactor(task): split ask into a gate wrapper and askImpl (CI diff bu…
myk1yt Oct 3, 2026
ebd2c53
Merge remote-tracking branch 'refs/remotes/upstream/main' into fix/93…
myk1yt Oct 3, 2026
4a2e675
fix(webview): post condenseTaskContextResponse on condense failure
myk1yt Oct 3, 2026
43a5ead
fix(api): contain rejected headless SendMessage inside the IPC boundary
myk1yt Oct 3, 2026
9a07b0c
fix(task): route queued-message edits through the task's pending subm…
myk1yt Oct 3, 2026
ed89a6e
fix(task): defer the ask-start queued claim to the handoff point
myk1yt Oct 3, 2026
f65eb12
fix(task): reserve a consumed drain submission through its durable ack
myk1yt Oct 3, 2026
d6d0b86
fix(task): track in-flight ask gates per ask
myk1yt Oct 3, 2026
a0f9770
fix(task): hand a consumed submission's ID only when its entry is res…
myk1yt Oct 3, 2026
e0b35ca
Merge remote-tracking branch 'origin/main' into fix/937-process-queue…
myk1yt Oct 5, 2026
db35152
fix(task): repair ask() abort re-check after main merge
myk1yt Oct 5, 2026
b3561ff
fix(task): persist queued edits made during a failed-save backoff
myk1yt Oct 5, 2026
6b801ce
test(tools): assert commandSubmitted in direct executeCommandInTermin…
myk1yt Oct 5, 2026
3e41611
chore: retrigger review-state reconciliation
myk1yt Oct 5, 2026
bc38df9
chore: re-run checks cancelled by infra
myk1yt Oct 5, 2026
d5872e3
chore: retrigger review-state reconciliation
myk1yt Oct 5, 2026
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
2 changes: 1 addition & 1 deletion packages/types/src/task.ts
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ export interface TaskLike {

approveAsk(options?: { text?: string; images?: string[] }): void
denyAsk(options?: { text?: string; images?: string[] }): void
submitUserMessage(text: string, images?: string[], mode?: string, providerProfile?: string): Promise<void>
submitUserMessage(text: string, images?: string[], mode?: string, providerProfile?: string): Promise<boolean>
abortTask(): Promise<void>
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ interface MockTask {
}
}
say: ReturnType<typeof vi.fn>
sayUserFeedbackAndAckQueued: ReturnType<typeof vi.fn>
ask: ReturnType<typeof vi.fn>
pushToolResultToUserContent: ReturnType<typeof vi.fn>
getTaskMode: ReturnType<typeof vi.fn>
Expand Down Expand Up @@ -119,6 +120,8 @@ function buildMockTask(): MockTask {
}),
},
say: vi.fn().mockResolvedValue(undefined),
// The merged askApproval routes feedback through the durable queued ack.
sayUserFeedbackAndAckQueued: vi.fn().mockResolvedValue(undefined),
ask: vi.fn().mockResolvedValue({ response: "yesButtonClicked" }),
pushToolResultToUserContent: vi.fn(),
getTaskMode: vi.fn().mockResolvedValue("code"),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

import type { Anthropic } from "@anthropic-ai/sdk"
import { describe, it, expect, beforeEach, vi, type Mock } from "vitest"
import { providerIdentifiers } from "@roo-code/types"
import { presentAssistantMessage } from "../presentAssistantMessage"
import { validateToolUse } from "../../tools/validateToolUse"
import { getModeBySlug } from "../../../shared/modes"
Expand Down Expand Up @@ -39,6 +40,7 @@ vi.mock("@roo-code/telemetry", () => ({
instance: {
captureToolUsage: vi.fn(),
captureConsecutiveMistakeError: vi.fn(),
captureException: vi.fn(),
captureEvent: vi.fn(),
},
},
Expand All @@ -64,6 +66,7 @@ interface MockTask {
api: { getModel: () => { id: string; info: Record<string, unknown> } }
recordToolUsage: ReturnType<typeof vi.fn>
recordToolError: ReturnType<typeof vi.fn>
apiConfiguration?: { apiProvider: string }
toolRepetitionDetector: { check: ReturnType<typeof vi.fn> }
providerRef: {
deref: () =>
Expand All @@ -74,6 +77,7 @@ interface MockTask {
| undefined
}
say: ReturnType<typeof vi.fn>
sayUserFeedbackAndAckQueued: ReturnType<typeof vi.fn>
ask: ReturnType<typeof vi.fn>
pushToolResultToUserContent: ReturnType<typeof vi.fn>
}
Expand All @@ -83,6 +87,10 @@ describe("presentAssistantMessage - tool usage attribution", () => {

beforeEach(() => {
vi.clearAllMocks()
// clearAllMocks keeps queued one-shot implementations alive; tests that skip the
// validation arm would otherwise leak a mockImplementationOnce throw into a
// later test that runs the validated arm.
vi.mocked(validateToolUse).mockReset()
vi.mocked(validateToolUse).mockImplementation(() => undefined)

mockTask = {
Expand Down Expand Up @@ -117,6 +125,7 @@ describe("presentAssistantMessage - tool usage attribution", () => {
}),
},
say: vi.fn().mockResolvedValue(undefined),
sayUserFeedbackAndAckQueued: vi.fn().mockResolvedValue(undefined),
ask: vi.fn().mockResolvedValue({ response: "yesButtonClicked" }),
pushToolResultToUserContent: vi.fn(),
}
Expand Down Expand Up @@ -651,6 +660,113 @@ describe("presentAssistantMessage - tool usage attribution", () => {
expect(mockTask.consecutiveMistakeCount).toBe(1)
expect(mockTask.didAlreadyUseTool).toBe(false)
})

it("routes MCP approval feedback through the queued-ack wrapper", async () => {
mockTask.providerRef = {
deref: () => ({
getState: vi.fn().mockResolvedValue({
mode: "code",
customModes: [],
}),
getMcpHub: () => ({
findServerNameBySanitizedName: () => "my_server",
getAllServers: () => [
{
name: "my_server",
tools: [{ name: "do_thing", enabledForPrompt: true }],
},
],
}),
}),
}
mockTask.assistantMessageContent = [
{
type: "mcp_tool_use",
id: "call_native_mcp_feedback",
name: "mcp_my_server_do_thing",
serverName: "my_server",
toolName: "do_thing",
arguments: {},
partial: false,
},
]
mockTask.ask = vi.fn().mockResolvedValue({ response: "yesButtonClicked", text: "Careful with this server" })

await presentAssistantMessage(mockTask as unknown as Task)

expect(mockTask.sayUserFeedbackAndAckQueued).toHaveBeenCalledExactlyOnceWith(
"Careful with this server",
undefined,
undefined,
)
expect(mockTask.say).not.toHaveBeenCalledWith("user_feedback", expect.anything(), expect.anything())
})

it("routes tool-repetition feedback through the queued-ack wrapper", async () => {
mockTask.toolRepetitionDetector.check = vi.fn().mockReturnValue({
allowExecution: false,
askUser: {
messageKey: "mistake_limit_reached",
messageDetail: "The tool {toolName} was called consecutively without progress.",
},
})
mockTask.apiConfiguration = { apiProvider: providerIdentifiers.anthropic }
mockTask.assistantMessageContent = [
{
type: "tool_use",
id: "call_repetition_feedback",
name: "read_file",
params: { path: "a.txt" },
nativeArgs: { path: "a.txt" },
partial: false,
},
]
mockTask.ask = vi.fn().mockResolvedValue({ response: "messageResponse", text: "Try another approach" })

await presentAssistantMessage(mockTask as unknown as Task)

expect(mockTask.sayUserFeedbackAndAckQueued).toHaveBeenCalledExactlyOnceWith(
"Try another approach",
undefined,
undefined,
)
expect(mockTask.say).not.toHaveBeenCalledWith("user_feedback", expect.anything(), expect.anything())
expect(mockTask.userMessageContent).toContainEqual(
expect.objectContaining({
type: "text",
text: expect.stringContaining("Try another approach"),
}),
)
})

it("acks tool-repetition feedback before it reaches the API turn", async () => {
mockTask.toolRepetitionDetector.check = vi.fn().mockReturnValue({
allowExecution: false,
askUser: {
messageKey: "mistake_limit_reached",
messageDetail: "The tool {toolName} was called consecutively without progress.",
},
})
mockTask.apiConfiguration = { apiProvider: providerIdentifiers.anthropic }
mockTask.assistantMessageContent = [
{
type: "tool_use",
id: "call_repetition_feedback",
name: "read_file",
params: { path: "a.txt" },
nativeArgs: { path: "a.txt" },
partial: false,
},
]
mockTask.ask = vi.fn().mockResolvedValue({ response: "messageResponse", text: "Try another approach" })
mockTask.sayUserFeedbackAndAckQueued = vi.fn().mockRejectedValue(new Error("persist failed"))

await expect(presentAssistantMessage(mockTask as unknown as Task)).rejects.toThrow("persist failed")

// The durable ack failed and the entry is re-queued for redelivery:
// the feedback must not already be part of the API turn.
expect(mockTask.userMessageContent).toHaveLength(0)
})
})

describe("undefined provider state", () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ describe("presentAssistantMessage - Unknown Tool Handling", () => {
}),
},
say: vi.fn().mockResolvedValue(undefined),
sayUserFeedbackAndAckQueued: vi.fn().mockResolvedValue(undefined),
ask: vi.fn().mockResolvedValue({ response: "yesButtonClicked" }),
}

Expand Down
17 changes: 10 additions & 7 deletions src/core/assistant-message/presentAssistantMessage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -218,7 +218,7 @@ export async function presentAssistantMessage(cline: Task) {
isProtected?: boolean,
autoApprovalContext?: AutoApprovalContext,
) => {
const { response, text, images, autoDenyDetail } = await cline.ask(
const { response, text, images, queuedMessageId, autoDenyDetail } = await cline.ask(
type,
partialMessage,
false,
Expand Down Expand Up @@ -250,8 +250,8 @@ export async function presentAssistantMessage(cline: Task) {
return false
}

await cline.sayUserFeedbackAndAckQueued(text, images, queuedMessageId)
if (text) {
await cline.say("user_feedback", text, images)
pushToolResult(formatResponse.toolResult(formatResponse.toolDeniedWithFeedback(text), images))
} else {
pushToolResult(formatResponse.toolDenied())
Expand All @@ -263,8 +263,8 @@ export async function presentAssistantMessage(cline: Task) {
// Store approval feedback to be merged into tool result (GitHub #10465)
// Don't push it as a separate tool_result here - that would create duplicates.
// The tool will call pushToolResult, which will merge the feedback into the actual result.
await cline.sayUserFeedbackAndAckQueued(text, images, queuedMessageId)
if (text) {
await cline.say("user_feedback", text, images)
approvalFeedback = { text, images }
}

Expand Down Expand Up @@ -787,12 +787,18 @@ export async function presentAssistantMessage(cline: Task) {
// If execution is not allowed, notify user and break.
if (!repetitionCheck.allowExecution && repetitionCheck.askUser) {
// Handle repetition similar to mistake_limit_reached pattern.
const { response, text, images } = await cline.ask(
const { response, text, images, queuedMessageId } = await cline.ask(
repetitionCheck.askUser.messageKey as ClineAsk,
repetitionCheck.askUser.messageDetail.replace("{toolName}", block.name),
)

if (response === "messageResponse") {
// Durable ack first: the feedback must reach history before
// the API turn. When the ack fails the entry is re-queued
// for redelivery, and the model must not already have seen
// the text (mirrors AttemptCompletionTool's ordering).
await cline.sayUserFeedbackAndAckQueued(text, images, queuedMessageId)

// Add user feedback to userContent.
cline.userMessageContent.push(
{
Expand All @@ -801,9 +807,6 @@ export async function presentAssistantMessage(cline: Task) {
},
...formatResponse.imageBlocks(images),
)

// Add user feedback to chat.
await cline.say("user_feedback", text, images)
}

// Track tool repetition in telemetry via PostHog exception tracking and event.
Expand Down
17 changes: 17 additions & 0 deletions src/core/message-queue/MessageQueueService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,23 @@ export class MessageQueueService extends EventEmitter<QueueEvents> {
return this._messages.length === 0
}

/**
* Reserve a specific queued message by ID when it is still present and
* unclaimed. Lets a consumer hold a consumed entry through its durable ack
* so neither another ask nor a background drain can claim it again
* mid-persistence.
*/
public claimMessage(id: string): boolean {
if (this.claimedMessageIds.has(id)) {
return false
}
if (!this._messages.some((message) => message.id === id)) {
return false
}
this.claimedMessageIds.add(id)
return true
}

/**
* Whether at least one queued message is still available to be claimed.
*
Expand Down
Loading
Loading