refactor(channels): share draft rotation and retirement policies

Share Slack generation and stale-create custody with the existing draft lifecycle while preserving human-reply retention, throttle deadlines, and deferred cleanup. Consolidate Mattermost accepted-delivery failure publication and remove the private Discord chunking wrapper.

Validated with focused and cross-shard tests, plugin contracts, changed checks, both zero-cycle checks, a full build, compatible SDK API diff, independent review, and exact-head hosted CI.
This commit is contained in:
Peter Steinberger 2026-09-30 18:18:22 -07:00 • committed by GitHub
parent 721efcefc0
commit c658913931
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 193 additions and 161 deletions

View file

@ -138,6 +138,13 @@ 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.
Pass `"keep"` as the throttle argument to `resetMessage("keep")` or
`reset("discard", "keep")` when rotation must preserve the existing throttle
window and scheduled flush. The default resets both. Transports that decide
stale-preview disposition during cleanup can pass `{ defer: true }` as the
third argument to `createMessage`; stale discarded receipts then enter deletion
custody without an immediate deletion attempt.
### Commentary delivery ownership
Set `commentaryPayloadsEnabled: true` when the channel supports durable commentary messages.

View file

@ -1,14 +0,0 @@
// Discord tests cover draft chunking plugin behavior.
import { describe, expect, it } from "vitest";
import { resolveDiscordDraftStreamingChunking } from "./draft-chunking.js";
import { EMPTY_DISCORD_TEST_CONFIG } from "./test-support/config.js";
describe("resolveDiscordDraftStreamingChunking", () => {
it("returns sane defaults when discord draft chunking is unset", () => {
expect(resolveDiscordDraftStreamingChunking(EMPTY_DISCORD_TEST_CONFIG)).toEqual({
minChars: 200,
maxChars: 800,
breakPreference: "paragraph",
});
});
});

View file

@ -1,15 +0,0 @@
import {
resolveChannelDraftStreamingChunking,
type ChannelDraftStreamingChunking,
} from "openclaw/plugin-sdk/channel-outbound";
import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts";
import { DISCORD_TEXT_CHUNK_LIMIT } from "./outbound-adapter.js";
export function resolveDiscordDraftStreamingChunking(
cfg: OpenClawConfig,
accountId?: string | null,
): ChannelDraftStreamingChunking {
return resolveChannelDraftStreamingChunking(cfg, "discord", accountId, {
fallbackLimit: DISCORD_TEXT_CHUNK_LIMIT,
});
}

View file

@ -49,10 +49,6 @@ export function createDiscordDraftStream(params: {
}: DiscordDraftUpdate): Promise<boolean> => {
const generation = lifecycle.generation;
const targetChannelId = channelId;
// Allow final flush even if stopped (e.g., after clear()).
if (streamState.stopped && !streamState.final) {
return false;
}
const trimmed = text.trimEnd();
if (!trimmed) {
return false;

View file

@ -3,6 +3,7 @@ import {
type ChannelProgressDraftLine,
createChannelProgressDraftCompositor,
createLivePreviewLifecycle,
resolveChannelDraftStreamingChunking,
resolveChannelStreamingBlockEnabled,
resolveChannelStreamingPreviewCommandText,
resolveChannelStreamingProgressNarration,
@ -18,10 +19,10 @@ import {
stripInlineDirectiveTagsForDelivery,
stripReasoningTagsFromText,
} from "openclaw/plugin-sdk/text-chunking";
import { resolveDiscordDraftStreamingChunking } from "../draft-chunking.js";
import { createDiscordDraftStream } from "../draft-stream.js";
import type { RequestClient } from "../internal/discord.js";
import { withDiscordRequestAuthority } from "../internal/request-authority.js";
import { DISCORD_TEXT_CHUNK_LIMIT } from "../outbound-adapter.js";
import { resolveDiscordPreviewStreamMode } from "../preview-streaming.js";
type DraftReplyReference = {
@ -82,7 +83,9 @@ export function createDiscordDraftPreviewController(params: {
: undefined;
const draftChunking =
draftStream && discordStreamMode === "block"
? resolveDiscordDraftStreamingChunking(params.cfg, params.accountId)
? resolveChannelDraftStreamingChunking(params.cfg, "discord", params.accountId, {
fallbackLimit: DISCORD_TEXT_CHUNK_LIMIT,
})
: undefined;
const shouldSplitPreviewMessages = discordStreamMode === "block";
const draftChunker = draftChunking ? new EmbeddedBlockChunker(draftChunking) : undefined;

View file

@ -29,21 +29,6 @@ type MattermostFinalTextResolution =
publishedParts: readonly MattermostDraftPublishedPart[];
};
type MattermostDraftStream = {
update: (text: string) => void;
updateAssistantText: (text: string) => void;
flush: () => Promise<void>;
postId: () => string | undefined;
clear: () => Promise<void>;
deleteCurrentMessage: () => Promise<void>;
discardPending: () => Promise<void>;
seal: () => Promise<void>;
stop: () => Promise<void>;
forceNewMessage: () => Promise<void>;
settleBoundaries: () => Promise<void>;
resolveFinalText: (text: string) => MattermostFinalTextResolution;
};
function normalizeMattermostDraftText(text: string, maxChars: number): string {
const trimmed = text.trim();
if (!trimmed) {
@ -100,7 +85,7 @@ export function createMattermostDraftStream(params: {
chunkText?: (text: string) => string[];
log?: (message: string) => void;
warn?: (message: string) => void;
}): MattermostDraftStream {
}) {
const maxChars = Math.min(
params.maxChars ?? MATTERMOST_STREAM_MAX_CHARS,
MATTERMOST_STREAM_MAX_CHARS,
@ -108,6 +93,16 @@ export function createMattermostDraftStream(params: {
const throttleMs = Math.max(250, params.throttleMs ?? DEFAULT_THROTTLE_MS);
const streamState = { stopped: false, final: false };
let terminalAcceptedDeliveryError: Error | undefined;
const retainAcceptedDeliveryFailure = (err: unknown) => {
if (isChannelPartialDeliveryError(err)) {
// Publish terminal state before warning hooks can re-enter delivery.
const acceptedDeliveryError = toErrorObject(err, "Mattermost accepted delivery failed");
streamState.stopped = true;
terminalAcceptedDeliveryError = acceptedDeliveryError;
return acceptedDeliveryError;
}
return undefined;
};
const assertNoAcceptedDeliveryFailure = () => {
if (terminalAcceptedDeliveryError !== undefined) {
throw terminalAcceptedDeliveryError;
@ -132,9 +127,6 @@ export function createMattermostDraftStream(params: {
const publishedAssistantParts = new Map<string, MattermostDraftPublishedPart>();
const sendOrEditStreamMessage = async (text: string): Promise<boolean> => {
if (streamState.stopped && !streamState.final) {
return false;
}
const target = currentGeneration;
const rendered = params.renderText?.(text) ?? text;
const normalized = normalizeMattermostDraftText(rendered, maxChars);
@ -168,13 +160,7 @@ export function createMattermostDraftStream(params: {
} catch (err) {
// Stop immediately so a discarded background failure cannot queue a second visible post.
streamState.stopped = true;
const acceptedDeliveryError = isChannelPartialDeliveryError(err)
? toErrorObject(err, "Mattermost accepted delivery failed")
: undefined;
if (acceptedDeliveryError) {
// Warning handlers can synchronously re-enter finalization; retain the failure first.
terminalAcceptedDeliveryError = acceptedDeliveryError;
}
const acceptedDeliveryError = retainAcceptedDeliveryFailure(err);
params.warn?.(
`mattermost stream preview failed: ${err instanceof Error ? err.message : String(err)}`,
);
@ -237,7 +223,6 @@ export function createMattermostDraftStream(params: {
await inFlightAtBoundary;
assertNoAcceptedDeliveryFailure();
if (streamState.stopped && !streamState.final) {
assertNoAcceptedDeliveryFailure();
return;
}
const sourceText = pendingText.trim() ? pendingText : sealed.latestSourceText;
@ -284,14 +269,7 @@ export function createMattermostDraftStream(params: {
sealedAssistantTexts.push({ text: assistantText, requiresBlockBoundary: true });
}
} catch (err) {
const acceptedDeliveryError = isChannelPartialDeliveryError(err)
? toErrorObject(err, "Mattermost accepted delivery failed")
: undefined;
if (acceptedDeliveryError) {
// Publish terminal state before warning hooks can re-enter update or forceNewMessage.
streamState.stopped = true;
terminalAcceptedDeliveryError = acceptedDeliveryError;
}
const acceptedDeliveryError = retainAcceptedDeliveryFailure(err);
const publishedAssistantPrefix = assistantText?.slice(0, publishedAssistantOffset).trim();
if (publishedAssistantPrefix) {
// A later physical chunk failed after this exact source prefix became durable.
@ -369,34 +347,34 @@ export function createMattermostDraftStream(params: {
await currentGeneration.ready;
assertNoAcceptedDeliveryFailure();
};
const resolveFinalText = (text: string) => {
const resolveFinalText = (text: string): MattermostFinalTextResolution => {
const publishedParts = [...publishedAssistantParts.values()];
if (sealedAssistantTexts.length === 0) {
return { kind: "full" as const, text, publishedParts };
return { kind: "full", text, publishedParts };
}
let remainingText = text.trim();
for (const sealedText of sealedAssistantTexts) {
const completed = sealedText.text.trim();
if (!completed || !remainingText.startsWith(completed)) {
return { kind: "full" as const, text, publishedParts };
return { kind: "full", text, publishedParts };
}
const suffix = remainingText.slice(completed.length);
// Canonical assistant block aggregation uses newline separators. A plain-space
// suffix can be a block-local final that merely shares the prior block's prefix.
if (sealedText.requiresBlockBoundary && suffix && !/^\r?\n/.test(suffix)) {
return { kind: "full" as const, text, publishedParts };
return { kind: "full", text, publishedParts };
}
remainingText = suffix.replace(sealedText.requiresBlockBoundary ? /^(?:\r?\n)+/ : /^\s+/, "");
}
const currentText = currentGeneration.latestAssistantText?.trim() ?? "";
const remaining = remainingText.trim();
if (currentText && !remaining.startsWith(currentText)) {
return { kind: "full" as const, text, publishedParts };
return { kind: "full", text, publishedParts };
}
return remaining
? { kind: "remaining" as const, text: remaining, publishedParts }
: { kind: "already-delivered" as const, publishedParts };
? { kind: "remaining", text: remaining, publishedParts }
: { kind: "already-delivered", publishedParts };
};
params.log?.(`mattermost stream preview ready (maxChars=${maxChars}, throttleMs=${throttleMs})`);

View file

@ -199,6 +199,79 @@ describe("createSlackDraftStream", () => {
expect(edit).toHaveBeenCalledTimes(0);
});
it.each(["turn", "human"])("keeps the throttle window after a %s rotation", async (rotation) => {
vi.useFakeTimers();
vi.setSystemTime(10_000);
const accountId = `throttled-${rotation}-rotation`;
const send = vi
.fn<DraftSendFn>()
.mockResolvedValueOnce(slackDraftSendResult("100.100"))
.mockResolvedValueOnce(slackDraftSendResult("100.300"));
const { stream, edit } = createDraftStreamHarness({
accountId,
threadTs: "100.000",
send,
});
try {
stream.update("first preview");
await stream.flush();
stream.update("queued old preview");
await vi.advanceTimersByTimeAsync(100);
if (rotation === "turn") {
stream.forceNewMessage();
} else {
noteSlackDraftConversationMessage({
accountId,
channelId: "C123",
threadTs: "100.000",
messageTs: "100.200",
userId: "U_HUMAN",
});
}
stream.update("replacement preview");
expect(send).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(149);
expect(send).toHaveBeenCalledOnce();
await vi.advanceTimersByTimeAsync(1);
expect(send).toHaveBeenCalledTimes(2);
expect(send).toHaveBeenLastCalledWith(
"channel:C123",
"replacement preview",
expect.any(Object),
);
expect(edit).not.toHaveBeenCalled();
} finally {
await stream.clear();
vi.useRealTimers();
}
});
it("defers a stale create even after cleanup admitted its generation", async () => {
const receipt = createDeferred<ReturnType<typeof slackDraftSendResult>>();
const send = vi.fn<DraftSendFn>().mockReturnValueOnce(receipt.promise);
const { stream, remove } = createDraftStreamHarness({ send });
try {
stream.update("old preview");
expect(send).toHaveBeenCalledOnce();
await stream.dropDetachedMessages();
stream.forceNewMessage();
receipt.resolve(slackDraftSendResult("100.100"));
await stream.flush();
expect(stream.messageId()).toBeUndefined();
expect(remove).not.toHaveBeenCalled();
await stream.dropDetachedMessages();
expect(remove).toHaveBeenCalledExactlyOnceWith("C123", "100.100", {
token: "xoxb-test",
accountId: undefined,
});
} finally {
receipt.resolve(slackDraftSendResult("100.100"));
await stream.clear();
}
});
it("drains past a failed preview and retries only the retained failure", async () => {
const send = vi
.fn<DraftSendFn>()

View file

@ -68,7 +68,6 @@ export function createSlackDraftStream(params: {
const remove = params.remove ?? deleteSlackMessage;
let streamMessage: SlackDraftMessage | undefined;
let streamGeneration = 0;
let cleanupGeneration = -1;
let untrackConversationBoundary: (() => void) | undefined;
let lastVisibleUpdate: { text: string; blocks?: (Block | KnownBlock)[] } | undefined;
@ -85,9 +84,6 @@ export function createSlackDraftStream(params: {
};
const sendOrEditStreamMessage = async (pending: SlackDraftStreamUpdate) => {
if (streamState.stopped) {
return;
}
const update = normalizeUpdate(pending);
const trimmed = update.text.trimEnd();
if (!trimmed) {
@ -107,7 +103,7 @@ export function createSlackDraftStream(params: {
return;
}
lastSentKey = sentKey;
const generation = streamGeneration;
const generation = lifecycle.generation;
try {
if (streamMessage) {
const message = streamMessage;
@ -124,57 +120,57 @@ export function createSlackDraftStream(params: {
return;
}
const threadTs = params.resolveThreadTs?.();
const pendingBoundary = params.conversationChannelId
? trackSlackDraftMessage({
accountId: params.accountId,
teamId: params.eventScope?.teamId,
channelId: params.conversationChannelId,
threadTs,
onInterveningMessage: () => forceNewMessage("human"),
})
: undefined;
untrackConversationBoundary = pendingBoundary?.stop;
const sent = await send(params.target, trimmed, {
cfg: params.cfg,
token: params.token,
accountId: params.accountId,
threadTs,
identity: params.identity,
eventScope: params.eventScope,
...(params.metadata ? { metadata: params.metadata } : {}),
...(blocks ? { blocks } : {}),
});
const sentMessage = { channelId: sent.channelId, messageId: sent.messageId, generation };
if (generation !== streamGeneration) {
if (sent.channelId && sent.messageId) {
void lifecycle.retire(sentMessage, { defer: true });
}
return;
}
if (!sent.channelId || !sent.messageId) {
stopTrackingConversationBoundary();
streamState.stopped = true;
params.warn?.("slack stream preview stopped (missing identifiers from sendMessage)");
return;
}
streamMessage = sentMessage;
lastVisibleUpdate = { text: trimmed, ...(blocks ? { blocks } : {}) };
if (pendingBoundary && params.conversationChannelId === streamMessage.channelId) {
pendingBoundary.setMessageTs(streamMessage.messageId);
} else {
stopTrackingConversationBoundary();
const tracker = trackSlackDraftMessage({
const trackMessage = (channelId: string, messageTs?: string) =>
trackSlackDraftMessage({
accountId: params.accountId,
teamId: params.eventScope?.teamId,
channelId: streamMessage.channelId,
channelId,
threadTs,
messageTs: streamMessage.messageId,
messageTs,
onInterveningMessage: () => forceNewMessage("human"),
});
untrackConversationBoundary = tracker.stop;
}
const pendingBoundary = params.conversationChannelId
? trackMessage(params.conversationChannelId)
: undefined;
untrackConversationBoundary = pendingBoundary?.stop;
await lifecycle.createMessage(
async () => {
const sent = await send(params.target, trimmed, {
cfg: params.cfg,
token: params.token,
accountId: params.accountId,
threadTs,
identity: params.identity,
eventScope: params.eventScope,
...(params.metadata ? { metadata: params.metadata } : {}),
...(blocks ? { blocks } : {}),
});
return sent.channelId && sent.messageId
? { channelId: sent.channelId, messageId: sent.messageId, generation }
: undefined;
},
(sentMessage) => {
if (!sentMessage) {
stopTrackingConversationBoundary();
streamState.stopped = true;
params.warn?.("slack stream preview stopped (missing identifiers from sendMessage)");
return false;
}
const { channelId, messageId } = sentMessage;
streamMessage = sentMessage;
lastVisibleUpdate = { text: trimmed, ...(blocks ? { blocks } : {}) };
if (pendingBoundary && params.conversationChannelId === channelId) {
pendingBoundary.setMessageTs(messageId);
} else {
stopTrackingConversationBoundary();
untrackConversationBoundary = trackMessage(channelId, messageId).stop;
}
return true;
},
{ defer: true },
);
} catch (err) {
if (generation === streamGeneration) {
if (generation === lifecycle.generation) {
stopTrackingConversationBoundary();
streamState.stopped = true;
}
@ -216,7 +212,7 @@ export function createSlackDraftStream(params: {
};
const dropDetachedMessages = () => {
const generation = streamGeneration;
const generation = lifecycle.generation;
return lifecycle.cleanupPending(() => {
preserveHumanReplies = false;
cleanupGeneration = generation;
@ -226,9 +222,9 @@ export function createSlackDraftStream(params: {
const discardPendingAndStopTracking = async () => {
// A human can reply before Slack returns the pending preview's identity.
// Reconcile that receipt before removing the conversation boundary tracker.
const generation = streamGeneration;
const generation = lifecycle.generation;
await discardPending();
if (generation === streamGeneration) {
if (generation === lifecycle.generation) {
stopTrackingConversationBoundary();
}
};
@ -236,13 +232,13 @@ export function createSlackDraftStream(params: {
const clear = async (options?: { preserveHumanReplies?: boolean }) => {
// Final delivery preserves human-replied context, while failed active
// deletions and explicit rotations remain eligible for cleanup.
const generation = streamGeneration;
const generation = lifecycle.generation;
let clearingMessage = streamMessage;
await lifecycle.clearWithStop(
async () => {
if (generation === streamGeneration) {
if (generation === lifecycle.generation) {
await discardPendingAndStopTracking();
if (generation === streamGeneration) {
if (generation === lifecycle.generation) {
clearingMessage = streamMessage;
clearMessageId();
}
@ -261,21 +257,20 @@ export function createSlackDraftStream(params: {
const forceNewMessage = (reason: "turn" | "human" = "turn") => {
stopTrackingConversationBoundary();
// Human boundaries change the target without reopening a stream that is
// closing. Only explicit admission of another turn resumes delivery.
if (reason === "turn") {
streamState.stopped = false;
streamGeneration += 1;
}
streamState.final = false;
if (streamMessage && !streamMessage.retained) {
// Slack decides whether human-replied context may be removed; the shared
// lifecycle retains deletion custody until that decision is made.
streamMessage.detachedByHuman = reason === "human";
void lifecycle.retire(streamMessage, { defer: true });
}
clearMessageId();
loop.resetPending();
// Human boundaries change the target without reopening a stream that is
// closing. Only explicit admission of another turn resumes delivery.
if (reason === "turn") {
lifecycle.reset("discard", "keep");
} else {
streamState.final = false;
lifecycle.resetMessage("keep");
}
};
const finalizeMessage = async (

View file

@ -261,12 +261,17 @@ export function createFinalizableDraftLifecycle<TMessageId, TUpdate = string>(
clearTail = stopRun;
return stopRun;
};
const resetMessage = () => {
const resetMessage = (throttle: "reset" | "keep" = "reset") => {
params.clearMessageId();
controls.loop.resetPending();
controls.loop.resetThrottleWindow();
if (throttle === "reset") {
controls.loop.resetThrottleWindow();
}
};
const reset = (mode: "preserve" | "discard" = "preserve") => {
const reset = (
mode: "preserve" | "discard" = "preserve",
throttle: "reset" | "keep" = "reset",
) => {
// A later rotation cannot revoke an earlier request to discard an in-flight create.
if (mode === "discard") {
discardThroughGeneration = generation;
@ -274,11 +279,12 @@ export function createFinalizableDraftLifecycle<TMessageId, TUpdate = string>(
generation += 1;
params.state.stopped = false;
params.state.final = false;
resetMessage();
resetMessage(throttle);
};
const createMessage = async (
send: () => Promise<TMessageId | undefined>,
publish: (messageId: TMessageId | undefined) => boolean,
retirementOptions?: { defer?: boolean },
): Promise<boolean> => {
const startedGeneration = generation;
const messageId = await send();
@ -287,7 +293,7 @@ export function createFinalizableDraftLifecycle<TMessageId, TUpdate = string>(
return publish(messageId);
}
if (startedGeneration <= discardThroughGeneration && params.isValidMessageId(messageId)) {
await retire(messageId);
await retire(messageId, retirementOptions);
}
return true;
};

View file

@ -1,20 +1,23 @@
// Tests shared channel draft chunking resolution exposed through plugin-sdk/channel-outbound.
import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts";
import { describe, expect, it } from "vitest";
import { resolveChannelDraftStreamingChunking } from "./channel-outbound.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { resolveChannelDraftStreamingChunking } from "./draft-streaming-chunking.js";
describe("resolveChannelDraftStreamingChunking", () => {
it("returns draft stream defaults when channel config is unset", () => {
expect(
resolveChannelDraftStreamingChunking(undefined, "telegram", "default", {
fallbackLimit: 4096,
}),
).toEqual({
minChars: 200,
maxChars: 800,
breakPreference: "paragraph",
});
});
it.each([
{ channelId: "discord", cfg: {}, accountId: undefined, fallbackLimit: 2000 },
{ channelId: "telegram", cfg: undefined, accountId: "default", fallbackLimit: 4096 },
] as const)(
"returns draft stream defaults when $channelId chunking is unset",
({ channelId, cfg, accountId, fallbackLimit }) => {
expect(
resolveChannelDraftStreamingChunking(cfg, channelId, accountId, { fallbackLimit }),
).toEqual({
minChars: 200,
maxChars: 800,
breakPreference: "paragraph",
});
},
);
it("clamps requested draft chunk sizes to the resolved text limit", () => {
const cfg: OpenClawConfig = {