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
5 changes: 5 additions & 0 deletions src/CodexAcpClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ import type {
ThreadGoal,
ThreadGoalStatus,
ThreadSourceKind,
ThreadTokenUsage,
TurnCompletedNotification,
TurnSteerResponse,
UserInput,
Expand Down Expand Up @@ -656,6 +657,10 @@ export class CodexAcpClient {
});
}

getThreadTokenUsage(sessionId: string): ThreadTokenUsage | null {
return this.codexClient.getThreadTokenUsage(sessionId);
}

async waitForSessionNotifications(sessionId: string): Promise<void> {
while (true) {
const queue = this.sessionNotificationQueues.get(sessionId);
Expand Down
74 changes: 48 additions & 26 deletions src/CodexAcpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ import {
REASONING_EFFORT_CONFIG_ID,
} from "./ModelConfigOption";
import type {TokenCount} from "./TokenCount";
import {toPromptUsage} from "./TokenCount";
import {subtractTokenCounts, toPromptUsage, toTokenCount} from "./TokenCount";
import {CodexCommands} from "./CodexCommands";
import {SteeringQueue} from "./SteeringQueue";
import type {QuotaMeta} from "./QuotaMeta";
Expand Down Expand Up @@ -1878,6 +1878,12 @@ export class CodexAcpServer {
prompt: params.prompt,
});
const sessionState = this.getSessionState(params.sessionId);
const latestThreadTokenUsage = this.codexAcpClient.getThreadTokenUsage(params.sessionId);
if (latestThreadTokenUsage != null) {
sessionState.totalTokenUsage = toTokenCount(latestThreadTokenUsage.total);
sessionState.modelContextWindow = latestThreadTokenUsage.modelContextWindow;
}
const promptStartTokenUsage = sessionState.totalTokenUsage;
sessionState.currentTurnId = null;
sessionState.lastTokenUsage = null;
const activePrompt = this.trackActivePrompt(params.sessionId);
Expand Down Expand Up @@ -1915,7 +1921,7 @@ export class CodexAcpServer {
elicitationHandler);

if (activePrompt.signal.aborted) {
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}

const commandPromise = this.availableCommands.tryHandleCommand(params.prompt, sessionState, {
Expand Down Expand Up @@ -1957,28 +1963,32 @@ export class CodexAcpServer {
this.cancelBeforeTurnStarted(activePrompt),
]);
if (commandResult === null) {
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}
if (commandResult.handled) {
logger.log("Prompt handled by a command");
await this.codexAcpClient.waitForSessionNotifications(params.sessionId);
if (commandResult.turnCompleted?.turn.status === "interrupted") {
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}
const error = eventHandler.getFailure();
if (error) {
// noinspection ExceptionCaughtLocallyJS
throw error;
}
const promptTokenUsage = subtractTokenCounts(
sessionState.totalTokenUsage,
promptStartTokenUsage,
);
return {
stopReason: "end_turn",
usage: this.buildPromptUsage(sessionState.lastTokenUsage),
_meta: this.buildQuotaMeta(sessionState),
usage: this.buildPromptUsage(promptTokenUsage),
_meta: this.buildQuotaMeta(sessionState, promptTokenUsage),
};
}

if (this.sessionIsClosing(params.sessionId)) {
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}

const modelId = ModelId.fromString(sessionState.currentModelId);
Expand Down Expand Up @@ -2036,14 +2046,14 @@ export class CodexAcpServer {
]);

if (turnCompleted === null) {
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}

await this.codexAcpClient.waitForSessionNotifications(params.sessionId);

if (turnCompleted.turn.status === "interrupted") {
await eventHandler.flushPendingPlanUpdates();
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}

const error = eventHandler.getFailure();
Expand All @@ -2065,7 +2075,7 @@ export class CodexAcpServer {
activePrompt.signal,
);
if (this.promptShouldStop(params.sessionId, activePrompt)) {
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}
if (approved && !this.promptShouldStop(params.sessionId, activePrompt)) {
await this.applyCollaborationModeChange(sessionState, DEFAULT_COLLABORATION_MODE);
Expand Down Expand Up @@ -2113,13 +2123,13 @@ export class CodexAcpServer {
]);

if (turnCompleted === null) {
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}

await this.codexAcpClient.waitForSessionNotifications(params.sessionId);
if (turnCompleted.turn.status === "interrupted") {
await eventHandler.flushPendingPlanUpdates();
return this.cancelledPromptResponse(sessionState);
return this.cancelledPromptResponse(sessionState, promptStartTokenUsage);
}

const implementationError = eventHandler.getFailure();
Expand All @@ -2134,10 +2144,14 @@ export class CodexAcpServer {
this.createPromptFallbackTitle(params.prompt),
);

const promptTokenUsage = subtractTokenCounts(
sessionState.totalTokenUsage,
promptStartTokenUsage,
);
return {
stopReason: "end_turn",
usage: this.buildPromptUsage(sessionState.lastTokenUsage),
_meta: this.buildQuotaMeta(sessionState),
usage: this.buildPromptUsage(promptTokenUsage),
_meta: this.buildQuotaMeta(sessionState, promptTokenUsage),
};
} catch (err) {
logger.error(`Prompt for session ${params.sessionId} failed`, err);
Expand Down Expand Up @@ -2215,38 +2229,46 @@ export class CodexAcpServer {
}
}

private cancelledPromptResponse(sessionState: SessionState): acp.PromptResponse {
private cancelledPromptResponse(
sessionState: SessionState,
promptStartTokenUsage: TokenCount | null,
): acp.PromptResponse {
const promptTokenUsage = subtractTokenCounts(
sessionState.totalTokenUsage,
promptStartTokenUsage,
);
return {
stopReason: "cancelled",
usage: this.buildPromptUsage(sessionState.lastTokenUsage),
_meta: this.buildQuotaMeta(sessionState),
usage: this.buildPromptUsage(promptTokenUsage),
_meta: this.buildQuotaMeta(sessionState, promptTokenUsage),
};
}

private buildQuotaMeta(sessionState: SessionState): { quota: QuotaMeta } {
const lastTokenUsage = sessionState.lastTokenUsage;

private buildQuotaMeta(
sessionState: SessionState,
promptTokenUsage: TokenCount | null,
): { quota: QuotaMeta } {
// Remove the "[reasoning-level]" suffix from currentModelId if present
const modelName = sessionState.currentModelId.replace(/\[.*?]$/, '');

// FIXME: currently all tokens are reported for the current model
const modelUsage = (lastTokenUsage != null)
? [{ model: modelName, token_count: lastTokenUsage }]
const modelUsage = (promptTokenUsage != null)
? [{ model: modelName, token_count: promptTokenUsage }]
: [];

return {
quota: {
token_count: sessionState.lastTokenUsage,
token_count: promptTokenUsage,
model_usage: modelUsage
}
};
}

private buildPromptUsage(lastTokenUsage: TokenCount | null): acp.Usage | null {
if (lastTokenUsage == null) {
private buildPromptUsage(promptTokenUsage: TokenCount | null): acp.Usage | null {
if (promptTokenUsage == null) {
return null;
}
return toPromptUsage(lastTokenUsage);
return toPromptUsage(promptTokenUsage);
}

private async runWithProcessCheck<T>(operation: () => Promise<T>): Promise<T> {
Expand Down
10 changes: 10 additions & 0 deletions src/CodexAppServerClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ import type {
ThreadSettings,
ThreadStartParams,
ThreadStartResponse,
ThreadTokenUsage,
ThreadUnsubscribeParams,
ThreadUnsubscribeResponse,
ToolRequestUserInputParams,
Expand Down Expand Up @@ -145,6 +146,7 @@ export class CodexAppServerClient {
private readonly threadGoalUpdateCaptures = new Map<string, Set<(event: ThreadGoalUpdatedNotification) => void>>();
private readonly threadGoalClearedCaptures = new Map<string, Set<() => void>>();
private readonly threadSettings = new Map<string, ThreadSettings>();
private readonly threadTokenUsage = new Map<string, ThreadTokenUsage>();
private readonly staleTurnIds = new Map<string, Set<string>>();

constructor(connection: MessageConnection) {
Expand Down Expand Up @@ -186,6 +188,9 @@ export class CodexAppServerClient {
if (this.handleStaleTurnNotification(serverNotification, routing)) {
return;
}
if (serverNotification.method === "thread/tokenUsage/updated") {
this.threadTokenUsage.set(serverNotification.params.threadId, serverNotification.params.tokenUsage);
}
this.notify(serverNotification);
for (const callback of this.codexEventHandlers) {
callback({ eventType: "notification", ...serverNotification });
Expand Down Expand Up @@ -260,6 +265,7 @@ export class CodexAppServerClient {
this.notificationHandlers.delete(threadId);
this.approvalHandlers.delete(threadId);
this.elicitationHandlers.delete(threadId);
this.threadTokenUsage.delete(threadId);
}

async initialize(params: InitializeParams): Promise<InitializeResponse> {
Expand Down Expand Up @@ -532,6 +538,10 @@ export class CodexAppServerClient {
return this.threadSettings.get(threadId);
}

getThreadTokenUsage(threadId: string): ThreadTokenUsage | null {
return this.threadTokenUsage.get(threadId) ?? null;
}

async threadSettingsUpdate(params: ExperimentalThreadSettingsUpdateParams): Promise<void> {
await this.connection.sendRequest("thread/settings/update", params);
}
Expand Down
40 changes: 40 additions & 0 deletions src/TokenCount.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,46 @@ export interface TokenCount {
reasoningOutputTokens: number;
}

/**
* Returns the category-by-category usage added to a cumulative token count.
* Each category is clamped independently because an upstream counter may be
* reset or corrected between snapshots.
*/
export function subtractTokenCounts(
current: TokenCount | null,
previous: TokenCount | null,
): TokenCount | null {
if (current == null) {
return null;
}

const difference = (currentValue: number, previousValue: number): number =>
Math.max(0, currentValue - previousValue);
const baseline = previous ?? {
totalTokens: 0,
inputTokens: 0,
cachedInputTokens: 0,
outputTokens: 0,
reasoningOutputTokens: 0,
};

const inputTokens = difference(current.inputTokens, baseline.inputTokens);
const cachedInputTokens = difference(current.cachedInputTokens, baseline.cachedInputTokens);
const outputTokens = difference(current.outputTokens, baseline.outputTokens);
const reasoningOutputTokens = Math.min(
outputTokens,
difference(current.reasoningOutputTokens, baseline.reasoningOutputTokens),
);

return {
totalTokens: inputTokens + cachedInputTokens + outputTokens,
inputTokens,
cachedInputTokens,
outputTokens,
reasoningOutputTokens,
};
}

/**
* Maps Codex's TokenUsageBreakdown to our TokenCount interface.
* This explicit mapping ensures compile-time errors if Codex changes their types.
Expand Down
18 changes: 9 additions & 9 deletions src/__tests__/CodexACPAgent/data/token-usage-cancelled.json
Original file line number Diff line number Diff line change
@@ -1,29 +1,29 @@
{
"stopReason": "cancelled",
"usage": {
"totalTokens": 1500,
"inputTokens": 1200,
"totalTokens": 3000,
"inputTokens": 2500,
"cachedReadTokens": 0,
"outputTokens": 300,
"outputTokens": 500,
"thoughtTokens": 0
},
"_meta": {
"quota": {
"token_count": {
"totalTokens": 1500,
"inputTokens": 1200,
"totalTokens": 3000,
"inputTokens": 2500,
"cachedInputTokens": 0,
"outputTokens": 300,
"outputTokens": 500,
"reasoningOutputTokens": 0
},
"model_usage": [
{
"model": "model-id",
"token_count": {
"totalTokens": 1500,
"inputTokens": 1200,
"totalTokens": 3000,
"inputTokens": 2500,
"cachedInputTokens": 0,
"outputTokens": 300,
"outputTokens": 500,
"reasoningOutputTokens": 0
}
}
Expand Down
30 changes: 15 additions & 15 deletions src/__tests__/CodexACPAgent/data/token-usage-end-turn.json
Original file line number Diff line number Diff line change
@@ -1,30 +1,30 @@
{
"stopReason": "end_turn",
"usage": {
"totalTokens": 2500,
"inputTokens": 1500,
"cachedReadTokens": 500,
"outputTokens": 450,
"thoughtTokens": 50
"totalTokens": 5000,
"inputTokens": 3100,
"cachedReadTokens": 1000,
"outputTokens": 900,
"thoughtTokens": 100
},
"_meta": {
"quota": {
"token_count": {
"totalTokens": 2500,
"inputTokens": 1500,
"cachedInputTokens": 500,
"outputTokens": 450,
"reasoningOutputTokens": 50
"totalTokens": 5000,
"inputTokens": 3100,
"cachedInputTokens": 1000,
"outputTokens": 900,
"reasoningOutputTokens": 100
},
"model_usage": [
{
"model": "model-id",
"token_count": {
"totalTokens": 2500,
"inputTokens": 1500,
"cachedInputTokens": 500,
"outputTokens": 450,
"reasoningOutputTokens": 50
"totalTokens": 5000,
"inputTokens": 3100,
"cachedInputTokens": 1000,
"outputTokens": 900,
"reasoningOutputTokens": 100
}
}
]
Expand Down
Loading