diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json index 65bc523e47e7..b33109166d5d 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness-runtime.json @@ -1 +1 @@ -{"contentHash":"31bcb6e7b685c67391cc16a005f4ff55ffd97d382a10f4b1796ddf2fd303a308","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"} +{"contentHash":"92726144940b7b1207374a945baeabf1a9b31058a4ac8ada3a0829db2ce925b9","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json index d8f24331b9b5..30565772ef44 100644 --- a/docs/.generated/plugin-sdk-api-baseline/agent-harness.json +++ b/docs/.generated/plugin-sdk-api-baseline/agent-harness.json @@ -1 +1 @@ -{"contentHash":"108af57a4b4bc57c70912fc40319876bb6ec994cce59d2b9c991980e18c4f80f","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} +{"contentHash":"ddb50f1e611d363742cf493741fad27f04884f318b1f68e3bc43e1a63ccca4cc","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-core.json b/docs/.generated/plugin-sdk-api-baseline/channel-core.json index 84a99d349f4f..9d6a7d300da0 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-core.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-core.json @@ -1 +1 @@ -{"contentHash":"dd255204bf11dccfe9d866e7a6725935bfead7378a52ddfa6b60b9a58c0cdf76","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} +{"contentHash":"e64d40bbc2f99f5834b46918359f1680ca40e22caac6642007fb7eb450dcab67","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json index 2c34141dd292..5b20ed737020 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-entry-contract.json @@ -1 +1 @@ -{"contentHash":"4bc2ad950266148cfbb2acc7996de31401d88dd5f80ea28754e76f3bd66ef0d2","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"} +{"contentHash":"64ade6b3025afc1547aea966303d92341431770286afeb19ca185f0a68023e77","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json b/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json index bb54e32fc363..a72172bb5642 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-inbound.json @@ -1 +1 @@ -{"contentHash":"f0be441226a760f75dc1634f93d7f28586a0b18f690e2dd0f87c5360f2cb7f03","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"} +{"contentHash":"2cf8445ec79b755119c0dd1efd24e4f06a0cb5f79e0f661619013a4c1961f7f8","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-ingress-runtime.json b/docs/.generated/plugin-sdk-api-baseline/channel-ingress-runtime.json index c0734dd626f8..a16bce455729 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-ingress-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-ingress-runtime.json @@ -1 +1 @@ -{"contentHash":"a8de2c1ed05d17353b0ffa7a6e94c802fba1acdcfea71ff0d415a7f1944fe7ba","entrypoint":"channel-ingress-runtime","importSpecifier":"openclaw/plugin-sdk/channel-ingress-runtime"} +{"contentHash":"bde82070f46bb731c62284e9b0e5631376d264d0903fa5e9e32ca7c5ac047a7d","entrypoint":"channel-ingress-runtime","importSpecifier":"openclaw/plugin-sdk/channel-ingress-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-message.json b/docs/.generated/plugin-sdk-api-baseline/channel-message.json index c6937ba3e11d..3026a93d6ced 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-message.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-message.json @@ -1 +1 @@ -{"contentHash":"8bf7d1fd1e21861fd472365e0ab10ca811800e1b0cb9e6aac0e4d6a41613a9a2","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} +{"contentHash":"3bb24706f7e76c556192dd100a856a490220490248892368bf05d1a6b3a67615","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json index 2bbe61f4c659..0a70ee16a584 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-outbound.json @@ -1 +1 @@ -{"contentHash":"d430c3def27acb48a5356dafac17b89b0cbfcb1d3f959136b8c7ae82586d87df","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} +{"contentHash":"7ba02e716d0039ef36ebbc74c524e818619714e6c7983c9cdf34e325a7d5fa95","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-pairing.json b/docs/.generated/plugin-sdk-api-baseline/channel-pairing.json index d9f559e8fb26..13654c0b935b 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-pairing.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-pairing.json @@ -1 +1 @@ -{"contentHash":"0d9fc3f9ef06aaa9995c5a670d62b99cfce25f0b2c42152aa33214f9dec78b19","entrypoint":"channel-pairing","importSpecifier":"openclaw/plugin-sdk/channel-pairing"} +{"contentHash":"8cf3a92d14f57bcd63a09f704a5b9dabb48a83d58df22dead6b6e859e4ba69cf","entrypoint":"channel-pairing","importSpecifier":"openclaw/plugin-sdk/channel-pairing"} diff --git a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json index 588ab76e1ecb..1dfdaf744351 100644 --- a/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json +++ b/docs/.generated/plugin-sdk-api-baseline/channel-plugin-common.json @@ -1 +1 @@ -{"contentHash":"aa9a880fa3abde607a3664e55cf6bd35966528d436c0438aa3c1b7607c8e46e7","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"} +{"contentHash":"7e333c654dbaa9943c83e9fa3e18947057481ecadd3962d907621724737b69ea","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"} diff --git a/docs/.generated/plugin-sdk-api-baseline/core.json b/docs/.generated/plugin-sdk-api-baseline/core.json index a97ddd3aecfe..3a93e0711caf 100644 --- a/docs/.generated/plugin-sdk-api-baseline/core.json +++ b/docs/.generated/plugin-sdk-api-baseline/core.json @@ -1 +1 @@ -{"contentHash":"4c667b7e88cd752c294bcc731bc157c1ccb2a7aa4afc58f4013409191f3a9b28","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} +{"contentHash":"9d6ee27bf084b6e4383bf6a4da070f1f2c0722b244776bd3c511703038df3714","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"} diff --git a/docs/.generated/plugin-sdk-api-baseline/discord.json b/docs/.generated/plugin-sdk-api-baseline/discord.json index 2c7c6e047633..299e11659a27 100644 --- a/docs/.generated/plugin-sdk-api-baseline/discord.json +++ b/docs/.generated/plugin-sdk-api-baseline/discord.json @@ -1 +1 @@ -{"contentHash":"027b2280171324b628865fb4cdf9f72c13c059385a2fbec5b650c69d66beecdc","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} +{"contentHash":"b8a7b3fb6c8d00dd5e944bd569ee22cf2ac994a453297fd1304d0443fd52bb86","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"} diff --git a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json index c746dcca3132..2b48dd1222ea 100644 --- a/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json +++ b/docs/.generated/plugin-sdk-api-baseline/inbound-reply-dispatch.json @@ -1 +1 @@ -{"contentHash":"94ee499e1ea0653d3e7daed7c6aaa2db8d7c3026433eba62b19a5f6cda93028b","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"} +{"contentHash":"b37f4e056c53f090a5df1468271bf2c556fd6caad66f0dece4146e4347c9fe75","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"} diff --git a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json index cfe816da9801..150c70d32f4d 100644 --- a/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/meeting-runtime.json @@ -1 +1 @@ -{"contentHash":"69fe3479e42b74771fb86d228264e8e19a3a16e1a07df4f89499111c2df2bc13","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} +{"contentHash":"a51fed7fe588e420896c747abc547a9d28c216ba682274bae683c308644d51d7","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json index 0bb25fd1bf8f..676dffd9f782 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-entry.json @@ -1 +1 @@ -{"contentHash":"f058be9cbffbcb6e1dafe69cbc701c7cc06d1ffc90323256685e31b111f91193","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} +{"contentHash":"c58b25311e4cbb436b35e084789748942fb12ebdf97a8ea8c085e5fc110ddf89","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"} diff --git a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json index 3eb80e8bec0e..df3031b89a2f 100644 --- a/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/plugin-runtime.json @@ -1 +1 @@ -{"contentHash":"e7ac61e29d41e43a4c780831a8aa0e5bed10e47630a36879f2555216eb1cce54","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} +{"contentHash":"8f95e88158576bbef982b4962b917c2b937e196a617a3b6a76cecef88a029921","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json index 1e1d21020f8a..9cc304bf6257 100644 --- a/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json +++ b/docs/.generated/plugin-sdk-api-baseline/provider-catalog-runtime.json @@ -1 +1 @@ -{"contentHash":"12f6874f2b2c49f6fa22b91acdae820a725564d381a5b26dc4b2b4e65d0ef803","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"} +{"contentHash":"e3b10b2c62e456729d096f10c45538b40a61a8568da6b88fa7455a2979668454","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"} diff --git a/docs/.generated/plugin-sdk-api-baseline/runtime-store.json b/docs/.generated/plugin-sdk-api-baseline/runtime-store.json index d7a05cfce8fd..0ed18655205a 100644 --- a/docs/.generated/plugin-sdk-api-baseline/runtime-store.json +++ b/docs/.generated/plugin-sdk-api-baseline/runtime-store.json @@ -1 +1 @@ -{"contentHash":"b965caa74f02ea0fa0b4b5a57bfacd0699df991452b5aa95f21237858930508f","entrypoint":"runtime-store","importSpecifier":"openclaw/plugin-sdk/runtime-store"} +{"contentHash":"d0af7cc6a491a50a90271814d47b91116eb7e600945ead8d3f54bfce6075bbdd","entrypoint":"runtime-store","importSpecifier":"openclaw/plugin-sdk/runtime-store"} diff --git a/docs/.generated/plugin-sdk-api-baseline/session-catalog.json b/docs/.generated/plugin-sdk-api-baseline/session-catalog.json index 6135612f8dec..76e0fbd9e836 100644 --- a/docs/.generated/plugin-sdk-api-baseline/session-catalog.json +++ b/docs/.generated/plugin-sdk-api-baseline/session-catalog.json @@ -1 +1 @@ -{"contentHash":"1c76ac2e6a40a56ae0dbf5f6e2b5d911883f8e11cc7c0f435f6cce0c185a60c2","entrypoint":"session-catalog","importSpecifier":"openclaw/plugin-sdk/session-catalog"} +{"contentHash":"2e56a2b98213e0f297084d33cce8880d9dfe931ae74b02f753c988fb82ec088d","entrypoint":"session-catalog","importSpecifier":"openclaw/plugin-sdk/session-catalog"} diff --git a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json index b7d31acd4088..4ec66ffff780 100644 --- a/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json +++ b/docs/.generated/plugin-sdk-api-baseline/tool-plugin.json @@ -1 +1 @@ -{"contentHash":"f76d5908b66bb49851cccc530dc63fc5f3352b327963f1091a0f34fe9ccc693b","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} +{"contentHash":"db4a41272079b955d5426e2bea3f1db87e39d4168abd1109ef9ec937a88d89b6","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"} diff --git a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json index f8e224cc6ced..e3591172b43a 100644 --- a/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json +++ b/docs/.generated/plugin-sdk-api-baseline/webhook-ingress.json @@ -1 +1 @@ -{"contentHash":"e9dd8d432098b26f2e102fe41edee7c556558ea9b26a0077da72b0967d6649d2","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} +{"contentHash":"0488f4915160e0b0ea4e574d6800ced82e88219e9a60aca0de5599879d0cbec4","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"} diff --git a/docs/channels/qa-channel.md b/docs/channels/qa-channel.md index 537c8f4b8099..ad3e658e9f17 100644 --- a/docs/channels/qa-channel.md +++ b/docs/channels/qa-channel.md @@ -78,6 +78,15 @@ Full repo-backed scenario suite: pnpm openclaw qa suite ``` +The isolated `channel-participant-identity-inspection` scenario enables +execution identity before startup, exercises DM, group, senderless, same- and +mixed-participant collect paths, proves an ingress rejection creates no audit +rows, and compares JSON plus human CLI inspection across Gateway restart: + +```bash +pnpm openclaw qa suite --scenario channel-participant-identity-inspection +``` + Runs scenarios in parallel against the QA gateway lane. See [QA overview](/concepts/qa-e2e-automation) for scenarios, profiles, and provider modes. Docker-backed QA site (gateway + QA Lab debugger UI in one stack): diff --git a/docs/cli/audit.md b/docs/cli/audit.md index a1234196f6a8..62c57f3b61ea 100644 --- a/docs/cli/audit.md +++ b/docs/cli/audit.md @@ -135,6 +135,15 @@ immutable connection-time audit fact. Ordinary session provenance stores no display label. An optional bounded, secret-redacted label can be retained only in execution identity after that audit storage is explicitly enabled. +An admitted channel run can also show a pseudonymized person invoker. That +identity comes from the channel's host-owned admission boundary, never from the +conversation, room, route, account, thread, message, transport, session key, or +display name. A collected run shows the person only when all queued inputs +carry valid evidence for the same participant; mixed or missing evidence shows +an unknown invoker. Its `channel/admission` receipt is enforced only when the +participant affected every contributing access decision; otherwise it is +attribution-only. + A terminal approval receipt shows `allowed` or `denied`, its stable reason code, enforcement state, authoritative source boundary, policy and grant references, context fields used, and remediation. Expired and cancelled diff --git a/docs/concepts/qa-e2e-automation.md b/docs/concepts/qa-e2e-automation.md index 22ef49be2455..d3206b65b4ba 100644 --- a/docs/concepts/qa-e2e-automation.md +++ b/docs/concepts/qa-e2e-automation.md @@ -1121,6 +1121,13 @@ Seed assets live in `qa/`: - `qa/scenarios/index.yaml` - `qa/scenarios//*.yaml` +Identity-sensitive channel changes use the isolated +`channel-participant-identity-inspection` QA Channel flow. It drives a real +ephemeral Gateway and mock provider, then inspects admitted runs with the same +`openclaw audit --run ... --explain` JSON and human surfaces operators use. +The flow includes lifecycle-owned restart and a row-count check for rejected +pre-run ingress. + These are intentionally in git so the QA plan is visible to both humans and the agent. diff --git a/docs/gateway/audit.md b/docs/gateway/audit.md index 67503a255864..d5cc7222c40b 100644 --- a/docs/gateway/audit.md +++ b/docs/gateway/audit.md @@ -86,10 +86,16 @@ state for these fields: - applicable grants and assurance evidence; - parent or child lineage when available. -The foundation records direct local CLI ingress and Gateway boot-system ingress -at their authoritative producers. Generic public ingress remains explicitly -unknown when its boundary cannot prove a more specific source; OpenClaw never -infers ingress or invoker identity from a session key. A direct local execution +The foundation records direct local CLI ingress, Gateway boot-system ingress, +and admitted channel participants at their authoritative producers. For a +channel run, the person invoker comes from host-minted admission evidence; the +room, route, account, thread, message, and transport remain non-principal +facts. Collected messages retain a person only when every contribution proves +the same participant. Mixed, missing, invalid, stale, or unminted evidence is +unknown, and an adapter that explicitly lacks support is unsupported. OpenClaw +never reconstructs a participant from `SenderId`, `From`, session keys, or +routing metadata. Other public ingress remains explicitly unknown when its +boundary cannot prove a more specific source. A direct local execution is `unattributed`: the Gateway cell, local CLI ingress, configured agent, and runtime binding are present, but no durable invoker principal is supplied at this boundary. A run becomes @@ -115,6 +121,13 @@ is `not-applicable`, its policy and grant references are empty, and its reason states that no identity-aware policy or grant evaluation was proven. This is an explanation of admission evidence, not an enforcement claim. +Admitted channel runs also project a `channel/admission` decision receipt after +their exact context/execution/run tuple is queued. Coverage is `enforced` only +when every contributing ingress decision was participant-aware and +outcome-affecting. Wildcard/open policy and explicit attribution-only adapters +remain `attribution-only`; mixed or missing evidence is `unknown`. Identity and +the corresponding decision share the existing audit-writer FIFO. + When the same `runId` has a retained terminal row in `operator_approvals`, the inspector also reads its owner-local `operator_approval_execution_identities` binding. Only an exact context, execution, and run tuple projects the approval diff --git a/docs/plugins/sdk-channel-ingress.md b/docs/plugins/sdk-channel-ingress.md index a4a1190ca9d6..def8dfa7efea 100644 --- a/docs/plugins/sdk-channel-ingress.md +++ b/docs/plugins/sdk-channel-ingress.md @@ -23,6 +23,7 @@ import { defineStableChannelIngressIdentity, resolveChannelMessageIngress, } from "openclaw/plugin-sdk/channel-ingress-runtime"; +import { buildChannelInboundEventContext } from "openclaw/plugin-sdk/channel-inbound"; const identity = defineStableChannelIngressIdentity({ key: "platform-user-id", @@ -49,6 +50,13 @@ const result = await resolveChannelMessageIngress({ readStoreAllowFrom, command: hasControlCommand ? { allowTextCommands: true, hasControlCommand } : undefined, }); + +const ctx = buildChannelInboundEventContext({ + // Pass the exact host result; do not rebuild participant evidence from + // SenderId, From, session keys, routes, rooms, or message metadata. + channelIngress: result, + // ...normalized channel facts +}); ``` Do not precompute effective allowlists, command owners, or command groups. @@ -74,6 +82,17 @@ Deprecated third-party SDK helpers may rebuild older shapes internally. New bundled receive paths should not translate modern results back into local DTOs. +When execution-identity audit collection is enabled, the inbound context +builder privately carries the exact admitted participant to run admission. +Queue collection retains attribution only when every contribution has valid +evidence for the same participant; mixed, missing, stale, or unminted evidence +is `unknown`. The carrier is opaque, bounded, one-shot, and diagnostic only. +Adapters whose access owner cannot pass the resolver result may use +`createChannelParticipantAdmissionEvidence(...)` on this same SDK subpath and +pass it as `channelParticipantEvidence`; that path is attribution-only, never +proof that participant identity affected access policy. Mark adapters that +cannot supply participant identity with `channelIngress: "unsupported"`. + ## Access groups `accessGroup:` entries stay redacted. Core resolves static diff --git a/docs/plugins/sdk-channel-plugins.md b/docs/plugins/sdk-channel-plugins.md index 84a7612b489e..bbd69bba3276 100644 --- a/docs/plugins/sdk-channel-plugins.md +++ b/docs/plugins/sdk-channel-plugins.md @@ -118,6 +118,13 @@ the resolved state or decision. See [Channel ingress API](/plugins/sdk-channel-ingress) for the API design, ownership boundary, and test expectations. +Pass the exact resolver result to `buildChannelInboundEventContext` as +`channelIngress`. This preserves host-minted participant evidence through +queued run admission without exposing it in message context fields. Never +reconstruct that evidence from sender, route, room, account, thread, message, +transport, or session values. Legacy adapters can explicitly pass +`channelIngress: "unsupported"`; absence remains unknown, not an allow signal. + ### Durable ingress and replay dedupe Channels adopting durable ingress should use `createChannelIngressMonitor` diff --git a/extensions/qa-channel/src/inbound.ts b/extensions/qa-channel/src/inbound.ts index ff441283ebdd..f2267c92dd89 100644 --- a/extensions/qa-channel/src/inbound.ts +++ b/extensions/qa-channel/src/inbound.ts @@ -400,6 +400,7 @@ export async function handleQaInbound(params: { }, message: { body, bodyForAgent: inbound.text, rawBody: inbound.text, commandBody: inbound.text }, media, + channelIngress: access, access: { commands: { authorized: true }, mentions: { canDetectMention: isGroup, wasMentioned: Boolean(wasMentioned) }, diff --git a/extensions/qa-lab/src/execution-identity-storage-inspection.test.ts b/extensions/qa-lab/src/execution-identity-storage-inspection.test.ts new file mode 100644 index 000000000000..58847b11f1e9 --- /dev/null +++ b/extensions/qa-lab/src/execution-identity-storage-inspection.test.ts @@ -0,0 +1,34 @@ +import fs from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { DatabaseSync } from "node:sqlite"; +import { describe, expect, it } from "vitest"; +import { inspectQaExecutionIdentityStorage } from "./execution-identity-storage-inspection.js"; + +describe("inspectQaExecutionIdentityStorage", () => { + it("returns only context and decision counts from the isolated QA database", async () => { + const stateDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-qa-identity-counts-")); + try { + const databasePath = path.join(stateDir, "state", "openclaw.sqlite"); + await fs.mkdir(path.dirname(databasePath), { recursive: true }); + const database = new DatabaseSync(databasePath); + database.exec(` + CREATE TABLE execution_identity_contexts (context_id TEXT PRIMARY KEY); + CREATE TABLE execution_decision_facts (receipt_id TEXT PRIMARY KEY); + INSERT INTO execution_identity_contexts VALUES ('context-1'), ('context-2'); + INSERT INTO execution_decision_facts VALUES ('receipt-1'); + `); + database.close(); + + expect( + inspectQaExecutionIdentityStorage({ + gateway: { + runtimeEnv: { OPENCLAW_STATE_DIR: stateDir }, + } as never, + }), + ).toEqual({ contextCount: 2, decisionCount: 1 }); + } finally { + await fs.rm(stateDir, { force: true, recursive: true }); + } + }); +}); diff --git a/extensions/qa-lab/src/execution-identity-storage-inspection.ts b/extensions/qa-lab/src/execution-identity-storage-inspection.ts new file mode 100644 index 000000000000..70570cdab10e --- /dev/null +++ b/extensions/qa-lab/src/execution-identity-storage-inspection.ts @@ -0,0 +1,45 @@ +import path from "node:path"; +import { DatabaseSync } from "node:sqlite"; +import type { QaSuiteRuntimeEnv } from "./suite-runtime-types.js"; + +function tableCount( + database: DatabaseSync, + tableName: "execution_identity_contexts" | "execution_decision_facts", +): number { + const table = database + .prepare("SELECT name FROM sqlite_schema WHERE type = 'table' AND name = ?") + .get(tableName); + if (!table) { + return 0; + } + const countSql = + tableName === "execution_identity_contexts" + ? "SELECT COUNT(*) AS count FROM execution_identity_contexts" + : "SELECT COUNT(*) AS count FROM execution_decision_facts"; + const row = database.prepare(countSql).get() as { + count: number; + }; + return row.count; +} + +/** Return only bounded row counts for deterministic no-synthetic-run proof. */ +export function inspectQaExecutionIdentityStorage(env: Pick): { + contextCount: number; + decisionCount: number; +} { + const stateDir = env.gateway.runtimeEnv.OPENCLAW_STATE_DIR?.trim(); + if (!stateDir) { + throw new Error("QA Gateway did not expose its isolated state directory"); + } + const database = new DatabaseSync(path.join(stateDir, "state", "openclaw.sqlite"), { + readOnly: true, + }); + try { + return { + contextCount: tableCount(database, "execution_identity_contexts"), + decisionCount: tableCount(database, "execution_decision_facts"), + }; + } finally { + database.close(); + } +} diff --git a/extensions/qa-lab/src/scenario-catalog-channels.test.ts b/extensions/qa-lab/src/scenario-catalog-channels.test.ts index b81fc546418e..2bff90c0cf5b 100644 --- a/extensions/qa-lab/src/scenario-catalog-channels.test.ts +++ b/extensions/qa-lab/src/scenario-catalog-channels.test.ts @@ -69,6 +69,23 @@ describe("qa scenario catalog channel contracts", () => { expect(readQaScenarioById("memory-tools-channel-context").execution.channel).toBe("qa-channel"); }); + it("keeps channel participant identity proof on isolated QA Channel lifecycle owners", () => { + const scenario = requireFlowScenario( + readQaScenarioById("channel-participant-identity-inspection"), + ); + const flow = JSON.stringify(scenario.execution.flow); + + expect(scenario.execution.channel).toBe("qa-channel"); + expect(scenario.execution.suiteIsolation).toBe("isolated"); + expect(scenario.gatewayConfigPatch).toMatchObject({ + logging: { audit: { executionIdentity: true } }, + messages: { queue: { mode: "collect", debounceMsByChannel: { "qa-channel": 1000 } } }, + channels: { "qa-channel": { groupPolicy: "allowlist" } }, + }); + expect(flow).toContain("inspectQaExecutionIdentityStorage"); + expect(flow).toContain("env.gateway.restartAfterStateMutation"); + }); + it("keeps stored inbound audio proof on the real QA Channel and Gateway flow", () => { const scenario = requireFlowScenario( readQaScenarioById("inbound-media-store-audio-transcription"), diff --git a/extensions/qa-lab/src/scenario-runtime-api.ts b/extensions/qa-lab/src/scenario-runtime-api.ts index efab719b2ae4..8ec8891d39dd 100644 --- a/extensions/qa-lab/src/scenario-runtime-api.ts +++ b/extensions/qa-lab/src/scenario-runtime-api.ts @@ -67,6 +67,7 @@ export type QaScenarioRuntimeDeps = { assertNoGatewayLogSentinels: QaScenarioRuntimeFunction; readSessionTranscriptSummary: QaScenarioRuntimeFunction; runQaCli: QaScenarioRuntimeFunction; + inspectQaExecutionIdentityStorage: QaScenarioRuntimeFunction; extractMediaPathFromText: QaScenarioRuntimeFunction; resolveGeneratedImagePath: QaScenarioRuntimeFunction; startAgentRun: QaScenarioRuntimeFunction; diff --git a/extensions/qa-lab/src/suite-runtime-agent.ts b/extensions/qa-lab/src/suite-runtime-agent.ts index b769b3275ba0..1ed898ddb0ba 100644 --- a/extensions/qa-lab/src/suite-runtime-agent.ts +++ b/extensions/qa-lab/src/suite-runtime-agent.ts @@ -18,6 +18,7 @@ export { waitForAgentRun, } from "./suite-runtime-agent-process.js"; export { runQaCli } from "./qa-cli-process.js"; +export { inspectQaExecutionIdentityStorage } from "./execution-identity-storage-inspection.js"; export { ensureImageGenerationConfigured, extractMediaPathFromText, diff --git a/extensions/telegram/src/bot-message-context.session.ts b/extensions/telegram/src/bot-message-context.session.ts index f75e9de511bf..209b53e4bb5a 100644 --- a/extensions/telegram/src/bot-message-context.session.ts +++ b/extensions/telegram/src/bot-message-context.session.ts @@ -10,6 +10,7 @@ import { type NormalizedLocation, type InboundEventKind, } from "openclaw/plugin-sdk/channel-inbound"; +import { createChannelParticipantAdmissionEvidence } from "openclaw/plugin-sdk/channel-ingress-runtime"; import { normalizeCommandBody } from "openclaw/plugin-sdk/command-surface"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import type { @@ -605,6 +606,11 @@ export async function buildTelegramInboundContextPayload(params: { : undefined; const ctxPayload = await sessionRuntime.buildChannelInboundEventContext({ channel: "telegram", + channelParticipantEvidence: createChannelParticipantAdmissionEvidence({ + channelId: "telegram", + accountId: route.accountId, + participantId: senderId, + }), resolveSupplementalMedia: true, accountId: route.accountId, messageId: options?.messageIdOverride ?? String(msg.message_id), diff --git a/qa/scenarios/channels/channel-participant-identity-inspection.yaml b/qa/scenarios/channels/channel-participant-identity-inspection.yaml new file mode 100644 index 000000000000..8e8d1ec16f9b --- /dev/null +++ b/qa/scenarios/channels/channel-participant-identity-inspection.yaml @@ -0,0 +1,320 @@ +title: Admitted channel participant identity inspection + +scenario: + id: channel-participant-identity-inspection + surface: qa-channel + coverage: + primary: + - gateway.identity-and-presence-apis + secondary: + - channels.qa-channel-final-reply + objective: Verify admitted QA Channel DM, group, senderless, and collect-configured turns retain only authoritative participant identity across restart while a rejected ingress creates no audit state. + gatewayConfigPatch: + logging: + audit: + executionIdentity: true + messages: + queue: + mode: collect + debounceMsByChannel: + qa-channel: 1000 + channels: + qa-channel: + allowFrom: ["*"] + groupPolicy: allowlist + groupAllowFrom: [qa-group-person, qa-same-person, qa-mixed-a, qa-mixed-b] + tools: + alsoAllow: [exec] + agents: + entries: + qa: + tools: + alsoAllow: [exec] + successCriteria: + - Direct and group runs project a person invoker while room and route identifiers never become principals. + - Senderless input is admitted with an unknown invoker, and an allowlist rejection creates no run, identity context, or decision fact. + - Same and mixed participants remain correctly distinguished across collect-configured QA Channel ingress. + - JSON and human CLI inspection remain stable after a Gateway replacement restart. + docsRefs: + - docs/gateway/audit.md + - docs/cli/audit.md + - docs/channels/qa-channel.md + - docs/concepts/qa-e2e-automation.md + codeRefs: + - extensions/qa-channel/src/inbound.ts + - src/channels/message-access/admission-evidence.ts + - src/auto-reply/reply/queue/drain.ts + - extensions/qa-lab/src/execution-identity-storage-inspection.ts + execution: + kind: flow + channel: qa-channel + suiteIsolation: isolated + isolationReason: Enables execution identity and collect mode, then inspects one isolated audit database across restart. + retryCount: 0 + summary: Run admitted channel turns through an ephemeral Gateway and mock provider, inspect exact identity/decision records, and prove restart stability plus rejection non-creation. + config: + requiredProviderMode: mock-openai + +flow: + steps: + - name: inspects direct, group, and senderless admitted participants + actions: + - assert: + expr: "env.providerMode === config.requiredProviderMode" + message: channel participant identity proof requires mock-openai + - call: waitForGatewayHealthy + args: [{ ref: env }, 60000] + - call: waitForQaChannelReady + args: [{ ref: env }, 60000] + - call: reset + - set: dmRunsBefore + value: + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string'))]" + - sendInbound: + conversation: { id: qa-dm-room, kind: direct } + senderId: qa-dm-person + senderName: QA DM Person + text: "Reply exactly: QA-CHANNEL-IDENTITY-DM-OK" + - waitForOutbound: + conversation: { id: qa-dm-room, kind: direct } + textIncludes: QA-CHANNEL-IDENTITY-DM-OK + timeoutMs: 60000 + - call: waitForCondition + saveAs: dmRunId + args: + - lambda: + async: true + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string' && !dmRunsBefore.includes(id)))][0]" + - 60000 + - 250 + - call: waitForCondition + saveAs: dmInspect + args: + - lambda: + async: true + expr: "runQaCli(env, ['audit', '--run', dmRunId, '--explain', '--json'], { timeoutMs: 60000, json: true }).then((value) => value.identity?.state === 'present' ? value : undefined).catch(() => undefined)" + - 60000 + - 250 + - set: dmText + value: + expr: "await runQaCli(env, ['audit', '--run', dmRunId, '--explain'], { timeoutMs: 60000 })" + - assert: + expr: "dmInspect.identity.context.invoker.state === 'present' && dmInspect.identity.context.invoker.principal?.kind === 'person' && dmInspect.identity.context.ingress.kind === 'channel' && !JSON.stringify(dmInspect).includes('qa-dm-person') && !JSON.stringify(dmInspect).includes('qa-dm-room')" + message: + expr: "`direct inspection must retain a pseudonymized person without room or raw sender material: ${JSON.stringify(dmInspect)}`" + - assert: + expr: "dmText.includes('Identity') && dmText.includes('Invoker [present]') && dmText.includes('Decisions')" + message: human audit output omitted direct participant evidence + - set: groupRunsBefore + value: + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string'))]" + - sendInbound: + conversation: { id: qa-group-room, kind: group } + senderId: qa-group-person + senderName: QA Group Person + text: "Reply exactly: QA-CHANNEL-IDENTITY-GROUP-OK" + - waitForOutbound: + conversation: { id: qa-group-room, kind: group } + textIncludes: QA-CHANNEL-IDENTITY-GROUP-OK + timeoutMs: 60000 + - call: waitForCondition + saveAs: groupRunId + args: + - lambda: + async: true + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string' && !groupRunsBefore.includes(id)))][0]" + - 60000 + - 250 + - call: waitForCondition + saveAs: groupInspect + args: + - lambda: + async: true + expr: "runQaCli(env, ['audit', '--run', groupRunId, '--explain', '--json'], { timeoutMs: 60000, json: true }).then((value) => value.identity?.state === 'present' ? value : undefined).catch(() => undefined)" + - 60000 + - 250 + - assert: + expr: "groupInspect.identity.context.invoker.state === 'present' && groupInspect.decisions.some((receipt) => receipt.action.family === 'channel' && receipt.enforcement.coverageState === 'enforced') && !JSON.stringify(groupInspect).includes('qa-group-room')" + message: group allowlist admission must record enforced participant evidence without the room as principal + - set: senderlessRunsBefore + value: + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string'))]" + - sendInbound: + conversation: { id: qa-senderless-room, kind: direct } + senderId: "" + text: "Reply exactly: QA-CHANNEL-IDENTITY-SENDERLESS-OK" + - waitForOutbound: + conversation: { id: qa-senderless-room, kind: direct } + textIncludes: QA-CHANNEL-IDENTITY-SENDERLESS-OK + timeoutMs: 60000 + - call: waitForCondition + saveAs: senderlessRunId + args: + - lambda: + async: true + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string' && !senderlessRunsBefore.includes(id)))][0]" + - 60000 + - 250 + - call: waitForCondition + saveAs: senderlessInspect + args: + - lambda: + async: true + expr: "runQaCli(env, ['audit', '--run', senderlessRunId, '--explain', '--json'], { timeoutMs: 60000, json: true }).then((value) => value.identity?.state === 'present' ? value : undefined).catch(() => undefined)" + - 60000 + - 250 + - assert: + expr: "senderlessInspect.identity.context.invoker.state === 'unknown' && senderlessInspect.decisions.some((receipt) => receipt.action.family === 'channel' && receipt.enforcement.coverageState === 'unknown')" + message: senderless admission must stay unknown + detailsExpr: "`dm=${dmRunId}; group=${groupRunId}; senderless=${senderlessRunId}`" + + - name: distinguishes same and mixed participants on collect-configured ingress + actions: + - sendInbound: + conversation: { id: qa-same-collect, kind: group } + senderId: qa-same-person + text: /queue collect debounce:1s + - waitForOutbound: + conversation: { id: qa-same-collect, kind: group } + textIncludes: Queue mode set to collect. + timeoutMs: 60000 + - set: sameRunsBefore + value: + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string'))]" + - sendInbound: + conversation: { id: qa-same-collect, kind: group } + senderId: qa-same-person + text: "Tool progress QA check: call the exec tool exactly once with this exact command before answering: `sleep 3`. After that command completes, reply exactly `QA-SAME-BUSY-OK`." + - call: sleep + args: [500] + - sendInbound: + conversation: { id: qa-same-collect, kind: group } + senderId: qa-same-person + text: repeated request queued reply gateway qa check same one + - sendInbound: + conversation: { id: qa-same-collect, kind: group } + senderId: qa-same-person + text: repeated request queued reply gateway qa check same two + - call: waitForCondition + saveAs: sameRunIds + args: + - lambda: + async: true + expr: "(ids => ids.length >= 3 ? ids : undefined)([...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string' && !sameRunsBefore.includes(id)))])" + - 60000 + - 250 + - assert: + expr: "sameRunIds.length === 3" + message: + expr: "`same-participant collect produced ${sameRunIds.length} runs: ${sameRunIds.join(',')}`" + - set: sameInspections + value: + expr: "await (async () => { const values = []; for (const id of sameRunIds) values.push(await runQaCli(env, ['audit', '--run', id, '--explain', '--json'], { timeoutMs: 60000, json: true })); return values; })()" + - assert: + expr: "sameInspections.every((value) => value.identity.state === 'present' && value.identity.context.invoker.state === 'present') && new Set(sameInspections.map((value) => value.identity.context.invoker.principal?.principalRef)).size === 1" + message: same-participant collect sources must retain one person pseudonym + - sendInbound: + conversation: { id: qa-mixed-collect, kind: group } + senderId: qa-mixed-a + text: /queue collect debounce:1s + - waitForOutbound: + conversation: { id: qa-mixed-collect, kind: group } + textIncludes: Queue mode set to collect. + timeoutMs: 60000 + - set: mixedRunsBefore + value: + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string'))]" + - sendInbound: + conversation: { id: qa-mixed-collect, kind: group } + senderId: qa-mixed-a + text: "Tool progress QA check: call the exec tool exactly once with this exact command before answering: `sleep 3`. After that command completes, reply exactly `QA-MIXED-BUSY-OK`." + - call: sleep + args: [500] + - sendInbound: + conversation: { id: qa-mixed-collect, kind: group } + senderId: qa-mixed-a + text: repeated request queued reply gateway qa check mixed one + - sendInbound: + conversation: { id: qa-mixed-collect, kind: group } + senderId: qa-mixed-b + text: repeated request queued reply gateway qa check mixed two + - call: waitForCondition + saveAs: mixedRunIds + args: + - lambda: + async: true + expr: "(ids => ids.length >= 3 ? ids : undefined)([...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string' && !mixedRunsBefore.includes(id)))])" + - 60000 + - 250 + - assert: + expr: "mixedRunIds.length === 3" + message: + expr: "`mixed-participant collect produced ${mixedRunIds.length} runs: ${mixedRunIds.join(',')}`" + - set: mixedInspections + value: + expr: "await (async () => { const values = []; for (const id of mixedRunIds) values.push(await runQaCli(env, ['audit', '--run', id, '--explain', '--json'], { timeoutMs: 60000, json: true })); return values; })()" + - assert: + expr: "mixedInspections.every((value) => value.identity.context?.invoker.state === 'present') && new Set(mixedInspections.map((value) => value.identity.context?.invoker.principal?.principalRef)).size === 2" + message: + expr: "`collect-configured mixed participants must preserve two distinct admitted people before aggregation: ${JSON.stringify(mixedInspections.map((value) => value.identity.context?.invoker))}`" + detailsExpr: "`same=${sameRunIds.join(',')}; mixed=${mixedRunIds.join(',')}`" + + - name: rejects pre-run ingress without synthetic audit state + actions: + - set: denialRunsBefore + value: + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string'))]" + - set: denialStorageBefore + value: + expr: "inspectQaExecutionIdentityStorage(env)" + - set: denialOutboundStart + value: + expr: "state.getSnapshot().messages.filter((message) => message.direction === 'outbound').length" + - sendInbound: + conversation: { id: qa-denied-group, kind: group } + senderId: qa-denied-person + text: "Reply exactly: QA-DENIED-MUST-NOT-RUN" + - waitForNoOutbound: + sinceIndex: { ref: denialOutboundStart } + quietMs: 1500 + - call: sleep + args: [1000] + - set: denialRunsAfter + value: + expr: "[...new Set((await runQaCli(env, ['audit', '--kind', 'agent_run', '--limit', '500', '--json'], { timeoutMs: 60000, json: true })).events.map((event) => event.runId).filter((id) => typeof id === 'string'))]" + - set: denialStorageAfter + value: + expr: "inspectQaExecutionIdentityStorage(env)" + - assert: + expr: "denialRunsAfter.every((id) => denialRunsBefore.includes(id)) && denialStorageAfter.contextCount === denialStorageBefore.contextCount && denialStorageAfter.decisionCount === denialStorageBefore.decisionCount" + message: pre-run rejection created a synthetic run, identity context, or decision fact + detailsExpr: "`contexts=${denialStorageAfter.contextCount}; decisions=${denialStorageAfter.decisionCount}`" + + - name: preserves exact CLI identity across Gateway restart + actions: + - set: dmContextBeforeRestart + value: + expr: "JSON.stringify(dmInspect.identity.context)" + - assert: + expr: "typeof env.gateway.restartAfterStateMutation === 'function'" + message: QA Gateway does not expose lifecycle-owned restart + - call: env.gateway.restartAfterStateMutation + args: + - lambda: + async: true + params: [ctx] + expr: "Promise.resolve()" + - call: waitForGatewayHealthy + args: [{ ref: env }, 60000] + - call: waitForQaChannelReady + args: [{ ref: env }, 60000] + - set: dmAfterRestart + value: + expr: "await runQaCli(env, ['audit', '--run', dmRunId, '--explain', '--json'], { timeoutMs: 60000, json: true })" + - set: dmTextAfterRestart + value: + expr: "await runQaCli(env, ['audit', '--run', dmRunId, '--explain'], { timeoutMs: 60000 })" + - assert: + expr: "JSON.stringify(dmAfterRestart.identity.context) === dmContextBeforeRestart && dmTextAfterRestart.includes('Invoker [present]') && dmTextAfterRestart.includes('Decisions')" + message: exact participant identity or human CLI projection changed across restart + detailsExpr: "`restart-stable run=${dmRunId}`" diff --git a/scripts/e2e/telegram-user-crabbox-proof.ts b/scripts/e2e/telegram-user-crabbox-proof.ts index 4d2bde682917..8e10c7995315 100644 --- a/scripts/e2e/telegram-user-crabbox-proof.ts +++ b/scripts/e2e/telegram-user-crabbox-proof.ts @@ -55,8 +55,10 @@ type Options = { crabboxClass: string; command: | "finish" + | "inspect" | "probe" | "publish" + | "restart" | "run" | "screenshot" | "send" @@ -64,6 +66,7 @@ type Options = { | "status" | "view"; crabboxBin: string; + chat?: string; desktopChatTitle: string; dryRun: boolean; envFile?: string; @@ -155,6 +158,7 @@ type SessionFile = { }; localRoot: string; localSut: { + configPath?: string; containerName?: string; sutAttestation?: { lane: "baseline" | "candidate"; sha: string }; gatewayLog: string; @@ -214,6 +218,8 @@ function usageText() { " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts [probe] [--text /status] [--expect OpenClaw]", " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts start [--tdlib-url ]", " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts send --session --text ", + " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts inspect --session ", + " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts restart --session ", " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts run --session -- ", " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts view --session --message-id ", " node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts screenshot --session ", @@ -223,6 +229,7 @@ function usageText() { "", "Useful options:", " --class Crabbox machine class. Default: standard.", + " --chat Telegram chat override for send (for example @bot for DM).", " --desktop-chat-title Telegram Desktop chat to select before recording.", " --human-delay-fixed-ms Set a fixed custom human delay before Gateway startup.", " --id Reuse an existing Crabbox desktop lease.", @@ -316,8 +323,10 @@ export function parseArgs(argvInput: string[]): Options { argv = argv[0] === "--" ? argv.slice(1) : argv; const commands = new Set([ "finish", + "inspect", "probe", "publish", + "restart", "run", "screenshot", "send", @@ -388,6 +397,8 @@ export function parseArgs(argvInput: string[]): Options { }; if (arg === "--class") { opts.crabboxClass = readValue(); + } else if (arg === "--chat") { + opts.chat = readValue(); } else if (arg === "--crabbox-bin") { opts.crabboxBin = readValue(); } else if (arg === "--desktop-chat-title") { @@ -496,7 +507,17 @@ export function parseArgs(argvInput: string[]): Options { throw new Error("run requires a remote command after --."); } if ( - ["finish", "publish", "run", "screenshot", "send", "status", "view"].includes(command) && + [ + "finish", + "inspect", + "publish", + "restart", + "run", + "screenshot", + "send", + "status", + "view", + ].includes(command) && !opts.sessionFile ) { throw new Error(`${command} requires --session.`); @@ -510,6 +531,9 @@ export function parseArgs(argvInput: string[]): Options { if (command !== "start" && opts.humanDelayFixedMs !== undefined) { throw new Error("--human-delay-fixed-ms is available only for start sessions."); } + if (command !== "send" && opts.chat) { + throw new Error("--chat is available only for held-session sends."); + } if (opts.mcpAppFixture && command !== "start") { throw new Error("--mcp-app-fixture is available only for start sessions."); } @@ -685,6 +709,42 @@ export function createOpenClawGatewaySpawnSpec(params: { }; } +export function createOpenClawCliSpawnSpec(params: { + args: string[]; + env: NodeJS.ProcessEnv; + repoRoot: string; + nodeExecPath?: string; + npmExecPath?: string; + pnpmExecPath?: string; + platform?: NodeJS.Platform; +}): GatewaySpawnSpec { + if (params.pnpmExecPath) { + return { + args: ["openclaw", ...params.args], + command: params.pnpmExecPath, + options: { cwd: params.repoRoot, env: params.env, shell: false }, + }; + } + const spec = createPnpmRunnerSpawnSpec({ + cwd: params.repoRoot, + env: params.env, + nodeExecPath: params.nodeExecPath, + npmExecPath: params.npmExecPath, + platform: params.platform, + pnpmArgs: ["openclaw", ...params.args], + }); + return { + args: spec.args, + command: spec.command, + options: { + cwd: spec.options.cwd, + env: spec.options.env, + shell: spec.options.shell, + windowsVerbatimArguments: spec.options.windowsVerbatimArguments, + }, + }; +} + function shellQuote(value: string) { return `'${value.replaceAll("'", "'\\''")}'`; } @@ -860,10 +920,12 @@ export function runCommand(params: { cwd: string; env?: NodeJS.ProcessEnv; outputFile?: string; + shell?: boolean | string; stdio?: "inherit" | "pipe"; stdin?: string; timeoutKillGraceMs?: number; timeoutMs?: number; + windowsVerbatimArguments?: boolean; }) { return new Promise((resolve, reject) => { if (params.outputFile) { @@ -873,7 +935,9 @@ export function runCommand(params: { cwd: params.cwd, detached: process.platform !== "win32", env: params.env ?? process.env, + shell: params.shell, stdio: ["pipe", "pipe", "pipe"], + windowsVerbatimArguments: params.windowsVerbatimArguments, }); activeCommandChildren.add(child); installCommandCleanupHandlers(); @@ -1195,6 +1259,47 @@ export async function waitForLog( ); } +export function readLogAfterOffset( + logPath: string, + offset: number, + maxBytes = LOG_READY_TAIL_BYTES, +) { + const size = fs.statSync(logPath).size; + if (size <= offset) { + return ""; + } + const start = Math.max(offset, size - Math.max(1, maxBytes)); + const buffer = Buffer.alloc(size - start); + const fd = fs.openSync(logPath, "r"); + try { + fs.readSync(fd, buffer, 0, buffer.length, start); + } finally { + fs.closeSync(fd); + } + return buffer.toString("utf8"); +} + +export async function waitForLogAfterOffset(params: { + label: string; + logPath: string; + offset: number; + pattern: RegExp; + timeoutMs: number; +}) { + const started = Date.now(); + while (Date.now() - started < params.timeoutMs) { + const text = readLogAfterOffset(params.logPath, params.offset); + if (params.pattern.test(text)) { + return text; + } + await sleep(250); + } + const text = readLogAfterOffset(params.logPath, params.offset); + throw new Error( + `${params.label} was not observed within ${params.timeoutMs}ms\n${sliceUtf16Safe(text, -4000)}`, + ); +} + async function telegram(token: string, method: string, body: JsonObject = {}) { return await telegramBotApi(token, method, body); } @@ -1289,7 +1394,7 @@ export function writeSutConfig(params: { }, // Exercise the opt-in message audit surface: the DM probe should produce // inbound/outbound rows under the privacy-sensitive "direct" mode. - logging: { audit: { enabled: true, messages: "direct" } }, + logging: { audit: { enabled: true, executionIdentity: true, messages: "direct" } }, channels: { telegram: { allowFrom: [params.testerId], @@ -2453,6 +2558,7 @@ sleep 1 } export function renderRemoteProbe(params: { + chat?: string; expect: string[]; outputPath?: string; sutUsername: string; @@ -2469,6 +2575,9 @@ export function renderRemoteProbe(params: { params.outputPath ?? `${REMOTE_ROOT}/probe.json`, "--json", ]; + if (params.chat) { + args.push("--chat", params.chat); + } for (const expected of params.expect) { args.push("--expect", expected); } @@ -2960,6 +3069,8 @@ async function startSession(root: string, opts: Options, outputDir: string) { webvnc: `${opts.crabboxBin} webvnc --provider ${opts.provider} --target ${opts.target} --id ${leaseId} --open`, commands: { send: `openclaw-telegram-user-crabbox-proof send --session ${path.relative(root, pathname)} --text '/status'`, + inspect: `openclaw-telegram-user-crabbox-proof inspect --session ${path.relative(root, pathname)}`, + restart: `openclaw-telegram-user-crabbox-proof restart --session ${path.relative(root, pathname)}`, view: `openclaw-telegram-user-crabbox-proof view --session ${path.relative(root, pathname)} --message-id `, run: `openclaw-telegram-user-crabbox-proof run --session ${path.relative(root, pathname)} -- bash -lc 'source ${REMOTE_ROOT}/env.sh && python3 ${REMOTE_ROOT}/user-driver.py transcript --limit 20 --json'`, finish: `openclaw-telegram-user-crabbox-proof finish --session ${path.relative(root, pathname)} --preview-crop telegram-window`, @@ -3019,6 +3130,7 @@ async function sendSessionProbe(root: string, opts: Options, outputDir: string) await writeExecutable( probeScript, renderRemoteProbe({ + chat: opts.chat?.replaceAll("{sut}", session.credential.sutUsername), expect: opts.expect, outputPath: remoteProbe, sutUsername: session.credential.sutUsername, @@ -3089,6 +3201,253 @@ async function statusSession(root: string, opts: Options, outputDir: string) { }; } +function sessionSutConfigPath(session: SessionFile) { + return session.localSut.configPath ?? path.join(session.localSut.tempRoot, "openclaw.json"); +} + +async function runSessionAuditCli( + root: string, + opts: Options, + session: SessionFile, + args: string[], +) { + const spec = createOpenClawCliSpawnSpec({ + args, + env: { + ...childProcessBaseEnv(), + OPENCLAW_CONFIG_PATH: sessionSutConfigPath(session), + OPENCLAW_STATE_DIR: session.localSut.stateDir, + }, + repoRoot: root, + nodeExecPath: opts.nodeBin, + pnpmExecPath: opts.pnpmBin, + }); + const cwd = spec.options.cwd; + return await runCommand({ + command: spec.command, + args: spec.args, + cwd: typeof cwd === "string" ? cwd : cwd ? fileURLToPath(cwd) : root, + env: spec.options.env, + shell: spec.options.shell, + timeoutMs: opts.timeoutMs, + windowsVerbatimArguments: spec.options.windowsVerbatimArguments, + }); +} + +function parseCommandJson(result: CommandResult, label: string): JsonObject { + try { + const parsed = JSON.parse(result.stdout) as unknown; + if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) { + throw new Error("expected a JSON object"); + } + return parsed as JsonObject; + } catch (error) { + throw new Error(`${label} returned invalid JSON: ${coerceErrorMessage(error)}`, { + cause: error, + }); + } +} + +function inspectIdentityContext(result: JsonObject): JsonObject | undefined { + const identity = result.identity; + if (!identity || typeof identity !== "object" || Array.isArray(identity)) { + return undefined; + } + const record = identity as JsonObject; + return record.state === "present" && record.context && typeof record.context === "object" + ? (record.context as JsonObject) + : undefined; +} + +async function inspectSessionIdentity(root: string, opts: Options, outputDir: string) { + const { session } = readSession(root, opts, outputDir); + const listed = parseCommandJson( + await runSessionAuditCli(root, opts, session, [ + "audit", + "--kind", + "agent_run", + "--limit", + "500", + "--json", + ]), + "audit activity list", + ); + const events = Array.isArray(listed.events) ? listed.events : []; + const runIds = [ + ...new Set( + events.flatMap((event) => { + if (!event || typeof event !== "object" || Array.isArray(event)) { + return []; + } + const runId = (event as JsonObject).runId; + return typeof runId === "string" && runId.trim() ? [runId] : []; + }), + ), + ]; + const inspections: Array<{ human: string; json: JsonObject; runId: string }> = []; + for (const runId of runIds) { + const json = parseCommandJson( + await runSessionAuditCli(root, opts, session, [ + "audit", + "--run", + runId, + "--explain", + "--json", + ]), + `audit inspection ${runId}`, + ); + const context = inspectIdentityContext(json); + if (!context) { + continue; + } + const ingress = context.ingress; + if ( + !ingress || + typeof ingress !== "object" || + Array.isArray(ingress) || + (ingress as JsonObject).kind !== "channel" + ) { + continue; + } + const human = ( + await runSessionAuditCli(root, opts, session, ["audit", "--run", runId, "--explain"]) + ).stdout; + inspections.push({ human, json, runId }); + } + if (inspections.length < 2) { + throw new Error( + `Telegram DM/group proof requires at least two admitted channel runs; found ${inspections.length}.`, + ); + } + const contextsByRun = Object.fromEntries( + inspections.map(({ json, runId }) => [runId, inspectIdentityContext(json)]), + ); + const serialized = JSON.stringify({ contextsByRun, inspections }); + for (const raw of [ + session.credential.groupId, + session.credential.testerUserId, + session.credential.testerUsername, + ]) { + if (raw && serialized.includes(raw)) { + throw new Error("Telegram audit inspection retained a raw participant or room identifier."); + } + } + const principalRefs = new Set(); + for (const inspection of inspections) { + const context = inspectIdentityContext(inspection.json); + const invoker = context?.invoker; + const principal = + invoker && typeof invoker === "object" && !Array.isArray(invoker) + ? (invoker as JsonObject).principal + : undefined; + const principalRef = + principal && typeof principal === "object" && !Array.isArray(principal) + ? (principal as JsonObject).principalRef + : undefined; + const decisions = Array.isArray(inspection.json.decisions) ? inspection.json.decisions : []; + const hasChannelDecision = decisions.some((decision) => { + if (!decision || typeof decision !== "object" || Array.isArray(decision)) { + return false; + } + const action = (decision as JsonObject).action; + return ( + action && + typeof action === "object" && + !Array.isArray(action) && + (action as JsonObject).family === "channel" && + (action as JsonObject).operation === "admission" + ); + }); + if ( + !invoker || + typeof invoker !== "object" || + Array.isArray(invoker) || + (invoker as JsonObject).state !== "present" || + !principal || + typeof principal !== "object" || + Array.isArray(principal) || + (principal as JsonObject).kind !== "person" || + typeof principalRef !== "string" || + !hasChannelDecision || + !inspection.human.includes("Invoker [present]") || + !inspection.human.includes("Decisions") + ) { + throw new Error(`Telegram run ${inspection.runId} omitted participant CLI evidence.`); + } + principalRefs.add(principalRef); + } + if (principalRefs.size !== 1) { + throw new Error("Telegram DM and group runs did not retain the same participant principal."); + } + + const jsonPath = path.join(session.outputDir, "telegram-execution-identity.private.json"); + const textPath = path.join(session.outputDir, "telegram-execution-identity.private.txt"); + const previous = readJsonFile(jsonPath); + const previousContexts = + previous.contextsByRun && + typeof previous.contextsByRun === "object" && + !Array.isArray(previous.contextsByRun) + ? (previous.contextsByRun as JsonObject) + : undefined; + const stableAcrossRestart = previousContexts + ? Object.entries(previousContexts).every( + ([runId, context]) => JSON.stringify(contextsByRun[runId]) === JSON.stringify(context), + ) + : undefined; + if (stableAcrossRestart === false) { + throw new Error("Telegram execution identity context changed across Gateway restart."); + } + fs.writeFileSync( + jsonPath, + `${JSON.stringify({ contextsByRun, runIds: inspections.map((item) => item.runId) }, null, 2)}\n`, + { mode: 0o600 }, + ); + fs.chmodSync(jsonPath, 0o600); + fs.writeFileSync( + textPath, + inspections.map((item) => `# ${item.runId}\n${item.human.trim()}\n`).join("\n"), + { mode: 0o600 }, + ); + fs.chmodSync(textPath, 0o600); + return { + inspectionCount: inspections.length, + json: path.relative(root, jsonPath), + runIds: inspections.map((item) => item.runId), + stableAcrossRestart: stableAcrossRestart ?? null, + status: "pass", + text: path.relative(root, textPath), + }; +} + +async function restartSessionGateway(root: string, opts: Options, outputDir: string) { + const { session } = readSession(root, opts, outputDir); + if (session.localSut.containerName) { + throw new Error( + "Held-session restart requires the lifecycle-owned host Gateway; container sessions are unsupported.", + ); + } + const pid = session.localSut.gatewayPid; + process.kill(pid, 0); + const offset = fs.statSync(session.localSut.gatewayLog).size; + process.kill(pid, "SIGUSR1"); + await waitForLogAfterOffset({ + label: "Gateway restart boundary", + logPath: session.localSut.gatewayLog, + offset, + pattern: /received SIGUSR1; restarting/u, + timeoutMs: opts.timeoutMs, + }); + await waitForLogAfterOffset({ + label: "Gateway restart readiness", + logPath: session.localSut.gatewayLog, + offset, + pattern: /gateway ready|restart trace: restart\.ready/u, + timeoutMs: opts.timeoutMs, + }); + process.kill(pid, 0); + return { gatewayPid: pid, logOffset: offset, status: "pass" }; +} + function telegramPrivatePostLink(groupId: string, messageId: string) { if (!/^-100\d+$/u.test(groupId)) { throw new Error(`Telegram privatepost links require a -100 group id, got ${groupId}.`); @@ -3386,6 +3745,14 @@ async function main() { console.log(JSON.stringify(await sendSessionProbe(root, opts, outputDir), null, 2)); return; } + if (opts.command === "inspect") { + console.log(JSON.stringify(await inspectSessionIdentity(root, opts, outputDir), null, 2)); + return; + } + if (opts.command === "restart") { + console.log(JSON.stringify(await restartSessionGateway(root, opts, outputDir), null, 2)); + return; + } if (opts.command === "run") { console.log(JSON.stringify(await runSessionCommand(root, opts, outputDir), null, 2)); return; diff --git a/scripts/plugin-sdk-surface-report.mts b/scripts/plugin-sdk-surface-report.mts index 4a9dd26ad799..8ed64474b4c0 100644 --- a/scripts/plugin-sdk-surface-report.mts +++ b/scripts/plugin-sdk-surface-report.mts @@ -275,7 +275,8 @@ export function readPluginSdkSurfaceBudgets(env: NodeJS.ProcessEnv = process.env // +2: narrow channel agent-run terminal reader and outcome contract. // +5: narrow string, record, and error coercion helpers. // +1: normalized Gateway public origin resolver for plugin-generated links. - 4879, + // +1: host-minted channel participant admission evidence for plugin ingress. + 4880, env, ), publicFunctionExports: readPluginSdkSurfaceBudgetEnv( @@ -342,7 +343,8 @@ export function readPluginSdkSurfaceBudgets(env: NodeJS.ProcessEnv = process.env // +1: narrow channel agent-run terminal reader. // +5: narrow string, record, and error coercion helpers. // +1: normalized Gateway public origin resolver for plugin-generated links. - 2932, + // +1: host-minted channel participant admission evidence for plugin ingress. + 2933, env, ), publicDeprecatedExports: readPluginSdkSurfaceBudgetEnv( diff --git a/src/audit/audit-recorder.ts b/src/audit/audit-recorder.ts index 30a50aa3aeb7..ce24e1cf38d6 100644 --- a/src/audit/audit-recorder.ts +++ b/src/audit/audit-recorder.ts @@ -16,6 +16,7 @@ let persistenceFailureWarned = false; type AuditEventRecorder = AgentEventAuditRecorder & { recordMessage: (event: TrustedMessageAuditEvent) => void; recordExecutionIdentity: (work: ExecutionIdentityAdmissionWork) => boolean; + recordExecutionDecision: AuditEventWriter["recordExecutionDecision"]; }; export function createAuditEventRecorder(options: { @@ -46,6 +47,7 @@ export function createAuditEventRecorder(options: { return { ...agentRecorder, recordExecutionIdentity: writer.recordExecutionIdentity, + recordExecutionDecision: writer.recordExecutionDecision, recordMessage: (event) => { if (options.messageMode === "off") { return; diff --git a/src/auto-reply/reply/agent-runner-execution-lifecycle.test.ts b/src/auto-reply/reply/agent-runner-execution-lifecycle.test.ts index e1311038247e..c531d0b9a1c8 100644 --- a/src/auto-reply/reply/agent-runner-execution-lifecycle.test.ts +++ b/src/auto-reply/reply/agent-runner-execution-lifecycle.test.ts @@ -2,6 +2,12 @@ import { describe, expect, it, vi } from "vitest"; import type { SessionMcpRuntime } from "../../agents/agent-bundle-mcp-types.js"; import { updateMcpAppModelContext } from "../../agents/mcp-app-model-context.js"; import { createAgentRunRestartAbortError } from "../../agents/run-termination.js"; +import { configureExecutionIdentityAdmissionSink } from "../../audit/execution-identity-admission.js"; +import { + configureChannelAdmissionDecisionSink, + configureChannelAdmissionEvidenceCollection, + createChannelParticipantAdmissionEvidence, +} from "../../channels/message-access/admission-evidence.js"; import { getDiagnosticSessionActivitySnapshot } from "../../logging/diagnostic-run-activity.js"; import { SILENT_REPLY_TOKEN } from "../tokens.js"; import type { GetReplyOptions } from "../types.js"; @@ -24,6 +30,71 @@ import { createReplyOperation, type ReplyOperation } from "./reply-run-registry. const state = setupAgentRunnerExecutionTestState(); describe("executeAgentTurn: run lifecycle and ownership", () => { + it("attributes one admitted channel participant before its admission decision", async () => { + const order: string[] = []; + const identityWork: unknown[] = []; + const decisionReceipts: unknown[] = []; + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + const clearIdentitySink = configureExecutionIdentityAdmissionSink((work) => { + order.push("identity"); + identityWork.push(work); + return true; + }); + const clearDecisionSink = configureChannelAdmissionDecisionSink((receipt) => { + order.push("decision"); + decisionReceipts.push(receipt); + return true; + }); + try { + const followupRun = createFollowupRun(); + followupRun.run.config = { logging: { audit: { executionIdentity: true } } }; + followupRun.channelAdmissionEvidence = createChannelParticipantAdmissionEvidence({ + channelId: "whatsapp", + accountId: "default", + participantId: "person-42", + }); + state.runEmbeddedAgentMock.mockImplementationOnce(async (params: EmbeddedAgentParams) => { + const admission = ( + params as EmbeddedAgentParams & { + preparedRunAdmission: { admit: (kind: "embedded") => Promise }; + } + ).preparedRunAdmission; + await admission.admit("embedded"); + return { payloads: [{ text: "ok" }], meta: {} }; + }); + + const executeAgentTurn = await getExecuteAgentTurnForTest(); + await executeAgentTurn({ + ...createMinimalRunAgentTurnParams({ followupRun }), + }); + + expect(order).toEqual(["identity", "decision"]); + expect(identityWork).toMatchObject([ + { + kind: "capture", + envelope: { + ingress: { kind: "channel", state: "present" }, + invoker: { + state: "present", + kind: "person", + rawPrincipalRef: '["whatsapp","default","person-42"]', + }, + }, + }, + ]); + expect(decisionReceipts).toMatchObject([ + { + action: { family: "channel", operation: "admission" }, + enforcement: { coverageState: "attribution-only" }, + }, + ]); + } finally { + clearDecisionSink(); + clearIdentitySink(); + clearCollection(); + } + }); + it("passes the reply abort signal to fallback orchestration and candidates", async () => { const { replyOperation } = createMockReplyOperation(); state.runEmbeddedAgentMock.mockResolvedValueOnce({ @@ -598,6 +669,62 @@ describe("executeAgentTurn: run lifecycle and ownership", () => { expect(state.runWithModelFallbackMock).not.toHaveBeenCalled(); }); + it("does not consume channel evidence until a retry reaches runtime admission", async () => { + const captured: unknown[] = []; + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + const clearSink = configureExecutionIdentityAdmissionSink((work) => { + captured.push(work); + return true; + }); + try { + const followupRun = createFollowupRun(); + followupRun.run.config = { logging: { audit: { executionIdentity: true } } }; + followupRun.channelAdmissionEvidence = createChannelParticipantAdmissionEvidence({ + channelId: "whatsapp", + participantId: "person-1", + }); + state.resolveCurrentTurnImagesMock.mockRejectedValueOnce(new Error("invalid image metadata")); + + const executeAgentTurn = await getExecuteAgentTurnForTest(); + await expect( + executeAgentTurn( + createMinimalRunAgentTurnParams({ + followupRun, + opts: { runId: "preflight-failure" }, + }), + ), + ).rejects.toThrow("invalid image metadata"); + expect(captured).toEqual([]); + + state.runEmbeddedAgentMock.mockImplementationOnce(async (params: EmbeddedAgentParams) => { + const admission = ( + params as EmbeddedAgentParams & { + preparedRunAdmission: { admit: (kind: "embedded") => Promise }; + } + ).preparedRunAdmission; + await admission.admit("embedded"); + return { payloads: [{ text: "ok" }], meta: {} }; + }); + await executeAgentTurn( + createMinimalRunAgentTurnParams({ + followupRun, + opts: { runId: "preflight-success" }, + }), + ); + + expect(captured).toHaveLength(1); + expect(captured).toMatchObject([ + { + kind: "capture", + envelope: { ingress: { state: "present" }, invoker: { state: "present" } }, + }, + ]); + } finally { + clearSink(); + clearCollection(); + } + }); + it("passes runtime toolsAllow to embedded agent runs", async () => { state.runEmbeddedAgentMock.mockResolvedValueOnce({ payloads: [{ text: "ok" }], diff --git a/src/auto-reply/reply/agent-runner-execution.ts b/src/auto-reply/reply/agent-runner-execution.ts index dd45e88391b5..216e22698e8b 100644 --- a/src/auto-reply/reply/agent-runner-execution.ts +++ b/src/auto-reply/reply/agent-runner-execution.ts @@ -6,11 +6,7 @@ import { } from "@openclaw/normalization-core/string-coerce"; import { hasOutboundReplyContent } from "openclaw/plugin-sdk/reply-payload"; import type { ChatRunStartupPhase } from "../../../packages/gateway-protocol/src/index.js"; -import { - createOperationalRunInstanceRef, - prepareAgentRunAdmission, - type PreparedAgentRunAdmission, -} from "../../agents/admitted-run-context.js"; +import type { PreparedAgentRunAdmission } from "../../agents/admitted-run-context.js"; import { peekSessionMcpRuntime } from "../../agents/agent-bundle-mcp-manager-api.js"; import { resolveBootstrapWarningSignaturesSeen } from "../../agents/bootstrap-budget.js"; import { @@ -66,6 +62,7 @@ import { import { createAgentTurnPresentation } from "./agent-runner-presentation.js"; import { createAgentTurnTimingTracker } from "./agent-runner-turn-timing.js"; import { resolveQueuedReplyRuntimeConfig } from "./agent-runner-utils.js"; +import { prepareChannelRunAdmission } from "./channel-run-admission.js"; import { shouldNotifyUserAboutCompaction } from "./compaction-notice.js"; import { resolveCurrentTurnImages } from "./current-turn-images.js"; import type { FollowupRun } from "./queue.js"; @@ -492,18 +489,13 @@ async function executeAgentTurnInternal( completed: false, }; const runId = params.opts?.runId ?? crypto.randomUUID(); - const preparedRunAdmission = prepareAgentRunAdmission({ + const preparedRunAdmission = prepareChannelRunAdmission({ cfg: resolveQueuedReplyRuntimeConfig(params.followupRun.run.config), - operationalRunInstance: createOperationalRunInstanceRef(runId), - facts: { - runId, - agentId: params.followupRun.run.agentId, - ingress: { - kind: "channel", - boundary: "auto-reply.agent-runner", - state: "present", - }, - }, + runId, + agentId: params.followupRun.run.agentId, + ingressKind: "channel", + boundary: "auto-reply.agent-runner", + evidence: params.followupRun.channelAdmissionEvidence, }); try { return await executeAgentTurnInternalWithRetryState( diff --git a/src/auto-reply/reply/channel-run-admission.test.ts b/src/auto-reply/reply/channel-run-admission.test.ts new file mode 100644 index 000000000000..1412ad99e982 --- /dev/null +++ b/src/auto-reply/reply/channel-run-admission.test.ts @@ -0,0 +1,115 @@ +import { describe, expect, it } from "vitest"; +import { + createOperationalRunInstanceRef, + prepareAgentRunAdmission, +} from "../../agents/admitted-run-context.js"; +import { configureExecutionIdentityAdmissionSink } from "../../audit/execution-identity-admission.js"; +import { + configureChannelAdmissionDecisionSink, + configureChannelAdmissionEvidenceCollection, + consumeChannelAdmissionEvidence, + createChannelParticipantAdmissionEvidence, +} from "../../channels/message-access/admission-evidence.js"; +import { prepareChannelRunAdmission } from "./channel-run-admission.js"; + +const identityConfig = { logging: { audit: { executionIdentity: true } } } as const; + +describe("channel run admission", () => { + it("consumes once across fallback admission and closes the exact prepared owner", async () => { + const identityWork: unknown[] = []; + const decisions: unknown[] = []; + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + const clearIdentitySink = configureExecutionIdentityAdmissionSink((work) => { + identityWork.push(work); + return true; + }); + const clearDecisionSink = configureChannelAdmissionDecisionSink((receipt) => { + decisions.push(receipt); + return true; + }); + try { + const evidence = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); + const prepared = prepareChannelRunAdmission({ + cfg: identityConfig, + runId: "run-1", + agentId: "main", + ingressKind: "channel", + boundary: "test.channel", + evidence, + }); + + const first = await prepared.admit("embedded"); + const fallback = await prepared.admit("embedded"); + + expect(fallback).toBe(first); + expect(identityWork).toHaveLength(1); + expect(decisions).toHaveLength(1); + expect(consumeChannelAdmissionEvidence(evidence)).toMatchObject({ + ingressState: "unknown", + }); + + prepared.close(); + await expect(prepared.admit("embedded")).rejects.toThrow( + "prepared execution context is already closed", + ); + } finally { + clearDecisionSink(); + clearIdentitySink(); + clearCollection(); + } + }); + + it("does not consume a cancelled pre-admission carrier or label internal ACP as a person", async () => { + const identityWork: unknown[] = []; + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + const clearIdentitySink = configureExecutionIdentityAdmissionSink((work) => { + identityWork.push(work); + return true; + }); + try { + const evidence = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); + const cancelled = prepareChannelRunAdmission({ + cfg: identityConfig, + runId: "cancelled-run", + agentId: "main", + ingressKind: "channel", + boundary: "test.channel", + evidence, + }); + cancelled.close(); + await expect(cancelled.admit("embedded")).rejects.toThrow( + "prepared execution context is already closed", + ); + expect(consumeChannelAdmissionEvidence(evidence)).toMatchObject({ + ingressState: "present", + }); + + const internalAcp = prepareAgentRunAdmission({ + cfg: identityConfig, + operationalRunInstance: createOperationalRunInstanceRef("internal-acp"), + facts: { + runId: "internal-acp", + agentId: "main", + ingress: { kind: "acp", boundary: "test.internal", state: "present" }, + }, + }); + await internalAcp.admit("acp"); + internalAcp.close(); + + expect(identityWork).toHaveLength(1); + expect(identityWork).toMatchObject([{ kind: "capture", envelope: {} }]); + expect( + (identityWork[0] as { envelope?: { invoker?: unknown } }).envelope?.invoker, + ).toBeUndefined(); + } finally { + clearIdentitySink(); + clearCollection(); + } + }); +}); diff --git a/src/auto-reply/reply/channel-run-admission.ts b/src/auto-reply/reply/channel-run-admission.ts new file mode 100644 index 000000000000..c13784ada0df --- /dev/null +++ b/src/auto-reply/reply/channel-run-admission.ts @@ -0,0 +1,96 @@ +import { + createOperationalRunInstanceRef, + prepareAgentRunAdmission, + type AdmittedRunContext, + type PreparedAgentRunAdmission, +} from "../../agents/admitted-run-context.js"; +import type { ExecutionIdentityAdmissionFacts } from "../../audit/execution-identity-admission.js"; +import { + consumeChannelAdmissionEvidence, + recordChannelAdmissionDecision, + type ChannelAdmissionEvidence, +} from "../../channels/message-access/admission-evidence.js"; +import type { OpenClawConfig } from "../../config/types.openclaw.js"; + +/** Adapt one opaque channel carrier to the canonical admitted-run facts and decision FIFO. */ +export function consumeChannelRunAdmission(evidence: ChannelAdmissionEvidence | undefined): { + ingressState: ExecutionIdentityAdmissionFacts["ingress"]["state"]; + facts: Pick; + onAdmitted: (context: AdmittedRunContext) => void; +} { + const admission = consumeChannelAdmissionEvidence(evidence); + return Object.freeze({ + ingressState: admission.ingressState, + facts: Object.freeze({ + invoker: admission.invoker, + ...(admission.assuranceRef + ? { + assurance: [ + { + kind: "channel-admission" as const, + rawEvidenceRef: admission.assuranceRef, + strength: "boundary-verified" as const, + }, + ], + } + : {}), + }), + onAdmitted: (context) => { + const token = context.executionIdentityToken; + if (token && admission.decisionCoverage) { + recordChannelAdmissionDecision({ + contextId: token.contextId, + executionId: token.executionId, + runId: token.runId, + occurredAt: token.createdAt, + coverageState: admission.decisionCoverage, + }); + } + }, + }); +} + +/** Defer evidence consumption until the selected runtime actually admits the run. */ +export function prepareChannelRunAdmission(params: { + cfg: OpenClawConfig; + runId: string; + agentId: string; + ingressKind: ExecutionIdentityAdmissionFacts["ingress"]["kind"]; + boundary: string; + evidence?: ChannelAdmissionEvidence; +}): PreparedAgentRunAdmission { + const operationalRunInstance = createOperationalRunInstanceRef(params.runId); + let prepared: PreparedAgentRunAdmission | undefined; + let closed = false; + return Object.freeze({ + operationalRunInstance, + admit: (runtimeKind, runtimeInstanceId) => { + if (closed) { + return Promise.reject(new Error("prepared execution context is already closed")); + } + if (!prepared) { + const channelAdmission = consumeChannelRunAdmission(params.evidence); + prepared = prepareAgentRunAdmission({ + cfg: params.cfg, + operationalRunInstance, + facts: { + runId: params.runId, + agentId: params.agentId, + ingress: { + kind: params.ingressKind, + boundary: params.boundary, + state: channelAdmission.ingressState, + }, + ...channelAdmission.facts, + }, + onAdmitted: channelAdmission.onAdmitted, + }); + } + return prepared.admit(runtimeKind, runtimeInstanceId); + }, + close: () => { + closed = true; + prepared?.close(); + }, + }); +} diff --git a/src/auto-reply/reply/commands-acp.test.ts b/src/auto-reply/reply/commands-acp.test.ts index c4d4c3e7d14b..0cbc7bbaa967 100644 --- a/src/auto-reply/reply/commands-acp.test.ts +++ b/src/auto-reply/reply/commands-acp.test.ts @@ -4,6 +4,12 @@ import os from "node:os"; import path from "node:path"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { AcpRuntimeError } from "../../acp/runtime/errors.js"; +import { configureExecutionIdentityAdmissionSink } from "../../audit/execution-identity-admission.js"; +import { + bindChannelContextAdmissionEvidence, + configureChannelAdmissionEvidenceCollection, + createChannelParticipantAdmissionEvidence, +} from "../../channels/message-access/admission-evidence.js"; import type { OpenClawConfig } from "../../config/config.js"; import type { SessionBindingRecord } from "../../infra/outbound/session-binding-service.js"; import { setActivePluginRegistry } from "../../plugins/runtime.js"; @@ -1658,6 +1664,62 @@ describe("/acp command", () => { expect(result?.reply?.text).toContain("Applied steering."); }); + it("admits ACP steer with the original channel participant", async () => { + const captured: unknown[] = []; + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + const clearSink = configureExecutionIdentityAdmissionSink((work) => { + captured.push(work); + return true; + }); + try { + hoisted.callGatewayMock.mockImplementation(async (request: { method?: string }) => { + if (request.method === "sessions.resolve") { + return { key: defaultAcpSessionKey }; + } + return { ok: true }; + }); + hoisted.readAcpSessionEntryMock.mockReturnValue(createAcpSessionEntry()); + hoisted.runTurnMock.mockImplementation(async function* () { + yield { type: "done" }; + }); + const cfg = { + ...baseCfg, + logging: { audit: { executionIdentity: true } }, + } satisfies OpenClawConfig; + const params = createDiscordParams( + `/acp steer --session ${defaultAcpSessionKey} tighten logging`, + cfg, + ); + const evidence = createChannelParticipantAdmissionEvidence({ + channelId: "discord", + accountId: "default", + participantId: "user-1", + }); + bindChannelContextAdmissionEvidence({ + context: params.ctx, + channelId: "discord", + accountId: "default", + evidence, + rawPrincipalRef: "user-1", + }); + + await handleAcpCommand(params, true); + + expect(captured).toMatchObject([ + { + kind: "capture", + envelope: { + ingress: { kind: "acp", state: "present" }, + invoker: { state: "present", kind: "person" }, + }, + }, + ]); + } finally { + clearSink(); + clearCollection(); + } + }); + it("keeps bounded ACP steer output UTF-16 safe", async () => { const prefix = "a".repeat(799); hoisted.callGatewayMock.mockImplementation(async (request: { method?: string }) => { diff --git a/src/auto-reply/reply/commands-acp/lifecycle.ts b/src/auto-reply/reply/commands-acp/lifecycle.ts index 67394c207533..0fe47c22deeb 100644 --- a/src/auto-reply/reply/commands-acp/lifecycle.ts +++ b/src/auto-reply/reply/commands-acp/lifecycle.ts @@ -25,6 +25,10 @@ import { resolveAcpSpawnRuntimePolicyError, resolveRuntimeCwdForAcpSpawn, } from "../../../agents/subagents/spawn/acp-spawn.js"; +import { + readChannelContextAdmissionEvidence, + type ChannelAdmissionEvidence, +} from "../../../channels/message-access/admission-evidence.js"; import { updateSessionEntry } from "../../../config/sessions/session-accessor.js"; import type { SessionAcpMeta } from "../../../config/sessions/types.js"; import type { OpenClawConfig } from "../../../config/types.openclaw.js"; @@ -34,6 +38,7 @@ import { type SessionBindingRecord, } from "../../../infra/outbound/session-binding-service.js"; import { resolveAgentIdFromSessionKey } from "../../../routing/session-key.js"; +import { consumeChannelRunAdmission } from "../channel-run-admission.js"; import { commandReply } from "../command-gates.js"; import type { CommandHandlerResult, HandleCommandsParams } from "../commands-types.js"; import { @@ -389,17 +394,25 @@ async function runAcpSteer(params: { sessionKey: string; instruction: string; requestId: string; + channelAdmissionEvidence?: ChannelAdmissionEvidence; }): Promise { const acpManager = getAcpSessionManager(); let output = ""; + const channelAdmission = consumeChannelRunAdmission(params.channelAdmissionEvidence); const admittedRunContext = await prepareAgentRunAdmission({ cfg: params.cfg, operationalRunInstance: createOperationalRunInstanceRef(params.requestId), facts: { runId: params.requestId, agentId: resolveAgentIdFromSessionKey(params.sessionKey), - ingress: { kind: "acp", boundary: "acp.command.steer", state: "present" }, + ingress: { + kind: "acp", + boundary: "acp.command.steer", + state: channelAdmission.ingressState, + }, + ...channelAdmission.facts, }, + onAdmitted: channelAdmission.onAdmitted, }).admit("acp"); try { @@ -477,6 +490,7 @@ export async function handleAcpSteerAction( sessionKey: target.sessionKey, instruction: parsed.value.instruction, requestId: `${resolveCommandRequestId(params)}:steer`, + channelAdmissionEvidence: readChannelContextAdmissionEvidence(params.rootCtx ?? params.ctx), }), fallbackCode: "ACP_TURN_FAILED", fallbackMessage: "ACP steer failed before completion.", diff --git a/src/auto-reply/reply/dispatch-acp.test.ts b/src/auto-reply/reply/dispatch-acp.test.ts index 727fb84ff65d..9f5ff42a5877 100644 --- a/src/auto-reply/reply/dispatch-acp.test.ts +++ b/src/auto-reply/reply/dispatch-acp.test.ts @@ -8,11 +8,18 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; import type { MediaUnderstandingSkipError } from "../../../packages/media-understanding-common/src/errors.js"; import { AcpRuntimeError } from "../../acp/runtime/errors.js"; import type { AcpSessionStoreEntry } from "../../acp/runtime/session-meta.js"; +import { configureExecutionIdentityAdmissionSink } from "../../audit/execution-identity-admission.js"; +import { buildChannelInboundEventContext } from "../../channels/inbound-event/context.js"; +import { + configureChannelAdmissionEvidenceCollection, + createChannelParticipantAdmissionEvidence, +} from "../../channels/message-access/admission-evidence.js"; import type { OpenClawConfig } from "../../config/config.js"; import type { SessionBindingRecord } from "../../infra/outbound/session-binding-service.js"; import type { ApplyMediaUnderstandingResult } from "../../media-understanding/apply.js"; import { isImageAttachment } from "../../media-understanding/attachments.normalize.js"; import { withFetchPreconnect } from "../../test-utils/fetch-mock.js"; +import type { FinalizedRuntimeMsgContext } from "../templating.js"; import { resolveAgentTurnAttachments, resolveInlineAgentImageAttachments, @@ -23,6 +30,7 @@ import { appendRecentHistoryImageContext, resolveRecentInboundHistoryImages, } from "./history-media.js"; +import { finalizeInboundContext } from "./inbound-context.js"; import { createReplyDispatcher } from "./reply-dispatcher.js"; import type { ReplyDispatcher } from "./reply-dispatcher.types.js"; import { buildTestCtx } from "./test-ctx.js"; @@ -368,16 +376,19 @@ async function runDispatch(params: { opts?: { reason?: string; error?: string }, ) => void; markIdle?: (reason: string) => void; + ctx?: FinalizedRuntimeMsgContext; }) { const targetSessionKey = params.sessionKeyOverride ?? sessionKey; return tryDispatchAcpReplyCore({ - ctx: buildTestCtx({ - Provider: "discord", - Surface: "discord", - SessionKey: targetSessionKey, - BodyForAgent: params.bodyForAgent, - ...params.ctxOverrides, - }), + ctx: + params.ctx ?? + buildTestCtx({ + Provider: "discord", + Surface: "discord", + SessionKey: targetSessionKey, + BodyForAgent: params.bodyForAgent, + ...params.ctxOverrides, + }), cfg: params.cfg ?? createAcpTestConfig(), dispatcher: params.dispatcher ?? createDispatcher().dispatcher, ...(params.runId ? { runId: params.runId } : {}), @@ -539,6 +550,56 @@ describe("tryDispatchAcpReplyCore", () => { globalThis.fetch = originalFetch; }); + it("admits ACP message turns with the original channel participant", async () => { + const captured: unknown[] = []; + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + const clearSink = configureExecutionIdentityAdmissionSink((work) => { + captured.push(work); + return true; + }); + try { + setReadyAcpResolution(); + const evidence = createChannelParticipantAdmissionEvidence({ + channelId: "discord", + accountId: "default", + participantId: "person-42", + }); + const ctx = finalizeInboundContext( + await buildChannelInboundEventContext({ + channel: "discord", + accountId: "default", + messageId: "msg-acp", + from: "discord:channel:room-1", + sender: { id: "person-42" }, + conversation: { kind: "group", id: "room-1" }, + route: { agentId: "main", routeSessionKey: sessionKey }, + reply: { to: "discord:channel:room-1" }, + message: { rawBody: "run acp", bodyForAgent: "run acp" }, + channelParticipantEvidence: evidence, + }), + ); + + await runDispatch({ + bodyForAgent: "run acp", + cfg: createAcpTestConfig({ logging: { audit: { executionIdentity: true } } }), + ctx, + }); + + expect(captured).toMatchObject([ + { + kind: "capture", + envelope: { + ingress: { kind: "acp", state: "present" }, + invoker: { state: "present", kind: "person" }, + }, + }, + ]); + } finally { + clearSink(); + clearCollection(); + } + }); + it("projects normal ACP dispatch lifecycle and tool events into audit diagnostics", async () => { setReadyAcpResolution(); mockToolLifecycleTurn("tool-audit"); diff --git a/src/auto-reply/reply/dispatch-acp.ts b/src/auto-reply/reply/dispatch-acp.ts index 0735ebcc3536..820d55c271e3 100644 --- a/src/auto-reply/reply/dispatch-acp.ts +++ b/src/auto-reply/reply/dispatch-acp.ts @@ -25,6 +25,7 @@ import { import { resolveAgentDir, resolveAgentWorkspaceDir } from "../../agents/agent-scope.js"; import { toolPolicyRestrictsTools } from "../../agents/tool-policy.js"; import type { ChatType } from "../../channels/chat-type.js"; +import { readChannelContextAdmissionEvidence } from "../../channels/message-access/admission-evidence.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import type { TtsAutoMode } from "../../config/types.tts.js"; import { logVerbose } from "../../globals.js"; @@ -52,6 +53,7 @@ import { resolveAgentTurnAttachments, resolveInlineAgentImageAttachments, } from "./agent-turn-attachments.js"; +import { consumeChannelRunAdmission } from "./channel-run-admission.js"; import { createAcpDispatchDeliveryCoordinator, type AcpDispatchDeliveryCoordinator, @@ -801,14 +803,23 @@ export async function tryDispatchAcpReplyCore(params: { } turnDispatched = true; + const channelAdmission = consumeChannelRunAdmission( + readChannelContextAdmissionEvidence(params.ctx), + ); admittedRunContext = await prepareAgentRunAdmission({ cfg: params.cfg, operationalRunInstance: createOperationalRunInstanceRef(requestId), facts: { runId: requestId, agentId: acpAgentId, - ingress: { kind: "acp", boundary: "auto-reply.acp", state: "present" }, + ingress: { + kind: "acp", + boundary: "auto-reply.acp", + state: channelAdmission.ingressState, + }, + ...channelAdmission.facts, }, + onAdmitted: channelAdmission.onAdmitted, }).admit("acp"); await acpManager.runTurn({ admittedRunContext, diff --git a/src/auto-reply/reply/get-reply-run-execute.ts b/src/auto-reply/reply/get-reply-run-execute.ts index 4a5628957ffe..69d464cb7e29 100644 --- a/src/auto-reply/reply/get-reply-run-execute.ts +++ b/src/auto-reply/reply/get-reply-run-execute.ts @@ -12,6 +12,7 @@ import { import { resolveFastModeState } from "../../agents/fast-mode.js"; import { runAgentHarnessBeforeMessageWriteHook } from "../../agents/harness/hook-helpers.js"; import { resolveOwnerPromptNumbers } from "../../agents/owner-display.js"; +import { readChannelContextAdmissionEvidence } from "../../channels/message-access/admission-evidence.js"; import { conversationIdentityFromMsgContext } from "../../config/sessions/conversation-identity.js"; import { resolveGroupSessionKey } from "../../config/sessions/group.js"; import { normalizeMediaFacts } from "../../media/media-facts.js"; @@ -329,6 +330,8 @@ export async function executePreparedReplyRun(state: PreparedReplyRunAdmission) ...(userTurnTranscriptRecorder ? { userTurnTranscriptRecorder } : {}), currentInboundEventKind: inboundEventKind, currentInboundAudio: hasInboundAudio(sessionCtx), + channelAdmissionEvidence: + readChannelContextAdmissionEvidence(ctx) ?? readChannelContextAdmissionEvidence(sessionCtx), currentInboundContext, ...(queuedFollowupAbortSignal ? { abortSignal: queuedFollowupAbortSignal } : {}), deliveryCorrelations: opts?.queuedDeliveryCorrelations, diff --git a/src/auto-reply/reply/queue.collect.test.ts b/src/auto-reply/reply/queue.collect.test.ts index 21d3dd12c763..42198c74be7d 100644 --- a/src/auto-reply/reply/queue.collect.test.ts +++ b/src/auto-reply/reply/queue.collect.test.ts @@ -4,6 +4,11 @@ import os from "node:os"; import path from "node:path"; import { describe, expect, it, vi } from "vitest"; import { createDeferred } from "../../../test/helpers/promise.js"; +import { + configureChannelAdmissionEvidenceCollection, + consumeChannelAdmissionEvidence, + createChannelParticipantAdmissionEvidence, +} from "../../channels/message-access/admission-evidence.js"; import { loadTranscriptEvents, replaceSessionEntry, @@ -2070,6 +2075,125 @@ describe("followup queue collect routing", () => { expect(calls[1]?.prompt).toContain("(from Owner)"); }); + it("preserves sender-scoped batching while identity collection is disabled", async () => { + const cleanup = configureChannelAdmissionEvidenceCollection(false); + try { + const { key, calls, done, runFollowup, settings } = createQueueCase( + `test-collect-identity-disabled-${Date.now()}`, + {}, + 2, + ); + for (const senderId of ["user-1", "user-2"]) { + const item = createRun({ + prompt: `from ${senderId}`, + originatingChannel: "slack", + originatingTo: "channel:A", + }); + enqueueFollowupRun( + key, + { + ...item, + run: { ...item.run, senderId, senderIsOwner: false }, + }, + settings, + ); + } + + await drainRecordedQueue(key, runFollowup, done); + await vi.waitFor(() => expect(getExistingFollowupQueue(key)).toBeUndefined()); + + expect(calls.map((call) => call.run.senderId)).toEqual(["user-1", "user-2"]); + } finally { + cleanup(); + } + }); + + it("keeps same-participant evidence and clears authority for a mixed-participant batch", async () => { + const cleanup = configureChannelAdmissionEvidenceCollection(true); + try { + const sameCase = createQueueCase(`test-collect-identity-same-${Date.now()}`); + for (const prompt of ["same one", "same two"]) { + const item = createRun({ + prompt, + originatingChannel: "slack", + originatingTo: "channel:A", + }); + enqueueFollowupRun( + sameCase.key, + { + ...item, + channelAdmissionEvidence: createChannelParticipantAdmissionEvidence({ + channelId: "slack", + accountId: "default", + participantId: "user-1", + }), + run: { ...item.run, senderId: "user-1", senderIsOwner: false }, + }, + sameCase.settings, + ); + } + await drainRecordedQueue(sameCase.key, sameCase.runFollowup, sameCase.done); + await vi.waitFor(() => expect(getExistingFollowupQueue(sameCase.key)).toBeUndefined()); + expect(sameCase.calls).toHaveLength(1); + expect(sameCase.calls[0]?.run.senderId).toBe("user-1"); + expect( + consumeChannelAdmissionEvidence(sameCase.calls[0]?.channelAdmissionEvidence), + ).toMatchObject({ + ingressState: "present", + invoker: { state: "present", kind: "person" }, + }); + + const mixedCase = createQueueCase(`test-collect-identity-mixed-${Date.now()}`); + for (const senderId of ["user-1", "user-2"]) { + const item = createRun({ + prompt: `mixed ${senderId}`, + originatingChannel: "slack", + originatingTo: "channel:A", + }); + enqueueFollowupRun( + mixedCase.key, + { + ...item, + channelAdmissionEvidence: createChannelParticipantAdmissionEvidence({ + channelId: "slack", + accountId: "default", + participantId: senderId, + }), + run: { + ...item.run, + senderId, + senderName: senderId, + senderE164: `+1555000${senderId.at(-1)}`, + senderIsOwner: false, + traceAuthorized: true, + ownerNumbers: ["+15550000000"], + }, + }, + mixedCase.settings, + ); + } + await drainRecordedQueue(mixedCase.key, mixedCase.runFollowup, mixedCase.done); + await vi.waitFor(() => expect(getExistingFollowupQueue(mixedCase.key)).toBeUndefined()); + + expect(mixedCase.calls).toHaveLength(1); + expect(mixedCase.calls[0]?.run).toMatchObject({ + senderIsOwner: false, + traceAuthorized: false, + ownerNumbers: [], + }); + expect(mixedCase.calls[0]?.run.senderId).toBeUndefined(); + expect(mixedCase.calls[0]?.run.senderE164).toBeUndefined(); + expect( + consumeChannelAdmissionEvidence(mixedCase.calls[0]?.channelAdmissionEvidence), + ).toMatchObject({ + ingressState: "unknown", + invoker: { state: "unknown" }, + }); + } finally { + cleanup(); + } + }); + it("splits collect batches when queued cancellation owners differ", async () => { const key = `test-collect-cancel-owner-split-${Date.now()}`; const { calls, done, runFollowup } = createDrainRecorder(2); diff --git a/src/auto-reply/reply/queue.media-carrier.test.ts b/src/auto-reply/reply/queue.media-carrier.test.ts index f27493a5cad2..5aba7384ee39 100644 --- a/src/auto-reply/reply/queue.media-carrier.test.ts +++ b/src/auto-reply/reply/queue.media-carrier.test.ts @@ -1,6 +1,12 @@ // Prompt media carrier tests cover collect batching, deferral, and retry identity. import { afterEach, describe, expect, it } from "vitest"; import { createDeferred } from "../../../test/helpers/promise.js"; +import { + compareChannelAdmissionParticipants, + configureChannelAdmissionEvidenceCollection, + consumeChannelAdmissionEvidence, + createChannelParticipantAdmissionEvidence, +} from "../../channels/message-access/admission-evidence.js"; import type { FollowupRun, QueueSettings } from "./queue.js"; import { enqueueFollowupRun, FollowupRunDeferredError, scheduleFollowupDrain } from "./queue.js"; import { createQueueTestRun } from "./queue.test-helpers.js"; @@ -8,16 +14,23 @@ import { createOverflowSummaryRetrySource } from "./queue/drain.js"; import { clearFollowupQueue } from "./queue/state.js"; const queueKeys = new Set(); +const evidenceCleanups = new Set<() => void>(); afterEach(() => { for (const key of queueKeys) { clearFollowupQueue(key); } queueKeys.clear(); + for (const cleanup of evidenceCleanups) { + cleanup(); + } + evidenceCleanups.clear(); }); describe("followup prompt media carrier", () => { it("keeps collected prompt bytes and ordered facts stable across deferred admission", async () => { + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + evidenceCleanups.add(clearCollection); const key = `prompt-media-collect-${Date.now()}`; queueKeys.add(key); const settings: QueueSettings = { mode: "collect", debounceMs: 0 }; @@ -30,6 +43,10 @@ describe("followup prompt media carrier", () => { ] as const) { const run = createQueueTestRun({ prompt }); run.media = [{ path, contentType }]; + run.channelAdmissionEvidence = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); enqueueFollowupRun(key, run, settings); } @@ -59,17 +76,31 @@ describe("followup prompt media carrier", () => { { path: "/tmp/b.pdf", contentType: "application/pdf" }, ], ]); + expect( + compareChannelAdmissionParticipants(calls.map((run) => run.channelAdmissionEvidence)), + ).toBe("same"); + expect(consumeChannelAdmissionEvidence(calls[1]?.channelAdmissionEvidence)).toMatchObject({ + ingressState: "present", + invoker: { state: "present", kind: "person" }, + }); }); it("preserves facts when an overflow source is rebuilt for retry", () => { + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + evidenceCleanups.add(clearCollection); const source = createQueueTestRun({ prompt: "[media attached: /tmp/retry.png (image/png)]\nretry me", }); source.media = [{ path: "/tmp/retry.png", contentType: "image/png" }]; + source.channelAdmissionEvidence = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); const retry = createOverflowSummaryRetrySource(source); expect(retry.prompt).toBe(source.prompt); expect(retry.media).toEqual(source.media); + expect(retry.channelAdmissionEvidence).toBe(source.channelAdmissionEvidence); }); }); diff --git a/src/auto-reply/reply/queue/drain.ts b/src/auto-reply/reply/queue/drain.ts index df971bf0b9eb..eb4a35a6703c 100644 --- a/src/auto-reply/reply/queue/drain.ts +++ b/src/auto-reply/reply/queue/drain.ts @@ -4,6 +4,10 @@ import { stableStringify } from "@openclaw/normalization-core"; import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce"; import { runAgentHarnessBeforeMessageWriteHook } from "../../../agents/harness/hook-helpers.js"; import { normalizeChatType } from "../../../channels/chat-type.js"; +import { + combineChannelAdmissionEvidence, + compareChannelAdmissionParticipants, +} from "../../../channels/message-access/admission-evidence.js"; import { resolveSessionStorePathCore } from "../../../config/sessions.js"; import { loadSessionEntryReadOnly } from "../../../config/sessions/session-accessor.js"; // Drains queued follow-up runs while preserving route and session identity. @@ -159,28 +163,60 @@ function resolveOriginRoutingMetadata(items: FollowupRun[]): OriginRoutingMetada // Keep this key aligned with the fields that affect per-message authorization or // exec-context propagation in collect-mode batching. Display-only sender fields // stay out of the key so profile/name drift does not force conservative splits. +// Raw sender keys remain the default; only verified opaque admission carriers +// may group participants before the aggregate clears mixed sender authority. // Fields like authProfileId, elevatedLevel, ownerNumbers, and config are // intentionally excluded because they are session-level or not consulted in // per-message authorization checks. -function resolveFollowupAuthorizationKey(run: FollowupRun["run"]): string { +function hasVerifiedAdmissionParticipant(run: FollowupRun): boolean { + return compareChannelAdmissionParticipants([run.channelAdmissionEvidence]) === "same"; +} + +function resolveFollowupAuthorizationKey(run: FollowupRun): string { + const execution = run.run; + const useOpaqueParticipant = hasVerifiedAdmissionParticipant(run); return JSON.stringify([ - run.senderId ?? "", - JSON.stringify(run.channelContext ?? null), - stableStringify(run.conversationToolPolicy ?? null), - run.senderE164 ?? "", - run.senderIsOwner === true, - run.execOverrides?.host ?? "", - run.execOverrides?.security ?? "", - run.execOverrides?.ask ?? "", - run.execOverrides?.node ?? "", - run.execOverrides?.nodeCwd ?? "", - run.bashElevated?.enabled === true, - run.bashElevated?.allowed === true, - run.bashElevated?.defaultLevel ?? "", - run.approvalReviewerDeviceId ?? "", + useOpaqueParticipant ? "" : (execution.senderId ?? ""), + JSON.stringify(execution.channelContext ?? null), + stableStringify(execution.conversationToolPolicy ?? null), + useOpaqueParticipant ? "" : (execution.senderE164 ?? ""), + execution.senderIsOwner === true, + execution.execOverrides?.host ?? "", + execution.execOverrides?.security ?? "", + execution.execOverrides?.ask ?? "", + execution.execOverrides?.node ?? "", + execution.execOverrides?.nodeCwd ?? "", + execution.bashElevated?.enabled === true, + execution.bashElevated?.allowed === true, + execution.bashElevated?.defaultLevel ?? "", + execution.approvalReviewerDeviceId ?? "", ]); } +function resolveCollectedRun(items: readonly FollowupRun[], source: FollowupRun["run"]) { + const participantComparison = compareChannelAdmissionParticipants( + items.map((item) => item.channelAdmissionEvidence), + ); + if ( + participantComparison === "same" || + !items.every((item) => hasVerifiedAdmissionParticipant(item)) + ) { + return source; + } + // Mixed or unverifiable people share no downstream sender authority. The + // opaque admission aggregate records unknown identity at the run boundary. + return { + ...source, + senderId: undefined, + senderName: undefined, + senderUsername: undefined, + senderE164: undefined, + senderIsOwner: false, + traceAuthorized: false, + ownerNumbers: [], + }; +} + export function resolveFollowupDeliveryContextKey(run: FollowupRun): string { const execution = run.run; const provenance = execution.inputProvenance; @@ -195,7 +231,7 @@ export function resolveFollowupDeliveryContextKey(run: FollowupRun): string { resolveFollowupReplyAnchor(run) ?? "", run.originatingReplyToMode ?? "", normalizeChatType(run.originatingChatType) ?? "", - resolveFollowupAuthorizationKey(execution), + resolveFollowupAuthorizationKey(run), run.turnAdoptionLifecycle?.ownerKey ?? "", normalizeOptionalString(execution.runtimePolicySessionKey ?? execution.sessionKey) ?? "", execution.messageProvider ?? "", @@ -318,6 +354,7 @@ type FollowupRuntimeMetadata = Pick< | "currentInboundEventKind" | "currentInboundAudio" | "currentInboundContext" + | "channelAdmissionEvidence" | "abortSignal" | "queueAbortSignal" | "deliveryCorrelations" @@ -558,6 +595,9 @@ function collectRuntimeMetadata( currentInboundEventKind: currentTurnSource?.currentInboundEventKind, currentInboundAudio: currentTurnSource?.currentInboundAudio, currentInboundContext: collectCurrentInboundContext(items), + channelAdmissionEvidence: combineChannelAdmissionEvidence( + items.map((item) => item.channelAdmissionEvidence), + ), abortSignal, queueAbortSignal: items.find((item) => item.queueAbortSignal)?.queueAbortSignal, deliveryCorrelations: deliveryCorrelations.length > 0 ? deliveryCorrelations : undefined, @@ -915,6 +955,7 @@ export function createOverflowSummaryRetrySource(source: FollowupRun): FollowupR queueAbortSignal: source.queueAbortSignal, transcriptPrompt: source.transcriptPrompt, media: source.media, + channelAdmissionEvidence: source.channelAdmissionEvidence, messageId: source.messageId, summaryLine: source.summaryLine, enqueuedAt: source.enqueuedAt, @@ -979,6 +1020,7 @@ async function runSyntheticOverflowSummary(params: { errorContext: "followup overflow summary transcript", }); const currentInboundEventKind = resolveOverflowSummaryInboundEventKind(params.sources); + const runtimeMetadata = collectRuntimeMetadata(params.sources); let admitted = false; await params.runFollowup({ prompt: params.prompt, @@ -986,10 +1028,11 @@ async function runSyntheticOverflowSummary(params: { transcriptPrompt: params.prompt, messageId: params.source.messageId, userTurnTranscriptRecorder, - run: params.source.run, + run: resolveCollectedRun(params.sources, params.source.run), enqueuedAt: Date.now(), abortSignal: params.abortSignal, - onReplyAdmissionWaitChange: collectRuntimeMetadata(params.sources).onReplyAdmissionWaitChange, + channelAdmissionEvidence: runtimeMetadata.channelAdmissionEvidence, + onReplyAdmissionWaitChange: runtimeMetadata.onReplyAdmissionWaitChange, ...(params.onAdmitted ? { turnAdoptionLifecycle: { @@ -1270,7 +1313,9 @@ export function scheduleFollowupDrain( } assertSingleAdmissionOwner(activeGroupItems); const groupSource = activeGroupItems.at(-1); - const run = groupSource?.run ?? queue.lastRun; + const run = groupSource + ? resolveCollectedRun(activeGroupItems, groupSource.run) + : queue.lastRun; if (!run) { break; } diff --git a/src/auto-reply/reply/queue/types.ts b/src/auto-reply/reply/queue/types.ts index 80b7dfd08d1f..7dc548e13ca8 100644 --- a/src/auto-reply/reply/queue/types.ts +++ b/src/auto-reply/reply/queue/types.ts @@ -9,6 +9,7 @@ import type { ModelFallbackRouteResolution } from "../../../agents/model-fallbac import type { SilentReplyPromptMode } from "../../../agents/system-prompt.types.js"; import type { ChatType } from "../../../channels/chat-type.js"; import type { InboundEventKind } from "../../../channels/inbound-event/kind.js"; +import type { ChannelAdmissionEvidence } from "../../../channels/message-access/admission-evidence.js"; import type { SessionEntry, SessionToolOverrides } from "../../../config/sessions.js"; import type { ReplyToMode } from "../../../config/types.base.js"; import type { OpenClawConfig } from "../../../config/types.openclaw.js"; @@ -81,6 +82,8 @@ export type FollowupRun = { currentInboundEventKind?: InboundEventKind; /** Whether the current inbound message contained audio for inbound-only TTS policy. */ currentInboundAudio?: boolean; + /** Host-minted participant evidence; raw channel identities never live on this object. */ + channelAdmissionEvidence?: ChannelAdmissionEvidence; /** Explicit current-turn context that should be visible for this run but not persisted as user text. */ currentInboundContext?: CurrentInboundPromptContext; /** Abort signal for turns that are canceled by their source-channel admission fence. */ diff --git a/src/channels/inbound-event/context.ts b/src/channels/inbound-event/context.ts index 7f78ebeb0ece..e17e6e92682c 100644 --- a/src/channels/inbound-event/context.ts +++ b/src/channels/inbound-event/context.ts @@ -24,7 +24,12 @@ import type { GroupToolPolicyConfig } from "../../config/types.tools.js"; import type { PluginHookChannelContext } from "../../plugins/hook-channel-context.types.js"; import { shouldIncludeSupplementalContext } from "../../security/context-visibility.js"; import type { InboundImplicitMentionKind } from "../mention-gating.js"; +import { + bindChannelContextAdmissionEvidence, + type ChannelAdmissionEvidence, +} from "../message-access/admission-evidence.js"; import type { ChannelIngressCommandAccess } from "../message-access/runtime-types.js"; +import type { ResolvedChannelMessageIngress } from "../message-access/runtime-types.js"; import type { CommandFacts, ConversationFacts, @@ -101,6 +106,10 @@ export type BuildChannelInboundEventContextParams = { finalize?: FinalizeInboundContextFn; finalizeOptions?: FinalizeInboundContextOptions; extra?: Record; + /** Host-resolved ingress result, or an explicit unsupported adapter marker. */ + channelIngress?: ResolvedChannelMessageIngress | "unsupported"; + /** Explicit attribution-only evidence minted by the host ingress SDK. */ + channelParticipantEvidence?: ChannelAdmissionEvidence; }; /** * @deprecated Prefer `BuildChannelInboundEventContextParams` with @@ -569,7 +578,17 @@ export function buildChannelInboundEventContext( suppressSelfQuoteMedia: params.suppressSelfQuoteMedia, }) : finalizeChannelInboundContextValue(finalizeParams); - return isPromiseLike(result) - ? result.then((finalized) => finalized.context as BuiltChannelInboundEventContext) - : (result.context as BuiltChannelInboundEventContext); + const bindEvidence = (finalized: FinalizeChannelInboundContextResult) => { + const built = finalized.context as BuiltChannelInboundEventContext; + bindChannelContextAdmissionEvidence({ + context: built, + channelId: params.channel, + accountId: params.accountId, + ingress: params.channelIngress, + evidence: params.channelParticipantEvidence, + rawPrincipalRef: params.sender.id, + }); + return built; + }; + return isPromiseLike(result) ? result.then(bindEvidence) : bindEvidence(result); } diff --git a/src/channels/message-access/admission-evidence.test.ts b/src/channels/message-access/admission-evidence.test.ts new file mode 100644 index 000000000000..e38eb321ef4d --- /dev/null +++ b/src/channels/message-access/admission-evidence.test.ts @@ -0,0 +1,254 @@ +import { describe, expect, it, vi } from "vitest"; +import { buildChannelInboundEventContext } from "../inbound-event/context.js"; +import { + combineChannelAdmissionEvidence, + configureChannelAdmissionEvidenceCollection, + consumeChannelAdmissionEvidence, + createChannelParticipantAdmissionEvidence, + readChannelContextAdmissionEvidence, + type ChannelAdmissionEvidence, +} from "./admission-evidence.js"; +import { resolveStableChannelMessageIngress } from "./runtime.js"; + +async function buildAdmittedContext(participantId: string, allowFrom = [participantId]) { + const channelIngress = await resolveStableChannelMessageIngress({ + channelId: "test", + accountId: "acct:primary", + subject: { stableId: participantId }, + conversation: { kind: "direct", id: "dm-1" }, + dmPolicy: "allowlist", + groupPolicy: "allowlist", + allowFrom, + }); + return buildChannelInboundEventContext({ + channel: "test", + accountId: "acct:primary", + messageId: "msg-1", + from: "test:route:dm-1", + sender: { id: participantId }, + conversation: { kind: "direct", id: "dm-1" }, + route: { agentId: "main", routeSessionKey: "agent:main:test:dm:dm-1" }, + reply: { to: "test:route:dm-1" }, + message: { rawBody: "hello" }, + channelIngress, + }); +} + +describe("channel admission evidence", () => { + it("carries the resolver participant to one run admission without route inference", async () => { + const cleanup = configureChannelAdmissionEvidenceCollection(true); + try { + const context = await buildAdmittedContext("person:42"); + const evidence = readChannelContextAdmissionEvidence(context); + const consumed = consumeChannelAdmissionEvidence(evidence); + + expect(consumed).toEqual({ + ingressState: "present", + invoker: { + state: "present", + kind: "person", + rawPrincipalRef: '["test","acct:primary","person:42"]', + }, + assuranceRef: "channel-admission", + decisionCoverage: "enforced", + }); + expect(Object.isFrozen(consumed)).toBe(true); + expect(Object.isFrozen(consumed.invoker)).toBe(true); + expect(consumeChannelAdmissionEvidence(evidence)).toMatchObject({ + ingressState: "unknown", + invoker: { state: "unknown" }, + decisionCoverage: "unknown", + }); + } finally { + cleanup(); + } + }); + + it("reports same-participant collection while mixed participants fail closed", () => { + const cleanup = configureChannelAdmissionEvidenceCollection(true); + try { + const first = createChannelParticipantAdmissionEvidence({ + channelId: "test", + accountId: "a:b", + participantId: "c", + }); + const same = createChannelParticipantAdmissionEvidence({ + channelId: "test", + accountId: "a:b", + participantId: "c", + }); + const tupleCollisionCandidate = createChannelParticipantAdmissionEvidence({ + channelId: "test", + accountId: "a", + participantId: "b:c", + }); + + expect( + consumeChannelAdmissionEvidence(combineChannelAdmissionEvidence([first, same])), + ).toEqual({ + ingressState: "present", + invoker: { + state: "present", + kind: "person", + rawPrincipalRef: '["test","a:b","c"]', + }, + assuranceRef: "channel-admission", + decisionCoverage: "attribution-only", + }); + expect( + consumeChannelAdmissionEvidence( + combineChannelAdmissionEvidence([ + createChannelParticipantAdmissionEvidence({ + channelId: "test", + accountId: "a:b", + participantId: "c", + }), + tupleCollisionCandidate, + ]), + ), + ).toMatchObject({ + ingressState: "unknown", + invoker: { state: "unknown" }, + decisionCoverage: "unknown", + }); + } finally { + cleanup(); + } + }); + + it("keeps wildcard admission attribution-only because identity did not affect the outcome", async () => { + const cleanup = configureChannelAdmissionEvidenceCollection(true); + try { + const context = await buildAdmittedContext("person-42", ["*"]); + expect( + consumeChannelAdmissionEvidence(readChannelContextAdmissionEvidence(context)), + ).toMatchObject({ + ingressState: "present", + invoker: { state: "present", kind: "person" }, + decisionCoverage: "attribution-only", + }); + } finally { + cleanup(); + } + }); + + it("rejects forged and prior-lifecycle carriers and stays empty when collection is disabled", () => { + const cleanup = configureChannelAdmissionEvidenceCollection(true); + const stale = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); + cleanup(); + + const nextCleanup = configureChannelAdmissionEvidenceCollection(true); + try { + expect(consumeChannelAdmissionEvidence(stale)).toMatchObject({ ingressState: "unknown" }); + expect( + consumeChannelAdmissionEvidence({ + kind: "channel-admission-evidence", + } as ChannelAdmissionEvidence), + ).toMatchObject({ ingressState: "unknown" }); + } finally { + nextCleanup(); + } + + expect( + createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }), + ).toBeUndefined(); + }); + + it("distinguishes unsupported, omitted, and mismatched adapter handoffs", () => { + const cleanup = configureChannelAdmissionEvidenceCollection(true); + try { + const base = { + channel: "legacy", + accountId: "default", + messageId: "msg-1", + from: "legacy:route:room-1", + sender: { id: "person-1" }, + conversation: { kind: "direct" as const, id: "room-1" }, + route: { agentId: "main", routeSessionKey: "agent:main:legacy:dm:room-1" }, + reply: { to: "legacy:route:room-1" }, + message: { rawBody: "hello" }, + }; + const unsupported = buildChannelInboundEventContext({ + ...base, + channelIngress: "unsupported", + }); + const omitted = buildChannelInboundEventContext(base); + const mismatched = buildChannelInboundEventContext({ + ...base, + channelParticipantEvidence: createChannelParticipantAdmissionEvidence({ + channelId: "legacy", + accountId: "default", + participantId: "someone-else", + }), + }); + + expect( + consumeChannelAdmissionEvidence(readChannelContextAdmissionEvidence(unsupported)), + ).toMatchObject({ ingressState: "unsupported", decisionCoverage: "unsupported" }); + expect( + consumeChannelAdmissionEvidence(readChannelContextAdmissionEvidence(omitted)), + ).toMatchObject({ + ingressState: "unknown", + decisionCoverage: "unknown", + }); + expect( + consumeChannelAdmissionEvidence(readChannelContextAdmissionEvidence(mismatched)), + ).toMatchObject({ ingressState: "unknown", decisionCoverage: "unknown" }); + } finally { + cleanup(); + } + }); + + it("expires a carrier at the bounded retention edge", () => { + vi.useFakeTimers(); + vi.setSystemTime(1_000); + const cleanup = configureChannelAdmissionEvidenceCollection(true); + try { + const evidence = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); + vi.setSystemTime(1_000 + 30 * 24 * 60 * 60_000 + 1); + expect(consumeChannelAdmissionEvidence(evidence)).toMatchObject({ + ingressState: "unknown", + decisionCoverage: "unknown", + }); + } finally { + cleanup(); + vi.useRealTimers(); + } + }); + + it("bounds aggregate fan-in and participant material", () => { + const cleanup = configureChannelAdmissionEvidenceCollection(true); + try { + const oversizedParticipant = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "x".repeat(4_097), + }); + expect(consumeChannelAdmissionEvidence(oversizedParticipant)).toMatchObject({ + ingressState: "unknown", + }); + + const sources = Array.from({ length: 17 }, (_, index) => + createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: `person-${index}`, + }), + ); + expect( + consumeChannelAdmissionEvidence(combineChannelAdmissionEvidence(sources)), + ).toMatchObject({ + ingressState: "unknown", + }); + } finally { + cleanup(); + } + }); +}); diff --git a/src/channels/message-access/admission-evidence.ts b/src/channels/message-access/admission-evidence.ts new file mode 100644 index 000000000000..6281799ff546 --- /dev/null +++ b/src/channels/message-access/admission-evidence.ts @@ -0,0 +1,462 @@ +import type { DecisionReceiptV1 } from "../../../packages/gateway-protocol/src/index.js"; +import { resolveGlobalSingleton } from "../../shared/global-singleton.js"; +import type { ResolvedChannelMessageIngress } from "./runtime-types.js"; + +export type ChannelAdmissionEvidence = Readonly<{ + kind: "channel-admission-evidence"; +}>; + +type ChannelAdmissionContribution = Readonly<{ + participant: + | { state: "present"; rawPrincipalRef: string } + | { state: "unknown" } + | { state: "unsupported" }; + decision?: Readonly<{ + participantAware: boolean; + outcomeAffecting: boolean; + }>; +}>; + +type ChannelAdmissionEvidencePayload = + | Readonly<{ + kind: "leaf"; + createdAt: number; + generation: number; + contribution: ChannelAdmissionContribution; + }> + | Readonly<{ + kind: "aggregate"; + createdAt: number; + generation: number; + sources: readonly (ChannelAdmissionEvidence | undefined)[]; + }>; + +type ConsumedChannelAdmissionEvidence = Readonly<{ + ingressState: "present" | "unknown" | "unsupported"; + invoker: { state: "present"; kind: "person"; rawPrincipalRef: string } | { state: "unknown" }; + assuranceRef?: string; + decisionCoverage?: "enforced" | "attribution-only" | "unknown" | "unsupported"; +}>; + +const CHANNEL_ADMISSION_EVIDENCE_MAX_CONTRIBUTIONS = 16; +const CHANNEL_ADMISSION_EVIDENCE_MAX_AGE_MS = 30 * 24 * 60 * 60_000; +const CHANNEL_ADMISSION_EVIDENCE_STATE_KEY = Symbol.for("openclaw.channelAdmissionEvidenceState"); +const state = resolveGlobalSingleton(CHANNEL_ADMISSION_EVIDENCE_STATE_KEY, () => ({ + collectionEnabled: false, + generation: 0, + payloadByEvidence: new WeakMap(), + evidenceByIngress: new WeakMap(), + evidenceByContext: new WeakMap(), + consumedEvidence: new WeakSet(), + decisionSink: undefined as ((receipt: DecisionReceiptV1) => boolean) | undefined, +})); + +export function configureChannelAdmissionEvidenceCollection(enabled: boolean): () => void { + const generation = ++state.generation; + state.collectionEnabled = enabled; + return () => { + if (state.generation === generation) { + state.collectionEnabled = false; + state.generation += 1; + } + }; +} + +export function configureChannelAdmissionDecisionSink( + sink: (receipt: DecisionReceiptV1) => boolean, +): () => void { + state.decisionSink = sink; + return () => { + if (state.decisionSink === sink) { + state.decisionSink = undefined; + } + }; +} + +function mintChannelAdmissionEvidence( + payload: + | Omit, "createdAt" | "generation"> + | Omit< + Extract, + "createdAt" | "generation" + >, +): ChannelAdmissionEvidence | undefined { + if (!state.collectionEnabled) { + return undefined; + } + const evidence = Object.freeze({ kind: "channel-admission-evidence" as const }); + state.payloadByEvidence.set( + evidence, + Object.freeze({ ...payload, createdAt: Date.now(), generation: state.generation }), + ); + return evidence; +} + +function scopedParticipantRef(params: { + channelId: string; + accountId?: string; + rawPrincipalRef: string | number | null | undefined; +}): string | undefined { + const channelId = params.channelId.trim(); + const accountId = params.accountId?.trim() || "default"; + const rawPrincipalRef = + params.rawPrincipalRef == null ? "" : String(params.rawPrincipalRef).trim(); + if (!channelId || !rawPrincipalRef) { + return undefined; + } + // Preserve tuple boundaries: channel, account, and participant identifiers may + // themselves contain colons or other separators. + const scoped = JSON.stringify([channelId, accountId, rawPrincipalRef]); + return scoped.length <= 4_096 ? scoped : undefined; +} + +function participantContribution(params: { + channelId: string; + accountId?: string; + rawPrincipalRef: string | number | null | undefined; +}): ChannelAdmissionContribution { + const rawPrincipalRef = scopedParticipantRef(params); + return Object.freeze( + rawPrincipalRef + ? { participant: Object.freeze({ state: "present" as const, rawPrincipalRef }) } + : { participant: Object.freeze({ state: "unknown" as const }) }, + ); +} + +/** Explicit attribution-only mint for adapters whose access result is owned elsewhere. */ +export function createChannelParticipantAdmissionEvidence(params: { + channelId: string; + accountId?: string; + participantId: string | number | null | undefined; +}): ChannelAdmissionEvidence | undefined { + return mintChannelAdmissionEvidence({ + kind: "leaf", + contribution: participantContribution({ + channelId: params.channelId, + accountId: params.accountId, + rawPrincipalRef: params.participantId, + }), + }); +} + +/** Bind an admitted resolver result to its host-owned participant without exposing that value. */ +export function bindChannelIngressAdmissionEvidence(params: { + result: ResolvedChannelMessageIngress; + channelId: string; + accountId?: string; + rawPrincipalRef: string | number | null | undefined; + participantOutcomeAffecting: boolean; +}): ResolvedChannelMessageIngress { + if (!state.collectionEnabled || params.result.ingress.admission !== "dispatch") { + return params.result; + } + const contribution = participantContribution(params); + const evidence = mintChannelAdmissionEvidence({ + kind: "leaf", + contribution: Object.freeze({ + ...contribution, + decision: Object.freeze({ + participantAware: contribution.participant.state === "present", + outcomeAffecting: params.participantOutcomeAffecting, + }), + }), + }); + if (evidence) { + state.evidenceByIngress.set(params.result, evidence); + } + return params.result; +} + +function evidenceMatchesContextParticipant(params: { + evidence: ChannelAdmissionEvidence; + channelId: string; + accountId?: string; + rawPrincipalRef: string | number | null | undefined; +}): boolean { + const expected = scopedParticipantRef(params); + const payload = state.payloadByEvidence.get(params.evidence); + return ( + payload?.kind === "leaf" && + payload.contribution.participant.state === "present" && + payload.contribution.participant.rawPrincipalRef === expected + ); +} + +/** Attach private evidence to the finalized context returned by the existing SDK builder. */ +export function bindChannelContextAdmissionEvidence(params: { + context: object; + channelId: string; + accountId?: string; + ingress?: ResolvedChannelMessageIngress | "unsupported"; + evidence?: ChannelAdmissionEvidence; + rawPrincipalRef: string | number | null | undefined; +}): void { + if (!state.collectionEnabled) { + return; + } + const ingressEvidence = + params.ingress && params.ingress !== "unsupported" + ? state.evidenceByIngress.get(params.ingress) + : undefined; + const evidence = params.evidence + ? evidenceMatchesContextParticipant({ ...params, evidence: params.evidence }) + ? params.evidence + : mintChannelAdmissionEvidence({ + kind: "leaf", + contribution: Object.freeze({ participant: { state: "unknown" as const } }), + }) + : params.ingress === "unsupported" + ? mintChannelAdmissionEvidence({ + kind: "leaf", + contribution: Object.freeze({ participant: { state: "unsupported" as const } }), + }) + : ingressEvidence && + evidenceMatchesContextParticipant({ ...params, evidence: ingressEvidence }) + ? ingressEvidence + : params.ingress + ? mintChannelAdmissionEvidence({ + kind: "leaf", + contribution: Object.freeze({ participant: { state: "unknown" as const } }), + }) + : mintChannelAdmissionEvidence({ + kind: "leaf", + contribution: Object.freeze({ participant: { state: "unknown" as const } }), + }); + if (evidence) { + state.evidenceByContext.set(params.context, evidence); + } +} + +export function readChannelContextAdmissionEvidence( + context: object, +): ChannelAdmissionEvidence | undefined { + return state.evidenceByContext.get(context); +} + +/** Preserve private evidence when an owner intentionally replaces a finalized context object. */ +export function copyChannelContextAdmissionEvidence(source: object, target: object): void { + const evidence = state.evidenceByContext.get(source); + if (evidence) { + state.evidenceByContext.set(target, evidence); + } +} + +function activePayload( + evidence: ChannelAdmissionEvidence | undefined, + now: number, +): ChannelAdmissionEvidencePayload | undefined { + if (!evidence || state.consumedEvidence.has(evidence)) { + return undefined; + } + const payload = state.payloadByEvidence.get(evidence); + return payload && + payload.generation === state.generation && + now - payload.createdAt <= CHANNEL_ADMISSION_EVIDENCE_MAX_AGE_MS + ? payload + : undefined; +} + +/** Preserve one source exactly; collected sources get one new bounded opaque aggregate. */ +export function combineChannelAdmissionEvidence( + evidence: readonly (ChannelAdmissionEvidence | undefined)[], +): ChannelAdmissionEvidence | undefined { + if (!state.collectionEnabled) { + return undefined; + } + if (evidence.length === 1) { + return evidence[0]; + } + if (evidence.length > CHANNEL_ADMISSION_EVIDENCE_MAX_CONTRIBUTIONS) { + return mintChannelAdmissionEvidence({ + kind: "leaf", + contribution: Object.freeze({ participant: { state: "unknown" } }), + }); + } + return mintChannelAdmissionEvidence({ kind: "aggregate", sources: Object.freeze([...evidence]) }); +} + +function inspectContributions(params: { + evidence: ChannelAdmissionEvidence | undefined; + now: number; + seen: Set; +}): ChannelAdmissionContribution[] { + const payload = activePayload(params.evidence, params.now); + if (!payload || !params.evidence || params.seen.has(params.evidence)) { + return [{ participant: { state: "unknown" } }]; + } + params.seen.add(params.evidence); + return payload.kind === "leaf" + ? [payload.contribution] + : payload.sources.flatMap((source) => inspectContributions({ ...params, evidence: source })); +} + +/** Compare opaque participants without exposing or consuming their raw references. */ +export function compareChannelAdmissionParticipants( + evidence: readonly (ChannelAdmissionEvidence | undefined)[], +): "same" | "mixed-or-unknown" { + const contributions = evidence.flatMap((candidate) => + inspectContributions({ evidence: candidate, now: Date.now(), seen: new Set() }), + ); + if ( + contributions.length === 0 || + contributions.length > CHANNEL_ADMISSION_EVIDENCE_MAX_CONTRIBUTIONS + ) { + return "mixed-or-unknown"; + } + const participants = contributions.map((item) => item.participant); + const first = participants[0]; + return first?.state === "present" && + participants.every( + (item) => item.state === "present" && item.rawPrincipalRef === first.rawPrincipalRef, + ) + ? "same" + : "mixed-or-unknown"; +} + +function consumeContributions(params: { + evidence: ChannelAdmissionEvidence | undefined; + now: number; + seen: Set; +}): ChannelAdmissionContribution[] { + const payload = activePayload(params.evidence, params.now); + if (!payload || !params.evidence || params.seen.has(params.evidence)) { + return [{ participant: { state: "unknown" } }]; + } + params.seen.add(params.evidence); + state.consumedEvidence.add(params.evidence); + if (payload.kind === "leaf") { + return [payload.contribution]; + } + const contributions = payload.sources.flatMap((source) => + consumeContributions({ ...params, evidence: source }), + ); + return contributions.length <= CHANNEL_ADMISSION_EVIDENCE_MAX_CONTRIBUTIONS + ? contributions + : [{ participant: { state: "unknown" } }]; +} + +function freezeConsumed( + value: Omit & { + invoker: ConsumedChannelAdmissionEvidence["invoker"]; + }, +): ConsumedChannelAdmissionEvidence { + return Object.freeze({ + ...value, + invoker: Object.freeze(value.invoker), + }); +} + +/** Consume one aggregate at run admission. Missing, forged, stale, or reused carriers are unknown. */ +export function consumeChannelAdmissionEvidence( + evidence: ChannelAdmissionEvidence | undefined, +): ConsumedChannelAdmissionEvidence { + const contributions = consumeContributions({ evidence, now: Date.now(), seen: new Set() }); + const participants = contributions.map((item) => item.participant); + const allUnsupported = + participants.length > 0 && participants.every((item) => item.state === "unsupported"); + if (allUnsupported) { + return freezeConsumed({ + ingressState: "unsupported", + invoker: { state: "unknown" }, + decisionCoverage: "unsupported", + }); + } + + const present = participants.filter( + (item): item is Extract<(typeof participants)[number], { state: "present" }> => + item.state === "present", + ); + const sameParticipant = + present.length === participants.length && + present.every((item) => item.rawPrincipalRef === present[0]?.rawPrincipalRef); + if (!sameParticipant || !present[0]) { + return freezeConsumed({ + ingressState: "unknown", + invoker: { state: "unknown" }, + decisionCoverage: "unknown", + }); + } + + const everyDecisionEnforced = contributions.every( + (item) => item.decision?.participantAware && item.decision.outcomeAffecting, + ); + return freezeConsumed({ + ingressState: "present", + invoker: { + state: "present", + kind: "person", + rawPrincipalRef: present[0].rawPrincipalRef, + }, + assuranceRef: "channel-admission", + decisionCoverage: everyDecisionEnforced ? "enforced" : "attribution-only", + }); +} + +/** Queue the channel decision after its exact identity tuple on the shared audit FIFO. */ +export function recordChannelAdmissionDecision(params: { + contextId: string; + executionId: string; + runId: string; + occurredAt: number; + coverageState: NonNullable; +}): boolean { + const missingEvidence = + params.coverageState === "unknown" + ? ["channel.admission_evidence"] + : params.coverageState === "unsupported" + ? ["channel.adapter_identity"] + : params.coverageState === "attribution-only" + ? ["decision.participant_effect"] + : []; + return ( + state.decisionSink?.({ + schemaVersion: 1, + receiptId: `${params.contextId}:channel-admission`, + contextId: params.contextId, + executionId: params.executionId, + runId: params.runId, + occurredAt: params.occurredAt, + action: { + family: "channel", + operation: "admission", + summary: "Channel ingress admitted this agent execution.", + }, + decision: { + outcome: + params.coverageState === "unknown" || params.coverageState === "unsupported" + ? "unknown" + : "allowed", + reasonCode: + params.coverageState === "enforced" + ? "channel_ingress_participant_enforced" + : params.coverageState === "attribution-only" + ? "channel_ingress_attribution_only" + : params.coverageState === "unsupported" + ? "channel_ingress_identity_unsupported" + : "channel_ingress_identity_unknown", + }, + enforcement: { + coverageState: params.coverageState, + evaluatorRef: "channel-ingress", + policyRefs: [], + grantRefs: [], + contextFieldsUsed: params.coverageState === "enforced" ? ["invoker.principal"] : [], + }, + source: { + owner: "channel-ingress", + recordRef: `${params.contextId}:channel-admission`, + decisionBoundary: "channel-ingress.run-admission", + }, + missingEvidence, + remediation: + params.coverageState === "enforced" + ? [] + : [ + { + code: "treat_as_diagnostic_provenance", + text: "Treat this receipt as diagnostic provenance, not authorization.", + }, + ], + }) ?? false + ); +} diff --git a/src/channels/message-access/index.ts b/src/channels/message-access/index.ts index bd8d7a594420..43868c52e461 100644 --- a/src/channels/message-access/index.ts +++ b/src/channels/message-access/index.ts @@ -1,6 +1,7 @@ // Public channel ingress/message-access barrel. Keep this as the narrow import // point for callers that need access decisions without plugin internals. export { defineStableChannelIngressIdentity } from "./runtime-identity.js"; +export { createChannelParticipantAdmissionEvidence } from "./admission-evidence.js"; export { channelIngressRoutes, createChannelIngressResolver, diff --git a/src/channels/message-access/runtime.ts b/src/channels/message-access/runtime.ts index d5e27ad63fdc..003c085f601f 100644 --- a/src/channels/message-access/runtime.ts +++ b/src/channels/message-access/runtime.ts @@ -8,6 +8,7 @@ import { uniqueStrings, } from "@openclaw/normalization-core/string-normalization"; import type { PairingChannel } from "../../pairing/pairing-store.types.js"; +import { bindChannelIngressAdmissionEvidence } from "./admission-evidence.js"; import { decideChannelIngress } from "./decision.js"; import { resolveChannelIngressEffectiveAllowFromLists } from "./effective-allow-from.js"; import { @@ -660,7 +661,7 @@ export async function resolveChannelMessageIngress( const routeAccess = projectRouteAccess({ ingress, route: params.route }); const commandAccess = projectCommandAccess({ ingress, policy }); const activationAccess = projectActivationAccess({ ingress }); - return { + const result: ResolvedChannelMessageIngress = { state, ingress, senderAccess, @@ -668,4 +669,17 @@ export async function resolveChannelMessageIngress( commandAccess, activationAccess, }; + return bindChannelIngressAdmissionEvidence({ + result, + channelId, + accountId: params.accountId, + rawPrincipalRef: params.subject.stableId, + participantOutcomeAffecting: + senderAccess.gate?.match?.matched === true && + (senderAccess.reasonCode === "dm_policy_allowlisted" || + senderAccess.reasonCode === "group_policy_allowed") && + !(isGroup + ? state.allowlists.group.hasWildcard + : state.allowlists.dm.hasWildcard || state.allowlists.pairingStore.hasWildcard), + }); } diff --git a/src/channels/turn/lifecycle.ts b/src/channels/turn/lifecycle.ts index 520e2cd2e677..80d5c7d47f0d 100644 --- a/src/channels/turn/lifecycle.ts +++ b/src/channels/turn/lifecycle.ts @@ -21,6 +21,7 @@ import { settlePendingFinalDelivery } from "../../infra/outbound/delivery-comple import { createMessageSentEmitter } from "../../infra/outbound/message-sent-hook.js"; import { summarizeOutboundPayloadForTransport } from "../../infra/outbound/payloads.js"; import { getGlobalHookRunner } from "../../plugins/hook-runner-global.js"; +import { copyChannelContextAdmissionEvidence } from "../message-access/admission-evidence.js"; import { resolveMessageReceiptPrimaryId } from "../message/receipt.js"; import { createChannelReplyPipeline } from "../message/reply-pipeline.js"; import { recordInboundSession } from "../session.js"; @@ -64,6 +65,20 @@ type RoutedAssembledChannelTurn = Omit< type DispatchableChannelTurn = AssembledChannelTurn | RoutedAssembledChannelTurn; type AnyChannelDeliveryAdapter = ChannelEventDeliveryAdapter | ChannelTurnDeliveryAdapter; +function applyRouteDmScope( + context: T, + dmScope: string | undefined, +): T { + if (!dmScope || context.DmScope === dmScope) { + return context; + } + const scoped = { ...context, DmScope: dmScope } as T; + // Finalized contexts carry identity evidence out-of-band; keep it attached + // when routing must replace the object to add the authoritative DM scope. + copyChannelContextAdmissionEvidence(context, scoped); + return scoped; +} + type PendingChannelDeliveryAttempt = { payload: ReplyPayload; info: ChannelDeliveryInfo; @@ -90,7 +105,7 @@ export function assembleResolvedChannelTurn< const { cfg, route, ...turn } = value; return { ...turn, - ctxPayload: route.dmScope ? { ...turn.ctxPayload, DmScope: route.dmScope } : turn.ctxPayload, + ctxPayload: applyRouteDmScope(turn.ctxPayload, route.dmScope), routeSessionKey: route.sessionKey, storePath: resolveSessionStorePathCore(cfg.session?.store, { agentId: route.agentId }), recordInboundSession, @@ -99,7 +114,7 @@ export function assembleResolvedChannelTurn< const { cfg, route, ...turn } = value; const assembled: RoutedAssembledChannelTurn = { ...turn, - ctxPayload: route.dmScope ? { ...turn.ctxPayload, DmScope: route.dmScope } : turn.ctxPayload, + ctxPayload: applyRouteDmScope(turn.ctxPayload, route.dmScope), cfg, agentId: route.agentId, routeSessionKey: route.sessionKey, diff --git a/src/channels/turn/run-channel-turn.finalize.test.ts b/src/channels/turn/run-channel-turn.finalize.test.ts index 75c7458497f8..dc488e4b9c17 100644 --- a/src/channels/turn/run-channel-turn.finalize.test.ts +++ b/src/channels/turn/run-channel-turn.finalize.test.ts @@ -6,6 +6,12 @@ import type { FinalizedMsgContext } from "../../auto-reply/templating.js"; import type { OpenClawConfig } from "../../config/types.openclaw.js"; import { resetDiagnosticEventsForTest } from "../../infra/diagnostic-events.js"; import { resetLogger, setLoggerOverride } from "../../logging/logger.js"; +import { + bindChannelContextAdmissionEvidence, + configureChannelAdmissionEvidenceCollection, + createChannelParticipantAdmissionEvidence, + readChannelContextAdmissionEvidence, +} from "../message-access/admission-evidence.js"; import { outboundMessageIdentities } from "../message/outbound-echo-state.js"; import { recordOutboundMessageIdentity } from "../message/outbound-echo.js"; import type { RecordInboundSession } from "../session.types.js"; @@ -767,6 +773,50 @@ describe("channel turn finalize", () => { expect(finalizedResult.routeSessionKey).toBe("agent:observer:test:peer"); }); + it("preserves private channel admission evidence when routing adds a DM scope", async () => { + const clearCollection = configureChannelAdmissionEvidenceCollection(true); + try { + const ctx = createCtx(); + const evidence = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); + bindChannelContextAdmissionEvidence({ + context: ctx, + channelId: "test", + evidence, + rawPrincipalRef: "person-1", + }); + dispatchReplyWithRoutedChannelDispatcherCore.mockImplementation(createDispatch()); + + await runChannelTurn({ + channel: "test", + raw: {}, + adapter: { + ingest: () => ({ id: "msg-evidence", rawText: "hello" }), + resolveTurn: () => ({ + cfg, + channel: "test", + route: { + agentId: "main", + dmScope: "per-channel-peer", + sessionKey: "agent:main:test:peer", + }, + ctxPayload: ctx, + delivery: { deliver: vi.fn() }, + record: { onRecordError: vi.fn() }, + }), + }, + }); + + const dispatched = dispatchReplyWithRoutedChannelDispatcherCore.mock.calls[0]?.[0]; + expect(dispatched?.ctx.DmScope).toBe("per-channel-peer"); + expect(readChannelContextAdmissionEvidence(dispatched?.ctx ?? {})).toBe(evidence); + } finally { + clearCollection(); + } + }); + it("finalizes failed dispatches before rethrowing", async () => { const onFinalize = vi.fn(); const dispatchError = new Error("dispatch failed"); diff --git a/src/config/sessions/session-accessor.sqlite-entry.ts b/src/config/sessions/session-accessor.sqlite-entry.ts index 01cd00271e66..6954fac585d0 100644 --- a/src/config/sessions/session-accessor.sqlite-entry.ts +++ b/src/config/sessions/session-accessor.sqlite-entry.ts @@ -675,12 +675,12 @@ export async function recordInboundSessionMeta(params: { if (context.existingEntry) { return metadataPatch; } - const senderId = params.ctx.From?.trim(); + const senderId = params.ctx.SenderId?.trim(); return { ...buildSessionCreationStamp( params.ctx.SessionCreation ?? { via: "channel", - actor: { type: "human", ...(senderId ? { id: senderId } : {}) }, + ...(senderId ? { actor: { type: "human", id: senderId } } : {}), }, ), ...metadataPatch, @@ -728,14 +728,12 @@ export async function updateSessionLastRoute(params: { if (context.existingEntry) { return routePatch; } - const senderId = params.ctx?.From?.trim(); + const senderId = params.ctx?.SenderId?.trim(); return { ...buildSessionCreationStamp( params.ctx?.SessionCreation ?? { via: "channel", - ...(params.ctx - ? { actor: { type: "human" as const, ...(senderId ? { id: senderId } : {}) } } - : {}), + ...(senderId ? { actor: { type: "human" as const, id: senderId } } : {}), }, ), ...routePatch, diff --git a/src/config/sessions/session-accessor.test.ts b/src/config/sessions/session-accessor.test.ts index 6e16e1e6f982..5a6de0eddc13 100644 --- a/src/config/sessions/session-accessor.test.ts +++ b/src/config/sessions/session-accessor.test.ts @@ -669,6 +669,7 @@ describe("session accessor seam", () => { Surface: "webchat", ChatType: "direct", From: "webchat:user-1", + SenderId: "webchat:user-1", To: "webchat:agent", SessionKey: sessionKey, OriginatingTo: "webchat:user-1", @@ -691,7 +692,7 @@ describe("session accessor seam", () => { await recordInboundSessionMeta({ storePath, sessionKey, - ctx: { ...ctx, From: "webchat:different-sender" }, + ctx: { ...ctx, From: "webchat:different-route", SenderId: "webchat:different-sender" }, }); expect(loadSessionEntry({ sessionKey, storePath })).toMatchObject(creationStamp); @@ -719,6 +720,23 @@ describe("session accessor seam", () => { createdActor: { type: "human", id: "profile-ada" }, createdAt: expect.any(Number), }); + + const senderlessKey = "agent:main:webchat:dm:senderless"; + const senderless = await recordInboundSessionMeta({ + storePath, + sessionKey: senderlessKey, + ctx: { + ...ctx, + SessionKey: senderlessKey, + From: "webchat:room-route", + SenderId: undefined, + }, + }); + expect(senderless).toMatchObject({ + createdVia: "channel", + createdAt: expect.any(Number), + }); + expect(senderless?.createdActor).toBeUndefined(); }); it("does not create sessions when inbound meta recording opts out of upsert", async () => { @@ -788,6 +806,32 @@ describe("session accessor seam", () => { expect(loadSessionEntry({ sessionKey, storePath })).toBeUndefined(); }); + it("stamps last-route creation from the participant, never the conversation route", async () => { + const participantKey = "agent:main:webchat:dm:route-participant"; + const participant = await updateSessionLastRoute({ + storePath, + sessionKey: participantKey, + channel: "webchat", + to: "webchat:room-1", + ctx: { From: "webchat:room-1", SenderId: "webchat:person-1" }, + }); + expect(participant).toMatchObject({ + createdVia: "channel", + createdActor: { type: "human", id: "webchat:person-1" }, + }); + + const senderlessKey = "agent:main:webchat:dm:route-senderless"; + const senderless = await updateSessionLastRoute({ + storePath, + sessionKey: senderlessKey, + channel: "webchat", + to: "webchat:room-2", + ctx: { From: "webchat:room-2" }, + }); + expect(senderless?.createdVia).toBe("channel"); + expect(senderless?.createdActor).toBeUndefined(); + }); + it("rejects alias targets and keeps canonical lifecycle mutations explicit", async () => { await replaceSessionEntry( { sessionKey: "agent:main:work", storePath }, diff --git a/src/gateway/server-runtime-subscriptions.test.ts b/src/gateway/server-runtime-subscriptions.test.ts index cff940ff28bf..77f37ebb56bb 100644 --- a/src/gateway/server-runtime-subscriptions.test.ts +++ b/src/gateway/server-runtime-subscriptions.test.ts @@ -5,6 +5,10 @@ import { enqueueExecutionIdentityContextAtAdmission, hasExecutionIdentityAdmissionSink, } from "../audit/execution-identity-admission.js"; +import { + consumeChannelAdmissionEvidence, + createChannelParticipantAdmissionEvidence, +} from "../channels/message-access/admission-evidence.js"; import type { CronServiceState } from "../cron/service/state.js"; import { tryFinishCronTaskRunWithoutHistory } from "../cron/service/task-runs.js"; import { @@ -68,6 +72,8 @@ const auditTestState = vi.hoisted(() => ({ created: 0, recorded: 0, identityRecorded: 0, + decisionRecorded: 0, + executionIdentityEnabled: false, stopped: 0, })); const agentEventHandlerMocks = vi.hoisted(() => ({ @@ -80,6 +86,7 @@ const transcriptBroadcastMocks = vi.hoisted(() => ({ vi.mock("../audit/audit-config.js", () => ({ isAuditLedgerEnabled: () => auditTestState.enabled, + isExecutionIdentityCollectionEnabled: () => auditTestState.executionIdentityEnabled, resolveAuditMessageMode: () => auditTestState.messageMode, })); @@ -96,6 +103,10 @@ vi.mock("../audit/audit-recorder.js", () => ({ auditTestState.identityRecorded += 1; return true; }), + recordExecutionDecision: vi.fn(() => { + auditTestState.decisionRecorded += 1; + return true; + }), stop: vi.fn(async () => { auditTestState.stopped += 1; }), @@ -169,6 +180,8 @@ describe("startGatewayEventSubscriptions", () => { auditTestState.created = 0; auditTestState.recorded = 0; auditTestState.identityRecorded = 0; + auditTestState.decisionRecorded = 0; + auditTestState.executionIdentityEnabled = false; auditTestState.stopped = 0; transcriptBroadcastMocks.useActualHandler = false; transcriptBroadcastMocks.readMessageCount.mockReset(); @@ -218,6 +231,26 @@ describe("startGatewayEventSubscriptions", () => { expect(hasExecutionIdentityAdmissionSink()).toBe(false); }); + it("owns channel evidence collection for the configured gateway lifecycle", async () => { + auditTestState.executionIdentityEnabled = true; + unsubs = startGatewayEventSubscriptions(createParams()); + + const evidence = createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-1", + }); + expect(evidence).toBeDefined(); + + await unsubs.agentUnsub(); + expect(consumeChannelAdmissionEvidence(evidence)).toMatchObject({ ingressState: "unknown" }); + expect( + createChannelParticipantAdmissionEvidence({ + channelId: "test", + participantId: "person-2", + }), + ).toBeUndefined(); + }); + it("keeps retention maintenance but creates no producers when audit.enabled is false", async () => { auditTestState.enabled = false; unsubs = startGatewayEventSubscriptions(createParams()); diff --git a/src/gateway/server-runtime-subscriptions.ts b/src/gateway/server-runtime-subscriptions.ts index 16bb3029f525..e2dcd5686139 100644 --- a/src/gateway/server-runtime-subscriptions.ts +++ b/src/gateway/server-runtime-subscriptions.ts @@ -1,9 +1,17 @@ // Gateway event subscription wiring for agent, heartbeat, transcript, and lifecycle broadcasts. import { resolveDefaultAgentId } from "../agents/agent-scope.js"; -import { isAuditLedgerEnabled, resolveAuditMessageMode } from "../audit/audit-config.js"; +import { + isAuditLedgerEnabled, + isExecutionIdentityCollectionEnabled, + resolveAuditMessageMode, +} from "../audit/audit-config.js"; import { createAuditEventRecorder } from "../audit/audit-recorder.js"; import { configureExecutionIdentityAdmissionSink } from "../audit/execution-identity-admission.js"; import { onTrustedMessageAuditEvent } from "../audit/message-audit-events.js"; +import { + configureChannelAdmissionDecisionSink, + configureChannelAdmissionEvidenceCollection, +} from "../channels/message-access/admission-evidence.js"; import { getRuntimeConfig } from "../config/io.js"; import { onAgentAuditEvent, onAgentRuntimeEvent } from "../infra/agent-events.js"; import { clearAgentRunContext } from "../infra/agent-run-registry.js"; @@ -90,6 +98,12 @@ export function startGatewayEventSubscriptions(params: { const clearExecutionIdentityAdmissionSink = configureExecutionIdentityAdmissionSink( auditRecorder.recordExecutionIdentity, ); + const clearChannelAdmissionEvidenceCollection = configureChannelAdmissionEvidenceCollection( + isExecutionIdentityCollectionEnabled(runtimeConfig), + ); + const clearChannelAdmissionDecisionSink = configureChannelAdmissionDecisionSink( + auditRecorder.recordExecutionDecision, + ); const sessionObserver = createSessionObserver({ getConfig: getRuntimeConfig, subscribers: params.sessionMessageSubscribers, @@ -348,6 +362,8 @@ export function startGatewayEventSubscriptions(params: { unsubscribeToolAuditEvents?.(); unsubscribeMessageAuditEvents?.(); clearExecutionIdentityAdmissionSink(); + clearChannelAdmissionEvidenceCollection(); + clearChannelAdmissionDecisionSink(); await agentEventHandlerLoader .peek() ?.then((handler) => handler.dispose()) diff --git a/src/plugin-sdk/channel-ingress-runtime.ts b/src/plugin-sdk/channel-ingress-runtime.ts index 3d958ac97912..f2f8c60eaf91 100644 --- a/src/plugin-sdk/channel-ingress-runtime.ts +++ b/src/plugin-sdk/channel-ingress-runtime.ts @@ -8,6 +8,7 @@ */ export { channelIngressRoutes, + createChannelParticipantAdmissionEvidence, createChannelIngressResolver, defineStableChannelIngressIdentity, readChannelIngressStoreAllowFromForDmPolicy, diff --git a/test/scripts/telegram-user-crabbox-proof.test.ts b/test/scripts/telegram-user-crabbox-proof.test.ts index 4066e327f87f..17ab6e4a4f2c 100644 --- a/test/scripts/telegram-user-crabbox-proof.test.ts +++ b/test/scripts/telegram-user-crabbox-proof.test.ts @@ -12,10 +12,12 @@ import { COMMAND_TIMEOUT_MS, createContainerizedSutSpawnSpec, createCrabboxWarmupArgs, + createOpenClawCliSpawnSpec, createOpenClawGatewaySpawnSpec, parseArgs, processTargetExists, readCodexProxyPort, + readLogAfterOffset, readLogTail, readTelegramUserProofLogTailBytes, recordProbeVideo, @@ -32,6 +34,7 @@ import { stageFullSessionArtifacts, startLocalSut, waitForLog, + waitForLogAfterOffset, writeSutConfig, } from "../../scripts/e2e/telegram-user-crabbox-proof.ts"; import { resolveWindowsTaskkillPath } from "../../scripts/lib/windows-taskkill.mjs"; @@ -143,6 +146,19 @@ describe("telegram user Crabbox proof log polling", () => { expect(spec.options.shell).toBe(false); }); + it("runs held-session audit inspection through the same pinned repo CLI", () => { + const spec = createOpenClawCliSpawnSpec({ + args: ["audit", "--run", "run-1", "--explain", "--json"], + env: { OPENCLAW_CONFIG_PATH: "/tmp/openclaw.json" }, + pnpmExecPath: "/opt/mantis-toolchain/pnpm", + repoRoot: "/repo", + }); + + expect(spec.command).toBe("/opt/mantis-toolchain/pnpm"); + expect(spec.args).toEqual(["openclaw", "audit", "--run", "run-1", "--explain", "--json"]); + expect(spec.options.env?.OPENCLAW_CONFIG_PATH).toBe("/tmp/openclaw.json"); + }); + it("routes fork SUT startup through the root-owned validating wrapper", () => { const repoRoot = makeTempDir(tempDirs, "openclaw-telegram-proof-"); const runtimeRoot = makeTempDir(tempDirs, "openclaw-telegram-proof-"); @@ -370,6 +386,19 @@ describe("telegram user Crabbox proof log polling", () => { expect(parseArgs(["--text", "-ping"]).text).toBe("-ping"); }); + it("requires held sessions for identity inspection and lifecycle restart", () => { + expect(() => parseArgs(["inspect"])).toThrow("inspect requires --session"); + expect(() => parseArgs(["restart"])).toThrow("restart requires --session"); + expect(parseArgs(["inspect", "--session", "session.json"]).command).toBe("inspect"); + expect(parseArgs(["restart", "--session", "session.json"]).command).toBe("restart"); + expect( + parseArgs(["send", "--session", "session.json", "--chat", "@sut", "--text", "hello"]).chat, + ).toBe("@sut"); + expect(() => parseArgs(["inspect", "--session", "session.json", "--chat", "@sut"])).toThrow( + "--chat is available only for held-session sends", + ); + }); + it("accepts an explicit Telegram link-preview setting", () => { expect(parseArgs(["start", "--link-preview", "false"]).linkPreview).toBe(false); expect(parseArgs(["start", "--link-preview", "true"]).linkPreview).toBe(true); @@ -486,6 +515,24 @@ describe("telegram user Crabbox proof log polling", () => { expect(JSON.stringify(config)).not.toContain("resource-ok"); }); + it("enables execution identity before Telegram Gateway startup", () => { + const configRoot = writeSutConfig({ + gatewayPort: 19042, + groupId: "group", + mockPort: 19043, + outputDir: makeTempDir(tempDirs, "openclaw-telegram-proof-"), + testerId: "tester", + }); + tempDirs.push(configRoot.tempRoot); + + const config = JSON.parse(fs.readFileSync(configRoot.configPath, "utf8")); + expect(config.logging.audit).toMatchObject({ + enabled: true, + executionIdentity: true, + messages: "direct", + }); + }); + it("injects the requested Telegram link-preview setting before startup", () => { const disabledConfigRoot = writeSutConfig({ gatewayPort: 19042, @@ -691,6 +738,25 @@ describe("telegram user Crabbox proof log polling", () => { expect(tail).not.toContain("old\nold\nold\nold\nold\nold\nold\nold\nold"); }); + it("observes restart readiness only after the lifecycle log boundary", async () => { + const logPath = path.join(makeTempDir(tempDirs, "openclaw-telegram-proof-"), "gateway.log"); + fs.writeFileSync(logPath, "[gateway] ready\n", "utf8"); + const offset = fs.statSync(logPath).size; + expect(readLogAfterOffset(logPath, offset)).toBe(""); + + fs.appendFileSync(logPath, "received SIGUSR1; restarting\ngateway ready\n", "utf8"); + + await expect( + waitForLogAfterOffset({ + label: "restart", + logPath, + offset, + pattern: /received SIGUSR1; restarting/u, + timeoutMs: 100, + }), + ).resolves.toContain("gateway ready"); + }); + it("keeps byte-cut log tails UTF-8 safe and reads at least one byte", () => { const logPath = path.join(makeTempDir(tempDirs, "openclaw-telegram-proof-"), "gateway.log"); fs.writeFileSync( @@ -854,6 +920,7 @@ fs.writeFileSync(process.env.OPENCLAW_TEST_ARGV_PATH, JSON.stringify(process.arg writeExecutable( scriptPath, renderRemoteProbe({ + chat: "@proof-bot", expect: [payload], sutUsername: payload, text: payload, @@ -874,6 +941,7 @@ fs.writeFileSync(process.env.OPENCLAW_TEST_ARGV_PATH, JSON.stringify(process.arg expect(result.status).toBe(0); expect(fs.existsSync(injectedPath)).toBe(false); expect(JSON.parse(fs.readFileSync(argvPath, "utf8"))).toContain(payload); + expect(JSON.parse(fs.readFileSync(argvPath, "utf8"))).toContain("@proof-bot"); }); it("clamps oversized command timeouts before arming timers", async () => {