Skip to content
Closed
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
7 changes: 7 additions & 0 deletions packages/types/src/message.ts
Original file line number Diff line number Diff line change
Expand Up @@ -294,11 +294,18 @@ export type TokenUsage = z.infer<typeof tokenUsageSchema>
* QueuedMessage
*/

export const queuedMessageDeliveryModeSchema = z.enum(["queue", "steer"])

export type QueuedMessageDeliveryMode = z.infer<typeof queuedMessageDeliveryModeSchema>

export const queuedMessageSchema = z.object({
timestamp: z.number(),
createdAt: z.number(),
updatedAt: z.number(),
id: z.string(),
text: z.string(),
images: z.array(z.string()).optional(),
deliveryMode: queuedMessageDeliveryModeSchema.default("queue"),
})

export type QueuedMessage = z.infer<typeof queuedMessageSchema>
4 changes: 3 additions & 1 deletion packages/types/src/task.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ export interface TaskProviderLike {
options?: CreateTaskOptions,
configuration?: RooCodeSettings,
): Promise<TaskLike>
cancelTask(): Promise<void>
cancelTask(taskId?: string): Promise<void>
clearTask(): Promise<void>
resumeTask(taskId: string): void

Expand Down Expand Up @@ -94,6 +94,8 @@ export interface CreateTaskOptions {
consecutiveMistakeLimit?: number
experiments?: Record<string, boolean>
initialTodos?: TodoItem[]
/** Whether the task should become the visible task in the UI (default: true). */
focus?: boolean
/** Initial status for the task's history item (e.g., "active" for child tasks) */
initialStatus?: "active" | "delegated" | "completed"
/** Whether to start the task loop immediately (default: true).
Expand Down
16 changes: 15 additions & 1 deletion packages/types/src/vscode-extension-host.ts
Original file line number Diff line number Diff line change
Expand Up @@ -245,6 +245,18 @@ export interface OpenAiCodexRateLimitsMessage {
error?: string
}

export interface ActiveConversationSummary {
rootTaskId: string
activeTaskId: string
rootTask: string
activeTask: string
ts: number
status: "running" | "interactive" | "resumable" | "idle" | "none"
parentTaskId?: string
queuedMessageCount: number
steerMessageCount: number
}

export type ExtensionState = Pick<
GlobalSettings,
| "currentApiConfigName"
Expand Down Expand Up @@ -315,6 +327,7 @@ export type ExtensionState = Pick<
apiConfiguration: ProviderSettings
uriScheme?: string
shouldShowAnnouncement: boolean
activeConversations?: ActiveConversationSummary[]

taskHistory: HistoryItem[]

Expand Down Expand Up @@ -411,7 +424,7 @@ export interface UpdateTodoListPayload {
todos: any[]
}

export type EditQueuedMessagePayload = Pick<QueuedMessage, "id" | "text" | "images">
export type EditQueuedMessagePayload = Pick<QueuedMessage, "id" | "text" | "images" | "deliveryMode">

export interface WebviewMessage {
type:
Expand Down Expand Up @@ -596,6 +609,7 @@ export interface WebviewMessage {
askResponse?: ClineAskResponse
apiConfiguration?: ProviderSettings
images?: string[]
deliveryMode?: QueuedMessage["deliveryMode"]
bool?: boolean
value?: number
stepIndex?: number
Expand Down
16 changes: 16 additions & 0 deletions src/core/assistant-message/presentAssistantMessage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -932,6 +932,22 @@ export async function presentAssistantMessage(cline: Task) {
// locked.
cline.presentAssistantMessageLocked = false

const completedToolBoundary =
(block.type === "tool_use" || block.type === "mcp_tool_use") &&
!block.partial &&
!cline.didRejectTool &&
!cline.didAlreadyUseTool

if (completedToolBoundary) {
const didInterruptForSteer = await cline.maybeInterruptForPendingSteerAtToolBoundary(
cline.currentStreamingContentIndex,
)

if (didInterruptForSteer) {
return
}
}

// NOTE: When tool is rejected, iterator stream is interrupted and it waits
// for `userMessageContentReady` to be true. Future calls to present will
// skip execution since `didRejectTool` and iterate until `contentIndex` is
Expand Down
188 changes: 163 additions & 25 deletions src/core/message-queue/MessageQueueService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { EventEmitter } from "events"

import { v4 as uuidv4 } from "uuid"

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

export interface MessageQueueState {
messages: QueuedMessage[]
Expand All @@ -15,84 +15,222 @@ export interface QueueEvents {
}

export class MessageQueueService extends EventEmitter<QueueEvents> {
private _messages: QueuedMessage[]
private readonly lanes: Record<QueuedMessageDeliveryMode, QueuedMessage[]>

constructor() {
super()

this._messages = []
this.lanes = {
queue: [],
steer: [],
}
}

private findMessage(id: string) {
const index = this._messages.findIndex((msg) => msg.id === id)
private normalizeMessage(
message: Pick<QueuedMessage, "id" | "text"> & Partial<Omit<QueuedMessage, "id" | "text">>,
): QueuedMessage {
const createdAt = message.createdAt ?? message.timestamp ?? Date.now()
const updatedAt = message.updatedAt ?? message.timestamp ?? createdAt

return {
timestamp: message.timestamp ?? updatedAt,
createdAt,
updatedAt,
id: message.id,
text: message.text,
images: message.images ? [...message.images] : undefined,
deliveryMode: message.deliveryMode ?? "queue",
}
}

private emitStateChanged(): void {
this.emit("stateChanged", this.messages)
}

if (index === -1) {
return { index, message: undefined }
private getLane(deliveryMode: QueuedMessageDeliveryMode): QueuedMessage[] {
return this.lanes[deliveryMode]
}

private findMessage(id: string) {
for (const deliveryMode of ["steer", "queue"] as const) {
const lane = this.getLane(deliveryMode)
const index = lane.findIndex((message) => message.id === id)

if (index !== -1) {
return {
deliveryMode,
index,
message: lane[index],
}
}
}

return { index, message: this._messages[index] }
return {
deliveryMode: undefined,
index: -1,
message: undefined,
}
}

public addMessage(text: string, images?: string[]): QueuedMessage | undefined {
public addMessage(
text: string,
images?: string[],
deliveryMode: QueuedMessageDeliveryMode = "queue",
): QueuedMessage | undefined {
if (!text && !images?.length) {
return undefined
}

const now = Date.now()
const message: QueuedMessage = {
timestamp: Date.now(),
timestamp: now,
createdAt: now,
updatedAt: now,
id: uuidv4(),
text,
images,
deliveryMode,
}

this._messages.push(message)
this.emit("stateChanged", this._messages)
this.getLane(deliveryMode).push(message)
this.emitStateChanged()

return message
}

public removeMessage(id: string): boolean {
const { index, message } = this.findMessage(id)
const { deliveryMode, index, message } = this.findMessage(id)

if (!message) {
if (!message || !deliveryMode) {
return false
}

this._messages.splice(index, 1)
this.emit("stateChanged", this._messages)
this.getLane(deliveryMode).splice(index, 1)
this.emitStateChanged()
return true
}

public updateMessage(id: string, text: string, images?: string[]): boolean {
public updateMessage(
id: string,
text: string,
images?: string[],
deliveryMode?: QueuedMessageDeliveryMode,
): boolean {
const now = Date.now()
const { message } = this.findMessage(id)

if (!message) {
return false
}

message.timestamp = Date.now()
const nextDeliveryMode = deliveryMode ?? message.deliveryMode

if (nextDeliveryMode !== message.deliveryMode) {
return this.setDeliveryMode(id, nextDeliveryMode, { text, images, updatedAt: now })
}

message.timestamp = now
message.updatedAt = now
message.text = text
message.images = images
this.emit("stateChanged", this._messages)
this.emitStateChanged()
return true
}

public dequeueMessage(): QueuedMessage | undefined {
const message = this._messages.shift()
this.emit("stateChanged", this._messages)
return this.dequeueNextMessage(["steer", "queue"])
}

public dequeueMessageByMode(deliveryMode: QueuedMessageDeliveryMode): QueuedMessage | undefined {
const message = this.getLane(deliveryMode).shift()

if (message) {
this.emitStateChanged()
}

return message
}

public dequeueNextMessage(priority: QueuedMessageDeliveryMode[] = ["steer", "queue"]): QueuedMessage | undefined {
for (const deliveryMode of priority) {
const message = this.dequeueMessageByMode(deliveryMode)

if (message) {
return message
}
}

return undefined
}

public setDeliveryMode(
id: string,
deliveryMode: QueuedMessageDeliveryMode,
overrides?: { text?: string; images?: string[]; updatedAt?: number },
): boolean {
const { deliveryMode: currentDeliveryMode, index, message } = this.findMessage(id)

if (!message || !currentDeliveryMode) {
return false
}

const now = overrides?.updatedAt ?? Date.now()
const nextMessage: QueuedMessage = {
...message,
timestamp: now,
updatedAt: now,
text: overrides?.text ?? message.text,
images: overrides?.images ?? message.images,
deliveryMode,
}

if (currentDeliveryMode === deliveryMode) {
this.getLane(currentDeliveryMode)[index] = nextMessage
this.emitStateChanged()
return true
}

this.getLane(currentDeliveryMode).splice(index, 1)
this.getLane(deliveryMode).push(nextMessage)
this.emitStateChanged()
return true
}

public restoreMessages(messages: QueuedMessage[]): void {
this.lanes.queue = []
this.lanes.steer = []

for (const message of messages) {
const normalized = this.normalizeMessage(message)
this.getLane(normalized.deliveryMode).push(normalized)
}

this.emitStateChanged()
}

public get messages(): QueuedMessage[] {
return this._messages
return [...this.lanes.steer, ...this.lanes.queue]
}

public getMessagesByMode(deliveryMode: QueuedMessageDeliveryMode): QueuedMessage[] {
return [...this.getLane(deliveryMode)]
}

public hasMessages(deliveryMode?: QueuedMessageDeliveryMode): boolean {
if (deliveryMode) {
return this.getLane(deliveryMode).length > 0
}

return this.messages.length > 0
}

public isEmpty(): boolean {
return this._messages.length === 0
public isEmpty(deliveryMode?: QueuedMessageDeliveryMode): boolean {
return !this.hasMessages(deliveryMode)
}

public dispose(): void {
this._messages = []
this.lanes.queue = []
this.lanes.steer = []
this.removeAllListeners()
}
}
37 changes: 37 additions & 0 deletions src/core/message-queue/__tests__/MessageQueueService.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
import { MessageQueueService } from "../MessageQueueService"

describe("MessageQueueService", () => {
it("keeps queue and steer lanes separate while exposing steer-first ordering", () => {
const service = new MessageQueueService()

const queued = service.addMessage("queued", undefined, "queue")
const steer = service.addMessage("steer", undefined, "steer")

expect(service.messages.map((message) => message.text)).toEqual(["steer", "queued"])
expect(service.getMessagesByMode("queue")).toEqual([queued])
expect(service.getMessagesByMode("steer")).toEqual([steer])

expect(service.dequeueMessageByMode("queue")?.text).toBe("queued")
expect(service.dequeueMessageByMode("steer")?.text).toBe("steer")
expect(service.isEmpty()).toBe(true)
})

it("can move messages between lanes without losing timestamps or images", () => {
const service = new MessageQueueService()
const message = service.addMessage("queued", ["image.png"], "queue")

expect(message).toBeDefined()
expect(service.updateMessage(message!.id, "steer now", ["image.png"], "steer")).toBe(true)

expect(service.getMessagesByMode("queue")).toEqual([])
expect(service.getMessagesByMode("steer")).toEqual([
expect.objectContaining({
id: message!.id,
text: "steer now",
images: ["image.png"],
deliveryMode: "steer",
createdAt: message!.createdAt,
}),
])
})
})
Loading