Skip to content
Open
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
92 changes: 92 additions & 0 deletions packages/core/src/session/tokens.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
export const DEFAULT_MIN_TPS_ELAPSED_MS = 250
export const DEFAULT_INCLUDE_REASONING = true

export interface TokenMetrics {
output: number
reasoning: number
}

export interface TimestampMetrics {
created: number
firstToken?: number
completed?: number
}

export interface TPSResult {
rate: number
totalTokens: number
elapsedMs: number
isValid: boolean
}

export function stampFirstToken(time: TimestampMetrics, now: number): number {
if (time.firstToken === undefined) time.firstToken = now
return time.firstToken
}

export function totalGeneratedTokens(tokens: TokenMetrics, includeReasoning = DEFAULT_INCLUDE_REASONING): number {
return tokens.output + (includeReasoning ? tokens.reasoning : 0)
}

type TPSMessage = {
summary?: boolean
finish?: string | null
tokens: TokenMetrics
time: TimestampMetrics
}

function tpsInputs(msg: TPSMessage): { totalTokens: number; elapsedMs: number } | undefined {
if (msg.summary) return undefined
if (!msg.finish) return undefined
if (["tool-calls", "unknown", "error"].includes(msg.finish)) return undefined

const totalTokens = totalGeneratedTokens(msg.tokens)
if (totalTokens <= 0) return undefined
const { firstToken, completed } = msg.time
if (firstToken === undefined || completed === undefined) return undefined

return { totalTokens, elapsedMs: completed - firstToken }
}

export function isValidForTPS(msg: TPSMessage & {
minElapsedMs?: number
}): boolean {
const inputs = tpsInputs(msg)
if (!inputs) return false
const minElapsedMs = msg.minElapsedMs ?? DEFAULT_MIN_TPS_ELAPSED_MS
return inputs.elapsedMs >= minElapsedMs
}

export function calculateTPS(
totalTokens: number,
elapsedMs: number,
minElapsedMs = DEFAULT_MIN_TPS_ELAPSED_MS,
): TPSResult | undefined {
if (totalTokens <= 0) return undefined
if (elapsedMs < minElapsedMs) return undefined

const rate = totalTokens / (elapsedMs / 1000)
if (!Number.isFinite(rate) || rate < 0) return undefined

return {
rate: Math.round(rate),
totalTokens,
elapsedMs,
isValid: true,
}
}

export function formatTPS(result: TPSResult): string {
return `${result.rate.toLocaleString()} tok/s`
}

export function getMessageTPS(msg: {
summary?: boolean
finish?: string | null
tokens: TokenMetrics
time: TimestampMetrics
}): TPSResult | undefined {
const inputs = tpsInputs(msg)
if (!inputs) return undefined
return calculateTPS(inputs.totalTokens, inputs.elapsedMs)
}
139 changes: 139 additions & 0 deletions packages/core/test/session-tps.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,139 @@
import { describe, expect, test } from "bun:test"
import {
calculateTPS,
DEFAULT_MIN_TPS_ELAPSED_MS,
getMessageTPS,
isValidForTPS,
stampFirstToken,
type TimestampMetrics,
} from "../src/session/tokens"

const validMessage = {
finish: "stop",
tokens: { output: 100, reasoning: 50 },
time: { created: 1000, firstToken: 1100, completed: 2100 },
}

describe("getMessageTPS", () => {
test("calculates rounded output and reasoning tokens per second", () => {
expect(getMessageTPS(validMessage)?.rate).toBe(150)
})

test("returns no value for summary messages", () => {
expect(getMessageTPS({ ...validMessage, summary: true })).toBeUndefined()
})

test("returns no value when finish is missing", () => {
expect(getMessageTPS({ ...validMessage, finish: undefined })).toBeUndefined()
})

test("returns no value for tool-call finishes", () => {
expect(getMessageTPS({ ...validMessage, finish: "tool-calls" })).toBeUndefined()
})

test("returns no value for unknown finishes", () => {
expect(getMessageTPS({ ...validMessage, finish: "unknown" })).toBeUndefined()
})

test("returns no value for error finishes", () => {
expect(getMessageTPS({ ...validMessage, finish: "error" })).toBeUndefined()
})

test("returns no value when token total is zero", () => {
expect(getMessageTPS({ ...validMessage, tokens: { output: 0, reasoning: 0 } })).toBeUndefined()
})

test("returns no value when first token timestamp is missing", () => {
expect(getMessageTPS({ ...validMessage, time: { created: 1000, completed: 2100 } })).toBeUndefined()
})

test("returns no value when completion timestamp is missing", () => {
expect(getMessageTPS({ ...validMessage, time: { created: 1000, firstToken: 1100 } })).toBeUndefined()
})

test("returns no value when elapsed time is below 250 milliseconds", () => {
expect(
getMessageTPS({
...validMessage,
time: {
...validMessage.time,
completed: validMessage.time.firstToken! + DEFAULT_MIN_TPS_ELAPSED_MS - 1,
},
}),
).toBeUndefined()
})

test("accepts elapsed time at the 250 millisecond threshold", () => {
expect(
getMessageTPS({
...validMessage,
time: {
...validMessage.time,
completed: validMessage.time.firstToken! + DEFAULT_MIN_TPS_ELAPSED_MS,
},
})?.rate,
).toBe(600)
})
})

describe("isValidForTPS", () => {
test("rejects negative token totals", () => {
expect(isValidForTPS({ ...validMessage, tokens: { output: -1, reasoning: 0 } })).toBe(false)
})

test("rejects zero token totals", () => {
expect(isValidForTPS({ ...validMessage, tokens: { output: 0, reasoning: 0 } })).toBe(false)
})

test("rejects missing finish values", () => {
expect(isValidForTPS({ ...validMessage, finish: null })).toBe(false)
})

test("rejects invalid finish values", () => {
for (const finish of ["tool-calls", "unknown", "error"]) {
expect(isValidForTPS({ ...validMessage, finish })).toBe(false)
}
})

test("rejects a missing first token timestamp", () => {
expect(isValidForTPS({ ...validMessage, time: { created: 1000, completed: 2100 } })).toBe(false)
})

test("rejects a missing completion timestamp", () => {
expect(isValidForTPS({ ...validMessage, time: { created: 1000, firstToken: 1100 } })).toBe(false)
})

test("rejects elapsed time below the configured threshold", () => {
expect(
isValidForTPS({
...validMessage,
time: { ...validMessage.time, completed: validMessage.time.firstToken! + 249 },
}),
).toBe(false)
})
})

describe("calculateTPS", () => {
test("rounds fractional rates to the nearest integer", () => {
expect(calculateTPS(100, 327)?.rate).toBe(306)
})

test("rejects zero and negative token totals", () => {
expect(calculateTPS(0, 1000)).toBeUndefined()
expect(calculateTPS(-1, 1000)).toBeUndefined()
})

test("rejects elapsed time below the configured threshold", () => {
expect(calculateTPS(100, 249)).toBeUndefined()
})
})

describe("stampFirstToken", () => {
test("stamps the first event and preserves it for later deltas", () => {
const time: TimestampMetrics = { created: 1000 }

expect(stampFirstToken(time, 1100)).toBe(1100)
expect(stampFirstToken(time, 1200)).toBe(1100)
expect(time.firstToken).toBe(1100)
})
})
15 changes: 15 additions & 0 deletions packages/opencode/src/session/processor.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { LayerNode } from "@opencode-ai/core/effect/layer-node"
import { stampFirstToken } from "@opencode-ai/core/session/tokens"
import { PermissionV1 } from "@opencode-ai/core/v1/permission"
import { Image } from "@/image/image"
import { SessionV1 } from "@opencode-ai/core/v1/session"
Expand Down Expand Up @@ -120,6 +121,10 @@ const layer = Layer.effect(
aborted,
})

const markFirstToken = () => {
stampFirstToken(ctx.assistantMessage.time, Date.now())
}

const settleToolCall = Effect.fn("SessionProcessor.settleToolCall")(function* (toolCallID: string) {
const done = ctx.toolcalls[toolCallID]?.done
delete ctx.toolcalls[toolCallID]
Expand Down Expand Up @@ -279,6 +284,7 @@ const layer = Layer.effect(
switch (value.type) {
case "reasoning-start":
if (value.id in ctx.reasoningMap) return
markFirstToken()
ctx.reasoningMap[value.id] = {
id: PartID.ascending(),
messageID: ctx.assistantMessage.id,
Expand All @@ -294,6 +300,7 @@ const layer = Layer.effect(
case "reasoning-delta":
// Match dev: silently drop orphan deltas (no preceding reasoning-start).
if (!(value.id in ctx.reasoningMap)) return
markFirstToken()
ctx.reasoningMap[value.id].text += value.text
if (value.providerMetadata) ctx.reasoningMap[value.id].metadata = value.providerMetadata
yield* session.updatePartDelta({
Expand All @@ -316,22 +323,26 @@ const layer = Layer.effect(
if (ctx.assistantMessage.summary) {
throw new Error(`Tool call not allowed while generating summary: ${value.name}`)
}
markFirstToken()
yield* ensureToolCall(value)
return

case "tool-input-delta":
yield* ensureToolCall(value)
markFirstToken()
return

case "tool-input-end": {
yield* ensureToolCall(value)
markFirstToken()
return
}

case "tool-call": {
if (ctx.assistantMessage.summary) {
throw new Error(`Tool call not allowed while generating summary: ${value.name}`)
}
markFirstToken()
yield* ensureToolCall(value)
const input = isRecord(value.input) ? value.input : { value: value.input }
yield* updateToolCall(value.id, (match) => ({
Expand Down Expand Up @@ -433,6 +444,7 @@ const layer = Layer.effect(
return

case "step-finish": {
markFirstToken()
const completedSnapshot = yield* snapshot.track()
yield* Effect.forEach(Object.keys(ctx.reasoningMap), finishReasoning)
const usage = Session.getUsage({
Expand Down Expand Up @@ -484,6 +496,7 @@ const layer = Layer.effect(
}

case "text-start":
markFirstToken()
ctx.currentText = {
id: PartID.ascending(),
messageID: ctx.assistantMessage.id,
Expand All @@ -498,6 +511,7 @@ const layer = Layer.effect(

case "text-delta":
if (!ctx.currentText) return
markFirstToken()
ctx.currentText.text += value.text
if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata
yield* session.updatePartDelta({
Expand Down Expand Up @@ -532,6 +546,7 @@ const layer = Layer.effect(
return

case "finish":
markFirstToken()
return
}
})
Expand Down
55 changes: 55 additions & 0 deletions packages/opencode/test/session/processor-effect.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -209,6 +209,27 @@ const providerErrorLLM = Layer.succeed(
const providerErrorEnv = LayerNode.compile(root, [...replacements, [LLM.node, providerErrorLLM]])
const itProviderError = testEffect(providerErrorEnv)

const toolOnlyLLM = Layer.succeed(
LLM.Service,
LLM.Service.of({
stream: () =>
Stream.make(
LLMEvent.stepStart({ index: 0 }),
LLMEvent.toolCall({ id: "call-only", name: "lookup", input: {}, providerExecuted: true }),
LLMEvent.toolResult({
id: "call-only",
name: "lookup",
result: { type: "json", value: { output: "tool result" } },
providerExecuted: true,
}),
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
LLMEvent.finish({ reason: "stop" }),
),
}),
)
const toolOnlyEnv = LayerNode.compile(root, [...replacements, [LLM.node, toolOnlyLLM]])
const itToolOnly = testEffect(toolOnlyEnv)

const fragmentFailureLLM = Layer.succeed(
LLM.Service,
LLM.Service.of({
Expand Down Expand Up @@ -1117,6 +1138,40 @@ itProviderError.live("session.processor effect tests fail provider-executed erro
),
)

itToolOnly.live("session.processor effect tests stamp first token for tool-only turns", () =>
provideTmpdirInstance(
(dir) =>
Effect.gen(function* () {
const { processors, session, provider } = yield* boot()
const chat = yield* session.create({})
const parent = yield* user(chat.id, "tool-only")
const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
const handle = yield* processors.create({ assistantMessage: msg, sessionID: chat.id, model: mdl })

yield* handle.process({
user: {
id: parent.id,
sessionID: chat.id,
role: "user",
time: parent.time,
agent: parent.agent,
model: { providerID: ref.providerID, modelID: ref.modelID },
} satisfies SessionV1.User,
sessionID: chat.id,
model: mdl,
agent: agent(),
system: [],
messages: [{ role: "user", content: "tool-only" }],
tools: {},
})

expect(handle.message.time.firstToken).toBeDefined()
}),
{ config: cfg },
),
)

itFragmentFailure.live("session.processor effect tests retain partial legacy parts without v2 events", () =>
provideTmpdirInstance(
(dir) =>
Expand Down
Loading
Loading