mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 17:53:39 +00:00
* refactor(state): register worker operations once per domain Infer shared-state worker command contracts from lazy per-domain handler tables. Migrate Web Push, APNs, worktree registry operations, and fleet registry while preserving the existing broker and transaction owners. * refactor(state): register transcript and skill worker operations Infer shared-state operation contracts from lazy handler tables for transcripts, skills, ACP, auth profiles, and plugin runtime state. Preserve native transactions, worker admission, lease checks, receipts, and existing update behavior. * fix(state): break registry type import cycles Keep inferred worker contracts downstream of leaf ACP and lease data types. Separate host lease dispatch from native lifecycle storage while preserving captured admission and package replacement preloading.
53 lines
1.9 KiB
TypeScript
53 lines
1.9 KiB
TypeScript
import { createDeferredCore } from "../shared/deferred.js";
|
|
import type { TranscriptAppendScheduler } from "./store-worker.types.js";
|
|
|
|
type AppendOutcome = { ok: true } | { ok: false; error: unknown };
|
|
|
|
/** Accepted speech belongs to its capture until the store has settled it. */
|
|
export function createTranscriptCaptureAppends(assertCurrent: () => void) {
|
|
let tail = Promise.resolve();
|
|
const pending = new Set<Promise<AppendOutcome>>();
|
|
return {
|
|
async run(prepare: (schedule: TranscriptAppendScheduler) => Promise<void>): Promise<void> {
|
|
const previous = tail;
|
|
const settled = createDeferredCore<AppendOutcome>();
|
|
let active = true;
|
|
// Reserve before metadata serialization can call user code or signal termination.
|
|
pending.add(settled.promise);
|
|
tail = previous.then(() => settled.promise).then(() => undefined);
|
|
const assertAccepted = () => {
|
|
if (!active) {
|
|
throw new Error("Transcript append has already settled");
|
|
}
|
|
assertCurrent();
|
|
};
|
|
const schedule: TranscriptAppendScheduler = (write) =>
|
|
previous.then(() => {
|
|
assertAccepted();
|
|
return write(assertAccepted);
|
|
});
|
|
const finish = (outcome: AppendOutcome) => {
|
|
active = false;
|
|
pending.delete(settled.promise);
|
|
settled.resolve(outcome);
|
|
};
|
|
try {
|
|
await prepare(schedule);
|
|
finish({ ok: true });
|
|
} catch (error) {
|
|
finish({ ok: false, error });
|
|
throw error;
|
|
}
|
|
},
|
|
async drain(): Promise<void> {
|
|
const outcomes = await Promise.all(pending);
|
|
const failures = outcomes.flatMap((outcome) => (outcome.ok ? [] : [outcome.error]));
|
|
if (failures.length === 1) {
|
|
throw failures[0];
|
|
}
|
|
if (failures.length > 1) {
|
|
throw new AggregateError(failures, "Accepted transcript appends failed");
|
|
}
|
|
},
|
|
};
|
|
}
|