


























@@ -26,7 +26,10 @@ import {
2626streamSessionTranscriptLinesReverse,
2727} from "./transcript-stream.js";
2828import { isCanonicalSessionTranscriptEntry } from "./transcript-tree.js";
29-import { resolveOwnedSessionTranscriptWriteLockRunner } from "./transcript-write-context.js";
29+import {
30+resolveOwnedSessionTranscriptWriteLockRunner,
31+type OwnedSessionTranscriptPublishedEntry,
32+} from "./transcript-write-context.js";
3033import { CURRENT_SESSION_VERSION } from "./version.js";
31343235const SESSION_MANAGER_APPEND_MAX_BYTES = 8 * 1024 * 1024;
@@ -385,6 +388,13 @@ export type AppendSessionTranscriptMessageResult<TMessage> = {
385388appended: boolean;
386389};
387390391+export type SessionTranscriptAppendTransactionContext = {
392+appendEvent: (event: unknown) => Promise<void>;
393+appendMessage: <TMessage>(
394+params: Omit<AppendSessionTranscriptMessageParams<TMessage>, "config" | "transcriptPath">,
395+) => Promise<AppendSessionTranscriptMessageResult<TMessage> | undefined>;
396+};
397+388398function isTranscriptAgentMessage(value: unknown): value is AgentMessage {
389399return (
390400typeof value === "object" &&
@@ -464,6 +474,57 @@ export async function appendSessionTranscriptMessageWithOwnedWriteLock<TMessage>
464474return await activeLockRunner(() => appendSessionTranscriptMessageLocked(params));
465475}
466476477+/**
478+ * Runs a group of transcript appends through one append queue and write lock.
479+ */
480+export async function runSessionTranscriptAppendTransaction<T>(
481+params: Pick<AppendSessionTranscriptMessageParams, "config" | "transcriptPath">,
482+run: (context: SessionTranscriptAppendTransactionContext) => Promise<T> | T,
483+): Promise<T> {
484+const publishedEntries: OwnedSessionTranscriptPublishedEntry[] = [];
485+const runTransaction = async (): Promise<T> =>
486+await run({
487+appendEvent: async (event) => {
488+const result = await appendSessionTranscriptEventLocked({
489+config: params.config,
490+ event,
491+transcriptPath: params.transcriptPath,
492+});
493+publishedEntries.push({ kind: "serialized", serialized: result.serializedEntry });
494+},
495+appendMessage: async (messageParams) => {
496+const result = await appendSessionTranscriptMessageLocked({
497+ ...messageParams,
498+config: params.config,
499+onHeaderCreated: (header) => {
500+publishedEntries.push({ kind: "header", serialized: header });
501+},
502+transcriptPath: params.transcriptPath,
503+});
504+if (result?.appended === true) {
505+publishedEntries.push({ kind: "id", id: result.messageId });
506+}
507+return result;
508+},
509+});
510+const activeLockRunner = resolveOwnedSessionTranscriptWriteLockRunner({
511+sessionFile: params.transcriptPath,
512+});
513+if (activeLockRunner) {
514+return await activeLockRunner(
515+() => withSessionTranscriptAppendQueue(params.transcriptPath, runTransaction),
516+{
517+publishOwnedWrite: true,
518+resolvePublishedEntries: () => publishedEntries,
519+resolvePublishedEntriesAfterFailure: () => publishedEntries,
520+},
521+);
522+}
523+return await withSessionTranscriptAppendQueue(params.transcriptPath, () =>
524+withSessionTranscriptWriteLock(params, runTransaction),
525+);
526+}
527+467528type AppendSessionTranscriptEventParams = {
468529config?: OpenClawConfig;
469530event: unknown;
此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。