mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
fix: preserve channel-owner revocation across recovery (#157709)
This commit is contained in:
parent
b98500ebcc
commit
83d24ab022
35 changed files with 1370 additions and 394 deletions
|
|
@ -40,7 +40,7 @@ public enum OpenClawNativeStateSQLiteValueType: Equatable, Sendable {
|
|||
/// One recursive connection lock serializes transactions and statement access.
|
||||
public final class OpenClawNativeStateSQLite: @unchecked Sendable {
|
||||
// Keep aligned with OPENCLAW_STATE_SCHEMA_VERSION. Native clients never upgrade this database.
|
||||
private static let maximumSupportedSchemaVersion: Int64 = 18
|
||||
private static let maximumSupportedSchemaVersion: Int64 = 19
|
||||
private static let defaultBusyTimeoutMilliseconds: Int32 = 5000
|
||||
|
||||
private struct SchemaObject: Hashable {
|
||||
|
|
|
|||
|
|
@ -30,6 +30,50 @@ Doctor completes recognized schema-1 databases that predate the audit ledger bef
|
|||
| 16 | Skill Workshop ownership moves from workspace/provenance columns to per-agent directory containment | Unreleased |
|
||||
| 17 | Prepared worker lifecycle facts and one-use node workspace bindings | Unreleased |
|
||||
| 18 | Original requesting authority retained with shared GitHub publication receipts | Unreleased |
|
||||
| 19 | Durable original channel-owner authorization and revocation continuity | Unreleased |
|
||||
|
||||
### State schema 19
|
||||
|
||||
Schema 19 adds nullable `authorization_id TEXT` and
|
||||
`authorization_basis_json TEXT` to the existing `user_profile_identities` rows.
|
||||
The profile owner reuses an uninterrupted channel link's opaque UUID and records
|
||||
only its original access-policy grant reference (or JSON `null`). It does not
|
||||
copy channel identities into another store. A reference is versioned JSON
|
||||
`{ "version": 1, "id": "<uuid>" }`, not a credential: recovery resolves it through
|
||||
the profile read worker and rechecks current roles or identity scopes and the original plugin grant.
|
||||
Unknown versions, malformed references, missing rows, or corrupt grant facts do
|
||||
not authorize work. This schema supplies the owner contract; consumers must
|
||||
retain and revalidate their own effect authority.
|
||||
|
||||
Profile authority mutations clear these fields in the same transaction as their
|
||||
role, alias, or ownership change. Unlinking removes the row. Restoring a role or
|
||||
link cannot restore a retired reference. The existing config machine-state store
|
||||
records the activated `gateway.roles` and `gateway.auth.identityScopes` policy
|
||||
snapshot under `operator.channelPolicy`. Identity-scope-only owners use the same
|
||||
reference contract; removing and restoring their configured scopes cannot revive
|
||||
a retired reference. Under the secrets activation lock, changes to either policy
|
||||
retire existing references durably before
|
||||
publishing the exact successor snapshot. Rollback uses the same operation with
|
||||
the merged target; publication failure reconciles the surviving active policy.
|
||||
No agent schema changes are involved.
|
||||
|
||||
Startup and Doctor add the nullable columns without assigning authority to legacy
|
||||
rows; unused profile tables remain absent until first use. The canonical schema
|
||||
and the profile owner's lazy ensure share the table and index definitions.
|
||||
Both published schema markers normally advance to 19. During the existing
|
||||
[older-updater publication deferral](/reference/database-schemas/versioning#schema-bumps-and-older-updaters),
|
||||
content can be upgraded while older published markers remain. References cannot
|
||||
be issued or recovered until version 19 is published.
|
||||
|
||||
Version 18 writers cannot preserve revocation history, so this is a versioned
|
||||
permission contract even though the added columns are nullable. Their normal
|
||||
worker admission checks foreign version changes under the state coordinator and
|
||||
refuses further writes, including deferred content newer than their support. Older readers also refuse schema 19.
|
||||
Create a verified, WAL-aware backup before upgrading. Downgrade requires restoring
|
||||
the matching pre-upgrade backup in a separate state directory; never lower the
|
||||
markers or remove authority columns. Restoring a backup loses later revocations
|
||||
and receipts and does not undo external effects; reconcile those with a compatible
|
||||
build before rollback.
|
||||
|
||||
### State schema 18
|
||||
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@
|
|||
"openclaw": {
|
||||
"updateAdmissionProtocol": 1,
|
||||
"schemaVersions": {
|
||||
"state": 18,
|
||||
"state": 19,
|
||||
"agent": 23
|
||||
}
|
||||
},
|
||||
|
|
|
|||
|
|
@ -7,7 +7,8 @@ import {
|
|||
getActiveSecretsRuntimeSnapshotState,
|
||||
getActiveSecretsRuntimeSnapshotRevisionState,
|
||||
graftActiveSecretsRuntimeAuthState,
|
||||
restoreSecretsRuntimeSnapshotStateIfCurrent,
|
||||
prepareSecretsRuntimeSnapshotRestoreState,
|
||||
activateSecretsRuntimeSnapshotStateIfCurrent,
|
||||
} from "../../secrets/runtime-state.js";
|
||||
import { withEnv } from "../../test-utils/env.js";
|
||||
import { resolveSharedAuthStorePath } from "./path-resolve.js";
|
||||
|
|
@ -249,13 +250,15 @@ describe("explicit auth state ownership", () => {
|
|||
activateSecretsRuntimeSnapshotState({ ...prepared, refreshHandler: null });
|
||||
} else {
|
||||
expect(
|
||||
restoreSecretsRuntimeSnapshotStateIfCurrent({
|
||||
snapshot: baseline,
|
||||
ownedSnapshot: owned,
|
||||
expectedRevision: revision,
|
||||
refreshContext: null,
|
||||
refreshHandler: null,
|
||||
}),
|
||||
activateSecretsRuntimeSnapshotStateIfCurrent(
|
||||
prepareSecretsRuntimeSnapshotRestoreState({
|
||||
snapshot: baseline,
|
||||
ownedSnapshot: owned,
|
||||
expectedRevision: revision,
|
||||
refreshContext: null,
|
||||
refreshHandler: null,
|
||||
})!,
|
||||
),
|
||||
).toBe(true);
|
||||
}
|
||||
expect(snapshotAt(resolveAuthProfileDatabasePath(agentDir))?.profiles.shared).toEqual(
|
||||
|
|
@ -356,13 +359,15 @@ describe("explicit auth state ownership", () => {
|
|||
expect(snapshotAt(first.agentPath)).toBeUndefined();
|
||||
vi.stubEnv("OPENCLAW_STATE_DIR", second.stateDir);
|
||||
expect(
|
||||
restoreSecretsRuntimeSnapshotStateIfCurrent({
|
||||
snapshot: captured.baseline,
|
||||
ownedSnapshot: captured.owned,
|
||||
expectedRevision: captured.revision,
|
||||
refreshContext: null,
|
||||
refreshHandler: null,
|
||||
}),
|
||||
activateSecretsRuntimeSnapshotStateIfCurrent(
|
||||
prepareSecretsRuntimeSnapshotRestoreState({
|
||||
snapshot: captured.baseline,
|
||||
ownedSnapshot: captured.owned,
|
||||
expectedRevision: captured.revision,
|
||||
refreshContext: null,
|
||||
refreshHandler: null,
|
||||
})!,
|
||||
),
|
||||
).toBe(true);
|
||||
expect(snapshotAt(first.agentPath)?.profiles.shared).toEqual(apiKey("first"));
|
||||
await updateAuthProfileStoreWithLock({
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
import type { OpenClawConfig } from "../../config/types.openclaw.js";
|
||||
import type { GatewayRequestContext } from "../../gateway/server-methods/types.js";
|
||||
import { linkUserChannelIdentity } from "../../state/user-channel-identities.js";
|
||||
import { publishCanonicalUserChannelPolicy } from "../../state/user-channel-identity-operations.js";
|
||||
import { ensureProfileForEmail, setUserProfileRole } from "../../state/user-profiles.js";
|
||||
import {
|
||||
withOpenClawTestState,
|
||||
|
|
@ -64,6 +65,12 @@ async function createFixture(state: OpenClawTestState, authority: "role" | "iden
|
|||
linkUserChannelIdentity(profile.id, identity);
|
||||
return { profile, identity };
|
||||
});
|
||||
const activatePolicy = async (update: Partial<NonNullable<OpenClawConfig["gateway"]>>) => {
|
||||
const gateway = { ...cfg.gateway, ...update };
|
||||
await publishCanonicalUserChannelPolicy(gateway);
|
||||
cfg.gateway = gateway;
|
||||
};
|
||||
await activatePolicy({});
|
||||
let live = true;
|
||||
const gateway = createCommandOwnerTestGateway(cfg);
|
||||
const owner = {
|
||||
|
|
@ -115,6 +122,7 @@ async function createFixture(state: OpenClawTestState, authority: "role" | "iden
|
|||
};
|
||||
return {
|
||||
cfg,
|
||||
activatePolicy,
|
||||
state,
|
||||
admins,
|
||||
context,
|
||||
|
|
|
|||
|
|
@ -9,11 +9,21 @@ import {
|
|||
} from "../../auto-reply/command-owner-authority.js";
|
||||
import { prepareChannelRunAdmission } from "../../auto-reply/reply/channel-run-admission.js";
|
||||
import { installDiscordRegistryHooks } from "../../auto-reply/test-helpers/command-auth-registry-fixture.js";
|
||||
import { prepareChannelOperatorAdmin } from "../../gateway/channel-operator-authority.js";
|
||||
import {
|
||||
createPluginRegistryFixture,
|
||||
registerVirtualTestPlugin,
|
||||
} from "../../plugin-sdk/test-helpers/contracts-testkit.js";
|
||||
import { stageActivePluginRegistry } from "../../plugins/runtime.js";
|
||||
import {
|
||||
closeOpenClawStateDatabaseAsync,
|
||||
openOpenClawStateDatabase,
|
||||
} from "../../state/openclaw-state-db.js";
|
||||
import {
|
||||
linkUserChannelIdentity,
|
||||
unlinkUserChannelIdentity,
|
||||
} from "../../state/user-channel-identities.js";
|
||||
import { setUserProfileRole } from "../../state/user-profiles.js";
|
||||
import { linkEmail, setUserProfileRole } from "../../state/user-profiles.js";
|
||||
import { withAdminIngress } from "./operator-authority.test-support.js";
|
||||
|
||||
installDiscordRegistryHooks();
|
||||
|
|
@ -138,6 +148,241 @@ it.each(["role", "role-scopes", "grant", "link", "reassign", "host"] as const)(
|
|||
},
|
||||
);
|
||||
|
||||
it.each(["definition", "default", "identity-scopes"] as const)(
|
||||
"does not revive an admitted channel owner after restoring its policy %s",
|
||||
async (change) => {
|
||||
await withAdminIngress(
|
||||
async ({ cfg, admins, context, activatePolicy }) => {
|
||||
const admin = admins[0]!;
|
||||
if (change === "default") {
|
||||
setUserProfileRole(admin.profile.id, null);
|
||||
await activatePolicy({ roles: { ...cfg.gateway!.roles!, default: "admin" } });
|
||||
}
|
||||
const ctx = await context(admin.identity.senderId);
|
||||
const assertCurrent = captureCommandOwnerAssertion(ctx);
|
||||
expect(assertCurrent).toBeTypeOf("function");
|
||||
expect(assertCurrent).not.toThrow();
|
||||
const restored = structuredClone(cfg.gateway!);
|
||||
const revoked = structuredClone(restored);
|
||||
if (change === "definition") {
|
||||
revoked.roles!.definitions.admin!.scopes = ["operator.read"];
|
||||
} else if (change === "default") {
|
||||
revoked.roles!.default = "member";
|
||||
} else {
|
||||
delete revoked.auth!.identityScopes;
|
||||
}
|
||||
await activatePolicy(revoked);
|
||||
await activatePolicy(restored);
|
||||
|
||||
const fresh = await context(admin.identity.senderId);
|
||||
expect(
|
||||
resolveCommandAuthorization({ cfg, ctx: fresh, commandAuthorized: true }).senderIsOwner,
|
||||
).toBe(true);
|
||||
expect(assertCurrent).toThrow();
|
||||
},
|
||||
change === "identity-scopes" ? "identity-grant" : "role",
|
||||
);
|
||||
},
|
||||
);
|
||||
|
||||
it.each(["role", "identity-grant"] as const)(
|
||||
"recovers an original %s owner from its exact JSON reference across a database lifecycle",
|
||||
async (authority) => {
|
||||
await withAdminIngress(async ({ cfg, admins }) => {
|
||||
const admitted = await prepareChannelOperatorAdmin(cfg, admins[0]!.identity);
|
||||
expect(admitted?.isCurrent(cfg)).toBe(true);
|
||||
expect(admitted?.recoveryReference).toEqual({ version: 1, id: expect.any(String) });
|
||||
const encoded = JSON.stringify(admitted!.recoveryReference);
|
||||
await closeOpenClawStateDatabaseAsync();
|
||||
expect(admitted?.isCurrent(cfg)).toBe(false);
|
||||
const reference = JSON.parse(encoded);
|
||||
const resumed = await prepareChannelOperatorAdmin(cfg, reference);
|
||||
expect(resumed?.isCurrent(cfg)).toBe(true);
|
||||
expect(resumed?.recoveryReference).toEqual(reference);
|
||||
}, authority);
|
||||
},
|
||||
);
|
||||
|
||||
it.each([
|
||||
"role",
|
||||
"role-restore",
|
||||
"unlink",
|
||||
"relink",
|
||||
"reassign",
|
||||
"merge",
|
||||
"definition",
|
||||
"default",
|
||||
"identity-scopes",
|
||||
] as const)(
|
||||
"never recovers the original owner after %s retires its durable reference",
|
||||
async (change) => {
|
||||
await withAdminIngress(
|
||||
async ({ cfg, admins, activatePolicy }) => {
|
||||
const admin = admins[0]!;
|
||||
if (change === "default") {
|
||||
setUserProfileRole(admin.profile.id, null);
|
||||
await activatePolicy({ roles: { ...cfg.gateway!.roles!, default: "admin" } });
|
||||
}
|
||||
const admitted = await prepareChannelOperatorAdmin(cfg, admin.identity);
|
||||
expect(admitted?.isCurrent(cfg)).toBe(true);
|
||||
const reference = admitted!.recoveryReference!;
|
||||
expect(reference).toBeDefined();
|
||||
if (change === "role" || change === "role-restore") {
|
||||
setUserProfileRole(admin.profile.id, "member");
|
||||
if (change === "role-restore") {
|
||||
setUserProfileRole(admin.profile.id, "admin");
|
||||
}
|
||||
} else if (change === "unlink" || change === "relink" || change === "reassign") {
|
||||
unlinkUserChannelIdentity(admin.profile.id, admin.identity);
|
||||
if (change !== "unlink") {
|
||||
linkUserChannelIdentity(
|
||||
change === "relink" ? admin.profile.id : admins[1]!.profile.id,
|
||||
admin.identity,
|
||||
);
|
||||
}
|
||||
} else if (change === "merge") {
|
||||
linkEmail("ada@example.test", admins[1]!.profile.id);
|
||||
} else if (change === "identity-scopes") {
|
||||
const original = structuredClone(cfg.gateway!.auth!);
|
||||
await activatePolicy({ auth: { ...original, identityScopes: undefined } });
|
||||
await activatePolicy({ auth: original });
|
||||
} else {
|
||||
const original = structuredClone(cfg.gateway!.roles!);
|
||||
const revoked = structuredClone(original);
|
||||
if (change === "definition") {
|
||||
revoked.definitions.admin!.scopes = ["operator.read"];
|
||||
} else {
|
||||
revoked.default = "member";
|
||||
}
|
||||
await activatePolicy({ roles: revoked });
|
||||
await activatePolicy({ roles: original });
|
||||
}
|
||||
expect(admitted?.isCurrent(cfg)).toBe(false);
|
||||
await closeOpenClawStateDatabaseAsync();
|
||||
await expect(prepareChannelOperatorAdmin(cfg, reference)).resolves.toBeUndefined();
|
||||
await expect(prepareChannelOperatorAdmin(cfg, reference)).resolves.toBeUndefined();
|
||||
},
|
||||
change === "identity-scopes" ? "identity-grant" : "role",
|
||||
);
|
||||
},
|
||||
);
|
||||
|
||||
it.each(["allowed", "revoked", "replaced", "unavailable"] as const)(
|
||||
"resumes only the original plugin grant when it is %s",
|
||||
async (change) => {
|
||||
await withAdminIngress(async ({ cfg, admins, activatePolicy }) => {
|
||||
const pluginId = "channel-owner-access";
|
||||
const originalId = "86633673-b1dd-4500-85e2-b6e6e490810f";
|
||||
let grantId: string | undefined = originalId;
|
||||
let lifetime = new AbortController();
|
||||
let unavailable = false;
|
||||
const { config, registry } = createPluginRegistryFixture(cfg);
|
||||
registerVirtualTestPlugin({
|
||||
registry,
|
||||
config,
|
||||
id: pluginId,
|
||||
name: "Channel owner access",
|
||||
register(api) {
|
||||
const current = () => {
|
||||
if (unavailable) {
|
||||
throw new Error("Policy store is unavailable");
|
||||
}
|
||||
const signal = lifetime.signal;
|
||||
return grantId
|
||||
? { grantId, signal, assertCurrent: () => signal.throwIfAborted() }
|
||||
: undefined;
|
||||
};
|
||||
api.registerGatewayAccessPolicy({
|
||||
authorize: current,
|
||||
resume: ({ grantId: requested }) => (requested === grantId ? current() : undefined),
|
||||
});
|
||||
},
|
||||
});
|
||||
stageActivePluginRegistry(registry.registry, null, "default");
|
||||
const roles = structuredClone(cfg.gateway!.roles!);
|
||||
roles.definitions.admin!.accessPolicyPlugin = pluginId;
|
||||
await activatePolicy({ roles });
|
||||
const original = await prepareChannelOperatorAdmin(cfg, admins[0]!.identity);
|
||||
const reference = original!.recoveryReference!;
|
||||
expect(reference).toBeDefined();
|
||||
if (change === "revoked" || change === "replaced") {
|
||||
lifetime.abort();
|
||||
grantId = change === "revoked" ? undefined : "78a7c3d0-c3a6-49a5-91e7-02c153e39ab5";
|
||||
lifetime = new AbortController();
|
||||
}
|
||||
unavailable = change === "unavailable";
|
||||
await closeOpenClawStateDatabaseAsync();
|
||||
if (change === "unavailable") {
|
||||
await expect(prepareChannelOperatorAdmin(cfg, reference)).rejects.toMatchObject({
|
||||
name: "GatewayOperatorAccessUnavailableError",
|
||||
});
|
||||
} else if (change === "allowed") {
|
||||
expect((await prepareChannelOperatorAdmin(cfg, reference))?.isCurrent(cfg)).toBe(true);
|
||||
} else {
|
||||
await expect(prepareChannelOperatorAdmin(cfg, reference)).resolves.toBeUndefined();
|
||||
// Restoring access creates a successor grant, never the retired grant's identity.
|
||||
grantId = "78a7c3d0-c3a6-49a5-91e7-02c153e39ab5";
|
||||
expect((await prepareChannelOperatorAdmin(cfg, admins[0]!.identity))?.isCurrent(cfg)).toBe(
|
||||
true,
|
||||
);
|
||||
await expect(prepareChannelOperatorAdmin(cfg, reference)).resolves.toBeUndefined();
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
it("keeps the active policy and its reference when durable policy retirement rolls back", async () => {
|
||||
await withAdminIngress(async ({ cfg, admins, activatePolicy }) => {
|
||||
const original = await prepareChannelOperatorAdmin(cfg, admins[0]!.identity);
|
||||
const activeRoles = structuredClone(cfg.gateway!.roles!);
|
||||
const changed = structuredClone(activeRoles);
|
||||
changed.definitions.admin!.scopes = ["operator.read"];
|
||||
const db = openOpenClawStateDatabase().db;
|
||||
db.exec(`CREATE TRIGGER fail_policy_publication BEFORE UPDATE ON config_machine_state
|
||||
WHEN NEW.state_key = 'operator.channelPolicy'
|
||||
BEGIN SELECT RAISE(ABORT, 'fixture policy write failed'); END;`);
|
||||
await expect(activatePolicy({ roles: changed })).rejects.toThrow("fixture policy write failed");
|
||||
expect(cfg.gateway!.roles).toEqual(activeRoles);
|
||||
db.exec("DROP TRIGGER fail_policy_publication");
|
||||
await closeOpenClawStateDatabaseAsync();
|
||||
expect(
|
||||
(await prepareChannelOperatorAdmin(cfg, original!.recoveryReference!))?.isCurrent(cfg),
|
||||
).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
it.each(["missing", "malformed", "version", "extra", "basis"] as const)(
|
||||
"fails closed on a %s recovery reference without treating its ID as authority",
|
||||
async (damage) => {
|
||||
await withAdminIngress(async ({ cfg, admins }) => {
|
||||
const admitted = await prepareChannelOperatorAdmin(cfg, admins[0]!.identity);
|
||||
const reference = { ...admitted!.recoveryReference! };
|
||||
expect(reference.id).toBeTypeOf("string");
|
||||
let encoded = JSON.stringify(reference);
|
||||
if (damage === "missing") {
|
||||
encoded = JSON.stringify({ ...reference, id: "00000000-0000-4000-8000-000000000000" });
|
||||
}
|
||||
if (damage === "malformed") {
|
||||
encoded = JSON.stringify({ ...reference, id: "corrupt" });
|
||||
}
|
||||
if (damage === "version") {
|
||||
encoded = JSON.stringify({ ...reference, version: 2 });
|
||||
}
|
||||
if (damage === "extra") {
|
||||
encoded = JSON.stringify({ ...reference, scopes: ["operator.admin"] });
|
||||
}
|
||||
if (damage === "basis") {
|
||||
openOpenClawStateDatabase()
|
||||
.db.prepare(
|
||||
"UPDATE user_profile_identities SET authorization_basis_json = ? WHERE authorization_id = ?",
|
||||
)
|
||||
.run("{", reference.id);
|
||||
}
|
||||
await expect(prepareChannelOperatorAdmin(cfg, JSON.parse(encoded))).resolves.toBeUndefined();
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
it.each([
|
||||
{ role: "maintainer", defaultRole: "member", identityGrant: false, owner: true },
|
||||
{ role: null, defaultRole: "maintainer", identityGrant: false, owner: true },
|
||||
|
|
|
|||
|
|
@ -1,7 +1,16 @@
|
|||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import type { OpenClawStateDatabaseOptions } from "../state/openclaw-state-db-contract.js";
|
||||
import { resolveUserChannelIdentity } from "../state/user-channel-identities.js";
|
||||
import { prepareUserChannelIdentityAuthority } from "../state/user-channel-identity-operations.js";
|
||||
import {
|
||||
parseUserChannelAuthorizationReference,
|
||||
resolveUserChannelAuthorizationPolicy,
|
||||
resolveUserChannelIdentity,
|
||||
type UserChannelAuthorization,
|
||||
type UserChannelAuthorizationReference,
|
||||
} from "../state/user-channel-identities.js";
|
||||
import {
|
||||
authorizeCanonicalUserChannelIdentity,
|
||||
prepareUserChannelIdentityAuthority,
|
||||
} from "../state/user-channel-identity-operations.js";
|
||||
import type {
|
||||
UserChannelIdentity,
|
||||
UserChannelIdentityAuthorityFacts,
|
||||
|
|
@ -10,6 +19,7 @@ import {
|
|||
GatewayOperatorAccessDeniedError,
|
||||
hasCurrentGatewayOperatorAccess,
|
||||
resolvePreparedGatewayOperatorAccessAuthority,
|
||||
resumeGatewayOperatorAccessGrant,
|
||||
} from "./operator-access-policy.js";
|
||||
import { resolveIdentityOperatorScopes } from "./operator-identity-scopes.js";
|
||||
import { resolveOperatorRolePolicyForAssignment } from "./operator-role-policy.js";
|
||||
|
|
@ -21,7 +31,9 @@ export function resolveChannelOperatorAdminAuthority(
|
|||
stateOptions: OpenClawStateDatabaseOptions = {},
|
||||
) {
|
||||
const prepared = resolveChannelOperatorIdentityFacts(cfg, identity, stateOptions);
|
||||
return prepared && captureLinkedOperatorAdmin(cfg, prepared.linked, prepared.isCurrent);
|
||||
return (
|
||||
prepared && captureLinkedOperatorAdmin(cfg, prepared.linked, prepared.isCurrent)?.authority
|
||||
);
|
||||
}
|
||||
|
||||
/** The update owner must prove accepted native custody before using identity-only checks. */
|
||||
|
|
@ -106,23 +118,41 @@ function captureLinkedOperatorAdmin(
|
|||
cfg: OpenClawConfig,
|
||||
linked: UserChannelIdentityAuthorityFacts,
|
||||
isIdentityCurrent: () => boolean,
|
||||
original?: UserChannelAuthorization,
|
||||
) {
|
||||
const identity = captureLinkedOperatorAdminIdentity(cfg, linked, isIdentityCurrent);
|
||||
if (!identity) {
|
||||
return undefined;
|
||||
}
|
||||
try {
|
||||
const access = resolvePreparedGatewayOperatorAccessAuthority(
|
||||
{ ...linked, isCurrent: isIdentityCurrent },
|
||||
cfg,
|
||||
);
|
||||
const access = original
|
||||
? null
|
||||
: resolvePreparedGatewayOperatorAccessAuthority(
|
||||
{ ...linked, isCurrent: isIdentityCurrent },
|
||||
cfg,
|
||||
);
|
||||
let current = true;
|
||||
const isCurrent = (currentCfg: OpenClawConfig) => {
|
||||
current &&= identity.isCurrent(currentCfg) && hasCurrentGatewayOperatorAccess(access);
|
||||
if (current && original) {
|
||||
resumeGatewayOperatorAccessGrant(
|
||||
{ profileId: linked.profileId, emails: linked.emails, assignedRole: linked.role },
|
||||
currentCfg,
|
||||
original.grant,
|
||||
);
|
||||
current &&= isIdentityCurrent();
|
||||
}
|
||||
return current;
|
||||
};
|
||||
return isCurrent(cfg)
|
||||
? { profileId: linked.profileId, isCurrent, ...(access ? { signal: access.signal } : {}) }
|
||||
? {
|
||||
authority: {
|
||||
profileId: linked.profileId,
|
||||
isCurrent,
|
||||
...(access ? { signal: access.signal } : {}),
|
||||
},
|
||||
grant: access === null ? null : access.gatewayAccessGrant,
|
||||
}
|
||||
: undefined;
|
||||
} catch (error) {
|
||||
if (error instanceof GatewayOperatorAccessDeniedError) {
|
||||
|
|
@ -134,12 +164,54 @@ function captureLinkedOperatorAdmin(
|
|||
|
||||
export async function prepareChannelOperatorAdmin(
|
||||
cfg: OpenClawConfig,
|
||||
identity: UserChannelIdentity,
|
||||
identity: UserChannelIdentity | UserChannelAuthorizationReference,
|
||||
stateOptions: OpenClawStateDatabaseOptions = {},
|
||||
) {
|
||||
if (!cfg.gateway?.roles && !cfg.gateway?.auth?.identityScopes) {
|
||||
return undefined;
|
||||
}
|
||||
const prepared = await prepareUserChannelIdentityAuthority(identity, stateOptions);
|
||||
return prepared && captureLinkedOperatorAdmin(cfg, prepared.linked, prepared.isCurrent);
|
||||
const policy = resolveUserChannelAuthorizationPolicy(cfg.gateway);
|
||||
const reference =
|
||||
"version" in identity ? parseUserChannelAuthorizationReference(identity) : undefined;
|
||||
if ("version" in identity && !reference) {
|
||||
return undefined;
|
||||
}
|
||||
const prepared = await prepareUserChannelIdentityAuthority(
|
||||
"version" in identity ? { authorizationId: identity.id, policy } : identity,
|
||||
stateOptions,
|
||||
);
|
||||
if (!prepared) {
|
||||
return undefined;
|
||||
}
|
||||
const captured = captureLinkedOperatorAdmin(
|
||||
cfg,
|
||||
prepared.linked,
|
||||
prepared.isCurrent,
|
||||
prepared.linked.authorization,
|
||||
);
|
||||
if (!captured) {
|
||||
return undefined;
|
||||
}
|
||||
const original = prepared.linked.authorization;
|
||||
const assertCurrent = () => {
|
||||
if (!captured.authority.isCurrent(cfg)) {
|
||||
throw new GatewayOperatorAccessDeniedError();
|
||||
}
|
||||
};
|
||||
const recoveryReference =
|
||||
original?.reference ??
|
||||
(!("version" in identity) && captured.grant !== undefined
|
||||
? await authorizeCanonicalUserChannelIdentity(
|
||||
{
|
||||
action: "authorize",
|
||||
identity,
|
||||
profileId: prepared.linked.profileId,
|
||||
policy,
|
||||
grant: captured.grant,
|
||||
},
|
||||
{ ...stateOptions, assertCurrent },
|
||||
)
|
||||
: undefined);
|
||||
assertCurrent();
|
||||
return { ...captured.authority, recoveryReference };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { createDeferred } from "../../test/helpers/promise.js";
|
||||
import { createModelProviderRouteOverrideResolver } from "../config/model-provider-config.js";
|
||||
import {
|
||||
getRuntimeConfigSnapshot,
|
||||
|
|
@ -6,6 +7,11 @@ import {
|
|||
} from "../config/runtime-snapshot.js";
|
||||
import { projectConfigOntoRuntimeSourceSnapshot } from "../config/runtime-source-projection.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import {
|
||||
activateSecretsRuntimeSnapshotStateIfCurrent,
|
||||
getActiveSecretsRuntimeSnapshotState,
|
||||
getActiveSecretsRuntimeSnapshotRevisionState,
|
||||
} from "../secrets/runtime-state.js";
|
||||
import {
|
||||
activateSecretsRuntimeSnapshot,
|
||||
activateSecretsRuntimeSnapshotWithSource,
|
||||
|
|
@ -61,6 +67,15 @@ const prepare = (config: OpenClawConfig) =>
|
|||
env: {},
|
||||
});
|
||||
|
||||
function activatorOptions() {
|
||||
return {
|
||||
logSecrets: { info: vi.fn(), warn: vi.fn(), error: vi.fn() },
|
||||
emitStateEvent: vi.fn(),
|
||||
prepareRuntimeSecretsSnapshot: ({ config }: { config: OpenClawConfig }) => prepare(config),
|
||||
activateRuntimeSecretsSnapshot: activateSecretsRuntimeSnapshot,
|
||||
};
|
||||
}
|
||||
|
||||
function expectAuthoredSource(source: OpenClawConfig) {
|
||||
const config = getRuntimeConfigSnapshot();
|
||||
expect(config).not.toBeNull();
|
||||
|
|
@ -77,12 +92,7 @@ async function createReload(commit: () => Promise<void>, beforePublication?: ()
|
|||
const initial = configPair("openclaw");
|
||||
activateSecretsRuntimeSnapshotWithSource(await prepare(initial.config), initial.source);
|
||||
expectAuthoredSource(initial.source);
|
||||
const activateRuntimeSecrets = createRuntimeSecretsActivator({
|
||||
logSecrets: { info: vi.fn(), warn: vi.fn(), error: vi.fn() },
|
||||
emitStateEvent: vi.fn(),
|
||||
prepareRuntimeSecretsSnapshot: ({ config }) => prepare(config),
|
||||
activateRuntimeSecretsSnapshot: activateSecretsRuntimeSnapshot,
|
||||
});
|
||||
const activateRuntimeSecrets = createRuntimeSecretsActivator(activatorOptions());
|
||||
// Only secret publication is exercised here; the service tail is injected at its existing seam.
|
||||
const params = {
|
||||
activateRuntimeSecrets,
|
||||
|
|
@ -141,6 +151,180 @@ async function createReload(commit: () => Promise<void>, beforePublication?: ()
|
|||
}
|
||||
|
||||
describe("managed reload authored source", () => {
|
||||
it.each([false, true])(
|
||||
"commits the exact target before publication (hook fails: %s)",
|
||||
async (fails) => {
|
||||
const initial = await prepare({ gateway: { port: 18789 } });
|
||||
const candidate = await prepare({ gateway: { port: 18790 } });
|
||||
activateSecretsRuntimeSnapshot(initial);
|
||||
const failure = new Error("durable publication failed");
|
||||
const beforeSnapshotPublication = vi.fn(async (config: OpenClawConfig | null) => {
|
||||
expect(getActiveSecretsRuntimeSnapshotState()?.config).toEqual(initial.config);
|
||||
if (fails && config === candidate.config) {
|
||||
throw failure;
|
||||
}
|
||||
});
|
||||
const options = {
|
||||
...activatorOptions(),
|
||||
activateRuntimeSecretsSnapshot: activateSecretsRuntimeSnapshot,
|
||||
beforeSnapshotPublication,
|
||||
};
|
||||
const activate = createRuntimeSecretsActivator(options);
|
||||
const publication = activate.activatePreparedSnapshotIfCurrent(
|
||||
candidate,
|
||||
getActiveSecretsRuntimeSnapshotRevisionState(),
|
||||
{ reason: "reload", activate: true },
|
||||
);
|
||||
if (fails) {
|
||||
await expect(publication).rejects.toBe(failure);
|
||||
expect(beforeSnapshotPublication.mock.calls.map(([config]) => config)).toEqual([
|
||||
candidate.config,
|
||||
initial.config,
|
||||
]);
|
||||
} else {
|
||||
await expect(publication).resolves.toBe(candidate);
|
||||
expect(beforeSnapshotPublication).toHaveBeenCalledExactlyOnceWith(candidate.config);
|
||||
}
|
||||
expect(getActiveSecretsRuntimeSnapshotState()?.config).toEqual(
|
||||
fails ? initial.config : candidate.config,
|
||||
);
|
||||
},
|
||||
);
|
||||
|
||||
it.each(["snapshot", "admission"] as const)(
|
||||
"reconciles the survivor when %s ownership changes during the durable hook",
|
||||
async (supersession) => {
|
||||
const initial = await prepare({ gateway: { port: 18789 } });
|
||||
const candidate = await prepare({ gateway: { port: 18790 } });
|
||||
const successor = await prepare({ gateway: { port: 18791 } });
|
||||
activateSecretsRuntimeSnapshot(initial);
|
||||
const entered = createDeferred();
|
||||
const release = createDeferred();
|
||||
const beforeSnapshotPublication = vi.fn(async (config: OpenClawConfig | null) => {
|
||||
if (config === candidate.config) {
|
||||
entered.resolve();
|
||||
await release.promise;
|
||||
}
|
||||
});
|
||||
const activator = createRuntimeSecretsActivator({
|
||||
...activatorOptions(),
|
||||
activateRuntimeSecretsSnapshot: activateSecretsRuntimeSnapshot,
|
||||
beforeSnapshotPublication,
|
||||
});
|
||||
let current = true;
|
||||
const published = vi.fn();
|
||||
const pending = activator.activatePreparedSnapshotIfCurrent(
|
||||
candidate,
|
||||
getActiveSecretsRuntimeSnapshotRevisionState(),
|
||||
{ reason: "reload", activate: true },
|
||||
published,
|
||||
() => current,
|
||||
);
|
||||
await entered.promise;
|
||||
expect(getActiveSecretsRuntimeSnapshotState()?.config).toEqual(initial.config);
|
||||
if (supersession === "snapshot") {
|
||||
activateSecretsRuntimeSnapshot(successor);
|
||||
} else {
|
||||
current = false;
|
||||
}
|
||||
release.resolve();
|
||||
await expect(pending).resolves.toBeNull();
|
||||
expect(published).not.toHaveBeenCalled();
|
||||
const survivor = supersession === "snapshot" ? successor : initial;
|
||||
expect(beforeSnapshotPublication.mock.calls.map(([config]) => config)).toEqual([
|
||||
candidate.config,
|
||||
survivor.config,
|
||||
]);
|
||||
expect(getActiveSecretsRuntimeSnapshotState()?.config).toEqual(survivor.config);
|
||||
},
|
||||
);
|
||||
|
||||
it.each([false, true])(
|
||||
"reconciles a throwing publisher (snapshot replaced: %s)",
|
||||
async (replaced) => {
|
||||
const initial = await prepare({ gateway: { port: 18789 } });
|
||||
const candidate = await prepare({ gateway: { port: 18790 } });
|
||||
activateSecretsRuntimeSnapshot(initial);
|
||||
const beforeSnapshotPublication = vi.fn(async (_config: OpenClawConfig | null) => {});
|
||||
const failure = new Error("snapshot publication failed");
|
||||
const activator = createRuntimeSecretsActivator({
|
||||
...activatorOptions(),
|
||||
beforeSnapshotPublication,
|
||||
activateRuntimeSecretsSnapshot: (snapshot) => {
|
||||
if (replaced) {
|
||||
activateSecretsRuntimeSnapshot(snapshot);
|
||||
}
|
||||
throw failure;
|
||||
},
|
||||
});
|
||||
await expect(
|
||||
activator.activatePreparedSnapshot(candidate, {
|
||||
reason: "reload",
|
||||
activate: true,
|
||||
}),
|
||||
).rejects.toBe(failure);
|
||||
const survivor = replaced ? candidate : initial;
|
||||
expect(beforeSnapshotPublication.mock.calls.map(([config]) => config)).toEqual([
|
||||
candidate.config,
|
||||
survivor.config,
|
||||
]);
|
||||
expect(getActiveSecretsRuntimeSnapshotState()?.config).toEqual(survivor.config);
|
||||
},
|
||||
);
|
||||
|
||||
it.each([false, true])(
|
||||
"durably publishes the exact three-way rollback (inside callback: %s)",
|
||||
async (insideCallback) => {
|
||||
const initial = await prepare({ gateway: { port: 18789 } });
|
||||
const candidate = await prepare({ gateway: { port: 18790 } });
|
||||
const descendant = await prepare({ gateway: { port: 18790, bind: "lan" } });
|
||||
const merged = { gateway: { port: 18789, bind: "lan" } };
|
||||
activateSecretsRuntimeSnapshot(initial);
|
||||
const beforeSnapshotPublication = vi.fn(async (config: OpenClawConfig | null) => {
|
||||
expect(getActiveSecretsRuntimeSnapshotState()?.config).toEqual(
|
||||
config === candidate.config ? initial.config : descendant.config,
|
||||
);
|
||||
});
|
||||
const activator = createRuntimeSecretsActivator({
|
||||
...activatorOptions(),
|
||||
activateRuntimeSecretsSnapshot: activateSecretsRuntimeSnapshot,
|
||||
beforeSnapshotPublication,
|
||||
});
|
||||
let publishedRevision = 0;
|
||||
const restoreMerged = async (restore: typeof activator.restoreSnapshotIfCurrent) => {
|
||||
expect(
|
||||
activateSecretsRuntimeSnapshotStateIfCurrent({
|
||||
snapshot: descendant,
|
||||
expectedRevision: publishedRevision,
|
||||
preserveActivationLineage: true,
|
||||
refreshContext: null,
|
||||
refreshHandler: null,
|
||||
}),
|
||||
).toBe(true);
|
||||
await expect(restore(initial, publishedRevision, candidate)).resolves.toBe(true);
|
||||
};
|
||||
await activator.activatePreparedSnapshotIfCurrent(
|
||||
candidate,
|
||||
getActiveSecretsRuntimeSnapshotRevisionState(),
|
||||
{ reason: "reload", activate: true },
|
||||
async (restore) => {
|
||||
publishedRevision = getActiveSecretsRuntimeSnapshotRevisionState();
|
||||
if (insideCallback) {
|
||||
await restoreMerged(restore);
|
||||
}
|
||||
},
|
||||
);
|
||||
if (!insideCallback) {
|
||||
await restoreMerged(activator.restoreSnapshotIfCurrent);
|
||||
}
|
||||
expect(beforeSnapshotPublication.mock.calls.map(([config]) => config)).toEqual([
|
||||
candidate.config,
|
||||
merged,
|
||||
]);
|
||||
expect(getActiveSecretsRuntimeSnapshotState()?.config).toEqual(merged);
|
||||
},
|
||||
);
|
||||
|
||||
it("rejects a closed plugin invoker before activating its prepared secrets", async () => {
|
||||
const failure = new Error("plugin invoker closed");
|
||||
const commit = vi.fn(async () => {});
|
||||
|
|
|
|||
|
|
@ -6,14 +6,12 @@ import {
|
|||
} from "../config/config.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import {
|
||||
clearSecretsRuntimeSnapshotState,
|
||||
getActiveSecretsRuntimeSnapshotState,
|
||||
getActiveSecretsRuntimeSnapshotRevisionState,
|
||||
hasActiveSecretsRuntimeSnapshotLineage,
|
||||
hasSameSecretReloadContract,
|
||||
restoreSecretsRuntimeSourceSnapshotIfLineageCurrent,
|
||||
setSecretsRuntimeSourceSnapshotIfCurrent,
|
||||
type PreparedSecretsRuntimeSnapshot,
|
||||
} from "../secrets/runtime-state.js";
|
||||
import { diffConfigPaths } from "./config-diff.js";
|
||||
import {
|
||||
|
|
@ -48,24 +46,6 @@ export function isRuntimeSecretsPreparationCurrent(
|
|||
return getActiveSecretsRuntimeSnapshotRevisionState() === preparation.expectedRevision;
|
||||
}
|
||||
|
||||
async function restoreSecretsRuntimeSnapshotIfCurrent(
|
||||
snapshot: PreparedSecretsRuntimeSnapshot,
|
||||
expectedRevision: number,
|
||||
ownedSnapshot: PreparedSecretsRuntimeSnapshot,
|
||||
options?: { onActivated?: () => void; runtimeSourceConfig?: OpenClawConfig },
|
||||
): Promise<boolean> {
|
||||
const runtime = await import("../secrets/runtime.js");
|
||||
if (
|
||||
!runtime.restoreSecretsRuntimeSnapshotIfCurrent(snapshot, expectedRevision, ownedSnapshot, {
|
||||
runtimeSourceConfig: options?.runtimeSourceConfig,
|
||||
})
|
||||
) {
|
||||
return false;
|
||||
}
|
||||
options?.onActivated?.();
|
||||
return true;
|
||||
}
|
||||
|
||||
type PrepareRuntimeCandidate = (
|
||||
runtimeConfig: OpenClawConfig,
|
||||
sourceConfig: OpenClawConfig,
|
||||
|
|
@ -244,7 +224,7 @@ export function createManagedReloadSecretHandlers(options: {
|
|||
const committedSecretsRevision = getActiveSecretsRuntimeSnapshotRevisionState();
|
||||
const rollbackPublishedSource = async () => {
|
||||
if (
|
||||
!(await restoreSecretsRuntimeSnapshotIfCurrent(
|
||||
!(await params.activateRuntimeSecrets.restoreSnapshotIfCurrent(
|
||||
previousSecretsSnapshot,
|
||||
committedSecretsRevision,
|
||||
activated,
|
||||
|
|
@ -345,7 +325,9 @@ export function createManagedReloadSecretHandlers(options: {
|
|||
null;
|
||||
let runtimePolicyReconciled = false;
|
||||
let applicationStatus: Awaited<ReturnType<typeof applyHotReload>>;
|
||||
const rollbackPublication = async () => {
|
||||
const rollbackPublication = async (
|
||||
restore: typeof params.activateRuntimeSecrets.restoreSnapshotIfCurrent,
|
||||
) => {
|
||||
const generationOwnership = publishedSharedGatewaySessionGeneration;
|
||||
if (
|
||||
!runtimeSecretsPublished ||
|
||||
|
|
@ -361,19 +343,12 @@ export function createManagedReloadSecretHandlers(options: {
|
|||
previousSharedGatewaySessionGeneration,
|
||||
);
|
||||
};
|
||||
let snapshotRestored = false;
|
||||
if (previousSnapshot) {
|
||||
snapshotRestored = await restoreSecretsRuntimeSnapshotIfCurrent(
|
||||
previousSnapshot,
|
||||
publishedSnapshotRevision,
|
||||
prepared,
|
||||
{ runtimeSourceConfig: previousRuntimeSourceConfig, onActivated: restoreGeneration },
|
||||
);
|
||||
} else if (getActiveSecretsRuntimeSnapshotRevisionState() === publishedSnapshotRevision) {
|
||||
clearSecretsRuntimeSnapshotState();
|
||||
snapshotRestored = true;
|
||||
restoreGeneration();
|
||||
}
|
||||
const snapshotRestored = await restore(
|
||||
previousSnapshot,
|
||||
publishedSnapshotRevision,
|
||||
prepared,
|
||||
{ runtimeSourceConfig: previousRuntimeSourceConfig, onActivated: restoreGeneration },
|
||||
);
|
||||
if (snapshotRestored) {
|
||||
if (previousSnapshot && shouldRefreshContextWindowCache(plan)) {
|
||||
await refreshContextWindowCache(previousSnapshot.config);
|
||||
|
|
@ -410,7 +385,9 @@ export function createManagedReloadSecretHandlers(options: {
|
|||
throw new GatewayHotReloadStaleSecretsError();
|
||||
}
|
||||
};
|
||||
const publishRuntime = async () => {
|
||||
const publishRuntime = async (
|
||||
restore: typeof params.activateRuntimeSecrets.restoreSnapshotIfCurrent,
|
||||
) => {
|
||||
runtimeSecretsPublished = true;
|
||||
publishedSnapshotRevision = getActiveSecretsRuntimeSnapshotRevisionState();
|
||||
// Claim the generation at the snapshot activation edge, but keep
|
||||
|
|
@ -445,7 +422,7 @@ export function createManagedReloadSecretHandlers(options: {
|
|||
}
|
||||
} catch (err) {
|
||||
if (!isCommitted()) {
|
||||
await rollbackPublication();
|
||||
await rollbackPublication(restore);
|
||||
}
|
||||
throw err;
|
||||
}
|
||||
|
|
@ -500,7 +477,7 @@ export function createManagedReloadSecretHandlers(options: {
|
|||
if (runtimeCommitted) {
|
||||
throw err;
|
||||
}
|
||||
await rollbackPublication();
|
||||
await rollbackPublication(params.activateRuntimeSecrets.restoreSnapshotIfCurrent);
|
||||
throw err;
|
||||
}
|
||||
// Runtime-secret refreshes can legitimately advance the snapshot
|
||||
|
|
|
|||
|
|
@ -60,23 +60,6 @@ export type GatewaySecretsReloaderParams = {
|
|||
logChannels: { info: (message: string) => void };
|
||||
};
|
||||
|
||||
async function restoreSnapshotIfCurrent(
|
||||
snapshot: PreparedSecretsRuntimeSnapshot,
|
||||
expectedRevision: number,
|
||||
ownedSnapshot: PreparedSecretsRuntimeSnapshot,
|
||||
onActivated: () => void,
|
||||
runtimeSourceConfig: OpenClawConfig | undefined,
|
||||
): Promise<void> {
|
||||
const runtime = await import("../secrets/runtime.js");
|
||||
if (
|
||||
runtime.restoreSecretsRuntimeSnapshotIfCurrent(snapshot, expectedRevision, ownedSnapshot, {
|
||||
runtimeSourceConfig,
|
||||
})
|
||||
) {
|
||||
onActivated();
|
||||
}
|
||||
}
|
||||
|
||||
/** Keeps snapshot CAS, generation ownership, and exact account recovery in one transaction. */
|
||||
export function createGatewaySecretsReloader(params: GatewaySecretsReloaderParams) {
|
||||
const buildReloadPlan = params.buildReloadPlan ?? buildGatewayReloadPlan;
|
||||
|
|
@ -359,32 +342,34 @@ export function createGatewaySecretsReloader(params: GatewaySecretsReloaderParam
|
|||
const failedTransaction = transaction;
|
||||
let restoration: SecretsReloadPublication | undefined;
|
||||
try {
|
||||
await restoreSnapshotIfCurrent(
|
||||
await params.activateRuntimeSecrets.restoreSnapshotIfCurrent(
|
||||
failedTransaction.previousSnapshot,
|
||||
failedTransaction.publishedSnapshotRevision,
|
||||
failedTransaction.prepared,
|
||||
() => {
|
||||
const generationRestored = params.sharedGatewaySessionGenerationState.replace(
|
||||
failedTransaction.generationOwnership,
|
||||
{
|
||||
current: failedTransaction.previousGeneration,
|
||||
required: failedTransaction.previousRequiredGeneration,
|
||||
},
|
||||
);
|
||||
if (generationRestored && failedTransaction.generationChanged) {
|
||||
disconnectStaleSharedGatewayAuthClients({
|
||||
state: params.sharedGatewaySessionGenerationState,
|
||||
clients: params.clients,
|
||||
expectedGeneration: failedTransaction.previousGeneration,
|
||||
});
|
||||
}
|
||||
// Restoration can preserve newer credential state; rebuild from what actually won,
|
||||
// not the predecessor snapshot. A newer config publication still fences this tail.
|
||||
restoration = capturePublication(
|
||||
params.sharedGatewaySessionGenerationState.capture(),
|
||||
);
|
||||
{
|
||||
onActivated: () => {
|
||||
const generationRestored = params.sharedGatewaySessionGenerationState.replace(
|
||||
failedTransaction.generationOwnership,
|
||||
{
|
||||
current: failedTransaction.previousGeneration,
|
||||
required: failedTransaction.previousRequiredGeneration,
|
||||
},
|
||||
);
|
||||
if (generationRestored && failedTransaction.generationChanged) {
|
||||
disconnectStaleSharedGatewayAuthClients({
|
||||
state: params.sharedGatewaySessionGenerationState,
|
||||
clients: params.clients,
|
||||
expectedGeneration: failedTransaction.previousGeneration,
|
||||
});
|
||||
}
|
||||
// Restoration can preserve newer credential state; rebuild from what actually won,
|
||||
// not the predecessor snapshot. A newer config publication still fences this tail.
|
||||
restoration = capturePublication(
|
||||
params.sharedGatewaySessionGenerationState.capture(),
|
||||
);
|
||||
},
|
||||
runtimeSourceConfig: failedTransaction.previousRuntimeSourceConfig,
|
||||
},
|
||||
failedTransaction.previousRuntimeSourceConfig,
|
||||
);
|
||||
await restoration?.modelPublication;
|
||||
} catch {
|
||||
|
|
|
|||
|
|
@ -258,6 +258,11 @@ export async function prepareGatewayServerBootstrap(input: {
|
|||
const activateRuntimeSecrets = createRuntimeSecretsActivator({
|
||||
logSecrets,
|
||||
emitStateEvent: emitSecretsStateEvent,
|
||||
beforeSnapshotPublication: async (config) => {
|
||||
const { publishCanonicalUserChannelPolicy } =
|
||||
await import("../state/user-channel-identity-operations.js");
|
||||
await publishCanonicalUserChannelPolicy(config?.gateway);
|
||||
},
|
||||
...(startupConfigLoad.pluginMetadataSnapshot
|
||||
? { pluginMetadataSnapshot: startupConfigLoad.pluginMetadataSnapshot }
|
||||
: {}),
|
||||
|
|
@ -265,7 +270,7 @@ export async function prepareGatewayServerBootstrap(input: {
|
|||
const startupActivationSourceConfig = configSnapshot.sourceConfig;
|
||||
const startupRuntimeConfig = captureConfigOverrideApplier()(startupConfigSnapshot.config);
|
||||
startupTrace.setConfig(startupRuntimeConfig);
|
||||
const { prepareGatewayStartupConfig } = await startupConfigModulePromise;
|
||||
const { prepareGatewayStartupConfig } = await import("./server-startup-config-helpers.js");
|
||||
const authBootstrap = await startupTrace.measure(
|
||||
"config.auth",
|
||||
() =>
|
||||
|
|
|
|||
|
|
@ -1,4 +1,5 @@
|
|||
// Shared validation, auth-surface, and config-load helpers for Gateway startup.
|
||||
import { isDeepStrictEqual } from "node:util";
|
||||
import {
|
||||
formatInvalidConfigRecoveryHint,
|
||||
formatPluginPackagingRuntimeOutputRecoveryHint,
|
||||
|
|
@ -10,6 +11,7 @@ import {
|
|||
} from "../config/io.js";
|
||||
import { renderConfigValidationIssueLines } from "../config/issue-location.js";
|
||||
import {
|
||||
inheritLegacyDefaultAgentId,
|
||||
retainLegacyDefaultAgentId,
|
||||
tryGetLegacyDefaultAgentId,
|
||||
} from "../config/legacy.default-agent-owner.js";
|
||||
|
|
@ -20,7 +22,9 @@ import { isPluginPackagingRuntimeOutputInvalidConfigSnapshot } from "../config/r
|
|||
import {
|
||||
copyConfigResolutionFacts,
|
||||
copyConfigResolutionFactsExcept,
|
||||
hasUnresolvedConfigPath,
|
||||
} from "../config/resolution-facts.js";
|
||||
import { applyConfigOverrides } from "../config/runtime-overrides.js";
|
||||
import type { GatewayAuthConfig, GatewayTailscaleConfig } from "../config/types.gateway.js";
|
||||
import type { ConfigFileSnapshot, OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import type { PluginMetadataSnapshot } from "../plugins/plugin-metadata-snapshot.js";
|
||||
|
|
@ -31,7 +35,13 @@ import {
|
|||
import { resolveGatewayAuthForConfig } from "./auth-resolve.js";
|
||||
import { assertGatewayAuthNotKnownWeak } from "./known-weak-gateway-secrets.js";
|
||||
import { mergeActivationSectionsIntoRuntimeConfig } from "./plugin-activation-runtime-config.js";
|
||||
import { mergeGatewayAuthConfig, mergeGatewayTailscaleConfig } from "./startup-auth.js";
|
||||
import type { ActivateRuntimeSecrets } from "./server-startup-config.js";
|
||||
import { resolveGatewayStartupSourceConfig } from "./server-startup-secret-surfaces.js";
|
||||
import {
|
||||
ensureGatewayStartupAuth,
|
||||
mergeGatewayAuthConfig,
|
||||
mergeGatewayTailscaleConfig,
|
||||
} from "./startup-auth.js";
|
||||
|
||||
export type GatewayStartupLog = {
|
||||
info: (message: string) => void;
|
||||
|
|
@ -51,7 +61,7 @@ export type GatewayStartupConfigSnapshotLoadResult = {
|
|||
};
|
||||
|
||||
/** Throw a formatted startup error when the loaded config snapshot is invalid. */
|
||||
export function assertValidGatewayStartupConfigSnapshot(
|
||||
function assertValidGatewayStartupConfigSnapshot(
|
||||
snapshot: ConfigFileSnapshot,
|
||||
options: { includeDoctorHint?: boolean } = {},
|
||||
): void {
|
||||
|
|
@ -231,3 +241,109 @@ export function applyGatewayAuthOverridesForStartupPreflight(
|
|||
]);
|
||||
return next;
|
||||
}
|
||||
|
||||
/** Prepare the effective Gateway startup config after auth, overrides, and secrets activation. */
|
||||
export async function prepareGatewayStartupConfig(params: {
|
||||
configSnapshot: ConfigFileSnapshot;
|
||||
authOverride?: GatewayAuthConfig;
|
||||
tailscaleOverride?: GatewayTailscaleConfig;
|
||||
activateRuntimeSecrets: ActivateRuntimeSecrets;
|
||||
log?: GatewayStartupLog;
|
||||
measure?: GatewayStartupConfigMeasure;
|
||||
}): Promise<Awaited<ReturnType<typeof ensureGatewayStartupAuth>>> {
|
||||
const measure = params.measure ?? (async (_name, run) => await run());
|
||||
await measure("config.auth.snapshot-validate", () =>
|
||||
assertValidGatewayStartupConfigSnapshot(params.configSnapshot),
|
||||
);
|
||||
|
||||
const runtimeConfig = await measure("config.auth.runtime-overrides", () =>
|
||||
applyConfigOverrides(params.configSnapshot.config),
|
||||
);
|
||||
copyConfigResolutionFacts(params.configSnapshot.config, runtimeConfig);
|
||||
const startupPreflightConfig = await measure("config.auth.startup-overrides", () =>
|
||||
applyGatewayAuthOverridesForStartupPreflight(runtimeConfig, {
|
||||
auth: params.authOverride,
|
||||
tailscale: params.tailscaleOverride,
|
||||
}),
|
||||
);
|
||||
const needsAuthSecretPreflight = await measure("config.auth.secret-surface", () =>
|
||||
hasActiveGatewayAuthSecretRef(startupPreflightConfig),
|
||||
);
|
||||
let preflightPrepared: Awaited<ReturnType<ActivateRuntimeSecrets>> | undefined;
|
||||
const preflightConfig = await measure(
|
||||
"config.auth.secret-preflight",
|
||||
async () => {
|
||||
if (!needsAuthSecretPreflight) {
|
||||
return startupPreflightConfig;
|
||||
}
|
||||
preflightPrepared = await params.activateRuntimeSecrets(startupPreflightConfig, {
|
||||
reason: "startup",
|
||||
activate: false,
|
||||
});
|
||||
return preflightPrepared.config;
|
||||
},
|
||||
{ omitErrorMessage: true },
|
||||
);
|
||||
const activateStartupSecrets = async (config: OpenClawConfig) => {
|
||||
// Reuse the preflight snapshot only if generated startup auth did not
|
||||
// change the secret-relevant source config.
|
||||
if (
|
||||
preflightPrepared &&
|
||||
isDeepStrictEqual(
|
||||
resolveGatewayStartupSourceConfig(config, process.env),
|
||||
preflightPrepared.sourceConfig,
|
||||
)
|
||||
) {
|
||||
return await params.activateRuntimeSecrets.activatePreparedSnapshot(preflightPrepared, {
|
||||
reason: "startup",
|
||||
activate: true,
|
||||
});
|
||||
}
|
||||
return await params.activateRuntimeSecrets(config, {
|
||||
reason: "startup",
|
||||
activate: true,
|
||||
});
|
||||
};
|
||||
const preflightAuthOverride = await measure("config.auth.preflight-override", () => {
|
||||
const token = preflightConfig.gateway?.auth?.token;
|
||||
const password = preflightConfig.gateway?.auth?.password;
|
||||
const resolvedToken =
|
||||
typeof token === "string" && !hasUnresolvedConfigPath(preflightConfig, "gateway.auth.token");
|
||||
const resolvedPassword =
|
||||
typeof password === "string" &&
|
||||
!hasUnresolvedConfigPath(preflightConfig, "gateway.auth.password");
|
||||
return resolvedToken || resolvedPassword
|
||||
? {
|
||||
...params.authOverride,
|
||||
...(resolvedToken ? { token } : {}),
|
||||
...(resolvedPassword ? { password } : {}),
|
||||
}
|
||||
: params.authOverride;
|
||||
});
|
||||
|
||||
const authBootstrap = await measure("config.auth.ensure", () =>
|
||||
ensureGatewayStartupAuth({
|
||||
cfg: runtimeConfig,
|
||||
env: process.env,
|
||||
authOverride: preflightAuthOverride,
|
||||
tailscaleOverride: params.tailscaleOverride,
|
||||
warn: params.log?.warn,
|
||||
}),
|
||||
);
|
||||
const runtimeStartupConfig = await measure("config.auth.runtime-startup-overrides", () =>
|
||||
applyGatewayAuthOverridesForStartupPreflight(authBootstrap.cfg, {
|
||||
auth: params.authOverride,
|
||||
tailscale: params.tailscaleOverride,
|
||||
}),
|
||||
);
|
||||
const activatedConfig = (
|
||||
await measure(
|
||||
"config.auth.secrets-activate",
|
||||
() => activateStartupSecrets(runtimeStartupConfig),
|
||||
{ omitErrorMessage: true },
|
||||
)
|
||||
).config;
|
||||
const config = inheritLegacyDefaultAgentId(params.configSnapshot.config, activatedConfig);
|
||||
copyConfigResolutionFacts(activatedConfig, config);
|
||||
return { ...authBootstrap, cfg: config };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -18,10 +18,8 @@ import {
|
|||
prepareSecretsRuntimeSnapshot,
|
||||
} from "../secrets/runtime.js";
|
||||
import { withEnvAsync } from "../test-utils/env.js";
|
||||
import {
|
||||
createRuntimeSecretsActivator,
|
||||
prepareGatewayStartupConfig,
|
||||
} from "./server-startup-config.js";
|
||||
import { prepareGatewayStartupConfig } from "./server-startup-config-helpers.js";
|
||||
import { createRuntimeSecretsActivator } from "./server-startup-config.js";
|
||||
import { buildTestConfigSnapshot } from "./test-helpers.config-snapshots.js";
|
||||
|
||||
const GATEWAY_TOKEN_ENV = "BREAKER_GATEWAY_AUTH_TOKEN";
|
||||
|
|
|
|||
|
|
@ -34,10 +34,8 @@ import {
|
|||
} from "../secrets/runtime-state.js";
|
||||
import type { PreparedSecretsRuntimeSnapshot, SecretResolverWarning } from "../secrets/runtime.js";
|
||||
import { withEnvAsync } from "../test-utils/env.js";
|
||||
import {
|
||||
createRuntimeSecretsActivator,
|
||||
prepareGatewayStartupConfig,
|
||||
} from "./server-startup-config.js";
|
||||
import { prepareGatewayStartupConfig } from "./server-startup-config-helpers.js";
|
||||
import { createRuntimeSecretsActivator } from "./server-startup-config.js";
|
||||
import { buildTestConfigSnapshot } from "./test-helpers.config-snapshots.js";
|
||||
|
||||
const KNOWN_WEAK_GATEWAY_TOKEN_PLACEHOLDERS = [
|
||||
|
|
@ -332,7 +330,15 @@ function createGatewayStartupSecretsRuntimeHarness(prefix: string) {
|
|||
vi.resetModules();
|
||||
const agentDir = mkdtempSync(path.join(tmpdir(), prefix));
|
||||
const runtimeImport = vi.fn();
|
||||
const prepareRuntimeSecretsSnapshot = vi.fn(async ({ config }) => preparedSnapshot(config));
|
||||
const prepareRuntimeSecretsSnapshot = vi.fn(async ({ config }) => {
|
||||
// Import-order fixtures must capture revisions from the same module generation as activation.
|
||||
const revisions = await import("../agents/auth-profiles/runtime-snapshots.js");
|
||||
return {
|
||||
...preparedSnapshot(config),
|
||||
authStoreCredentialsRevision: revisions.getRuntimeAuthProfileStoreCredentialsRevision(),
|
||||
authStoreSnapshotsRevision: revisions.getRuntimeAuthProfileStoreSnapshotsRevision(),
|
||||
};
|
||||
});
|
||||
const activateRuntimeSecretsSnapshot = vi.fn();
|
||||
return {
|
||||
activateRuntimeSecretsSnapshot,
|
||||
|
|
@ -354,7 +360,7 @@ function createGatewayStartupSecretsRuntimeHarness(prefix: string) {
|
|||
};
|
||||
}
|
||||
|
||||
async function activateImportedStartupConfig(config: OpenClawConfig) {
|
||||
async function activateImportedStartupConfig(config: OpenClawConfig, env?: NodeJS.ProcessEnv) {
|
||||
const { createRuntimeSecretsActivator: createActivator } =
|
||||
await import("./server-startup-config.js");
|
||||
return await createActivator(runtimeSecretsActivatorOptionsForTest())(
|
||||
|
|
@ -362,21 +368,11 @@ async function activateImportedStartupConfig(config: OpenClawConfig) {
|
|||
{
|
||||
reason: "startup",
|
||||
activate: true,
|
||||
env,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
async function activateStartupConfigWithEnv(config: OpenClawConfig, env: NodeJS.ProcessEnv) {
|
||||
const activateRuntimeSecrets = createRuntimeSecretsActivator(
|
||||
runtimeSecretsActivatorOptionsForTest(),
|
||||
);
|
||||
return await activateRuntimeSecrets(gatewayTokenConfig(config), {
|
||||
reason: "startup",
|
||||
activate: true,
|
||||
env,
|
||||
});
|
||||
}
|
||||
|
||||
function writePersistedOpenAiProfile(agentDir: string, key: string): void {
|
||||
writePersistedAuthProfileStoreRaw(
|
||||
createAuthProfileStoreFixture({
|
||||
|
|
@ -3087,7 +3083,7 @@ describe("gateway startup config secret preflight", () => {
|
|||
};
|
||||
|
||||
try {
|
||||
await activateStartupConfigWithEnv(
|
||||
await activateImportedStartupConfig(
|
||||
{ agents: { list: [{ id: "main", agentDir: relocatedMainAgentDir }] } },
|
||||
activationEnv,
|
||||
);
|
||||
|
|
@ -3128,7 +3124,7 @@ describe("gateway startup config secret preflight", () => {
|
|||
};
|
||||
|
||||
try {
|
||||
await activateStartupConfigWithEnv(
|
||||
await activateImportedStartupConfig(
|
||||
{
|
||||
agents: {
|
||||
list: [{ id: "main", agentDir: "~/configured-agent" }],
|
||||
|
|
|
|||
|
|
@ -1,12 +1,8 @@
|
|||
// Gateway startup config loads, repairs, validates, and activates runtime config
|
||||
// plus secrets snapshots before the server exposes user-facing surfaces.
|
||||
import { isDeepStrictEqual } from "node:util";
|
||||
import { hasLegacyAuthProfileSourcesForStartup } from "../agents/auth-profiles/legacy-source-diagnostic.js";
|
||||
import { inheritLegacyDefaultAgentId } from "../config/legacy.default-agent-owner.js";
|
||||
import { copyConfigResolutionFacts, hasUnresolvedConfigPath } from "../config/resolution-facts.js";
|
||||
import { applyConfigOverrides } from "../config/runtime-overrides.js";
|
||||
import type { GatewayAuthConfig, GatewayTailscaleConfig } from "../config/types.gateway.js";
|
||||
import type { ConfigFileSnapshot, OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { getRuntimeAuthProfileStoreSnapshotsRevision } from "../agents/auth-profiles/runtime-snapshots.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { measureDiagnosticsTimelineSpan } from "../infra/diagnostics-timeline.js";
|
||||
import type { PluginManifestRegistry } from "../plugins/manifest-registry.js";
|
||||
import type { PluginMetadataSnapshot } from "../plugins/plugin-metadata-snapshot.js";
|
||||
|
|
@ -32,6 +28,7 @@ import {
|
|||
} from "../secrets/runtime-provider-auth-scope.js";
|
||||
import {
|
||||
activateSecretsRuntimeSnapshotState,
|
||||
clearSecretsRuntimeSnapshotState,
|
||||
graftActiveSecretsRuntimeAuthState,
|
||||
getActiveSecretsRuntimeSnapshotState,
|
||||
getActiveSecretsRuntimeSnapshotRevisionState,
|
||||
|
|
@ -43,12 +40,9 @@ import {
|
|||
import { logRuntimeSecretWarnings } from "../secrets/runtime-warning-log.js";
|
||||
import { createLazyPromise } from "../shared/lazy-runtime.js";
|
||||
import {
|
||||
applyGatewayAuthOverridesForStartupPreflight,
|
||||
assertRuntimeGatewayAuthNotKnownWeak,
|
||||
assertValidGatewayStartupConfigSnapshot,
|
||||
hasActiveGatewayAuthSecretRef,
|
||||
logGatewayAuthSurfaceDiagnostics,
|
||||
type GatewayStartupConfigMeasure,
|
||||
type GatewayStartupLog,
|
||||
} from "./server-startup-config-helpers.js";
|
||||
import {
|
||||
|
|
@ -56,7 +50,6 @@ import {
|
|||
logThrownSecretDegradations,
|
||||
} from "./server-startup-secret-diagnostics.js";
|
||||
import { resolveGatewayStartupSourceConfig } from "./server-startup-secret-surfaces.js";
|
||||
import { ensureGatewayStartupAuth } from "./startup-auth.js";
|
||||
export {
|
||||
applyGatewayAuthOverridesForStartupPreflight,
|
||||
loadGatewayStartupConfigSnapshot,
|
||||
|
|
@ -107,10 +100,18 @@ export type ActivateRuntimeSecrets = ((
|
|||
snapshot: PreparedRuntimeSecretsSnapshot,
|
||||
expectedRevision: number,
|
||||
params: RuntimeSecretsActivationParams,
|
||||
onActivated?: () => void | Promise<void>,
|
||||
onActivated?: (
|
||||
restore: ActivateRuntimeSecrets["restoreSnapshotIfCurrent"],
|
||||
) => void | Promise<void>,
|
||||
canActivate?: () => boolean,
|
||||
checkpoint?: () => Promise<void>,
|
||||
) => Promise<PreparedRuntimeSecretsSnapshot | null>;
|
||||
restoreSnapshotIfCurrent: (
|
||||
snapshot: PreparedRuntimeSecretsSnapshot | null,
|
||||
expectedRevision: number,
|
||||
ownedSnapshot: PreparedRuntimeSecretsSnapshot,
|
||||
options?: { onActivated?: () => void; runtimeSourceConfig?: OpenClawConfig },
|
||||
) => Promise<boolean>;
|
||||
publishStateTransition: (
|
||||
snapshot: PreparedRuntimeSecretsSnapshot,
|
||||
options?: { sourceOnly?: boolean; expectedRevision?: number },
|
||||
|
|
@ -127,6 +128,8 @@ export function createRuntimeSecretsActivator(params: {
|
|||
) => void;
|
||||
prepareRuntimeSecretsSnapshot?: PrepareRuntimeSecretsSnapshot;
|
||||
activateRuntimeSecretsSnapshot?: ActivateRuntimeSecretsSnapshot;
|
||||
/** Commit owner-held durable policy before publication; also reconcile the survivor on failure. */
|
||||
beforeSnapshotPublication?: (config: OpenClawConfig | null) => Promise<void>;
|
||||
manifestRegistry?: Pick<PluginManifestRegistry, "plugins">;
|
||||
pluginMetadataSnapshot?: Pick<PluginMetadataSnapshot, "plugins" | "manifestRegistry">;
|
||||
}): ActivateRuntimeSecrets {
|
||||
|
|
@ -158,13 +161,83 @@ export function createRuntimeSecretsActivator(params: {
|
|||
return await run;
|
||||
};
|
||||
|
||||
const loadActivateRuntimeSecretsSnapshot = async () => {
|
||||
const loadActivateRuntimeSecretsSnapshot = async (source?: OpenClawConfig) => {
|
||||
if (source) {
|
||||
const runtime = await loadSecretsRuntime();
|
||||
return (snapshot: PreparedRuntimeSecretsSnapshot) =>
|
||||
runtime.activateSecretsRuntimeSnapshotWithSource(snapshot, source);
|
||||
}
|
||||
if (params.activateRuntimeSecretsSnapshot) {
|
||||
return params.activateRuntimeSecretsSnapshot;
|
||||
}
|
||||
return (await loadSecretsRuntime()).activateSecretsRuntimeSnapshot;
|
||||
};
|
||||
|
||||
const supersededActivation = new Error("Secrets runtime publication was superseded.");
|
||||
const publishSnapshot = async (
|
||||
config: OpenClawConfig | null,
|
||||
isCurrent: () => boolean,
|
||||
publish: () => void,
|
||||
): Promise<boolean> => {
|
||||
if (!isCurrent()) {
|
||||
return false;
|
||||
}
|
||||
let published = false;
|
||||
try {
|
||||
await params.beforeSnapshotPublication?.(config);
|
||||
if (!isCurrent()) {
|
||||
return false;
|
||||
}
|
||||
publish();
|
||||
published = true;
|
||||
return true;
|
||||
} finally {
|
||||
if (!published) {
|
||||
await params.beforeSnapshotPublication?.(
|
||||
getActiveSecretsRuntimeSnapshotState()?.config ?? null,
|
||||
);
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
const restoreSnapshot = async (
|
||||
isOwned: () => boolean,
|
||||
...[snapshot, expectedRevision, ownedSnapshot, options]: Parameters<
|
||||
ActivateRuntimeSecrets["restoreSnapshotIfCurrent"]
|
||||
>
|
||||
) => {
|
||||
const restoration = snapshot
|
||||
? (await loadSecretsRuntime()).prepareSecretsRuntimeSnapshotRestore(
|
||||
snapshot,
|
||||
expectedRevision,
|
||||
ownedSnapshot,
|
||||
options,
|
||||
)
|
||||
: null;
|
||||
if (snapshot && !restoration) {
|
||||
return false;
|
||||
}
|
||||
return await publishSnapshot(
|
||||
restoration?.snapshot.config ?? null,
|
||||
() =>
|
||||
isOwned() &&
|
||||
getActiveSecretsRuntimeSnapshotRevisionState() ===
|
||||
(restoration?.expectedRevision ?? expectedRevision) &&
|
||||
(!restoration ||
|
||||
(hasCurrentAuthStoreCredentialsRevision(restoration.snapshot) &&
|
||||
getRuntimeAuthProfileStoreSnapshotsRevision() ===
|
||||
restoration.snapshot.authStoreSnapshotsRevision)),
|
||||
() => {
|
||||
if (restoration) {
|
||||
activateSecretsRuntimeSnapshotState(restoration);
|
||||
} else {
|
||||
clearSecretsRuntimeSnapshotState();
|
||||
}
|
||||
options?.onActivated?.();
|
||||
},
|
||||
);
|
||||
};
|
||||
|
||||
const publishRecovery = (
|
||||
config: OpenClawConfig,
|
||||
expectedGeneration?: number,
|
||||
|
|
@ -226,6 +299,7 @@ export function createRuntimeSecretsActivator(params: {
|
|||
options?: {
|
||||
activateRuntimeSecretsSnapshot?: (snapshot: PreparedRuntimeSecretsSnapshot) => void;
|
||||
onActivated?: () => void;
|
||||
canActivate?: () => boolean;
|
||||
alreadyActivated?: boolean;
|
||||
stateScope?: SecretsStateScope;
|
||||
stateDegradedOwners?: PreparedRuntimeSecretsSnapshot["degradedOwners"];
|
||||
|
|
@ -234,13 +308,26 @@ export function createRuntimeSecretsActivator(params: {
|
|||
assertRuntimeGatewayAuthNotKnownWeak(prepared.config);
|
||||
if (activationParams.activate && !options?.alreadyActivated) {
|
||||
const activateRuntimeSecretsSnapshot =
|
||||
options?.activateRuntimeSecretsSnapshot ?? (await loadActivateRuntimeSecretsSnapshot());
|
||||
activateRuntimeSecretsSnapshot(prepared);
|
||||
options?.activateRuntimeSecretsSnapshot ??
|
||||
(await loadActivateRuntimeSecretsSnapshot(activationParams.runtimeSourceConfig));
|
||||
const revision = getActiveSecretsRuntimeSnapshotRevisionState();
|
||||
if (
|
||||
!(await publishSnapshot(
|
||||
prepared.config,
|
||||
() =>
|
||||
getActiveSecretsRuntimeSnapshotRevisionState() === revision &&
|
||||
hasCurrentAuthStoreCredentialsRevision(prepared) &&
|
||||
(options?.canActivate?.() ?? true),
|
||||
() => {
|
||||
activateRuntimeSecretsSnapshot(prepared);
|
||||
options?.onActivated?.();
|
||||
},
|
||||
))
|
||||
) {
|
||||
throw supersededActivation;
|
||||
}
|
||||
}
|
||||
if (activationParams.activate) {
|
||||
// Invoke publication at the activation edge so no microtask can replace
|
||||
// the candidate before its runtime commit begins.
|
||||
options?.onActivated?.();
|
||||
logGatewayAuthSurfaceDiagnostics(prepared, params.logSecrets);
|
||||
}
|
||||
logRuntimeSecretWarnings({
|
||||
|
|
@ -462,23 +549,8 @@ export function createRuntimeSecretsActivator(params: {
|
|||
|
||||
const activatePreparedSnapshotIfCurrent: ActivateRuntimeSecrets["activatePreparedSnapshotIfCurrent"] =
|
||||
async (snapshot, expectedRevision, activationParams, onActivated, canActivate, checkpoint) => {
|
||||
// Resolve the lazy activator before entering the compare-and-activate
|
||||
// section so no await separates revision ownership from state publication.
|
||||
const runtimeSourceConfig = activationParams.runtimeSourceConfig;
|
||||
const activateRuntimeSecretsSnapshot = activationParams.activate
|
||||
? runtimeSourceConfig
|
||||
? (
|
||||
(runtime) => (preparedSnapshot: PreparedRuntimeSecretsSnapshot) =>
|
||||
runtime.activateSecretsRuntimeSnapshotWithSource(
|
||||
preparedSnapshot,
|
||||
runtimeSourceConfig,
|
||||
)
|
||||
)(await loadSecretsRuntime())
|
||||
: await loadActivateRuntimeSecretsSnapshot()
|
||||
: undefined;
|
||||
return await runWithSecretsActivationLock(async () => {
|
||||
// Resolve source observations inside the lock, then recheck every revision.
|
||||
// No await may separate these final guards from activation/publication.
|
||||
// Source observations and durable preparation share activation ownership.
|
||||
await checkpoint?.();
|
||||
if (
|
||||
getActiveSecretsRuntimeSnapshotRevisionState() !== expectedRevision ||
|
||||
|
|
@ -489,28 +561,33 @@ export function createRuntimeSecretsActivator(params: {
|
|||
}
|
||||
let activated: PreparedRuntimeSecretsSnapshot;
|
||||
let publication: Promise<void> | undefined;
|
||||
let callbackOpen = true;
|
||||
try {
|
||||
activated = await finishPreparedSnapshot(
|
||||
snapshot,
|
||||
activationParams,
|
||||
activateRuntimeSecretsSnapshot
|
||||
? {
|
||||
activateRuntimeSecretsSnapshot,
|
||||
...(onActivated
|
||||
? {
|
||||
onActivated: () => {
|
||||
publication = Promise.resolve(onActivated());
|
||||
},
|
||||
}
|
||||
: {}),
|
||||
activated = await finishPreparedSnapshot(snapshot, activationParams, {
|
||||
canActivate: () =>
|
||||
getActiveSecretsRuntimeSnapshotRevisionState() === expectedRevision &&
|
||||
(canActivate?.() ?? true),
|
||||
onActivated: onActivated
|
||||
? () => {
|
||||
publication = Promise.resolve(
|
||||
onActivated((...args) => restoreSnapshot(() => callbackOpen, ...args)),
|
||||
);
|
||||
}
|
||||
: undefined,
|
||||
);
|
||||
});
|
||||
} catch (err) {
|
||||
callbackOpen = false;
|
||||
if (err === supersededActivation) {
|
||||
return null;
|
||||
}
|
||||
return handleSecretsActivationError(err, activationParams, snapshot.sourceConfig);
|
||||
}
|
||||
await publication;
|
||||
return activated;
|
||||
try {
|
||||
await publication;
|
||||
return activated;
|
||||
} finally {
|
||||
callbackOpen = false;
|
||||
}
|
||||
});
|
||||
};
|
||||
|
||||
|
|
@ -604,112 +681,9 @@ export function createRuntimeSecretsActivator(params: {
|
|||
return Object.assign(prepareRuntimeSecrets, {
|
||||
activatePreparedSnapshot,
|
||||
activatePreparedSnapshotIfCurrent,
|
||||
restoreSnapshotIfCurrent: (
|
||||
...args: Parameters<ActivateRuntimeSecrets["restoreSnapshotIfCurrent"]>
|
||||
) => runWithSecretsActivationLock(() => restoreSnapshot(() => true, ...args)),
|
||||
publishStateTransition,
|
||||
});
|
||||
}
|
||||
|
||||
/** Prepare the effective Gateway startup config after auth, overrides, and secrets activation. */
|
||||
export async function prepareGatewayStartupConfig(params: {
|
||||
configSnapshot: ConfigFileSnapshot;
|
||||
authOverride?: GatewayAuthConfig;
|
||||
tailscaleOverride?: GatewayTailscaleConfig;
|
||||
activateRuntimeSecrets: ActivateRuntimeSecrets;
|
||||
log?: GatewayStartupLog;
|
||||
measure?: GatewayStartupConfigMeasure;
|
||||
}): Promise<Awaited<ReturnType<typeof ensureGatewayStartupAuth>>> {
|
||||
const measure = params.measure ?? (async (_name, run) => await run());
|
||||
await measure("config.auth.snapshot-validate", () =>
|
||||
assertValidGatewayStartupConfigSnapshot(params.configSnapshot),
|
||||
);
|
||||
|
||||
const runtimeConfig = await measure("config.auth.runtime-overrides", () =>
|
||||
applyConfigOverrides(params.configSnapshot.config),
|
||||
);
|
||||
copyConfigResolutionFacts(params.configSnapshot.config, runtimeConfig);
|
||||
const startupPreflightConfig = await measure("config.auth.startup-overrides", () =>
|
||||
applyGatewayAuthOverridesForStartupPreflight(runtimeConfig, {
|
||||
auth: params.authOverride,
|
||||
tailscale: params.tailscaleOverride,
|
||||
}),
|
||||
);
|
||||
const needsAuthSecretPreflight = await measure("config.auth.secret-surface", () =>
|
||||
hasActiveGatewayAuthSecretRef(startupPreflightConfig),
|
||||
);
|
||||
let preflightPrepared: PreparedRuntimeSecretsSnapshot | undefined;
|
||||
const preflightConfig = await measure(
|
||||
"config.auth.secret-preflight",
|
||||
async () => {
|
||||
if (!needsAuthSecretPreflight) {
|
||||
return startupPreflightConfig;
|
||||
}
|
||||
preflightPrepared = await params.activateRuntimeSecrets(startupPreflightConfig, {
|
||||
reason: "startup",
|
||||
activate: false,
|
||||
});
|
||||
return preflightPrepared.config;
|
||||
},
|
||||
{ omitErrorMessage: true },
|
||||
);
|
||||
const activateStartupSecrets = async (config: OpenClawConfig) => {
|
||||
// Reuse the preflight snapshot only if generated startup auth did not
|
||||
// change the secret-relevant source config.
|
||||
if (
|
||||
preflightPrepared &&
|
||||
isDeepStrictEqual(
|
||||
resolveGatewayStartupSourceConfig(config, process.env),
|
||||
preflightPrepared.sourceConfig,
|
||||
)
|
||||
) {
|
||||
return await params.activateRuntimeSecrets.activatePreparedSnapshot(preflightPrepared, {
|
||||
reason: "startup",
|
||||
activate: true,
|
||||
});
|
||||
}
|
||||
return await params.activateRuntimeSecrets(config, {
|
||||
reason: "startup",
|
||||
activate: true,
|
||||
});
|
||||
};
|
||||
const preflightAuthOverride = await measure("config.auth.preflight-override", () => {
|
||||
const token = preflightConfig.gateway?.auth?.token;
|
||||
const password = preflightConfig.gateway?.auth?.password;
|
||||
const resolvedToken =
|
||||
typeof token === "string" && !hasUnresolvedConfigPath(preflightConfig, "gateway.auth.token");
|
||||
const resolvedPassword =
|
||||
typeof password === "string" &&
|
||||
!hasUnresolvedConfigPath(preflightConfig, "gateway.auth.password");
|
||||
return resolvedToken || resolvedPassword
|
||||
? {
|
||||
...params.authOverride,
|
||||
...(resolvedToken ? { token } : {}),
|
||||
...(resolvedPassword ? { password } : {}),
|
||||
}
|
||||
: params.authOverride;
|
||||
});
|
||||
|
||||
const authBootstrap = await measure("config.auth.ensure", () =>
|
||||
ensureGatewayStartupAuth({
|
||||
cfg: runtimeConfig,
|
||||
env: process.env,
|
||||
authOverride: preflightAuthOverride,
|
||||
tailscaleOverride: params.tailscaleOverride,
|
||||
warn: params.log?.warn,
|
||||
}),
|
||||
);
|
||||
const runtimeStartupConfig = await measure("config.auth.runtime-startup-overrides", () =>
|
||||
applyGatewayAuthOverridesForStartupPreflight(authBootstrap.cfg, {
|
||||
auth: params.authOverride,
|
||||
tailscale: params.tailscaleOverride,
|
||||
}),
|
||||
);
|
||||
const activatedConfig = (
|
||||
await measure(
|
||||
"config.auth.secrets-activate",
|
||||
() => activateStartupSecrets(runtimeStartupConfig),
|
||||
{ omitErrorMessage: true },
|
||||
)
|
||||
).config;
|
||||
const config = inheritLegacyDefaultAgentId(params.configSnapshot.config, activatedConfig);
|
||||
copyConfigResolutionFacts(activatedConfig, config);
|
||||
return { ...authBootstrap, cfg: config };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -48,7 +48,7 @@ import {
|
|||
getActiveSecretsRuntimeSnapshotRevisionState,
|
||||
hasSameSecretReloadContract,
|
||||
restoreSecretsRuntimeSourceSnapshotIfLineageCurrent,
|
||||
restoreSecretsRuntimeSnapshotStateIfCurrent,
|
||||
prepareSecretsRuntimeSnapshotRestoreState,
|
||||
setSecretsRuntimeSourceSnapshotIfCurrent,
|
||||
type PreparedSecretsRuntimeSnapshot,
|
||||
} from "./runtime-state.js";
|
||||
|
|
@ -155,17 +155,12 @@ function activateSnapshotIfCurrent(
|
|||
});
|
||||
}
|
||||
|
||||
type RestoreIfCurrentOptions = Omit<
|
||||
Parameters<typeof restoreSecretsRuntimeSnapshotStateIfCurrent>[0],
|
||||
"snapshot" | "ownedSnapshot" | "expectedRevision" | "refreshContext" | "refreshHandler"
|
||||
> & { expectedRevision?: number };
|
||||
|
||||
function restoreSnapshotIfCurrent(
|
||||
snapshot: PreparedSecretsRuntimeSnapshot,
|
||||
ownedSnapshot: PreparedSecretsRuntimeSnapshot,
|
||||
options: RestoreIfCurrentOptions = {},
|
||||
options: ActivateIfCurrentOptions = {},
|
||||
): boolean {
|
||||
return restoreSecretsRuntimeSnapshotStateIfCurrent({
|
||||
const restoration = prepareSecretsRuntimeSnapshotRestoreState({
|
||||
snapshot,
|
||||
ownedSnapshot,
|
||||
expectedRevision: options.expectedRevision ?? getActiveSecretsRuntimeSnapshotRevisionState(),
|
||||
|
|
@ -173,6 +168,7 @@ function restoreSnapshotIfCurrent(
|
|||
refreshHandler: null,
|
||||
...options,
|
||||
});
|
||||
return restoration !== null && activateSecretsRuntimeSnapshotStateIfCurrent(restoration);
|
||||
}
|
||||
|
||||
describe("secrets runtime state", () => {
|
||||
|
|
|
|||
|
|
@ -982,15 +982,15 @@ export function activateSecretsRuntimeSnapshotStateIfCurrent(
|
|||
return true;
|
||||
}
|
||||
|
||||
/** Restores an owned predecessor while retaining changes after candidate preparation. */
|
||||
export function restoreSecretsRuntimeSnapshotStateIfCurrent(
|
||||
/** Computes the owned predecessor while retaining changes after candidate preparation. */
|
||||
export function prepareSecretsRuntimeSnapshotRestoreState(
|
||||
params: Parameters<typeof activateSecretsRuntimeSnapshotState>[0] & {
|
||||
expectedRevision: number;
|
||||
ownedSnapshot: PreparedSecretsRuntimeSnapshot;
|
||||
},
|
||||
): boolean {
|
||||
) {
|
||||
if (!activeSnapshot || activeSnapshotLineageStartRevision !== params.expectedRevision) {
|
||||
return false;
|
||||
return null;
|
||||
}
|
||||
const currentEntries = listOwnedRuntimeAuthProfileStoreSnapshots();
|
||||
// A later owner is outside this activation's rollback authority, even when its bytes match.
|
||||
|
|
@ -1055,7 +1055,7 @@ export function restoreSecretsRuntimeSnapshotStateIfCurrent(
|
|||
restoredSourceConfig,
|
||||
activeSnapshot.sourceConfig,
|
||||
) as OpenClawConfig;
|
||||
return activateSecretsRuntimeSnapshotStateIfCurrent({
|
||||
return {
|
||||
...params,
|
||||
snapshot: {
|
||||
...params.snapshot,
|
||||
|
|
@ -1077,7 +1077,7 @@ export function restoreSecretsRuntimeSnapshotStateIfCurrent(
|
|||
mergeLiveAuthBookkeeping: false,
|
||||
preserveActivationLineage: false,
|
||||
expectedRevision: activeSnapshotRevision,
|
||||
});
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -50,7 +50,7 @@ import {
|
|||
getActiveSecretsRuntimeSnapshotRevisionState,
|
||||
graftActiveSecretsRuntimeAuthState,
|
||||
getPreparedSecretsRuntimeSnapshotRefreshContext,
|
||||
restoreSecretsRuntimeSnapshotStateIfCurrent,
|
||||
prepareSecretsRuntimeSnapshotRestoreState,
|
||||
setPreparedSecretsRuntimeSnapshotRefreshContext,
|
||||
type PreparedSecretsRuntimeSnapshot,
|
||||
type SecretsRuntimeRefreshContext,
|
||||
|
|
@ -341,14 +341,14 @@ export function activateSecretsRuntimeSnapshotIfCurrent(
|
|||
});
|
||||
}
|
||||
|
||||
/** Restores an owned predecessor while retaining changes after candidate preparation. */
|
||||
export function restoreSecretsRuntimeSnapshotIfCurrent(
|
||||
/** Prepares the exact rollback successor before the activation owner publishes it. */
|
||||
export function prepareSecretsRuntimeSnapshotRestore(
|
||||
snapshot: PreparedSecretsRuntimeSnapshot,
|
||||
expectedRevision: number,
|
||||
ownedSnapshot: PreparedSecretsRuntimeSnapshot,
|
||||
options?: { runtimeSourceConfig?: OpenClawConfig },
|
||||
): boolean {
|
||||
return restoreSecretsRuntimeSnapshotStateIfCurrent({
|
||||
) {
|
||||
return prepareSecretsRuntimeSnapshotRestoreState({
|
||||
...createSecretsRuntimeSnapshotActivation(snapshot),
|
||||
expectedRevision,
|
||||
ownedSnapshot,
|
||||
|
|
|
|||
|
|
@ -8,6 +8,8 @@ type LazyColumn = readonly [
|
|||
|
||||
// Added after v6 shipped; first-use-only columns stay absent until their feature writes.
|
||||
const lazyColumns = [
|
||||
["user_profile_identities", "authorization_id", "TEXT"],
|
||||
["user_profile_identities", "authorization_basis_json", "TEXT"],
|
||||
["claw_installs", "bootstrap_content_digest", "TEXT"],
|
||||
["claw_installs", "bootstrap_source_path", "TEXT"],
|
||||
["worker_environments", "desktop_json", "TEXT"],
|
||||
|
|
|
|||
|
|
@ -4,6 +4,7 @@ import type { DatabasePathIdentity } from "../infra/sqlite-worker-identity.js";
|
|||
|
||||
export type OpenClawStateSchemaReadAdmission = (database: DatabaseSync) => (() => void) | undefined;
|
||||
|
||||
// v19 preserves original channel-owner authorization across recovery.
|
||||
// v18 binds shared GitHub publication to its original requesting authority.
|
||||
// v17 records one-use prepared worker capacity and node workspace ownership.
|
||||
// v16 makes Skill Workshop ownership directory-based instead of row-provenance-based.
|
||||
|
|
@ -18,13 +19,14 @@ export type OpenClawStateSchemaReadAdmission = (database: DatabaseSync) => (() =
|
|||
// v7 retires the inert shared commitments table.
|
||||
// v6 makes every committed shared-state table part of the canonical runtime schema.
|
||||
// v5 records durable cloud-worker result refs on pending workspace fences.
|
||||
export const OPENCLAW_STATE_SCHEMA_VERSION = 18;
|
||||
export const OPENCLAW_STATE_SCHEMA_VERSION = 19;
|
||||
export const OPENCLAW_STATE_STRICT_SCHEMA_VERSION = 3;
|
||||
// Absence records lost history; only Doctor may reconstruct these on existing state.
|
||||
export const DOCTOR_OWNED_STATE_TABLES = ["agent_deletion_journal"] as const;
|
||||
// Privacy-sensitive feature tables remain absent even in fresh databases until
|
||||
// their feature-local first write. The canonical SQL still owns their shape.
|
||||
export const FIRST_USE_STATE_TABLES = [
|
||||
"user_profile_identities",
|
||||
"local_workspace_projections",
|
||||
"update_runs",
|
||||
"session_repository_workspaces",
|
||||
|
|
@ -54,6 +56,8 @@ export const FIRST_USE_STATE_TABLES = [
|
|||
"outbound_message_progress",
|
||||
] as const;
|
||||
export const FIRST_USE_STATE_INDEXES = [
|
||||
"idx_user_profile_identities_profile_id",
|
||||
"idx_user_profile_identities_authorization",
|
||||
"idx_update_runs_created",
|
||||
"idx_update_runs_active",
|
||||
"idx_github_repository_publication_shared_request",
|
||||
|
|
|
|||
11
src/state/openclaw-state-db.generated.d.ts
vendored
11
src/state/openclaw-state-db.generated.d.ts
vendored
|
|
@ -1504,6 +1504,16 @@ export interface UserPreferences {
|
|||
value_json: string;
|
||||
}
|
||||
|
||||
export interface UserProfileIdentities {
|
||||
authorization_basis_json: string | null;
|
||||
authorization_id: string | null;
|
||||
canonical_login: string | null;
|
||||
created_at: number;
|
||||
profile_id: string;
|
||||
provider: string;
|
||||
subject: string;
|
||||
}
|
||||
|
||||
export interface WebPushApprovalDeliveries {
|
||||
approval_id: string;
|
||||
device_id: string;
|
||||
|
|
@ -1884,6 +1894,7 @@ export interface DB {
|
|||
task_runs: TaskRuns;
|
||||
update_runs: UpdateRuns;
|
||||
user_preferences: UserPreferences;
|
||||
user_profile_identities: UserProfileIdentities;
|
||||
web_push_approval_deliveries: WebPushApprovalDeliveries;
|
||||
web_push_subscriptions: WebPushSubscriptions;
|
||||
worker_environment_credentials: WorkerEnvironmentCredentials;
|
||||
|
|
|
|||
|
|
@ -96,7 +96,7 @@ function captureCommand(command: OpenClawStateReadCommand): OpenClawStateReadCom
|
|||
return structuredClone(command);
|
||||
}
|
||||
if (command.type === "userProfiles.channelIdentity.resolve") {
|
||||
return { type: command.type, identity: { ...command.identity } };
|
||||
return { type: command.type, identity: structuredClone(command.identity) };
|
||||
}
|
||||
if (command.type === "userProfiles.githubAttribution.resolve") {
|
||||
return { type: command.type, profileIds: [...command.profileIds] };
|
||||
|
|
@ -400,13 +400,7 @@ function commandBytes(command: OpenClawStateReadRequest["command"]): number {
|
|||
);
|
||||
}
|
||||
if (command.type === "userProfiles.channelIdentity.resolve") {
|
||||
return (
|
||||
bytes +
|
||||
Object.values(command.identity).reduce(
|
||||
(total, value) => total + Buffer.byteLength(value, "utf8"),
|
||||
0,
|
||||
)
|
||||
);
|
||||
return bytes + Buffer.byteLength(JSON.stringify(command.identity), "utf8");
|
||||
}
|
||||
if (command.type === "userProfiles.email.resolve") {
|
||||
return bytes + Buffer.byteLength(command.email, "utf8");
|
||||
|
|
|
|||
|
|
@ -81,7 +81,7 @@ import type { OpenClawStateWorkerContext } from "./openclaw-state-worker-context
|
|||
import type { OpenClawStateWorkerErrorPayload } from "./openclaw-state-worker-error.js";
|
||||
import type { SessionRepositoryWorkspaceRecord } from "./session-repository-workspaces.types.js";
|
||||
import type {
|
||||
UserChannelIdentity,
|
||||
UserChannelIdentitySelector,
|
||||
UserChannelIdentityLink,
|
||||
UserChannelIdentityAuthorityFacts,
|
||||
UserChannelIdentityResult,
|
||||
|
|
@ -148,7 +148,7 @@ export type OpenClawStateReadCommand =
|
|||
| { type: "onboardingRecommendations.read"; configKey: string }
|
||||
| { type: "userProfiles.reconcile"; profileId: string }
|
||||
| { type: "userProfiles.channelIdentity.list"; profileId: string }
|
||||
| { type: "userProfiles.channelIdentity.resolve"; identity: UserChannelIdentity }
|
||||
| { type: "userProfiles.channelIdentity.resolve"; identity: UserChannelIdentitySelector }
|
||||
| { type: "userProfiles.authority.resolve"; profileId: string }
|
||||
| { type: "userProfiles.githubIdentity.cached"; accountId: number; email: string }
|
||||
| { type: "userProfiles.githubAttribution.resolve"; profileIds: readonly string[] }
|
||||
|
|
|
|||
|
|
@ -161,7 +161,10 @@ export function isReadRequest(input: unknown): input is OpenClawStateReadRequest
|
|||
Array.isArray(input.command.profileIds) &&
|
||||
input.command.profileIds.every((profileId) => typeof profileId === "string")) ||
|
||||
(input.command.type === "userProfiles.channelIdentity.resolve" &&
|
||||
Check(UserChannelIdentitySchema, input.command.identity)) ||
|
||||
(Check(UserChannelIdentitySchema, input.command.identity) ||
|
||||
(isRecord(input.command.identity) &&
|
||||
typeof input.command.identity.authorizationId === "string" &&
|
||||
isRecord(input.command.identity.policy)))) ||
|
||||
(input.command.type === "userProfiles.email.resolve" &&
|
||||
typeof input.command.email === "string") ||
|
||||
(input.command.type === "audit.run.inspect" &&
|
||||
|
|
|
|||
|
|
@ -2690,6 +2690,21 @@ CREATE TABLE IF NOT EXISTS claw_mcp_server_refs (
|
|||
PRIMARY KEY (agent_id, name)
|
||||
) STRICT;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS user_profile_identities (
|
||||
provider TEXT NOT NULL,
|
||||
subject TEXT NOT NULL,
|
||||
profile_id TEXT NOT NULL,
|
||||
canonical_login TEXT,
|
||||
created_at INTEGER NOT NULL,
|
||||
authorization_id TEXT,
|
||||
authorization_basis_json TEXT,
|
||||
PRIMARY KEY (provider, subject)
|
||||
) STRICT;
|
||||
CREATE INDEX IF NOT EXISTS idx_user_profile_identities_profile_id
|
||||
ON user_profile_identities(profile_id);
|
||||
CREATE UNIQUE INDEX IF NOT EXISTS idx_user_profile_identities_authorization
|
||||
ON user_profile_identities(authorization_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS outbound_media_provenance (
|
||||
realpath TEXT NOT NULL PRIMARY KEY,
|
||||
kind TEXT NOT NULL,
|
||||
|
|
|
|||
|
|
@ -1,12 +1,20 @@
|
|||
import type { DatabaseSync } from "node:sqlite";
|
||||
import { isDeepStrictEqual } from "node:util";
|
||||
import { Check } from "typebox/value";
|
||||
import { z } from "zod";
|
||||
import {
|
||||
GATEWAY_OWNER_PROFILE_ID,
|
||||
UserChannelIdentitySchema,
|
||||
} from "../../packages/gateway-protocol/src/schema/users.js";
|
||||
import type { GatewayConfig } from "../config/types.gateway.js";
|
||||
import { executeSqliteQuerySync, executeSqliteQueryTakeFirstSync } from "../infra/kysely-sync.js";
|
||||
import { generateSecureUuid } from "../infra/secure-random.js";
|
||||
import { deferSqlitePostCommitPublication } from "../infra/sqlite-post-commit.js";
|
||||
import { getAdmittedSqliteSchemaFacts } from "../infra/sqlite-schema-facts.js";
|
||||
import { runSqliteDeferredTransactionSync } from "../infra/sqlite-transaction.js";
|
||||
import type { GatewayAccessGrantRef } from "../plugins/gateway-access-policy.types.js";
|
||||
import { updateConfigMachineStateInDatabase } from "./config-machine-state-write.js";
|
||||
import { readConfigMachineStateRowInDatabase } from "./config-machine-state.js";
|
||||
import { withExistingOpenClawStateDatabaseReadOnly } from "./openclaw-state-db-readonly.js";
|
||||
import { tableExists } from "./openclaw-state-db-schema-helpers.js";
|
||||
import {
|
||||
|
|
@ -16,6 +24,7 @@ import {
|
|||
import {
|
||||
publishUserProfileAliasChange,
|
||||
publishUserChannelIdentityAuthorityChange,
|
||||
publishUserProfileAuthorityChange,
|
||||
} from "./user-profile-events.js";
|
||||
import { selectStoredGitHubIdentities } from "./user-profile-github-identity.js";
|
||||
import { publishUserProfilesChange } from "./user-profile-list.js";
|
||||
|
|
@ -30,10 +39,142 @@ import type {
|
|||
UserChannelIdentity,
|
||||
UserChannelIdentityLink,
|
||||
UserChannelIdentityAuthorityFacts,
|
||||
UserChannelIdentitySelector,
|
||||
} from "./user-profiles.types.js";
|
||||
|
||||
// The dot keeps administrator-attested channel links outside Tailscale login namespaces.
|
||||
const CHANNEL_IDENTITY_PROVIDER = "channel.identity";
|
||||
const POLICY_KEY = "operator.channelPolicy";
|
||||
const referenceSchema = z.strictObject({ version: z.literal(1), id: z.uuid() });
|
||||
const grantSchema = z
|
||||
.strictObject({ pluginId: z.string().min(1).max(128), grantId: z.uuid() })
|
||||
.nullable();
|
||||
export type UserChannelAuthorizationReference = Readonly<z.infer<typeof referenceSchema>>;
|
||||
export type UserChannelAuthorization = {
|
||||
reference: UserChannelAuthorizationReference;
|
||||
subject: string;
|
||||
grant: GatewayAccessGrantRef | null;
|
||||
};
|
||||
|
||||
export function parseUserChannelAuthorizationReference(value: unknown) {
|
||||
return referenceSchema.safeParse(value).data;
|
||||
}
|
||||
|
||||
export function resolveUserChannelAuthorizationPolicy(
|
||||
gateway: Pick<GatewayConfig, "roles" | "auth"> | undefined,
|
||||
) {
|
||||
return { roles: gateway?.roles ?? null, identityScopes: gateway?.auth?.identityScopes ?? null };
|
||||
}
|
||||
export type UserChannelAuthorizationPolicy = ReturnType<
|
||||
typeof resolveUserChannelAuthorizationPolicy
|
||||
>;
|
||||
|
||||
function matchesPolicy(db: DatabaseSync, policy: UserChannelAuthorizationPolicy): boolean {
|
||||
const row = readConfigMachineStateRowInDatabase(db, POLICY_KEY);
|
||||
return row !== undefined && isDeepStrictEqual(JSON.parse(row.value_json), policy);
|
||||
}
|
||||
|
||||
/** Activation and rollback retire old references before their exact authorization policy is published. */
|
||||
export function publishUserChannelPolicyInDatabase(
|
||||
db: DatabaseSync,
|
||||
policy: UserChannelAuthorizationPolicy,
|
||||
): string[] {
|
||||
let profiles: string[] = [];
|
||||
updateConfigMachineStateInDatabase(
|
||||
db,
|
||||
POLICY_KEY,
|
||||
(current: unknown) => {
|
||||
if (
|
||||
!isDeepStrictEqual(current, policy) &&
|
||||
getAdmittedSqliteSchemaFacts(db)?.tables.has("user_profile_identities")
|
||||
) {
|
||||
profiles = executeSqliteQuerySync(
|
||||
db,
|
||||
userProfilesDb(db)
|
||||
.selectFrom("user_profile_identities")
|
||||
.select("profile_id")
|
||||
.distinct()
|
||||
.where("provider", "=", CHANNEL_IDENTITY_PROVIDER),
|
||||
).rows.map((row) => row.profile_id);
|
||||
publishUserProfileAuthorityChange(db, ...profiles);
|
||||
}
|
||||
return policy;
|
||||
},
|
||||
Date.now(),
|
||||
);
|
||||
return profiles;
|
||||
}
|
||||
|
||||
/** The reference names this uninterrupted channel link; it carries no permission by itself. */
|
||||
export function authorizeUserChannelIdentityInDatabase(
|
||||
db: DatabaseSync,
|
||||
input: {
|
||||
identity: UserChannelIdentity;
|
||||
profileId: string;
|
||||
policy: UserChannelAuthorizationPolicy;
|
||||
grant: GatewayAccessGrantRef | null;
|
||||
},
|
||||
): UserChannelAuthorizationReference | undefined {
|
||||
if (
|
||||
(getAdmittedSqliteSchemaFacts(db)?.userVersion ?? 0) < 19 ||
|
||||
!matchesPolicy(db, input.policy)
|
||||
) {
|
||||
return undefined;
|
||||
}
|
||||
const subject = userChannelIdentitySubject(input.identity);
|
||||
const row = executeSqliteQueryTakeFirstSync(
|
||||
db,
|
||||
userProfilesDb(db)
|
||||
.selectFrom("user_profile_identities")
|
||||
.select(["authorization_id", "authorization_basis_json"])
|
||||
.where("provider", "=", CHANNEL_IDENTITY_PROVIDER)
|
||||
.where("subject", "=", subject)
|
||||
.where("profile_id", "=", input.profileId),
|
||||
);
|
||||
if (!row) {
|
||||
return undefined;
|
||||
}
|
||||
const basis = JSON.stringify(grantSchema.parse(input.grant));
|
||||
const original = referenceSchema.safeParse({ version: 1, id: row.authorization_id }).data;
|
||||
if (row.authorization_basis_json === basis && original) {
|
||||
return original;
|
||||
}
|
||||
const id = generateSecureUuid();
|
||||
executeSqliteQuerySync(
|
||||
db,
|
||||
userProfilesDb(db)
|
||||
.updateTable("user_profile_identities")
|
||||
.set({ authorization_id: id, authorization_basis_json: basis })
|
||||
.where("provider", "=", CHANNEL_IDENTITY_PROVIDER)
|
||||
.where("subject", "=", subject),
|
||||
);
|
||||
return { version: 1, id };
|
||||
}
|
||||
|
||||
function readAuthorization(db: DatabaseSync, id: string, policy: UserChannelAuthorizationPolicy) {
|
||||
if ((getAdmittedSqliteSchemaFacts(db)?.userVersion ?? 0) < 19 || !matchesPolicy(db, policy)) {
|
||||
return undefined;
|
||||
}
|
||||
const row = executeSqliteQueryTakeFirstSync(
|
||||
db,
|
||||
userProfilesDb(db)
|
||||
.selectFrom("user_profile_identities")
|
||||
.select(["subject", "authorization_basis_json"])
|
||||
.where("provider", "=", CHANNEL_IDENTITY_PROVIDER)
|
||||
.where("authorization_id", "=", id),
|
||||
);
|
||||
if (!row?.authorization_basis_json || row.authorization_basis_json.length > 1024) {
|
||||
return undefined;
|
||||
}
|
||||
try {
|
||||
const grant = grantSchema.safeParse(JSON.parse(row.authorization_basis_json));
|
||||
return grant.success
|
||||
? { subject: row.subject, reference: { version: 1 as const, id }, grant: grant.data }
|
||||
: undefined;
|
||||
} catch {
|
||||
return undefined;
|
||||
}
|
||||
}
|
||||
|
||||
export class UserChannelIdentityConflictError extends Error {
|
||||
constructor() {
|
||||
|
|
@ -190,13 +331,23 @@ export function listUserChannelIdentitiesInDatabase(
|
|||
/** Reads the current person and login grant subjects; channel links never become login aliases. */
|
||||
export function resolveUserChannelIdentityInDatabase(
|
||||
db: DatabaseSync,
|
||||
identity: UserChannelIdentity,
|
||||
identity: UserChannelIdentitySelector,
|
||||
): UserChannelIdentityAuthorityFacts | undefined {
|
||||
const subject = userChannelIdentitySubject(identity);
|
||||
return runSqliteDeferredTransactionSync(db, () => {
|
||||
if (!hasIdentityTables(db)) {
|
||||
return undefined;
|
||||
}
|
||||
let authorization: UserChannelAuthorization | undefined;
|
||||
let subject: string;
|
||||
if ("authorizationId" in identity) {
|
||||
authorization = readAuthorization(db, identity.authorizationId, identity.policy);
|
||||
if (!authorization) {
|
||||
return undefined;
|
||||
}
|
||||
subject = authorization.subject;
|
||||
} else {
|
||||
subject = userChannelIdentitySubject(identity);
|
||||
}
|
||||
const link = selectLink(db, subject);
|
||||
const profile = link ? selectResolvedUserProfileMetadataById(db, link.profile_id) : undefined;
|
||||
if (!profile || profile.id === GATEWAY_OWNER_PROFILE_ID) {
|
||||
|
|
@ -237,6 +388,7 @@ export function resolveUserChannelIdentityInDatabase(
|
|||
.get(profile.id)
|
||||
?.accounts.map((account) => `${account.login.toLowerCase()}@github`) ?? [];
|
||||
return {
|
||||
...(authorization ? { authorization } : {}),
|
||||
profileId: profile.id,
|
||||
role: profile.role ?? null,
|
||||
emails,
|
||||
|
|
|
|||
|
|
@ -1,4 +1,3 @@
|
|||
import type { DatabaseSync } from "node:sqlite";
|
||||
import {
|
||||
deferSqliteWorkerCommitReceipt,
|
||||
requestSqliteWorkerOperationAdmission,
|
||||
|
|
@ -12,6 +11,8 @@ import {
|
|||
unlinkUserChannelIdentity,
|
||||
UserChannelIdentityConflictError,
|
||||
userChannelIdentitySubject,
|
||||
authorizeUserChannelIdentityInDatabase,
|
||||
publishUserChannelPolicyInDatabase,
|
||||
} from "./user-channel-identities.js";
|
||||
import {
|
||||
ensureUserProfilesSchema,
|
||||
|
|
@ -44,42 +45,64 @@ export function executeUserChannelIdentityChange(
|
|||
input: UserChannelIdentityWorkerOperations["userProfiles.channelIdentity.change"]["input"],
|
||||
options: OpenClawStateDatabaseOptions,
|
||||
): UserChannelIdentityWorkerOperations["userProfiles.channelIdentity.change"]["output"] {
|
||||
const subject = userChannelIdentitySubject(input.identity);
|
||||
const subject = "identity" in input ? userChannelIdentitySubject(input.identity) : undefined;
|
||||
const profiles: string[] = [];
|
||||
const facts = {
|
||||
kind: "channel-identity",
|
||||
action: input.action,
|
||||
subject,
|
||||
profiles,
|
||||
};
|
||||
let changed = false;
|
||||
const mutationOptions = {
|
||||
...options,
|
||||
beforeChange(db: DatabaseSync) {
|
||||
beforeChange() {
|
||||
changed = true;
|
||||
requestSqliteWorkerOperationAdmission({
|
||||
stage: "transaction",
|
||||
facts: { kind: "channel-identity", subject },
|
||||
facts,
|
||||
});
|
||||
deferSqliteWorkerCommitReceipt(db, { kind: "channel-identity", subject });
|
||||
},
|
||||
};
|
||||
ensureUserProfilesSchema(options);
|
||||
if (input.action !== "policy") {
|
||||
ensureUserProfilesSchema(options);
|
||||
}
|
||||
return readUserChannelIdentityResult(() =>
|
||||
runOpenClawStateWriteTransaction(
|
||||
() => {
|
||||
({ db }) => {
|
||||
if (input.action === "policy" || input.action === "authorize") {
|
||||
mutationOptions.beforeChange();
|
||||
}
|
||||
if (input.action === "policy") {
|
||||
facts.profiles = publishUserChannelPolicyInDatabase(db, input.policy);
|
||||
}
|
||||
const value =
|
||||
input.action === "link"
|
||||
? {
|
||||
kind: "linked" as const,
|
||||
link: linkUserChannelIdentity(input.profileId, input.identity, mutationOptions),
|
||||
}
|
||||
: {
|
||||
kind: "unlinked" as const,
|
||||
removed: unlinkUserChannelIdentity(
|
||||
input.profileId,
|
||||
input.identity,
|
||||
mutationOptions,
|
||||
),
|
||||
};
|
||||
input.action === "policy"
|
||||
? { kind: "policy" as const }
|
||||
: input.action === "authorize"
|
||||
? {
|
||||
kind: "authorized" as const,
|
||||
reference: authorizeUserChannelIdentityInDatabase(db, input),
|
||||
}
|
||||
: input.action === "link"
|
||||
? {
|
||||
kind: "linked" as const,
|
||||
link: linkUserChannelIdentity(input.profileId, input.identity, mutationOptions),
|
||||
}
|
||||
: {
|
||||
kind: "unlinked" as const,
|
||||
removed: unlinkUserChannelIdentity(
|
||||
input.profileId,
|
||||
input.identity,
|
||||
mutationOptions,
|
||||
),
|
||||
};
|
||||
if (changed) {
|
||||
requestSqliteWorkerOperationAdmission({
|
||||
stage: "commit",
|
||||
facts: { kind: "channel-identity", subject },
|
||||
facts,
|
||||
});
|
||||
deferSqliteWorkerCommitReceipt(db, facts);
|
||||
}
|
||||
return value;
|
||||
},
|
||||
|
|
|
|||
|
|
@ -10,6 +10,7 @@ import { runOpenClawStateWorkerOperation } from "./openclaw-state-worker-store.j
|
|||
import {
|
||||
UserChannelIdentityConflictError,
|
||||
userChannelIdentitySubject,
|
||||
resolveUserChannelAuthorizationPolicy,
|
||||
} from "./user-channel-identities.js";
|
||||
import {
|
||||
captureUserProfileAuthorityRead,
|
||||
|
|
@ -22,6 +23,8 @@ import type {
|
|||
UserChannelIdentity,
|
||||
UserChannelIdentityLink,
|
||||
UserChannelIdentityResult,
|
||||
UserChannelIdentitySelector,
|
||||
UserChannelIdentityWorkerOperations,
|
||||
} from "./user-profiles.types.js";
|
||||
|
||||
type IdentityOptions = Pick<OpenClawStateDatabaseOptions, "path" | "env">;
|
||||
|
|
@ -66,14 +69,16 @@ export async function listCanonicalUserChannelIdentities(
|
|||
return unwrapIdentityResult(reply.result, profileId);
|
||||
}
|
||||
|
||||
export async function changeCanonicalUserChannelIdentity(
|
||||
action: "link" | "unlink",
|
||||
profileId: string,
|
||||
identity: UserChannelIdentity,
|
||||
type IdentityChange =
|
||||
UserChannelIdentityWorkerOperations["userProfiles.channelIdentity.change"]["input"];
|
||||
|
||||
async function changeIdentity(
|
||||
input: IdentityChange,
|
||||
options: IdentityOptions & { assertCurrent?: () => void } = {},
|
||||
) {
|
||||
const capturedIdentity = { ...identity };
|
||||
const subject = userChannelIdentitySubject(capturedIdentity);
|
||||
const captured = structuredClone(input);
|
||||
const subject =
|
||||
"identity" in captured ? userChannelIdentitySubject(captured.identity) : undefined;
|
||||
const assertCurrent = options.assertCurrent;
|
||||
const context = captureOpenClawStateWorkerContext(options);
|
||||
const result = await runOpenClawStateWorkerOperation(
|
||||
|
|
@ -81,7 +86,7 @@ export async function changeCanonicalUserChannelIdentity(
|
|||
(scope) =>
|
||||
scope.execute({
|
||||
type: "userProfiles.channelIdentity.change",
|
||||
input: { action, profileId, identity: capturedIdentity },
|
||||
input: captured,
|
||||
}),
|
||||
{
|
||||
assertCurrent,
|
||||
|
|
@ -92,17 +97,22 @@ export async function changeCanonicalUserChannelIdentity(
|
|||
(request.stage !== "transaction" && request.stage !== "commit") ||
|
||||
!isRecord(request.facts) ||
|
||||
request.facts.kind !== "channel-identity" ||
|
||||
request.facts.subject !== subject
|
||||
request.facts.action !== captured.action ||
|
||||
request.facts.subject !== subject ||
|
||||
!Array.isArray(request.facts.profiles) ||
|
||||
!request.facts.profiles.every(
|
||||
(profileId): profileId is string => typeof profileId === "string",
|
||||
)
|
||||
) {
|
||||
throw new Error("Channel identity mutation requires exact transaction admission");
|
||||
}
|
||||
context.admission.assertCurrent();
|
||||
assertCurrent?.();
|
||||
if (request.stage === "commit") {
|
||||
if (request.stage === "commit" && captured.action !== "authorize") {
|
||||
fence ??= fenceUserProfileMutationAuthority(context.admission, {
|
||||
profiles: [],
|
||||
profiles: request.facts.profiles,
|
||||
identities: [],
|
||||
channels: [subject],
|
||||
channels: subject === undefined ? [] : [subject],
|
||||
});
|
||||
}
|
||||
grant();
|
||||
|
|
@ -110,6 +120,7 @@ export async function changeCanonicalUserChannelIdentity(
|
|||
void operation.settled.then((settlement) => {
|
||||
const committed = admission.committed;
|
||||
if (
|
||||
captured.action !== "authorize" &&
|
||||
committed &&
|
||||
isRecord(committed.facts) &&
|
||||
committed.facts.kind === "channel-identity" &&
|
||||
|
|
@ -124,16 +135,52 @@ export async function changeCanonicalUserChannelIdentity(
|
|||
},
|
||||
},
|
||||
);
|
||||
return unwrapIdentityResult(result, profileId);
|
||||
return unwrapIdentityResult(result, "profileId" in captured ? captured.profileId : "");
|
||||
}
|
||||
|
||||
export async function changeCanonicalUserChannelIdentity(
|
||||
action: "link" | "unlink",
|
||||
profileId: string,
|
||||
identity: UserChannelIdentity,
|
||||
options: IdentityOptions & { assertCurrent?: () => void } = {},
|
||||
) {
|
||||
const result = await changeIdentity({ action, profileId, identity }, options);
|
||||
if (result.kind !== "linked" && result.kind !== "unlinked") {
|
||||
throw new Error("Unexpected channel identity change");
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
export async function publishCanonicalUserChannelPolicy(
|
||||
gateway: Parameters<typeof resolveUserChannelAuthorizationPolicy>[0],
|
||||
) {
|
||||
await changeIdentity({
|
||||
action: "policy",
|
||||
policy: resolveUserChannelAuthorizationPolicy(gateway),
|
||||
});
|
||||
}
|
||||
|
||||
export async function authorizeCanonicalUserChannelIdentity(
|
||||
input: Extract<IdentityChange, { action: "authorize" }>,
|
||||
options: IdentityOptions & { assertCurrent: () => void },
|
||||
) {
|
||||
const result = await changeIdentity(input, options);
|
||||
if (result.kind !== "authorized") {
|
||||
throw new Error("Unexpected channel authorization result");
|
||||
}
|
||||
return result.reference;
|
||||
}
|
||||
|
||||
/** Qualify worker-read facts against the same physical profile owner's mutation lifetime. */
|
||||
export async function prepareUserChannelIdentityAuthority(
|
||||
identity: UserChannelIdentity,
|
||||
identity: UserChannelIdentitySelector,
|
||||
options: IdentityOptions = {},
|
||||
) {
|
||||
const capturedIdentity = { ...identity };
|
||||
const subject = userChannelIdentitySubject(capturedIdentity);
|
||||
const capturedIdentity = structuredClone(identity);
|
||||
const subject =
|
||||
"authorizationId" in capturedIdentity
|
||||
? undefined
|
||||
: userChannelIdentitySubject(capturedIdentity);
|
||||
const context = captureAuthorityContext(options);
|
||||
for (let attempt = 0; attempt < 3; attempt += 1) {
|
||||
const read = await captureUserProfileAuthorityRead(context.admission, subject);
|
||||
|
|
@ -154,7 +201,10 @@ export async function prepareUserChannelIdentityAuthority(
|
|||
if (!reply.linked) {
|
||||
return undefined;
|
||||
}
|
||||
const isCurrent = read.bind(reply.linked.profileId);
|
||||
const isCurrent = read.bind(
|
||||
reply.linked.profileId,
|
||||
reply.linked.authorization?.subject ?? subject,
|
||||
);
|
||||
if (isCurrent) {
|
||||
return { linked: reply.linked, isCurrent };
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,7 +1,6 @@
|
|||
import type { DatabaseSync } from "node:sqlite";
|
||||
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
|
||||
import { executeSqliteQueryTakeFirstSync } from "../infra/kysely-sync.js";
|
||||
import { publishUserProfileAuthorityChange } from "./user-profile-events.js";
|
||||
import { publishUserProfilesChange } from "./user-profile-list.js";
|
||||
import type { UserProfileMutationContext } from "./user-profile-mutation.js";
|
||||
import {
|
||||
|
|
@ -43,7 +42,6 @@ export function ensureProfileForEmailInDatabase(
|
|||
const row = insertUserProfile(db, displayName, now, mutation);
|
||||
setUserProfileEmailBinding(db, email, row.id, now);
|
||||
mutation?.authority(row.id);
|
||||
publishUserProfileAuthorityChange(db, row.id);
|
||||
mutation?.publish(row.id);
|
||||
publishUserProfilesChange(db, row.id);
|
||||
return toUserProfile(row);
|
||||
|
|
|
|||
|
|
@ -1,5 +1,7 @@
|
|||
import type { DatabaseSync } from "node:sqlite";
|
||||
import { executeSqliteQuerySync, getNodeSqliteKysely } from "../infra/kysely-sync.js";
|
||||
import { stageSqliteTransactionState } from "../infra/sqlite-post-commit.js";
|
||||
import { getAdmittedSqliteSchemaFacts } from "../infra/sqlite-schema-facts.js";
|
||||
import type { DatabasePathIdentity } from "../infra/sqlite-worker-identity.js";
|
||||
import { sessionChanges } from "../sessions/session-row-changes.js";
|
||||
import { createDeferredCore } from "../shared/deferred.js";
|
||||
|
|
@ -7,7 +9,7 @@ import { notifyListeners, registerListener } from "../shared/listeners.js";
|
|||
import type { OpenClawStateDatabaseReadAdmission } from "./openclaw-state-db-async-lifecycle.js";
|
||||
import { registerOpenClawStateDatabaseLifecycleListener } from "./openclaw-state-db-cache.js";
|
||||
import type { UserProfileMutationChanges } from "./user-profile-mutation.js";
|
||||
import type { UserProfileEmailBinding } from "./user-profiles.types.js";
|
||||
import type { UserProfileEmailBinding, UserProfilesDatabase } from "./user-profiles.types.js";
|
||||
|
||||
type EmailBindingChange = {
|
||||
db: DatabaseSync;
|
||||
|
|
@ -81,6 +83,20 @@ function observeAuthorityLifecycle(): void {
|
|||
|
||||
/** Authority revisions belong to the profile writer, independently of display notifications. */
|
||||
export function publishUserProfileAuthorityChange(db: DatabaseSync, ...profileIds: string[]): void {
|
||||
// Native and worker mutations retire recovery custody in their original transaction.
|
||||
const schema = profileIds.length ? getAdmittedSqliteSchemaFacts(db) : undefined;
|
||||
if (profileIds.length && !schema) {
|
||||
throw new Error("Profile authority mutation requires admitted schema facts");
|
||||
}
|
||||
if (schema && schema.userVersion >= 19 && schema.tables.has("user_profile_identities")) {
|
||||
executeSqliteQuerySync(
|
||||
db,
|
||||
getNodeSqliteKysely<UserProfilesDatabase>(db)
|
||||
.updateTable("user_profile_identities")
|
||||
.set({ authorization_id: null, authorization_basis_json: null })
|
||||
.where("profile_id", "in", profileIds),
|
||||
);
|
||||
}
|
||||
observeAuthorityLifecycle();
|
||||
const store = changes.authorityHandles.get(db);
|
||||
if (!store || profileIds.length === 0) {
|
||||
|
|
@ -205,9 +221,8 @@ export async function captureUserProfileAuthorityRead(
|
|||
await Promise.all(pending);
|
||||
}
|
||||
admission.assertCurrent();
|
||||
const subjectIsSettled = () =>
|
||||
subjectKey === undefined ||
|
||||
(!store.uncertain.has(subjectKey) && !store.pending.get(subjectKey)?.size);
|
||||
const subjectIsSettled = (key = subjectKey) =>
|
||||
key === undefined || (!store.uncertain.has(key) && !store.pending.get(key)?.size);
|
||||
if (!subjectIsSettled()) {
|
||||
throw new UserProfileMutationUnsettledError("pending");
|
||||
}
|
||||
|
|
@ -235,7 +250,10 @@ export async function captureUserProfileAuthorityRead(
|
|||
throw new UserProfileMutationUnsettledError("pending");
|
||||
}
|
||||
},
|
||||
bind(profileIds: string | readonly string[]): (() => boolean) | undefined {
|
||||
bind(
|
||||
profileIds: string | readonly string[],
|
||||
boundSubject = subject,
|
||||
): (() => boolean) | undefined {
|
||||
admission.assertCurrent();
|
||||
if (
|
||||
changes.authorityStores.get(admission.identity.key) !== store ||
|
||||
|
|
@ -255,7 +273,13 @@ export async function captureUserProfileAuthorityRead(
|
|||
if (profiles.some(({ key }) => store.pending.get(key)?.size)) {
|
||||
return undefined;
|
||||
}
|
||||
const identity = subject === undefined ? undefined : store.channelIdentities.get(subject);
|
||||
const identity =
|
||||
boundSubject === undefined ? undefined : store.channelIdentities.get(boundSubject);
|
||||
const boundKey =
|
||||
boundSubject === undefined ? undefined : mutationKey("channel", boundSubject);
|
||||
if (!subjectIsSettled(boundKey)) {
|
||||
return undefined;
|
||||
}
|
||||
return () => {
|
||||
try {
|
||||
admission.assertCurrent();
|
||||
|
|
@ -267,8 +291,9 @@ export async function captureUserProfileAuthorityRead(
|
|||
!store.uncertain.has(profile.key) &&
|
||||
!store.pending.get(profile.key)?.size,
|
||||
) &&
|
||||
(subject === undefined || store.channelIdentities.get(subject) === identity) &&
|
||||
subjectIsSettled()
|
||||
(boundSubject === undefined ||
|
||||
store.channelIdentities.get(boundSubject) === identity) &&
|
||||
subjectIsSettled(boundKey)
|
||||
);
|
||||
} catch {
|
||||
return false;
|
||||
|
|
|
|||
|
|
@ -2,12 +2,14 @@ import type { DatabaseSync } from "node:sqlite";
|
|||
import { executeSqliteQuerySync, getNodeSqliteKysely } from "../infra/kysely-sync.js";
|
||||
import { generateSecureUuid } from "../infra/secure-random.js";
|
||||
import { stageSqliteTransactionState } from "../infra/sqlite-post-commit.js";
|
||||
import { extractSqliteTableSchema } from "../infra/sqlite-schema-sql.js";
|
||||
import { ensureColumn, tableHasColumn } from "./openclaw-state-db-schema-helpers.js";
|
||||
import {
|
||||
openOpenClawStateDatabase,
|
||||
runOpenClawStateWriteTransaction,
|
||||
type OpenClawStateDatabaseOptions,
|
||||
} from "./openclaw-state-db.js";
|
||||
import { OPENCLAW_STATE_SCHEMA_SQL } from "./openclaw-state-schema.js";
|
||||
import { stageUserProfileEmailBindingChange } from "./user-profile-events.js";
|
||||
import {
|
||||
runUserProfileWriteTransaction,
|
||||
|
|
@ -41,17 +43,7 @@ CREATE TABLE IF NOT EXISTS user_profile_emails (
|
|||
CREATE INDEX IF NOT EXISTS idx_user_profile_emails_profile_id
|
||||
ON user_profile_emails(profile_id);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS user_profile_identities (
|
||||
provider TEXT NOT NULL,
|
||||
subject TEXT NOT NULL,
|
||||
profile_id TEXT NOT NULL,
|
||||
canonical_login TEXT,
|
||||
created_at INTEGER NOT NULL,
|
||||
PRIMARY KEY (provider, subject)
|
||||
) STRICT;
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_user_profile_identities_profile_id
|
||||
ON user_profile_identities(profile_id);
|
||||
${extractSqliteTableSchema(OPENCLAW_STATE_SCHEMA_SQL, "user_profile_identities")}
|
||||
`;
|
||||
|
||||
export class UserProfileNotFoundError extends Error {
|
||||
|
|
@ -110,6 +102,13 @@ export function ensureUserProfilesSchema(
|
|||
({ db }) => {
|
||||
db.exec(USER_PROFILES_SCHEMA_SQL); // sqlite-allow-raw -- Canonical feature-local additive DDL.
|
||||
ensureColumn(db, "user_profile_identities", "canonical_login TEXT");
|
||||
ensureColumn(db, "user_profile_identities", "authorization_id TEXT");
|
||||
ensureColumn(db, "user_profile_identities", "authorization_basis_json TEXT");
|
||||
db.exec(
|
||||
extractSqliteTableSchema(OPENCLAW_STATE_SCHEMA_SQL, "user_profile_identities", {
|
||||
endMarker: "ON user_profile_identities(authorization_id);",
|
||||
}),
|
||||
); // sqlite-allow-raw -- Canonical first-use channel-link indexes after additive columns.
|
||||
ensureColumn(db, "user_profiles", "primary_github_account_id INTEGER");
|
||||
ensureColumn(db, "user_profile_emails", "binding_id TEXT");
|
||||
const kysely = getNodeSqliteKysely<UserProfilesDatabase>(db);
|
||||
|
|
|
|||
|
|
@ -1,12 +1,21 @@
|
|||
import { join } from "node:path";
|
||||
import { DatabaseSync } from "node:sqlite";
|
||||
import { afterEach, describe, expect, it } from "vitest";
|
||||
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
|
||||
import { createUpdateRun } from "../infra/update-run-ledger.js";
|
||||
import { OPENCLAW_STATE_SCHEMA_VERSION } from "./openclaw-state-db-contract.js";
|
||||
import { tableHasColumn } from "./openclaw-state-db-schema-helpers.js";
|
||||
import {
|
||||
closeOpenClawStateDatabaseForTest,
|
||||
openOpenClawStateDatabase,
|
||||
runOpenClawStateWriteTransaction,
|
||||
} from "./openclaw-state-db.js";
|
||||
import {
|
||||
linkUserChannelIdentity,
|
||||
authorizeUserChannelIdentityInDatabase,
|
||||
publishUserChannelPolicyInDatabase,
|
||||
resolveUserChannelAuthorizationPolicy,
|
||||
} from "./user-channel-identities.js";
|
||||
import {
|
||||
listUserProfilesSync,
|
||||
readUserProfileEmailBindings,
|
||||
|
|
@ -232,3 +241,71 @@ describe("user profile role schema", () => {
|
|||
},
|
||||
);
|
||||
});
|
||||
|
||||
it.each([false, true])(
|
||||
"upgrades channel links without granting legacy recovery custody (deferred: %s)",
|
||||
(deferred) => {
|
||||
const options = stateOptions();
|
||||
const profile = ensureProfileForEmail("upgrade@example.test", options);
|
||||
const identity = { channelId: "discord", accountId: "team", senderId: "100" };
|
||||
linkUserChannelIdentity(profile.id, identity, options);
|
||||
const runId = "ed099411-cfbd-4304-a6b7-d3e504a48505";
|
||||
if (deferred) {
|
||||
createUpdateRun({ runId, trigger: "cli", before: { version: "2026.9.2" } }, options);
|
||||
}
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
const legacy = new DatabaseSync(options.path);
|
||||
legacy.exec(`
|
||||
DROP INDEX idx_user_profile_identities_authorization;
|
||||
ALTER TABLE user_profile_identities DROP COLUMN authorization_id;
|
||||
ALTER TABLE user_profile_identities DROP COLUMN authorization_basis_json;
|
||||
PRAGMA user_version = 18;
|
||||
UPDATE schema_meta SET schema_version = 18;
|
||||
`);
|
||||
const before = legacy.prepare("SELECT * FROM user_profile_identities").all();
|
||||
legacy.close();
|
||||
let db = openOpenClawStateDatabase(options).db;
|
||||
expect(db.prepare("SELECT * FROM user_profile_identities").all()).toEqual(
|
||||
before.map((row) =>
|
||||
Object.assign(row, {
|
||||
authorization_id: null,
|
||||
authorization_basis_json: null,
|
||||
}),
|
||||
),
|
||||
);
|
||||
const policy = resolveUserChannelAuthorizationPolicy({
|
||||
roles: {
|
||||
default: "admin",
|
||||
definitions: {
|
||||
admin: { scopes: ["operator.admin"], agents: "*", sessions: { others: "write" } },
|
||||
},
|
||||
},
|
||||
});
|
||||
const mint = () =>
|
||||
runOpenClawStateWriteTransaction(({ db: writer }) => {
|
||||
publishUserChannelPolicyInDatabase(writer, policy);
|
||||
return authorizeUserChannelIdentityInDatabase(writer, {
|
||||
identity,
|
||||
profileId: profile.id,
|
||||
policy,
|
||||
grant: null,
|
||||
});
|
||||
}, options);
|
||||
if (deferred) {
|
||||
expect(db.prepare("PRAGMA user_version").get()).toEqual({ user_version: 18 });
|
||||
expect(mint()).toBeUndefined();
|
||||
db.prepare(
|
||||
"UPDATE update_runs SET status = 'succeeded', phase = 'finished', finished_at_ms = ? WHERE run_id = ?",
|
||||
).run(Date.now() - 300_001, runId);
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
db = openOpenClawStateDatabase(options).db;
|
||||
}
|
||||
expect(db.prepare("PRAGMA user_version").get()).toEqual({
|
||||
user_version: OPENCLAW_STATE_SCHEMA_VERSION,
|
||||
});
|
||||
const reference = mint();
|
||||
expect(reference).toEqual({ version: 1, id: expect.any(String) });
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
expect(mint()).toEqual(reference);
|
||||
},
|
||||
);
|
||||
|
|
|
|||
|
|
@ -316,7 +316,6 @@ function ensureProfileForProviderIdentity(params: {
|
|||
}),
|
||||
);
|
||||
options.mutation?.authority(row.id);
|
||||
publishUserProfileAuthorityChange(db, row.id);
|
||||
options.mutation?.publish(row.id);
|
||||
publishUserProfilesChange(db, row.id);
|
||||
return toUserProfile(row);
|
||||
|
|
|
|||
|
|
@ -1,5 +1,12 @@
|
|||
import type { SqlBool } from "kysely";
|
||||
import type { GatewayAccessGrantRef } from "../plugins/gateway-access-policy.types.js";
|
||||
import type { USER_PROFILE_AVATAR_MIME_TYPES } from "../shared/avatar-limits.js";
|
||||
import type { DB } from "./openclaw-state-db.generated.js";
|
||||
import type {
|
||||
UserChannelAuthorization,
|
||||
UserChannelAuthorizationPolicy,
|
||||
UserChannelAuthorizationReference,
|
||||
} from "./user-channel-identities.js";
|
||||
|
||||
export const MAX_USER_PROFILE_DISPLAY_NAME_LENGTH = 256;
|
||||
|
||||
|
|
@ -26,8 +33,12 @@ export type UserProfileGitHubAttributionRead = {
|
|||
};
|
||||
|
||||
export type UserChannelIdentity = { channelId: string; accountId: string; senderId: string };
|
||||
export type UserChannelIdentitySelector =
|
||||
| UserChannelIdentity
|
||||
| { authorizationId: string; policy: UserChannelAuthorizationPolicy };
|
||||
export type UserChannelIdentityLink = { profileId: string; identity: UserChannelIdentity };
|
||||
export type UserChannelIdentityAuthorityFacts = {
|
||||
authorization?: UserChannelAuthorization;
|
||||
profileId: string;
|
||||
role: string | null;
|
||||
emails: string[];
|
||||
|
|
@ -41,9 +52,21 @@ export type UserChannelIdentityResult<T> =
|
|||
|
||||
export type UserChannelIdentityWorkerOperations = {
|
||||
"userProfiles.channelIdentity.change": {
|
||||
input: { action: "link" | "unlink"; profileId: string; identity: UserChannelIdentity };
|
||||
input:
|
||||
| { action: "link" | "unlink"; profileId: string; identity: UserChannelIdentity }
|
||||
| { action: "policy"; policy: UserChannelAuthorizationPolicy }
|
||||
| {
|
||||
action: "authorize";
|
||||
profileId: string;
|
||||
identity: UserChannelIdentity;
|
||||
policy: UserChannelAuthorizationPolicy;
|
||||
grant: GatewayAccessGrantRef | null;
|
||||
};
|
||||
output: UserChannelIdentityResult<
|
||||
{ kind: "linked"; link: UserChannelIdentityLink } | { kind: "unlinked"; removed: boolean }
|
||||
| { kind: "linked"; link: UserChannelIdentityLink }
|
||||
| { kind: "unlinked"; removed: boolean }
|
||||
| { kind: "policy" }
|
||||
| { kind: "authorized"; reference: UserChannelAuthorizationReference | undefined }
|
||||
>;
|
||||
};
|
||||
};
|
||||
|
|
@ -95,13 +118,7 @@ export type UserProfilesDatabase = {
|
|||
binding_id: string | null;
|
||||
created_at: number;
|
||||
};
|
||||
user_profile_identities: {
|
||||
provider: string;
|
||||
subject: string;
|
||||
profile_id: string;
|
||||
canonical_login: string | null;
|
||||
created_at: number;
|
||||
};
|
||||
user_profile_identities: DB["user_profile_identities"];
|
||||
};
|
||||
|
||||
export type ProfileDisplayRow = Pick<
|
||||
|
|
|
|||
|
|
@ -62,6 +62,8 @@ export const databaseWorkerCoreTestFiles = [
|
|||
"src/node-host/invoke.test.ts",
|
||||
"src/node-host/worker-runtime.test.ts",
|
||||
"src/skills/workshop/store.test.ts",
|
||||
"src/channels/message-access/operator-authority.test.ts",
|
||||
"src/auto-reply/reply/commands-allowlist.owner.test.ts",
|
||||
"src/channels/message-access/discord-native-acp-owner.test.ts",
|
||||
"src/channels/message-access/telegram-native-acp-owner.test.ts",
|
||||
"src/auto-reply/reply/commands-acp.owner.test.ts",
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue