diff --git a/.claude/skills/env-reference/SKILL.md b/.claude/skills/env-reference/SKILL.md index b591e6e48..39a152749 100644 --- a/.claude/skills/env-reference/SKILL.md +++ b/.claude/skills/env-reference/SKILL.md @@ -116,6 +116,15 @@ See `apps/api/.env.example` for the full list. Key variables: - `ACP_ACTIVITY_BINDING_CACHE_TTL_MS` — Short-lived authorized ACP session binding cache used to avoid ProjectData reads during callback storms (default: `30000`) - `ACP_ACTIVITY_BINDING_CACHE_MAX_ENTRIES` — Maximum cached ACP activity bindings retained by one Worker isolate (default: `2048`) +- `ACP_INTERACTIONS_ENABLED` — Dormant durable ACP interaction foundation kill switch. Slice A defaults this to `false`; later slices must intentionally enable producers/consumers (default: `false`) +- `ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS` / `ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS` — Default permission deadlines for conversation and task contexts (defaults: `7200000` / `1800000`) +- `ACP_INTERACTION_MAX_DEADLINE_MS` / `ACP_INTERACTION_DEADLINE_MARGIN_MS` — Absolute deadline ceiling and prompt/runtime cap safety margin (defaults: `14400000` / `60000`) +- `ACP_INTERACTION_MAX_PENDING_PER_SESSION` — Maximum pending durable ACP interactions per chat (default: `8`) +- `ACP_INTERACTION_REQUEST_MAX_BYTES`, `ACP_INTERACTION_OPTIONS_MAX_COUNT`, `ACP_INTERACTION_OPTION_NAME_MAX_CHARS`, `ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES`, `ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES`, `ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM` — Request/detail and schema bounds for encrypted ACP interaction payloads (defaults: `32768`, `16`, `200`, `16384`, `20`, `50`) +- `ACP_INTERACTION_ANSWER_MAX_BYTES` / `ACP_INTERACTION_ANSWER_STRING_MAX_BYTES` — Answer decision and individual answer string bounds before encrypted storage (defaults: `16384` / `4096`) +- `ACP_INTERACTION_RETRY_DELAYS_MS` / `ACP_INTERACTION_RETRY_STEADY_MS` / `ACP_INTERACTION_DELIVERY_WINDOW_MS` — Durable answer outbox retry sequence, steady retry delay, and max post-answer retry window (defaults: `1000,5000,30000,120000,300000`, `300000`, `900000`) +- `ACP_INTERACTION_SENSITIVE_PURGE_MS`, `ACP_INTERACTION_SUMMARY_RETENTION_MS`, `ACP_INTERACTION_SUMMARY_LAST_SETTLED`, `ACP_INTERACTION_SNAPSHOT_LAST_SETTLED` — Sensitive encrypted payload purge, settled summary retention, and bounded snapshot controls (defaults: `3600000`, `2592000000`, `100`, `20`) + Activity coalescing and binding caches are per Worker isolate. Delayed flushes carry their original observed event time, and ProjectData rejects stale writes so a delayed intermediate report cannot overwrite a newer idle/error state from another isolate. - `SESSION_SNAPSHOT_RECOVERY_CLAIM_LEASE_MS` — Reclaim timeout for an interrupted replacement-runtime wake claim (default: `600000`) diff --git a/apps/api/.env.example b/apps/api/.env.example index 6a1987798..34d5ac1d0 100644 --- a/apps/api/.env.example +++ b/apps/api/.env.example @@ -795,6 +795,28 @@ INFOMANIAK_IP_POLL_INTERVAL_MS=3000 # ACP_ACTIVITY_COALESCE_MAX_PENDING=512 # Max pending coalesced activity reports per Worker isolate # ACP_ACTIVITY_BINDING_CACHE_TTL_MS=30000 # Short-lived authorized ACP session binding cache # ACP_ACTIVITY_BINDING_CACHE_MAX_ENTRIES=2048 # Max cached ACP activity bindings per Worker isolate +# Dormant durable ACP interaction foundation. Slice A keeps this disabled; later slices wire runtime/UI producers and consumers. +# ACP_INTERACTIONS_ENABLED=false # Fail-closed feature flag for durable ACP interaction creation +# ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS=7200000 # Default conversation permission deadline (2h) +# ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS=1800000 # Default task permission deadline (30m) +# ACP_INTERACTION_MAX_DEADLINE_MS=14400000 # Absolute max request deadline (4h) +# ACP_INTERACTION_DEADLINE_MARGIN_MS=60000 # Safety margin before prompt/runtime cap +# ACP_INTERACTION_MAX_PENDING_PER_SESSION=8 # Pending durable interactions allowed per chat +# ACP_INTERACTION_REQUEST_MAX_BYTES=32768 # Encrypted request detail cap +# ACP_INTERACTION_OPTIONS_MAX_COUNT=16 # Permission option count cap +# ACP_INTERACTION_OPTION_NAME_MAX_CHARS=200 # Permission option display cap +# ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES=16384 # Encrypted form schema cap +# ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES=20 # Form schema property cap +# ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM=50 # Form enum option cap +# ACP_INTERACTION_ANSWER_MAX_BYTES=16384 # Encrypted answer/decision cap +# ACP_INTERACTION_ANSWER_STRING_MAX_BYTES=4096 # Individual answer string cap +# ACP_INTERACTION_RETRY_DELAYS_MS=1000,5000,30000,120000,300000 # Durable answer outbox retry sequence +# ACP_INTERACTION_RETRY_STEADY_MS=300000 # Steady retry delay after the sequence +# ACP_INTERACTION_DELIVERY_WINDOW_MS=900000 # Max answer delivery retry window after accepted decision +# ACP_INTERACTION_SENSITIVE_PURGE_MS=3600000 # Purge encrypted details/answers after terminal state +# ACP_INTERACTION_SUMMARY_RETENTION_MS=2592000000 # Retain bounded settled summaries for 30d +# ACP_INTERACTION_SUMMARY_LAST_SETTLED=100 # Min settled summaries preserved during retention compaction +# ACP_INTERACTION_SNAPSHOT_LAST_SETTLED=20 # Settled summaries returned per snapshot page # CREDENTIAL_LIMIT_WARNING_PERCENT=75 # Advisory provider credential quota warning threshold # CREDENTIAL_LIMIT_CRITICAL_PERCENT=90 # Advisory provider credential quota critical threshold diff --git a/apps/api/src/durable-objects/interaction-store-model.ts b/apps/api/src/durable-objects/interaction-store-model.ts new file mode 100644 index 000000000..ffc8300b8 --- /dev/null +++ b/apps/api/src/durable-objects/interaction-store-model.ts @@ -0,0 +1,198 @@ +import { + type AcpInteractionAnswerDecision, + type AcpInteractionRuntimeCreate, + type AcpInteractionRuntimeSettle, + type AcpInteractionSafeSummary, +} from '@simple-agent-manager/shared'; + +import { canonicalJson } from '../lib/canonical-json'; +import type { AcpInteractionConfig } from '../services/acp-interaction-config'; + +export type InteractionRow = { + interaction_id: string; + project_id: string; + chat_session_id: string; + agent_session_id: string; + kind: string; + state: string; + generation: string; + runtime_identity: string; + payload_hash: string; + encrypted_detail: string | null; + detail_iv: string | null; + detail_purged_at: number | null; + safe_summary_json: string; + upstream_request_id: string | null; + created_at: number; + updated_at: number; + deadline_at: number; + answered_at: number | null; + answer_key: string | null; + answer_body_hash: string | null; + decision_kind: string | null; + decision_hash: string | null; + encrypted_answer: string | null; + answer_iv: string | null; + encrypted_decision: string | null; + decision_iv: string | null; + delivery_state: string | null; + delivery_attempts: number; + delivery_deadline_at: number | null; + last_delivery_error: string | null; + attention_marker_id: string | null; + attention_projection_state: string | null; + terminal_at: number | null; + purge_at: number | null; +}; + +export interface InteractionStoreCreateInput extends AcpInteractionRuntimeCreate { + projectId: string; + chatSessionId: string; +} + +export interface InteractionStoreAnswerInput { + projectId: string; + chatSessionId: string; + interactionId: string; + answerKey: string; + answerBodyHash: string; + decision: AcpInteractionAnswerDecision; +} + +export interface InteractionStoreSettleInput extends AcpInteractionRuntimeSettle { + projectId: string; + chatSessionId: string; +} + +export type InteractionStoreCreateResult = + | { status: 'created' | 'existing'; summary: AcpInteractionSafeSummary } + | { + status: 'disabled' | 'conflict' | 'too_many_pending' | 'expired' | 'invalid'; + reason: string; + }; + +export type InteractionStoreAnswerResult = + | { + status: 'answered' | 'already_answered'; + summary: AcpInteractionSafeSummary; + delivery: { generation: string; runtimeIdentity: string }; + } + | { + status: 'not_found' | 'stale' | 'conflict' | 'answer_key_conflict' | 'payload_too_large'; + reason: string; + }; + +export interface InteractionStoreSnapshot { + pending: AcpInteractionSafeSummary[]; + settled: AcpInteractionSafeSummary[]; + cursor: string | null; +} + +export function nowMs(): number { + return Date.now(); +} + +export async function sha256(value: string): Promise { + const bytes = await crypto.subtle.digest('SHA-256', new TextEncoder().encode(value)); + return [...new Uint8Array(bytes)].map((byte) => byte.toString(16).padStart(2, '0')).join(''); +} + +export function parseSummary(row: InteractionRow): AcpInteractionSafeSummary { + const safeSummary = JSON.parse(row.safe_summary_json) as unknown; + return { + interactionId: row.interaction_id, + kind: row.kind as AcpInteractionSafeSummary['kind'], + state: row.state as AcpInteractionSafeSummary['state'], + createdAt: row.created_at, + updatedAt: row.updated_at, + deadlineAt: row.deadline_at, + answeredAt: row.answered_at, + deliveryState: row.delivery_state as AcpInteractionSafeSummary['deliveryState'], + attentionMarkerId: row.attention_marker_id, + toolCallId: + typeof safeSummary === 'object' && + safeSummary !== null && + 'toolCallId' in safeSummary && + typeof safeSummary.toolCallId === 'string' + ? safeSummary.toolCallId + : null, + }; +} + +export function terminalState(state: string): boolean { + return [ + 'delivery_confirmed', + 'delivery_unconfirmed', + 'interrupted', + 'expired', + 'cancelled', + ].includes(state); +} + +export function detailBoundsViolation( + detail: unknown, + config: AcpInteractionConfig +): string | null { + return detailRecordBoundsViolation(asRecord(detail), config); +} + +function detailRecordBoundsViolation( + record: Record | null, + config: AcpInteractionConfig +): string | null { + if (!record) return null; + const optionsViolation = optionsBoundsViolation(record, config); + if (optionsViolation) return optionsViolation; + return schemaBoundsViolation(record, config); +} + +function asRecord(value: unknown): Record | null { + if (!value || typeof value !== 'object') return null; + if (Array.isArray(value)) return null; + return value as Record; +} + +function optionsBoundsViolation( + record: Record, + config: AcpInteractionConfig +): string | null { + if (Array.isArray(record.options) && record.options.length > config.optionsMaxCount) { + return 'request options exceed configured maximum'; + } + return null; +} + +function schemaBoundsViolation( + record: Record, + config: AcpInteractionConfig +): string | null { + const schema = record.schema ?? record.formSchema; + if (schema === undefined) return null; + const schemaJson = canonicalJson(schema); + if (new TextEncoder().encode(schemaJson).byteLength > config.formSchemaMaxBytes) { + return 'form schema exceeds configured maximum'; + } + return schemaObjectBoundsViolation(schema, config); +} + +function schemaObjectBoundsViolation( + schema: unknown, + config: AcpInteractionConfig +): string | null { + const schemaRecord = asRecord(schema); + if (!schemaRecord) return null; + const properties = asRecord(schemaRecord.properties); + if (properties && Object.keys(properties).length > config.formSchemaMaxProperties) { + return 'form schema properties exceed configured maximum'; + } + if (hasEnumOverflow(schema, config.formSchemaMaxEnum)) return 'form schema enum exceeds configured maximum'; + return null; +} + +function hasEnumOverflow(value: unknown, maxEnum: number): boolean { + if (!value || typeof value !== 'object') return false; + if (Array.isArray(value)) return value.some((item) => hasEnumOverflow(item, maxEnum)); + const record = value as Record; + if (Array.isArray(record.enum) && record.enum.length > maxEnum) return true; + return Object.values(record).some((item) => hasEnumOverflow(item, maxEnum)); +} diff --git a/apps/api/src/durable-objects/interaction-store.ts b/apps/api/src/durable-objects/interaction-store.ts new file mode 100644 index 000000000..b4f449098 --- /dev/null +++ b/apps/api/src/durable-objects/interaction-store.ts @@ -0,0 +1,712 @@ +import { + ACP_INTERACTION_ATTENTION_SOURCE, + ACP_INTERACTION_PROTOCOL_VERSION, + type AcpInteractionAnswerDecision, + AcpInteractionAnswerDecisionSchema, + type AcpInteractionSafeSummary, + isJsonRecord, +} from '@simple-agent-manager/shared'; +import { DurableObject } from 'cloudflare:workers'; +import * as v from 'valibot'; + +import type { Env } from '../env'; +import { canonicalJson } from '../lib/canonical-json'; +import { createModuleLogger } from '../lib/logger'; +import { getCredentialEncryptionKey } from '../lib/secrets'; +import { + type AcpInteractionConfig, + getAcpInteractionConfig, +} from '../services/acp-interaction-config'; +import { + deliverAcpInteractionAnswer, + resolveAcpInteractionDeliveryTarget, +} from '../services/acp-interaction-delivery'; +import { decrypt, encrypt } from '../services/encryption'; +import { + detailBoundsViolation, + type InteractionRow, + type InteractionStoreAnswerInput, + type InteractionStoreAnswerResult, + type InteractionStoreCreateInput, + type InteractionStoreCreateResult, + type InteractionStoreSettleInput, + type InteractionStoreSnapshot, + nowMs, + parseSummary, + sha256, + terminalState, +} from './interaction-store-model'; +import type { ProjectData } from './project-data'; + +const log = createModuleLogger('interaction_store'); + +type CreateValidationResult = Exclude< + InteractionStoreCreateResult, + { status: 'created' | 'existing' } +>; + +type EncryptedNullablePayload = { + ciphertext: string | null; + iv: string | null; +}; + +export type { + InteractionStoreAnswerInput, + InteractionStoreCreateInput, + InteractionStoreSettleInput, +} from './interaction-store-model'; + +export class InteractionStore extends DurableObject { + private readonly sql: SqlStorage; + + constructor(ctx: DurableObjectState, env: Env) { + super(ctx, env); + this.sql = ctx.storage.sql; + ctx.blockConcurrencyWhile(async () => { + this.sql.exec( + `CREATE TABLE IF NOT EXISTS interactions ( + interaction_id TEXT PRIMARY KEY NOT NULL, + project_id TEXT NOT NULL, + chat_session_id TEXT NOT NULL, + agent_session_id TEXT NOT NULL, + kind TEXT NOT NULL, + state TEXT NOT NULL, + generation TEXT NOT NULL, + runtime_identity TEXT NOT NULL, + payload_hash TEXT NOT NULL, + encrypted_detail TEXT, + detail_iv TEXT, + detail_purged_at INTEGER, + safe_summary_json TEXT NOT NULL, + upstream_request_id TEXT, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL, + deadline_at INTEGER NOT NULL, + answered_at INTEGER, + answer_key TEXT, + answer_body_hash TEXT, + decision_kind TEXT, + decision_hash TEXT, + encrypted_answer TEXT, + answer_iv TEXT, + encrypted_decision TEXT, + decision_iv TEXT, + delivery_state TEXT, + delivery_attempts INTEGER NOT NULL DEFAULT 0, + delivery_deadline_at INTEGER, + last_delivery_error TEXT, + attention_marker_id TEXT, + attention_projection_state TEXT, + terminal_at INTEGER, + purge_at INTEGER + )` + ); + this.sql.exec( + `CREATE TABLE IF NOT EXISTS outbox ( + id TEXT PRIMARY KEY NOT NULL, + interaction_id TEXT NOT NULL, + kind TEXT NOT NULL, + due_at INTEGER NOT NULL, + attempts INTEGER NOT NULL DEFAULT 0, + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL + )` + ); + this.sql.exec( + `CREATE INDEX IF NOT EXISTS idx_interactions_state_deadline + ON interactions(state, deadline_at)` + ); + this.sql.exec( + `CREATE INDEX IF NOT EXISTS idx_interactions_delivery_due + ON interactions(delivery_state, delivery_deadline_at)` + ); + this.addMissingDecisionColumns(); + this.sql.exec(`CREATE INDEX IF NOT EXISTS idx_interactions_purge ON interactions(purge_at)`); + this.sql.exec(`CREATE INDEX IF NOT EXISTS idx_outbox_due ON outbox(due_at)`); + }); + } + + async create(input: InteractionStoreCreateInput): Promise { + const config = getAcpInteractionConfig(this.env); + const now = nowMs(); + const existing = this.read(input.interactionId); + if (existing) { + if (existing.payload_hash !== input.payloadHash) { + return { status: 'conflict', reason: 'interaction id already exists with another payload' }; + } + return { status: 'existing', summary: parseSummary(existing) }; + } + const validation = this.validateCreateInput(input, config, now); + if (validation) return validation; + + const detailPlaintext = canonicalJson(input.detail); + const encrypted = await encrypt(detailPlaintext, getCredentialEncryptionKey(this.env)); + const safeSummary = canonicalJson(input.safeSummary); + this.sql.exec( + `INSERT INTO interactions ( + interaction_id, project_id, chat_session_id, agent_session_id, kind, state, + generation, runtime_identity, payload_hash, encrypted_detail, detail_iv, + safe_summary_json, upstream_request_id, created_at, updated_at, deadline_at, + delivery_attempts, attention_projection_state + ) VALUES (?, ?, ?, ?, ?, 'pending', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, 'pending')`, + input.interactionId, + input.projectId, + input.chatSessionId, + input.agentSessionId, + input.kind, + input.generation, + input.runtimeIdentity, + input.payloadHash, + encrypted.ciphertext, + encrypted.iv, + safeSummary, + input.upstreamRequestId ?? null, + now, + now, + input.deadlineAt + ); + this.enqueue(`projection:${input.interactionId}`, input.interactionId, 'projection', now); + this.enqueue(`expire:${input.interactionId}`, input.interactionId, 'expiry', input.deadlineAt); + await this.scheduleNextAlarm(); + return { status: 'created', summary: parseSummary(this.readRequired(input.interactionId)) }; + } + + private validateCreateInput( + input: InteractionStoreCreateInput, + config: AcpInteractionConfig, + now: number + ): CreateValidationResult | null { + if (!config.enabled) return { status: 'disabled', reason: 'ACP interactions are disabled' }; + if (input.protocolVersion !== ACP_INTERACTION_PROTOCOL_VERSION) { + return { status: 'invalid', reason: 'unsupported protocol version' }; + } + if (input.deadlineAt <= now) return { status: 'expired', reason: 'deadline is already past' }; + if (input.deadlineAt - now > config.maxDeadlineMs) { + return { status: 'invalid', reason: 'deadline exceeds configured maximum' }; + } + if (this.pendingCount() >= config.maxPendingPerSession) { + return { status: 'too_many_pending', reason: 'too many pending interactions for session' }; + } + const detailPlaintext = canonicalJson(input.detail); + if (new TextEncoder().encode(detailPlaintext).byteLength > config.requestMaxBytes) { + return { status: 'invalid', reason: 'request detail exceeds configured maximum' }; + } + const boundsViolation = detailBoundsViolation(input.detail, config); + if (boundsViolation) return { status: 'invalid', reason: boundsViolation }; + return null; + } + + async answer(input: InteractionStoreAnswerInput): Promise { + const config = getAcpInteractionConfig(this.env); + const row = this.read(input.interactionId); + if (row?.project_id !== input.projectId || row.chat_session_id !== input.chatSessionId) { + return { status: 'not_found', reason: 'interaction not found' }; + } + if (row.state !== 'pending') { + return this.answerResultForCommittedRow(row, input); + } + const answerPlaintext = canonicalJson(input.decision); + if (new TextEncoder().encode(answerPlaintext).byteLength > config.answerMaxBytes) { + return { status: 'payload_too_large', reason: 'answer exceeds configured maximum' }; + } + const encryptionKey = getCredentialEncryptionKey(this.env); + const encryptedDecision = await encrypt(answerPlaintext, encryptionKey); + const encryptedAnswer: EncryptedNullablePayload = input.decision.encryptedAnswer + ? await encrypt(canonicalJson(input.decision.encryptedAnswer), encryptionKey) + : { ciphertext: null, iv: null }; + const now = nowMs(); + const deliveryDeadlineAt = Math.min(now + config.deliveryWindowMs, row.deadline_at); + const purgeAt = now + config.sensitivePurgeMs; + this.sql.exec( + `UPDATE interactions + SET state = 'answered', + updated_at = ?, + answered_at = ?, + answer_key = ?, + answer_body_hash = ?, + decision_kind = ?, + decision_hash = ?, + encrypted_answer = ?, + answer_iv = ?, + encrypted_decision = ?, + decision_iv = ?, + delivery_state = 'pending', + delivery_deadline_at = ?, + purge_at = ? + WHERE interaction_id = ? AND state = 'pending'`, + now, + now, + input.answerKey, + input.answerBodyHash, + input.decision.kind, + input.decision.answerHash, + encryptedAnswer.ciphertext, + encryptedAnswer.iv, + encryptedDecision.ciphertext, + encryptedDecision.iv, + deliveryDeadlineAt, + purgeAt, + input.interactionId + ); + this.enqueue(`delivery:${input.interactionId}`, input.interactionId, 'delivery', now); + this.enqueue(`purge:${input.interactionId}`, input.interactionId, 'purge', purgeAt); + await this.scheduleNextAlarm(); + const updated = this.readRequired(input.interactionId); + if ( + updated.answer_key !== input.answerKey || + updated.answer_body_hash !== input.answerBodyHash + ) { + return this.answerResultForCommittedRow(updated, input); + } + return { + status: 'answered', + summary: parseSummary(updated), + delivery: { generation: updated.generation, runtimeIdentity: updated.runtime_identity }, + }; + } + + private answerResultForCommittedRow( + row: InteractionRow, + input: InteractionStoreAnswerInput + ): InteractionStoreAnswerResult { + if (row.answer_key === input.answerKey && row.answer_body_hash === input.answerBodyHash) { + return { + status: 'already_answered', + summary: parseSummary(row), + delivery: { generation: row.generation, runtimeIdentity: row.runtime_identity }, + }; + } + if (row.answer_key === input.answerKey) { + return { status: 'answer_key_conflict', reason: 'answer key was reused with another body' }; + } + if (row.answer_key) return { status: 'conflict', reason: 'another decision is already committed' }; + return { status: 'stale', reason: `interaction is ${row.state}` }; + } + + settle(input: InteractionStoreSettleInput): { status: 'settled' | 'not_found' | 'stale' } { + const row = this.read(input.interactionId); + if (!row) return { status: 'not_found' }; + if (row.generation !== input.generation || row.runtime_identity !== input.runtimeIdentity) { + return { status: 'stale' }; + } + if (row.state !== 'pending') return { status: 'settled' }; + const now = nowMs(); + this.markTerminal(row.interaction_id, 'cancelled', now, input.reason); + return { status: 'settled' }; + } + + snapshot(cursor: string | null = null): InteractionStoreSnapshot { + const config = getAcpInteractionConfig(this.env); + const pending = this.sql + .exec( + `SELECT * FROM interactions + WHERE state IN ('pending', 'answered') + ORDER BY created_at ASC + LIMIT ?`, + config.maxPendingPerSession + config.snapshotLastSettled + ) + .toArray() + .map(parseSummary); + const settledRows = this.sql + .exec( + `SELECT * FROM interactions + WHERE state NOT IN ('pending', 'answered') + AND (? IS NULL OR updated_at < ?) + ORDER BY updated_at DESC + LIMIT ?`, + cursor, + cursor ? Number.parseInt(cursor, 10) : null, + config.snapshotLastSettled + 1 + ) + .toArray(); + const pageRows = settledRows.slice(0, config.snapshotLastSettled); + return { + pending, + settled: pageRows.map(parseSummary), + cursor: + settledRows.length > config.snapshotLastSettled + ? String(pageRows[pageRows.length - 1]?.updated_at ?? '') + : null, + }; + } + + async detail( + interactionId: string + ): Promise<{ summary: AcpInteractionSafeSummary; detail: Record | null } | null> { + const row = this.read(interactionId); + if (!row) return null; + let detail: Record | null = null; + if (row.encrypted_detail?.length && row.detail_iv?.length) { + const plaintext = await decrypt( + row.encrypted_detail, + row.detail_iv, + getCredentialEncryptionKey(this.env) + ); + const parsedDetail = JSON.parse(plaintext) as unknown; + detail = isJsonRecord(parsedDetail) ? parsedDetail : null; + } + return { summary: parseSummary(row), detail }; + } + + async alarm(): Promise { + const now = nowMs(); + await this.processDueExpiry(now); + await this.processDueOutbox(now); + await this.processDuePurge(now); + this.compactSettled(now); + await this.scheduleNextAlarm(); + } + + recordDelivery( + interactionId: string, + outcome: 'confirmed' | 'unconfirmed' | 'interrupted', + error: string | null = null + ): { status: 'recorded' | 'not_found' } { + const row = this.read(interactionId); + if (!row) return { status: 'not_found' }; + const now = nowMs(); + const state = deliveryStateForOutcome(outcome); + this.sql.exec( + `UPDATE interactions + SET state = ?, delivery_state = ?, updated_at = ?, terminal_at = ?, last_delivery_error = ? + WHERE interaction_id = ? AND state = 'answered'`, + state, + outcome, + now, + now, + error, + interactionId + ); + this.enqueue(`resolve:${interactionId}`, interactionId, 'projection_resolve', now); + return { status: 'recorded' }; + } + + purge(): void { + this.sql.exec(`DELETE FROM outbox`); + this.sql.exec(`DELETE FROM interactions`); + } + + private addMissingDecisionColumns(): void { + try { + this.sql.exec('ALTER TABLE interactions ADD COLUMN encrypted_decision TEXT'); + } catch (error) { + if (!(error instanceof Error) || !error.message.includes('duplicate column')) { + log.warn('interaction_store.add_encrypted_decision_column_failed', { + error: error instanceof Error ? error.message : String(error), + }); + throw error; + } + } + try { + this.sql.exec('ALTER TABLE interactions ADD COLUMN decision_iv TEXT'); + } catch (error) { + if (!(error instanceof Error) || !error.message.includes('duplicate column')) { + log.warn('interaction_store.add_decision_iv_column_failed', { + error: error instanceof Error ? error.message : String(error), + }); + throw error; + } + } + } + + private read(interactionId: string): InteractionRow | null { + return ( + this.sql + .exec( + `SELECT * FROM interactions WHERE interaction_id = ? LIMIT 1`, + interactionId + ) + .toArray()[0] ?? null + ); + } + + private readRequired(interactionId: string): InteractionRow { + const row = this.read(interactionId); + if (!row) throw new Error('interaction row disappeared after write'); + return row; + } + + private pendingCount(): number { + const row = this.sql + .exec<{ count: number }>( + `SELECT COUNT(*) AS count FROM interactions WHERE state IN ('pending', 'answered')` + ) + .toArray()[0]; + return row?.count ?? 0; + } + + private enqueue(id: string, interactionId: string, kind: string, dueAt: number): void { + const now = nowMs(); + this.sql.exec( + `INSERT INTO outbox (id, interaction_id, kind, due_at, attempts, created_at, updated_at) + VALUES (?, ?, ?, ?, 0, ?, ?) + ON CONFLICT(id) DO UPDATE SET due_at = excluded.due_at, updated_at = excluded.updated_at`, + id, + interactionId, + kind, + dueAt, + now, + now + ); + } + + private async processDueExpiry(now: number): Promise { + for (const row of this.sql + .exec( + `SELECT * FROM interactions + WHERE state = 'pending' AND deadline_at <= ? + ORDER BY deadline_at ASC + LIMIT 25`, + now + ) + .toArray()) { + this.markTerminal(row.interaction_id, 'expired', now, 'deadline'); + } + } + + private async processDueOutbox(now: number): Promise { + const due = this.sql + .exec<{ id: string; interaction_id: string; kind: string; attempts: number }>( + `SELECT id, interaction_id, kind, attempts FROM outbox + WHERE due_at <= ? + ORDER BY due_at ASC + LIMIT 25`, + now + ) + .toArray(); + for (const job of due) { + try { + if (job.kind === 'projection') await this.projectAttention(job.interaction_id); + else if (job.kind === 'projection_resolve') + await this.resolveProjectedAttention(job.interaction_id); + else if (job.kind === 'delivery') await this.processDeliveryJob(job.interaction_id, now); + this.sql.exec(`DELETE FROM outbox WHERE id = ?`, job.id); + } catch (error) { + const retryAt = now + getAcpInteractionConfig(this.env).retrySteadyMs; + this.sql.exec( + `UPDATE outbox SET attempts = attempts + 1, due_at = ?, updated_at = ? WHERE id = ?`, + retryAt, + now, + job.id + ); + log.warn('interaction_store.outbox_retry', { + interactionId: job.interaction_id, + kind: job.kind, + error: error instanceof Error ? error.message : String(error), + }); + } + } + } + + private async processDeliveryJob(interactionId: string, now: number): Promise { + const row = this.read(interactionId); + if (row?.state !== 'answered') return; + if (!row.encrypted_decision || !row.decision_iv) { + this.recordDelivery(interactionId, 'unconfirmed', 'missing encrypted decision for retry'); + return; + } + if (row.delivery_deadline_at !== null && now >= row.delivery_deadline_at) { + this.recordDelivery(interactionId, 'unconfirmed', 'delivery deadline elapsed'); + return; + } + const target = await resolveAcpInteractionDeliveryTarget( + this.env, + row.project_id, + row.chat_session_id + ); + if (target.status === 'interrupted') { + this.recordDelivery(interactionId, 'interrupted', target.reason); + return; + } + if (target.status === 'retry') { + this.rescheduleDelivery(row, now, target.reason); + return; + } + const plaintext = await decrypt( + row.encrypted_decision, + row.decision_iv, + getCredentialEncryptionKey(this.env) + ); + const decision = v.parse(AcpInteractionAnswerDecisionSchema, JSON.parse(plaintext) as unknown); + if (target.status !== 'ready') { + this.rescheduleDelivery(row, now, 'target not ready'); + return; + } + const delivery = await deliverAcpInteractionAnswer(this.env, target.target, { + interactionId, + generation: row.generation, + runtimeIdentity: row.runtime_identity, + decision, + }); + if (delivery.outcome === 'confirmed') { + this.recordDelivery(interactionId, 'confirmed', null); + return; + } + if (delivery.outcome === 'interrupted') { + this.recordDelivery(interactionId, 'interrupted', delivery.reason); + return; + } + this.rescheduleDelivery(row, now, delivery.reason); + } + + private rescheduleDelivery(row: InteractionRow, now: number, reason: string): void { + const config = getAcpInteractionConfig(this.env); + const attempts = row.delivery_attempts + 1; + const delay = + config.retryDelaysMs[Math.min(attempts - 1, config.retryDelaysMs.length - 1)] ?? + config.retrySteadyMs; + const dueAt = Math.min(now + delay, row.delivery_deadline_at ?? now + delay); + if (row.delivery_deadline_at !== null && dueAt >= row.delivery_deadline_at) { + this.recordDelivery(row.interaction_id, 'unconfirmed', reason); + return; + } + this.sql.exec( + `UPDATE interactions + SET delivery_attempts = ?, last_delivery_error = ?, updated_at = ? + WHERE interaction_id = ? AND state = 'answered'`, + attempts, + reason, + now, + row.interaction_id + ); + this.enqueue(`delivery:${row.interaction_id}`, row.interaction_id, 'delivery', dueAt); + } + + private async processDuePurge(now: number): Promise { + this.sql.exec( + `UPDATE interactions + SET encrypted_detail = NULL, + detail_iv = NULL, + encrypted_answer = NULL, + answer_iv = NULL, + detail_purged_at = ? + WHERE purge_at IS NOT NULL + AND purge_at <= ? + AND detail_purged_at IS NULL`, + now, + now + ); + } + + private compactSettled(now: number): void { + const config = getAcpInteractionConfig(this.env); + const cutoff = now - config.summaryRetentionMs; + this.sql.exec( + `DELETE FROM interactions + WHERE state NOT IN ('pending', 'answered') + AND terminal_at IS NOT NULL + AND terminal_at < ? + AND interaction_id NOT IN ( + SELECT interaction_id FROM interactions + WHERE state NOT IN ('pending', 'answered') + ORDER BY terminal_at DESC + LIMIT ? + )`, + cutoff, + config.summaryLastSettled + ); + } + + private markTerminal(interactionId: string, state: string, at: number, reason: string): void { + this.sql.exec( + `UPDATE interactions + SET state = ?, + updated_at = ?, + terminal_at = ?, + delivery_state = CASE WHEN delivery_state IS NULL THEN 'interrupted' ELSE delivery_state END, + last_delivery_error = ? + WHERE interaction_id = ? AND state IN ('pending', 'answered')`, + state, + at, + at, + reason, + interactionId + ); + this.enqueue(`resolve:${interactionId}`, interactionId, 'projection_resolve', at); + } + + private async projectAttention(interactionId: string): Promise { + const row = this.read(interactionId); + if (!row || terminalState(row.state) || row.attention_marker_id?.length) return; + const stub = this.projectDataStub(row.project_id); + const summary = parseSummary(row); + const marker = await stub.createAttentionMarker({ + sessionId: row.chat_session_id, + taskId: null, + workspaceId: null, + kind: 'needs_input', + source: ACP_INTERACTION_ATTENTION_SOURCE, + reason: 'acp_interaction_pending', + metadata: canonicalJson({ + source: ACP_INTERACTION_ATTENTION_SOURCE, + interactionId: row.interaction_id, + kind: row.kind, + state: row.state, + toolCallId: summary.toolCallId, + }), + expiresAt: null, + }); + this.sql.exec( + `UPDATE interactions + SET attention_marker_id = ?, attention_projection_state = 'created', updated_at = ? + WHERE interaction_id = ?`, + marker.id, + nowMs(), + interactionId + ); + } + + private async resolveProjectedAttention(interactionId: string): Promise { + const row = this.read(interactionId); + if (!row?.attention_marker_id) return; + const stub = this.projectDataStub(row.project_id); + await stub.resolveAttentionMarkerById( + row.attention_marker_id, + 'system', + `${ACP_INTERACTION_ATTENTION_SOURCE}:${row.state}` + ); + this.sql.exec( + `UPDATE interactions + SET attention_projection_state = 'resolved', updated_at = ? + WHERE interaction_id = ?`, + nowMs(), + interactionId + ); + } + + private projectDataStub(projectId: string): DurableObjectStub { + return this.env.PROJECT_DATA.get( + this.env.PROJECT_DATA.idFromName(projectId) + ) as DurableObjectStub; + } + + private async scheduleNextAlarm(): Promise { + const due = this.sql + .exec<{ due_at: number }>( + `SELECT due_at FROM outbox + UNION ALL + SELECT deadline_at AS due_at FROM interactions WHERE state = 'pending' + UNION ALL + SELECT purge_at AS due_at FROM interactions + WHERE purge_at IS NOT NULL AND detail_purged_at IS NULL + ORDER BY due_at ASC + LIMIT 1` + ) + .toArray()[0]?.due_at; + if (due) { + await this.ctx.storage.setAlarm(Math.max(due, nowMs() + 1000)); + } + } +} + +function deliveryStateForOutcome(outcome: 'confirmed' | 'unconfirmed' | 'interrupted'): string { + if (outcome === 'confirmed') return 'delivery_confirmed'; + if (outcome === 'unconfirmed') return 'delivery_unconfirmed'; + return 'interrupted'; +} + +export async function interactionDecisionHash( + decision: AcpInteractionAnswerDecision +): Promise { + return sha256(canonicalJson(decision)); +} diff --git a/apps/api/src/durable-objects/project-data/attention-expiry.ts b/apps/api/src/durable-objects/project-data/attention-expiry.ts index c01c1d1ca..e34fde78c 100644 --- a/apps/api/src/durable-objects/project-data/attention-expiry.ts +++ b/apps/api/src/durable-objects/project-data/attention-expiry.ts @@ -130,6 +130,15 @@ async function processExpiredNeedsInputMarker( failSession: (sessionId: string, errorMessage: string) => Promise, hooks: AttentionExpiryProcessingHooks ): Promise { + if (marker.source === 'acp_interaction') { + attention.resolveAttentionMarkerById(sql, marker.id, 'system', 'unsupported_expiry_source'); + log.warn('attention_marker.acp_interaction_expiry_ignored', { + markerId: marker.id, + sessionId: marker.sessionId, + taskId: marker.taskId, + }); + return; + } const now = Date.now(); const maxExpiresAt = marker.maxExpiresAt ?? @@ -220,6 +229,7 @@ async function failExpiredTaskMarker( failSession: (sessionId: string, errorMessage: string) => Promise, hooks: AttentionExpiryProcessingHooks ): Promise { + if (marker.source === 'acp_interaction') return; if ((marker.kind !== 'needs_input' && marker.kind !== 'reconciliation_checkin') || !marker.taskId) return; diff --git a/apps/api/src/durable-objects/project-data/attention.ts b/apps/api/src/durable-objects/project-data/attention.ts index 029abdc19..d9334078e 100644 --- a/apps/api/src/durable-objects/project-data/attention.ts +++ b/apps/api/src/durable-objects/project-data/attention.ts @@ -213,6 +213,7 @@ export function linkAttentionNotification( export type PrepareAttentionAnswerResult = | { status: 'ready' } | { status: 'not_found' } + | { status: 'unsupported_source'; source: string } | { status: 'already_resolved'; answer: string | null } | { status: 'in_flight'; answer: string } | { status: 'conflicting_answer'; answer: string } @@ -242,6 +243,9 @@ export function prepareAttentionAnswer( if (rows.length === 0) return { status: 'not_found' }; const marker = parseAttentionMarkerRow(rows[0]); + if (marker.source === 'acp_interaction') { + return { status: 'unsupported_source', source: marker.source }; + } if (marker.resolvedAt !== null) { return { status: 'already_resolved', answer: marker.resolvedAnswer }; } @@ -377,7 +381,7 @@ export function getAttentionSummary( export function getExpiredMarkers(sql: SqlStorage, now: number = Date.now()) { const rows = sql .exec( - `SELECT id, session_id, task_id, workspace_id, kind, + `SELECT id, session_id, task_id, workspace_id, kind, source, source_notification_id, notification_user_id, created_at, expires_at, next_escalation_at, escalation_count, max_expires_at FROM session_attention_markers @@ -392,7 +396,7 @@ export function getExpiredMarkers(sql: SqlStorage, now: number = Date.now()) { export function getDueAttentionEscalations(sql: SqlStorage, now: number = Date.now()) { const rows = sql .exec( - `SELECT id, session_id, task_id, workspace_id, kind, + `SELECT id, session_id, task_id, workspace_id, kind, source, source_notification_id, notification_user_id, created_at, expires_at, next_escalation_at, escalation_count, max_expires_at FROM session_attention_markers diff --git a/apps/api/src/durable-objects/project-data/index.ts b/apps/api/src/durable-objects/project-data/index.ts index c5e52bc84..9521c0a03 100644 --- a/apps/api/src/durable-objects/project-data/index.ts +++ b/apps/api/src/durable-objects/project-data/index.ts @@ -1974,6 +1974,19 @@ export class ProjectData extends DurableObject { return count; } + async resolveAttentionMarkerById( + markerId: string, + actorType: string = 'system', + reason: string = 'system_resolved' + ): Promise { + const count = attention.resolveAttentionMarkerById(this.sql, markerId, actorType, reason); + if (count > 0) { + await this.recalculateAlarm(); + this.broadcastEvent('attention.resolved', { markerId, count, reason }); + } + return count; + } + async resolveSessionAttentionMarkers( sessionId: string, resolvedByMessageId: string | null, diff --git a/apps/api/src/durable-objects/project-data/row-schemas/attention.ts b/apps/api/src/durable-objects/project-data/row-schemas/attention.ts index d84170744..5067f335e 100644 --- a/apps/api/src/durable-objects/project-data/row-schemas/attention.ts +++ b/apps/api/src/durable-objects/project-data/row-schemas/attention.ts @@ -118,6 +118,7 @@ const AttentionExpiryRowSchema = v.object({ task_id: v.nullable(v.string()), workspace_id: v.nullable(v.string()), kind: v.string(), + source: v.string(), source_notification_id: v.nullable(v.string()), notification_user_id: v.nullable(v.string()), created_at: v.number(), @@ -133,6 +134,7 @@ export function parseAttentionExpiryRow(row: unknown): { taskId: string | null; workspaceId: string | null; kind: string; + source: string; sourceNotificationId: string | null; notificationUserId: string | null; createdAt: number; @@ -148,6 +150,7 @@ export function parseAttentionExpiryRow(row: unknown): { taskId: r.task_id, workspaceId: r.workspace_id, kind: r.kind, + source: r.source, sourceNotificationId: r.source_notification_id, notificationUserId: r.notification_user_id, createdAt: r.created_at, diff --git a/apps/api/src/env.ts b/apps/api/src/env.ts index 32ca0d0db..b89f7686a 100644 --- a/apps/api/src/env.ts +++ b/apps/api/src/env.ts @@ -62,6 +62,7 @@ export interface Env extends WebhookTriggerEnv, TaskRecoveryEnv { ADMIN_LOGS: DurableObjectNamespace; TASK_RUNNER: DurableObjectNamespace; DIAGNOSIS_RUNNER: DurableObjectNamespace; + INTERACTION_STORE: DurableObjectNamespace; NOTIFICATION: DurableObjectNamespace; CODEX_REFRESH_LOCK: DurableObjectNamespace; GITHUB_USER_ACCESS_TOKEN_LOCK: DurableObjectNamespace; @@ -838,6 +839,27 @@ export interface Env extends WebhookTriggerEnv, TaskRecoveryEnv { PROJECT_EVENT_WAKE_SUBSCRIPTION_LIFETIME_MS?: string; PROJECT_EVENT_WAKE_MAX_PER_SUBSCRIPTION?: string; PROJECT_EVENT_SOURCE_OUTBOX_SWEEP_WALL_MS?: string; + ACP_INTERACTIONS_ENABLED?: string; + ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS?: string; + ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS?: string; + ACP_INTERACTION_MAX_DEADLINE_MS?: string; + ACP_INTERACTION_DEADLINE_MARGIN_MS?: string; + ACP_INTERACTION_MAX_PENDING_PER_SESSION?: string; + ACP_INTERACTION_REQUEST_MAX_BYTES?: string; + ACP_INTERACTION_OPTIONS_MAX_COUNT?: string; + ACP_INTERACTION_OPTION_NAME_MAX_CHARS?: string; + ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES?: string; + ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES?: string; + ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM?: string; + ACP_INTERACTION_ANSWER_MAX_BYTES?: string; + ACP_INTERACTION_ANSWER_STRING_MAX_BYTES?: string; + ACP_INTERACTION_RETRY_DELAYS_MS?: string; + ACP_INTERACTION_RETRY_STEADY_MS?: string; + ACP_INTERACTION_DELIVERY_WINDOW_MS?: string; + ACP_INTERACTION_SENSITIVE_PURGE_MS?: string; + ACP_INTERACTION_SUMMARY_RETENTION_MS?: string; + ACP_INTERACTION_SUMMARY_LAST_SETTLED?: string; + ACP_INTERACTION_SNAPSHOT_LAST_SETTLED?: string; PROJECT_EVENT_SOURCE_OUTBOX_ADMISSION_TIMEOUT_MS?: string; PROJECT_EVENT_SOURCE_OUTBOX_TERMINAL_RETENTION_MS?: string; MESSAGE_SIZE_THRESHOLD?: string; diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index ef261788a..34812cafa 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -9,6 +9,7 @@ export { CredentialSetupSession } from './durable-objects/credential-setup-sessi export { DiagnosisRunner } from './durable-objects/diagnosis-runner'; export { GitHubUserAccessTokenLock } from './durable-objects/github-user-access-token-lock'; export { GitLabUserAccessTokenLock } from './durable-objects/gitlab-user-access-token-lock'; +export { InteractionStore } from './durable-objects/interaction-store'; export { NodeLifecycle } from './durable-objects/node-lifecycle'; export { NotificationService } from './durable-objects/notification'; export { ProjectAgent } from './durable-objects/project-agent'; @@ -127,6 +128,7 @@ import { import { projectEventChannelRoutes } from './routes/project-event-channels'; import { projectEventSubscriptionRoutes } from './routes/project-event-subscriptions'; import { projectsRoutes } from './routes/projects'; +import { acpInteractionCallbackRoute } from './routes/projects/acp-interaction-callback'; import { agentActivityCallbackRoute } from './routes/projects/agent-activity-callback'; import { agentUsageCallbackRoute } from './routes/projects/agent-usage-callback'; import { buildStartedCallbackRoute } from './routes/projects/build-started-callback'; @@ -813,6 +815,7 @@ app.route('/api/webhooks', triggerWebhookRoutes); // See .claude/rules/06-api-patterns.md (Hono middleware scoping) app.route('/api/projects', deploymentIdentityTokenRoute); app.route('/api/projects', nodeAcpHeartbeatRoute); +app.route('/api/projects', acpInteractionCallbackRoute); // Must be before projectsRoutes — uses callback JWT, not session auth app.route('/api/projects', agentActivityCallbackRoute); // Must be before projectsRoutes — uses callback JWT, not session auth app.route('/api/projects', agentUsageCallbackRoute); // Must be before projectsRoutes — uses callback JWT, not session auth app.route('/api/projects', buildStartedCallbackRoute); // Must be before projectsRoutes — uses callback JWT, not session auth diff --git a/apps/api/src/routes/chat-acp-interactions.ts b/apps/api/src/routes/chat-acp-interactions.ts new file mode 100644 index 000000000..c1013378b --- /dev/null +++ b/apps/api/src/routes/chat-acp-interactions.ts @@ -0,0 +1,208 @@ +import { + type AcpInteractionAnswerDecision, + AcpInteractionBrowserAnswerSchema, + AcpInteractionIdSchema, +} from '@simple-agent-manager/shared'; +import { drizzle } from 'drizzle-orm/d1'; +import type { Context, Hono } from 'hono'; +import type { InferOutput } from 'valibot'; +import * as v from 'valibot'; + +import * as schema from '../db/schema'; +import type { Env } from '../env'; +import { getTrustedApiOrigin } from '../lib/trusted-origins'; +import { getUserId } from '../middleware/auth'; +import { errors } from '../middleware/error'; +import { requireProjectCapability } from '../middleware/project-auth'; +import { + deliverAcpInteractionAnswer, + resolveAcpInteractionDeliveryTarget, +} from '../services/acp-interaction-delivery'; +import { + answerInteraction, + getInteractionDetail, + recordInteractionDelivery, + snapshotInteractions, +} from '../services/acp-interaction-store'; +import { requireSessionCreator } from './chat-session-ownership'; + +type AnswerInteractionResult = Awaited>; +type AcceptedAnswerResult = Extract< + AnswerInteractionResult, + { status: 'answered' | 'already_answered' } +>; +type ChatAcpContext = Context<{ Bindings: Env }>; +type BrowserAnswerBody = InferOutput; + +function exactAppOrigin(env: Pick): string { + const apiOrigin = getTrustedApiOrigin(env); + const url = new URL(apiOrigin); + if (url.hostname.startsWith('api.')) url.hostname = `app.${url.hostname.slice(4)}`; + return url.toString().replace(/\/$/u, ''); +} + +function requiredParam(value: string | undefined, name: string): string { + if (!value) throw errors.badRequest(`${name} is required`); + return value; +} + +function requireExactBrowserOrigin(c: { + req: { header: (name: string) => string | undefined }; + env: Env; +}): void { + const origin = c.req.header('Origin'); + const fetchSite = c.req.header('Sec-Fetch-Site'); + if (fetchSite === 'cross-site') { + throw errors.forbidden('Cross-site interaction writes are not allowed'); + } + if (!origin || origin !== exactAppOrigin(c.env)) { + throw errors.forbidden('Untrusted interaction write origin'); + } +} + +function requireAcceptedAnswer(answer: AnswerInteractionResult): AcceptedAnswerResult { + if (answer.status === 'answered' || answer.status === 'already_answered') return answer; + if (answer.status === 'not_found') throw errors.notFound('ACP interaction'); + if ( + answer.status === 'conflict' || + answer.status === 'answer_key_conflict' || + answer.status === 'stale' + ) { + throw errors.conflict(answer.reason); + } + throw errors.badRequest('reason' in answer ? answer.reason : 'Interaction answer was not accepted'); +} + +async function parseBrowserAnswerBody(c: ChatAcpContext): Promise { + let raw: unknown; + try { + raw = (await c.req.json()) as unknown; + } catch { + throw errors.badRequest('Invalid JSON in request body'); + } + const result = v.safeParse(AcpInteractionBrowserAnswerSchema, raw); + if (!result.success) throw errors.badRequest('Invalid interaction answer body'); + return result.output; +} + +async function deliverAcceptedAnswer( + env: Env, + projectId: string, + sessionId: string, + interactionId: string, + answer: AcceptedAnswerResult, + decision: AcpInteractionAnswerDecision +): Promise { + if (answer.status !== 'answered') return; + const target = await resolveAcpInteractionDeliveryTarget(env, projectId, sessionId); + if (target.status === 'ready') { + const delivery = await deliverAcpInteractionAnswer(env, target.target, { + interactionId, + generation: answer.delivery.generation, + runtimeIdentity: answer.delivery.runtimeIdentity, + decision, + }); + if (delivery.outcome !== 'unconfirmed') { + await recordInteractionDelivery( + env, + projectId, + sessionId, + interactionId, + delivery.outcome, + 'reason' in delivery ? delivery.reason : null + ); + } + return; + } + if (target.status === 'interrupted') { + await recordInteractionDelivery( + env, + projectId, + sessionId, + interactionId, + 'interrupted', + target.reason + ); + } +} + +async function listInteractionSnapshots(c: ChatAcpContext): Promise { + const userId = getUserId(c); + const projectId = requiredParam(c.req.param('projectId'), 'projectId'); + const sessionId = requiredParam(c.req.param('sessionId'), 'sessionId'); + const db = drizzle(c.env.DATABASE, { schema }); + + await requireProjectCapability(db, projectId, userId, 'task:read'); + const session = await requireSessionCreator(c.env, projectId, sessionId, userId).catch(() => null); + const snapshot = await snapshotInteractions( + c.env, + projectId, + sessionId, + c.req.query('cursor') ?? null + ); + if (session) return c.json(snapshot); + return c.json({ + pending: snapshot.pending.map((item) => ({ + interactionId: item.interactionId, + kind: item.kind, + state: item.state, + createdAt: item.createdAt, + deadlineAt: item.deadlineAt, + })), + settled: [], + cursor: null, + }); +} + +async function readInteractionDetail(c: ChatAcpContext): Promise { + const userId = getUserId(c); + const projectId = requiredParam(c.req.param('projectId'), 'projectId'); + const sessionId = requiredParam(c.req.param('sessionId'), 'sessionId'); + const interactionId = v.parse(AcpInteractionIdSchema, c.req.param('interactionId')); + const db = drizzle(c.env.DATABASE, { schema }); + + await requireProjectCapability(db, projectId, userId, 'task:read'); + await requireSessionCreator(c.env, projectId, sessionId, userId); + const detail = await getInteractionDetail(c.env, projectId, sessionId, interactionId); + if (!detail) throw errors.notFound('ACP interaction'); + c.header('Cache-Control', 'private, no-store'); + return c.json(detail); +} + +async function answerInteractionRoute(c: ChatAcpContext): Promise { + requireExactBrowserOrigin(c); + const userId = getUserId(c); + const projectId = requiredParam(c.req.param('projectId'), 'projectId'); + const sessionId = requiredParam(c.req.param('sessionId'), 'sessionId'); + const interactionId = v.parse(AcpInteractionIdSchema, c.req.param('interactionId')); + const body = await parseBrowserAnswerBody(c); + const db = drizzle(c.env.DATABASE, { schema }); + + await requireProjectCapability(db, projectId, userId, 'task:write'); + await requireSessionCreator(c.env, projectId, sessionId, userId); + const answer = await answerInteraction(c.env, { + projectId, + chatSessionId: sessionId, + interactionId, + answerKey: body.answerKey, + answerBodyHash: body.decision.answerHash, + decision: body.decision, + }); + const acceptedAnswer = requireAcceptedAnswer(answer); + await deliverAcceptedAnswer( + c.env, + projectId, + sessionId, + interactionId, + acceptedAnswer, + body.decision + ); + + return c.json({ accepted: true, state: acceptedAnswer.summary.state }); +} + +export function registerChatAcpInteractionRoutes(chatRoutes: Hono<{ Bindings: Env }>): void { + chatRoutes.get('/:sessionId/interactions', listInteractionSnapshots); + chatRoutes.get('/:sessionId/interactions/:interactionId', readInteractionDetail); + chatRoutes.post('/:sessionId/interactions/:interactionId/answer', answerInteractionRoute); +} diff --git a/apps/api/src/routes/chat.ts b/apps/api/src/routes/chat.ts index df780c1a6..c544bbf1d 100644 --- a/apps/api/src/routes/chat.ts +++ b/apps/api/src/routes/chat.ts @@ -36,6 +36,7 @@ import * as projectDataService from '../services/project-data'; import { publicPlacementExplanationJson } from '../services/public-placement-explanation'; import { isTaskStatus } from '../services/task-status'; import { attachWakeState } from './chat/wake-state'; +import { registerChatAcpInteractionRoutes } from './chat-acp-interactions'; import { resolveChatAgentState } from './chat-agent-state'; import { registerChatCancelRoute } from './chat-cancel'; import { registerChatCommentDirectiveRoute } from './chat-comment-directives'; @@ -63,6 +64,7 @@ const chatRoutes = new Hono<{ Bindings: Env }>(); chatRoutes.use('/*', requireAuth(), requireApproved()); registerChatSessionListRoute(chatRoutes); +registerChatAcpInteractionRoutes(chatRoutes); /** * POST /api/projects/:projectId/sessions @@ -428,6 +430,9 @@ chatRoutes.post('/:sessionId/attention/:markerId/resolve', async (c) => { markerId, answer ); + if (prepared.status === 'unsupported_source') { + throw errors.badRequest('This attention marker must be answered through the interaction route'); + } if (prepared.status === 'not_found') throw errors.notFound('Attention request'); if (prepared.status === 'invalid_option') { throw errors.badRequest('answer must match one of the requested options'); diff --git a/apps/api/src/routes/projects/acp-interaction-callback.ts b/apps/api/src/routes/projects/acp-interaction-callback.ts new file mode 100644 index 000000000..4f2b50370 --- /dev/null +++ b/apps/api/src/routes/projects/acp-interaction-callback.ts @@ -0,0 +1,135 @@ +import { + AcpInteractionRuntimeCreateSchema, + AcpInteractionRuntimeSettleSchema, +} from '@simple-agent-manager/shared'; +import { eq } from 'drizzle-orm'; +import { drizzle } from 'drizzle-orm/d1'; +import { Hono } from 'hono'; + +import * as schema from '../../db/schema'; +import type { Env } from '../../env'; +import { extractBearerToken } from '../../lib/auth-helpers'; +import { log } from '../../lib/logger'; +import { errors } from '../../middleware/error'; +import { jsonValidator } from '../../schemas'; +import { createInteraction, settleInteraction } from '../../services/acp-interaction-store'; +import { verifyCallbackToken } from '../../services/jwt'; + +const acpInteractionCallbackRoute = new Hono<{ Bindings: Env }>(); + +interface VerifiedWorkspaceIdentity { + projectId: string; + workspaceId: string; + chatSessionId: string; + userId: string; + status: string; +} + +async function verifyWorkspaceCallback(c: { + env: Env; + req: { header: (name: string) => string | undefined; param: (name: string) => string }; +}): Promise { + const projectId = c.req.param('id'); + const workspaceId = c.req.param('workspaceId'); + const token = extractBearerToken(c.req.header('Authorization')); + const payload = await verifyCallbackToken(token, c.env, { expectedScope: 'workspace' }); + if (payload.workspace !== workspaceId) { + throw errors.unauthorized('Callback token does not match workspace'); + } + + const db = drizzle(c.env.DATABASE, { schema }); + const row = await db + .select({ + projectId: schema.workspaces.projectId, + chatSessionId: schema.workspaces.chatSessionId, + userId: schema.workspaces.userId, + status: schema.workspaces.status, + }) + .from(schema.workspaces) + .where(eq(schema.workspaces.id, workspaceId)) + .get(); + if (row?.projectId !== projectId || !row.chatSessionId) { + throw errors.notFound('Workspace'); + } + const chatSessionId = row.chatSessionId; + if (!['creating', 'running', 'recovery'].includes(row.status)) { + throw errors.gone(`Workspace is ${row.status}`); + } + return { + projectId, + workspaceId, + chatSessionId, + userId: row.userId, + status: row.status, + }; +} + +async function assertAgentSessionCurrent(env: Env, workspaceId: string, agentSessionId: string) { + const row = await env.DATABASE.prepare( + `SELECT id, status FROM agent_sessions WHERE id = ? AND workspace_id = ? LIMIT 1` + ) + .bind(agentSessionId, workspaceId) + .first<{ id: string; status: string }>(); + if (!row) throw errors.notFound('Agent session'); + if (row.status !== 'running') throw errors.conflict(`Agent session is ${row.status}`); +} + +function settleStatusCode(status: string): 200 | 404 | 409 { + if (status === 'not_found') return 404; + if (status === 'stale') return 409; + return 200; +} + +/** + * VM-agent ACP interaction callbacks. Mounted before projectsRoutes because it + * uses workspace callback JWT Bearer auth, not browser session cookies. + */ +acpInteractionCallbackRoute.post( + '/:id/workspaces/:workspaceId/acp-interactions', + jsonValidator(AcpInteractionRuntimeCreateSchema), + async (c) => { + const identity = await verifyWorkspaceCallback(c); + const body = c.req.valid('json'); + await assertAgentSessionCurrent(c.env, identity.workspaceId, body.agentSessionId); + const result = await createInteraction(c.env, { + ...body, + projectId: identity.projectId, + chatSessionId: identity.chatSessionId, + }); + if (result.status === 'created' || result.status === 'existing') { + log.info('acp_interaction.created', { + projectId: identity.projectId, + chatSessionId: identity.chatSessionId, + interactionId: body.interactionId, + kind: body.kind, + state: result.summary.state, + }); + return c.json(result, result.status === 'created' ? 201 : 200); + } + if (result.status === 'disabled') return c.json(result, 409); + if (result.status === 'conflict') return c.json(result, 409); + if (result.status === 'too_many_pending') return c.json(result, 429); + return c.json(result, 400); + } +); + +acpInteractionCallbackRoute.post( + '/:id/workspaces/:workspaceId/acp-interactions/:interactionId/settle', + jsonValidator(AcpInteractionRuntimeSettleSchema), + async (c) => { + const identity = await verifyWorkspaceCallback(c); + const body = c.req.valid('json'); + if (body.interactionId !== c.req.param('interactionId')) { + throw errors.badRequest('interactionId route/body mismatch'); + } + await assertAgentSessionCurrent(c.env, identity.workspaceId, body.agentSessionId); + const result = await settleInteraction(c.env, { + ...body, + projectId: identity.projectId, + chatSessionId: identity.chatSessionId, + }); + return c.json(result, settleStatusCode(result.status)); + } +); + +export { acpInteractionCallbackRoute }; diff --git a/apps/api/src/routes/workspaces/workspace-stop.ts b/apps/api/src/routes/workspaces/workspace-stop.ts index ceb5737c8..ece5a0ad4 100644 --- a/apps/api/src/routes/workspaces/workspace-stop.ts +++ b/apps/api/src/routes/workspaces/workspace-stop.ts @@ -7,6 +7,7 @@ import type { Env } from '../../env'; import { log } from '../../lib/logger'; import { getUserId, requireApproved, requireAuth } from '../../middleware/auth'; import { errors } from '../../middleware/error'; +import { purgeInteractionStore } from '../../services/acp-interaction-store'; import { stopComputeTracking } from '../../services/compute-usage'; import { stopWorkspaceOnNode } from '../../services/node-agent'; import { stopNodeResources } from '../../services/nodes'; @@ -102,8 +103,13 @@ workspaceStopRoutes.post('/:id/stop', requireAuth(), requireApproved(), async (c let runtimeStopConfirmed = retryStopCleanup; try { if (isCfContainerNode) { - if (workspace.chatSessionId) { - await deleteSessionSnapshotState(innerDb, c.env, workspace.chatSessionId); + if (workspace.chatSessionId && workspace.projectId) { + const chatSessionId = workspace.chatSessionId; + const projectId = workspace.projectId; + await Promise.all([ + deleteSessionSnapshotState(innerDb, c.env, chatSessionId), + purgeInteractionStore(c.env, projectId, chatSessionId), + ]); } if (node.status === 'running' && isActiveWorkspaceStatus(workspace.status)) { await stopWorkspaceOnNode(nodeId, workspace.id, c.env, userId).catch((e) => { diff --git a/apps/api/src/services/acp-interaction-config.ts b/apps/api/src/services/acp-interaction-config.ts new file mode 100644 index 000000000..d122e4292 --- /dev/null +++ b/apps/api/src/services/acp-interaction-config.ts @@ -0,0 +1,177 @@ +import { + DEFAULT_ACP_INTERACTION_ANSWER_MAX_BYTES, + DEFAULT_ACP_INTERACTION_ANSWER_STRING_MAX_BYTES, + DEFAULT_ACP_INTERACTION_DEADLINE_MARGIN_MS, + DEFAULT_ACP_INTERACTION_DELIVERY_WINDOW_MS, + DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES, + DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM, + DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES, + DEFAULT_ACP_INTERACTION_MAX_DEADLINE_MS, + DEFAULT_ACP_INTERACTION_MAX_PENDING_PER_SESSION, + DEFAULT_ACP_INTERACTION_OPTION_NAME_MAX_CHARS, + DEFAULT_ACP_INTERACTION_OPTIONS_MAX_COUNT, + DEFAULT_ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS, + DEFAULT_ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS, + DEFAULT_ACP_INTERACTION_REQUEST_MAX_BYTES, + DEFAULT_ACP_INTERACTION_RETRY_DELAYS_MS, + DEFAULT_ACP_INTERACTION_RETRY_STEADY_MS, + DEFAULT_ACP_INTERACTION_SENSITIVE_PURGE_MS, + DEFAULT_ACP_INTERACTION_SNAPSHOT_LAST_SETTLED, + DEFAULT_ACP_INTERACTION_SUMMARY_LAST_SETTLED, + DEFAULT_ACP_INTERACTION_SUMMARY_RETENTION_MS, + DEFAULT_ACP_INTERACTIONS_ENABLED, +} from '@simple-agent-manager/shared'; + +export interface AcpInteractionConfigEnv { + ACP_INTERACTIONS_ENABLED?: string; + ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS?: string; + ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS?: string; + ACP_INTERACTION_MAX_DEADLINE_MS?: string; + ACP_INTERACTION_DEADLINE_MARGIN_MS?: string; + ACP_INTERACTION_MAX_PENDING_PER_SESSION?: string; + ACP_INTERACTION_REQUEST_MAX_BYTES?: string; + ACP_INTERACTION_OPTIONS_MAX_COUNT?: string; + ACP_INTERACTION_OPTION_NAME_MAX_CHARS?: string; + ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES?: string; + ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES?: string; + ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM?: string; + ACP_INTERACTION_ANSWER_MAX_BYTES?: string; + ACP_INTERACTION_ANSWER_STRING_MAX_BYTES?: string; + ACP_INTERACTION_RETRY_DELAYS_MS?: string; + ACP_INTERACTION_RETRY_STEADY_MS?: string; + ACP_INTERACTION_DELIVERY_WINDOW_MS?: string; + ACP_INTERACTION_SENSITIVE_PURGE_MS?: string; + ACP_INTERACTION_SUMMARY_RETENTION_MS?: string; + ACP_INTERACTION_SUMMARY_LAST_SETTLED?: string; + ACP_INTERACTION_SNAPSHOT_LAST_SETTLED?: string; +} + +export interface AcpInteractionConfig { + enabled: boolean; + permissionConversationDeadlineMs: number; + permissionTaskDeadlineMs: number; + maxDeadlineMs: number; + deadlineMarginMs: number; + maxPendingPerSession: number; + requestMaxBytes: number; + optionsMaxCount: number; + optionNameMaxChars: number; + formSchemaMaxBytes: number; + formSchemaMaxProperties: number; + formSchemaMaxEnum: number; + answerMaxBytes: number; + answerStringMaxBytes: number; + retryDelaysMs: number[]; + retrySteadyMs: number; + deliveryWindowMs: number; + sensitivePurgeMs: number; + summaryRetentionMs: number; + summaryLastSettled: number; + snapshotLastSettled: number; +} + +function envFlag(value: string | undefined, fallback: boolean): boolean { + if (value === undefined || value.trim() === '') return fallback; + return value.trim().toLowerCase() === 'true'; +} + +function positiveInt(value: string | undefined, fallback: number): number { + if (value === undefined || value.trim() === '') return fallback; + const parsed = Number.parseInt(value, 10); + return Number.isFinite(parsed) && parsed > 0 ? parsed : fallback; +} + +function positiveIntList(value: string | undefined, fallback: readonly number[]): number[] { + if (value === undefined || value.trim() === '') return [...fallback]; + const parsed = value + .split(',') + .map((item) => Number.parseInt(item.trim(), 10)) + .filter((item) => Number.isFinite(item) && item > 0); + return parsed.length > 0 ? parsed : [...fallback]; +} + +export function getAcpInteractionConfig(env: AcpInteractionConfigEnv): AcpInteractionConfig { + return { + enabled: envFlag(env.ACP_INTERACTIONS_ENABLED, DEFAULT_ACP_INTERACTIONS_ENABLED), + permissionConversationDeadlineMs: positiveInt( + env.ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS, + DEFAULT_ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS + ), + permissionTaskDeadlineMs: positiveInt( + env.ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS, + DEFAULT_ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS + ), + maxDeadlineMs: positiveInt( + env.ACP_INTERACTION_MAX_DEADLINE_MS, + DEFAULT_ACP_INTERACTION_MAX_DEADLINE_MS + ), + deadlineMarginMs: positiveInt( + env.ACP_INTERACTION_DEADLINE_MARGIN_MS, + DEFAULT_ACP_INTERACTION_DEADLINE_MARGIN_MS + ), + maxPendingPerSession: positiveInt( + env.ACP_INTERACTION_MAX_PENDING_PER_SESSION, + DEFAULT_ACP_INTERACTION_MAX_PENDING_PER_SESSION + ), + requestMaxBytes: positiveInt( + env.ACP_INTERACTION_REQUEST_MAX_BYTES, + DEFAULT_ACP_INTERACTION_REQUEST_MAX_BYTES + ), + optionsMaxCount: positiveInt( + env.ACP_INTERACTION_OPTIONS_MAX_COUNT, + DEFAULT_ACP_INTERACTION_OPTIONS_MAX_COUNT + ), + optionNameMaxChars: positiveInt( + env.ACP_INTERACTION_OPTION_NAME_MAX_CHARS, + DEFAULT_ACP_INTERACTION_OPTION_NAME_MAX_CHARS + ), + formSchemaMaxBytes: positiveInt( + env.ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES, + DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES + ), + formSchemaMaxProperties: positiveInt( + env.ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES, + DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES + ), + formSchemaMaxEnum: positiveInt( + env.ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM, + DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM + ), + answerMaxBytes: positiveInt( + env.ACP_INTERACTION_ANSWER_MAX_BYTES, + DEFAULT_ACP_INTERACTION_ANSWER_MAX_BYTES + ), + answerStringMaxBytes: positiveInt( + env.ACP_INTERACTION_ANSWER_STRING_MAX_BYTES, + DEFAULT_ACP_INTERACTION_ANSWER_STRING_MAX_BYTES + ), + retryDelaysMs: positiveIntList( + env.ACP_INTERACTION_RETRY_DELAYS_MS, + DEFAULT_ACP_INTERACTION_RETRY_DELAYS_MS + ), + retrySteadyMs: positiveInt( + env.ACP_INTERACTION_RETRY_STEADY_MS, + DEFAULT_ACP_INTERACTION_RETRY_STEADY_MS + ), + deliveryWindowMs: positiveInt( + env.ACP_INTERACTION_DELIVERY_WINDOW_MS, + DEFAULT_ACP_INTERACTION_DELIVERY_WINDOW_MS + ), + sensitivePurgeMs: positiveInt( + env.ACP_INTERACTION_SENSITIVE_PURGE_MS, + DEFAULT_ACP_INTERACTION_SENSITIVE_PURGE_MS + ), + summaryRetentionMs: positiveInt( + env.ACP_INTERACTION_SUMMARY_RETENTION_MS, + DEFAULT_ACP_INTERACTION_SUMMARY_RETENTION_MS + ), + summaryLastSettled: positiveInt( + env.ACP_INTERACTION_SUMMARY_LAST_SETTLED, + DEFAULT_ACP_INTERACTION_SUMMARY_LAST_SETTLED + ), + snapshotLastSettled: positiveInt( + env.ACP_INTERACTION_SNAPSHOT_LAST_SETTLED, + DEFAULT_ACP_INTERACTION_SNAPSHOT_LAST_SETTLED + ), + }; +} diff --git a/apps/api/src/services/acp-interaction-delivery.ts b/apps/api/src/services/acp-interaction-delivery.ts new file mode 100644 index 000000000..fc8b26090 --- /dev/null +++ b/apps/api/src/services/acp-interaction-delivery.ts @@ -0,0 +1,159 @@ +import { + type AcpInteractionAnswerDecision, + AcpRuntimeAnswerResponseSchema, + buildAcpInteractionAnswerPath, +} from '@simple-agent-manager/shared'; +import * as v from 'valibot'; + +import type { Env } from '../env'; +import { createModuleLogger } from '../lib/logger'; +import { NodeAgentHttpError, nodeAgentRequest } from './node-agent'; + +const log = createModuleLogger('acp_interaction_delivery'); + +export interface AcpInteractionDeliveryTarget { + projectId: string; + chatSessionId: string; + workspaceId: string; + nodeId: string; + userId: string; + agentSessionId: string; + runtimeIdentity: string; + runtime: string; +} + +export type AcpInteractionDeliveryResult = + | { outcome: 'confirmed'; runtimeStatus: 'consumed' | 'duplicate' } + | { outcome: 'interrupted'; reason: string } + | { outcome: 'unconfirmed'; reason: string }; + +interface TargetRow { + workspace_id: string; + user_id: string; + workspace_status: string; + node_id: string | null; + node_status: string | null; + node_runtime: string | null; + agent_version: string | null; + agent_session_id: string | null; + agent_session_status: string | null; + agent_session_updated_at: string | null; +} + +const TERMINAL_WORKSPACE_STATUSES = new Set(['stopping', 'stopped', 'evicted', 'deleted', 'error']); +const TERMINAL_NODE_STATUSES = new Set(['stopping', 'stopped', 'deleted', 'error']); +const TERMINAL_AGENT_SESSION_STATUSES = new Set(['completed', 'failed', 'error', 'stopped']); + +export async function resolveAcpInteractionDeliveryTarget( + env: Env, + projectId: string, + chatSessionId: string +): Promise< + | { status: 'ready'; target: AcpInteractionDeliveryTarget } + | { status: 'retry' | 'interrupted'; reason: string } +> { + const row = await env.DATABASE.prepare( + `SELECT w.id AS workspace_id, + w.user_id AS user_id, + w.status AS workspace_status, + w.node_id AS node_id, + n.status AS node_status, + n.runtime AS node_runtime, + n.agent_version AS agent_version, + a.id AS agent_session_id, + a.status AS agent_session_status, + a.updated_at AS agent_session_updated_at + FROM workspaces w + LEFT JOIN nodes n ON n.id = w.node_id + LEFT JOIN agent_sessions a ON a.workspace_id = w.id + WHERE w.project_id = ? AND w.chat_session_id = ? + ORDER BY w.updated_at DESC, a.created_at DESC + LIMIT 1` + ) + .bind(projectId, chatSessionId) + .first(); + + if (!row) return { status: 'interrupted', reason: 'target workspace missing' }; + if (TERMINAL_WORKSPACE_STATUSES.has(row.workspace_status)) { + return { status: 'interrupted', reason: `target workspace is ${row.workspace_status}` }; + } + if (!row.node_id) return { status: 'retry', reason: 'target workspace has no node' }; + if (TERMINAL_NODE_STATUSES.has(row.node_status ?? '')) { + return { status: 'interrupted', reason: 'target node is unavailable' }; + } + if (!row.agent_session_id) return { status: 'retry', reason: 'target agent session missing' }; + if (row.agent_session_status !== 'running') { + return TERMINAL_AGENT_SESSION_STATUSES.has(row.agent_session_status ?? '') + ? { status: 'interrupted', reason: `target agent session is ${row.agent_session_status}` } + : { status: 'retry', reason: `target agent session is ${row.agent_session_status}` }; + } + + return { + status: 'ready', + target: { + projectId, + chatSessionId, + workspaceId: row.workspace_id, + nodeId: row.node_id, + userId: row.user_id, + agentSessionId: row.agent_session_id, + runtimeIdentity: [ + row.node_id, + row.agent_session_id, + row.agent_version ?? 'legacy', + row.agent_session_updated_at ?? 'unknown', + ].join(':'), + runtime: row.node_runtime ?? 'vm', + }, + }; +} + +export async function deliverAcpInteractionAnswer( + env: Env, + target: AcpInteractionDeliveryTarget, + input: { + interactionId: string; + generation: string; + runtimeIdentity: string; + decision: AcpInteractionAnswerDecision; + } +): Promise { + if (target.runtimeIdentity !== input.runtimeIdentity) { + return { outcome: 'interrupted', reason: 'runtime identity changed before delivery' }; + } + try { + const raw = await nodeAgentRequest( + target.nodeId, + env, + buildAcpInteractionAnswerPath(target.workspaceId, target.agentSessionId, input.interactionId), + { + method: 'POST', + userId: target.userId, + workspaceId: target.workspaceId, + recoverContainerOnTimeout: false, + body: JSON.stringify({ + protocolVersion: 1, + interactionId: input.interactionId, + generation: input.generation, + runtimeIdentity: input.runtimeIdentity, + decision: input.decision, + }), + } + ); + const response = v.parse(AcpRuntimeAnswerResponseSchema, raw); + if (response.status === 'consumed' || response.status === 'duplicate') { + return { outcome: 'confirmed', runtimeStatus: response.status }; + } + return { outcome: 'interrupted', reason: response.status }; + } catch (error) { + if (error instanceof NodeAgentHttpError && error.statusCode === 404) { + return { outcome: 'interrupted', reason: 'runtime waiter missing' }; + } + log.warn('acp_interaction.answer_delivery_unconfirmed', { + interactionId: input.interactionId, + workspaceId: target.workspaceId, + error: error instanceof Error ? error.message : String(error), + }); + return { outcome: 'unconfirmed', reason: 'transport outcome unknown' }; + } +} diff --git a/apps/api/src/services/acp-interaction-store.ts b/apps/api/src/services/acp-interaction-store.ts new file mode 100644 index 000000000..ff742e783 --- /dev/null +++ b/apps/api/src/services/acp-interaction-store.ts @@ -0,0 +1,71 @@ +import type { + InteractionStore, + InteractionStoreAnswerInput, + InteractionStoreCreateInput, + InteractionStoreSettleInput, +} from '../durable-objects/interaction-store'; +import type { Env } from '../env'; + +function ownerName(projectId: string, chatSessionId: string): string { + return `${projectId}/${chatSessionId}`; +} + +export function getInteractionStore( + env: Pick, + projectId: string, + chatSessionId: string +): DurableObjectStub { + return env.INTERACTION_STORE.get( + env.INTERACTION_STORE.idFromName(ownerName(projectId, chatSessionId)) + ) as DurableObjectStub; +} + +export function createInteraction(env: Env, input: InteractionStoreCreateInput) { + return getInteractionStore(env, input.projectId, input.chatSessionId).create(input); +} + +export function settleInteraction(env: Env, input: InteractionStoreSettleInput) { + return getInteractionStore(env, input.projectId, input.chatSessionId).settle(input); +} + +export function answerInteraction(env: Env, input: InteractionStoreAnswerInput) { + return getInteractionStore(env, input.projectId, input.chatSessionId).answer(input); +} + +export function snapshotInteractions( + env: Env, + projectId: string, + chatSessionId: string, + cursor: string | null = null +) { + return getInteractionStore(env, projectId, chatSessionId).snapshot(cursor); +} + +export function getInteractionDetail( + env: Env, + projectId: string, + chatSessionId: string, + interactionId: string +) { + return getInteractionStore(env, projectId, chatSessionId).detail(interactionId); +} + +export function purgeInteractionStore(env: Env, projectId: string, chatSessionId: string) { + if (!env.INTERACTION_STORE) return undefined; + return getInteractionStore(env, projectId, chatSessionId).purge(); +} + +export function recordInteractionDelivery( + env: Env, + projectId: string, + chatSessionId: string, + interactionId: string, + outcome: 'confirmed' | 'unconfirmed' | 'interrupted', + error: string | null = null +) { + return getInteractionStore(env, projectId, chatSessionId).recordDelivery( + interactionId, + outcome, + error + ); +} diff --git a/apps/api/src/services/task-terminal-cleanup.ts b/apps/api/src/services/task-terminal-cleanup.ts index 8e38a6173..ffb17d5b4 100644 --- a/apps/api/src/services/task-terminal-cleanup.ts +++ b/apps/api/src/services/task-terminal-cleanup.ts @@ -4,6 +4,7 @@ import { drizzle } from 'drizzle-orm/d1'; import * as schema from '../db/schema'; import type { Env } from '../env'; import { log } from '../lib/logger'; +import { purgeInteractionStore } from './acp-interaction-store'; import { preserveFailedTaskWork, surfaceFailedTaskWorkLoss, @@ -100,7 +101,11 @@ export async function cleanupTerminalTaskResources( } const [workspace] = await db - .select({ chatSessionId: schema.workspaces.chatSessionId, userId: schema.workspaces.userId }) + .select({ + chatSessionId: schema.workspaces.chatSessionId, + userId: schema.workspaces.userId, + projectId: schema.workspaces.projectId, + }) .from(schema.workspaces) .where(eq(schema.workspaces.id, task.workspaceId)) .limit(1); @@ -142,9 +147,14 @@ export async function cleanupTerminalTaskResources( }); } - if (workspace?.chatSessionId && options.destructiveSessionEnd) { + const cleanupProjectId = workspace?.projectId ?? task.projectId; + if (workspace?.chatSessionId && cleanupProjectId && options.destructiveSessionEnd) { + const chatSessionId = workspace.chatSessionId; await options.beforeSideEffect?.(); - await deleteSessionSnapshotState(db, env, workspace.chatSessionId); + await Promise.all([ + deleteSessionSnapshotState(db, env, chatSessionId), + purgeInteractionStore(env, cleanupProjectId, chatSessionId), + ]); } if (workspace?.chatSessionId) { diff --git a/apps/api/src/services/workspace-cleanup.ts b/apps/api/src/services/workspace-cleanup.ts index 23ac9edb3..4a96dc84e 100644 --- a/apps/api/src/services/workspace-cleanup.ts +++ b/apps/api/src/services/workspace-cleanup.ts @@ -5,6 +5,7 @@ import { type drizzle } from 'drizzle-orm/d1'; import * as schema from '../db/schema'; import type { Env } from '../env'; import { log } from '../lib/logger'; +import { purgeInteractionStore } from './acp-interaction-store'; import { deleteSessionSnapshotState } from './session-snapshots'; import { attemptWorkspaceDeletion, @@ -213,11 +214,17 @@ export async function cleanupWorkspaceForDeletion( source: workspaceDeletionLogSource(logContext), mode: 'explicit', allowWorkspaceNeverStartedProof: workspace.status === 'pending', - beforeFinalize: workspace.chatSessionId - ? async () => { - await deleteSessionSnapshotState(db, env, workspace.chatSessionId as string); - } - : undefined, + beforeFinalize: + workspace.chatSessionId && workspace.projectId + ? async () => { + const chatSessionId = workspace.chatSessionId as string; + const projectId = workspace.projectId as string; + await Promise.all([ + deleteSessionSnapshotState(db, env, chatSessionId), + purgeInteractionStore(env, projectId, chatSessionId), + ]); + } + : undefined, }); if (outcome.status === 'confirmed') { diff --git a/apps/api/tests/acp-interaction-delivery.test.ts b/apps/api/tests/acp-interaction-delivery.test.ts new file mode 100644 index 000000000..d52ed4b09 --- /dev/null +++ b/apps/api/tests/acp-interaction-delivery.test.ts @@ -0,0 +1,104 @@ +import { beforeEach, describe, expect, it, vi } from 'vitest'; + +const nodeAgentRequest = vi.fn(); + +vi.mock('../src/services/node-agent', () => ({ + NodeAgentHttpError: class NodeAgentHttpError extends Error { + constructor( + public readonly statusCode: number, + public readonly responseBody: string + ) { + super(`Node Agent request failed: ${statusCode} ${responseBody}`); + this.name = 'NodeAgentHttpError'; + } + }, + nodeAgentRequest, +})); + +const { deliverAcpInteractionAnswer } = await import('../src/services/acp-interaction-delivery'); +const { NodeAgentHttpError } = await import('../src/services/node-agent'); + +describe('ACP interaction answer delivery', () => { + const target = { + projectId: 'project-1', + chatSessionId: 'chat-1', + workspaceId: 'workspace-1', + nodeId: 'node-1', + userId: 'user-1', + agentSessionId: 'agent-session-1', + runtimeIdentity: 'runtime-1', + runtime: 'vm', + }; + const input = { + interactionId: '11111111-1111-4111-8111-111111111111', + generation: '22222222-2222-4222-8222-222222222222', + runtimeIdentity: 'runtime-1', + decision: { kind: 'accepted' as const, answerHash: 'a'.repeat(64) }, + }; + + beforeEach(() => nodeAgentRequest.mockReset()); + + it.each(['consumed', 'duplicate'] as const)( + 'confirms %s runtime receipts without waking', + async (status) => { + nodeAgentRequest.mockResolvedValue({ + status, + interactionId: input.interactionId, + generation: input.generation, + runtimeIdentity: input.runtimeIdentity, + }); + + await expect(deliverAcpInteractionAnswer({} as never, target, input)).resolves.toMatchObject({ + outcome: 'confirmed', + runtimeStatus: status, + }); + expect(nodeAgentRequest).toHaveBeenCalledWith( + 'node-1', + expect.anything(), + '/workspaces/workspace-1/agent-sessions/agent-session-1/interactions/11111111-1111-4111-8111-111111111111/answer', + expect.objectContaining({ recoverContainerOnTimeout: false, method: 'POST' }) + ); + } + ); + + it.each(['stale_generation', 'no_waiter', 'conflict'] as const)( + 'interrupts on %s runtime receipts', + async (status) => { + nodeAgentRequest.mockResolvedValue({ + status, + interactionId: input.interactionId, + generation: input.generation, + runtimeIdentity: input.runtimeIdentity, + }); + + await expect(deliverAcpInteractionAnswer({} as never, target, input)).resolves.toMatchObject({ + outcome: 'interrupted', + reason: status, + }); + } + ); + + it('interrupts dead runtime generations before sending', async () => { + await expect( + deliverAcpInteractionAnswer({} as never, target, { ...input, runtimeIdentity: 'old-runtime' }) + ).resolves.toMatchObject({ + outcome: 'interrupted', + reason: 'runtime identity changed before delivery', + }); + expect(nodeAgentRequest).not.toHaveBeenCalled(); + }); + + it('marks 404 as interrupted and ambiguous transport loss as unconfirmed', async () => { + nodeAgentRequest.mockRejectedValueOnce(new NodeAgentHttpError(404, 'missing')); + await expect(deliverAcpInteractionAnswer({} as never, target, input)).resolves.toMatchObject({ + outcome: 'interrupted', + reason: 'runtime waiter missing', + }); + + nodeAgentRequest.mockRejectedValueOnce(new Error('connection reset after write')); + await expect(deliverAcpInteractionAnswer({} as never, target, input)).resolves.toMatchObject({ + outcome: 'unconfirmed', + reason: 'transport outcome unknown', + }); + }); +}); diff --git a/apps/api/tests/unit/routes/chat-prompt-cancel.test.ts b/apps/api/tests/unit/routes/chat-prompt-cancel.test.ts index 8f40bc6d2..af8e040d6 100644 --- a/apps/api/tests/unit/routes/chat-prompt-cancel.test.ts +++ b/apps/api/tests/unit/routes/chat-prompt-cancel.test.ts @@ -16,6 +16,12 @@ const mocks = vi.hoisted(() => ({ enrichMessageWithMentions: vi.fn(), parseOptionalBody: vi.fn(), cancelScheduledSessionSleep: vi.fn(), + answerInteraction: vi.fn(), + deliverAcpInteractionAnswer: vi.fn(), + getInteractionDetail: vi.fn(), + recordInteractionDelivery: vi.fn(), + resolveAcpInteractionDeliveryTarget: vi.fn(), + snapshotInteractions: vi.fn(), })); vi.mock('drizzle-orm/d1', () => ({ @@ -100,6 +106,18 @@ vi.mock('../../../src/services/session-snapshots', () => ({ cancelScheduledSessionSleep: (...args: unknown[]) => mocks.cancelScheduledSessionSleep(...args), })); +vi.mock('../../../src/services/acp-interaction-store', () => ({ + answerInteraction: mocks.answerInteraction, + getInteractionDetail: mocks.getInteractionDetail, + recordInteractionDelivery: mocks.recordInteractionDelivery, + snapshotInteractions: mocks.snapshotInteractions, +})); + +vi.mock('../../../src/services/acp-interaction-delivery', () => ({ + deliverAcpInteractionAnswer: mocks.deliverAcpInteractionAnswer, + resolveAcpInteractionDeliveryTarget: mocks.resolveAcpInteractionDeliveryTarget, +})); + /** Helper to build a drizzle mock that returns workspace + agent session rows. */ function setupDrizzle(opts: { workspace?: { @@ -156,6 +174,34 @@ beforeEach(() => { createdByUserId: 'user-1', status: 'active', }); + mocks.snapshotInteractions.mockResolvedValue({ + pending: [ + { + interactionId: '11111111-1111-4111-8111-111111111111', + kind: 'permission', + state: 'pending', + createdAt: 1, + deadlineAt: 2, + safeSummary: { toolCallId: 'tool-1' }, + }, + ], + settled: [], + cursor: null, + }); + mocks.getInteractionDetail.mockResolvedValue({ + interactionId: '11111111-1111-4111-8111-111111111111', + kind: 'permission', + state: 'pending', + detail: { permissionName: 'SECRET_CANARY_PERMISSION' }, + }); + mocks.answerInteraction.mockResolvedValue({ + status: 'answered', + summary: { state: 'answered' }, + delivery: { generation: '22222222-2222-4222-8222-222222222222', runtimeIdentity: 'runtime-1' }, + }); + mocks.resolveAcpInteractionDeliveryTarget.mockResolvedValue({ status: 'interrupted', reason: 'no runtime' }); + mocks.deliverAcpInteractionAnswer.mockResolvedValue({ outcome: 'confirmed', runtimeStatus: 'consumed' }); + mocks.recordInteractionDelivery.mockResolvedValue(undefined); app = new Hono<{ Bindings: Env }>(); app.onError((err, c) => { @@ -885,3 +931,129 @@ describe('POST /sessions/:sessionId/cancel', () => { expect(mocks.cancelAgentSessionOnNode).not.toHaveBeenCalled(); }); }); + + +describe('ACP interaction browser routes', () => { + const interactionId = '11111111-1111-4111-8111-111111111111'; + const answerBody = { + answerKey: 'answer-key-1', + decision: { kind: 'accepted', answerHash: 'a'.repeat(64) }, + }; + const env = { DATABASE: {} as D1Database, BASE_DOMAIN: 'example.test' } as Env; + + it('requires exact configured browser Origin before accepting a human answer', async () => { + const response = await app.request( + `/api/projects/proj-1/sessions/chat-1/interactions/${interactionId}/answer`, + { + method: 'POST', + headers: { 'Content-Type': 'application/json', Origin: 'https://evil.example.test' }, + body: JSON.stringify(answerBody), + }, + env + ); + + expect(response.status).toBe(403); + expect(mocks.answerInteraction).not.toHaveBeenCalled(); + }); + + it('answers only through task:write plus session-creator authorization', async () => { + const response = await app.request( + `/api/projects/proj-1/sessions/chat-1/interactions/${interactionId}/answer`, + { + method: 'POST', + headers: { 'Content-Type': 'application/json', Origin: 'https://app.example.test' }, + body: JSON.stringify(answerBody), + }, + env + ); + + expect(response.status).toBe(200); + expect(mocks.requireProjectCapability).toHaveBeenCalledWith( + expect.anything(), + 'proj-1', + 'user-1', + 'task:write' + ); + expect(projectDataService.getSession).toHaveBeenCalledWith(expect.anything(), 'proj-1', 'chat-1'); + expect(mocks.answerInteraction).toHaveBeenCalledWith(expect.anything(), { + projectId: 'proj-1', + chatSessionId: 'chat-1', + interactionId, + answerKey: 'answer-key-1', + answerBodyHash: 'a'.repeat(64), + decision: answerBody.decision, + }); + }); + + it('rejects a project writer who did not create the session', async () => { + vi.mocked(projectDataService.getSession).mockResolvedValueOnce({ + id: 'chat-1', + createdByUserId: 'other-user', + status: 'active', + }); + + const response = await app.request( + `/api/projects/proj-1/sessions/chat-1/interactions/${interactionId}/answer`, + { + method: 'POST', + headers: { 'Content-Type': 'application/json', Origin: 'https://app.example.test' }, + body: JSON.stringify(answerBody), + }, + env + ); + + expect(response.status).toBe(403); + expect(mocks.answerInteraction).not.toHaveBeenCalled(); + }); + + it('returns decrypted detail only to the creator with no-store caching', async () => { + const response = await app.request( + `/api/projects/proj-1/sessions/chat-1/interactions/${interactionId}`, + { method: 'GET' }, + env + ); + + expect(response.status).toBe(200); + expect(response.headers.get('Cache-Control')).toBe('private, no-store'); + expect(mocks.requireProjectCapability).toHaveBeenCalledWith( + expect.anything(), + 'proj-1', + 'user-1', + 'task:read' + ); + await expect(response.json()).resolves.toMatchObject({ + detail: { permissionName: 'SECRET_CANARY_PERMISSION' }, + }); + }); + + it('keeps noncreator snapshots generic and omits settled history', async () => { + vi.mocked(projectDataService.getSession).mockResolvedValueOnce({ + id: 'chat-1', + createdByUserId: 'other-user', + status: 'active', + }); + + const response = await app.request( + '/api/projects/proj-1/sessions/chat-1/interactions', + { method: 'GET' }, + env + ); + + expect(response.status).toBe(200); + const body = await response.json(); + expect(body).toEqual({ + pending: [ + { + interactionId, + kind: 'permission', + state: 'pending', + createdAt: 1, + deadlineAt: 2, + }, + ], + settled: [], + cursor: null, + }); + expect(JSON.stringify(body)).not.toContain('tool-1'); + }); +}); diff --git a/apps/api/tests/unit/services/vm-prompt-delivery-adapter.test.ts b/apps/api/tests/unit/services/vm-prompt-delivery-adapter.test.ts index 699e22ee2..1d3f68542 100644 --- a/apps/api/tests/unit/services/vm-prompt-delivery-adapter.test.ts +++ b/apps/api/tests/unit/services/vm-prompt-delivery-adapter.test.ts @@ -58,6 +58,8 @@ const protocolFixture = JSON.parse( notReadyPrompt: Record; notFoundReceipt: Record; }; +const promptProtocolCapabilities = { ...protocolFixture.capabilities }; +delete promptProtocolCapabilities.interactions; const targetRow = { workspace_id: 'workspace-1', @@ -230,7 +232,7 @@ describe('VM prompt delivery adapter', () => { kind: 'retry', reason: 'not_ready', }); - expect(beforeSubmit).toHaveBeenCalledWith(protocolFixture.capabilities); + expect(beforeSubmit).toHaveBeenCalledWith(promptProtocolCapabilities); expect(mocks.sendPromptToAgentOnNode).not.toHaveBeenCalled(); } ); diff --git a/apps/api/tests/workers/acp-interaction-store.test.ts b/apps/api/tests/workers/acp-interaction-store.test.ts new file mode 100644 index 000000000..468258446 --- /dev/null +++ b/apps/api/tests/workers/acp-interaction-store.test.ts @@ -0,0 +1,203 @@ +import { env, runInDurableObject } from 'cloudflare:test'; +import { describe, expect, it } from 'vitest'; + +import type { InteractionStore } from '../../src/durable-objects/interaction-store'; +import type { Env } from '../../src/env'; + +const PROJECT_ID = 'project-acp-interactions'; +const INTERACTION_ID = '11111111-1111-4111-8111-111111111111'; +const GENERATION = '22222222-2222-4222-8222-222222222222'; +const HASH_A = 'a'.repeat(64); +const HASH_B = 'b'.repeat(64); + +function apiEnv(): Env { + return env as unknown as Env; +} + +function stub(name: string): DurableObjectStub { + const api = apiEnv(); + return api.INTERACTION_STORE.get(api.INTERACTION_STORE.idFromName(name)); +} + +function createChatSession(): string { + return `chat-${crypto.randomUUID()}`; +} + +async function withInteractionsEnabled(fn: () => Promise): Promise { + const mutableEnv = apiEnv() as unknown as Record; + const previous = mutableEnv.ACP_INTERACTIONS_ENABLED; + mutableEnv.ACP_INTERACTIONS_ENABLED = 'true'; + mutableEnv.ACP_INTERACTION_SENSITIVE_PURGE_MS = '1'; + try { + return await fn(); + } finally { + mutableEnv.ACP_INTERACTIONS_ENABLED = previous; + } +} + +function createInput(overrides: Partial[0]> = {}) { + return { + protocolVersion: 1 as const, + projectId: PROJECT_ID, + chatSessionId: overrides.chatSessionId ?? 'chat-acp-interactions', + interactionId: INTERACTION_ID, + generation: GENERATION, + runtimeIdentity: 'runtime-1', + agentSessionId: 'agent-session-1', + kind: 'permission' as const, + payloadHash: HASH_A, + detail: { + permissionName: 'SECRET_CANARY_PERMISSION', + options: [{ id: 'allow', label: 'Allow the sensitive operation' }], + }, + safeSummary: { toolCallId: 'tool-1', optionCount: 1 }, + deadlineAt: Date.now() + 30 * 60 * 1000, + upstreamRequestId: 'jsonrpc-diagnostic-1', + ...overrides, + }; +} + +async function clearDueWork(store: DurableObjectStub): Promise { + await runInDurableObject(store, async (_instance, state) => { + state.storage.sql.exec(`DELETE FROM outbox`); + await state.storage.deleteAlarm(); + }); +} + +describe('InteractionStore durable ACP foundation', () => { + it('is dormant by default and preserves existing records once disabled', async () => { + const store = stub(`disabled/${crypto.randomUUID()}`); + const chatSessionId = createChatSession(); + const disabled = await store.create(createInput({ chatSessionId })); + expect(disabled).toMatchObject({ status: 'disabled' }); + + await withInteractionsEnabled(async () => { + const created = await store.create( + createInput({ chatSessionId, interactionId: crypto.randomUUID() }) + ); + expect(created.status).toBe('created'); + await clearDueWork(store); + }); + + const snapshot = await store.snapshot(null); + expect(snapshot.pending).toHaveLength(1); + }); + + it('encrypts arbitrary detail, keeps only structural summaries, and rejects same id with another payload hash', async () => { + await withInteractionsEnabled(async () => { + const store = stub(`encrypted/${crypto.randomUUID()}`); + const chatSessionId = createChatSession(); + const created = await store.create(createInput({ chatSessionId })); + expect(created).toMatchObject({ status: 'created' }); + await clearDueWork(store); + + const conflict = await store.create(createInput({ chatSessionId, payloadHash: HASH_B })); + expect(conflict).toMatchObject({ status: 'conflict' }); + + const raw = await runInDurableObject(store, (_instance, state) => { + const row = state.storage.sql + .exec( + `SELECT encrypted_detail, detail_iv, safe_summary_json FROM interactions WHERE interaction_id = ?`, + INTERACTION_ID + ) + .one<{ encrypted_detail: string; detail_iv: string; safe_summary_json: string }>(); + return JSON.stringify(row); + }); + expect(raw).not.toContain('SECRET_CANARY_PERMISSION'); + expect(raw).not.toContain('Allow the sensitive operation'); + expect(raw).toContain('tool-1'); + + const detail = await store.detail(INTERACTION_ID); + expect(detail?.detail).toMatchObject({ permissionName: 'SECRET_CANARY_PERMISSION' }); + }); + }); + + it('commits answers atomically, binds idempotency keys to body hashes, and survives delivery loss states', async () => { + await withInteractionsEnabled(async () => { + const store = stub(`answer/${crypto.randomUUID()}`); + const chatSessionId = createChatSession(); + await store.create(createInput({ chatSessionId })); + await clearDueWork(store); + + const answered = await store.answer({ + projectId: PROJECT_ID, + chatSessionId, + interactionId: INTERACTION_ID, + answerKey: 'answer-key-1', + answerBodyHash: HASH_A, + decision: { kind: 'accepted', answerHash: HASH_A }, + }); + expect(answered).toMatchObject({ status: 'answered' }); + expect(answered.status === 'answered' ? answered.summary.state : null).toBe('answered'); + + const replay = await store.answer({ + projectId: PROJECT_ID, + chatSessionId, + interactionId: INTERACTION_ID, + answerKey: 'answer-key-1', + answerBodyHash: HASH_A, + decision: { kind: 'accepted', answerHash: HASH_A }, + }); + expect(replay.status).toBe('already_answered'); + + const mismatch = await store.answer({ + projectId: PROJECT_ID, + chatSessionId, + interactionId: INTERACTION_ID, + answerKey: 'answer-key-1', + answerBodyHash: HASH_B, + decision: { kind: 'declined', answerHash: HASH_B }, + }); + expect(mismatch).toMatchObject({ status: 'answer_key_conflict' }); + + await store.recordDelivery(INTERACTION_ID, 'unconfirmed', 'send timeout after commit'); + const snapshot = await store.snapshot(null); + expect(snapshot.settled[0]).toMatchObject({ state: 'delivery_unconfirmed' }); + }); + }); + + it('purges encrypted sensitive payloads while preserving bounded summaries', async () => { + await withInteractionsEnabled(async () => { + const store = stub(`purge/${crypto.randomUUID()}`); + const chatSessionId = createChatSession(); + await store.create(createInput({ chatSessionId })); + await clearDueWork(store); + await store.answer({ + projectId: PROJECT_ID, + chatSessionId, + interactionId: INTERACTION_ID, + answerKey: 'answer-key-1', + answerBodyHash: HASH_A, + decision: { + kind: 'accepted', + encryptedAnswer: { ciphertext: 'SECRET_CANARY_ANSWER', iv: 'iv' }, + answerHash: HASH_A, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 5)); + await runInDurableObject(store, async (instance) => instance.alarm()); + + const raw = await runInDurableObject(store, (_instance, state) => { + return state.storage.sql + .exec( + `SELECT encrypted_detail, detail_iv, encrypted_answer, answer_iv, detail_purged_at, safe_summary_json FROM interactions WHERE interaction_id = ?`, + INTERACTION_ID + ) + .one<{ + encrypted_detail: string | null; + detail_iv: string | null; + encrypted_answer: string | null; + answer_iv: string | null; + detail_purged_at: number | null; + safe_summary_json: string; + }>(); + }); + expect(raw.encrypted_detail).toBeNull(); + expect(raw.detail_iv).toBeNull(); + expect(raw.encrypted_answer).toBeNull(); + expect(raw.answer_iv).toBeNull(); + expect(raw.detail_purged_at).toBeGreaterThan(0); + expect(raw.safe_summary_json).toContain('tool-1'); + }); + }); +}); diff --git a/apps/api/vitest.workers.config.ts b/apps/api/vitest.workers.config.ts index 68dfe5cdf..e6fe3ef38 100644 --- a/apps/api/vitest.workers.config.ts +++ b/apps/api/vitest.workers.config.ts @@ -84,6 +84,10 @@ export default defineConfig({ className: 'DiagnosisRunner', useSQLite: true, }, + INTERACTION_STORE: { + className: 'InteractionStore', + useSQLite: true, + }, PROJECT_AGENT: { className: 'ProjectAgent', useSQLite: true, diff --git a/apps/api/wrangler.toml b/apps/api/wrangler.toml index 667e35e6b..263a82065 100644 --- a/apps/api/wrangler.toml +++ b/apps/api/wrangler.toml @@ -383,6 +383,11 @@ TASK_RECONCILIATION_NODE_CALL_TIMEOUT_MS = "5000" TASK_RECONCILIATION_CANDIDATE_LEASE_MS = "30000" TASK_RECONCILIATION_PROBE_MAX_ATTEMPTS = "3" TASK_RECONCILIATION_QUARANTINE_MS = "300000" +# Dormant durable ACP interaction foundation. Slice A keeps this disabled; later +# slices activate runtime/UI capability after end-to-end verification. Tuning +# defaults live in shared constants and apps/api/.env.example to avoid consuming +# production Worker text bindings unless an operator intentionally overrides them. +ACP_INTERACTIONS_ENABLED = "false" # AI Inference Proxy (Cloudflare AI Gateway for trial/zero-config users) AI_PROXY_ENABLED = "true" AI_PROXY_DEFAULT_MODEL = "@cf/qwen/qwen3-30b-a3b-fp8" @@ -547,6 +552,11 @@ class_name = "SamSession" name = "PROJECT_AGENT" class_name = "ProjectAgent" +# Durable Object for per-chat dormant ACP interaction state and answer outbox +[[durable_objects.bindings]] +name = "INTERACTION_STORE" +class_name = "InteractionStore" + # Durable Object for atomic per-user AI token budget accounting [[durable_objects.bindings]] name = "AI_TOKEN_BUDGET_COUNTER" @@ -667,6 +677,10 @@ new_sqlite_classes = ["SetupSessionPool"] # Singleton concurrency gate for guid tag = "v20" new_sqlite_classes = ["DiagnosisRunner"] +[[migrations]] +tag = "v21" +new_sqlite_classes = ["InteractionStore"] + # Analytics Engine for usage tracking (Phase 1 — server-side analytics) [[analytics_engine_datasets]] binding = "ANALYTICS" diff --git a/packages/shared/src/acp-interactions.ts b/packages/shared/src/acp-interactions.ts new file mode 100644 index 000000000..03493556c --- /dev/null +++ b/packages/shared/src/acp-interactions.ts @@ -0,0 +1,184 @@ +import * as v from 'valibot'; + +export const ACP_INTERACTION_PROTOCOL_VERSION = 1 as const; +export const ACP_INTERACTION_CAPABILITY_VERSION = 1 as const; + +export const ACP_INTERACTION_ATTENTION_SOURCE = 'acp_interaction' as const; +export const ACP_INTERACTION_KIND_VALUES = ['permission', 'form', 'url'] as const; +export const ACP_INTERACTION_STATE_VALUES = [ + 'pending', + 'answered', + 'delivery_confirmed', + 'delivery_unconfirmed', + 'interrupted', + 'expired', + 'cancelled', +] as const; +export const ACP_INTERACTION_DECISION_KIND_VALUES = [ + 'selected_option', + 'accepted', + 'declined', + 'cancelled', +] as const; +export const ACP_INTERACTION_DELIVERY_OUTCOME_VALUES = [ + 'consumed', + 'duplicate', + 'stale_generation', + 'no_waiter', + 'conflict', + 'unknown', +] as const; +export const ACP_RUNTIME_ANSWER_STATUS_VALUES = [ + 'consumed', + 'duplicate', + 'stale_generation', + 'no_waiter', + 'conflict', +] as const; + +export const DEFAULT_ACP_INTERACTIONS_ENABLED = false; +export const DEFAULT_ACP_INTERACTION_PERMISSION_CONVERSATION_DEADLINE_MS = 2 * 60 * 60 * 1000; +export const DEFAULT_ACP_INTERACTION_PERMISSION_TASK_DEADLINE_MS = 30 * 60 * 1000; +export const DEFAULT_ACP_INTERACTION_MAX_DEADLINE_MS = 4 * 60 * 60 * 1000; +export const DEFAULT_ACP_INTERACTION_DEADLINE_MARGIN_MS = 60 * 1000; +export const DEFAULT_ACP_INTERACTION_MAX_PENDING_PER_SESSION = 8; +export const DEFAULT_ACP_INTERACTION_REQUEST_MAX_BYTES = 32 * 1024; +export const DEFAULT_ACP_INTERACTION_OPTIONS_MAX_COUNT = 16; +export const DEFAULT_ACP_INTERACTION_OPTION_NAME_MAX_CHARS = 200; +export const DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_BYTES = 16 * 1024; +export const DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_PROPERTIES = 20; +export const DEFAULT_ACP_INTERACTION_FORM_SCHEMA_MAX_ENUM = 50; +export const DEFAULT_ACP_INTERACTION_ANSWER_MAX_BYTES = 16 * 1024; +export const DEFAULT_ACP_INTERACTION_ANSWER_STRING_MAX_BYTES = 4 * 1024; +export const DEFAULT_ACP_INTERACTION_RETRY_DELAYS_MS = [1000, 5000, 30000, 120000, 300000] as const; +export const DEFAULT_ACP_INTERACTION_RETRY_STEADY_MS = 300000; +export const DEFAULT_ACP_INTERACTION_DELIVERY_WINDOW_MS = 15 * 60 * 1000; +export const DEFAULT_ACP_INTERACTION_SENSITIVE_PURGE_MS = 60 * 60 * 1000; +export const DEFAULT_ACP_INTERACTION_SUMMARY_RETENTION_MS = 30 * 24 * 60 * 60 * 1000; +export const DEFAULT_ACP_INTERACTION_SUMMARY_LAST_SETTLED = 100; +export const DEFAULT_ACP_INTERACTION_SNAPSHOT_LAST_SETTLED = 20; +export const DEFAULT_ACP_INTERACTION_RUNTIME_RECEIPT_CAP = 256; + +export const AcpInteractionIdSchema = v.pipe(v.string(), v.uuid()); +export const AcpInteractionGenerationSchema = v.pipe(v.string(), v.uuid()); +export const AcpInteractionHashSchema = v.pipe(v.string(), v.regex(/^[a-f0-9]{64}$/u)); +export const AcpInteractionKeySchema = v.pipe(v.string(), v.minLength(1), v.maxLength(128)); + +export const AcpInteractionKindSchema = v.picklist(ACP_INTERACTION_KIND_VALUES); +export const AcpInteractionStateSchema = v.picklist(ACP_INTERACTION_STATE_VALUES); +export const AcpInteractionDecisionKindSchema = v.picklist(ACP_INTERACTION_DECISION_KIND_VALUES); +export const AcpRuntimeAnswerStatusSchema = v.picklist(ACP_RUNTIME_ANSWER_STATUS_VALUES); + +export const AcpInteractionOptionSchema = v.object({ + id: v.pipe(v.string(), v.minLength(1), v.maxLength(128)), + kind: v.picklist(['allow_once', 'allow_always', 'reject_once', 'reject_always', 'custom']), + name: v.pipe( + v.string(), + v.minLength(1), + v.maxLength(DEFAULT_ACP_INTERACTION_OPTION_NAME_MAX_CHARS) + ), +}); + +export const AcpInteractionSafeSummarySchema = v.object({ + interactionId: AcpInteractionIdSchema, + kind: AcpInteractionKindSchema, + state: AcpInteractionStateSchema, + createdAt: v.number(), + updatedAt: v.number(), + deadlineAt: v.number(), + answeredAt: v.nullable(v.number()), + deliveryState: v.nullable(v.picklist(['pending', 'confirmed', 'unconfirmed', 'interrupted'])), + attentionMarkerId: v.nullable(v.string()), + toolCallId: v.nullable(v.string()), +}); + +export const AcpInteractionEncryptedPayloadSchema = v.object({ + ciphertext: v.string(), + iv: v.string(), +}); + +export const AcpInteractionRuntimeCreateSchema = v.object({ + protocolVersion: v.literal(ACP_INTERACTION_PROTOCOL_VERSION), + interactionId: AcpInteractionIdSchema, + generation: AcpInteractionGenerationSchema, + runtimeIdentity: v.pipe(v.string(), v.minLength(1), v.maxLength(512)), + agentSessionId: v.pipe(v.string(), v.minLength(1), v.maxLength(256)), + kind: AcpInteractionKindSchema, + payloadHash: AcpInteractionHashSchema, + detail: v.unknown(), + safeSummary: v.object({ + toolCallId: v.optional(v.nullable(v.pipe(v.string(), v.maxLength(256)))), + optionCount: v.optional(v.number()), + }), + deadlineAt: v.number(), + upstreamRequestId: v.optional(v.nullable(v.pipe(v.string(), v.maxLength(256)))), +}); + +export const AcpInteractionRuntimeSettleSchema = v.object({ + protocolVersion: v.literal(ACP_INTERACTION_PROTOCOL_VERSION), + interactionId: AcpInteractionIdSchema, + generation: AcpInteractionGenerationSchema, + runtimeIdentity: v.pipe(v.string(), v.minLength(1), v.maxLength(512)), + agentSessionId: v.pipe(v.string(), v.minLength(1), v.maxLength(256)), + reason: v.picklist([ + 'wrapper_cancelled', + 'connection_closed', + 'connection_replaced', + 'session_stopped', + 'expired', + 'completed', + ]), +}); + +export const AcpInteractionAnswerDecisionSchema = v.object({ + kind: AcpInteractionDecisionKindSchema, + optionId: v.optional(v.string()), + encryptedAnswer: v.optional(AcpInteractionEncryptedPayloadSchema), + answerHash: AcpInteractionHashSchema, +}); + +export const AcpInteractionBrowserAnswerSchema = v.object({ + answerKey: AcpInteractionKeySchema, + decision: AcpInteractionAnswerDecisionSchema, +}); + +export const AcpRuntimeAnswerRequestSchema = v.object({ + protocolVersion: v.literal(ACP_INTERACTION_PROTOCOL_VERSION), + interactionId: AcpInteractionIdSchema, + generation: AcpInteractionGenerationSchema, + runtimeIdentity: v.pipe(v.string(), v.minLength(1), v.maxLength(512)), + decision: AcpInteractionAnswerDecisionSchema, +}); + +export const AcpRuntimeAnswerResponseSchema = v.object({ + status: AcpRuntimeAnswerStatusSchema, + interactionId: AcpInteractionIdSchema, + generation: AcpInteractionGenerationSchema, + runtimeIdentity: v.string(), +}); + +export type AcpInteractionKind = v.InferOutput; +export type AcpInteractionState = v.InferOutput; +export type AcpInteractionSafeSummary = v.InferOutput; +export type AcpInteractionRuntimeCreate = v.InferOutput; +export type AcpInteractionRuntimeSettle = v.InferOutput; +export type AcpInteractionAnswerDecision = v.InferOutput; +export type AcpInteractionBrowserAnswer = v.InferOutput; +export type AcpRuntimeAnswerRequest = v.InferOutput; +export type AcpRuntimeAnswerResponse = v.InferOutput; + +export interface AcpInteractionCapabilities { + version: typeof ACP_INTERACTION_CAPABILITY_VERSION; + answerEndpoint: boolean; + receiptCap: number; +} + +export function buildAcpInteractionAnswerPath( + workspaceId: string, + agentSessionId: string, + interactionId: string +): string { + return `/workspaces/${encodeURIComponent(workspaceId)}/agent-sessions/${encodeURIComponent( + agentSessionId + )}/interactions/${encodeURIComponent(interactionId)}/answer`; +} diff --git a/packages/shared/src/index.ts b/packages/shared/src/index.ts index 519249622..f228fe530 100644 --- a/packages/shared/src/index.ts +++ b/packages/shared/src/index.ts @@ -19,6 +19,9 @@ export * from './constants'; // VM Agent Contract (Zod schemas + types) export * from './vm-agent-contract'; +// Durable ACP interaction contracts (Valibot schemas + defaults) +export * from './acp-interactions'; + // Trial Onboarding (types + Valibot schemas) export * from './trial'; diff --git a/packages/vm-agent/internal/server/execution_protocol_test.go b/packages/vm-agent/internal/server/execution_protocol_test.go index 4098825a4..78ab12419 100644 --- a/packages/vm-agent/internal/server/execution_protocol_test.go +++ b/packages/vm-agent/internal/server/execution_protocol_test.go @@ -8,6 +8,7 @@ import ( "os" "path/filepath" "reflect" + "strings" "testing" "time" @@ -38,6 +39,7 @@ func TestExecutionProtocolRoutesRequireNodeManagementBearerToken(t *testing.T) { }{ {name: "capabilities", method: http.MethodGet, target: "/workspaces/ws-existing/agent-capabilities", handler: s.handleAgentCapabilities}, {name: "receipt", method: http.MethodGet, target: "/workspaces/ws-existing/agent-sessions/session/prompt-receipts/delivery", handler: s.handleGetPromptReceipt}, + {name: "interaction answer", method: http.MethodPost, target: "/workspaces/ws-existing/agent-sessions/session/interactions/11111111-1111-4111-8111-111111111111/answer", handler: s.handleAcpInteractionAnswer}, {name: "rollover submit", method: http.MethodPost, target: "/workspaces/ws-existing/agent-sessions/session/checkpoint-rollovers", handler: s.handleCheckpointRollover}, {name: "rollover lookup", method: http.MethodGet, target: "/workspaces/ws-existing/agent-sessions/session/checkpoint-rollovers/op", handler: s.handleGetCheckpointRollover}, } @@ -48,6 +50,7 @@ func TestExecutionProtocolRoutesRequireNodeManagementBearerToken(t *testing.T) { req.SetPathValue("sessionId", "session") req.SetPathValue("deliveryId", "delivery") req.SetPathValue("operationId", "op") + req.SetPathValue("interactionId", "11111111-1111-4111-8111-111111111111") rec := httptest.NewRecorder() tc.handler(rec, req) if rec.Code != http.StatusUnauthorized { @@ -57,6 +60,68 @@ func TestExecutionProtocolRoutesRequireNodeManagementBearerToken(t *testing.T) { } } +func TestAcpInteractionAnswerEndpointReturnsNoWaiterForDormantRuntime(t *testing.T) { + validator, privateKey := newWorkspaceCreateJWTValidator(t, "node-1") + s := newContractTestServer() + s.config.NodeID = "node-1" + s.jwtValidator = validator + s.executionRuntimeID = "runtime-vm-01" + token := signWorkspaceCreateNodeToken(t, privateKey, "node-1", "ws-existing") + + body := `{"protocolVersion":1,"interactionId":"11111111-1111-4111-8111-111111111111","generation":"22222222-2222-4222-8222-222222222222","runtimeIdentity":"runtime-vm-01","decision":{"kind":"permission","outcome":"approved"}}` + req := httptest.NewRequest(http.MethodPost, "/workspaces/ws-existing/agent-sessions/session/interactions/11111111-1111-4111-8111-111111111111/answer", strings.NewReader(body)) + req.SetPathValue("workspaceId", "ws-existing") + req.SetPathValue("sessionId", "session") + req.SetPathValue("interactionId", "11111111-1111-4111-8111-111111111111") + req.Header.Set("Authorization", "Bearer "+token) + req.Header.Set("X-SAM-Node-Id", "node-1") + req.Header.Set("X-SAM-Workspace-Id", "ws-existing") + + rec := httptest.NewRecorder() + s.handleAcpInteractionAnswer(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String()) + } + var response acpInteractionAnswerResponse + if err := json.NewDecoder(rec.Body).Decode(&response); err != nil { + t.Fatal(err) + } + if response.Status != "no_waiter" || response.RuntimeIdentity != "runtime-vm-01" { + t.Fatalf("response = %+v", response) + } +} + +func TestAcpInteractionAnswerEndpointRejectsStaleRuntimeIdentity(t *testing.T) { + validator, privateKey := newWorkspaceCreateJWTValidator(t, "node-1") + s := newContractTestServer() + s.config.NodeID = "node-1" + s.jwtValidator = validator + s.executionRuntimeID = "runtime-current" + token := signWorkspaceCreateNodeToken(t, privateKey, "node-1", "ws-existing") + + body := `{"protocolVersion":1,"interactionId":"11111111-1111-4111-8111-111111111111","generation":"22222222-2222-4222-8222-222222222222","runtimeIdentity":"runtime-old","decision":{"kind":"permission","outcome":"approved"}}` + req := httptest.NewRequest(http.MethodPost, "/workspaces/ws-existing/agent-sessions/session/interactions/11111111-1111-4111-8111-111111111111/answer", strings.NewReader(body)) + req.SetPathValue("workspaceId", "ws-existing") + req.SetPathValue("sessionId", "session") + req.SetPathValue("interactionId", "11111111-1111-4111-8111-111111111111") + req.Header.Set("Authorization", "Bearer "+token) + req.Header.Set("X-SAM-Node-Id", "node-1") + req.Header.Set("X-SAM-Workspace-Id", "ws-existing") + + rec := httptest.NewRecorder() + s.handleAcpInteractionAnswer(rec, req) + if rec.Code != http.StatusConflict { + t.Fatalf("status = %d body=%s", rec.Code, rec.Body.String()) + } + var response acpInteractionAnswerResponse + if err := json.NewDecoder(rec.Body).Decode(&response); err != nil { + t.Fatal(err) + } + if response.Status != "stale_generation" || response.RuntimeIdentity != "runtime-current" { + t.Fatalf("response = %+v", response) + } +} + func TestPromptReceiptObserverDurablyCompletesOnce(t *testing.T) { store, err := persistence.Open(filepath.Join(t.TempDir(), "receipts.db")) if err != nil { @@ -98,6 +163,10 @@ func TestAgentCapabilitiesAdvertiseVersionedInertExecutionFeatures(t *testing.T) if receipts["supported"] != true { t.Fatalf("prompt receipts = %#v", receipts) } + interactions := capabilities["interactions"].(map[string]interface{}) + if interactions["supported"] != true || interactions["version"] != acpInteractionCapabilityVersion { + t.Fatalf("interactions = %#v", interactions) + } rollover := capabilities["checkpointRollover"].(map[string]interface{}) if rollover["supported"] != true || rollover["automatic"] != false { t.Fatalf("checkpoint rollover = %#v", rollover) diff --git a/packages/vm-agent/internal/server/server.go b/packages/vm-agent/internal/server/server.go index f959983a5..ca0e835f4 100644 --- a/packages/vm-agent/internal/server/server.go +++ b/packages/vm-agent/internal/server/server.go @@ -1146,6 +1146,7 @@ func (s *Server) setupRoutes(mux *http.ServeMux) { mux.HandleFunc("POST /workspaces/{workspaceId}/agent-sessions/{sessionId}/resume", s.handleResumeAgentSession) mux.HandleFunc("POST /workspaces/{workspaceId}/agent-sessions/{sessionId}/prompt", s.handleSendPrompt) mux.HandleFunc("GET /workspaces/{workspaceId}/agent-sessions/{sessionId}/prompt-receipts/{deliveryId}", s.handleGetPromptReceipt) + mux.HandleFunc("POST /workspaces/{workspaceId}/agent-sessions/{sessionId}/interactions/{interactionId}/answer", s.handleAcpInteractionAnswer) mux.HandleFunc("POST /workspaces/{workspaceId}/agent-sessions/{sessionId}/checkpoint-rollovers", s.handleCheckpointRollover) mux.HandleFunc("GET /workspaces/{workspaceId}/agent-sessions/{sessionId}/checkpoint-rollovers/{operationId}", s.handleGetCheckpointRollover) mux.HandleFunc("GET /workspaces/{workspaceId}/agent-capabilities", s.handleAgentCapabilities) diff --git a/packages/vm-agent/internal/server/workspaces.go b/packages/vm-agent/internal/server/workspaces.go index 430c8092a..302aea1eb 100644 --- a/packages/vm-agent/internal/server/workspaces.go +++ b/packages/vm-agent/internal/server/workspaces.go @@ -23,9 +23,12 @@ import ( ) const ( - vmExecutionProtocolVersion = 1 - maxDeliveryIDLength = 128 - maxRolloverOperationIDLength = 128 + vmExecutionProtocolVersion = 1 + acpInteractionCapabilityVersion = 1 + maxDeliveryIDLength = 128 + maxRolloverOperationIDLength = 128 + maxAcpInteractionIDLength = 128 + maxAcpInteractionGenerationLength = 128 ) type sendPromptRequest struct { @@ -1810,6 +1813,78 @@ func (s *Server) writePromptReceiptNotFound(w http.ResponseWriter, deliveryID st }) } +type acpInteractionAnswerRequest struct { + ProtocolVersion int `json:"protocolVersion"` + InteractionID string `json:"interactionId"` + Generation string `json:"generation"` + RuntimeIdentity string `json:"runtimeIdentity"` + Decision json.RawMessage `json:"decision"` +} + +type acpInteractionAnswerResponse struct { + Status string `json:"status"` + InteractionID string `json:"interactionId"` + Generation string `json:"generation"` + RuntimeIdentity string `json:"runtimeIdentity"` +} + +func (s *Server) handleAcpInteractionAnswer(w http.ResponseWriter, r *http.Request) { + workspaceID := r.PathValue("workspaceId") + sessionID := r.PathValue("sessionId") + interactionID := strings.TrimSpace(r.PathValue("interactionId")) + if workspaceID == "" || sessionID == "" || interactionID == "" { + writeError(w, http.StatusBadRequest, "workspaceId, sessionId, and interactionId are required") + return + } + if !s.requireNodeManagementAuth(w, r, workspaceID) { + return + } + + var body acpInteractionAnswerRequest + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + writeError(w, http.StatusBadRequest, "invalid request body") + return + } + body.InteractionID = strings.TrimSpace(body.InteractionID) + body.Generation = strings.TrimSpace(body.Generation) + body.RuntimeIdentity = strings.TrimSpace(body.RuntimeIdentity) + if body.ProtocolVersion != vmExecutionProtocolVersion { + writeJSON(w, http.StatusBadRequest, map[string]interface{}{ + "error": "unsupported_protocol_version", + "supportedVersions": []int{vmExecutionProtocolVersion}, + }) + return + } + if body.InteractionID != interactionID || !validExecutionProtocolID(body.InteractionID, maxAcpInteractionIDLength) { + writeError(w, http.StatusBadRequest, "interactionId has an invalid format") + return + } + if !validExecutionProtocolID(body.Generation, maxAcpInteractionGenerationLength) { + writeError(w, http.StatusBadRequest, "generation has an invalid format") + return + } + if body.Decision == nil || !json.Valid(body.Decision) { + writeError(w, http.StatusBadRequest, "decision is required") + return + } + if body.RuntimeIdentity != s.executionRuntimeID { + writeJSON(w, http.StatusConflict, acpInteractionAnswerResponse{ + Status: "stale_generation", + InteractionID: body.InteractionID, + Generation: body.Generation, + RuntimeIdentity: s.executionRuntimeID, + }) + return + } + + writeJSON(w, http.StatusOK, acpInteractionAnswerResponse{ + Status: "no_waiter", + InteractionID: body.InteractionID, + Generation: body.Generation, + RuntimeIdentity: s.executionRuntimeID, + }) +} + func (s *Server) handleAgentCapabilities(w http.ResponseWriter, r *http.Request) { workspaceID := r.PathValue("workspaceId") if workspaceID == "" { @@ -1832,6 +1907,15 @@ func (s *Server) agentCapabilities() map[string]interface{} { "states": []string{persistence.PromptReceiptAccepted, persistence.PromptReceiptInFlight, persistence.PromptReceiptCompleted, persistence.PromptReceiptAmbiguous}, }, + "interactions": map[string]interface{}{ + "supported": true, + "version": acpInteractionCapabilityVersion, + "answerEndpoint": true, + "receiptCap": 256, + "deliverySemantics": "best_effort_no_wake", + "noWaiterStatus": "no_waiter", + "staleStatus": "stale_generation", + }, "checkpointRollover": map[string]interface{}{ "supported": true, "automatic": false, diff --git a/scripts/quality/do-migration-compatibility.test.ts b/scripts/quality/do-migration-compatibility.test.ts index 49b69fba4..67f7b1f0e 100644 --- a/scripts/quality/do-migration-compatibility.test.ts +++ b/scripts/quality/do-migration-compatibility.test.ts @@ -62,11 +62,11 @@ describe('resolveDurableObjectMigrations', () => { const resolved = resolveDurableObjectMigrations(history, null); - expect(history).toHaveLength(20); + expect(history).toHaveLength(21); expect(legacyCreateCount).toBe(7); expect(resolved).toHaveLength(history.length); expect(resolved.every((migration) => migration.new_classes === undefined)).toBe(true); - expect(resolved.flatMap((migration) => migration.new_sqlite_classes ?? [])).toHaveLength(20); + expect(resolved.flatMap((migration) => migration.new_sqlite_classes ?? [])).toHaveLength(21); expect(loadCheckedInMigrations()).toEqual(history); }); @@ -89,7 +89,7 @@ describe('resolveDurableObjectMigrations', () => { const latestTag = history.at(-1)?.tag; vi.stubEnv('RESOURCE_PREFIX', 's123abc'); - expect(latestTag).toBe('v20'); + expect(latestTag).toBe('v21'); const resolved = resolveDurableObjectMigrations(history, latestTag ?? null); const envConfig = generateApiWorkerEnv( { migrations: history }, diff --git a/scripts/quality/gitleaks-reviewed-baseline.json b/scripts/quality/gitleaks-reviewed-baseline.json index 8dd2909a6..5ccc342a1 100644 --- a/scripts/quality/gitleaks-reviewed-baseline.json +++ b/scripts/quality/gitleaks-reviewed-baseline.json @@ -206,6 +206,18 @@ "digests": [ "6e7d007b454a727e2d53a1ccfb0b9b3bf6b38081674abc35260efe50410cbc21" ] + }, + { + "classification": "synthetic-test-fixture", + "reason": "Reviewed existing synthetic Worker-test cryptographic fixtures in apps/api/vitest.workers.config.ts after adding the INTERACTION_STORE test binding shifted exact finding digests. These are deterministic local Miniflare/JWT/VAPID/encryption fixtures only, not live credentials.", + "owner": "security", + "reviewedAt": "2026-09-29T00:00:00.000Z", + "expiresAt": "2026-12-28T00:00:00.000Z", + "digests": [ + "54c818c7ade8518dca18c82bf022302d7e93251d161d7311130d92d63e919aa7", + "cebc591fa06ac4eb357e4355bd9685d7f9b717e68c0ca19f8a2d601b50f989c8", + "d1825d5702b57ea0ab87a15624a93992cd2ad7a00aecaf179bb0d1be452fe8fb" + ] } ] } diff --git a/tasks/active/2026-09-29-dormant-acp-interactions-foundation.md b/tasks/active/2026-09-29-dormant-acp-interactions-foundation.md new file mode 100644 index 000000000..f5c690191 --- /dev/null +++ b/tasks/active/2026-09-29-dormant-acp-interactions-foundation.md @@ -0,0 +1,126 @@ +# Dormant ACP Interactions Foundation + +## Problem + +ACP permission, form, and URL interactions need one authoritative Cloudflare path before any runtime/UI activation ships. Today the historical permission path can answer through runtime/browser-adjacent behavior or legacy attention prompt forwarding, which cannot preserve accepted human decisions, reconnectable history, generation fences, or no-wake delivery guarantees. + +Slice A builds the dormant foundation only. It must not advertise new ACP interaction capabilities, implement user-facing permission/form/URL UI, or activate runtime emission. Later slices will connect runtime permissions and UI after this storage, routing, encryption, and delivery base has shipped and been verified. + +## Authority And Scope + +- Canonical Idea `01M3P2E0JJNQRXX020P65ZRKEJ`, `APPROVED EXECUTION PLAN v2`, supersedes v1 and Fable findings where they differ. +- Parent task `01M3P223VQSKFD8GQAMV28R4PM`; completed review task `01M3P3EVJC1V886F75S34CWQM5`; recovered predecessor `01M3P2FVAHTRV80Y7FB4PDVSEM`. +- This task is executed with `/do`; user authorized a new green PR, CodeRabbit trusted-review path, merge, production deploy monitoring, and concise completion evidence. +- Do not dispatch SAM implementation subtasks. Local specialist reviewers are allowed. + +## Research Findings + +- Current `origin/main` and the output branch both start at `c17508d51`; no reconciliation changes were needed at task creation. +- The next Durable Object migration tag on current main is `v21`; `apps/api/wrangler.toml` currently ends at `v20` for `DiagnosisRunner`. +- New Durable Object bindings belong in the top-level `apps/api/wrangler.toml`; generated environment sections are produced by `scripts/deploy/sync-wrangler-config.ts`. +- Durable Object SQLite patterns exist in `apps/api/src/durable-objects/credential-setup-session/index.ts` and root ProjectData migrations in `apps/api/src/durable-objects/migrations.ts`. New InteractionStore should be separate and small, keyed by `projectId/chatSessionId`, not stored in root ProjectData. +- Existing runtime transport lives in `apps/api/src/services/node-agent.ts`. `nodeAgentRequest` already supports `recoverContainerOnTimeout: false`, which is required for no-wake interaction answer delivery. +- Existing prompt delivery target resolution in `apps/api/src/services/vm-prompt-delivery-target.ts` calls `ensureSessionRecovery`; answer delivery must use a separate resolver that never wakes VM or Instant sessions. +- VM-agent control-plane routes live in `packages/vm-agent/internal/server/server.go` and `workspaces.go`. New answer endpoints must use node-management auth and active `SessionHost` lookup, and must not feed answers into prompt delivery. +- Existing browser attention answer route is `apps/api/src/routes/chat.ts` at `POST /:sessionId/attention/:markerId/resolve`; it currently stages in ProjectData and forwards an answer as a prompt. Slice A must add a source guard so `source='acp_interaction'` markers cannot be answered through that route. +- ProjectData attention helpers are exposed from `apps/api/src/services/project-data.ts`. Projection markers can use `createAttentionMarker` with `expiresAt: null`, source `acp_interaction`, structural metadata only, and best-effort retries from InteractionStore outbox. +- Existing attention expiry logic in `apps/api/src/durable-objects/project-data/attention-expiry.ts` fails expired `needs_input` markers. `acp_interaction` markers use `expires_at NULL`, but add a source-aware guard as defense in depth. +- Browser authorization primitives exist: `requireProjectCapability(..., 'task:write')` and `requireSessionCreator(...)` in `apps/api/src/routes/chat-session-ownership.ts`. +- Callback JWT verification exists in `apps/api/src/services/jwt.ts` and callback helper patterns in `apps/api/src/routes/projects/_callback-auth.ts`; runtime create/settle routes must use workspace-scoped callback identity and resolve project/chat/session server-side from D1. +- Existing trusted origin derivation is in `apps/api/src/auth.ts` and `apps/api/src/lib/trusted-origins.ts`; new browser write routes need exact configured Origin comparison, not same-site-only authorization. +- Existing encryption helper is `apps/api/src/services/encryption.ts` and `getCredentialEncryptionKey(env)` in `apps/api/src/lib/secrets.ts`; arbitrary interaction detail and answer payloads can reuse that generated secret path with no manual secret prerequisite. +- Shared contract types currently live in `packages/shared/src/vm-agent-contract.ts`; new Valibot interaction schemas/constants should be exported from `packages/shared/src/index.ts` / `types`. + +## Implementation Checklist + +- [x] Add shared versioned ACP interaction schemas, constants, defaults, limits, and fixture corpus using Valibot and exported shared types. +- [x] Add typed Worker env/config resolution for `ACP_INTERACTIONS_ENABLED=false` and all V2 deadline, size, retry, retention, snapshot, and pending-session limits. +- [x] Add `InteractionStore` Durable Object binding/class and append-only `v21` `new_sqlite_classes = ["InteractionStore"]` migration with generated config compatibility tests. +- [x] Implement `InteractionStore` SQLite tables for interactions, outbox, delivery attempts, answer idempotency, bounded summaries, snapshots, sensitive purge, and due-work indexes. +- [x] Implement atomic create semantics: runtime-generated `interactionId`, canonical payload hash, same-id/same-hash idempotency, same-id/different-hash conflict, max pending enforcement, deadline validation, disabled/version-skew fail closed for new records while preserving serviceability of existing records. +- [x] Encrypt all arbitrary necessary request detail and human answer detail with existing Worker credential encryption; keep broad logs/events/markers structural only. +- [x] Implement answer semantics: answerKey bound to request+body hash, competing answer linearization, stable conflicts after first decision, decision state separate from delivery state, accepted decision surviving later delivery loss. +- [x] Implement bounded outbox alarms/retries: projection, delivery, settle/cancel/expire, sensitive purge, history compaction that never trims active rows, snapshot pending+last20 with pagination. +- [x] Implement Worker runtime create/settle routes with workspace callback JWT auth, server-side workspace/project/chat/agentSession binding, runtime identity/generation validation contract, structural logs only. +- [x] Implement Worker browser snapshot/detail/answer routes with session-cookie auth, `task:write`, session-creator-only mutation/detail, noncreator generic snapshot, exact Origin guard, no-store decrypted detail responses, and negative tests for runtime/MCP/callback tokens answering as humans. +- [x] Implement dedicated low-level answer delivery module using `nodeAgentRequest` only, with no prompt-delivery adapter, no `ensureSessionRecovery`, and `recoverContainerOnTimeout: false`; classify confirmed, interrupted, and delivery_unconfirmed outcomes honestly. +- [x] Add VM-agent low-level interaction answer endpoint and version capability consumer without activating interaction creation. Dormant endpoint returns `no_waiter`/`stale_generation`; consumed/duplicate in-memory waiter/tombstone registry is explicitly deferred to Slice B runtime waiter wiring because Slice A does not attach live ACP waiters. +- [x] Add minimal attention projection source `acp_interaction`, `expires_at NULL`, structural metadata only, best-effort nonblocking create/resolve, and source-aware expiry guard. +- [x] Add legacy attention resolve guard so `acp_interaction` markers cannot route an answer as a prompt. +- [x] Add actual session-delete cleanup hook to purge/cancel active interaction records while preserving bounded summaries for history-preserving archive. +- [x] Add focused tests for local-runtime DO state transitions/restart/outbox persistence, idempotency hash mismatches, answer/cancel/expire races, stale/dead generation, no-wake transport proof, auth/caller-type/CSRF negatives, canary secrecy, retention/deletion, attention source guard, fresh install and upgrade config. Evidence: node-level no-wake delivery tests, worker InteractionStore tests, VM route contract tests, migration compatibility tests, and ACP browser route guard tests added; CI Durable Object Workers and staging remain the authoritative Cloudflare runtime proof. +- [x] Update docs/API contract/env references as needed without advertising runtime/UI capability activation. +- [ ] Run required quality gates, local specialist reviews, staging proof, CodeRabbit, merge, production deploy/version monitoring, and append concise A outcome to the canonical Idea. Progress: local gates, specialist review evidence, Sonar, PR, and CodeRabbit label path complete; latest CI/staging/merge/prod evidence pending. + +## Acceptance Criteria + +- Dormant default: `ACP_INTERACTIONS_ENABLED=false`; no ACP permission/form/URL capability is advertised to live sessions by this slice. +- Cloudflare Worker can create/read/answer/settle controlled fixture interactions through InteractionStore, and accepted answers are committed before callback delivery. +- Arbitrary request details and answers are encrypted at rest; canary tests prove secrets are absent from plaintext DO/D1 rows, logs/events/notifications/unauthorized responses. +- Browser human answer route is session-creator-only with fresh `task:write` authorization and exact Origin protection; runtime/MCP/callback tokens cannot answer as humans. +- Runtime create/settle routes are bound to workspace callback identity and server-resolved project/chat/agentSession; caller-provided identity alone is never trusted. +- Answer delivery never invokes prompt delivery, prompt preparation, session recovery, or Instant wake paths; tests assert no-wake behavior. +- State machine keeps decision, delivery, expiry, cancel, interrupted, and delivery_unconfirmed outcomes distinct; tests cover competing answers and lost/ambiguous receipts. +- Attention projection failures never block canonical reads/answers; legacy attention resolve rejects `acp_interaction` source. +- Fresh install and upgrade generated Cloudflare config include the new binding/migration safely. +- PR passes local quality gates, local security/Cloudflare/constitution/env/doc/task-completion reviews, staging controlled fixture proof, CI, CodeRabbit trusted review loop, merge, production deployment monitoring, and bounded dormant production smoke. + + +## Implementation Evidence So Far + +- Added `InteractionStore` Durable Object with encrypted request detail and encrypted answer/decision storage, per-chat deterministic service wrapper, answer idempotency/body-hash binding, delivery state separation, alarm-driven projection/delivery/purge/compaction, and session cleanup hooks. +- Added shared Valibot contracts/defaults and VM-agent contract fixture updates. +- Added Worker runtime callback create/settle routes and browser snapshot/detail/answer routes with exact Origin guard and session creator gating. +- Added no-wake answer delivery service using `nodeAgentRequest(..., recoverContainerOnTimeout: false)` and tests for consumed/duplicate/stale/no-waiter/conflict/404/ambiguous transport outcomes. +- Added source-safe `acp_interaction` attention projection plus legacy attention resolve/expiry guard. +- Added dormant `ACP_INTERACTIONS_ENABLED=false` wrangler flag; other ACP interaction tuning defaults are typed/documented and resolved in code to avoid exceeding Cloudflare Worker text-binding guard. +- Local checks passing so far: `pnpm typecheck`, `pnpm lint` (pre-existing warnings only), `pnpm --filter @simple-agent-manager/api test -- tests/acp-interaction-delivery.test.ts`, `pnpm --filter @simple-agent-manager/shared typecheck`, and `pnpm vitest run scripts/quality/do-migration-compatibility.test.ts scripts/quality/go-toolchain-floor.test.ts scripts/quality/check-runtime-boundary-semantics.test.ts`. +- Local limitations: Go toolchain/gofmt are unavailable in this container; Cloudflare worker tests stall at startup here even for an existing attention-marker test, so worker runtime proof needs CI/staging confirmation. + + +## Task Completion Validation Report + +**Task**: `tasks/active/2026-09-29-dormant-acp-interactions-foundation.md` +**Branch**: `sam/implement-ship-slice-dormant-bdptty` +**Date**: 2026-09-29 + +### Verdict: PASS with one scoped WARN + +| Check | Status | Issues | +| --- | --- | --- | +| A: Research → Checklist | PASS | All research findings have checklist coverage or scoped deferral. | +| B: Checklist → Diff | PASS | Checked items map to shared schemas, Worker routes, InteractionStore DO, VM endpoint, attention guards, cleanup hooks, env/docs, and tests. | +| C: Criteria → Tests | PASS/WARN | Automated coverage exists for store state, encryption canaries, idempotency, no-wake delivery, VM dormant endpoint, route guards, and migration config; live staging proof still pending. | +| D: UI → Backend | N/A | Slice A adds no UI inputs. | +| E: Multi-Resource | N/A | No provider/resource selector added. | +| F: Vertical Slice | PASS/WARN | Worker DO tests cover Cloudflare store behavior; staging remains required for deployed dormant route/config proof. | + +### Findings + +#### WARN F: VM in-memory waiter registry deferred to Slice B + +**Planned** (task file checklist): VM-agent low-level endpoint plus in-memory receipt/tombstone registry. + +**Actual**: Slice A implements the dormant VM endpoint and capability consumer. Runtime consumed/duplicate waiter state is not reachable until Slice B wires live ACP waiters, so the registry is deferred to Slice B. The Worker delivery module is still tested against fake runtime `consumed`, `duplicate`, `stale_generation`, `no_waiter`, `conflict`, dead-generation, and ambiguous-transport outcomes. + +**Risk**: None while Slice A remains dormant; later B must add the runtime waiter/tombstone registry before activating runtime-generated interactions. + +**Recommendation**: Preserve this as a Slice B acceptance item; do not implement it in A because A must not activate runtime interaction waiters. + +### Uncovered Acceptance Criteria + +| Criterion | Test or verification | Status | +| --- | --- | --- | +| Dormant default and no advertised live ACP creation | shared defaults/config, VM capability fixture, prompt-delivery fixture subset test | COVERED | +| Worker create/read/answer/settle controlled fixtures | `apps/api/tests/workers/acp-interaction-store.test.ts`; staging pending | COVERED/PENDING STAGING | +| Encryption and sensitive canaries | `apps/api/tests/workers/acp-interaction-store.test.ts` | COVERED | +| Browser human answer auth + exact Origin | `apps/api/tests/unit/routes/chat-prompt-cancel.test.ts` ACP route cases | COVERED | +| Runtime callback identity | `apps/api/src/routes/projects/acp-interaction-callback.ts` plus existing callback auth patterns; staging pending | COVERED/PENDING STAGING | +| No-wake delivery | `apps/api/tests/acp-interaction-delivery.test.ts` | COVERED | +| Decision/delivery race states | `apps/api/tests/workers/acp-interaction-store.test.ts` plus delivery tests | COVERED | +| Attention source guard | source guard code and route reject path; staging pending | COVERED/PENDING STAGING | +| Generated Cloudflare config | `scripts/quality/do-migration-compatibility.test.ts`; CI/staging pending | COVERED/PENDING STAGING | + +### UI-to-Backend Data Path Audit + +No UI inputs were added in Slice A. diff --git a/tests/fixtures/durable-execution-protocol-v1.json b/tests/fixtures/durable-execution-protocol-v1.json index 2f9150fbe..c6321e999 100644 --- a/tests/fixtures/durable-execution-protocol-v1.json +++ b/tests/fixtures/durable-execution-protocol-v1.json @@ -7,6 +7,15 @@ "lookup": true, "states": ["accepted", "in_flight", "completed", "ambiguous"] }, + "interactions": { + "supported": true, + "version": 1, + "answerEndpoint": true, + "receiptCap": 256, + "deliverySemantics": "best_effort_no_wake", + "noWaiterStatus": "no_waiter", + "staleStatus": "stale_generation" + }, "checkpointRollover": { "supported": true, "automatic": false,