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";