Skip to content

Commit fa2b2ff

Browse files
authored
fix(channels): recover failed progress draft starts (#88749)
1 parent a6f4de4 commit fa2b2ff

7 files changed

Lines changed: 147 additions & 22 deletions

File tree

extensions/discord/src/monitor/message-handler.draft-preview.ts

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -256,12 +256,14 @@ export function createDiscordDraftPreviewController(params: {
256256
);
257257
}
258258
const alreadyStarted = progressDraftGate.hasStarted;
259+
let progressActive = false;
259260
if (shouldStartDiscordProgressDraftNow(line)) {
260261
await progressDraftGate.startNow();
262+
progressActive = progressDraftGate.hasStarted;
261263
} else {
262-
await progressDraftGate.noteWork();
264+
progressActive = await progressDraftGate.noteWork();
263265
}
264-
if (alreadyStarted && progressDraftGate.hasStarted) {
266+
if ((alreadyStarted || progressActive) && progressDraftGate.hasStarted) {
265267
await renderProgressDraft();
266268
}
267269
},
@@ -294,9 +296,8 @@ export function createDiscordDraftPreviewController(params: {
294296
}
295297
lastReasoningProgressLine = normalized;
296298
}
297-
const alreadyStarted = progressDraftGate.hasStarted;
298-
await progressDraftGate.noteWork();
299-
if (alreadyStarted && progressDraftGate.hasStarted) {
299+
const progressActive = await progressDraftGate.noteWork();
300+
if (progressActive && progressDraftGate.hasStarted) {
300301
await renderProgressDraft();
301302
}
302303
},

extensions/matrix/src/matrix/monitor/handler.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1650,8 +1650,8 @@ export function createMatrixRoomMessageHandler(params: MatrixMonitorHandlerParam
16501650
);
16511651
}
16521652
const alreadyStarted = progressDraftGate.hasStarted;
1653-
await progressDraftGate.noteWork();
1654-
if (alreadyStarted && progressDraftGate.hasStarted) {
1653+
const progressActive = await progressDraftGate.noteWork();
1654+
if ((alreadyStarted || progressActive) && progressDraftGate.hasStarted) {
16551655
renderProgressDraft();
16561656
}
16571657
};

extensions/msteams/src/reply-stream-controller.ts

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -215,10 +215,10 @@ export function createTeamsReplyStreamController(params: {
215215
return;
216216
}
217217
const hadStarted = progressDraftGate.hasStarted;
218-
await progressDraftGate.noteWork();
218+
const progressActive = await progressDraftGate.noteWork();
219219
// If the gate was already started, the call above is a no-op — refresh
220220
// the informative line manually so the latest progress lines render.
221-
if (hadStarted && progressDraftGate.hasStarted) {
221+
if ((hadStarted || progressActive) && progressDraftGate.hasStarted) {
222222
renderInformativeUpdate();
223223
}
224224
},
@@ -252,8 +252,8 @@ export function createTeamsReplyStreamController(params: {
252252
}
253253
}
254254
const hadStarted = progressDraftGate.hasStarted;
255-
await progressDraftGate.noteWork();
256-
if (hadStarted && progressDraftGate.hasStarted) {
255+
const progressActive = await progressDraftGate.noteWork();
256+
if ((hadStarted || progressActive) && progressDraftGate.hasStarted) {
257257
renderInformativeUpdate();
258258
}
259259
},

extensions/slack/src/monitor/message-handler/dispatch.ts

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1492,8 +1492,8 @@ export async function dispatchPreparedSlackMessage(prepared: PreparedSlackMessag
14921492
return;
14931493
}
14941494
const alreadyStarted = progressDraftGate.hasStarted;
1495-
await progressDraftGate.noteWork();
1496-
if (alreadyStarted && progressDraftGate.hasStarted) {
1495+
const progressActive = await progressDraftGate.noteWork();
1496+
if ((alreadyStarted || progressActive) && progressDraftGate.hasStarted) {
14971497
await refreshStartedProgressDraft();
14981498
}
14991499
return;
@@ -1533,12 +1533,15 @@ export async function dispatchPreparedSlackMessage(prepared: PreparedSlackMessag
15331533
await updateNativeProgressStream();
15341534
} else {
15351535
await progressDraftGate.startNow();
1536+
if (progressDraftGate.hasStarted) {
1537+
await updateNativeProgressStream();
1538+
}
15361539
}
15371540
return;
15381541
}
15391542
const alreadyStarted = progressDraftGate.hasStarted;
1540-
await progressDraftGate.noteWork();
1541-
if (alreadyStarted && progressDraftGate.hasStarted) {
1543+
const progressActive = await progressDraftGate.noteWork();
1544+
if ((alreadyStarted || progressActive) && progressDraftGate.hasStarted) {
15421545
await refreshStartedProgressDraft();
15431546
}
15441547
};

extensions/telegram/src/bot-message-dispatch.ts

Lines changed: 3 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1046,17 +1046,16 @@ export const dispatchTelegramMessage = async ({
10461046
}
10471047
streamToolProgressLines = nextLines;
10481048
if (options?.startImmediately) {
1049-
const alreadyStarted = progressDraftGate.hasStarted;
10501049
await progressDraftGate.startNow();
1051-
if (alreadyStarted && progressDraftGate.hasStarted) {
1050+
if (progressDraftGate.hasStarted) {
10521051
await renderProgressDraft();
10531052
return true;
10541053
}
10551054
return progressDraftGate.hasStarted;
10561055
}
10571056
const alreadyStarted = progressDraftGate.hasStarted;
1058-
await progressDraftGate.noteWork();
1059-
if (alreadyStarted && progressDraftGate.hasStarted) {
1057+
const progressActive = await progressDraftGate.noteWork();
1058+
if ((alreadyStarted || progressActive) && progressDraftGate.hasStarted) {
10601059
await renderProgressDraft();
10611060
return true;
10621061
}

src/channels/streaming.ts

Lines changed: 28 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -539,9 +539,29 @@ export function createChannelProgressDraftGate(params: {
539539
if (disposed || started) {
540540
return startPromise ?? Promise.resolve();
541541
}
542-
started = true;
542+
if (startPromise) {
543+
return startPromise;
544+
}
543545
clearTimer();
544-
startPromise = Promise.resolve().then(params.onStart);
546+
started = true;
547+
const nextStart = Promise.resolve()
548+
.then(params.onStart)
549+
.then(() => {
550+
if (disposed) {
551+
started = false;
552+
}
553+
if (startPromise === nextStart) {
554+
startPromise = undefined;
555+
}
556+
})
557+
.catch((error: unknown) => {
558+
if (startPromise === nextStart) {
559+
startPromise = undefined;
560+
}
561+
started = false;
562+
throw error;
563+
});
564+
startPromise = nextStart;
545565
return startPromise;
546566
};
547567

@@ -567,12 +587,16 @@ export function createChannelProgressDraftGate(params: {
567587
return false;
568588
}
569589
workEvents += 1;
590+
if (startPromise) {
591+
await startPromise;
592+
return started;
593+
}
570594
if (started) {
571595
return true;
572596
}
573597
if (workEvents > 1) {
574598
await start();
575-
return true;
599+
return started;
576600
}
577601
schedule();
578602
return false;
@@ -582,6 +606,7 @@ export function createChannelProgressDraftGate(params: {
582606
},
583607
cancel(): void {
584608
disposed = true;
609+
started = false;
585610
clearTimer();
586611
},
587612
};

src/plugin-sdk/channel-streaming.test.ts

Lines changed: 97 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -590,6 +590,103 @@ describe("channel-streaming", () => {
590590
expect(onStart).toHaveBeenCalledTimes(1);
591591
});
592592

593+
it("does not report started when delayed progress startup rejects", async () => {
594+
vi.useFakeTimers();
595+
const onStart = vi
596+
.fn<() => Promise<void>>()
597+
.mockRejectedValueOnce(new Error("draft unavailable"))
598+
.mockResolvedValueOnce(undefined);
599+
const gate = createChannelProgressDraftGate({ onStart });
600+
601+
await expect(gate.noteWork()).resolves.toBe(false);
602+
await vi.advanceTimersByTimeAsync(5_000);
603+
604+
expect(onStart).toHaveBeenCalledTimes(1);
605+
expect(gate.hasStarted).toBe(false);
606+
607+
await expect(gate.noteWork()).resolves.toBe(true);
608+
609+
expect(onStart).toHaveBeenCalledTimes(2);
610+
expect(gate.hasStarted).toBe(true);
611+
});
612+
613+
it("keeps concurrent progress startup single-flight until onStart resolves", async () => {
614+
vi.useFakeTimers();
615+
let resolveStart: (() => void) | undefined;
616+
const onStart = vi.fn(
617+
() =>
618+
new Promise<void>((resolve) => {
619+
resolveStart = resolve;
620+
}),
621+
);
622+
const gate = createChannelProgressDraftGate({ onStart });
623+
624+
await gate.noteWork();
625+
const firstStart = gate.noteWork();
626+
const secondStart = gate.startNow();
627+
await Promise.resolve();
628+
629+
expect(onStart).toHaveBeenCalledTimes(1);
630+
expect(gate.hasStarted).toBe(true);
631+
632+
resolveStart?.();
633+
await expect(firstStart).resolves.toBe(true);
634+
await expect(secondStart).resolves.toBeUndefined();
635+
636+
expect(onStart).toHaveBeenCalledTimes(1);
637+
expect(gate.hasStarted).toBe(true);
638+
});
639+
640+
it("does not report active when cancel wins the startup race", async () => {
641+
vi.useFakeTimers();
642+
let resolveStart: (() => void) | undefined;
643+
const onStart = vi.fn(
644+
() =>
645+
new Promise<void>((resolve) => {
646+
resolveStart = resolve;
647+
}),
648+
);
649+
const gate = createChannelProgressDraftGate({ onStart });
650+
651+
await gate.noteWork();
652+
const startResult = gate.noteWork();
653+
await Promise.resolve();
654+
655+
expect(onStart).toHaveBeenCalledTimes(1);
656+
gate.cancel();
657+
658+
resolveStart?.();
659+
660+
await expect(startResult).resolves.toBe(false);
661+
expect(gate.hasStarted).toBe(false);
662+
});
663+
664+
it("joins explicit startup before applying the first-work delay", async () => {
665+
vi.useFakeTimers();
666+
let resolveStart: (() => void) | undefined;
667+
const onStart = vi.fn(
668+
() =>
669+
new Promise<void>((resolve) => {
670+
resolveStart = resolve;
671+
}),
672+
);
673+
const gate = createChannelProgressDraftGate({ onStart });
674+
675+
const explicitStart = gate.startNow();
676+
await Promise.resolve();
677+
const workDuringStart = gate.noteWork();
678+
679+
expect(onStart).toHaveBeenCalledTimes(1);
680+
expect(gate.hasStarted).toBe(true);
681+
682+
resolveStart?.();
683+
684+
await expect(explicitStart).resolves.toBeUndefined();
685+
await expect(workDuringStart).resolves.toBe(true);
686+
expect(onStart).toHaveBeenCalledTimes(1);
687+
expect(gate.hasStarted).toBe(true);
688+
});
689+
593690
it("ignores message-like tools for progress draft work", () => {
594691
expect(isChannelProgressDraftWorkToolName("message")).toBe(false);
595692
expect(isChannelProgressDraftWorkToolName("react")).toBe(false);

0 commit comments

Comments
 (0)