Skip to content

Commit 24a0038

Browse files
committed
fix: publish transcript turn owned entries
1 parent 9428e97 commit 24a0038

3 files changed

Lines changed: 25 additions & 5 deletions

File tree

src/config/sessions/session-accessor.test.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -764,14 +764,17 @@ describe("session accessor file-backed seam", () => {
764764
storePath,
765765
};
766766
const publishOptions: Array<boolean | undefined> = [];
767+
const publishedEntryBatches: unknown[][] = [];
767768

768769
await withOwnedSessionTranscriptWrites(
769770
{
770771
sessionFile: transcriptPath,
771772
sessionKey: scope.sessionKey,
772773
withSessionWriteLock: async (run, options) => {
773774
publishOptions.push(options?.publishOwnedWrite);
774-
return await run();
775+
const result = await run();
776+
publishedEntryBatches.push([...(options?.resolvePublishedEntries?.(result) ?? [])]);
777+
return result;
775778
},
776779
},
777780
async () =>
@@ -793,6 +796,11 @@ describe("session accessor file-backed seam", () => {
793796
);
794797

795798
expect(publishOptions).toEqual([true]);
799+
expect(publishedEntryBatches).toHaveLength(1);
800+
expect(publishedEntryBatches[0]).toEqual([
801+
expect.objectContaining({ kind: "header" }),
802+
expect.objectContaining({ kind: "id" }),
803+
]);
796804
await expect(loadTranscriptEvents(scope)).resolves.toEqual([
797805
expect.objectContaining({ type: "session" }),
798806
expect.objectContaining({

src/config/sessions/session-accessor.ts

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ import {
4242
import { resolveSessionTranscriptFile } from "./transcript-file-resolve.js";
4343
import { streamSessionTranscriptLines } from "./transcript-stream.js";
4444
import {
45+
type OwnedSessionTranscriptPublishedEntry,
4546
resolveOwnedSessionTranscriptWriteLockRunner,
4647
withOwnedSessionTranscriptWrites,
4748
} from "./transcript-write-context.js";
@@ -511,6 +512,7 @@ async function appendTranscriptTurnMessages(
511512
options: SessionTranscriptTurnPersistOptions,
512513
): Promise<TranscriptMessageAppendResult<unknown>[]> {
513514
const appendedMessages: TranscriptMessageAppendResult<unknown>[] = [];
515+
const publishedEntries: OwnedSessionTranscriptPublishedEntry[] = [];
514516
const appendMessages = async (appendMessage: SessionTranscriptTurnAppendRunner) => {
515517
for (const append of options.messages) {
516518
const shouldAppend = append.shouldAppend
@@ -535,12 +537,18 @@ async function appendTranscriptTurnMessages(
535537
...(append.prepareMessageAfterIdempotencyCheck
536538
? { prepareMessageAfterIdempotencyCheck: append.prepareMessageAfterIdempotencyCheck }
537539
: {}),
540+
onHeaderCreated: (header) => {
541+
publishedEntries.push({ kind: "header", serialized: header });
542+
},
538543
...(append.useRawWhenLinear !== undefined
539544
? { useRawWhenLinear: append.useRawWhenLinear }
540545
: {}),
541546
});
542547
if (result) {
543548
appendedMessages.push(result);
549+
if (result.appended) {
550+
publishedEntries.push({ kind: "id", id: result.messageId });
551+
}
544552
}
545553
}
546554
};
@@ -560,7 +568,11 @@ async function appendTranscriptTurnMessages(
560568
if (activeLockRunner) {
561569
await activeLockRunner(
562570
() => withSessionTranscriptAppendQueue(target.sessionFile, runBatchWithOwnedLock),
563-
{ publishOwnedWrite: true },
571+
{
572+
publishOwnedWrite: true,
573+
resolvePublishedEntries: () => publishedEntries,
574+
resolvePublishedEntriesAfterFailure: () => publishedEntries,
575+
},
564576
);
565577
} else {
566578
await withSessionTranscriptAppendQueue(target.sessionFile, async () => {

src/config/sessions/transcript-append.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -375,6 +375,8 @@ export type AppendSessionTranscriptMessageParams<TMessage = unknown> = {
375375
/** Runs under the transcript write lock after idempotency replay checks and before append. */
376376
prepareMessageAfterIdempotencyCheck?: (message: TMessage) => TMessage | undefined;
377377
config?: OpenClawConfig;
378+
/** Internal owned-batch hook for publishing a newly created transcript header. */
379+
onHeaderCreated?: (serializedHeader: string) => void;
378380
};
379381

380382
export type AppendSessionTranscriptMessageResult<TMessage> = {
@@ -521,9 +523,7 @@ async function appendSessionTranscriptEventLocked(
521523
}
522524

523525
async function appendSessionTranscriptMessageLocked<TMessage>(
524-
params: AppendSessionTranscriptMessageParams<TMessage> & {
525-
onHeaderCreated?: (serializedHeader: string) => void;
526-
},
526+
params: AppendSessionTranscriptMessageParams<TMessage>,
527527
): Promise<AppendSessionTranscriptMessageResult<TMessage> | undefined> {
528528
const now = params.now ?? Date.now();
529529
const serializedHeader = await ensureTranscriptHeader(params.transcriptPath, {

0 commit comments

Comments
 (0)