mirror of
https://github.com/openclaw/openclaw.git
synced 2026-05-06 12:41:37 +00:00
fix(auto-reply): coalesce block replies and document streaming toggles (#536) (thanks @mcinteerj)
This commit is contained in:
125
src/auto-reply/reply/block-reply-coalescer.ts
Normal file
125
src/auto-reply/reply/block-reply-coalescer.ts
Normal file
@@ -0,0 +1,125 @@
|
||||
import type { ReplyPayload } from "../types.js";
|
||||
import type { BlockStreamingCoalescing } from "./block-streaming.js";
|
||||
|
||||
export type BlockReplyCoalescer = {
|
||||
enqueue: (payload: ReplyPayload) => void;
|
||||
flush: (options?: { force?: boolean }) => Promise<void>;
|
||||
hasBuffered: () => boolean;
|
||||
stop: () => void;
|
||||
};
|
||||
|
||||
export function createBlockReplyCoalescer(params: {
|
||||
config: BlockStreamingCoalescing;
|
||||
shouldAbort: () => boolean;
|
||||
onFlush: (payload: ReplyPayload) => Promise<void> | void;
|
||||
}): BlockReplyCoalescer {
|
||||
const { config, shouldAbort, onFlush } = params;
|
||||
const minChars = Math.max(1, Math.floor(config.minChars));
|
||||
const maxChars = Math.max(minChars, Math.floor(config.maxChars));
|
||||
const idleMs = Math.max(0, Math.floor(config.idleMs));
|
||||
const joiner = config.joiner ?? "";
|
||||
|
||||
let bufferText = "";
|
||||
let bufferReplyToId: ReplyPayload["replyToId"];
|
||||
let bufferAudioAsVoice: ReplyPayload["audioAsVoice"];
|
||||
let idleTimer: NodeJS.Timeout | undefined;
|
||||
|
||||
const clearIdleTimer = () => {
|
||||
if (!idleTimer) return;
|
||||
clearTimeout(idleTimer);
|
||||
idleTimer = undefined;
|
||||
};
|
||||
|
||||
const resetBuffer = () => {
|
||||
bufferText = "";
|
||||
bufferReplyToId = undefined;
|
||||
bufferAudioAsVoice = undefined;
|
||||
};
|
||||
|
||||
const scheduleIdleFlush = () => {
|
||||
if (idleMs <= 0) return;
|
||||
clearIdleTimer();
|
||||
idleTimer = setTimeout(() => {
|
||||
void flush({ force: false });
|
||||
}, idleMs);
|
||||
};
|
||||
|
||||
const flush = async (options?: { force?: boolean }) => {
|
||||
clearIdleTimer();
|
||||
if (shouldAbort()) {
|
||||
resetBuffer();
|
||||
return;
|
||||
}
|
||||
if (!bufferText) return;
|
||||
if (!options?.force && bufferText.length < minChars) {
|
||||
scheduleIdleFlush();
|
||||
return;
|
||||
}
|
||||
const payload: ReplyPayload = {
|
||||
text: bufferText,
|
||||
replyToId: bufferReplyToId,
|
||||
audioAsVoice: bufferAudioAsVoice,
|
||||
};
|
||||
resetBuffer();
|
||||
await onFlush(payload);
|
||||
};
|
||||
|
||||
const enqueue = (payload: ReplyPayload) => {
|
||||
if (shouldAbort()) return;
|
||||
const hasMedia =
|
||||
Boolean(payload.mediaUrl) || (payload.mediaUrls?.length ?? 0) > 0;
|
||||
const text = payload.text ?? "";
|
||||
const hasText = text.trim().length > 0;
|
||||
if (hasMedia) {
|
||||
void flush({ force: true });
|
||||
void onFlush(payload);
|
||||
return;
|
||||
}
|
||||
if (!hasText) return;
|
||||
|
||||
if (
|
||||
bufferText &&
|
||||
(bufferReplyToId !== payload.replyToId ||
|
||||
bufferAudioAsVoice !== payload.audioAsVoice)
|
||||
) {
|
||||
void flush({ force: true });
|
||||
}
|
||||
|
||||
if (!bufferText) {
|
||||
bufferReplyToId = payload.replyToId;
|
||||
bufferAudioAsVoice = payload.audioAsVoice;
|
||||
}
|
||||
|
||||
const nextText = bufferText ? `${bufferText}${joiner}${text}` : text;
|
||||
if (nextText.length > maxChars) {
|
||||
if (bufferText) {
|
||||
void flush({ force: true });
|
||||
bufferReplyToId = payload.replyToId;
|
||||
bufferAudioAsVoice = payload.audioAsVoice;
|
||||
if (text.length >= maxChars) {
|
||||
void onFlush(payload);
|
||||
return;
|
||||
}
|
||||
bufferText = text;
|
||||
scheduleIdleFlush();
|
||||
return;
|
||||
}
|
||||
void onFlush(payload);
|
||||
return;
|
||||
}
|
||||
|
||||
bufferText = nextText;
|
||||
if (bufferText.length >= maxChars) {
|
||||
void flush({ force: true });
|
||||
return;
|
||||
}
|
||||
scheduleIdleFlush();
|
||||
};
|
||||
|
||||
return {
|
||||
enqueue,
|
||||
flush,
|
||||
hasBuffered: () => Boolean(bufferText),
|
||||
stop: () => clearIdleTimer(),
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user