Skip to content

Commit af8d1ad

Browse files
authored
Merge pull request #45 from baily-zhang/fix-cascade-reuse-offsets
fix: prevent cascade reuse from replaying old context
2 parents 20dea9e + 8b82045 commit af8d1ad

4 files changed

Lines changed: 75 additions & 9 deletions

File tree

src/client.js

Lines changed: 45 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -266,7 +266,7 @@ export class WindsurfClient {
266266
* @param {object} opts - { onChunk, onEnd, onError }
267267
*/
268268
async cascadeChat(messages, modelEnum, modelUid, opts = {}) {
269-
const { onChunk, onEnd, onError, signal, reuseEntry, toolPreamble } = opts;
269+
let { onChunk, onEnd, onError, signal, reuseEntry, toolPreamble } = opts;
270270
const aborted = () => signal?.aborted;
271271
const inputChars = messages.reduce((n, m) => n + contentToString(m?.content).length, 0);
272272

@@ -311,6 +311,42 @@ export class WindsurfClient {
311311
cascadeId = await openCascade();
312312
}
313313

314+
// A resumed cascade already contains every prior turn in its trajectory.
315+
// If we poll from step offset 0 again, the old planner-response steps are
316+
// replayed as fresh output and both text and usage grow cumulatively
317+
// across turns (`alpha` -> `alphabeta` -> ...). Store absolute offsets in
318+
// the conversation pool and reuse them here; fall back to a one-shot
319+
// snapshot so entries created before this fix still resume safely.
320+
let stepOffset = Number.isInteger(reuseEntry?.stepOffset) && reuseEntry.stepOffset >= 0
321+
? reuseEntry.stepOffset
322+
: 0;
323+
let generatorOffset = Number.isInteger(reuseEntry?.generatorOffset) && reuseEntry.generatorOffset >= 0
324+
? reuseEntry.generatorOffset
325+
: 0;
326+
if (reuseEntry?.cascadeId && (!Number.isInteger(reuseEntry?.stepOffset) || !Number.isInteger(reuseEntry?.generatorOffset))) {
327+
try {
328+
if (!Number.isInteger(reuseEntry?.stepOffset)) {
329+
const resumeStepsResp = await grpcUnary(
330+
this.port, this.csrfToken,
331+
`${LS_SERVICE}/GetCascadeTrajectorySteps`,
332+
grpcFrame(buildGetTrajectoryStepsRequest(cascadeId, 0))
333+
);
334+
stepOffset = parseTrajectorySteps(resumeStepsResp).length;
335+
}
336+
if (!Number.isInteger(reuseEntry?.generatorOffset)) {
337+
const resumeMetaResp = await grpcUnary(
338+
this.port, this.csrfToken,
339+
`${LS_SERVICE}/GetCascadeTrajectoryGeneratorMetadata`,
340+
grpcFrame(buildGetGeneratorMetadataRequest(cascadeId, 0)),
341+
5000
342+
);
343+
generatorOffset = parseGeneratorMetadata(resumeMetaResp)?.entryCount || 0;
344+
}
345+
} catch (e) {
346+
log.warn(`Cascade resume snapshot failed: ${e.message}`);
347+
}
348+
}
349+
314350
let text;
315351
let images = [];
316352
const systemMsgs = messages.filter(m => m.role === 'system');
@@ -443,7 +479,7 @@ export class WindsurfClient {
443479
pollCount++;
444480

445481
// Get steps
446-
const stepsProto = buildGetTrajectoryStepsRequest(cascadeId, 0);
482+
const stepsProto = buildGetTrajectoryStepsRequest(cascadeId, stepOffset);
447483
const stepsResp = await grpcUnary(
448484
this.port, this.csrfToken, `${LS_SERVICE}/GetCascadeTrajectorySteps`, grpcFrame(stepsProto)
449485
);
@@ -615,6 +651,7 @@ export class WindsurfClient {
615651
this.port, this.csrfToken, `${LS_SERVICE}/GetCascadeTrajectorySteps`, grpcFrame(stepsProto)
616652
);
617653
const finalSteps = parseTrajectorySteps(finalResp);
654+
lastStepCount = finalSteps.length;
618655
for (let i = 0; i < finalSteps.length; i++) {
619656
const step = finalSteps[i];
620657
const responseText = step.responseText || '';
@@ -662,7 +699,7 @@ export class WindsurfClient {
662699
polls: pollCount,
663700
textLen: totalYielded,
664701
thinkingLen: totalThinking,
665-
stepCount: Math.max(yieldedByStep.size, thinkingByStep.size, lastStepCount),
702+
stepCount: stepOffset + Math.max(yieldedByStep.size, thinkingByStep.size, lastStepCount),
666703
toolCalls: seenToolCallIds.size,
667704
sawActive,
668705
sawText,
@@ -686,7 +723,7 @@ export class WindsurfClient {
686723
// itself is already formed.
687724
let serverUsage = null;
688725
try {
689-
const metaReq = buildGetGeneratorMetadataRequest(cascadeId, 0);
726+
const metaReq = buildGetGeneratorMetadataRequest(cascadeId, generatorOffset);
690727
const metaResp = await grpcUnary(
691728
this.port, this.csrfToken,
692729
`${LS_SERVICE}/GetCascadeTrajectoryGeneratorMetadata`,
@@ -722,6 +759,10 @@ export class WindsurfClient {
722759
// that iterate over it keep working.
723760
chunks.cascadeId = cascadeId;
724761
chunks.sessionId = sessionId;
762+
chunks.stepOffset = stepOffset + Math.max(yieldedByStep.size, thinkingByStep.size, lastStepCount);
763+
chunks.generatorOffset = serverUsage?.entryCount != null
764+
? generatorOffset + serverUsage.entryCount
765+
: null;
725766
chunks.toolCalls = toolCalls;
726767
chunks.usage = serverUsage;
727768
if (serverUsage) {

src/conversation-pool.js

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,11 @@ function positiveIntEnv(name, fallback) {
3232
const POOL_TTL_MS = positiveIntEnv('CASCADE_POOL_TTL_MS', 30 * 60 * 1000);
3333
const POOL_MAX = 500;
3434

35-
// fingerprint -> { cascadeId, sessionId, lsPort, apiKey, createdAt, lastAccess }
35+
// fingerprint -> {
36+
// cascadeId, sessionId, lsPort, apiKey,
37+
// stepOffset, generatorOffset,
38+
// createdAt, lastAccess
39+
// }
3640
const _pool = new Map();
3741

3842
const stats = { hits: 0, misses: 0, stores: 0, evictions: 0, expired: 0 };
@@ -188,6 +192,8 @@ export function checkin(fingerprint, entry) {
188192
sessionId: entry.sessionId,
189193
lsPort: entry.lsPort,
190194
apiKey: entry.apiKey,
195+
stepOffset: Number.isFinite(entry.stepOffset) ? entry.stepOffset : 0,
196+
generatorOffset: Number.isFinite(entry.generatorOffset) ? entry.generatorOffset : 0,
191197
createdAt: entry.createdAt || now,
192198
lastAccess: now,
193199
});

src/handlers/chat.js

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -501,7 +501,12 @@ async function nonStreamResponse(client, id, created, model, modelKey, messages,
501501
if (c.text) allText += c.text;
502502
if (c.thinking) allThinking += c.thinking;
503503
}
504-
cascadeMeta = { cascadeId: chunks.cascadeId, sessionId: chunks.sessionId };
504+
cascadeMeta = {
505+
cascadeId: chunks.cascadeId,
506+
sessionId: chunks.sessionId,
507+
stepOffset: chunks.stepOffset,
508+
generatorOffset: chunks.generatorOffset,
509+
};
505510
serverUsage = chunks.usage || null;
506511
// Always strip <tool_call>/<tool_result> blocks from Cascade text.
507512
// - emulateTools=true: parsed tool_calls become OpenAI-format tool_calls.
@@ -543,6 +548,8 @@ async function nonStreamResponse(client, id, created, model, modelKey, messages,
543548
sessionId: cascadeMeta.sessionId,
544549
lsPort: poolCtx.lsPort,
545550
apiKey: poolCtx.apiKey,
551+
stepOffset: Number.isFinite(cascadeMeta.stepOffset) ? cascadeMeta.stepOffset : poolCtx.reuseEntry?.stepOffset,
552+
generatorOffset: Number.isFinite(cascadeMeta.generatorOffset) ? cascadeMeta.generatorOffset : poolCtx.reuseEntry?.generatorOffset,
546553
createdAt: poolCtx.reuseEntry?.createdAt,
547554
});
548555
}
@@ -921,6 +928,8 @@ function streamResponse(id, created, model, modelKey, messages, cascadeMessages,
921928
sessionId: cascadeResult.sessionId,
922929
lsPort: ls.port,
923930
apiKey: currentApiKey,
931+
stepOffset: Number.isFinite(cascadeResult.stepOffset) ? cascadeResult.stepOffset : reuseEntry?.stepOffset,
932+
generatorOffset: Number.isFinite(cascadeResult.generatorOffset) ? cascadeResult.generatorOffset : reuseEntry?.generatorOffset,
924933
createdAt: reuseEntry?.createdAt,
925934
});
926935
}

src/windsurf.js

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -573,8 +573,12 @@ export function buildGetGeneratorMetadataRequest(cascadeId, offset = 0) {
573573
* }
574574
*
575575
* Returns null if nothing reported; otherwise an aggregated
576-
* {inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens} summed
577-
* across every generator invocation (multi-model trajectories sum).
576+
* {inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens, entryCount}
577+
* summed across every generator invocation (multi-model trajectories sum).
578+
*
579+
* `entryCount` is the number of generator-metadata records returned by this
580+
* response. On resumed cascades we use it as the next offset so prior-turn
581+
* usage is not counted again.
578582
*/
579583
export function parseGeneratorMetadata(buf) {
580584
const fields = parseFields(buf);
@@ -609,7 +613,13 @@ export function parseGeneratorMetadata(buf) {
609613
}
610614
}
611615
if (!found) return null;
612-
return { inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens };
616+
return {
617+
inputTokens,
618+
outputTokens,
619+
cacheReadTokens,
620+
cacheWriteTokens,
621+
entryCount: metaEntries.length,
622+
};
613623
}
614624

615625
// ─── Cascade response parsers ──────────────────────────────

0 commit comments

Comments
 (0)