mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
refactor(channels): finish shared draft-stream ownership (#162174)
Move draft generation, atomic receipt publication, sticky stale-create retirement, and current-message reset/retirement into the existing shared lifecycle. Discord uses the generation operations; Matrix reuses message reset while keeping its live-marker policy.
Preserve existing SDK parameters and operations; the API report confirms only additive members on the existing helper. Validation: 273 draft/entry-point tests, 1153 plugin contract tests, changed-file checks, both cycle checks at zero, and independent review through P2.
Hosted CI hit the pre-existing chat.abort-errors fixture deadlock also present in main run 36780090668. Main already contains the owning fixture fix, 79f0e9bf12. The PR records this inherited failure and its incidental SDK barrel import; security checks passed.
This commit is contained in:
parent
5ad2c6b4b4
commit
5b9e29c20c
5 changed files with 113 additions and 133 deletions
|
|
@ -127,6 +127,17 @@ When transport cleanup policy must change, pass a synchronous `prepareCleanup`
|
|||
callback to `cleanupPending`; it runs in order with clears, before deletion.
|
||||
Rejected deletions remain owned for a later cleanup attempt.
|
||||
|
||||
For synchronous turn rotation, `reset()` advances `generation`, reopens delivery,
|
||||
and resets the current message and pending updates. Published messages remain the
|
||||
adapter's responsibility. `reset("discard")` also retires
|
||||
creates from earlier generations when they settle. Call `createMessage(send, publish)`
|
||||
inside the serialized send loop; its synchronous `publish` callback installs only
|
||||
current-generation receipts. Capture `generation` before edits and recheck it before
|
||||
publishing their results. `retireCurrent(stopForClear)` similarly fences awaited
|
||||
cleanup. Adapters that already settle their sends before rotation can use
|
||||
`resetMessage()` to reset identity and pending/throttle state without reopening delivery
|
||||
or advancing the generation.
|
||||
|
||||
### Commentary delivery ownership
|
||||
|
||||
Set `commentaryPayloadsEnabled: true` when the channel supports durable commentary messages.
|
||||
|
|
|
|||
|
|
@ -343,57 +343,38 @@ describe("createDiscordDraftStream", () => {
|
|||
expect(stream.messageId()).toBe("1002");
|
||||
});
|
||||
|
||||
it("preserves an in-flight block while starting the next block", async () => {
|
||||
let finishFirstCreate: ((value: { id: string }) => void) | undefined;
|
||||
const firstCreate = new Promise<{ id: string }>((resolve) => {
|
||||
finishFirstCreate = resolve;
|
||||
});
|
||||
it.each([
|
||||
{ modes: ["preserve"] as const, discarded: false },
|
||||
{ modes: ["discard"] as const, discarded: true },
|
||||
{ modes: ["discard", "preserve"] as const, discarded: true },
|
||||
{ modes: ["preserve", "discard"] as const, discarded: true },
|
||||
])("settles an in-flight create after $modes rotations", async ({ modes, discarded }) => {
|
||||
const firstCreate = createDeferred<{ id: string }>();
|
||||
const rest = {
|
||||
post: vi.fn().mockReturnValueOnce(firstCreate).mockResolvedValueOnce({ id: "1002" }),
|
||||
post: vi.fn().mockReturnValueOnce(firstCreate.promise).mockResolvedValueOnce({ id: "1002" }),
|
||||
patch: vi.fn(async () => undefined),
|
||||
delete: vi.fn(async () => undefined),
|
||||
};
|
||||
const stream = createDraftStream(rest);
|
||||
|
||||
stream.update("old turn draft");
|
||||
await vi.waitFor(() => expect(rest.post).toHaveBeenCalledTimes(1));
|
||||
stream.forceNewMessage();
|
||||
expect(rest.post).toHaveBeenCalledTimes(1);
|
||||
for (const mode of modes) {
|
||||
stream.forceNewMessage(mode);
|
||||
}
|
||||
stream.update("queued turn draft");
|
||||
finishFirstCreate?.({ id: "1001" });
|
||||
firstCreate.resolve({ id: "1001" });
|
||||
await stream.flush();
|
||||
|
||||
expect(rest.post).toHaveBeenCalledTimes(2);
|
||||
expect(rest.post.mock.calls[1]?.[1]).toMatchObject({
|
||||
body: { content: "queued turn draft" },
|
||||
});
|
||||
expect(rest.delete).not.toHaveBeenCalled();
|
||||
expect(stream.messageId()).toBe("1002");
|
||||
});
|
||||
|
||||
it("discards an in-flight progress draft while starting the queued turn", async () => {
|
||||
let finishFirstCreate: ((value: { id: string }) => void) | undefined;
|
||||
const firstCreate = new Promise<{ id: string }>((resolve) => {
|
||||
finishFirstCreate = resolve;
|
||||
});
|
||||
const rest = {
|
||||
post: vi.fn().mockReturnValueOnce(firstCreate).mockResolvedValueOnce({ id: "1002" }),
|
||||
patch: vi.fn(async () => undefined),
|
||||
delete: vi.fn(async () => undefined),
|
||||
};
|
||||
const stream = createDraftStream(rest);
|
||||
|
||||
stream.update("old progress draft");
|
||||
await vi.waitFor(() => expect(rest.post).toHaveBeenCalledTimes(1));
|
||||
stream.forceNewMessage("discard");
|
||||
stream.update("queued turn draft");
|
||||
finishFirstCreate?.({ id: "1001" });
|
||||
await stream.flush();
|
||||
|
||||
expect(rest.post).toHaveBeenCalledTimes(2);
|
||||
expect(rest.post.mock.calls[1]?.[1]).toMatchObject({
|
||||
body: { content: "queued turn draft" },
|
||||
});
|
||||
expect(rest.delete).toHaveBeenCalledWith(Routes.channelMessage("c1", "1001"));
|
||||
if (discarded) {
|
||||
expect(rest.delete).toHaveBeenCalledExactlyOnceWith(Routes.channelMessage("c1", "1001"));
|
||||
} else {
|
||||
expect(rest.delete).not.toHaveBeenCalled();
|
||||
}
|
||||
expect(stream.messageId()).toBe("1002");
|
||||
});
|
||||
|
||||
|
|
|
|||
|
|
@ -42,15 +42,12 @@ export function createDiscordDraftStream(params: {
|
|||
const streamState = { stopped: false, final: false };
|
||||
let streamMessage: DiscordDraftMessage | undefined;
|
||||
let lastSentText = "";
|
||||
let streamGeneration = 0;
|
||||
let activeCreateGeneration: number | undefined;
|
||||
let discardActiveCreate = false;
|
||||
|
||||
const sendOrEditStreamMessage = async ({
|
||||
text,
|
||||
complete,
|
||||
}: DiscordDraftUpdate): Promise<boolean> => {
|
||||
const generation = streamGeneration;
|
||||
const generation = lifecycle.generation;
|
||||
const targetChannelId = channelId;
|
||||
// Allow final flush even if stopped (e.g., after clear()).
|
||||
if (streamState.stopped && !streamState.final) {
|
||||
|
|
@ -88,7 +85,7 @@ export function createDiscordDraftStream(params: {
|
|||
await editChannelMessage(rest, streamMessage.channelId, streamMessage.messageId, {
|
||||
body,
|
||||
});
|
||||
if (generation === streamGeneration) {
|
||||
if (generation === lifecycle.generation) {
|
||||
lastSentText = trimmed;
|
||||
}
|
||||
return true;
|
||||
|
|
@ -97,37 +94,31 @@ export function createDiscordDraftStream(params: {
|
|||
const messageReference = replyToMessageId
|
||||
? { message_id: replyToMessageId, fail_if_not_exists: false }
|
||||
: undefined;
|
||||
activeCreateGeneration = generation;
|
||||
const sent = (await rest.post(Routes.channelMessages(targetChannelId), {
|
||||
body: {
|
||||
...body,
|
||||
...(messageReference ? { message_reference: messageReference } : {}),
|
||||
return await lifecycle.createMessage(
|
||||
async () => {
|
||||
const sent = (await rest.post(Routes.channelMessages(targetChannelId), {
|
||||
body: {
|
||||
...body,
|
||||
...(messageReference ? { message_reference: messageReference } : {}),
|
||||
},
|
||||
})) as { id?: string }; // SAFETY: The create response's ID is checked before use.
|
||||
return typeof sent?.id === "string" && sent.id
|
||||
? { channelId: targetChannelId, messageId: sent.id }
|
||||
: undefined;
|
||||
},
|
||||
})) as { id?: string }; // SAFETY: The create response's ID is checked before use.
|
||||
const sentMessageId = sent?.id;
|
||||
const shouldDiscardStaleCreate = activeCreateGeneration === generation && discardActiveCreate;
|
||||
activeCreateGeneration = undefined;
|
||||
discardActiveCreate = false;
|
||||
if (generation !== streamGeneration) {
|
||||
if (shouldDiscardStaleCreate && typeof sentMessageId === "string" && sentMessageId) {
|
||||
await lifecycle.retire({ channelId: targetChannelId, messageId: sentMessageId });
|
||||
}
|
||||
return true;
|
||||
}
|
||||
if (typeof sentMessageId !== "string" || !sentMessageId) {
|
||||
streamState.stopped = true;
|
||||
params.warn?.("discord stream preview stopped (missing message id from send)");
|
||||
return false;
|
||||
}
|
||||
streamMessage = { channelId: targetChannelId, messageId: sentMessageId };
|
||||
lastSentText = trimmed;
|
||||
return true;
|
||||
(message) => {
|
||||
if (!message) {
|
||||
streamState.stopped = true;
|
||||
params.warn?.("discord stream preview stopped (missing message id from send)");
|
||||
return false;
|
||||
}
|
||||
streamMessage = message;
|
||||
lastSentText = trimmed;
|
||||
return true;
|
||||
},
|
||||
);
|
||||
} catch (err) {
|
||||
if (activeCreateGeneration === generation) {
|
||||
activeCreateGeneration = undefined;
|
||||
discardActiveCreate = false;
|
||||
}
|
||||
if (generation !== streamGeneration) {
|
||||
if (generation !== lifecycle.generation) {
|
||||
return true;
|
||||
}
|
||||
streamState.stopped = true;
|
||||
|
|
@ -159,22 +150,6 @@ export function createDiscordDraftStream(params: {
|
|||
const update = (text: string, options?: { complete?: boolean }) =>
|
||||
updateDraft({ text, complete: options?.complete === true });
|
||||
|
||||
/** Reset the draft so the next update creates a new message. */
|
||||
const forceNewMessage = (mode: "preserve" | "discard" = "preserve") => {
|
||||
// In-flight REST calls may finish after a turn boundary. Advance identity
|
||||
// synchronously so their result cannot overwrite the next turn's state.
|
||||
// Block mode preserves the prior block; progress mode discards its draft.
|
||||
if (mode === "discard" && activeCreateGeneration !== undefined) {
|
||||
discardActiveCreate = true;
|
||||
}
|
||||
streamGeneration += 1;
|
||||
streamState.stopped = false;
|
||||
streamState.final = false;
|
||||
streamMessage = undefined;
|
||||
lastSentText = "";
|
||||
loop.resetPending();
|
||||
loop.resetThrottleWindow();
|
||||
};
|
||||
/** Move the draft to another channel, preserving its current text. */
|
||||
const retarget = async (nextChannelId: string) => {
|
||||
const normalized = nextChannelId.trim();
|
||||
|
|
@ -185,13 +160,8 @@ export function createDiscordDraftStream(params: {
|
|||
const pending = loop.takePending();
|
||||
const previousMessage = streamMessage;
|
||||
const previousText = pending.text || lastSentText;
|
||||
streamGeneration += 1;
|
||||
channelId = normalized;
|
||||
streamMessage = undefined;
|
||||
lastSentText = "";
|
||||
streamState.stopped = false;
|
||||
streamState.final = false;
|
||||
loop.resetThrottleWindow();
|
||||
lifecycle.reset();
|
||||
if (previousText) {
|
||||
update(previousText, { complete: pending.text ? pending.complete : true });
|
||||
await loop.flush();
|
||||
|
|
@ -204,18 +174,6 @@ export function createDiscordDraftStream(params: {
|
|||
await lifecycle.retire(previousMessage);
|
||||
}
|
||||
};
|
||||
const retireCurrentMessage = async (stopForClear: () => Promise<void>) => {
|
||||
const generation = streamGeneration;
|
||||
await stopForClear();
|
||||
if (generation !== streamGeneration) {
|
||||
return;
|
||||
}
|
||||
const message = streamMessage;
|
||||
clearMessageId();
|
||||
if (message) {
|
||||
await lifecycle.retire(message);
|
||||
}
|
||||
};
|
||||
|
||||
params.log?.(`discord stream preview ready (maxChars=${maxChars}, throttleMs=${throttleMs})`);
|
||||
|
||||
|
|
@ -224,9 +182,9 @@ export function createDiscordDraftStream(params: {
|
|||
flush: loop.flush,
|
||||
messageId: () => streamMessage?.messageId,
|
||||
lastDeliveredText: () => lastSentText,
|
||||
clear: () => retireCurrentMessage(discardPending),
|
||||
clear: () => lifecycle.retireCurrent(discardPending),
|
||||
deleteCurrentMessage: () =>
|
||||
retireCurrentMessage(async () => {
|
||||
lifecycle.retireCurrent(async () => {
|
||||
loop.resetPending();
|
||||
await loop.waitForInFlight();
|
||||
}),
|
||||
|
|
@ -234,9 +192,7 @@ export function createDiscordDraftStream(params: {
|
|||
seal,
|
||||
stop,
|
||||
retarget,
|
||||
cleanupPendingMessages: async () => {
|
||||
await lifecycle.cleanupPending();
|
||||
},
|
||||
forceNewMessage,
|
||||
cleanupPendingMessages: lifecycle.cleanupPending,
|
||||
forceNewMessage: lifecycle.reset,
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -114,6 +114,7 @@ export function createMatrixDraftStream(params: {
|
|||
clear,
|
||||
retire,
|
||||
cleanupPending,
|
||||
resetMessage,
|
||||
} = createFinalizableDraftLifecycle({
|
||||
throttleMs: DEFAULT_THROTTLE_MS,
|
||||
state: streamState,
|
||||
|
|
@ -166,14 +167,10 @@ export function createMatrixDraftStream(params: {
|
|||
};
|
||||
|
||||
const resetCurrentMessage = (): void => {
|
||||
currentEventId = undefined;
|
||||
lastSentText = "";
|
||||
lastSentContent = "";
|
||||
resetMessage();
|
||||
sendFailed = false;
|
||||
finalizeInPlaceBlocked = false;
|
||||
liveFinalized = false;
|
||||
loop.resetPending();
|
||||
loop.resetThrottleWindow();
|
||||
};
|
||||
const reset = (): void => {
|
||||
// A new block consumes the first-only reply reference; retraction does not.
|
||||
|
|
|
|||
|
|
@ -1,9 +1,6 @@
|
|||
import { formatErrorMessage } from "../infra/errors.js";
|
||||
import { createDraftStreamLoop } from "./draft-stream-loop.js";
|
||||
|
||||
/**
|
||||
* Mutable finalization flags shared by draft stream controls and channel adapters.
|
||||
*/
|
||||
export type FinalizableDraftStreamState = {
|
||||
stopped: boolean;
|
||||
final: boolean;
|
||||
|
|
@ -101,17 +98,13 @@ export function createFinalizableDraftStreamControls<T = string>(
|
|||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Creates finalizable draft controls backed by a shared mutable state object.
|
||||
*/
|
||||
export function createFinalizableDraftStreamControlsForState<T = string>(
|
||||
params: DraftStreamOptions<T> & {
|
||||
state: FinalizableDraftStreamState;
|
||||
},
|
||||
) {
|
||||
return createFinalizableDraftStreamControls<T>({
|
||||
throttleMs: params.throttleMs,
|
||||
coalesceInFlight: params.coalesceInFlight,
|
||||
...params,
|
||||
isStopped: () => params.state.stopped,
|
||||
isFinal: () => params.state.final,
|
||||
markStopped: () => {
|
||||
|
|
@ -120,15 +113,9 @@ export function createFinalizableDraftStreamControlsForState<T = string>(
|
|||
markFinal: () => {
|
||||
params.state.final = true;
|
||||
},
|
||||
sendOrEditStreamMessage: params.sendOrEditStreamMessage,
|
||||
...(params.emptyValue !== undefined ? { emptyValue: params.emptyValue } : {}),
|
||||
...(params.isEmpty ? { isEmpty: params.isEmpty } : {}),
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Stops a draft stream, reads the current preview message id, then clears the stored id.
|
||||
*/
|
||||
export async function takeMessageIdAfterStop<T>(
|
||||
params: StopAndClearMessageIdParams<T>,
|
||||
): Promise<T | undefined> {
|
||||
|
|
@ -180,9 +167,6 @@ export async function clearFinalizableDraftMessage<T>(
|
|||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Builds the standard draft lifecycle used by channel streaming preview implementations.
|
||||
*/
|
||||
export function createFinalizableDraftLifecycle<TMessageId, TUpdate = string>(
|
||||
params: FinalizableDraftLifecycleParams<TMessageId, TUpdate>,
|
||||
) {
|
||||
|
|
@ -193,6 +177,8 @@ export function createFinalizableDraftLifecycle<TMessageId, TUpdate = string>(
|
|||
};
|
||||
const pending = new Map<TMessageId, Retirement>();
|
||||
let clearTail = Promise.resolve();
|
||||
let generation = 0;
|
||||
let discardThroughGeneration = -1;
|
||||
|
||||
const claim = (
|
||||
messageId: TMessageId,
|
||||
|
|
@ -275,8 +261,57 @@ export function createFinalizableDraftLifecycle<TMessageId, TUpdate = string>(
|
|||
clearTail = stopRun;
|
||||
return stopRun;
|
||||
};
|
||||
const resetMessage = () => {
|
||||
params.clearMessageId();
|
||||
controls.loop.resetPending();
|
||||
controls.loop.resetThrottleWindow();
|
||||
};
|
||||
const reset = (mode: "preserve" | "discard" = "preserve") => {
|
||||
// A later rotation cannot revoke an earlier request to discard an in-flight create.
|
||||
if (mode === "discard") {
|
||||
discardThroughGeneration = generation;
|
||||
}
|
||||
generation += 1;
|
||||
params.state.stopped = false;
|
||||
params.state.final = false;
|
||||
resetMessage();
|
||||
};
|
||||
const createMessage = async (
|
||||
send: () => Promise<TMessageId | undefined>,
|
||||
publish: (messageId: TMessageId | undefined) => boolean,
|
||||
): Promise<boolean> => {
|
||||
const startedGeneration = generation;
|
||||
const messageId = await send();
|
||||
if (startedGeneration === generation) {
|
||||
// Publish in the same turn as the generation check so rotation cannot lose the receipt.
|
||||
return publish(messageId);
|
||||
}
|
||||
if (startedGeneration <= discardThroughGeneration && params.isValidMessageId(messageId)) {
|
||||
await retire(messageId);
|
||||
}
|
||||
return true;
|
||||
};
|
||||
const retireCurrent = async (stopForClear: () => Promise<void>): Promise<void> => {
|
||||
const startedGeneration = generation;
|
||||
await stopForClear();
|
||||
if (startedGeneration !== generation) {
|
||||
return;
|
||||
}
|
||||
const messageId = params.readMessageId();
|
||||
params.clearMessageId();
|
||||
if (params.isValidMessageId(messageId)) {
|
||||
await retire(messageId);
|
||||
}
|
||||
};
|
||||
return {
|
||||
...controls,
|
||||
get generation() {
|
||||
return generation;
|
||||
},
|
||||
reset,
|
||||
resetMessage,
|
||||
createMessage,
|
||||
retireCurrent,
|
||||
stop,
|
||||
clear: () => clearWithStop(controls.stopForClear),
|
||||
clearWithStop,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue