mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
refactor(channels): discover ingress accounts in worker (#158477)
* refactor(channels): discover ingress accounts in worker * docs(state): refresh worker access inventory
This commit is contained in:
parent
9ee2b0b3b6
commit
15a0f34487
11 changed files with 194 additions and 165 deletions
|
|
@ -64,6 +64,11 @@ runners and registries. These helpers reuse their core owners; register the
|
|||
session fixture lifecycle explicitly. Use published runtime subpaths when
|
||||
they already expose the needed operation.
|
||||
|
||||
Await `listChannelIngressQueueAccountIdsForTests` from
|
||||
`channel-ingress-test-runtime` or `plugin-state-test-runtime`. It uses the shared
|
||||
read-only worker and leaves missing state uncreated. Join asynchronous database
|
||||
cleanup before removing a fixture's state directory.
|
||||
|
||||
### Available exports
|
||||
|
||||
| Export | Purpose |
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ title: "Database worker migration inventory"
|
|||
|
||||
<!-- Generated by scripts/database-worker-inventory.mjs. Edit its classification evidence, then regenerate. -->
|
||||
|
||||
This snapshot contains **540 non-test files and 2536 call expressions** for the five primitives below. The campaign previously reported 404 files; that is a historical estimate, not a fixed target or a count of call expressions. This inventory follows current source and excludes import-only matches, comments, tests, fixtures, and test support. Its scan scope and exclusions are explicit below.
|
||||
This snapshot contains **542 non-test files and 2536 call expressions** for the five primitives below. The campaign previously reported 404 files; that is a historical estimate, not a fixed target or a count of call expressions. This inventory follows current source and excludes import-only matches, comments, tests, fixtures, and test support. Its scan scope and exclusions are explicit below.
|
||||
|
||||
Regenerate with `pnpm db:worker-inventory:gen`; verify with `pnpm db:worker-inventory:check`. `node scripts/database-worker-inventory.mjs --json` emits every call's primitive, line, column, file owner, tier, and classification evidence. The script uses the repository's TypeScript parser and `rg`; it does not load application code or open a database.
|
||||
|
||||
|
|
@ -30,10 +30,10 @@ The scan covers JavaScript/TypeScript files under `src/`, `extensions/`, `packag
|
|||
|
||||
| Tier | Files | Call expressions |
|
||||
| ---- | ----: | ---------------: |
|
||||
| T1 | 403 | 2042 |
|
||||
| T1 | 403 | 2038 |
|
||||
| T2 | 54 | 250 |
|
||||
| T3 | 23 | 92 |
|
||||
| W | 60 | 152 |
|
||||
| W | 62 | 156 |
|
||||
|
||||
## Profile priority and current cutover status
|
||||
|
||||
|
|
@ -73,7 +73,7 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| Owner / file | Calls (Q/F/S/A/R) | First line | Exposure evidence |
|
||||
| ------------------------------------------------------------------------------------------------------------ | ----------------: | ---------: | ------------------------------------------------------------------------ |
|
||||
| **src/infra** · `src/infra/exec-approvals-sqlite.ts` | 2/3/0/0/0 | 42 | Approval-policy writes; write-coordination cutover owned separately |
|
||||
| **src/infra** · `src/infra/exec-approvals-store.ts` | 0/0/4/0/0 | 235 | Approval-policy writes; write-coordination cutover owned separately |
|
||||
| **src/infra** · `src/infra/exec-approvals-store.ts` | 0/0/4/0/0 | 231 | Approval-policy writes; write-coordination cutover owned separately |
|
||||
| **src/state** · `src/state/user-profiles.ts` | 12/4/2/0/0 | 96 | Profile creation; write-coordination cutover owned separately |
|
||||
| **src/config/sessions** · `src/config/sessions/session-accessor.sqlite-entry-read.ts` | 2/1/0/0/0 | 159 | Session-entry read kernel; inspect each caller's execution context |
|
||||
| **src/gateway** · `src/gateway/session-row-projection-materialize.ts` | 0/0/0/0/1 | 247 | Session-list row entries and membership; process-held incognito path |
|
||||
|
|
@ -85,7 +85,7 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| **src/agents** · `src/agents/plugin-model-catalog.ts` | 8/0/0/3/1 | 63 | Persisted catalog reads in prepared model runtime; also Doctor migration |
|
||||
| **src/config/sessions** · `src/config/sessions/session-sharing-store.kernel.ts` | 1/1/0/0/0 | 19 | Member-row kernel shared by session readers |
|
||||
| **src/config/sessions** · `src/config/sessions/session-sharing-store.ts` | 0/0/0/0/1 | 29 | Session-list member reads and membership mutations |
|
||||
| **src/infra** · `src/infra/device-pairing-store.ts` | 19/4/6/0/0 | 113 | Device/node RPC pairing reads and writes |
|
||||
| **src/infra** · `src/infra/device-pairing-store.ts` | 18/4/6/0/0 | 106 | Device/node RPC pairing reads and writes |
|
||||
| **extensions/memory-core** · `extensions/memory-core/src/memory-entry-origins.ts` | 6/0/0/0/4 | 67 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **extensions/memory-core** · `extensions/memory-core/src/memory-forget-index-sources.ts` | 5/0/0/0/2 | 78 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **extensions/memory-core** · `extensions/memory-core/src/memory-forget.ts` | 4/0/0/0/0 | 529 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
|
|
@ -141,7 +141,7 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| **src/audit** · `src/audit/message-execution-binding.ts` | 1/1/1/0/0 | 77 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/boards** · `src/boards/sqlite-board-store.kernel.ts` | 9/0/0/0/0 | 149 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/boards** · `src/boards/sqlite-board-store.ts` | 0/0/0/1/3 | 124 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/channels** · `src/channels/message/ingress-queue.ts` | 26/1/12/0/0 | 320 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/channels** · `src/channels/message/ingress-queue.ts` | 24/1/12/0/0 | 322 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/claws** · `src/claws/cron.ts` | 4/0/4/0/0 | 131 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/claws** · `src/claws/lifecycle-config-removal.ts` | 0/0/3/0/0 | 186 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/claws** · `src/claws/lifecycle-delete-support.ts` | 3/0/1/0/0 | 523 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
|
|
@ -272,7 +272,7 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| **src/config/sessions** · `src/config/sessions/transcript-payload.ts` | 0/1/0/0/0 | 101 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/cron** · `src/cron/scratch-store.ts` | 6/0/2/0/0 | 64 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/cron** · `src/cron/service/runtime-store.ts` | 0/0/1/0/0 | 65 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/cron** · `src/cron/store.ts` | 0/0/3/0/0 | 107 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/cron** · `src/cron/store.ts` | 0/0/2/0/0 | 167 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/cron** · `src/cron/store/doctor-inventory.ts` | 1/0/0/0/0 | 14 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/cron** · `src/cron/store/doctor.ts` | 0/0/1/0/0 | 94 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/cron** · `src/cron/store/job-name.ts` | 1/0/0/0/0 | 18 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
|
|
@ -305,7 +305,7 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| **src/gateway** · `src/gateway/progress-card-store.ts` | 0/0/0/1/1 | 42 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway** · `src/gateway/session-group-catalog.kernel.ts` | 10/0/1/0/0 | 33 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway** · `src/gateway/session-group-registration.kernel.ts` | 4/0/1/0/0 | 11 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway** · `src/gateway/session-row-membership-read.ts` | 0/0/0/0/1 | 215 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway** · `src/gateway/session-row-membership-read.ts` | 0/0/0/0/1 | 226 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/server-methods** · `src/gateway/server-methods/session-placement-read-projection.ts` | 0/0/0/0/1 | 130 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/inference-store.kernel.ts` | 7/3/0/0/0 | 121 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/local-workspace-store.ts` | 0/6/3/0/0 | 20 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
|
|
@ -315,8 +315,8 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-read-projection.ts` | 1/0/0/0/0 | 37 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-row-codec.ts` | 4/1/0/0/0 | 123 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-session-tool-operations.ts` | 11/0/0/0/0 | 99 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-store.ts` | 8/0/1/0/0 | 107 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-turn-claims.ts` | 12/0/0/0/0 | 115 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-store.ts` | 8/0/1/0/0 | 108 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-turn-claims.ts` | 12/0/0/0/0 | 110 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-workspace-journal.ts` | 9/0/0/0/0 | 52 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-workspace-reservation.ts` | 0/3/0/0/0 | 24 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-workspace-result.ts` | 15/0/0/0/0 | 112 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
|
|
@ -459,7 +459,7 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| **src/state** · `src/state/user-preferences.ts` | 0/0/1/0/0 | 58 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/state** · `src/state/user-profile-email.kernel.ts` | 0/1/0/0/0 | 32 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/state** · `src/state/user-profile-events.ts` | 1/0/0/0/0 | 92 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/state** · `src/state/user-profile-github-identity.ts` | 6/4/0/0/0 | 117 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/state** · `src/state/user-profile-github-identity.ts` | 6/4/0/0/0 | 114 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/state** · `src/state/user-profile-identity.read.ts` | 8/1/0/0/0 | 51 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/state** · `src/state/user-profile-mutation.ts` | 0/0/1/0/0 | 44 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
| **src/state** · `src/state/user-profiles-internal.ts` | 3/1/0/0/0 | 60 | Runtime/mixed candidate; main-thread reachability needs tracing |
|
||||
|
|
@ -499,7 +499,7 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
| **src/infra** · `src/infra/state-migrations.commitments.ts` | 0/0/1/0/0 | 115 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
| **src/infra** · `src/infra/state-migrations.debug-proxy.ts` | 0/0/1/0/0 | 421 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
| **src/infra** · `src/infra/state-migrations.device-auth.ts` | 1/2/1/0/0 | 100 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
| **src/infra** · `src/infra/state-migrations.device-identity.ts` | 2/1/1/0/0 | 149 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
| **src/infra** · `src/infra/state-migrations.device-identity.ts` | 2/1/1/0/0 | 150 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
| **src/infra** · `src/infra/state-migrations.exec-approvals.ts` | 0/0/1/0/0 | 147 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
| **src/infra** · `src/infra/state-migrations.managed-outgoing-images.ts` | 2/3/2/0/0 | 324 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
| **src/infra** · `src/infra/state-migrations.mcp-oauth.ts` | 2/2/1/0/0 | 149 | Schema/startup/migration candidate; verify no runtime caller |
|
||||
|
|
@ -565,65 +565,67 @@ Counts use `Q/F/S/A/R` in that order. Source locations are available in `--json`
|
|||
|
||||
### W
|
||||
|
||||
| Owner / file | Calls (Q/F/S/A/R) | First line | Exposure evidence |
|
||||
| ---------------------------------------------------------------------------------------------------- | ----------------: | ---------: | --------------------------------------------- |
|
||||
| **extensions/imessage** · `extensions/imessage/src/chat-db.worker.ts` | 2/2/0/0/0 | 42 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/logbook** · `extensions/logbook/src/store.worker.ts` | 10/7/0/0/0 | 267 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/memory-core** · `extensions/memory-core/src/memory/manager-search.worker.ts` | 1/0/0/0/1 | 79 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/qa-lab** · `extensions/qa-lab/src/execution-identity-storage-inspection.worker.ts` | 0/2/0/0/0 | 24 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/team-reports** · `extensions/team-reports/src/store.worker.ts` | 14/5/0/0/0 | 140 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents** · `src/agents/mcp-oauth-store.worker.ts` | 5/0/1/0/0 | 73 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/harness** · `src/agents/harness/native-hook-relay-store.worker.ts` | 0/0/1/0/0 | 26 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/sandbox** · `src/agents/sandbox/registry-import.worker.ts` | 0/0/1/0/0 | 15 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/sandbox** · `src/agents/sandbox/registry-write.worker.ts` | 0/0/1/0/0 | 12 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/sessions** · `src/agents/sessions/session-manager-metadata.worker.ts` | 0/0/0/2/0 | 221 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/worktrees** · `src/agents/worktrees/registry-retirement.worker.ts` | 2/0/1/0/0 | 31 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/worktrees** · `src/agents/worktrees/run-lease-store.worker.ts` | 0/0/1/0/0 | 15 | Worker implementation; keep SQL in this owner |
|
||||
| **src/channels** · `src/channels/message/ingress-queue-health.kernel.ts` | 2/0/0/0/0 | 17 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/provider-review-store.worker.ts` | 0/0/0/1/0 | 18 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-accessor.sqlite-archive.worker.ts` | 1/0/0/0/0 | 390 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-accessor.sqlite-mutation-worker.runtime.ts` | 0/0/0/1/0 | 323 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-accessor.sqlite-transcript-reports.worker.ts` | 0/0/0/1/0 | 137 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-entry-read.worker.ts` | 1/0/0/0/2 | 47 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-history-archive-pruning.worker.ts` | 3/0/0/2/1 | 31 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-transcript-projection-publication.worker.ts` | 0/1/0/0/0 | 80 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-transcript.worker.ts` | 0/0/0/0/8 | 144 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/dispatch.worker.ts` | 0/0/2/0/0 | 116 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/load.worker.ts` | 0/0/1/0/0 | 16 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/run-admission.worker.ts` | 0/0/4/0/0 | 51 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/run-recovery.worker.ts` | 0/0/1/0/0 | 18 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/runtime-maintenance.worker.ts` | 0/0/2/0/0 | 23 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/save.worker.ts` | 0/0/1/0/0 | 20 | Worker implementation; keep SQL in this owner |
|
||||
| **src/fleet** · `src/fleet/registry.worker.ts` | 0/0/1/0/0 | 22 | Worker implementation; keep SQL in this owner |
|
||||
| **src/gateway** · `src/gateway/operator-approval-store.worker.ts` | 0/0/1/0/0 | 35 | Worker implementation; keep SQL in this owner |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/inference-store.worker.ts` | 0/0/1/0/0 | 14 | Worker implementation; keep SQL in this owner |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/store.worker.ts` | 0/0/1/0/0 | 22 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/device-pairing-dispatch.worker.ts` | 0/0/1/0/0 | 74 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/exec-approvals-authorization.worker.ts` | 0/0/1/0/0 | 70 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/promotions-feed.worker.ts` | 0/0/2/0/0 | 34 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/push-apns-store.worker.ts` | 2/2/1/0/0 | 36 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/session-cost-usage-worker.ts` | 0/0/0/0/3 | 231 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/current-conversation-bindings.worker.ts` | 0/0/2/0/0 | 110 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-ack.worker.ts` | 0/0/1/0/0 | 12 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-enqueue.worker.ts` | 0/0/1/0/0 | 41 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-pending-failure.worker.ts` | 0/0/1/0/0 | 19 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-platform-lease.worker.ts` | 0/0/1/0/0 | 20 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-storage.worker.ts` | 0/0/1/0/0 | 252 | Worker implementation; keep SQL in this owner |
|
||||
| **src/plugin-state** · `src/plugin-state/plugin-blob-store.worker.ts` | 0/0/1/0/0 | 30 | Worker implementation; keep SQL in this owner |
|
||||
| **src/plugin-state** · `src/plugin-state/plugin-state.worker.ts` | 0/0/1/0/0 | 142 | Worker implementation; keep SQL in this owner |
|
||||
| **src/projects** · `src/projects/project-registry.worker.ts` | 0/0/1/0/0 | 79 | Worker implementation; keep SQL in this owner |
|
||||
| **src/sessions** · `src/sessions/session-state-events.worker.ts` | 0/0/2/0/0 | 37 | Worker implementation; keep SQL in this owner |
|
||||
| **src/skills/workshop** · `src/skills/workshop/store.worker.ts` | 0/0/2/0/0 | 86 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-agent-execution-cleanup.worker.ts` | 0/0/1/0/0 | 16 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-agent-execution.worker.ts` | 0/0/0/4/0 | 356 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-state-read.worker.ts` | 1/0/0/0/0 | 565 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-state-worker-runtime.ts` | 0/0/9/0/0 | 328 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/user-channel-identities.worker.ts` | 0/0/1/0/0 | 72 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/user-preferences.worker.ts` | 0/0/1/0/0 | 45 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/user-profiles.worker.ts` | 1/0/1/0/0 | 97 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-flow-maintenance.worker.ts` | 0/0/1/0/0 | 29 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-initial.worker.ts` | 0/0/1/0/0 | 51 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry-agent-event.worker.ts` | 0/0/1/0/0 | 29 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry-live-flow.worker.ts` | 0/0/1/0/0 | 23 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry-restore.worker.ts` | 0/0/2/0/0 | 88 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry.worker.ts` | 0/0/3/0/0 | 74 | Worker implementation; keep SQL in this owner |
|
||||
| Owner / file | Calls (Q/F/S/A/R) | First line | Exposure evidence |
|
||||
| ------------------------------------------------------------------------------------------------------- | ----------------: | ---------: | --------------------------------------------- |
|
||||
| **extensions/imessage** · `extensions/imessage/src/chat-db.worker.ts` | 2/2/0/0/0 | 42 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/logbook** · `extensions/logbook/src/store.worker.ts` | 10/7/0/0/0 | 267 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/memory-core** · `extensions/memory-core/src/memory/manager-search.worker.ts` | 1/0/0/0/1 | 85 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/qa-lab** · `extensions/qa-lab/src/execution-identity-storage-inspection.worker.ts` | 0/2/0/0/0 | 24 | Worker implementation; keep SQL in this owner |
|
||||
| **extensions/team-reports** · `extensions/team-reports/src/store.worker.ts` | 14/5/0/0/0 | 140 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents** · `src/agents/mcp-oauth-store.worker.ts` | 5/0/1/0/0 | 73 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/harness** · `src/agents/harness/native-hook-relay-store.worker.ts` | 0/0/1/0/0 | 26 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/sandbox** · `src/agents/sandbox/registry-import.worker.ts` | 0/0/1/0/0 | 15 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/sandbox** · `src/agents/sandbox/registry-write.worker.ts` | 0/0/1/0/0 | 12 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/sessions** · `src/agents/sessions/session-manager-metadata.worker.ts` | 0/0/0/2/0 | 221 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/worktrees** · `src/agents/worktrees/registry-retirement.worker.ts` | 2/0/1/0/0 | 31 | Worker implementation; keep SQL in this owner |
|
||||
| **src/agents/worktrees** · `src/agents/worktrees/run-lease-store.worker.ts` | 0/0/1/0/0 | 15 | Worker implementation; keep SQL in this owner |
|
||||
| **src/channels** · `src/channels/message/ingress-queue-health.kernel.ts` | 2/0/0/0/0 | 17 | Worker implementation; keep SQL in this owner |
|
||||
| **src/channels** · `src/channels/message/ingress-queue.kernel.ts` | 1/0/0/0/0 | 9 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/provider-review-store.worker.ts` | 0/0/0/1/0 | 18 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-accessor.sqlite-archive.worker.ts` | 1/0/0/0/0 | 390 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-accessor.sqlite-mutation-worker.runtime.ts` | 0/0/0/1/0 | 323 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-accessor.sqlite-transcript-reports.worker.ts` | 0/0/0/1/0 | 137 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-entry-read.worker.ts` | 1/0/0/0/3 | 53 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-history-archive-pruning.worker.ts` | 3/0/0/2/1 | 31 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-transcript-projection-publication.worker.ts` | 0/1/0/0/0 | 80 | Worker implementation; keep SQL in this owner |
|
||||
| **src/config/sessions** · `src/config/sessions/session-transcript.worker.ts` | 0/0/0/0/8 | 144 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/dispatch.worker.ts` | 0/0/2/0/0 | 116 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/load.worker.ts` | 0/0/1/0/0 | 16 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/run-admission.worker.ts` | 0/0/4/0/0 | 51 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/run-recovery.worker.ts` | 0/0/1/0/0 | 18 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/runtime-maintenance.worker.ts` | 0/0/2/0/0 | 23 | Worker implementation; keep SQL in this owner |
|
||||
| **src/cron** · `src/cron/store/save.worker.ts` | 0/0/1/0/0 | 20 | Worker implementation; keep SQL in this owner |
|
||||
| **src/fleet** · `src/fleet/registry.worker.ts` | 0/0/1/0/0 | 22 | Worker implementation; keep SQL in this owner |
|
||||
| **src/gateway** · `src/gateway/operator-approval-store.worker.ts` | 0/0/1/0/0 | 35 | Worker implementation; keep SQL in this owner |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/inference-store.worker.ts` | 0/0/1/0/0 | 14 | Worker implementation; keep SQL in this owner |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/placement-turn-claims.worker.ts` | 0/0/1/0/0 | 21 | Worker implementation; keep SQL in this owner |
|
||||
| **src/gateway/worker-environments** · `src/gateway/worker-environments/store.worker.ts` | 0/0/1/0/0 | 22 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/device-pairing-dispatch.worker.ts` | 0/0/1/0/0 | 74 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/exec-approvals-authorization.worker.ts` | 0/0/1/0/0 | 70 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/promotions-feed.worker.ts` | 0/0/2/0/0 | 34 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/push-apns-store.worker.ts` | 2/2/1/0/0 | 36 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra** · `src/infra/session-cost-usage-worker.ts` | 0/0/0/0/3 | 231 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/current-conversation-bindings.worker.ts` | 0/0/2/0/0 | 110 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-ack.worker.ts` | 0/0/1/0/0 | 12 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-enqueue.worker.ts` | 0/0/1/0/0 | 41 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-pending-failure.worker.ts` | 0/0/1/0/0 | 19 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-platform-lease.worker.ts` | 0/0/1/0/0 | 20 | Worker implementation; keep SQL in this owner |
|
||||
| **src/infra/outbound** · `src/infra/outbound/delivery-queue-storage.worker.ts` | 0/0/1/0/0 | 252 | Worker implementation; keep SQL in this owner |
|
||||
| **src/plugin-state** · `src/plugin-state/plugin-blob-store.worker.ts` | 0/0/1/0/0 | 30 | Worker implementation; keep SQL in this owner |
|
||||
| **src/plugin-state** · `src/plugin-state/plugin-state.worker.ts` | 0/0/1/0/0 | 142 | Worker implementation; keep SQL in this owner |
|
||||
| **src/projects** · `src/projects/project-registry.worker.ts` | 0/0/1/0/0 | 79 | Worker implementation; keep SQL in this owner |
|
||||
| **src/sessions** · `src/sessions/session-state-events.worker.ts` | 0/0/2/0/0 | 37 | Worker implementation; keep SQL in this owner |
|
||||
| **src/skills/workshop** · `src/skills/workshop/store.worker.ts` | 0/0/2/0/0 | 86 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-agent-execution-cleanup.worker.ts` | 0/0/1/0/0 | 16 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-agent-execution.worker.ts` | 0/0/0/4/0 | 356 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-state-read.worker.ts` | 1/0/0/0/0 | 565 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/openclaw-state-worker-runtime.ts` | 0/0/9/0/0 | 333 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/user-channel-identities.worker.ts` | 0/0/1/0/0 | 72 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/user-preferences.worker.ts` | 0/0/1/0/0 | 45 | Worker implementation; keep SQL in this owner |
|
||||
| **src/state** · `src/state/user-profiles.worker.ts` | 2/0/1/0/0 | 52 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-flow-maintenance.worker.ts` | 0/0/1/0/0 | 29 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-initial.worker.ts` | 0/0/1/0/0 | 51 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry-agent-event.worker.ts` | 0/0/1/0/0 | 29 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry-live-flow.worker.ts` | 0/0/1/0/0 | 23 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry-restore.worker.ts` | 0/0/2/0/0 | 88 | Worker implementation; keep SQL in this owner |
|
||||
| **src/tasks** · `src/tasks/task-registry.worker.ts` | 0/0/3/0/0 | 74 | Worker implementation; keep SQL in this owner |
|
||||
|
|
|
|||
|
|
@ -13,6 +13,7 @@ import type {
|
|||
PluginDoctorChannelIngressQueueAccess,
|
||||
PluginDoctorStateMigrationContext,
|
||||
} from "openclaw/plugin-sdk/runtime-doctor-migrations";
|
||||
import { closeOpenClawStateDatabaseAsync } from "openclaw/plugin-sdk/sqlite-runtime-testing";
|
||||
import { afterEach, describe, expect, it } from "vitest";
|
||||
import { stateMigrations } from "./doctor-contract-api.js";
|
||||
|
||||
|
|
@ -74,6 +75,7 @@ async function withStateDir<T>(fn: (stateDir: string) => Promise<T>): Promise<T>
|
|||
try {
|
||||
return await fn(stateDir);
|
||||
} finally {
|
||||
await closeOpenClawStateDatabaseAsync();
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
await fs.rm(stateDir, { recursive: true, force: true });
|
||||
}
|
||||
|
|
|
|||
|
|
@ -107,6 +107,7 @@ const reviewed = new Map([
|
|||
],
|
||||
]);
|
||||
const workerModules = new Set([
|
||||
"src/channels/message/ingress-queue.kernel.ts",
|
||||
"src/channels/message/ingress-queue-health.kernel.ts",
|
||||
"src/state/openclaw-state-worker-runtime.ts",
|
||||
"src/config/sessions/session-accessor.sqlite-mutation-worker.runtime.ts",
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ export type ChannelIngressPressureHealth = {
|
|||
};
|
||||
|
||||
type ChannelIngressReadOperations = {
|
||||
"channelIngress.accounts": { input: { channelId: string }; output: string[] };
|
||||
"channelIngress.failedHealth": { input: undefined; output: ChannelIngressFailedHealth[] };
|
||||
"channelIngress.pressureHealth": {
|
||||
input: { now: number };
|
||||
|
|
@ -40,6 +41,9 @@ export function isChannelIngressReadCommand(value: unknown): value is ChannelIng
|
|||
if (value.type === "channelIngress.failedHealth") {
|
||||
return value.input === undefined;
|
||||
}
|
||||
if (value.type === "channelIngress.accounts") {
|
||||
return isRecord(value.input) && typeof value.input.channelId === "string";
|
||||
}
|
||||
return (
|
||||
value.type === "channelIngress.pressureHealth" &&
|
||||
isRecord(value.input) &&
|
||||
|
|
|
|||
|
|
@ -7,11 +7,20 @@ import type {
|
|||
ChannelIngressReadCommand,
|
||||
ChannelIngressReadReply,
|
||||
} from "./ingress-queue-read-contract.js";
|
||||
import { listChannelIngressAccountsInDatabase } from "./ingress-queue.kernel.js";
|
||||
|
||||
export function readChannelIngressInDatabase(
|
||||
db: DatabaseSync,
|
||||
command: ChannelIngressReadCommand,
|
||||
): ChannelIngressReadReply {
|
||||
if (command.type === "channelIngress.accounts") {
|
||||
return {
|
||||
ok: true,
|
||||
sourceAdmitted: true,
|
||||
type: command.type,
|
||||
result: listChannelIngressAccountsInDatabase(db, command.input),
|
||||
};
|
||||
}
|
||||
if (command.type === "channelIngress.failedHealth") {
|
||||
return {
|
||||
ok: true,
|
||||
|
|
|
|||
18
src/channels/message/ingress-queue.kernel.ts
Normal file
18
src/channels/message/ingress-queue.kernel.ts
Normal file
|
|
@ -0,0 +1,18 @@
|
|||
import type { DatabaseSync } from "node:sqlite";
|
||||
import { executeSqliteQuerySync, getNodeSqliteKysely } from "../../infra/kysely-sync.js";
|
||||
import type { DB } from "../../state/openclaw-state-db.generated.js";
|
||||
|
||||
export function listChannelIngressAccountsInDatabase(
|
||||
db: DatabaseSync,
|
||||
input: { channelId: string },
|
||||
): string[] {
|
||||
return executeSqliteQuerySync(
|
||||
db,
|
||||
getNodeSqliteKysely<Pick<DB, "channel_ingress_events">>(db)
|
||||
.selectFrom("channel_ingress_events")
|
||||
.select("account_id")
|
||||
.distinct()
|
||||
.where("channel_id", "=", input.channelId)
|
||||
.orderBy("account_id", "asc"),
|
||||
).rows.map((row) => row.account_id);
|
||||
}
|
||||
|
|
@ -1,9 +1,9 @@
|
|||
// Read-only ingress listing tests cover access that must not create shared state.
|
||||
import fs from "node:fs/promises";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { closeOpenClawStateDatabaseForTest } from "../../state/openclaw-state-db.js";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import * as sqliteQueries from "../../infra/kysely-sync.js";
|
||||
import { closeOpenClawStateDatabaseByPathAsync } from "../../state/openclaw-state-db.js";
|
||||
import { withOpenClawTestState } from "../../test-utils/openclaw-test-state.js";
|
||||
import {
|
||||
createChannelIngressQueue,
|
||||
listChannelIngressQueueAccountIdsReadOnly,
|
||||
|
|
@ -11,55 +11,66 @@ import {
|
|||
|
||||
describe("read-only listing access", () => {
|
||||
it("lists without creating the shared state database", async () => {
|
||||
const stateDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-ingress-readonly-"));
|
||||
try {
|
||||
const sqlitePath = path.join(stateDir, "state", "openclaw.sqlite");
|
||||
await expect(fs.access(sqlitePath)).rejects.toThrow();
|
||||
await withOpenClawTestState(
|
||||
{ layout: "state-only", prefix: "openclaw-ingress-readonly-", applyEnv: false },
|
||||
async ({ stateDir, statePath }) => {
|
||||
const sqlitePath = statePath("state", "openclaw.sqlite");
|
||||
await expect(fs.access(sqlitePath)).rejects.toThrow();
|
||||
|
||||
const reader = createChannelIngressQueue<{ text: string }>({
|
||||
channelId: "line",
|
||||
accountId: "default",
|
||||
stateDir,
|
||||
access: "read-only",
|
||||
});
|
||||
// The read-only opener never creates, migrates or configures the file, so a
|
||||
// caller that runs before it owns the state cannot bring the store into being.
|
||||
// Account discovery runs before the inspection facade is even opened, so it is
|
||||
// the first thing that could create the store.
|
||||
expect(
|
||||
await listChannelIngressQueueAccountIdsReadOnly({ channelId: "line", stateDir }),
|
||||
).toEqual([]);
|
||||
await expect(fs.access(sqlitePath)).rejects.toThrow();
|
||||
const reader = createChannelIngressQueue<{ text: string }>({
|
||||
channelId: "line",
|
||||
accountId: "default",
|
||||
stateDir,
|
||||
access: "read-only",
|
||||
});
|
||||
// The read-only opener never creates, migrates or configures the file, so a
|
||||
// caller that runs before it owns the state cannot bring the store into being.
|
||||
// Account discovery runs before the inspection facade is even opened, so it is
|
||||
// the first thing that could create the store.
|
||||
expect(
|
||||
await listChannelIngressQueueAccountIdsReadOnly({ channelId: "line", stateDir }),
|
||||
).toEqual([]);
|
||||
await expect(fs.access(sqlitePath)).rejects.toThrow();
|
||||
|
||||
expect(await reader.listPending({ limit: "all" })).toEqual([]);
|
||||
expect(await reader.listClaims()).toEqual([]);
|
||||
expect(await reader.listFailed?.({ limit: "all" })).toEqual([]);
|
||||
await expect(fs.access(sqlitePath)).rejects.toThrow();
|
||||
expect(await reader.listPending({ limit: "all" })).toEqual([]);
|
||||
expect(await reader.listClaims()).toEqual([]);
|
||||
expect(await reader.listFailed?.({ limit: "all" })).toEqual([]);
|
||||
await expect(fs.access(sqlitePath)).rejects.toThrow();
|
||||
|
||||
// A read-write queue is what actually creates it, and the read-only reader then
|
||||
// sees the same rows - so the empty results above are the access mode, not a
|
||||
// broken reader.
|
||||
await createChannelIngressQueue<{ text: string }>({
|
||||
channelId: "line",
|
||||
accountId: "default",
|
||||
stateDir,
|
||||
}).enqueue("evt-1", { text: "hello" });
|
||||
await fs.access(sqlitePath);
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
// A read-write queue is what actually creates it, and the read-only reader then
|
||||
// sees the same rows - so the empty results above are the access mode, not a
|
||||
// broken reader.
|
||||
await createChannelIngressQueue<{ text: string }>({
|
||||
channelId: "line",
|
||||
accountId: "default",
|
||||
stateDir,
|
||||
}).enqueue("evt-1", { text: "hello" });
|
||||
await fs.access(sqlitePath);
|
||||
await closeOpenClawStateDatabaseByPathAsync(sqlitePath);
|
||||
|
||||
const after = createChannelIngressQueue<{ text: string }>({
|
||||
channelId: "line",
|
||||
accountId: "default",
|
||||
stateDir,
|
||||
access: "read-only",
|
||||
});
|
||||
expect((await after.listPending({ limit: "all" })).map((row) => row.id)).toEqual(["evt-1"]);
|
||||
expect(
|
||||
await listChannelIngressQueueAccountIdsReadOnly({ channelId: "line", stateDir }),
|
||||
).toEqual(["default"]);
|
||||
} finally {
|
||||
closeOpenClawStateDatabaseForTest();
|
||||
await fs.rm(stateDir, { recursive: true, force: true });
|
||||
}
|
||||
const after = createChannelIngressQueue<{ text: string }>({
|
||||
channelId: "line",
|
||||
accountId: "default",
|
||||
stateDir,
|
||||
access: "read-only",
|
||||
});
|
||||
expect((await after.listPending({ limit: "all" })).map((row) => row.id)).toEqual(["evt-1"]);
|
||||
const hostQueries = vi
|
||||
.spyOn(sqliteQueries, "executeSqliteQuerySync")
|
||||
.mockImplementation(() => {
|
||||
throw new Error(
|
||||
"Ingress account discovery must not query SQLite on the calling thread",
|
||||
);
|
||||
});
|
||||
try {
|
||||
expect(
|
||||
await listChannelIngressQueueAccountIdsReadOnly({ channelId: "line", stateDir }),
|
||||
).toEqual(["default"]);
|
||||
expect(hostQueries).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
hostQueries.mockRestore();
|
||||
}
|
||||
},
|
||||
);
|
||||
});
|
||||
});
|
||||
|
|
|
|||
|
|
@ -10,12 +10,14 @@ import {
|
|||
executeSqliteQueryTakeFirstSync,
|
||||
getNodeSqliteKysely,
|
||||
} from "../../infra/kysely-sync.js";
|
||||
import { executeExistingOpenClawStateRead } from "../../state/openclaw-state-db-readonly.js";
|
||||
import type { DB as OpenClawStateKyselyDatabase } from "../../state/openclaw-state-db.generated.js";
|
||||
import {
|
||||
openExistingOpenClawStateDatabaseReadOnly,
|
||||
openOpenClawStateDatabase,
|
||||
runOpenClawStateWriteTransaction,
|
||||
} from "../../state/openclaw-state-db.js";
|
||||
import { resolveChannelIngressStateEnv } from "./ingress-queue-client.js";
|
||||
import {
|
||||
FAILED_NULL_PAYLOAD_SENTINEL,
|
||||
parseFailedPayload,
|
||||
|
|
@ -448,53 +450,25 @@ function queueNameForParts(channelId: string, accountId: string): string {
|
|||
return JSON.stringify([channelId, accountId]);
|
||||
}
|
||||
|
||||
/** Lists account ids that hold any ingress rows for a channel, so doctor
|
||||
* migrations can sweep durable state whose account is gone from config. */
|
||||
export function listChannelIngressQueueAccountIds(params: {
|
||||
channelId: string;
|
||||
stateDir?: string;
|
||||
}): string[] {
|
||||
const channelId = normalizePart(params.channelId, "unknown");
|
||||
const database = openChannelIngressDatabase(params.stateDir);
|
||||
const rows = executeSqliteQuerySync(
|
||||
database.db,
|
||||
getChannelIngressKysely(database.db)
|
||||
.selectFrom("channel_ingress_events")
|
||||
.select("account_id")
|
||||
.distinct()
|
||||
.where("channel_id", "=", channelId)
|
||||
.orderBy("account_id", "asc"),
|
||||
).rows;
|
||||
return rows.map((row) => row.account_id);
|
||||
}
|
||||
|
||||
/**
|
||||
* Account discovery for callers that must not touch durable state yet. Uses the
|
||||
* non-creating read-only opener, so an absent store yields no accounts instead of
|
||||
* being created and migrated by the lookup itself.
|
||||
*/
|
||||
/** Account discovery never creates or migrates a missing database. */
|
||||
export async function listChannelIngressQueueAccountIdsReadOnly(params: {
|
||||
channelId: string;
|
||||
stateDir?: string;
|
||||
}): Promise<string[]> {
|
||||
const channelId = normalizePart(params.channelId, "unknown");
|
||||
const handle = await openChannelIngressDatabaseForListing(params.stateDir, "read-only");
|
||||
if (!handle) {
|
||||
const reply = await executeExistingOpenClawStateRead(
|
||||
{ env: resolveChannelIngressStateEnv(params.stateDir) },
|
||||
{
|
||||
type: "channelIngress.accounts",
|
||||
input: { channelId: normalizePart(params.channelId, "unknown") },
|
||||
},
|
||||
);
|
||||
if (!reply) {
|
||||
return [];
|
||||
}
|
||||
try {
|
||||
return executeSqliteQuerySync(
|
||||
handle.db,
|
||||
getChannelIngressKysely(handle.db)
|
||||
.selectFrom("channel_ingress_events")
|
||||
.select("account_id")
|
||||
.distinct()
|
||||
.where("channel_id", "=", channelId)
|
||||
.orderBy("account_id", "asc"),
|
||||
).rows.map((row) => row.account_id);
|
||||
} finally {
|
||||
handle.release();
|
||||
if (!reply.ok || reply.type !== "channelIngress.accounts") {
|
||||
throw new Error("Channel ingress account reader returned an unexpected result");
|
||||
}
|
||||
return reply.result;
|
||||
}
|
||||
|
||||
/** Creates a durable channel/account-scoped ingress queue backed by the OpenClaw state database. */
|
||||
|
|
|
|||
|
|
@ -8,7 +8,7 @@ export {
|
|||
export { createHostChannelIngressRuntime } from "../channels/message-access/runtime.js";
|
||||
export {
|
||||
createChannelIngressQueue as createChannelIngressQueueForTests,
|
||||
listChannelIngressQueueAccountIds as listChannelIngressQueueAccountIdsForTests,
|
||||
listChannelIngressQueueAccountIdsReadOnly as listChannelIngressQueueAccountIdsForTests,
|
||||
} from "../channels/message/ingress-queue.js";
|
||||
export { closeOpenClawStateDatabaseForTest } from "../state/openclaw-state-db.js";
|
||||
export { withRegisteredChannelIngress } from "./test-helpers/registered-channel-ingress.js";
|
||||
|
|
|
|||
|
|
@ -96,6 +96,9 @@ function captureCommand(command: OpenClawStateReadCommand): OpenClawStateReadCom
|
|||
if (command.type === "channelIngress.pressureHealth") {
|
||||
return { type: command.type, input: { now: command.input.now } };
|
||||
}
|
||||
if (command.type === "channelIngress.accounts") {
|
||||
return { type: command.type, input: { channelId: command.input.channelId } };
|
||||
}
|
||||
if (command.type === "cron.jobNames") {
|
||||
return { ...command, jobIds: [...command.jobIds] };
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue