feat: audit admitted channel participant identity

This commit is contained in:
joshavant
2026-08-12 18:26:35 -05:00
parent 7b73d33b60
commit bf041d2e3c
64 changed files with 2606 additions and 87 deletions
@@ -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"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"108af57a4b4bc57c70912fc40319876bb6ec994cce59d2b9c991980e18c4f80f","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
{"contentHash":"ddb50f1e611d363742cf493741fad27f04884f318b1f68e3bc43e1a63ccca4cc","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"dd255204bf11dccfe9d866e7a6725935bfead7378a52ddfa6b60b9a58c0cdf76","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
{"contentHash":"e64d40bbc2f99f5834b46918359f1680ca40e22caac6642007fb7eb450dcab67","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
@@ -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"}
@@ -1 +1 @@
{"contentHash":"f0be441226a760f75dc1634f93d7f28586a0b18f690e2dd0f87c5360f2cb7f03","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"}
{"contentHash":"2cf8445ec79b755119c0dd1efd24e4f06a0cb5f79e0f661619013a4c1961f7f8","entrypoint":"channel-inbound","importSpecifier":"openclaw/plugin-sdk/channel-inbound"}
@@ -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"}
@@ -1 +1 @@
{"contentHash":"8bf7d1fd1e21861fd472365e0ab10ca811800e1b0cb9e6aac0e4d6a41613a9a2","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
{"contentHash":"3bb24706f7e76c556192dd100a856a490220490248892368bf05d1a6b3a67615","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
@@ -1 +1 @@
{"contentHash":"d430c3def27acb48a5356dafac17b89b0cbfcb1d3f959136b8c7ae82586d87df","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
{"contentHash":"7ba02e716d0039ef36ebbc74c524e818619714e6c7983c9cdf34e325a7d5fa95","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
@@ -1 +1 @@
{"contentHash":"0d9fc3f9ef06aaa9995c5a670d62b99cfce25f0b2c42152aa33214f9dec78b19","entrypoint":"channel-pairing","importSpecifier":"openclaw/plugin-sdk/channel-pairing"}
{"contentHash":"8cf3a92d14f57bcd63a09f704a5b9dabb48a83d58df22dead6b6e859e4ba69cf","entrypoint":"channel-pairing","importSpecifier":"openclaw/plugin-sdk/channel-pairing"}
@@ -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"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"4c667b7e88cd752c294bcc731bc157c1ccb2a7aa4afc58f4013409191f3a9b28","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
{"contentHash":"9d6ee27bf084b6e4383bf6a4da070f1f2c0722b244776bd3c511703038df3714","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"027b2280171324b628865fb4cdf9f72c13c059385a2fbec5b650c69d66beecdc","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
{"contentHash":"b8a7b3fb6c8d00dd5e944bd569ee22cf2ac994a453297fd1304d0443fd52bb86","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
@@ -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"}
@@ -1 +1 @@
{"contentHash":"69fe3479e42b74771fb86d228264e8e19a3a16e1a07df4f89499111c2df2bc13","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
{"contentHash":"a51fed7fe588e420896c747abc547a9d28c216ba682274bae683c308644d51d7","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"f058be9cbffbcb6e1dafe69cbc701c7cc06d1ffc90323256685e31b111f91193","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
{"contentHash":"c58b25311e4cbb436b35e084789748942fb12ebdf97a8ea8c085e5fc110ddf89","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"e7ac61e29d41e43a4c780831a8aa0e5bed10e47630a36879f2555216eb1cce54","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
{"contentHash":"8f95e88158576bbef982b4962b917c2b937e196a617a3b6a76cecef88a029921","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
@@ -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"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"b965caa74f02ea0fa0b4b5a57bfacd0699df991452b5aa95f21237858930508f","entrypoint":"runtime-store","importSpecifier":"openclaw/plugin-sdk/runtime-store"}
{"contentHash":"d0af7cc6a491a50a90271814d47b91116eb7e600945ead8d3f54bfce6075bbdd","entrypoint":"runtime-store","importSpecifier":"openclaw/plugin-sdk/runtime-store"}
@@ -1 +1 @@
{"contentHash":"1c76ac2e6a40a56ae0dbf5f6e2b5d911883f8e11cc7c0f435f6cce0c185a60c2","entrypoint":"session-catalog","importSpecifier":"openclaw/plugin-sdk/session-catalog"}
{"contentHash":"2e56a2b98213e0f297084d33cce8880d9dfe931ae74b02f753c988fb82ec088d","entrypoint":"session-catalog","importSpecifier":"openclaw/plugin-sdk/session-catalog"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"f76d5908b66bb49851cccc530dc63fc5f3352b327963f1091a0f34fe9ccc693b","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
{"contentHash":"db4a41272079b955d5426e2bea3f1db87e39d4168abd1109ef9ec937a88d89b6","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
@@ -1 +1 @@
{"contentHash":"e9dd8d432098b26f2e102fe41edee7c556558ea9b26a0077da72b0967d6649d2","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
{"contentHash":"0488f4915160e0b0ea4e574d6800ced82e88219e9a60aca0de5599879d0cbec4","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
+9
View File
@@ -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):
+9
View File
@@ -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
+7
View File
@@ -1121,6 +1121,13 @@ Seed assets live in `qa/`:
- `qa/scenarios/index.yaml`
- `qa/scenarios/<theme>/*.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.
+17 -4
View File
@@ -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
+19
View File
@@ -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:<name>` entries stay redacted. Core resolves static
+7
View File
@@ -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`
+1
View File
@@ -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) },
@@ -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 });
}
});
});
@@ -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<QaSuiteRuntimeEnv, "gateway">): {
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();
}
}
@@ -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"),
@@ -67,6 +67,7 @@ export type QaScenarioRuntimeDeps = {
assertNoGatewayLogSentinels: QaScenarioRuntimeFunction;
readSessionTranscriptSummary: QaScenarioRuntimeFunction;
runQaCli: QaScenarioRuntimeFunction;
inspectQaExecutionIdentityStorage: QaScenarioRuntimeFunction;
extractMediaPathFromText: QaScenarioRuntimeFunction;
resolveGeneratedImagePath: QaScenarioRuntimeFunction;
startAgentRun: QaScenarioRuntimeFunction;
@@ -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,
@@ -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),
@@ -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}`"
+369 -2
View File
@@ -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 <url>]",
" node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts send --session <session.json> --text <text>",
" node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts inspect --session <session.json>",
" node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts restart --session <session.json>",
" node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts run --session <session.json> -- <remote command>",
" node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts view --session <session.json> --message-id <id>",
" node --import tsx scripts/e2e/telegram-user-crabbox-proof.ts screenshot --session <session.json>",
@@ -223,6 +229,7 @@ function usageText() {
"",
"Useful options:",
" --class <name> Crabbox machine class. Default: standard.",
" --chat <id|username> Telegram chat override for send (for example @bot for DM).",
" --desktop-chat-title <name> Telegram Desktop chat to select before recording.",
" --human-delay-fixed-ms <ms> Set a fixed custom human delay before Gateway startup.",
" --id <cbx_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<CommandResult>((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 <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<string>();
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;
+4 -2
View File
@@ -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(
+2
View File
@@ -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;
@@ -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<unknown> };
}
).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<unknown> };
}
).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" }],
+8 -16
View File
@@ -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(
@@ -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();
}
});
});
@@ -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<ExecutionIdentityAdmissionFacts, "invoker" | "assurance">;
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();
},
});
}
+62
View File
@@ -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 }) => {
+15 -1
View File
@@ -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<string> {
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.",
+68 -7
View File
@@ -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");
+12 -1
View File
@@ -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,
@@ -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,
+124
View File
@@ -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);
@@ -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<string>();
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);
});
});
+64 -19
View File
@@ -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;
}
+3
View File
@@ -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. */
+22 -3
View File
@@ -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<string, unknown>;
/** 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<typeof context>) => {
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);
}
@@ -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();
}
});
});
@@ -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<object, ChannelAdmissionEvidencePayload>(),
evidenceByIngress: new WeakMap<object, ChannelAdmissionEvidence>(),
evidenceByContext: new WeakMap<object, ChannelAdmissionEvidence>(),
consumedEvidence: new WeakSet<object>(),
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<Extract<ChannelAdmissionEvidencePayload, { kind: "leaf" }>, "createdAt" | "generation">
| Omit<
Extract<ChannelAdmissionEvidencePayload, { kind: "aggregate" }>,
"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<object>;
}): 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<object>;
}): 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<ConsumedChannelAdmissionEvidence, "invoker"> & {
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<ConsumedChannelAdmissionEvidence["decisionCoverage"]>;
}): 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
);
}
+1
View File
@@ -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,
+15 -1
View File
@@ -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),
});
}
+17 -2
View File
@@ -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<T extends AssembledChannelTurn["ctxPayload"]>(
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,
@@ -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");
@@ -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,
+45 -1
View File
@@ -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 },
@@ -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());
+17 -1
View File
@@ -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())
@@ -8,6 +8,7 @@
*/
export {
channelIngressRoutes,
createChannelParticipantAdmissionEvidence,
createChannelIngressResolver,
defineStableChannelIngressIdentity,
readChannelIngressStoreAllowFromForDmPolicy,
@@ -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 () => {