Skip to content
28 changes: 26 additions & 2 deletions packages/telemetry/src/TelemetryService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import {
type TelemetryPropertiesProvider,
TelemetryEventName,
type TelemetrySetting,
type ToolUsage,
} from "@roo-code/types"

/**
Expand Down Expand Up @@ -86,8 +87,31 @@ export class TelemetryService {
this.captureEvent(TelemetryEventName.TASK_RESTARTED, { taskId })
}

public captureTaskCompleted(taskId: string): void {
this.captureEvent(TelemetryEventName.TASK_COMPLETED, { taskId })
/**
* Captures task completion, summarizing per-task tool and message counts that
* were previously reported as separate per-turn events to reduce event volume.
*
* A single task may emit this more than once (e.g. an "idle" or "shutdown"
* installment followed by a final "attempt_completion" one). toolsUsed and
* messageCount are always deltas since the previous emission for that task,
* not running totals -- summing installments for a taskId reconstructs the
* full-task counts without double-counting.
*
* Note: "attempt_completion" means the model called that tool, not that the
* user accepted the result.
*/
public captureTaskCompleted(
taskId: string,
toolsUsed?: ToolUsage,
messageCount?: { user: number; assistant: number },
completionReason: "attempt_completion" | "idle" | "shutdown" = "attempt_completion",
): void {
this.captureEvent(TelemetryEventName.TASK_COMPLETED, {
taskId,
completionReason,
...(toolsUsed !== undefined && { toolsUsed }),
...(messageCount !== undefined && { messageCount }),
})
}

public captureConversationMessage(taskId: string, source: "user" | "assistant"): void {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
// pnpm --filter @roo-code/telemetry test src/__tests__/TelemetryService.task-completed.test.ts

import { TelemetryEventName, type TelemetryClient } from "@roo-code/types"

import { TelemetryService } from "../TelemetryService"

describe("TelemetryService.captureTaskCompleted", () => {
let mockClient: TelemetryClient

beforeEach(() => {
mockClient = {
setProvider: vi.fn(),
capture: vi.fn().mockResolvedValue(undefined),
captureException: vi.fn().mockResolvedValue(undefined),
updateTelemetryState: vi.fn(),
isTelemetryEnabled: vi.fn().mockReturnValue(true),
shutdown: vi.fn().mockResolvedValue(undefined),
}
})

it("captures Task Completed with the taskId and a default 'attempt_completion' completionReason when no summary is provided", () => {
const service = new TelemetryService([mockClient])

service.captureTaskCompleted("task_1")

expect(mockClient.capture).toHaveBeenCalledWith({
event: TelemetryEventName.TASK_COMPLETED,
properties: { taskId: "task_1", completionReason: "attempt_completion" },
})
})

it("includes toolsUsed and messageCount summaries when provided", () => {
const service = new TelemetryService([mockClient])

service.captureTaskCompleted(
"task_1",
{ read_file: { attempts: 3, failures: 0 }, apply_diff: { attempts: 1, failures: 1 } },
{ user: 4, assistant: 5 },
)

expect(mockClient.capture).toHaveBeenCalledWith({
event: TelemetryEventName.TASK_COMPLETED,
properties: {
taskId: "task_1",
completionReason: "attempt_completion",
toolsUsed: { read_file: { attempts: 3, failures: 0 }, apply_diff: { attempts: 1, failures: 1 } },
messageCount: { user: 4, assistant: 5 },
},
})
})

it("includes the given completionReason for idle/shutdown installments", () => {
const service = new TelemetryService([mockClient])

service.captureTaskCompleted("task_1", { read_file: { attempts: 1, failures: 0 } }, undefined, "idle")

expect(mockClient.capture).toHaveBeenCalledWith({
event: TelemetryEventName.TASK_COMPLETED,
properties: {
taskId: "task_1",
completionReason: "idle",
toolsUsed: { read_file: { attempts: 1, failures: 0 } },
},
})
})
})
1 change: 1 addition & 0 deletions src/__tests__/history-resume-delegation.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1255,6 +1255,7 @@ describe("History resume delegation - parent metadata transitions", () => {
userMessageContent: [],
consecutiveMistakeCount: 0,
emitFinalTokenUsageUpdate: vi.fn(),
flushTelemetryInstallment: vi.fn(),
} as unknown as import("../core/task/Task").Task

const block = {
Expand Down
2 changes: 2 additions & 0 deletions src/__tests__/nested-delegation-resume.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -203,6 +203,7 @@ describe("Nested delegation resume (A → B → C)", () => {
userMessageContent: [],
consecutiveMistakeCount: 0,
emitFinalTokenUsageUpdate: vi.fn(),
flushTelemetryInstallment: vi.fn(),
} as unknown as Task

const blockC = {
Expand Down Expand Up @@ -250,6 +251,7 @@ describe("Nested delegation resume (A → B → C)", () => {
userMessageContent: [],
consecutiveMistakeCount: 0,
emitFinalTokenUsageUpdate: vi.fn(),
flushTelemetryInstallment: vi.fn(),
} as unknown as Task

const blockB = {
Expand Down
144 changes: 131 additions & 13 deletions src/core/task/Task.ts
Original file line number Diff line number Diff line change
Expand Up @@ -135,6 +135,7 @@ import { MessageManager } from "../message-manager"
import { validateAndFixToolResultIds } from "./validateToolResultIds"
import { mergeConsecutiveApiMessages } from "./mergeConsecutiveApiMessages"
import { prepareApiConversationMessage } from "./apiConversationHistory"
import { shouldAddUserMessageToHistory } from "./messageCounting"

const MAX_EXPONENTIAL_BACKOFF_SECONDS = 600 // 10 minutes
const DEFAULT_USAGE_COLLECTION_TIMEOUT_MS = 5000 // 5 seconds
Expand Down Expand Up @@ -322,6 +323,25 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
consecutiveNoAssistantMessagesCount: number = 0
toolUsage: ToolUsage = {}

// Conversation message counts, summarized once per Task Completed
// installment instead of emitting a separate telemetry event per turn.
messageCounts: { user: number; assistant: number } = { user: 0, assistant: 0 }

// Idle/shutdown telemetry flush: reports toolUsage/messageCounts for tasks that
// go quiet or get torn down without the model ever calling attempt_completion
// (or without the user accepting it), so long-running/abandoned tasks aren't
// invisible to telemetry. Each flush reports only what changed since the previous
// one, tracked via telemetryToolUsageBaseline/telemetryMessageCountsBaseline --
// task.toolUsage/messageCounts themselves are never mutated by this, since they're
// also read as running totals by the public TaskCompleted API event and the UI.
// Checked on an interval rather than hooked into every say()/ask() call site.
private static readonly IDLE_TELEMETRY_CHECK_INTERVAL_MS = 5 * 60 * 1000
private static readonly IDLE_TELEMETRY_THRESHOLD_MS = 30 * 60 * 1000
private idleTelemetryCheckInterval?: NodeJS.Timeout
private lastTelemetryFlushAt: number = Date.now()
private telemetryToolUsageBaseline: ToolUsage = {}
private telemetryMessageCountsBaseline: { user: number; assistant: number } = { user: 0, assistant: 0 }

// Checkpoints
enableCheckpoints: boolean
checkpointTimeout: number
Expand Down Expand Up @@ -601,6 +621,7 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {

if (startTask) {
this._started = true
this.startIdleTelemetryCheck()
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if (task || images) {
void this.startTask(task, images).catch((error) => {
console.error("[Task#constructor] startTask failed:", error)
Expand Down Expand Up @@ -1878,6 +1899,7 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
return
}
this._started = true
this.startIdleTelemetryCheck()

const { task, images } = this.metadata

Expand Down Expand Up @@ -2270,6 +2292,17 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
public dispose(): void {
console.log(`[Task#dispose] disposing task ${this.taskId}.${this.instanceId}`)

// Stop the idle telemetry check and report any unflushed activity as a
// shutdown installment, so a task torn down mid-work (panel closed, task
// switched, extension deactivated) isn't invisible to telemetry.
try {
clearInterval(this.idleTelemetryCheckInterval)
this.idleTelemetryCheckInterval = undefined
this.flushTelemetryInstallment("shutdown")
} catch (error) {
console.error("Error flushing shutdown telemetry:", error)
}

// Cancel any in-progress HTTP request
try {
this.cancelCurrentRequest()
Expand Down Expand Up @@ -2629,18 +2662,16 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
// Add environment details as its own text block, separate from tool
// results.
const finalUserContent = [...contentWithoutEnvDetails, { type: "text" as const, text: environmentDetails }]
// Only add user message to conversation history if:
// 1. This is the first attempt (retryAttempt === 0), AND
// 2. The original userContent was not empty (empty signals delegation resume where
// the user message with tool_result and env details is already in history), OR
// 3. The message was removed in a previous iteration (userMessageWasRemoved === true)
// This prevents consecutive user messages while allowing re-add when needed
// See shouldAddUserMessageToHistory for the full add/skip rules (retry/empty/removed).
const isEmptyUserContent = currentUserContent.length === 0
const shouldAddUserMessage =
((currentItem.retryAttempt ?? 0) === 0 && !isEmptyUserContent) || currentItem.userMessageWasRemoved
const shouldAddUserMessage = shouldAddUserMessageToHistory({
retryAttempt: currentItem.retryAttempt,
isEmptyUserContent,
userMessageWasRemoved: currentItem.userMessageWasRemoved,
})
if (shouldAddUserMessage) {
await this.addToApiConversationHistory({ role: "user", content: finalUserContent })
TelemetryService.instance.captureConversationMessage(this.taskId, "user")
this.messageCounts.user++
}

// Since we sent off a placeholder api_req_started message to update the
Expand Down Expand Up @@ -3557,7 +3588,7 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
)
this.assistantMessageSavedToHistory = true

TelemetryService.instance.captureConversationMessage(this.taskId, "assistant")
this.messageCounts.assistant++
}

// Present any partial blocks that were just completed.
Expand Down Expand Up @@ -3661,8 +3692,19 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
if (this.apiConversationHistory.length > 0) {
const lastMessage = this.apiConversationHistory[this.apiConversationHistory.length - 1]
if (lastMessage.role === "user") {
// Remove the last user message that we added earlier
// Remove the last user message that we added earlier. Decrement
// messageCounts.user to match -- both retry branches below mark
// userMessageWasRemoved so the message (and its count) is restored
// exactly once when the retry succeeds, keeping the total symmetric
// regardless of how many empty-response cycles occur first.
// Guard: only reverse a count this iteration actually added. The
// popped message may predate this turn (resumed history, or a
// message appended by flushPendingToolResultsToHistory, neither
// of which incremented messageCounts.user).
this.apiConversationHistory.pop()
if (shouldAddUserMessage) {
this.messageCounts.user--
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated
}
}

Expand Down Expand Up @@ -3706,32 +3748,45 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
if (response === "yesButtonClicked") {
await this.say("api_req_retried")

// Push the same content back to retry
// Push the same content back to retry. Mark that user message was
// removed (same as the auto-retry path above) so it gets re-added --
// and messageCounts.user re-incremented -- on the retried attempt;
// otherwise shouldAddUserMessageToHistory sees retryAttempt > 0 with
// userMessageWasRemoved unset and skips re-adding it entirely.
stack.push({
userContent: currentUserContent,
includeFileDetails: false,
retryAttempt: (currentItem.retryAttempt ?? 0) + 1,
userMessageWasRemoved: true,
})

// Continue to retry the request
continue
} else {
// User declined to retry
// Re-add the user message we removed.
// Re-add the user message we removed (see messageCounts.user-- above)
// and increment messageCounts.user to match, same as the normal
// add-to-history path -- otherwise this abandoned-task path
// permanently undercounts by one.
await this.addToApiConversationHistory({
role: "user",
content: currentUserContent,
})
this.messageCounts.user++

await this.say(
"error",
"Unexpected API Response: The language model did not provide any assistant messages. This may indicate an issue with the API or the model's output.",
)

// Synthetic assistant message recording the failure -- increment
// messageCounts.assistant to match, same as the normal
// assistant-message-saved path.
await this.addToApiConversationHistory({
role: "assistant",
content: [{ type: "text", text: "Failure: I did not provide a response." }],
})
this.messageCounts.assistant++
}
}
}
Expand Down Expand Up @@ -4679,6 +4734,69 @@ export class Task extends EventEmitter<TaskEvents> implements TaskLike {
}
}

/**
* Emits a Task Completed installment for whatever toolUsage/messageCounts have
* changed since the previous installment (from any reason), then advances the
* telemetry baseline so a later installment reports only its own delta. Does
* NOT touch task.toolUsage/messageCounts themselves -- those stay running totals
* for the public TaskCompleted API event and the UI. No-ops if nothing changed
* since the last installment, so idle/shutdown checks don't emit empty events
* for tasks that were already fully reported (e.g. right after attempt_completion).
*/
public flushTelemetryInstallment(reason: "attempt_completion" | "idle" | "shutdown"): void {
const toolUsageDelta: ToolUsage = {}

for (const [toolName, usage] of Object.entries(this.toolUsage) as [ToolName, ToolUsage[ToolName]][]) {
if (!usage) {
continue
}

const baseline = this.telemetryToolUsageBaseline[toolName]
const attempts = usage.attempts - (baseline?.attempts ?? 0)
const failures = usage.failures - (baseline?.failures ?? 0)

if (attempts > 0 || failures > 0) {
toolUsageDelta[toolName] = { attempts, failures }
}
}

const messageCountDelta = {
user: Math.max(0, this.messageCounts.user - this.telemetryMessageCountsBaseline.user),
assistant: Math.max(0, this.messageCounts.assistant - this.telemetryMessageCountsBaseline.assistant),
}

const hasToolUsageDelta = Object.keys(toolUsageDelta).length > 0
const hasMessageDelta = messageCountDelta.user > 0 || messageCountDelta.assistant > 0

if (!hasToolUsageDelta && !hasMessageDelta) {
return
}

this.emitFinalTokenUsageUpdate()
TelemetryService.instance.captureTaskCompleted(this.taskId, toolUsageDelta, messageCountDelta, reason)

this.telemetryToolUsageBaseline = JSON.parse(JSON.stringify(this.toolUsage))
this.telemetryMessageCountsBaseline = { ...this.messageCounts }
this.lastTelemetryFlushAt = Date.now()
}

startIdleTelemetryCheck(): void {
this.idleTelemetryCheckInterval = setInterval(() => {
// Measure idleness from the later of the last activity and the last flush.
// Using lastMessageTs alone would keep the condition true forever after the
// first idle flush, re-running the empty-delta check on every interval tick.
const lastEventAt = Math.max(this.lastMessageTs ?? 0, this.lastTelemetryFlushAt)
const idleForMs = Date.now() - lastEventAt

if (idleForMs >= Task.IDLE_TELEMETRY_THRESHOLD_MS) {
this.flushTelemetryInstallment("idle")
}
}, Task.IDLE_TELEMETRY_CHECK_INTERVAL_MS)

// Don't hold the process open just for this timer.
this.idleTelemetryCheckInterval?.unref?.()
}

// Getters

public get taskStatus(): TaskStatus {
Expand Down
Loading
Loading