Skip to content
Merged
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
49 changes: 45 additions & 4 deletions src/client.js
Original file line number Diff line number Diff line change
Expand Up @@ -249,7 +249,7 @@ export class WindsurfClient {
* @param {object} opts - { onChunk, onEnd, onError }
*/
async cascadeChat(messages, modelEnum, modelUid, opts = {}) {
const { onChunk, onEnd, onError, signal, reuseEntry, toolPreamble } = opts;
let { onChunk, onEnd, onError, signal, reuseEntry, toolPreamble } = opts;
const aborted = () => signal?.aborted;
const inputChars = messages.reduce((n, m) => n + contentToString(m?.content).length, 0);

Expand Down Expand Up @@ -294,6 +294,42 @@ export class WindsurfClient {
cascadeId = await openCascade();
}

// A resumed cascade already contains every prior turn in its trajectory.
// If we poll from step offset 0 again, the old planner-response steps are
// replayed as fresh output and both text and usage grow cumulatively
// across turns (`alpha` -> `alphabeta` -> ...). Store absolute offsets in
// the conversation pool and reuse them here; fall back to a one-shot
// snapshot so entries created before this fix still resume safely.
let stepOffset = Number.isInteger(reuseEntry?.stepOffset) && reuseEntry.stepOffset >= 0
? reuseEntry.stepOffset
: 0;
let generatorOffset = Number.isInteger(reuseEntry?.generatorOffset) && reuseEntry.generatorOffset >= 0
? reuseEntry.generatorOffset
: 0;
if (reuseEntry?.cascadeId && (!Number.isInteger(reuseEntry?.stepOffset) || !Number.isInteger(reuseEntry?.generatorOffset))) {
try {
if (!Number.isInteger(reuseEntry?.stepOffset)) {
const resumeStepsResp = await grpcUnary(
this.port, this.csrfToken,
`${LS_SERVICE}/GetCascadeTrajectorySteps`,
grpcFrame(buildGetTrajectoryStepsRequest(cascadeId, 0))
);
stepOffset = parseTrajectorySteps(resumeStepsResp).length;
}
if (!Number.isInteger(reuseEntry?.generatorOffset)) {
const resumeMetaResp = await grpcUnary(
this.port, this.csrfToken,
`${LS_SERVICE}/GetCascadeTrajectoryGeneratorMetadata`,
grpcFrame(buildGetGeneratorMetadataRequest(cascadeId, 0)),
5000
);
generatorOffset = parseGeneratorMetadata(resumeMetaResp)?.entryCount || 0;
}
} catch (e) {
log.warn(`Cascade resume snapshot failed: ${e.message}`);
}
}

let text;
let images = [];
const systemMsgs = messages.filter(m => m.role === 'system');
Expand Down Expand Up @@ -433,7 +469,7 @@ export class WindsurfClient {
pollCount++;

// Get steps
const stepsProto = buildGetTrajectoryStepsRequest(cascadeId, 0);
const stepsProto = buildGetTrajectoryStepsRequest(cascadeId, stepOffset);
const stepsResp = await grpcUnary(
this.port, this.csrfToken, `${LS_SERVICE}/GetCascadeTrajectorySteps`, grpcFrame(stepsProto)
);
Expand Down Expand Up @@ -605,6 +641,7 @@ export class WindsurfClient {
this.port, this.csrfToken, `${LS_SERVICE}/GetCascadeTrajectorySteps`, grpcFrame(stepsProto)
);
const finalSteps = parseTrajectorySteps(finalResp);
lastStepCount = finalSteps.length;
for (let i = 0; i < finalSteps.length; i++) {
const step = finalSteps[i];
const responseText = step.responseText || '';
Expand Down Expand Up @@ -652,7 +689,7 @@ export class WindsurfClient {
polls: pollCount,
textLen: totalYielded,
thinkingLen: totalThinking,
stepCount: Math.max(yieldedByStep.size, thinkingByStep.size, lastStepCount),
stepCount: stepOffset + Math.max(yieldedByStep.size, thinkingByStep.size, lastStepCount),
toolCalls: seenToolCallIds.size,
sawActive,
sawText,
Expand All @@ -676,7 +713,7 @@ export class WindsurfClient {
// itself is already formed.
let serverUsage = null;
try {
const metaReq = buildGetGeneratorMetadataRequest(cascadeId, 0);
const metaReq = buildGetGeneratorMetadataRequest(cascadeId, generatorOffset);
const metaResp = await grpcUnary(
this.port, this.csrfToken,
`${LS_SERVICE}/GetCascadeTrajectoryGeneratorMetadata`,
Expand Down Expand Up @@ -712,6 +749,10 @@ export class WindsurfClient {
// that iterate over it keep working.
chunks.cascadeId = cascadeId;
chunks.sessionId = sessionId;
chunks.stepOffset = stepOffset + Math.max(yieldedByStep.size, thinkingByStep.size, lastStepCount);
chunks.generatorOffset = serverUsage?.entryCount != null
? generatorOffset + serverUsage.entryCount
: null;
chunks.toolCalls = toolCalls;
chunks.usage = serverUsage;
if (serverUsage) {
Expand Down
8 changes: 7 additions & 1 deletion src/conversation-pool.js
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,11 @@ function positiveIntEnv(name, fallback) {
const POOL_TTL_MS = positiveIntEnv('CASCADE_POOL_TTL_MS', 30 * 60 * 1000);
const POOL_MAX = 500;

// fingerprint -> { cascadeId, sessionId, lsPort, apiKey, createdAt, lastAccess }
// fingerprint -> {
// cascadeId, sessionId, lsPort, apiKey,
// stepOffset, generatorOffset,
// createdAt, lastAccess
// }
const _pool = new Map();

const stats = { hits: 0, misses: 0, stores: 0, evictions: 0, expired: 0 };
Expand Down Expand Up @@ -181,6 +185,8 @@ export function checkin(fingerprint, entry) {
sessionId: entry.sessionId,
lsPort: entry.lsPort,
apiKey: entry.apiKey,
stepOffset: Number.isFinite(entry.stepOffset) ? entry.stepOffset : 0,
generatorOffset: Number.isFinite(entry.generatorOffset) ? entry.generatorOffset : 0,
createdAt: entry.createdAt || now,
lastAccess: now,
});
Expand Down
11 changes: 10 additions & 1 deletion src/handlers/chat.js
Original file line number Diff line number Diff line change
Expand Up @@ -498,7 +498,12 @@ async function nonStreamResponse(client, id, created, model, modelKey, messages,
if (c.text) allText += c.text;
if (c.thinking) allThinking += c.thinking;
}
cascadeMeta = { cascadeId: chunks.cascadeId, sessionId: chunks.sessionId };
cascadeMeta = {
cascadeId: chunks.cascadeId,
sessionId: chunks.sessionId,
stepOffset: chunks.stepOffset,
generatorOffset: chunks.generatorOffset,
};
serverUsage = chunks.usage || null;
// Always strip <tool_call>/<tool_result> blocks from Cascade text.
// - emulateTools=true: parsed tool_calls become OpenAI-format tool_calls.
Expand Down Expand Up @@ -540,6 +545,8 @@ async function nonStreamResponse(client, id, created, model, modelKey, messages,
sessionId: cascadeMeta.sessionId,
lsPort: poolCtx.lsPort,
apiKey: poolCtx.apiKey,
stepOffset: Number.isFinite(cascadeMeta.stepOffset) ? cascadeMeta.stepOffset : poolCtx.reuseEntry?.stepOffset,
generatorOffset: Number.isFinite(cascadeMeta.generatorOffset) ? cascadeMeta.generatorOffset : poolCtx.reuseEntry?.generatorOffset,
createdAt: poolCtx.reuseEntry?.createdAt,
});
}
Expand Down Expand Up @@ -916,6 +923,8 @@ function streamResponse(id, created, model, modelKey, messages, cascadeMessages,
sessionId: cascadeResult.sessionId,
lsPort: ls.port,
apiKey: currentApiKey,
stepOffset: Number.isFinite(cascadeResult.stepOffset) ? cascadeResult.stepOffset : reuseEntry?.stepOffset,
generatorOffset: Number.isFinite(cascadeResult.generatorOffset) ? cascadeResult.generatorOffset : reuseEntry?.generatorOffset,
createdAt: reuseEntry?.createdAt,
});
}
Expand Down
16 changes: 13 additions & 3 deletions src/windsurf.js
Original file line number Diff line number Diff line change
Expand Up @@ -573,8 +573,12 @@ export function buildGetGeneratorMetadataRequest(cascadeId, offset = 0) {
* }
*
* Returns null if nothing reported; otherwise an aggregated
* {inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens} summed
* across every generator invocation (multi-model trajectories sum).
* {inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens, entryCount}
* summed across every generator invocation (multi-model trajectories sum).
*
* `entryCount` is the number of generator-metadata records returned by this
* response. On resumed cascades we use it as the next offset so prior-turn
* usage is not counted again.
*/
export function parseGeneratorMetadata(buf) {
const fields = parseFields(buf);
Expand Down Expand Up @@ -609,7 +613,13 @@ export function parseGeneratorMetadata(buf) {
}
}
if (!found) return null;
return { inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens };
return {
inputTokens,
outputTokens,
cacheReadTokens,
cacheWriteTokens,
entryCount: metaEntries.length,
};
}

// ─── Cascade response parsers ──────────────────────────────
Expand Down
Loading