From 5b4fdfcb22cf93cb395fa7849dcbc2fa4d30e016 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 1 Oct 2026 03:05:44 -0700 Subject: [PATCH] refactor(state): register transcript, skill and runtime worker operations (#162524) * 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. --- .../control-plane/manager.accepted-turns.ts | 6 +- .../control-plane/manager.cancel-session.ts | 2 +- .../manager.runtime-handle-ensure.ts | 2 +- .../manager.runtime-resume-state.ts | 2 +- src/acp/control-plane/manager.types.ts | 4 +- src/acp/runtime/session-control-owner.ts | 18 +- src/acp/runtime/session-meta-control.ts | 6 +- src/acp/runtime/session-meta-control.types.ts | 20 +- src/acp/runtime/session-meta-doctor.ts | 2 +- src/acp/runtime/session-meta-entry.kernel.ts | 6 +- src/acp/runtime/session-meta-entry.ts | 2 +- src/acp/runtime/session-meta-entry.types.ts | 2 +- src/acp/runtime/session-meta-keys.ts | 18 +- .../runtime/session-meta-legacy-cleanup.ts | 2 +- src/acp/runtime/session-meta-read.types.ts | 13 + src/acp/runtime/session-meta-readonly.ts | 3 +- src/acp/runtime/session-meta-write.kernel.ts | 2 +- src/acp/runtime/session-meta-write.native.ts | 2 +- src/acp/runtime/session-meta-write.types.ts | 33 +- .../session-meta-write.worker-contract.ts | 4 + src/acp/runtime/session-meta-write.worker.ts | 22 +- .../auth-profiles/store.worker-contract.ts | 4 + src/agents/auth-profiles/store.worker.ts | 51 +++ src/agents/mcp-oauth-store.ts | 2 +- src/agents/mcp-oauth-store.types.ts | 2 +- src/infra/deferred-plugin-migrations.ts | 4 +- .../deferred-plugin-migrations.worker.ts | 55 --- ...-catalog-snapshot-store.worker-contract.ts | 12 - src/plugins/state.worker-contract.ts | 4 + src/plugins/state.worker.ts | 104 ++++++ .../project-clone-registration.test.ts | 2 +- src/projects/project-environment.test.ts | 2 +- src/projects/project-registration.ts | 2 +- src/projects/project-registry.ts | 4 +- .../project-registry.worker-contract.ts | 2 +- src/skills/lifecycle/upload-store.ts | 2 +- .../lifecycle/upload-store.worker-contract.ts | 4 + src/skills/lifecycle/upload-store.worker.ts | 109 ++---- src/skills/workshop/store-client.ts | 4 +- src/skills/workshop/store.worker-contract.ts | 75 +--- src/skills/workshop/store.worker.ts | 345 ++++++++++-------- src/state/openclaw-state-lease-acquisition.ts | 2 +- ...claw-state-lease-async.maintenance.test.ts | 2 +- src/state/openclaw-state-lease-context.ts | 2 +- .../openclaw-state-lease-heartbeat-shared.ts | 2 +- src/state/openclaw-state-lease-storage.ts | 5 +- src/state/openclaw-state-lease-store.ts | 8 +- .../openclaw-state-lease-worker-group.test.ts | 4 +- .../openclaw-state-lease-worker-operation.ts | 66 ++++ .../openclaw-state-lease-worker-owner.ts | 2 +- .../openclaw-state-lease-worker-storage.ts | 64 +--- ...ate-lease-worker.transaction-group.test.ts | 2 +- src/state/openclaw-state-lease-worker.ts | 2 +- src/state/openclaw-state-lease.types.ts | 4 + .../openclaw-state-read-worker.mcp.test.ts | 2 +- src/state/openclaw-state-read.types.ts | 2 +- src/state/openclaw-state-worker-contract.ts | 75 +--- src/state/openclaw-state-worker-registry.ts | 30 +- src/state/openclaw-state-worker-runtime.ts | 138 ------- src/transcripts/capture-appends.ts | 2 +- src/transcripts/store-worker-client.ts | 10 +- src/transcripts/store-worker-contract.ts | 57 --- src/transcripts/store-worker-write.ts | 219 ++++++----- src/transcripts/store-worker.types.ts | 8 + .../store-write.worker-contract.ts | 4 + src/transcripts/store.ts | 10 +- 66 files changed, 771 insertions(+), 911 deletions(-) create mode 100644 src/acp/runtime/session-meta-read.types.ts create mode 100644 src/acp/runtime/session-meta-write.worker-contract.ts create mode 100644 src/agents/auth-profiles/store.worker-contract.ts create mode 100644 src/agents/auth-profiles/store.worker.ts delete mode 100644 src/infra/deferred-plugin-migrations.worker.ts delete mode 100644 src/plugins/official-external-plugin-catalog-snapshot-store.worker-contract.ts create mode 100644 src/plugins/state.worker-contract.ts create mode 100644 src/plugins/state.worker.ts create mode 100644 src/skills/lifecycle/upload-store.worker-contract.ts create mode 100644 src/state/openclaw-state-lease-worker-operation.ts create mode 100644 src/state/openclaw-state-lease.types.ts create mode 100644 src/transcripts/store-worker.types.ts create mode 100644 src/transcripts/store-write.worker-contract.ts diff --git a/src/acp/control-plane/manager.accepted-turns.ts b/src/acp/control-plane/manager.accepted-turns.ts index 2e3fe7e85deb..427fbefe02f1 100644 --- a/src/acp/control-plane/manager.accepted-turns.ts +++ b/src/acp/control-plane/manager.accepted-turns.ts @@ -1,7 +1,9 @@ /** Owns each admitted turn from actor admission through final settlement. */ import { createDeferredCore } from "../../shared/deferred.js"; -import type { AcpSessionRuntimeLocator } from "../runtime/session-control-owner.js"; -import type { AcpSessionControlConstraint } from "../runtime/session-meta-control.types.js"; +import type { + AcpSessionRuntimeLocator, + AcpSessionControlConstraint, +} from "../runtime/session-meta-control.types.js"; import type { AcpRunTurnInput, ActiveTurnState, diff --git a/src/acp/control-plane/manager.cancel-session.ts b/src/acp/control-plane/manager.cancel-session.ts index 8623e8d30932..87570cde33d5 100644 --- a/src/acp/control-plane/manager.cancel-session.ts +++ b/src/acp/control-plane/manager.cancel-session.ts @@ -7,8 +7,8 @@ import { import { matchesAcpSessionRuntimeLocator, resolveAcpSessionControlOwner, - type AcpSessionRuntimeLocator, } from "../runtime/session-control-owner.js"; +import type { AcpSessionRuntimeLocator } from "../runtime/session-meta-control.types.js"; import type { AcceptedTurnState, AcceptedTurns } from "./manager.accepted-turns.js"; import type { ManagerRuntimeHandleCache } from "./manager.runtime-handle-cache.js"; import type { diff --git a/src/acp/control-plane/manager.runtime-handle-ensure.ts b/src/acp/control-plane/manager.runtime-handle-ensure.ts index f15315f8d810..742151a88d1b 100644 --- a/src/acp/control-plane/manager.runtime-handle-ensure.ts +++ b/src/acp/control-plane/manager.runtime-handle-ensure.ts @@ -17,8 +17,8 @@ import { import { matchesAcpSessionControlBinding, resolveAcpSessionControlOwner, - type AcpSessionRuntimeLocator, } from "../runtime/session-control-owner.js"; +import type { AcpSessionRuntimeLocator } from "../runtime/session-meta-control.types.js"; import { assertAcpSessionMutationEntry } from "../runtime/session-meta-entry.kernel.js"; import type { ManagerRuntimeHandleCache } from "./manager.runtime-handle-cache.js"; import { diff --git a/src/acp/control-plane/manager.runtime-resume-state.ts b/src/acp/control-plane/manager.runtime-resume-state.ts index ed5e502c108a..fc44cee253ab 100644 --- a/src/acp/control-plane/manager.runtime-resume-state.ts +++ b/src/acp/control-plane/manager.runtime-resume-state.ts @@ -4,7 +4,7 @@ import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { logVerbose } from "../../globals.js"; import { formatErrorMessage, toErrorObject } from "../../infra/errors.js"; import type { AcpRuntimeError } from "../runtime/errors.js"; -import type { AcpSessionControlBinding } from "../runtime/session-control-owner.js"; +import type { AcpSessionControlBinding } from "../runtime/session-meta-control.types.js"; import type { ManagerRuntimeHandleCache } from "./manager.runtime-handle-cache.js"; import { assertAcpRuntimeOwnerSupport, diff --git a/src/acp/control-plane/manager.types.ts b/src/acp/control-plane/manager.types.ts index 1ba299fdc45f..180706869d03 100644 --- a/src/acp/control-plane/manager.types.ts +++ b/src/acp/control-plane/manager.types.ts @@ -21,9 +21,9 @@ import type { AcpRuntimeError } from "../runtime/errors.js"; import { getAcpRuntimeBackend, requireAcpRuntimeBackend } from "../runtime/registry.js"; import type { AcpSessionControlBinding, + AcpSessionControlConstraint, AcpSessionRuntimeLocator, -} from "../runtime/session-control-owner.js"; -import type { AcpSessionControlConstraint } from "../runtime/session-meta-control.types.js"; +} from "../runtime/session-meta-control.types.js"; import { listAcpSessionEntries, readAcpSessionEntry, diff --git a/src/acp/runtime/session-control-owner.ts b/src/acp/runtime/session-control-owner.ts index 1eb0dfae8581..7a35dc896b09 100644 --- a/src/acp/runtime/session-control-owner.ts +++ b/src/acp/runtime/session-control-owner.ts @@ -1,8 +1,8 @@ -import type { SessionAcpMeta, SessionEntry } from "../../config/sessions/types.js"; - -export type AcpSessionRuntimeLocator = Readonly< - Pick ->; +import type { SessionEntry } from "../../config/sessions/types.js"; +import type { + AcpSessionRuntimeLocator, + AcpSessionControlBinding, +} from "./session-meta-control.types.js"; /** Runtime names are opaque backend locators; ordinary metadata enrichment may continue. */ export function matchesAcpSessionRuntimeLocator( @@ -22,14 +22,6 @@ export function resolveAcpSessionControlOwner( return entry?.spawnedBy?.trim() || entry?.parentSessionKey?.trim(); } -/** A cleanup target constraint; live task and actor authority remain separate. */ -export type AcpSessionControlBinding = Readonly<{ - sessionId: string; - lifecycleRevision?: string; - sessionStartedAt?: number; - ownerKey: string; -}>; - export function matchesAcpSessionControlBinding( entry: SessionEntry | undefined, expected: AcpSessionControlBinding, diff --git a/src/acp/runtime/session-meta-control.ts b/src/acp/runtime/session-meta-control.ts index c261f0e95643..9e029568686c 100644 --- a/src/acp/runtime/session-meta-control.ts +++ b/src/acp/runtime/session-meta-control.ts @@ -7,9 +7,11 @@ import { captureOpenClawStateReadContext } from "../../state/openclaw-state-work import { matchesAcpSessionRuntimeLocator, resolveAcpSessionControlOwner, - type AcpSessionRuntimeLocator, } from "./session-control-owner.js"; -import type { AcpSessionControlConstraint } from "./session-meta-control.types.js"; +import type { + AcpSessionRuntimeLocator, + AcpSessionControlConstraint, +} from "./session-meta-control.types.js"; import { assertAcpSessionMutationEntry, captureAcpSessionEntryBinding, diff --git a/src/acp/runtime/session-meta-control.types.ts b/src/acp/runtime/session-meta-control.types.ts index 317dffd1b374..f9c3da77deb7 100644 --- a/src/acp/runtime/session-meta-control.types.ts +++ b/src/acp/runtime/session-meta-control.types.ts @@ -1,10 +1,18 @@ -import type { SessionEntry } from "../../config/sessions/types.js"; +import type { SessionAcpMeta, SessionEntry } from "../../config/sessions/types.js"; import type { DatabasePathIdentity } from "../../infra/sqlite-worker-identity.js"; -import type { - AcpSessionControlBinding, - AcpSessionRuntimeLocator, -} from "./session-control-owner.js"; -import type { AcpSessionReadInput } from "./session-meta-keys.js"; +import type { AcpSessionReadInput } from "./session-meta-read.types.js"; + +export type AcpSessionRuntimeLocator = Readonly< + Pick +>; + +/** A cleanup target constraint; live task and actor authority remain separate. */ +export type AcpSessionControlBinding = Readonly<{ + sessionId: string; + lifecycleRevision?: string; + sessionStartedAt?: number; + ownerKey: string; +}>; export type AcpSessionSourceReadInput = { source: { diff --git a/src/acp/runtime/session-meta-doctor.ts b/src/acp/runtime/session-meta-doctor.ts index 88a8a1c9207a..ee026c9bf735 100644 --- a/src/acp/runtime/session-meta-doctor.ts +++ b/src/acp/runtime/session-meta-doctor.ts @@ -31,8 +31,8 @@ import { selectAcpSessionRow, selectLegacyFreeAcpSessionRows, upsertAcpSessionMetaRow, - type AcpSessionRow, } from "./session-meta-keys.js"; +import type { AcpSessionRow } from "./session-meta-read.types.js"; import { rowToAcpSessionMeta } from "./session-meta-readonly.js"; import { resolveSessionStorePathForAcp } from "./session-meta-store.js"; diff --git a/src/acp/runtime/session-meta-entry.kernel.ts b/src/acp/runtime/session-meta-entry.kernel.ts index cd19b2c97ccf..34b464fa2a06 100644 --- a/src/acp/runtime/session-meta-entry.kernel.ts +++ b/src/acp/runtime/session-meta-entry.kernel.ts @@ -1,8 +1,6 @@ import type { SessionEntry } from "../../config/sessions/types.js"; -import { - matchesAcpSessionControlBinding, - type AcpSessionControlBinding, -} from "./session-control-owner.js"; +import { matchesAcpSessionControlBinding } from "./session-control-owner.js"; +import type { AcpSessionControlBinding } from "./session-meta-control.types.js"; export type AcpSessionEntryExpectation = Pick< SessionEntry, diff --git a/src/acp/runtime/session-meta-entry.ts b/src/acp/runtime/session-meta-entry.ts index 154bb805fa67..148f8c60ac8d 100644 --- a/src/acp/runtime/session-meta-entry.ts +++ b/src/acp/runtime/session-meta-entry.ts @@ -15,7 +15,7 @@ import type { SqliteWorkerOperationAdmission } from "../../infra/sqlite-worker-o import type { RetainedWorkerTransactionAdmission } from "../../infra/sqlite-worker-operation-settlement.js"; import type { OpenClawAgentDatabaseOptions } from "../../state/openclaw-agent-db.js"; import type { OpenClawAgentDatabaseExecution } from "../../state/openclaw-agent-execution.js"; -import type { AcpSessionControlBinding } from "./session-control-owner.js"; +import type { AcpSessionControlBinding } from "./session-meta-control.types.js"; import { captureAcpSessionEntryBinding, type AcpSessionEntryExpectation, diff --git a/src/acp/runtime/session-meta-entry.types.ts b/src/acp/runtime/session-meta-entry.types.ts index 4e3b4b2405fd..15b9c84b18dc 100644 --- a/src/acp/runtime/session-meta-entry.types.ts +++ b/src/acp/runtime/session-meta-entry.types.ts @@ -1,6 +1,6 @@ import type { SessionEntryReplacementPublication } from "../../config/sessions/session-accessor.sqlite-entry-cache.js"; import type { SessionEntry } from "../../config/sessions/types.js"; -import type { AcpSessionControlBinding } from "./session-control-owner.js"; +import type { AcpSessionControlBinding } from "./session-meta-control.types.js"; import type { AcpSessionEntryExpectation } from "./session-meta-entry.kernel.js"; export type AcpSessionEntryMutationInput = { diff --git a/src/acp/runtime/session-meta-keys.ts b/src/acp/runtime/session-meta-keys.ts index 8c5d3c8d1ae7..902628d5ffdb 100644 --- a/src/acp/runtime/session-meta-keys.ts +++ b/src/acp/runtime/session-meta-keys.ts @@ -1,9 +1,8 @@ import type { DatabaseSync } from "node:sqlite"; import { normalizeLowercaseStringOrEmpty } from "@openclaw/normalization-core/string-coerce"; -import type { Insertable, Selectable } from "kysely"; +import type { Insertable } from "kysely"; import { tryResolveLegacyDataOwnerAgentId } from "../../agents/agent-scope-config.js"; import { resolvePersistedSessionStoreOwnerForKey } from "../../config/sessions/session-store-owner.js"; -import type { SessionEntry } from "../../config/sessions/types.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { executeSqliteQuerySync, @@ -12,17 +11,14 @@ import { } from "../../infra/kysely-sync.js"; import { normalizeAgentId, parseAgentSessionKey } from "../../routing/session-key.js"; import type { DB as OpenClawStateKyselyDatabase } from "../../state/openclaw-state-db.generated.js"; +import type { + AcpSessionsTable, + AcpSessionRow, + AcpSessionEntryBinding, + AcpSessionReadInput, +} from "./session-meta-read.types.js"; -export type AcpSessionsTable = OpenClawStateKyselyDatabase["acp_sessions"]; type AcpSessionMetaDatabase = Pick; -export type AcpSessionRow = Selectable; -export type AcpSessionEntryBinding = Pick & - Partial>; -export type AcpSessionReadInput = { - keys: readonly string[]; - legacyKey?: string; - entry?: AcpSessionEntryBinding; -}; export function getAcpSessionKysely(db: DatabaseSync) { return getNodeSqliteKysely(db); diff --git a/src/acp/runtime/session-meta-legacy-cleanup.ts b/src/acp/runtime/session-meta-legacy-cleanup.ts index 2d62f289e23c..00e17711c9a3 100644 --- a/src/acp/runtime/session-meta-legacy-cleanup.ts +++ b/src/acp/runtime/session-meta-legacy-cleanup.ts @@ -1,5 +1,5 @@ import { patchSessionEntryWithKey } from "../../config/sessions/session-accessor.js"; -import type { AcpSessionControlBinding } from "./session-control-owner.js"; +import type { AcpSessionControlBinding } from "./session-meta-control.types.js"; import { assertAcpSessionMutationEntry, type AcpSessionEntryExpectation, diff --git a/src/acp/runtime/session-meta-read.types.ts b/src/acp/runtime/session-meta-read.types.ts new file mode 100644 index 000000000000..f3fa092ea813 --- /dev/null +++ b/src/acp/runtime/session-meta-read.types.ts @@ -0,0 +1,13 @@ +import type { Selectable } from "kysely"; +import type { SessionEntry } from "../../config/sessions/types.js"; +import type { DB as OpenClawStateKyselyDatabase } from "../../state/openclaw-state-db.generated.js"; + +export type AcpSessionsTable = OpenClawStateKyselyDatabase["acp_sessions"]; +export type AcpSessionRow = Selectable; +export type AcpSessionEntryBinding = Pick & + Partial>; +export type AcpSessionReadInput = { + keys: readonly string[]; + legacyKey?: string; + entry?: AcpSessionEntryBinding; +}; diff --git a/src/acp/runtime/session-meta-readonly.ts b/src/acp/runtime/session-meta-readonly.ts index e2486ac5eb0e..3a991617ee8e 100644 --- a/src/acp/runtime/session-meta-readonly.ts +++ b/src/acp/runtime/session-meta-readonly.ts @@ -11,14 +11,13 @@ import { withExistingOpenClawStateDatabaseCurrentReadOnly, } from "../../state/openclaw-state-db-readonly.js"; import { - type AcpSessionEntryBinding, - type AcpSessionRow, buildAcpDatabaseSessionKey, legacyAcpDatabaseSessionKeys, resolveLegacyFreeAcpSessionKey, resolveReadableAcpSessionRow, selectAcpSessionRowForStoreEntry, } from "./session-meta-keys.js"; +import type { AcpSessionEntryBinding, AcpSessionRow } from "./session-meta-read.types.js"; /** Each result stays bound to the entry lifecycle captured by the row reader. */ export async function readAcpSessionMetaForEntries( diff --git a/src/acp/runtime/session-meta-write.kernel.ts b/src/acp/runtime/session-meta-write.kernel.ts index 93231ee807f1..56e98165ba67 100644 --- a/src/acp/runtime/session-meta-write.kernel.ts +++ b/src/acp/runtime/session-meta-write.kernel.ts @@ -10,8 +10,8 @@ import { selectAcpSessionRow, selectLegacyFreeAcpSessionRows, upsertAcpSessionMetaRow, - type AcpSessionsTable, } from "./session-meta-keys.js"; +import type { AcpSessionsTable } from "./session-meta-read.types.js"; import type { AcpSessionMutationCommit } from "./session-meta-write.types.js"; export function bindAcpSessionMeta(params: { diff --git a/src/acp/runtime/session-meta-write.native.ts b/src/acp/runtime/session-meta-write.native.ts index 754762d78e45..fa1157838fe9 100644 --- a/src/acp/runtime/session-meta-write.native.ts +++ b/src/acp/runtime/session-meta-write.native.ts @@ -15,7 +15,7 @@ import { } from "../../infra/legacy-acp-migration-source.js"; import { sessionChanges } from "../../sessions/session-row-changes.js"; import { runOpenClawStateWriteTransaction } from "../../state/openclaw-state-db.js"; -import type { AcpSessionControlBinding } from "./session-control-owner.js"; +import type { AcpSessionControlBinding } from "./session-meta-control.types.js"; import { assertAcpSessionMutationEntry } from "./session-meta-entry.kernel.js"; import { selectAcpSessionRowForStoreEntry } from "./session-meta-keys.js"; import { clearLegacyEmbeddedAcpMetadata } from "./session-meta-legacy-cleanup.js"; diff --git a/src/acp/runtime/session-meta-write.types.ts b/src/acp/runtime/session-meta-write.types.ts index 06583662cc70..cd93b61d641a 100644 --- a/src/acp/runtime/session-meta-write.types.ts +++ b/src/acp/runtime/session-meta-write.types.ts @@ -1,10 +1,10 @@ import type { SessionAcpMeta, SessionEntry } from "../../config/sessions/types.js"; -import type { AcpSessionControlBinding } from "./session-control-owner.js"; import type { + AcpSessionControlBinding, AcpSessionControlConstraint, AcpSessionSourceReadInput, } from "./session-meta-control.types.js"; -import type { AcpSessionReadInput } from "./session-meta-keys.js"; +import type { AcpSessionReadInput } from "./session-meta-read.types.js"; export type AcpSessionMutationDecision = | { kind: "keep" } @@ -33,23 +33,14 @@ export type AcpSessionMutationCommit = { control?: AcpSessionControlConstraint; }; -export type AcpSessionWriteOperations = { - "acp.prepareMutation": { - input: { - nonce: string; - read: AcpSessionReadInput; - entry?: SessionEntry; - updatedAt: number; - source: AcpSessionSourceReadInput["source"]; - sessionKey: string; - agentId: string; - expectedControlBinding?: AcpSessionControlBinding; - control?: AcpSessionControlConstraint; - }; - output: AcpSessionMutationPreparation; - }; - "acp.commitMutation": { - input: AcpSessionMutationCommit & { nonce: string }; - output: { nonce: string }; - }; +export type AcpSessionMutationPrepareInput = { + nonce: string; + read: AcpSessionReadInput; + entry?: SessionEntry; + updatedAt: number; + source: AcpSessionSourceReadInput["source"]; + sessionKey: string; + agentId: string; + expectedControlBinding?: AcpSessionControlBinding; + control?: AcpSessionControlConstraint; }; diff --git a/src/acp/runtime/session-meta-write.worker-contract.ts b/src/acp/runtime/session-meta-write.worker-contract.ts new file mode 100644 index 000000000000..22c65adb49ab --- /dev/null +++ b/src/acp/runtime/session-meta-write.worker-contract.ts @@ -0,0 +1,4 @@ +import type { WorkerOperations } from "../../state/worker-operation-registry.js"; +import type { acpSessionOperations } from "./session-meta-write.worker.js"; + +export type AcpSessionWriteOperations = WorkerOperations; diff --git a/src/acp/runtime/session-meta-write.worker.ts b/src/acp/runtime/session-meta-write.worker.ts index c22a96da5f04..d68cd955adc0 100644 --- a/src/acp/runtime/session-meta-write.worker.ts +++ b/src/acp/runtime/session-meta-write.worker.ts @@ -5,7 +5,6 @@ import { legacyAcpMigrationBindingMatches, recordLegacyAcpMigrationCompletion, } from "../../infra/legacy-acp-migration-source.js"; -import type { SqliteWorkerCommand } from "../../infra/sqlite-worker-contract.js"; import { deferSqliteWorkerCommitReceipt, requestSqliteWorkerOperationAdmission, @@ -15,6 +14,7 @@ import { runOpenClawStateWriteTransaction, type OpenClawStateDatabase, } from "../../state/openclaw-state-db.js"; +import type { WorkerOperationHandlers } from "../../state/worker-operation-registry.js"; import type { AcpSessionControlConstraint } from "./session-meta-control.types.js"; import { assertAcpSessionMutationEntry } from "./session-meta-entry.kernel.js"; import { @@ -33,17 +33,15 @@ import type { AcpSessionMutationCommit, AcpSessionMutationDecision, AcpSessionMutationPreparation, - AcpSessionWriteOperations, + AcpSessionMutationPrepareInput, } from "./session-meta-write.types.js"; -export function executeAcpSessionMutationInWorker( - database: OpenClawStateDatabase, - command: SqliteWorkerCommand, -) { - return command.type === "acp.prepareMutation" - ? prepareAcpSessionMutationInWorker(database, command.input) - : commitAcpSessionMutationInWorker(database, command.input); -} +export const acpSessionOperations = { + "acp.prepareMutation": (input: AcpSessionMutationPrepareInput, { open }) => + prepareAcpSessionMutationInWorker(open(), input), + "acp.commitMutation": (input: AcpSessionMutationCommit & { nonce: string }, { open }) => + commitAcpSessionMutationInWorker(open(), input), +} satisfies WorkerOperationHandlers; function readControlledAcpSessionMutation( database: OpenClawStateDatabase, @@ -67,7 +65,7 @@ function readControlledAcpSessionMutation( function prepareAcpSessionMutationInWorker( database: OpenClawStateDatabase, - input: AcpSessionWriteOperations["acp.prepareMutation"]["input"], + input: AcpSessionMutationPrepareInput, ): AcpSessionMutationPreparation { return runOpenClawStateWriteTransaction( ({ db }) => { @@ -133,7 +131,7 @@ function consumeSources(database: OpenClawStateDatabase, input: AcpSessionMutati function commitAcpSessionMutationInWorker( database: OpenClawStateDatabase, - input: AcpSessionWriteOperations["acp.commitMutation"]["input"], + input: AcpSessionMutationCommit & { nonce: string }, ) { return runOpenClawStateWriteTransaction( (current) => { diff --git a/src/agents/auth-profiles/store.worker-contract.ts b/src/agents/auth-profiles/store.worker-contract.ts new file mode 100644 index 000000000000..37cb1b83f1c7 --- /dev/null +++ b/src/agents/auth-profiles/store.worker-contract.ts @@ -0,0 +1,4 @@ +import type { WorkerOperations } from "../../state/worker-operation-registry.js"; +import type { authProfileOperations } from "./store.worker.js"; + +export type AuthProfileWorkerOperations = WorkerOperations; diff --git a/src/agents/auth-profiles/store.worker.ts b/src/agents/auth-profiles/store.worker.ts new file mode 100644 index 000000000000..c6fa5e4b7a2b --- /dev/null +++ b/src/agents/auth-profiles/store.worker.ts @@ -0,0 +1,51 @@ +import { readConfigMachineState } from "../../state/config-machine-state.js"; +import { + withArtifactPreservingStateReads, + withExistingOpenClawStateDatabaseReadOnly, +} from "../../state/openclaw-state-db-readonly.js"; +import { readUserModelAuthProfile } from "../../state/user-model-accounts.js"; +import type { WorkerOperationHandlers } from "../../state/worker-operation-registry.js"; +import { readAuthProfileRows, SHARED_AUTH_STORE_STATE_KEY } from "./sqlite-json.js"; +import { isMissingDatabasePath } from "./sqlite-read-pool.js"; +import type { AuthProfileRowRead } from "./types.js"; + +export const authProfileOperations = { + "authProfiles.read": (input: { artifactPreserving: boolean }, { stateOptions }) => { + const read = (): AuthProfileRowRead => { + const options = stateOptions(); + const missing: AuthProfileRowRead = { + store: { status: "missing", reason: "database" }, + state: { status: "missing", reason: "database" }, + cacheable: false, + }; + try { + return ( + withExistingOpenClawStateDatabaseReadOnly( + ({ db }) => readAuthProfileRows(db, options.path, "shared-state"), + options, + ) ?? missing + ); + } catch { + return isMissingDatabasePath(options.path) + ? missing + : { + store: { status: "unreadable" }, + state: { status: "unreadable" }, + cacheable: false, + }; + } + }; + return input.artifactPreserving ? withArtifactPreservingStateReads(read) : read(); + }, + "authProfiles.sharedOwnership": (input: { artifactPreserving: boolean }, { stateOptions }) => { + const read = () => readConfigMachineState(SHARED_AUTH_STORE_STATE_KEY, stateOptions()); + return input.artifactPreserving ? withArtifactPreservingStateReads(read) : read(); + }, + "authProfiles.personal": ( + input: { profileId: string; artifactPreserving: boolean }, + { stateOptions }, + ) => { + const read = () => readUserModelAuthProfile(input.profileId, stateOptions()); + return input.artifactPreserving ? withArtifactPreservingStateReads(read) : read(); + }, +} satisfies WorkerOperationHandlers; diff --git a/src/agents/mcp-oauth-store.ts b/src/agents/mcp-oauth-store.ts index ee14e5fda347..0d095bb5c88c 100644 --- a/src/agents/mcp-oauth-store.ts +++ b/src/agents/mcp-oauth-store.ts @@ -2,8 +2,8 @@ import { createSqliteWorkerWriteAdmission } from "../infra/sqlite-worker-store.js"; import { executeExistingOpenClawStateRead } from "../state/openclaw-state-db-readonly.js"; import type { OpenClawStateAsyncLeaseContext } from "../state/openclaw-state-lease-context.js"; +import { runWithOpenClawStateLeaseWorker } from "../state/openclaw-state-lease-worker-operation.js"; import type { OpenClawStateLeaseWorkerAuthority } from "../state/openclaw-state-lease-worker-owner.js"; -import { runWithOpenClawStateLeaseWorker } from "../state/openclaw-state-lease-worker-storage.js"; import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.js"; import type { OpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.types.js"; import { diff --git a/src/agents/mcp-oauth-store.types.ts b/src/agents/mcp-oauth-store.types.ts index d621b9e970d0..7a4a456cee83 100644 --- a/src/agents/mcp-oauth-store.types.ts +++ b/src/agents/mcp-oauth-store.types.ts @@ -3,7 +3,7 @@ import type { OAuthClientInformationMixed, OAuthTokens, } from "@modelcontextprotocol/sdk/shared/auth.js"; -import type { OpenClawStateLeaseIdentity } from "../state/openclaw-state-lease-store.js"; +import type { OpenClawStateLeaseIdentity } from "../state/openclaw-state-lease.types.js"; type McpOAuthAuthorizationChallenge = { resourceMetadataUrl?: string; diff --git a/src/infra/deferred-plugin-migrations.ts b/src/infra/deferred-plugin-migrations.ts index 8bfad88c6746..e077de920b98 100644 --- a/src/infra/deferred-plugin-migrations.ts +++ b/src/infra/deferred-plugin-migrations.ts @@ -211,7 +211,7 @@ export function readDeferredPluginMigrations( /** Prime resolution in this module before the updater replaces its package. */ export async function prepareDeferredPluginMigrationRuntime(): Promise { const { prepareOpenClawStateLeaseWorkerRuntime } = - await import("../state/openclaw-state-lease-worker-storage.js"); + await import("../state/openclaw-state-lease-worker-operation.js"); await Promise.all([ import("../state/openclaw-state-worker-store.js"), import("../plugins/plugin-lifecycle-lease.js"), @@ -488,7 +488,7 @@ export async function recordDeferredPluginMigrations( }); const { withPluginLifecycleLease } = await import("../plugins/plugin-lifecycle-lease.js"); const { runWithOpenClawStateLeaseWorker } = - await import("../state/openclaw-state-lease-worker-storage.js"); + await import("../state/openclaw-state-lease-worker-operation.js"); return withPluginLifecycleLease({ env: params.env }, async (lease) => { const context = captureOpenClawStateWorkerContext({ env: params.env, diff --git a/src/infra/deferred-plugin-migrations.worker.ts b/src/infra/deferred-plugin-migrations.worker.ts deleted file mode 100644 index 945dc251101a..000000000000 --- a/src/infra/deferred-plugin-migrations.worker.ts +++ /dev/null @@ -1,55 +0,0 @@ -import { ZodError } from "zod"; -import { - runOpenClawStateWriteTransaction, - type OpenClawStateDatabaseOptions, -} from "../state/openclaw-state-db.js"; -import { assertOpenClawStateLeaseWorkerOwnedInTransaction } from "../state/openclaw-state-lease-worker.js"; -import type { OpenClawStateWorkerOperations } from "../state/openclaw-state-worker-contract.js"; -import { - DeferredPluginMigrationConflictError, - recordDeferredPluginMigrationsInTransaction, - readDeferredPluginMigrationCompletions, - readDeferredPluginMigrations, -} from "./deferred-plugin-migrations.js"; -import type { SqliteWorkerCommand } from "./sqlite-worker-contract.js"; - -export function readDeferredPluginMigrationsInWorker( - command: Extract< - SqliteWorkerCommand, - { type: "plugins.deferredMigrations.read" | "plugins.deferredMigrations.completions.read" } - >, - options: { path: string; env: NodeJS.ProcessEnv }, -) { - return command.type === "plugins.deferredMigrations.read" - ? readDeferredPluginMigrations({ - ...options, - artifactPreservingReadOnly: command.input.artifactPreservingReadOnly, - }) - : readDeferredPluginMigrationCompletions(options); -} - -export function recordDeferredPluginMigrationsInWorker( - input: OpenClawStateWorkerOperations["plugins.deferredMigrations.record"]["input"], - options: OpenClawStateDatabaseOptions, -): OpenClawStateWorkerOperations["plugins.deferredMigrations.record"]["output"] { - try { - return runOpenClawStateWriteTransaction( - ({ db }) => { - assertOpenClawStateLeaseWorkerOwnedInTransaction(db, input.identity); - const transitions = recordDeferredPluginMigrationsInTransaction(db, input); - assertOpenClawStateLeaseWorkerOwnedInTransaction(db, input.identity, "write", "commit"); - return { kind: "recorded" as const, transitions }; - }, - options, - { operationLabel: "state.plugin-migration-deferral" }, - ); - } catch (error) { - if (error instanceof DeferredPluginMigrationConflictError) { - return { kind: "conflict", pending: error.pending }; - } - if (error instanceof ZodError) { - return { kind: "invalid", issues: error.issues }; - } - throw error; - } -} diff --git a/src/plugins/official-external-plugin-catalog-snapshot-store.worker-contract.ts b/src/plugins/official-external-plugin-catalog-snapshot-store.worker-contract.ts deleted file mode 100644 index 102ce52c72cd..000000000000 --- a/src/plugins/official-external-plugin-catalog-snapshot-store.worker-contract.ts +++ /dev/null @@ -1,12 +0,0 @@ -import type { HostedOfficialExternalPluginCatalogSnapshot } from "./official-external-plugin-catalog.types.js"; - -export type HostedCatalogSnapshotWorkerOperations = { - "plugins.catalogSnapshot.read": { - input: { url: string }; - output: HostedOfficialExternalPluginCatalogSnapshot | null; - }; - "plugins.catalogSnapshot.write": { - input: { snapshot: HostedOfficialExternalPluginCatalogSnapshot; now: number }; - output: { ok: true } | { ok: false; message: string }; - }; -}; diff --git a/src/plugins/state.worker-contract.ts b/src/plugins/state.worker-contract.ts new file mode 100644 index 000000000000..f0e52824fb64 --- /dev/null +++ b/src/plugins/state.worker-contract.ts @@ -0,0 +1,4 @@ +import type { WorkerOperations } from "../state/worker-operation-registry.js"; +import type { pluginRuntimeOperations } from "./state.worker.js"; + +export type PluginRuntimeWorkerOperations = WorkerOperations; diff --git a/src/plugins/state.worker.ts b/src/plugins/state.worker.ts new file mode 100644 index 000000000000..8c970dab4cd4 --- /dev/null +++ b/src/plugins/state.worker.ts @@ -0,0 +1,104 @@ +import { ZodError } from "zod"; +import { + DeferredPluginMigrationConflictError, + readDeferredPluginMigrationCompletions, + readDeferredPluginMigrations, + recordDeferredPluginMigrationsInTransaction, + type DeferredPluginMigrationRecordInput, +} from "../infra/deferred-plugin-migrations.js"; +import { runOpenClawStateWriteTransaction } from "../state/openclaw-state-db.js"; +import { assertOpenClawStateLeaseWorkerOwnedInTransaction } from "../state/openclaw-state-lease-worker.js"; +import type { OpenClawStateLeaseIdentity } from "../state/openclaw-state-lease.types.js"; +import type { WorkerOperationHandlers } from "../state/worker-operation-registry.js"; +import { + readPluginBindingApprovalsInDatabase, + upsertPluginBindingApprovalInDatabase, +} from "./conversation-binding-state.kernel.js"; +import type { PluginBindingApprovalEntry } from "./conversation-binding-state.types.js"; +import { publishPluginSourceAdmissionInDatabase } from "./installed-plugin-index-store-write.js"; +import { + readHostedCatalogSnapshotInDatabase, + writeHostedCatalogSnapshotInDatabase, +} from "./official-external-plugin-catalog-snapshot-store.kernel.js"; +import { HostedCatalogSignedFeedMonotonicityError } from "./official-external-plugin-catalog-source.js"; +import type { HostedOfficialExternalPluginCatalogSnapshot } from "./official-external-plugin-catalog.types.js"; +import type { PluginSourceAdmissionPublication } from "./plugin-source-admission.types.js"; + +export const pluginRuntimeOperations = { + "plugins.conversationBindingApprovals.read": (_input: undefined, { open }) => + readPluginBindingApprovalsInDatabase(open().db), + "plugins.conversationBindingApprovals.upsert": ( + input: PluginBindingApprovalEntry, + { open, stateOptions }, + ) => + runOpenClawStateWriteTransaction(({ db }) => upsertPluginBindingApprovalInDatabase(db, input), { + database: open(), + ...stateOptions(), + }), + "plugins.catalogSnapshot.read": (input: { url: string }, { open }) => + readHostedCatalogSnapshotInDatabase(open().db, input.url), + "plugins.catalogSnapshot.write": ( + input: { snapshot: HostedOfficialExternalPluginCatalogSnapshot; now: number }, + { open, stateOptions }, + ) => { + const options = { database: open(), ...stateOptions() }; + try { + runOpenClawStateWriteTransaction( + ({ db }) => writeHostedCatalogSnapshotInDatabase(db, input.snapshot, input.now), + options, + ); + return { ok: true as const }; + } catch (error) { + if (error instanceof HostedCatalogSignedFeedMonotonicityError) { + return { ok: false as const, message: error.message }; + } + throw error; + } + }, + "plugins.metadata.sourceAdmission.publish": ( + input: PluginSourceAdmissionPublication, + { open, stateOptions }, + ) => + runOpenClawStateWriteTransaction( + ({ db }) => publishPluginSourceAdmissionInDatabase(db, input), + { database: open(), ...stateOptions() }, + ), + "plugins.deferredMigrations.record": ( + input: Omit & { + identity: OpenClawStateLeaseIdentity; + }, + { open, stateOptions }, + ) => { + const options = { database: open(), ...stateOptions() }; + try { + return runOpenClawStateWriteTransaction( + ({ db }) => { + assertOpenClawStateLeaseWorkerOwnedInTransaction(db, input.identity); + const transitions = recordDeferredPluginMigrationsInTransaction(db, input); + assertOpenClawStateLeaseWorkerOwnedInTransaction(db, input.identity, "write", "commit"); + return { kind: "recorded" as const, transitions }; + }, + options, + { operationLabel: "state.plugin-migration-deferral" }, + ); + } catch (error) { + if (error instanceof DeferredPluginMigrationConflictError) { + return { kind: "conflict" as const, pending: error.pending }; + } + if (error instanceof ZodError) { + return { kind: "invalid" as const, issues: error.issues }; + } + throw error; + } + }, + "plugins.deferredMigrations.read": ( + input: { artifactPreservingReadOnly: boolean }, + { stateOptions }, + ) => + readDeferredPluginMigrations({ + ...stateOptions(), + artifactPreservingReadOnly: input.artifactPreservingReadOnly, + }), + "plugins.deferredMigrations.completions.read": (_input: undefined, { stateOptions }) => + readDeferredPluginMigrationCompletions(stateOptions()), +} satisfies WorkerOperationHandlers; diff --git a/src/projects/project-clone-registration.test.ts b/src/projects/project-clone-registration.test.ts index 855437fb2b2d..901d245d5b61 100644 --- a/src/projects/project-clone-registration.test.ts +++ b/src/projects/project-clone-registration.test.ts @@ -59,7 +59,7 @@ vi.mock("./project-clone-runtime.js", () => ({ cloneProjectCheckout: fixture.clone, })); -vi.mock("../state/openclaw-state-lease-worker-storage.js", () => ({ +vi.mock("../state/openclaw-state-lease-worker-operation.js", () => ({ runWithOpenClawStateLeaseWorker: async ( _lease: OpenClawStateLeaseContext, _context: unknown, diff --git a/src/projects/project-environment.test.ts b/src/projects/project-environment.test.ts index bc7f96a1a23a..3ccc368b2b8e 100644 --- a/src/projects/project-environment.test.ts +++ b/src/projects/project-environment.test.ts @@ -50,7 +50,7 @@ vi.mock("../state/openclaw-state-lease.js", () => ({ vi.mock("../state/openclaw-state-worker-store.js", () => ({ executeOpenClawStateWorker: mocks.execute, })); -vi.mock("../state/openclaw-state-lease-worker-storage.js", () => ({ +vi.mock("../state/openclaw-state-lease-worker-operation.js", () => ({ runWithOpenClawStateLeaseWorker: async ( _lease: unknown, context: unknown, diff --git a/src/projects/project-registration.ts b/src/projects/project-registration.ts index c665e1d6c2df..bd7611931d80 100644 --- a/src/projects/project-registration.ts +++ b/src/projects/project-registration.ts @@ -54,7 +54,7 @@ export async function registerPreparedProjectRegistry( ); } const { runWithOpenClawStateLeaseWorker } = - await import("../state/openclaw-state-lease-worker-storage.js"); + await import("../state/openclaw-state-lease-worker-operation.js"); return await runWithOpenClawStateLeaseWorker(lease, context, async (scope, identity) => { const project = await scope.execute({ type: "projects.insert", diff --git a/src/projects/project-registry.ts b/src/projects/project-registry.ts index 45191ecfe4eb..25e60ab3778b 100644 --- a/src/projects/project-registry.ts +++ b/src/projects/project-registry.ts @@ -243,7 +243,7 @@ export async function resolveProjectCloneRefreshOwner( context: OpenClawStateWorkerContext, ): Promise { const { runWithOpenClawStateLeaseWorker } = - await import("../state/openclaw-state-lease-worker-storage.js"); + await import("../state/openclaw-state-lease-worker-operation.js"); return await runWithOpenClawStateLeaseWorker(lease, context, (scope, identity) => scope.execute({ type: "projects.resolveRefreshOwner", @@ -285,7 +285,7 @@ export async function removeProjectRegistry( { path: context.admission.databasePath, env }, async (lease) => { const { runWithOpenClawStateLeaseWorker } = - await import("../state/openclaw-state-lease-worker-storage.js"); + await import("../state/openclaw-state-lease-worker-operation.js"); return await runWithOpenClawStateLeaseWorker(lease, context, (scope, identity) => scope.execute({ type: "projects.remove", diff --git a/src/projects/project-registry.worker-contract.ts b/src/projects/project-registry.worker-contract.ts index d59fba135aff..2d96f9069f56 100644 --- a/src/projects/project-registry.worker-contract.ts +++ b/src/projects/project-registry.worker-contract.ts @@ -1,4 +1,4 @@ -import type { OpenClawStateLeaseIdentity } from "../state/openclaw-state-lease-store.js"; +import type { OpenClawStateLeaseIdentity } from "../state/openclaw-state-lease.types.js"; import type { ProjectRegistryIdentity, ProjectRegistryInsert, diff --git a/src/skills/lifecycle/upload-store.ts b/src/skills/lifecycle/upload-store.ts index a3d881fd81b5..a06e403e26ad 100644 --- a/src/skills/lifecycle/upload-store.ts +++ b/src/skills/lifecycle/upload-store.ts @@ -18,7 +18,7 @@ import { resolveSkillUploadDatabaseOptions, type SkillUploadMetadataRow, } from "./upload-store.sqlite.js"; -import type { SkillUploadWorkerOperations } from "./upload-store.worker.js"; +import type { SkillUploadWorkerOperations } from "./upload-store.worker-contract.js"; type SkillUploadScope = Pick, "execute">; /** Time window in which uploaded skill archive chunks may be committed. */ diff --git a/src/skills/lifecycle/upload-store.worker-contract.ts b/src/skills/lifecycle/upload-store.worker-contract.ts new file mode 100644 index 000000000000..aa27288aa03d --- /dev/null +++ b/src/skills/lifecycle/upload-store.worker-contract.ts @@ -0,0 +1,4 @@ +import type { WorkerOperations } from "../../state/worker-operation-registry.js"; +import type { skillUploadOperations } from "./upload-store.worker.js"; + +export type SkillUploadWorkerOperations = WorkerOperations; diff --git a/src/skills/lifecycle/upload-store.worker.ts b/src/skills/lifecycle/upload-store.worker.ts index a7c00520e544..f8163f960dd5 100644 --- a/src/skills/lifecycle/upload-store.worker.ts +++ b/src/skills/lifecycle/upload-store.worker.ts @@ -1,9 +1,5 @@ -import type { SqliteWorkerCommand } from "../../infra/sqlite-worker-contract.js"; import { requestSqliteWorkerOperationAdmission } from "../../infra/sqlite-worker-operation-admission.js"; -import type { - OpenClawStateDatabase, - OpenClawStateDatabaseOptions, -} from "../../state/openclaw-state-db.js"; +import type { WorkerOperationHandlers } from "../../state/worker-operation-registry.js"; import { commitSkillUploadInDatabase } from "./upload-store-commit.js"; import { appendSkillUploadChunkInDatabase, @@ -15,69 +11,42 @@ import { } from "./upload-store.kernel.js"; import { deleteOwnedSkillUpload, renewSkillUploadInstallLease } from "./upload-store.sqlite.js"; -type Operation unknown> = { - input: Parameters[0]; - output: ReturnType; -}; -export type SkillUploadWorkerOperations = { - "skillUploads.begin": Operation; - "skillUploads.chunk": Operation; - "skillUploads.commit": Operation; - "skillUploads.expired": Operation; - "skillUploads.deleteExpired": Operation; - "skillUploads.claim": Operation; - "skillUploads.renew": { - input: Omit[0], "options">; - output: boolean; - }; - "skillUploads.consume": { - input: { uploadId: string; owner: string }; - output: ReturnType; - }; - "skillUploads.release": Operation; -}; +const admit = (stage: "transaction" | "commit") => + requestSqliteWorkerOperationAdmission({ stage, facts: undefined }); -export function isSkillUploadCommand(command: { - type: string; - input: unknown; -}): command is SqliteWorkerCommand { - return ( - command.type === "skillUploads.begin" || - command.type === "skillUploads.chunk" || - command.type === "skillUploads.commit" || - command.type === "skillUploads.expired" || - command.type === "skillUploads.deleteExpired" || - command.type === "skillUploads.claim" || - command.type === "skillUploads.renew" || - command.type === "skillUploads.consume" || - command.type === "skillUploads.release" - ); -} - -export function executeSkillUploadCommand( - command: SqliteWorkerCommand, - options: OpenClawStateDatabaseOptions & { database: OpenClawStateDatabase }, -): SkillUploadWorkerOperations[keyof SkillUploadWorkerOperations]["output"] { - const admit = (stage: "transaction" | "commit") => - requestSqliteWorkerOperationAdmission({ stage, facts: undefined }); - switch (command.type) { - case "skillUploads.begin": - return beginSkillUploadInDatabase(command.input, options, admit); - case "skillUploads.chunk": - return appendSkillUploadChunkInDatabase(command.input, options, admit); - case "skillUploads.commit": - return commitSkillUploadInDatabase(command.input, options, admit); - case "skillUploads.expired": - return listExpiredSkillUploadsInDatabase(command.input, options); - case "skillUploads.deleteExpired": - return deleteExpiredSkillUploadInDatabase(command.input, options); - case "skillUploads.claim": - return claimSkillUploadInDatabase(command.input, options); - case "skillUploads.renew": - return renewSkillUploadInstallLease({ ...command.input, options }); - case "skillUploads.consume": - return deleteOwnedSkillUpload(command.input.uploadId, command.input.owner, options); - case "skillUploads.release": - return releaseSkillUploadInDatabase(command.input, options); - } -} +export const skillUploadOperations = { + "skillUploads.begin": ( + input: Parameters[0], + { open, stateOptions }, + ) => beginSkillUploadInDatabase(input, { database: open(), ...stateOptions() }, admit), + "skillUploads.chunk": ( + input: Parameters[0], + { open, stateOptions }, + ) => appendSkillUploadChunkInDatabase(input, { database: open(), ...stateOptions() }, admit), + "skillUploads.commit": ( + input: Parameters[0], + { open, stateOptions }, + ) => commitSkillUploadInDatabase(input, { database: open(), ...stateOptions() }, admit), + "skillUploads.expired": ( + input: Parameters[0], + { open, stateOptions }, + ) => listExpiredSkillUploadsInDatabase(input, { database: open(), ...stateOptions() }), + "skillUploads.deleteExpired": ( + input: Parameters[0], + { open, stateOptions }, + ) => deleteExpiredSkillUploadInDatabase(input, { database: open(), ...stateOptions() }), + "skillUploads.claim": ( + input: Parameters[0], + { open, stateOptions }, + ) => claimSkillUploadInDatabase(input, { database: open(), ...stateOptions() }), + "skillUploads.renew": ( + input: Omit[0], "options">, + { open, stateOptions }, + ) => renewSkillUploadInstallLease({ ...input, options: { database: open(), ...stateOptions() } }), + "skillUploads.consume": (input: { uploadId: string; owner: string }, { open, stateOptions }) => + deleteOwnedSkillUpload(input.uploadId, input.owner, { database: open(), ...stateOptions() }), + "skillUploads.release": ( + input: Parameters[0], + { open, stateOptions }, + ) => releaseSkillUploadInDatabase(input, { database: open(), ...stateOptions() }), +} satisfies WorkerOperationHandlers; diff --git a/src/skills/workshop/store-client.ts b/src/skills/workshop/store-client.ts index c005dff7cd73..ba3f6a7d5243 100644 --- a/src/skills/workshop/store-client.ts +++ b/src/skills/workshop/store-client.ts @@ -1,8 +1,8 @@ import { cloneEnvWithPlatformSemantics } from "../../config/config-env-vars.js"; import type { SqliteWorkerStore } from "../../infra/sqlite-worker-contract.js"; import { createSqliteWorkerWriteAdmission } from "../../infra/sqlite-worker-store.js"; -import type { OpenClawStateLeaseIdentity } from "../../state/openclaw-state-lease-store.js"; -import { runWithOpenClawStateLeasesWorker } from "../../state/openclaw-state-lease-worker-storage.js"; +import { runWithOpenClawStateLeasesWorker } from "../../state/openclaw-state-lease-worker-operation.js"; +import type { OpenClawStateLeaseIdentity } from "../../state/openclaw-state-lease.types.js"; import { captureOpenClawStateWorkerContext } from "../../state/openclaw-state-worker-context.js"; import { runOpenClawStateWorkerOperation } from "../../state/openclaw-state-worker-store.js"; import { resolveSkillWorkshopStateDir } from "./proposal-generation.js"; diff --git a/src/skills/workshop/store.worker-contract.ts b/src/skills/workshop/store.worker-contract.ts index c1463a4de080..a9f3efefb352 100644 --- a/src/skills/workshop/store.worker-contract.ts +++ b/src/skills/workshop/store.worker-contract.ts @@ -1,69 +1,10 @@ -import type { OpenClawStateLeaseIdentity } from "../../state/openclaw-state-lease-store.js"; -import type { - RecordSkillExperienceReviewOutcomeInput, - ReadSkillCollectionBackupDropsInput, - SkillCollectionReviewOutcome, -} from "./collection-review.kernel.js"; -import type { RecordSkillProposalEvaluationInput } from "./store-evaluation.kernel.js"; -import type { - CreateSkillProposalInput, - ImportLegacySkillProposalInput, - ListStoredSkillProposalsInput, - UpdateSkillProposalRecordInput, -} from "./store-proposal.kernel.js"; -import type { StoredSkillProposal } from "./store-sqlite-record.js"; -import type { - ClearSkillProposalRollbackInput, - WriteSkillProposalRollbackInput, -} from "./store-sqlite-rollback.js"; -import type { - CommitPendingSkillProposalTransitionInput, - PendingSkillProposalTransitionCommit, - ReadCommittedSkillProposalTransitionInput, -} from "./store-sqlite-transition.js"; -import type { SkillProposalEvent, SkillProposalRecord, SkillProposalRollback } from "./types.js"; +import type { WorkerOperations } from "../../state/worker-operation-registry.js"; +import type { skillWorkshopOperations, skillCuratorOperations } from "./store.worker.js"; -type WorkshopOperation = { - input: { - value: Input; - agentId?: string; - leaseIdentities?: readonly OpenClawStateLeaseIdentity[]; - }; - output: Output; -}; +export type SkillWorkshopWorkerOperations = WorkerOperations; +export type SkillWorkshopExecutionOperations = Omit< + SkillWorkshopWorkerOperations, + "workshop.events.list" +>; -export type SkillWorkshopExecutionOperations = { - "workshop.schema.ensure": WorkshopOperation; - "workshop.proposal.read": WorkshopOperation; - "workshop.proposals.list": WorkshopOperation< - ListStoredSkillProposalsInput, - StoredSkillProposal[] - >; - "workshop.proposal.create": WorkshopOperation; - "workshop.proposal.update": WorkshopOperation< - UpdateSkillProposalRecordInput, - SkillProposalEvent | undefined - >; - "workshop.proposal.import": WorkshopOperation< - ImportLegacySkillProposalInput, - "imported" | "already-imported" - >; - "workshop.proposal.evaluate": WorkshopOperation< - RecordSkillProposalEvaluationInput, - { record: SkillProposalRecord; event: SkillProposalEvent } - >; - "workshop.transition.commit": WorkshopOperation< - CommitPendingSkillProposalTransitionInput & { operationLabel: string }, - PendingSkillProposalTransitionCommit - >; - "workshop.transition.committed": WorkshopOperation< - ReadCommittedSkillProposalTransitionInput, - Extract | null - >; - "workshop.rollback.read": WorkshopOperation; - "workshop.rollback.write": WorkshopOperation; - "workshop.rollback.clear": WorkshopOperation; - "workshop.collection.list": WorkshopOperation; - "workshop.collection.drops": WorkshopOperation>; - "workshop.experience.record": WorkshopOperation; -}; +export type SkillCuratorOperations = WorkerOperations; diff --git a/src/skills/workshop/store.worker.ts b/src/skills/workshop/store.worker.ts index aa0b74ce2f31..4a3df012578f 100644 --- a/src/skills/workshop/store.worker.ts +++ b/src/skills/workshop/store.worker.ts @@ -1,11 +1,13 @@ import type { DatabaseSync } from "node:sqlite"; -import type { SqliteWorkerCommand } from "../../infra/sqlite-worker-contract.js"; import { requestSqliteWorkerOperationAdmission } from "../../infra/sqlite-worker-operation-admission.js"; -import { getSqliteWorkerStateContext } from "../../infra/sqlite-worker-state-context.js"; import type { OpenClawStateDatabase } from "../../state/openclaw-state-db-contract.js"; import { runOpenClawStateWriteTransaction } from "../../state/openclaw-state-db.js"; import { assertOpenClawStateLeasesWorkerOwnedInTransaction } from "../../state/openclaw-state-lease-worker.js"; -import type { OpenClawStateWorkerOperations } from "../../state/openclaw-state-worker-contract.js"; +import type { OpenClawStateLeaseIdentity } from "../../state/openclaw-state-lease.types.js"; +import type { + WorkerOperationContext, + WorkerOperationHandlers, +} from "../../state/worker-operation-registry.js"; import { listSkillCollectionReviewOutcomesInDatabase, readSkillCollectionBackupDropsInDatabase, @@ -32,162 +34,197 @@ import { commitPendingSkillProposalTransitionInDatabase, readCommittedSkillProposalTransitionInDatabase, } from "./store-sqlite-transition.js"; -import type { SkillWorkshopExecutionOperations } from "./store.worker-contract.js"; -type Operations = SkillWorkshopExecutionOperations & - Pick< - OpenClawStateWorkerOperations, - "skills.curator.read" | "skills.usage.record" | "workshop.events.list" - >; +type WorkshopInput = { + value: Value; + agentId?: string; + leaseIdentities?: readonly OpenClawStateLeaseIdentity[]; +}; -export function isSkillWorkshopCommand(command: { - type: string; - input: unknown; -}): command is SqliteWorkerCommand { - switch (command.type) { - case "skills.curator.read": - case "skills.usage.record": - case "workshop.events.list": - case "workshop.schema.ensure": - case "workshop.proposal.read": - case "workshop.proposals.list": - case "workshop.proposal.create": - case "workshop.proposal.update": - case "workshop.proposal.import": - case "workshop.proposal.evaluate": - case "workshop.transition.commit": - case "workshop.transition.committed": - case "workshop.rollback.read": - case "workshop.rollback.write": - case "workshop.rollback.clear": - case "workshop.collection.list": - case "workshop.collection.drops": - case "workshop.experience.record": - return true; - default: - return false; +function assertWrite( + db: DatabaseSync, + input: WorkshopInput, + stage: "transaction" | "commit", +) { + if (input.leaseIdentities) { + assertOpenClawStateLeasesWorkerOwnedInTransaction(db, input.leaseIdentities, stage); + } else { + requestSqliteWorkerOperationAdmission({ stage, facts: undefined }); } } -export function executeSkillWorkshopCommand( - command: SqliteWorkerCommand, - database: OpenClawStateDatabase, - databasePath: string, -): Operations[keyof Operations]["output"] { - if (command.type === "skills.curator.read") { - return readSkillCuratorStateInDatabase(database, command.input.skillFiles); - } - const options = { - database, - path: databasePath, - env: getSqliteWorkerStateContext().environment, - }; - if (command.type === "skills.usage.record") { - return runOpenClawStateWriteTransaction( - (current) => recordSkillUsageInDatabase(current, command.input), - options, - ); - } - if (command.type === "workshop.events.list") { - ensureSkillWorkshopSchemaInDatabase(database, options); - return listStoredSkillProposalEventsInDatabase(database.db, command.input); - } - const { leaseIdentities } = command.input; - const assertWrite = (db: DatabaseSync, stage: "transaction" | "commit") => { - if (leaseIdentities) { - assertOpenClawStateLeasesWorkerOwnedInTransaction(db, leaseIdentities, stage); - } else { - requestSqliteWorkerOperationAdmission({ stage, facts: undefined }); - } - }; - if (command.type === "workshop.schema.ensure") { - return ensureSkillWorkshopSchemaInDatabase(database, options, assertWrite); - } - const write = (operationLabel: string, operation: () => T, proposalId?: string): T => - runOpenClawStateWriteTransaction( - ({ db }) => { - assertWrite(db, "transaction"); - if (leaseIdentities) { - const stored = proposalId ? readStoredProposalInDatabase(db, proposalId) : null; - const agentId = command.input.agentId; - for (const identity of leaseIdentities) { - if ( - !agentId || - (identity.scope === "skill-collection" - ? identity.key !== agentId - : identity.scope !== "skill-workshop-target" || - !stored || - identity.key !== - `${agentId}:${hashSkillProposalContent(stored.record.target.skillFile)}`) || - (stored !== null && - stored.row.owner_agent_id !== null && - stored.row.owner_agent_id !== agentId) - ) { - throw new Error("Skill Workshop lease does not match its proposal target."); - } +function write( + input: WorkshopInput, + { open, stateOptions }: WorkerOperationContext, + operationLabel: string, + operation: (database: OpenClawStateDatabase) => T, + proposalId?: string, +): T { + const database = open(); + return runOpenClawStateWriteTransaction( + ({ db }) => { + assertWrite(db, input, "transaction"); + if (input.leaseIdentities) { + const stored = proposalId ? readStoredProposalInDatabase(db, proposalId) : null; + const { agentId } = input; + for (const identity of input.leaseIdentities) { + if ( + !agentId || + (identity.scope === "skill-collection" + ? identity.key !== agentId + : identity.scope !== "skill-workshop-target" || + !stored || + identity.key !== + `${agentId}:${hashSkillProposalContent(stored.record.target.skillFile)}`) || + (stored !== null && + stored.row.owner_agent_id !== null && + stored.row.owner_agent_id !== agentId) + ) { + throw new Error("Skill Workshop lease does not match its proposal target."); } } - const result = operation(); - assertWrite(db, "commit"); - return result; - }, - options, - { operationLabel }, - ); - switch (command.type) { - case "workshop.proposal.read": - return readStoredProposalInDatabase(database.db, command.input.value); - case "workshop.proposals.list": - return listStoredSkillProposalsInDatabase(database.db, command.input.value); - case "workshop.proposal.create": - return write("skill-workshop.proposal.create", () => - createSkillProposalInDatabase(database.db, command.input.value), - ); - case "workshop.proposal.update": - return write( - "skill-workshop.proposal.update", - () => updateSkillProposalRecordInDatabase(database.db, command.input.value), - command.input.value.record.id, - ); - case "workshop.proposal.import": - return write("doctor.skill-workshop.import", () => - importLegacySkillProposalInDatabase(database.db, command.input.value), - ); - case "workshop.proposal.evaluate": - return write( - "skill-workshop.proposal.evaluate", - () => recordSkillProposalEvaluationInDatabase(database.db, command.input.value), - command.input.value.proposalId, - ); - case "workshop.transition.commit": - return write( - command.input.value.operationLabel, - () => commitPendingSkillProposalTransitionInDatabase(database.db, command.input.value), - command.input.value.expected.id, - ); - case "workshop.transition.committed": - return readCommittedSkillProposalTransitionInDatabase(database.db, command.input.value); - case "workshop.rollback.read": - return readSkillProposalRollbackInDatabase(database.db, command.input.value); - case "workshop.rollback.write": - return write( - "skill-workshop.rollback.write", - () => writeSkillProposalRollbackInDatabase(database.db, command.input.value), - command.input.value.proposalId, - ); - case "workshop.rollback.clear": - return write( - "skill-workshop.rollback.clear", - () => clearSkillProposalRollbackInDatabase(database.db, command.input.value), - command.input.value.proposalId, - ); - case "workshop.collection.list": - return listSkillCollectionReviewOutcomesInDatabase(database.db, command.input.value); - case "workshop.collection.drops": - return readSkillCollectionBackupDropsInDatabase(database.db, command.input.value); - case "workshop.experience.record": - return write("skill-workshop.experience.record", () => - recordSkillExperienceReviewOutcomeInDatabase(database, command.input.value), - ); - } + } + const result = operation(database); + assertWrite(db, input, "commit"); + return result; + }, + { database, ...stateOptions() }, + { operationLabel }, + ); } + +export const skillCuratorOperations = { + "skills.curator.read": (input: { skillFiles: readonly string[] }, { open }) => + readSkillCuratorStateInDatabase(open(), input.skillFiles), + "skills.usage.record": ( + input: Parameters[1], + { open, stateOptions }, + ) => + runOpenClawStateWriteTransaction((current) => recordSkillUsageInDatabase(current, input), { + database: open(), + ...stateOptions(), + }), +} satisfies WorkerOperationHandlers; + +export const skillWorkshopOperations = { + "workshop.events.list": ( + input: Parameters[1], + { open, stateOptions }, + ) => { + const database = open(); + ensureSkillWorkshopSchemaInDatabase(database, { database, ...stateOptions() }); + return listStoredSkillProposalEventsInDatabase(database.db, input); + }, + "workshop.schema.ensure": (input: WorkshopInput, { open, stateOptions }) => { + const database = open(); + return ensureSkillWorkshopSchemaInDatabase( + database, + { database, ...stateOptions() }, + (db, stage) => assertWrite(db, input, stage), + ); + }, + "workshop.proposal.read": ( + input: WorkshopInput[1]>, + { open }, + ) => readStoredProposalInDatabase(open().db, input.value), + "workshop.proposals.list": ( + input: WorkshopInput[1]>, + { open }, + ) => listStoredSkillProposalsInDatabase(open().db, input.value), + "workshop.proposal.create": ( + input: WorkshopInput[1]>, + context, + ) => + write(input, context, "skill-workshop.proposal.create", (database) => + createSkillProposalInDatabase(database.db, input.value), + ), + "workshop.proposal.update": ( + input: WorkshopInput[1]>, + context, + ) => + write( + input, + context, + "skill-workshop.proposal.update", + (database) => updateSkillProposalRecordInDatabase(database.db, input.value), + input.value.record.id, + ), + "workshop.proposal.import": ( + input: WorkshopInput[1]>, + context, + ) => + write(input, context, "doctor.skill-workshop.import", (database) => + importLegacySkillProposalInDatabase(database.db, input.value), + ), + "workshop.proposal.evaluate": ( + input: WorkshopInput[1]>, + context, + ) => + write( + input, + context, + "skill-workshop.proposal.evaluate", + (database) => recordSkillProposalEvaluationInDatabase(database.db, input.value), + input.value.proposalId, + ), + "workshop.transition.commit": ( + input: WorkshopInput< + Parameters[1] & { + operationLabel: string; + } + >, + context, + ) => + write( + input, + context, + input.value.operationLabel, + (database) => commitPendingSkillProposalTransitionInDatabase(database.db, input.value), + input.value.expected.id, + ), + "workshop.transition.committed": ( + input: WorkshopInput[1]>, + { open }, + ) => readCommittedSkillProposalTransitionInDatabase(open().db, input.value), + "workshop.rollback.read": ( + input: WorkshopInput[1]>, + { open }, + ) => readSkillProposalRollbackInDatabase(open().db, input.value), + "workshop.rollback.write": ( + input: WorkshopInput[1]>, + context, + ) => + write( + input, + context, + "skill-workshop.rollback.write", + (database) => writeSkillProposalRollbackInDatabase(database.db, input.value), + input.value.proposalId, + ), + "workshop.rollback.clear": ( + input: WorkshopInput[1]>, + context, + ) => + write( + input, + context, + "skill-workshop.rollback.clear", + (database) => clearSkillProposalRollbackInDatabase(database.db, input.value), + input.value.proposalId, + ), + "workshop.collection.list": ( + input: WorkshopInput[1]>, + { open }, + ) => listSkillCollectionReviewOutcomesInDatabase(open().db, input.value), + "workshop.collection.drops": ( + input: WorkshopInput[1]>, + { open }, + ) => readSkillCollectionBackupDropsInDatabase(open().db, input.value), + "workshop.experience.record": ( + input: WorkshopInput[1]>, + context, + ) => + write(input, context, "skill-workshop.experience.record", (database) => + recordSkillExperienceReviewOutcomeInDatabase(database, input.value), + ), +} satisfies WorkerOperationHandlers; diff --git a/src/state/openclaw-state-lease-acquisition.ts b/src/state/openclaw-state-lease-acquisition.ts index 4dd8e2e98bc6..a7657872ad23 100644 --- a/src/state/openclaw-state-lease-acquisition.ts +++ b/src/state/openclaw-state-lease-acquisition.ts @@ -11,7 +11,7 @@ import { OpenClawStateLeaseError, } from "./openclaw-state-lease-error.js"; import { STATE_LEASE_WRITE_BACKOFF } from "./openclaw-state-lease-storage.js"; -import type { OpenClawStateLeaseAcquisition } from "./openclaw-state-lease-store.js"; +import type { OpenClawStateLeaseAcquisition } from "./openclaw-state-lease.types.js"; const log = createSubsystemLogger("state/lease"); diff --git a/src/state/openclaw-state-lease-async.maintenance.test.ts b/src/state/openclaw-state-lease-async.maintenance.test.ts index e0eeb8073bad..d2f0584523d0 100644 --- a/src/state/openclaw-state-lease-async.maintenance.test.ts +++ b/src/state/openclaw-state-lease-async.maintenance.test.ts @@ -7,9 +7,9 @@ import { } from "./openclaw-state-db-async-lifecycle.js"; import type { OpenClawStateAsyncLeaseContext } from "./openclaw-state-lease-context.js"; import type { LeaseHeartbeatCleanup } from "./openclaw-state-lease-heartbeat.js"; -import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease-store.js"; import { withOpenClawStateLeaseWorkerAdmission } from "./openclaw-state-lease-worker-owner.js"; import { withOpenClawStateLeaseAsync } from "./openclaw-state-lease.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context.types.js"; type CreateStorage = diff --git a/src/state/openclaw-state-lease-context.ts b/src/state/openclaw-state-lease-context.ts index 8200e8142ea7..f795784bca1f 100644 --- a/src/state/openclaw-state-lease-context.ts +++ b/src/state/openclaw-state-lease-context.ts @@ -2,7 +2,7 @@ import type { DatabaseSync } from "node:sqlite"; import type { OpenClawStateLeaseAcquisition, OpenClawStateLeaseIdentity, -} from "./openclaw-state-lease-store.js"; +} from "./openclaw-state-lease.types.js"; export type OpenClawStateLeaseContext = { signal: AbortSignal; diff --git a/src/state/openclaw-state-lease-heartbeat-shared.ts b/src/state/openclaw-state-lease-heartbeat-shared.ts index cb9d79e9f32c..f9f3cd9ee3c9 100644 --- a/src/state/openclaw-state-lease-heartbeat-shared.ts +++ b/src/state/openclaw-state-lease-heartbeat-shared.ts @@ -1,5 +1,5 @@ import type { StateLeaseProcessOwner } from "../infra/state-lease-process-owner.js"; -import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease-store.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; import type { OpenClawStateWorkerErrorPayload } from "./openclaw-state-worker-error.js"; // Allow headroom over observed 38 s cold Gateway boots under load; committed lease expiry still bounds startup. diff --git a/src/state/openclaw-state-lease-storage.ts b/src/state/openclaw-state-lease-storage.ts index 169012fa2f71..8c070f87dba0 100644 --- a/src/state/openclaw-state-lease-storage.ts +++ b/src/state/openclaw-state-lease-storage.ts @@ -25,10 +25,9 @@ import { releaseOpenClawStateLeaseInTransaction, renewOpenClawStateLeaseInTransaction, } from "./openclaw-state-lease-store.js"; -import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease-store.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; import { OPENCLAW_STATE_SCHEMA_SQL } from "./openclaw-state-schema.js"; import { captureOpenClawStateWorkerContext } from "./openclaw-state-worker-context.js"; -import { runOpenClawStateWorkerOperation } from "./openclaw-state-worker-store.js"; export type OpenClawStateLeaseDatabase = { scope: "shared"; @@ -98,6 +97,8 @@ export async function acquireLease( }); } }; + const { runOpenClawStateWorkerOperation } = await import("./openclaw-state-worker-store.js"); + assertAdmission(); const result = await runOpenClawStateWorkerOperation( context, (scope) => diff --git a/src/state/openclaw-state-lease-store.ts b/src/state/openclaw-state-lease-store.ts index fcc31ee5001e..46103385bee7 100644 --- a/src/state/openclaw-state-lease-store.ts +++ b/src/state/openclaw-state-lease-store.ts @@ -11,11 +11,11 @@ import { type StateLeaseProcessOwner, } from "../infra/state-lease-process-owner.js"; import type { DB } from "./openclaw-state-db.generated.js"; +import type { + OpenClawStateLeaseIdentity, + OpenClawStateLeaseAcquisition, +} from "./openclaw-state-lease.types.js"; -export type OpenClawStateLeaseIdentity = { scope: string; key: string; owner: string }; -export type OpenClawStateLeaseAcquisition = - | { kind: "acquired"; expiresAt: number } - | { kind: "held"; holder: { owner: string; epoch: number; expiresAt: number | null } }; type LeaseDatabase = Pick; /** The caller owns the write transaction; only absent or expired leases can be acquired. */ diff --git a/src/state/openclaw-state-lease-worker-group.test.ts b/src/state/openclaw-state-lease-worker-group.test.ts index afdcc524789b..d69e3908df62 100644 --- a/src/state/openclaw-state-lease-worker-group.test.ts +++ b/src/state/openclaw-state-lease-worker-group.test.ts @@ -5,13 +5,13 @@ import type { SqliteWorkerOperationSettlement } from "../infra/sqlite-worker-ope import { createDeferredCore } from "../shared/deferred.js"; import { createOpenClawDatabaseMaintenanceScope } from "./openclaw-state-db-async-lifecycle.js"; import type { OpenClawStateAsyncLeaseContext } from "./openclaw-state-lease-context.js"; -import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease-store.js"; +import { runWithOpenClawStateLeasesWorker } from "./openclaw-state-lease-worker-operation.js"; import { createOpenClawStateLeaseWorkerOwner, withOpenClawStateLeaseWorkerAdmission, withOpenClawStateLeasesWorkerAdmission, } from "./openclaw-state-lease-worker-owner.js"; -import { runWithOpenClawStateLeasesWorker } from "./openclaw-state-lease-worker-storage.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context.types.js"; const runWorkerOperation = vi.hoisted(() => vi.fn()); diff --git a/src/state/openclaw-state-lease-worker-operation.ts b/src/state/openclaw-state-lease-worker-operation.ts new file mode 100644 index 000000000000..1a9d563abe65 --- /dev/null +++ b/src/state/openclaw-state-lease-worker-operation.ts @@ -0,0 +1,66 @@ +import type { OpenClawStateWorkerLeaseContext } from "./openclaw-state-lease-context.js"; +import { + withOpenClawStateLeaseWorkerAdmission, + withOpenClawStateLeasesWorkerAdmission, + type OpenClawStateLeaseWorkerAuthority, +} from "./openclaw-state-lease-worker-owner.js"; +import { prepareOpenClawStateLeaseStorageRuntime } from "./openclaw-state-lease-worker-storage.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; +import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context.types.js"; +import type { DomainScope } from "./openclaw-state-worker-store.types.js"; + +/** Both importers must resolve their transport before package replacement. */ +export async function prepareOpenClawStateLeaseWorkerRuntime(): Promise { + await Promise.all([ + prepareOpenClawStateLeaseStorageRuntime(), + import("./openclaw-state-worker-store.js"), + ]); +} + +/** Retain the actual lease until every admitted worker transaction has settled. */ +export function runWithOpenClawStateLeaseWorker( + lease: OpenClawStateWorkerLeaseContext, + context: OpenClawStateWorkerContext, + operation: (scope: DomainScope, identity: OpenClawStateLeaseIdentity) => Promise, + authority?: OpenClawStateLeaseWorkerAuthority, +): Promise { + return withOpenClawStateLeaseWorkerAdmission( + lease, + context.admission.databasePath, + async (admission) => { + const { runOpenClawStateWorkerOperation } = await import("./openclaw-state-worker-store.js"); + return runOpenClawStateWorkerOperation( + context, + (scope) => operation(scope, admission.identity), + { assertCurrent: admission.assertCurrent, createAdmission: admission.createAdmission }, + ); + }, + authority, + ); +} + +/** Share one actor operation while every original lease retains its native settlement. */ +export function runWithOpenClawStateLeasesWorker( + leases: readonly OpenClawStateWorkerLeaseContext[], + context: OpenClawStateWorkerContext, + operation: (scope: DomainScope, identities: readonly OpenClawStateLeaseIdentity[]) => Promise, + authority?: OpenClawStateLeaseWorkerAuthority, +): Promise { + return withOpenClawStateLeasesWorkerAdmission( + leases, + context, + async (admission) => { + const { runOpenClawStateWorkerOperation } = await import("./openclaw-state-worker-store.js"); + admission.assertCurrent(); + return runOpenClawStateWorkerOperation( + context, + (scope) => operation(scope, admission.identities), + { + assertCurrent: admission.assertCurrent, + createAdmission: admission.createAdmission, + }, + ); + }, + authority, + ); +} diff --git a/src/state/openclaw-state-lease-worker-owner.ts b/src/state/openclaw-state-lease-worker-owner.ts index b068549639a4..c233a1e30d01 100644 --- a/src/state/openclaw-state-lease-worker-owner.ts +++ b/src/state/openclaw-state-lease-worker-owner.ts @@ -16,7 +16,7 @@ import { createDeferredCore } from "../shared/deferred.js"; import { resolveGlobalSingleton } from "../shared/global-singleton.js"; import type { OpenClawStateWorkerLeaseContext } from "./openclaw-state-lease-context.js"; import { OpenClawStateLeaseError } from "./openclaw-state-lease-error.js"; -import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease-store.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context.types.js"; export type OpenClawStateLeaseWorkerPurpose = "write" | "acquire" | "verify" | "renew" | "release"; diff --git a/src/state/openclaw-state-lease-worker-storage.ts b/src/state/openclaw-state-lease-worker-storage.ts index 2728a107f8dd..cc235b1ae673 100644 --- a/src/state/openclaw-state-lease-worker-storage.ts +++ b/src/state/openclaw-state-lease-worker-storage.ts @@ -1,31 +1,27 @@ import { throwSqliteLifecycleErrors } from "../infra/sqlite-lifecycle-errors.js"; import type { SqliteWorkerStore } from "../infra/sqlite-worker-store.js"; -import type { OpenClawStateWorkerLeaseContext } from "./openclaw-state-lease-context.js"; +import type { OpenClawStateLeaseLifecycleOperations } from "./openclaw-state-lease-context.js"; import { OpenClawStateLeaseError } from "./openclaw-state-lease-error.js"; import { leaseHeartbeatState } from "./openclaw-state-lease-heartbeat-shared.js"; import { startOpenClawStateLeaseTimer } from "./openclaw-state-lease-heartbeat.js"; +import type { + createOpenClawStateLeaseWorkerOwner, + WorkerLeaseScope, +} from "./openclaw-state-lease-worker-owner.js"; import type { OpenClawStateLeaseAcquisition, OpenClawStateLeaseIdentity, -} from "./openclaw-state-lease-store.js"; -import { - withOpenClawStateLeaseWorkerAdmission, - withOpenClawStateLeasesWorkerAdmission, - type createOpenClawStateLeaseWorkerOwner, - type OpenClawStateLeaseWorkerAuthority, - type WorkerLeaseScope, -} from "./openclaw-state-lease-worker-owner.js"; +} from "./openclaw-state-lease.types.js"; import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context.types.js"; -import type { OpenClawStateWorkerOperations } from "./openclaw-state-worker-contract.js"; type LeaseWorkerOwner = ReturnType; type LeaseWorkerOperation = ( - scope: Pick, "execute">, + scope: Pick, "execute">, identity: OpenClawStateLeaseIdentity, ) => Promise; /** Keep this importer's admission and release code available across package replacement. */ -export async function prepareOpenClawStateLeaseWorkerRuntime(): Promise { +export async function prepareOpenClawStateLeaseStorageRuntime(): Promise { await Promise.all([ import("./openclaw-state-worker-store.js"), import("../infra/sqlite-worker-identity.js"), @@ -221,47 +217,3 @@ export function createOpenClawStateLeaseWorkerStorage( }; return storage; } - -/** Retain the actual lease until every admitted worker transaction has settled. */ -export function runWithOpenClawStateLeaseWorker( - lease: OpenClawStateWorkerLeaseContext, - context: OpenClawStateWorkerContext, - operation: LeaseWorkerOperation, - authority?: OpenClawStateLeaseWorkerAuthority, -): Promise { - return withOpenClawStateLeaseWorkerAdmission( - lease, - context.admission.databasePath, - admittedWorkerOperation(context, operation), - authority, - ); -} - -/** Share one actor operation while every original lease retains its native settlement. */ -export function runWithOpenClawStateLeasesWorker( - leases: readonly OpenClawStateWorkerLeaseContext[], - context: OpenClawStateWorkerContext, - operation: ( - scope: Pick, "execute">, - identities: readonly OpenClawStateLeaseIdentity[], - ) => Promise, - authority?: OpenClawStateLeaseWorkerAuthority, -): Promise { - return withOpenClawStateLeasesWorkerAdmission( - leases, - context, - async (admission) => { - const { runOpenClawStateWorkerOperation } = await import("./openclaw-state-worker-store.js"); - admission.assertCurrent(); - return runOpenClawStateWorkerOperation( - context, - (scope) => operation(scope, admission.identities), - { - assertCurrent: admission.assertCurrent, - createAdmission: admission.createAdmission, - }, - ); - }, - authority, - ); -} diff --git a/src/state/openclaw-state-lease-worker.transaction-group.test.ts b/src/state/openclaw-state-lease-worker.transaction-group.test.ts index 6f77c0e8fc36..b4bcc471cf2e 100644 --- a/src/state/openclaw-state-lease-worker.transaction-group.test.ts +++ b/src/state/openclaw-state-lease-worker.transaction-group.test.ts @@ -2,8 +2,8 @@ import { DatabaseSync } from "node:sqlite"; import { afterEach, describe, expect, it, vi } from "vitest"; import { disposeNodeSqliteDependents } from "../infra/kysely-sync-cache-state.js"; import { requestSqliteWorkerOperationAdmission } from "../infra/sqlite-worker-operation-admission.js"; -import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease-store.js"; import { assertOpenClawStateLeasesWorkerOwnedInTransaction } from "./openclaw-state-lease-worker.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; vi.mock("../infra/sqlite-worker-operation-admission.js", () => ({ requestSqliteWorkerOperationAdmission: vi.fn(), diff --git a/src/state/openclaw-state-lease-worker.ts b/src/state/openclaw-state-lease-worker.ts index 20aff93bb049..bd9934a1d7da 100644 --- a/src/state/openclaw-state-lease-worker.ts +++ b/src/state/openclaw-state-lease-worker.ts @@ -32,8 +32,8 @@ import { reclaimDeadOpenClawStateLeaseInTransaction, releaseOpenClawStateLeaseInTransaction, renewOpenClawStateLeaseInTransaction, - type OpenClawStateLeaseIdentity, } from "./openclaw-state-lease-store.js"; +import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease.types.js"; function takeLeaseExpiryObservation(identity: OpenClawStateLeaseIdentity): BigInt64Array { const attachment = takeSqliteWorkerOperationAdmissionAttachment(); diff --git a/src/state/openclaw-state-lease.types.ts b/src/state/openclaw-state-lease.types.ts new file mode 100644 index 000000000000..f69f89d2f20e --- /dev/null +++ b/src/state/openclaw-state-lease.types.ts @@ -0,0 +1,4 @@ +export type OpenClawStateLeaseIdentity = { scope: string; key: string; owner: string }; +export type OpenClawStateLeaseAcquisition = + | { kind: "acquired"; expiresAt: number } + | { kind: "held"; holder: { owner: string; epoch: number; expiresAt: number | null } }; diff --git a/src/state/openclaw-state-read-worker.mcp.test.ts b/src/state/openclaw-state-read-worker.mcp.test.ts index 56e79d94503a..741c3787b50d 100644 --- a/src/state/openclaw-state-read-worker.mcp.test.ts +++ b/src/state/openclaw-state-read-worker.mcp.test.ts @@ -2,7 +2,7 @@ // oxfmt-ignore import { emptyReply, queueTask, source } from "./openclaw-state-read-worker.test-harness.js"; import { expect, it } from "vitest"; -import type { AcpSessionReadInput } from "../acp/runtime/session-meta-keys.js"; +import type { AcpSessionReadInput } from "../acp/runtime/session-meta-read.types.js"; import { createDeferredCore } from "../shared/deferred.js"; import { captureOpenClawStateReadSource } from "./openclaw-state-read-worker.js"; import { captureOpenClawStateWorkerContext } from "./openclaw-state-worker-context.js"; diff --git a/src/state/openclaw-state-read.types.ts b/src/state/openclaw-state-read.types.ts index 8ebff6b8ab79..beb1a9bf54da 100644 --- a/src/state/openclaw-state-read.types.ts +++ b/src/state/openclaw-state-read.types.ts @@ -1,6 +1,6 @@ import type { DatabaseSync } from "node:sqlite"; import type { Selectable } from "kysely"; -import type { AcpSessionReadInput, AcpSessionRow } from "../acp/runtime/session-meta-keys.js"; +import type { AcpSessionReadInput, AcpSessionRow } from "../acp/runtime/session-meta-read.types.js"; import type { McpOAuthReadOnlyOperations } from "../agents/mcp-oauth-store.kernel.js"; import type { SandboxBrowserRegistryEntry, diff --git a/src/state/openclaw-state-worker-contract.ts b/src/state/openclaw-state-worker-contract.ts index 5b5202caa662..2812ceebbd09 100644 --- a/src/state/openclaw-state-worker-contract.ts +++ b/src/state/openclaw-state-worker-contract.ts @@ -1,6 +1,3 @@ -import type { ZodIssue } from "zod"; -import type { AcpSessionWriteOperations } from "../acp/runtime/session-meta-write.types.js"; -import type { AuthProfileRowRead, UserModelAuthProfile } from "../agents/auth-profiles/types.js"; import type { NativeHookRelayStoreWorkerOperations } from "../agents/harness/native-hook-relay-store.worker-contract.js"; import type { McpOAuthReadOperations } from "../agents/mcp-oauth-store.kernel.js"; import type { McpOAuthWriteOperations } from "../agents/mcp-oauth-store.types.js"; @@ -41,12 +38,6 @@ import type { PlacementSessionToolWorkerOperations } from "../gateway/worker-env import type { PlacementTurnClaimWorkerOperations } from "../gateway/worker-environments/placement-turn-claims.worker-contract.js"; import type { WorkspaceJournalWorkerOperations } from "../gateway/worker-environments/placement-workspace-journal.worker-contract.js"; import type { WorkerEnvironmentWorkerOperations } from "../gateway/worker-environments/store-worker-contract.js"; -import type { - DeferredPluginMigration, - DeferredPluginMigrationRecordInput, - recordDeferredPluginMigrationsInTransaction, - readDeferredPluginMigrationCompletions, -} from "../infra/deferred-plugin-migrations.js"; import type { DeliveryQueueWorkerOperations } from "../infra/delivery-queue.worker-contract.js"; import type * as deviceAuth from "../infra/device-auth-store.kernel.js"; import type { DeviceIdentity } from "../infra/device-identity-store.js"; @@ -75,10 +66,7 @@ import type { readRemoteModelCatalog } from "../model-catalog/remote-store.js"; import type { NodeWorkerJournalWorkerOperations } from "../node-host/node-worker-journal.worker-contract.js"; import type { PluginBlobWorkerOperations } from "../plugin-state/plugin-blob-worker-contract.js"; import type { PluginStateWorkerOperations } from "../plugin-state/plugin-state-worker-contract.js"; -import type { PluginBindingApprovalEntry } from "../plugins/conversation-binding-state.types.js"; import type { PluginMetadataStateSelector } from "../plugins/installed-plugin-index-row.js"; -import type { HostedCatalogSnapshotWorkerOperations } from "../plugins/official-external-plugin-catalog-snapshot-store.worker-contract.js"; -import type { PluginSourceAdmissionPublication } from "../plugins/plugin-source-admission.types.js"; import type { ProjectRegistryWorkerOperations } from "../projects/project-registry.worker-contract.js"; import type { CaptureWorkerOperations } from "../proxy-capture/store.worker-contract.js"; import type { SecretStoreConfigRefWrite } from "../secrets/store/secret-store-config-ref.kernel.js"; @@ -87,15 +75,9 @@ import type { SessionStateWorkerOperations } from "../sessions/session-state-eve import type { SessionUpstreamLink } from "../sessions/session-upstream-links.kernel.js"; import type { SessionUpstreamWorkerOperations } from "../sessions/session-upstream-links.worker-contract.js"; import type { DeviceAuthEntry } from "../shared/device-auth.js"; -import type { SkillUploadWorkerOperations } from "../skills/lifecycle/upload-store.worker.js"; -import type * as curator from "../skills/workshop/curator.kernel.js"; -import type { listStoredSkillProposalEventsInDatabase } from "../skills/workshop/store-sqlite-event.js"; -import type { SkillWorkshopExecutionOperations } from "../skills/workshop/store.worker-contract.js"; +import type { SkillUploadWorkerOperations } from "../skills/lifecycle/upload-store.worker-contract.js"; import type { SkillProposalEvent, SkillProposalRecord } from "../skills/workshop/types.js"; -import type { - TranscriptReadOperations, - TranscriptWriteOperations, -} from "../transcripts/store-worker-contract.js"; +import type { TranscriptReadOperations } from "../transcripts/store-worker-contract.js"; import type { TuiLastSessionWorkerOperations } from "../tui/tui-last-session.contract.js"; import type { AgentProvenance } from "./agent-provenance.types.js"; import type { PreparedBackupRunRecord } from "./backup-run-records.kernel.js"; @@ -106,7 +88,6 @@ import type { import type { OnboardingRecommendationWriteOperations } from "./onboarding-recommendations.contract.js"; import type { OpenClawAgentDatabaseWorkerLeaseReceipt } from "./openclaw-agent-db-lease.js"; import type { OpenClawStateLeaseLifecycleOperations } from "./openclaw-state-lease-context.js"; -import type { OpenClawStateLeaseIdentity } from "./openclaw-state-lease-store.js"; import type { RegisteredStateWorkerOperations } from "./openclaw-state-worker-registry.js"; import type { RepositoryWorkspaceWorkerOperations } from "./session-repository-workspaces.types.js"; import type { UserPreferenceWorkerOperations } from "./user-preferences.types.js"; @@ -121,11 +102,9 @@ export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations & RepositoryWorkspaceWorkerOperations & CaptureWorkerOperations & TuiLastSessionWorkerOperations & - AcpSessionWriteOperations & SessionStateWorkerOperations & SessionUpstreamWorkerOperations & McpOAuthReadOperations & - SkillWorkshopExecutionOperations & CurrentConversationBindingWorkerOperations & McpOAuthWriteOperations & LegacyMcpOAuthWorkerOperations & @@ -135,7 +114,6 @@ export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations & AuditWriterOperations & NativeHookRelayStoreWorkerOperations & TelemetryWorkerOperations & - HostedCatalogSnapshotWorkerOperations & PluginStateWorkerOperations & PluginBlobWorkerOperations & UserPreferenceWorkerOperations & @@ -153,9 +131,7 @@ export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations & SessionDeliveryWorkerOperations & DeliveryQueueWorkerOperations & TranscriptReadOperations & - TranscriptWriteOperations & NodeWorkerJournalWorkerOperations & - SkillUploadWorkerOperations & OpenClawStateLeaseLifecycleOperations & ManagedImageRecordWorkerOperations & { "database.walMaintenance": { input: SqliteWalPeriodicRequest; output: SqliteWalPeriodicResult }; @@ -224,12 +200,6 @@ export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations & output: ReturnType; }; - "authProfiles.read": { input: { artifactPreserving: boolean }; output: AuthProfileRowRead }; - "authProfiles.sharedOwnership": { input: { artifactPreserving: boolean }; output: unknown }; - "authProfiles.personal": { - input: { profileId: string; artifactPreserving: boolean }; - output: UserModelAuthProfile | undefined; - }; "agentProvenance.readBatch": { input: { agentIds: readonly string[] }; output: AgentProvenance[]; @@ -253,15 +223,6 @@ export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations & input: SessionGroupCatalogMutation; output: SessionGroupCatalogMutationResult; }; - "skills.curator.read": { - input: { skillFiles: readonly string[] }; - output: ReturnType; - }; - "skills.usage.record": { input: curator.PreparedSkillUsage; output: void }; - "workshop.events.list": { - input: Parameters[1]; - output: ReturnType; - }; "doctor.workshopMigrationRecords.read": { input: { includeEvents: boolean }; output: @@ -275,42 +236,10 @@ export type OpenClawStateWorkerOperations = RegisteredStateWorkerOperations & input: { artifactPreservingReadOnly: boolean }; output: ReturnType; }; - "plugins.conversationBindingApprovals.read": { - input: undefined; - output: PluginBindingApprovalEntry[]; - }; - "plugins.conversationBindingApprovals.upsert": { - input: PluginBindingApprovalEntry; - output: void; - }; "plugins.metadata.read": { input: { selector: PluginMetadataStateSelector; artifactPreservingReadOnly?: boolean }; output: { value_json: string } | undefined; }; - "plugins.metadata.sourceAdmission.publish": { - input: PluginSourceAdmissionPublication; - output: boolean; - }; - "plugins.deferredMigrations.record": { - input: Omit & { - identity: OpenClawStateLeaseIdentity; - }; - output: - | { - kind: "recorded"; - transitions: ReturnType; - } - | { kind: "conflict"; pending: readonly DeferredPluginMigration[] } - | { kind: "invalid"; issues: ZodIssue[] }; - }; - "plugins.deferredMigrations.read": { - input: { artifactPreservingReadOnly: boolean }; - output: readonly DeferredPluginMigration[]; - }; - "plugins.deferredMigrations.completions.read": { - input: undefined; - output: ReturnType; - }; "claws.install-schema-versions": { input: { artifactPreservingReadOnly: boolean }; output: ClawInstallSchemaVersionRow[] | undefined; diff --git a/src/state/openclaw-state-worker-registry.ts b/src/state/openclaw-state-worker-registry.ts index c12c23739c14..318946ac9d1c 100644 --- a/src/state/openclaw-state-worker-registry.ts +++ b/src/state/openclaw-state-worker-registry.ts @@ -1,15 +1,43 @@ +import type { AcpSessionWriteOperations } from "../acp/runtime/session-meta-write.worker-contract.js"; +import type { AuthProfileWorkerOperations } from "../agents/auth-profiles/store.worker-contract.js"; import type { WorktreeWorkerOperations } from "../agents/worktrees/dispatch.worker.js"; import type { FleetRegistryWriteOperations } from "../fleet/registry.worker-contract.js"; import type { ApnsRegistrationWorkerOperations } from "../infra/push-apns-store.worker-contract.js"; import type { WebPushWorkerOperations } from "../infra/push-web-store.worker-contract.js"; +import type { PluginRuntimeWorkerOperations } from "../plugins/state.worker-contract.js"; +import type { SkillUploadWorkerOperations } from "../skills/lifecycle/upload-store.worker-contract.js"; +import type { + SkillWorkshopWorkerOperations, + SkillCuratorOperations, +} from "../skills/workshop/store.worker-contract.js"; +import type { TranscriptWriteOperations } from "../transcripts/store-write.worker-contract.js"; import { createWorkerOperationRegistry } from "./worker-operation-registry.js"; export type RegisteredStateWorkerOperations = WebPushWorkerOperations & ApnsRegistrationWorkerOperations & WorktreeWorkerOperations & - FleetRegistryWriteOperations; + FleetRegistryWriteOperations & + AcpSessionWriteOperations & + SkillUploadWorkerOperations & + SkillWorkshopWorkerOperations & + SkillCuratorOperations & + TranscriptWriteOperations & + AuthProfileWorkerOperations & + PluginRuntimeWorkerOperations; export const stateWorkerRegistry = createWorkerOperationRegistry({ + authProfiles: () => + import("../agents/auth-profiles/store.worker.js").then((m) => m.authProfileOperations), + plugins: () => import("../plugins/state.worker.js").then((m) => m.pluginRuntimeOperations), + acp: () => + import("../acp/runtime/session-meta-write.worker.js").then((m) => m.acpSessionOperations), + skillUploads: () => + import("../skills/lifecycle/upload-store.worker.js").then((m) => m.skillUploadOperations), + workshop: () => + import("../skills/workshop/store.worker.js").then((m) => m.skillWorkshopOperations), + skills: () => import("../skills/workshop/store.worker.js").then((m) => m.skillCuratorOperations), + transcripts: () => + import("../transcripts/store-worker-write.js").then((m) => m.transcriptWriteOperations), webPush: () => import("../infra/push-web-store.worker.js").then((m) => m.webPushOperations), apns: () => import("../infra/push-apns-store.worker.js").then((m) => m.apnsOperations), worktrees: () => diff --git a/src/state/openclaw-state-worker-runtime.ts b/src/state/openclaw-state-worker-runtime.ts index 30b24cf3c10b..56cbb6cd4d53 100644 --- a/src/state/openclaw-state-worker-runtime.ts +++ b/src/state/openclaw-state-worker-runtime.ts @@ -1,10 +1,3 @@ -import { executeAcpSessionMutationInWorker } from "../acp/runtime/session-meta-write.worker.js"; -import { - readAuthProfileRows, - SHARED_AUTH_STORE_STATE_KEY, -} from "../agents/auth-profiles/sqlite-json.js"; -import { isMissingDatabasePath } from "../agents/auth-profiles/sqlite-read-pool.js"; -import type { AuthProfileRowRead } from "../agents/auth-profiles/types.js"; import { readNativeHookRelayBridgeSnapshotFromDatabase, listNativeHookRelayBridgeSnapshotsInDatabase, @@ -58,10 +51,6 @@ import { isWorkspaceJournalWriteCommand } from "../gateway/worker-environments/p import { executeWorkspaceJournalCommand } from "../gateway/worker-environments/placement-workspace-journal.worker.js"; import { isWorkerEnvironmentCommand } from "../gateway/worker-environments/store-worker-contract.js"; import { executeWorkerEnvironmentCommand } from "../gateway/worker-environments/store.worker.js"; -import { - readDeferredPluginMigrationsInWorker, - recordDeferredPluginMigrationsInWorker, -} from "../infra/deferred-plugin-migrations.worker.js"; import * as deliveryQueue from "../infra/delivery-queue.worker.js"; import * as deviceAuth from "../infra/device-auth-store.kernel.js"; import { executeDevicePairingMutationInWorker } from "../infra/device-pairing-dispatch.worker.js"; @@ -95,15 +84,6 @@ import { isNodeWorkerJournalCommand } from "../node-host/node-worker-journal.wor import { executeNodeWorkerJournalCommand } from "../node-host/node-worker-journal.worker.js"; import { executePluginBlobCommand } from "../plugin-state/plugin-blob-store.worker.js"; import { isPluginBlobWorkerCommand } from "../plugin-state/plugin-blob-worker-contract.js"; -import { - readPluginBindingApprovalsInDatabase, - upsertPluginBindingApprovalInDatabase, -} from "../plugins/conversation-binding-state.kernel.js"; -import { - readHostedCatalogSnapshotInDatabase, - writeHostedCatalogSnapshotInDatabase, -} from "../plugins/official-external-plugin-catalog-snapshot-store.kernel.js"; -import { HostedCatalogSignedFeedMonotonicityError } from "../plugins/official-external-plugin-catalog-source.js"; import { executeProjectRegistryCommand, isProjectRegistryCommand, @@ -113,17 +93,7 @@ import { purgeExpiredSecretStoreEntriesInDatabase } from "../secrets/store/secre import { executeSessionStateCommand } from "../sessions/session-state-events.worker.js"; import { listWatchedSessionUpstreamLinksInDatabase } from "../sessions/session-upstream-links.kernel.js"; import { executeSessionUpstreamCommand } from "../sessions/session-upstream-links.worker.js"; -import { createLazyRuntimeModule } from "../shared/lazy-runtime.js"; -import { - isSkillUploadCommand, - executeSkillUploadCommand, -} from "../skills/lifecycle/upload-store.worker.js"; -import * as skillWorkshop from "../skills/workshop/store.worker.js"; import { executeTranscriptRead } from "../transcripts/store-worker-read.js"; -import { - executeTranscriptWrite, - isTranscriptWriteCommand, -} from "../transcripts/store-worker-write.js"; import { clearRetiredTuiPointers } from "../tui/tui-last-session.kernel.js"; import { listAgentProvenanceInDatabase, @@ -132,7 +102,6 @@ import { import { ensureAgentProvenanceSchema } from "./agent-provenance.schema.js"; import { recordBackupRunInDatabase } from "./backup-run-records.kernel.js"; import { writeConfigMachineState } from "./config-machine-state-write.js"; -import { readConfigMachineState } from "./config-machine-state.js"; import { deletePersonalGitHubSessionReceiptsInDatabase, readSessionReceiptDeletionIdentitiesInDatabase, @@ -157,23 +126,12 @@ import { executeRepositoryWorkspaceCommand, isRepositoryWorkspaceCommand, } from "./session-repository-workspaces.worker.js"; -import { readUserModelAuthProfile } from "./user-model-accounts.js"; import { executeUserPreferenceCommand } from "./user-preferences.worker.js"; import { executeUserProfileCommand, isUserProfileCommand } from "./user-profiles.worker.js"; const log = createSubsystemLogger("state/worker"); -const loadPluginIndexWriter = createLazyRuntimeModule( - () => import("../plugins/installed-plugin-index-store-write.js"), -); -let pluginIndexWriter: Awaited> | undefined; - export function prepareSharedStateCommand(type: PropertyKey): Promise | undefined { - if (type === "plugins.metadata.sourceAdmission.publish" && !pluginIndexWriter) { - return loadPluginIndexWriter().then((loaded) => { - pluginIndexWriter = loaded; - }); - } return stateWorkerRegistry.prepare(type) ?? prepareCronStateWorkerCommand(type); } @@ -238,43 +196,6 @@ export function executeSharedStateCommand( if (command.type === "audit.writer.process" || command.type === "audit.writer.prune") { return executeAuditWriterCommand(command, stateOptions(), open); } - if ( - command.type === "authProfiles.read" || - command.type === "authProfiles.sharedOwnership" || - command.type === "authProfiles.personal" - ) { - const read = () => { - const options = stateOptions(); - if (command.type === "authProfiles.sharedOwnership") { - return readConfigMachineState(SHARED_AUTH_STORE_STATE_KEY, options); - } - if (command.type === "authProfiles.personal") { - return readUserModelAuthProfile(command.input.profileId, options); - } - const missing: AuthProfileRowRead = { - store: { status: "missing", reason: "database" }, - state: { status: "missing", reason: "database" }, - cacheable: false, - }; - try { - return ( - withExistingOpenClawStateDatabaseReadOnly( - ({ db }) => readAuthProfileRows(db, context.databasePath, "shared-state"), - options, - ) ?? missing - ); - } catch { - return isMissingDatabasePath(context.databasePath) - ? missing - : { - store: { status: "unreadable" as const }, - state: { status: "unreadable" as const }, - cacheable: false, - }; - } - }; - return command.input.artifactPreserving ? withArtifactPreservingStateReads(read) : read(); - } if (command.type === "promotions.markNotified" || command.type === "promotions.recordClaim") { return executePromotionCommand(command, stateOptions(), open); } @@ -314,19 +235,6 @@ export function executeSharedStateCommand( ? withArtifactPreservingStateReads(read) : read(); } - if (command.type === "acp.prepareMutation" || command.type === "acp.commitMutation") { - return executeAcpSessionMutationInWorker(open(), command); - } - if (command.type === "plugins.conversationBindingApprovals.read") { - return readPluginBindingApprovalsInDatabase(open().db); - } - if (command.type === "plugins.conversationBindingApprovals.upsert") { - const database = open(); - return runOpenClawStateWriteTransaction( - ({ db }) => upsertPluginBindingApprovalInDatabase(db, command.input), - { database, path: context.databasePath, env: getSqliteWorkerStateContext().environment }, - ); - } if (command.type === "updateRuns.recordStep" || command.type === "updateRuns.recordPhase") { return recordUpdateRunMutationInWorker(command, stateOptions(), (stage) => requestSqliteWorkerOperationAdmission({ stage, facts: undefined }), @@ -342,12 +250,6 @@ export function executeSharedStateCommand( requestSqliteWorkerOperationAdmission({ stage, facts: undefined }), ); } - if ( - command.type === "plugins.deferredMigrations.read" || - command.type === "plugins.deferredMigrations.completions.read" - ) { - return readDeferredPluginMigrationsInWorker(command, stateOptions()); - } if (command.type === "claws.install-schema-versions") { const read = command.input.artifactPreservingReadOnly ? withExistingOpenClawStateDatabaseArtifactPreservingReadOnly @@ -431,15 +333,9 @@ export function executeSharedStateCommand( if (command.type === "githubRepository.personalPending") { return readPendingRepositoryGitHubPublicationInDatabase(database.db, command.input); } - if (skillWorkshop.isSkillWorkshopCommand(command)) { - return skillWorkshop.executeSkillWorkshopCommand(command, database, context.databasePath); - } if (command.type === "deviceAuth.list") { return deviceAuth.readDeviceAuthTokensFromDatabase(database.db, command.input); } - if (isTranscriptWriteCommand(command)) { - return executeTranscriptWrite(command, { database, path: context.databasePath }); - } switch (command.type) { case "transcripts.canonicalSessionRow": case "transcripts.readEntries": @@ -467,9 +363,6 @@ export function executeSharedStateCommand( if (isManagedImageRecordCommand(command)) { return executeManagedImageRecordCommand(command, database); } - if (command.type === "plugins.catalogSnapshot.read") { - return readHostedCatalogSnapshotInDatabase(database.db, command.input.url); - } if (command.type === "nativeHookRelay.listSnapshots") { return listNativeHookRelayBridgeSnapshotsInDatabase(database); } @@ -486,9 +379,6 @@ export function executeSharedStateCommand( database, ...stateOptions(), }; - if (command.type === "plugins.deferredMigrations.record") { - return recordDeferredPluginMigrationsInWorker(command.input, writeOptions); - } if ( command.type === "nativeHookRelay.write" || command.type === "nativeHookRelay.renew" || @@ -538,9 +428,6 @@ export function executeSharedStateCommand( if (deliveryQueue.isDeliveryQueueCommand(command)) { return deliveryQueue.executeDeliveryQueueCommand(command, writeOptions); } - if (isSkillUploadCommand(command)) { - return executeSkillUploadCommand(command, writeOptions); - } if ( command.type === "deviceAuth.store" || command.type === "deviceAuth.storeOrigin" || @@ -581,31 +468,6 @@ export function executeSharedStateCommand( if (command.type === "sessionState.record" || command.type === "sessionState.prune") { return executeSessionStateCommand(command, writeOptions); } - if (command.type === "plugins.catalogSnapshot.write") { - try { - runOpenClawStateWriteTransaction( - ({ db }) => - writeHostedCatalogSnapshotInDatabase(db, command.input.snapshot, command.input.now), - writeOptions, - ); - return { ok: true }; - } catch (error) { - if (error instanceof HostedCatalogSignedFeedMonotonicityError) { - return { ok: false, message: error.message }; - } - throw error; - } - } - if (command.type === "plugins.metadata.sourceAdmission.publish") { - if (!pluginIndexWriter) { - throw new Error("Plugin source admission writer is not prepared"); - } - const { publishPluginSourceAdmissionInDatabase } = pluginIndexWriter; - return runOpenClawStateWriteTransaction( - ({ db }) => publishPluginSourceAdmissionInDatabase(db, command.input), - writeOptions, - ); - } if (command.type === "subagents.persistChanges") { const { writeId, values, deleteRunIds } = command.input; let committed = false; diff --git a/src/transcripts/capture-appends.ts b/src/transcripts/capture-appends.ts index cfe490b24fa6..d16530d9ac80 100644 --- a/src/transcripts/capture-appends.ts +++ b/src/transcripts/capture-appends.ts @@ -1,5 +1,5 @@ import { createDeferredCore } from "../shared/deferred.js"; -import type { TranscriptAppendScheduler } from "./store-worker-contract.js"; +import type { TranscriptAppendScheduler } from "./store-worker.types.js"; type AppendOutcome = { ok: true } | { ok: false; error: unknown }; diff --git a/src/transcripts/store-worker-client.ts b/src/transcripts/store-worker-client.ts index 1ceee48be6e7..4d835aec1d9a 100644 --- a/src/transcripts/store-worker-client.ts +++ b/src/transcripts/store-worker-client.ts @@ -1,17 +1,15 @@ import { createSqliteWorkerWriteAdmission } from "../infra/sqlite-worker-store.js"; import type { OpenClawStateDatabaseOptions } from "../state/openclaw-state-db.js"; import type { OpenClawStateLeaseContext } from "../state/openclaw-state-lease-context.js"; -import { runWithOpenClawStateLeaseWorker } from "../state/openclaw-state-lease-worker-storage.js"; +import { runWithOpenClawStateLeaseWorker } from "../state/openclaw-state-lease-worker-operation.js"; import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.js"; import type { OpenClawStateWorkerOperations } from "../state/openclaw-state-worker-contract.js"; import { runOpenClawStateWorkerOperation } from "../state/openclaw-state-worker-store.js"; import { prepareTranscriptDateReader } from "./store-date-preparation.js"; import { TranscriptLibraryError } from "./store-read.js"; -import type { - TranscriptExportWriteKey, - TranscriptReadRequests, - TranscriptWriteOperations, -} from "./store-worker-contract.js"; +import type { TranscriptReadRequests } from "./store-worker-contract.js"; +import type { TranscriptExportWriteKey } from "./store-worker.types.js"; +import type { TranscriptWriteOperations } from "./store-write.worker-contract.js"; /** One captured database generation spans planning, filesystem work, and persistence. */ export function createTranscriptStoreOperation( diff --git a/src/transcripts/store-worker-contract.ts b/src/transcripts/store-worker-contract.ts index ce7305867cef..2663445b327c 100644 --- a/src/transcripts/store-worker-contract.ts +++ b/src/transcripts/store-worker-contract.ts @@ -1,5 +1,4 @@ import type { SqliteWorkerCommand } from "../infra/sqlite-worker-contract.js"; -import type { OpenClawStateLeaseIdentity } from "../state/openclaw-state-lease-store.js"; import type { TranscriptSessionDescriptor, TranscriptSourceLocator } from "./provider-types.js"; import type { queryTranscriptReadEntries, @@ -24,68 +23,12 @@ import type { readTranscriptJsonlDigest, } from "./store-sqlite-read.js"; import type { - writeMeetingTranscriptSessionInDatabase, - writeMeetingTranscriptSummaryInDatabase, -} from "./store-sqlite-write.js"; -import type { - appendMeetingTranscriptUtterance, readRecentStoppedTranscriptSession, readTranscriptSummaryInputRevision, } from "./store-sqlite.js"; type SessionIdentity = Pick; -/** Host-only capture scheduling; functions never cross the worker boundary. */ -export type TranscriptAppendScheduler = ( - write: (assertCurrent: () => void) => Promise, -) => Promise; - -export type TranscriptWriteOperations = { - "transcripts.writeSession": { - input: Parameters[1] & { readOnly?: boolean }; - output: { ok: true } | { ok: false; reason: "changed" | "conflict" }; - }; - "transcripts.markPendingExports": { - input: { - session: SessionIdentity; - fileNames: string[]; - readOnly?: boolean; - lease?: OpenClawStateLeaseIdentity; - }; - output: void; - }; - "transcripts.recordExportManifest": { - input: { - session: SessionIdentity; - exportedHashes: Record; - removedExports: string[]; - readOnly?: boolean; - lease?: OpenClawStateLeaseIdentity; - }; - output: void; - }; - "transcripts.append": { - input: Omit[0], "database"> & { - readOnly?: boolean; - }; - output: void; - }; - "transcripts.writeSummary": { - input: { - session: SessionIdentity; - summaryValues: Parameters[2]; - guard?: Parameters[3]; - readOnly?: boolean; - }; - output: { ok: true } | { ok: false; reason: "changed" }; - }; -}; - -export type TranscriptWriteCommand = SqliteWorkerCommand; -export type TranscriptExportWriteKey = - | "transcripts.markPendingExports" - | "transcripts.recordExportManifest"; - export type TranscriptReadRequests = { "transcripts.canonicalSessionRow": { input: { selector: string }; diff --git a/src/transcripts/store-worker-write.ts b/src/transcripts/store-worker-write.ts index 9971b7f01ecd..d73e1f65a5d6 100644 --- a/src/transcripts/store-worker-write.ts +++ b/src/transcripts/store-worker-write.ts @@ -1,10 +1,13 @@ +import type { DatabaseSync } from "node:sqlite"; import { requestSqliteWorkerOperationAdmission } from "../infra/sqlite-worker-operation-admission.js"; -import { getSqliteWorkerStateContext } from "../infra/sqlite-worker-state-context.js"; -import { - runOpenClawStateWriteTransaction, - type OpenClawStateDatabase, -} from "../state/openclaw-state-db.js"; +import { runOpenClawStateWriteTransaction } from "../state/openclaw-state-db.js"; import { assertOpenClawStateLeaseWorkerOwnedInTransaction } from "../state/openclaw-state-lease-worker.js"; +import type { OpenClawStateLeaseIdentity } from "../state/openclaw-state-lease.types.js"; +import type { + WorkerOperationContext, + WorkerOperationHandlers, +} from "../state/worker-operation-registry.js"; +import type { TranscriptSessionDescriptor } from "./provider-types.js"; import { ensureMeetingTranscriptsSchema } from "./sqlite-schema.js"; import { transcriptSessionExportKey } from "./store-artifacts.js"; import { TranscriptSessionConflictError, TranscriptsSummaryChangedError } from "./store-errors.js"; @@ -15,112 +18,126 @@ import { writeMeetingTranscriptSummaryInDatabase, } from "./store-sqlite-write.js"; import { appendMeetingTranscriptUtterance } from "./store-sqlite.js"; -import type { TranscriptWriteCommand, TranscriptWriteOperations } from "./store-worker-contract.js"; -const operationLabels: Record = { - "transcripts.append": "meeting-transcripts.utterance.append", - "transcripts.writeSummary": "meeting-transcripts.summary.write", - "transcripts.writeSession": "meeting-transcripts.session.write", - "transcripts.markPendingExports": "meeting-transcripts.export.pending", - "transcripts.recordExportManifest": "meeting-transcripts.export.record", +type SessionIdentity = Pick; +type ExportInput = { + session: SessionIdentity; + lease?: OpenClawStateLeaseIdentity; + readOnly?: boolean; }; -export function isTranscriptWriteCommand(command: { - type: string; -}): command is TranscriptWriteCommand { - return Object.hasOwn(operationLabels, command.type); -} - -export function executeTranscriptWrite( - command: TranscriptWriteCommand, - target: { database: OpenClawStateDatabase; path: string }, -): TranscriptWriteOperations[keyof TranscriptWriteOperations]["output"] { - const options = { - ...target, - env: getSqliteWorkerStateContext().environment, - readOnly: command.input.readOnly, - }; +function writeTranscript( + input: { readOnly?: boolean }, + { open, stateOptions }: WorkerOperationContext, + operationLabel: string, + write: (db: DatabaseSync) => void, + exportInput?: ExportInput, +) { + const options = { database: open(), ...stateOptions(), readOnly: input.readOnly }; ensureMeetingTranscriptsSchema(options); - try { - runOpenClawStateWriteTransaction( - ({ db }) => { - const lease = - command.type === "transcripts.markPendingExports" || - command.type === "transcripts.recordExportManifest" - ? command.input.lease - : undefined; - const assertLease = () => { - if (!lease) { - return; - } + runOpenClawStateWriteTransaction( + ({ db }) => { + const assertWrite = (stage: "transaction" | "commit") => { + if (exportInput?.lease) { + const { lease, session } = exportInput; if ( lease.scope !== "meeting-transcript.export" || - lease.key !== transcriptSessionExportKey(command.input.session) + lease.key !== transcriptSessionExportKey(session) ) { throw new Error("Transcript export lease does not match its session"); } assertOpenClawStateLeaseWorkerOwnedInTransaction(db, lease); - }; - if (lease) { - assertLease(); } else { - requestSqliteWorkerOperationAdmission({ stage: "transaction", facts: undefined }); + requestSqliteWorkerOperationAdmission({ stage, facts: undefined }); } - switch (command.type) { - case "transcripts.append": - appendMeetingTranscriptUtterance({ ...command.input, database: db }); - break; - case "transcripts.writeSummary": { - const { session, summaryValues, guard } = command.input; - writeMeetingTranscriptSummaryInDatabase(db, session, summaryValues, guard); - break; - } - case "transcripts.writeSession": - writeMeetingTranscriptSessionInDatabase(db, command.input); - break; - case "transcripts.markPendingExports": - markMeetingTranscriptPendingExportsInDatabase( - db, - command.input.session, - command.input.fileNames, - ); - break; - case "transcripts.recordExportManifest": - updateMeetingTranscriptExportManifestInDatabase( - db, - command.input.session, - command.input.exportedHashes, - new Set(command.input.removedExports), - ); - break; - } - if (lease) { - assertLease(); - } else { - requestSqliteWorkerOperationAdmission({ stage: "commit", facts: undefined }); - } - }, - options, - { operationLabel: operationLabels[command.type] }, - ); - return command.type === "transcripts.writeSummary" || - command.type === "transcripts.writeSession" - ? { ok: true } - : undefined; - } catch (error) { - if ( - (command.type === "transcripts.writeSummary" || - command.type === "transcripts.writeSession") && - error instanceof TranscriptsSummaryChangedError - ) { - return { ok: false, reason: "changed" }; - } - if ( - command.type === "transcripts.writeSession" && - error instanceof TranscriptSessionConflictError - ) { - return { ok: false, reason: "conflict" }; - } - throw error; - } + }; + assertWrite("transaction"); + write(db); + assertWrite("commit"); + }, + options, + { operationLabel }, + ); } + +export const transcriptWriteOperations = { + "transcripts.append": ( + input: Omit[0], "database"> & { + readOnly?: boolean; + }, + context, + ) => + writeTranscript(input, context, "meeting-transcripts.utterance.append", (db) => + appendMeetingTranscriptUtterance({ ...input, database: db }), + ), + "transcripts.writeSummary": ( + input: { + session: SessionIdentity; + summaryValues: Parameters[2]; + guard?: Parameters[3]; + readOnly?: boolean; + }, + context, + ) => { + try { + writeTranscript(input, context, "meeting-transcripts.summary.write", (db) => + writeMeetingTranscriptSummaryInDatabase( + db, + input.session, + input.summaryValues, + input.guard, + ), + ); + return { ok: true } as const; + } catch (error) { + if (error instanceof TranscriptsSummaryChangedError) { + return { ok: false, reason: "changed" } as const; + } + throw error; + } + }, + "transcripts.writeSession": ( + input: Parameters[1] & { readOnly?: boolean }, + context, + ) => { + try { + writeTranscript(input, context, "meeting-transcripts.session.write", (db) => + writeMeetingTranscriptSessionInDatabase(db, input), + ); + return { ok: true } as const; + } catch (error) { + if (error instanceof TranscriptsSummaryChangedError) { + return { ok: false, reason: "changed" } as const; + } + if (error instanceof TranscriptSessionConflictError) { + return { ok: false, reason: "conflict" } as const; + } + throw error; + } + }, + "transcripts.markPendingExports": (input: ExportInput & { fileNames: string[] }, context) => + writeTranscript( + input, + context, + "meeting-transcripts.export.pending", + (db) => markMeetingTranscriptPendingExportsInDatabase(db, input.session, input.fileNames), + input, + ), + "transcripts.recordExportManifest": ( + input: ExportInput & { exportedHashes: Record; removedExports: string[] }, + context, + ) => + writeTranscript( + input, + context, + "meeting-transcripts.export.record", + (db) => + updateMeetingTranscriptExportManifestInDatabase( + db, + input.session, + input.exportedHashes, + new Set(input.removedExports), + ), + input, + ), +} satisfies WorkerOperationHandlers; diff --git a/src/transcripts/store-worker.types.ts b/src/transcripts/store-worker.types.ts new file mode 100644 index 000000000000..6474f8261701 --- /dev/null +++ b/src/transcripts/store-worker.types.ts @@ -0,0 +1,8 @@ +/** Host-only capture scheduling; functions never cross the worker boundary. */ +export type TranscriptAppendScheduler = ( + write: (assertCurrent: () => void) => Promise, +) => Promise; + +export type TranscriptExportWriteKey = + | "transcripts.markPendingExports" + | "transcripts.recordExportManifest"; diff --git a/src/transcripts/store-write.worker-contract.ts b/src/transcripts/store-write.worker-contract.ts new file mode 100644 index 000000000000..17505611211a --- /dev/null +++ b/src/transcripts/store-write.worker-contract.ts @@ -0,0 +1,4 @@ +import type { WorkerOperations } from "../state/worker-operation-registry.js"; +import type { transcriptWriteOperations } from "./store-worker-write.js"; + +export type TranscriptWriteOperations = WorkerOperations; diff --git a/src/transcripts/store.ts b/src/transcripts/store.ts index bd0ff22c2b1a..106e6d4b9cea 100644 --- a/src/transcripts/store.ts +++ b/src/transcripts/store.ts @@ -1,4 +1,3 @@ -// Stores meeting-capture transcripts in the shared SQLite state database. import fs from "node:fs/promises"; import path from "node:path"; import type { TranscriptUtterance as ProjectedTranscriptUtterance } from "../../packages/gateway-protocol/src/schema/transcripts.js"; @@ -50,11 +49,10 @@ import { createTranscriptStoreOperation, type TranscriptStoreOperation, } from "./store-worker-client.js"; -import type { - TranscriptAppendScheduler, - TranscriptReadRequests, - TranscriptWriteOperations, -} from "./store-worker-contract.js"; +import type { TranscriptReadRequests } from "./store-worker-contract.js"; +import type { TranscriptAppendScheduler } from "./store-worker.types.js"; +// Stores meeting-capture transcripts in the shared SQLite state database. +import type { TranscriptWriteOperations } from "./store-write.worker-contract.js"; import type { TranscriptsSummary } from "./summary.js"; import { renderTranscriptsMarkdown } from "./summary.js";