import assert from 'node:assert/strict' import { mkdtemp, rm } from 'node:fs/promises' import net from 'node:net' import { tmpdir } from 'node:os' import { join } from 'node:path' import { it } from 'node:test' import { CredentialClient, createSession } from '../../atproto/dist/index.js' import { COLLECTIONS } from '../../core/dist/index.js' import { LocalPds } from '../../atproto/test/local-pds.mjs' import { TurnSocketServer, answerRkey, implArtifactRkey, planArtifactRkey, reviewRkey } from '../dist/index.js' const REQUEST = { uri: 'at://did:plc:human/com.disnetdev.radial.artifactRequest/req1', cid: 'reqcid', goal: { uri: 'at://did:plc:human/com.disnetdev.radial.goal/g1', cid: 'goalcid' }, } const TOKEN = 'turn-token-1' const NOW = '2026-07-19T00:00:00.000Z' const SUBJECT = { uri: 'at://did:plc:agent/com.disnetdev.radial.artifact/impl-1', cid: 'subjcid' } const REVIEW_MODE = { mode: { kind: 'review', subject: SUBJECT } } function sendLine(socketPath, obj) { return new Promise((resolve, reject) => { const socket = net.createConnection(socketPath) let data = '' socket.on('connect', () => socket.write(`${JSON.stringify(obj)}\n`)) socket.on('data', (chunk) => { data += chunk.toString('utf8') }) socket.on('end', () => { try { resolve(JSON.parse(data)) } catch (error) { reject(new Error(`Non-JSON response: ${data} (${error.message})`)) } }) socket.on('error', reject) }) } async function makeClient(pds) { const session = await createSession(pds.service, pds.handle, 'pw', pds.fetch.bind(pds)) return new CredentialClient(session, pds.fetch.bind(pds)) } async function startServer(directory, client, options = {}) { const socketPath = join(directory, `turn-${Math.random().toString(36).slice(2)}.sock`) // Default to a plan-style artifact turn; individual tests override `mode` via options.context. const context = { token: TOKEN, request: REQUEST, client, now: () => NOW, mode: { kind: 'artifact', type: 'plan' }, ...options.context } const server = new TurnSocketServer(socketPath, context, options.limits ? { limits: options.limits } : {}) await server.listen() return { server, socketPath, context } } async function withTempDir(run) { const directory = await mkdtemp(join(tmpdir(), 'radial-turn-socket-')) try { await run(directory) } finally { await rm(directory, { recursive: true, force: true }) } } it('submitArtifact writes a plan artifact anchored to the request and goal', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent1') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { model: () => 'anthropic/claude-opus-4-8' } }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'the plan body', criteria: ['criterion one'], }) assert.equal(response.ok, true) const expectedRkey = planArtifactRkey(REQUEST.uri, REQUEST.cid) assert.equal(response.ref.uri, `at://${pds.did}/${COLLECTIONS.artifact}/${expectedRkey}`) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.artifact) assert.equal(stored.length, 1) const record = stored[0].value assert.deepEqual(record.request, { uri: REQUEST.uri, cid: REQUEST.cid }) assert.deepEqual(record.goal, REQUEST.goal) assert.equal(record.type, 'plan') assert.equal(record.title, 'A short title') assert.equal(record.body, 'the plan body') assert.deepEqual(record.links, {}) assert.deepEqual(record.criteria, ['criterion one']) assert.equal(record.model, 'anthropic/claude-opus-4-8') assert.equal(record.createdAt, NOW) assert.ok(stored[0].uri.endsWith(`/${expectedRkey}`)) } finally { await server.close() } }) }) it('requires a usable title at the untrusted edge, and writes nothing without one', async () => { // The lexicon leaves `title` optional so historical records stay valid; the WRITER requires it, and // the socket is the writer for every turn. Each refusal here leaves the turn retryable: nothing is // written, and the in-flight reservation is not taken, so a corrected submit still lands. await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-title') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client) try { const missing = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, body: 'no title' }) assert.equal(missing.ok, false) assert.match(missing.error, /requires --title/) const blank = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: ' \n ', body: 'x' }) assert.equal(blank.ok, false) assert.match(blank.error, /non-blank/) const long = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'x'.repeat(201), body: 'x', }) assert.equal(long.ok, false) assert.match(long.error, /200 characters/) // A title that is not a string at all is a malformed envelope, refused before anything else. const wrongType = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: { short: 'no' }, body: 'x', }) assert.equal(wrongType.ok, false) assert.match(wrongType.error, /invalid turn RPC envelope/) assert.equal(pds.records.size, 0) // Not a spent turn: the same connection sequence ends in a real artifact once a title is given. const good = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: ' Four moving parts \n', body: 'x', }) assert.equal(good.ok, true) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.artifact) assert.equal(stored.length, 1) // Stored trimmed — the trimmed form is what every reader compares and displays. assert.equal(stored[0].value.title, 'Four moving parts') assert.equal(stored[0].value.model, undefined) } finally { await server.close() } }) }) it('refuses to adopt a deterministic record whose title or model differs', async () => { // The title is part of the payload, so it is part of what makes the record at the request's rkey // the one this submit would have written. Adopting on a title mismatch would report an artifact as // this turn's own while the space showed a different name for it. await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-title2') const client = await makeClient(pds) const first = await startServer(directory, client, { context: { model: () => 'model-a' } }) const firstResponse = await sendLine(first.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'The title that landed', body: 'same body', }) await first.server.close() assert.equal(firstResponse.ok, true) const same = await startServer(directory, client, { context: { model: () => 'model-a' } }) const sameResponse = await sendLine(same.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'The title that landed', body: 'same body', }) await same.server.close() assert.equal(sameResponse.ok, true) assert.deepEqual(sameResponse.ref, firstResponse.ref) const renamed = await startServer(directory, client, { context: { model: () => 'model-a' } }) const renamedResponse = await sendLine(renamed.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A different title', body: 'same body', }) await renamed.server.close() assert.equal(renamedResponse.ok, false) assert.match(renamedResponse.error, /refusing to adopt tampered record/) const rerouted = await startServer(directory, client, { context: { model: () => 'model-b' } }) const reroutedResponse = await sendLine(rerouted.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'The title that landed', body: 'same body', }) await rerouted.server.close() assert.equal(reroutedResponse.ok, false) assert.match(reroutedResponse.error, /refusing to adopt tampered record/) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.artifact) assert.equal(stored.length, 1) assert.equal(stored[0].value.title, 'The title that landed') assert.equal(stored[0].value.model, 'model-a') }) }) it('askQuestion writes a message anchored to the request goal', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent2') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { model: () => 'openai/gpt-5.6-sol' } }) try { const response = await sendLine(socketPath, { method: 'askQuestion', token: TOKEN, body: 'why?' }) assert.equal(response.ok, true) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.message) assert.equal(stored.length, 1) const record = stored[0].value assert.deepEqual(record.goal, REQUEST.goal) assert.deepEqual(record.re, { uri: REQUEST.uri, cid: REQUEST.cid }) assert.deepEqual(record.mentions, []) assert.equal(record.body, 'why?') assert.equal(record.model, 'openai/gpt-5.6-sol') assert.equal(record.createdAt, NOW) } finally { await server.close() } }) }) it('rejects an oversized artifact body and writes no record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent3') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { limits: { maxFrameBytes: 1_000_000, frameTimeoutMs: 5000, artifactBodyChars: 10, messageBodyChars: 10 }, }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'this body is definitely longer than ten characters', }) assert.equal(response.ok, false) assert.equal(pds.records.size, 0) } finally { await server.close() } }) }) it('rejects a frame that exceeds the hard byte cap and writes no record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent4') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { limits: { maxFrameBytes: 40, frameTimeoutMs: 5000, artifactBodyChars: 100_000, messageBodyChars: 100_000 }, }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'a modestly sized body that is well under the per-field limit but the frame as a whole exceeds the tiny cap', }) assert.equal(response.ok, false) assert.match(response.error, /frame exceeds/) assert.equal(pds.records.size, 0) } finally { await server.close() } }) }) it('accepts only one terminal artifact per request', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent5') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client) try { const first = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'first' }) assert.equal(first.ok, true) const second = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'second' }) assert.equal(second.ok, false) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.artifact) assert.equal(stored.length, 1) assert.equal(stored[0].value.body, 'first') } finally { await server.close() } }) }) it('adopts the existing record on RecordAlreadyExists instead of duplicating it', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent6') const client = await makeClient(pds) const first = await startServer(directory, client) const firstResponse = await sendLine(first.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'idempotent body', }) await first.server.close() assert.equal(firstResponse.ok, true) // A fresh server/context for the same request simulates a retried turn: the daemon // recomputes the same deterministic rkey and must adopt rather than error. const second = await startServer(directory, client) const secondResponse = await sendLine(second.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'idempotent body', }) await second.server.close() assert.equal(secondResponse.ok, true) assert.deepEqual(secondResponse.ref, firstResponse.ref) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.artifact) assert.equal(stored.length, 1) }) }) it('rejects a mismatched turn token and writes no record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent7') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: 'wrong-token', title: 'A short title', body: 'should not land', }) assert.equal(response.ok, false) assert.match(response.error, /token/) assert.equal(pds.records.size, 0) } finally { await server.close() } }) }) it('tcp transport: submitArtifact over a raw TCP connection to boundPort writes a plan artifact (no Docker needed)', async () => { const pds = new LocalPds('did:plc:agent-tcp') const client = await makeClient(pds) const context = { token: TOKEN, request: REQUEST, client, now: () => NOW, mode: { kind: 'artifact', type: 'plan' } } const server = new TurnSocketServer({ kind: 'tcp', host: '127.0.0.1', port: 0 }, context) await server.listen() try { const port = server.boundPort assert.equal(typeof port, 'number') const response = await new Promise((resolve, reject) => { const socket = net.createConnection({ host: '127.0.0.1', port }) let data = '' socket.on('connect', () => socket.write( `${JSON.stringify({ method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'the plan body', criteria: ['criterion one'] })}\n`, ), ) socket.on('data', (chunk) => { data += chunk.toString('utf8') }) socket.on('end', () => { try { resolve(JSON.parse(data)) } catch (error) { reject(new Error(`Non-JSON response: ${data} (${error.message})`)) } }) socket.on('error', reject) }) assert.equal(response.ok, true) const expectedRkey = planArtifactRkey(REQUEST.uri, REQUEST.cid) assert.equal(response.ref.uri, `at://${pds.did}/${COLLECTIONS.artifact}/${expectedRkey}`) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.artifact) assert.equal(stored.length, 1) assert.equal(stored[0].value.body, 'the plan body') assert.equal(stored[0].value.type, 'plan') } finally { await server.close() } }) it('planArtifactRkey is deterministic and a valid record key', () => { const rkey = planArtifactRkey(REQUEST.uri, REQUEST.cid) assert.match(rkey, /^plan-[a-z2-7]+$/) assert.equal(rkey, planArtifactRkey(REQUEST.uri, REQUEST.cid)) assert.notEqual(rkey, planArtifactRkey(REQUEST.uri, 'other-cid')) }) it('implementation submit synchronously writes the deterministic artifact with validated links and prev', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-impl') const client = await makeClient(pds) const prev = { uri: 'at://did:plc:agent/com.disnetdev.radial.artifact/impl-old', cid: 'oldcid' } const branch = 'radial/impl-123456789abc' const { server, socketPath } = await startServer(directory, client, { context: { mode: { kind: 'artifact', type: 'implementation' }, implementation: { gitUrl: 'https://github.com/acme/widget.git', branch }, prev }, }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'done', branch, commit: 'a'.repeat(40), pr: 'https://github.com/acme/widget/pull/7' }) assert.equal(response.ok, true) assert.equal(response.ref.uri, `at://${pds.did}/${COLLECTIONS.artifact}/${implArtifactRkey(REQUEST.uri, REQUEST.cid)}`) const record = [...pds.records.values()].find((entry) => entry.value.$type === COLLECTIONS.artifact).value assert.deepEqual(record.links, { branch, commit: 'a'.repeat(40), pr: 'https://github.com/acme/widget/pull/7' }) assert.deepEqual(record.prev, prev) } finally { await server.close() } }) }) it('implementation submit rejects the wrong branch, noncanonical PR URLs, and cross-repository PRs', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-impl-validation') const client = await makeClient(pds) const branch = 'radial/impl-123456789abc' const { server, socketPath } = await startServer(directory, client, { context: { mode: { kind: 'artifact', type: 'implementation' }, implementation: { gitUrl: 'https://github.com/acme/widget.git', branch }, }, }) const submit = (overrides) => sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'done', branch, commit: 'a'.repeat(40), pr: 'https://github.com/acme/widget/pull/7', ...overrides, }) try { const wrongBranch = await submit({ branch: 'radial/impl-wrong' }) assert.equal(wrongBranch.ok, false) assert.match(wrongBranch.error, /daemon-selected/) const noncanonical = await submit({ pr: 'https://github.com/acme/widget/pull/7?diff=split' }) assert.equal(noncanonical.ok, false) assert.match(noncanonical.error, /canonical/) const crossRepo = await submit({ pr: 'https://github.com/attacker/widget/pull/7' }) assert.equal(crossRepo.ok, false) assert.match(crossRepo.error, /not for this project/) assert.equal(pds.records.size, 0) } finally { await server.close() } }) }) it('implementation adoption refuses an existing deterministic record with different criteria', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-impl-adoption') const client = await makeClient(pds) const branch = 'radial/impl-123456789abc' const context = { mode: { kind: 'artifact', type: 'implementation' }, implementation: { gitUrl: 'https://github.com/acme/widget.git', branch }, } const frame = { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'done', branch, commit: 'a'.repeat(40), pr: 'https://github.com/acme/widget/pull/7', } const first = await startServer(directory, client, { context }) const firstResponse = await sendLine(first.socketPath, { ...frame, criteria: ['original'] }) await first.server.close() assert.equal(firstResponse.ok, true) const second = await startServer(directory, client, { context }) try { const secondResponse = await sendLine(second.socketPath, { ...frame, criteria: ['tampered'] }) assert.equal(secondResponse.ok, false) assert.match(secondResponse.error, /refusing to adopt tampered record/) } finally { await second.server.close() } }) }) // --- Review turns --------------------------------------------------------- const findings = [{ severity: 'error', path: 'src/x.ts', line: 12, body: 'Bug here.' }] it('submitReview: writes a verdict anchored to the context subject and request backref, at the deterministic review- rkey', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-review') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { ...REVIEW_MODE, model: () => 'openai/gpt-5.6-sol' } }) try { const response = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'request_changes', findings, }) assert.equal(response.ok, true) const expectedRkey = reviewRkey(REQUEST.uri, REQUEST.cid) assert.match(expectedRkey, /^review-[a-z2-7]+$/) assert.equal(response.ref.uri, `at://${pds.did}/${COLLECTIONS.review}/${expectedRkey}`) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.review) assert.equal(stored.length, 1) const record = stored[0].value // subject + request come from CONTEXT, never the envelope. assert.deepEqual(record.subject, SUBJECT) assert.deepEqual(record.request, { uri: REQUEST.uri, cid: REQUEST.cid }) assert.equal(record.verdict, 'request_changes') assert.deepEqual(record.findings, findings) assert.equal(record.model, 'openai/gpt-5.6-sol') assert.equal(record.createdAt, NOW) } finally { await server.close() } }) }) it('submitReview: an envelope-supplied subject/request is ignored — the daemon anchors from context', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-review-anchor') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: REVIEW_MODE }) try { const response = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [], subject: { uri: 'at://did:plc:attacker/x/y', cid: 'evil' }, request: { uri: 'at://did:plc:attacker/a/b', cid: 'evil2' }, }) assert.equal(response.ok, true) const record = [...pds.records.values()].find((r) => r.value.$type === COLLECTIONS.review).value assert.deepEqual(record.subject, SUBJECT) assert.deepEqual(record.request, { uri: REQUEST.uri, cid: REQUEST.cid }) } finally { await server.close() } }) }) it('mutual exclusion: submitReview is rejected on an artifact turn, submitArtifact is rejected on a review turn', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-mutex') const client = await makeClient(pds) // Artifact (plan) turn: submitReview rejected. const artifactTurn = await startServer(directory, client) const reviewOnArtifact = await sendLine(artifactTurn.socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [], }) await artifactTurn.server.close() assert.equal(reviewOnArtifact.ok, false) assert.match(reviewOnArtifact.error, /not a review turn/) // Review turn: submitArtifact rejected. const reviewTurn = await startServer(directory, client, { context: REVIEW_MODE }) const artifactOnReview = await sendLine(reviewTurn.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'sneaky artifact', }) await reviewTurn.server.close() assert.equal(artifactOnReview.ok, false) assert.match(artifactOnReview.error, /review turn/) assert.equal(pds.records.size, 0) }) }) it('submitReview: rejects a bad verdict, an invalid severity, and >100 findings; writes no record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-review-validation') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: REVIEW_MODE }) try { const badVerdict = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'lgtm', findings: [] }) assert.equal(badVerdict.ok, false) const badSeverity = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [{ severity: 'critical', body: 'x' }], }) assert.equal(badSeverity.ok, false) const oversizedBody = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [{ severity: 'info', body: 'x'.repeat(10_001) }], }) assert.equal(oversizedBody.ok, false) const tooMany = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: Array.from({ length: 101 }, () => ({ severity: 'info', body: 'x' })), }) assert.equal(tooMany.ok, false) assert.equal(pds.records.size, 0) } finally { await server.close() } }) }) it('submitReview: adoption on retry with DIFFERENT findings succeeds (same subject/request), converging on one record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-review-adopt') const client = await makeClient(pds) const first = await startServer(directory, client, { context: REVIEW_MODE }) const firstResponse = await sendLine(first.socketPath, { method: 'submitReview', token: TOKEN, verdict: 'request_changes', findings: [{ severity: 'warning', body: 'first pass' }], }) await first.server.close() assert.equal(firstResponse.ok, true) // A retried review turn legitimately produces a different verdict + findings; adoption keys on // subject/request only, so it converges on the already-written record. const second = await startServer(directory, client, { context: REVIEW_MODE }) const secondResponse = await sendLine(second.socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [{ severity: 'info', body: 'second pass, different' }], }) await second.server.close() assert.equal(secondResponse.ok, true) assert.deepEqual(secondResponse.ref, firstResponse.ref) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.review) assert.equal(stored.length, 1) // The adopted record is the ORIGINAL (findings/verdict are not overwritten on adoption). assert.equal(stored[0].value.verdict, 'request_changes') }) }) it('submitReview: adoption refuses an existing review- record with a mismatched subject', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-review-mismatch') const client = await makeClient(pds) const first = await startServer(directory, client, { context: REVIEW_MODE }) const firstResponse = await sendLine(first.socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [] }) await first.server.close() assert.equal(firstResponse.ok, true) // A second review turn for the same request but a DIFFERENT subject would collide on the same // deterministic rkey — it must refuse to adopt rather than serve a verdict of the wrong subject. const otherSubject = { uri: 'at://did:plc:agent/com.disnetdev.radial.artifact/impl-2', cid: 'othercid' } const second = await startServer(directory, client, { context: { mode: { kind: 'review', subject: otherSubject } } }) try { const secondResponse = await sendLine(second.socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [] }) assert.equal(secondResponse.ok, false) assert.match(secondResponse.error, /mismatched subject\/request/) } finally { await second.server.close() } const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.review) assert.equal(stored.length, 1) assert.deepEqual(stored[0].value.subject, SUBJECT) }) }) it('submitReview: rejects an over-long path (>2000) and a line < 1; writes no record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-review-finding-validation') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: REVIEW_MODE }) try { const longPath = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [{ severity: 'info', body: 'x', path: 'p'.repeat(2001) }], }) assert.equal(longPath.ok, false) const zeroLine = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [{ severity: 'info', body: 'x', line: 0 }], }) assert.equal(zeroLine.ok, false) // A valid boundary (path exactly 2000, line exactly 1) is accepted. const ok = await sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [{ severity: 'info', body: 'x', path: 'p'.repeat(2000), line: 1 }], }) assert.equal(ok.ok, true) } finally { await server.close() } }) }) // --- Project-scoped (system artifact) turns ------------------------------- const PROJECT = { uri: 'at://did:plc:human/com.disnetdev.radial.project/proj1', cid: 'projcid' } const PROJECT_REQUEST = { uri: 'at://did:plc:human/com.disnetdev.radial.artifactRequest/adr-req', cid: 'adr-req-cid', project: PROJECT, } const SOURCE_GOAL = { uri: 'at://did:plc:human/com.disnetdev.radial.goal/src', cid: 'srcgoalcid' } const BASED_ON_ARTIFACT = { uri: 'at://did:plc:agent/com.disnetdev.radial.artifact/aaa', cid: 'aaacid' } it('project-scoped submitArtifact writes an artifact anchored to the project (no goal)', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-adr') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { request: PROJECT_REQUEST, mode: { kind: 'artifact', type: 'adr' } }, }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'the adr' }) assert.equal(response.ok, true) const record = [...pds.records.values()].find((r) => r.value.$type === COLLECTIONS.artifact).value assert.deepEqual(record.project, PROJECT) assert.equal(record.goal, undefined) assert.equal(record.type, 'adr') assert.equal(record.body, 'the adr') } finally { await server.close() } }) }) it('project-scoped submitArtifact adopts on retry, comparing the project anchor', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-adr-adopt') const client = await makeClient(pds) const context = { request: PROJECT_REQUEST, mode: { kind: 'artifact', type: 'adr' } } const first = await startServer(directory, client, { context }) const firstResponse = await sendLine(first.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'adr body' }) await first.server.close() assert.equal(firstResponse.ok, true) const second = await startServer(directory, client, { context }) const secondResponse = await sendLine(second.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'adr body' }) await second.server.close() assert.equal(secondResponse.ok, true) assert.deepEqual(secondResponse.ref, firstResponse.ref) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.artifact) assert.equal(stored.length, 1) }) }) it('project-scoped submitArtifact refuses to adopt an existing record carrying a stray second anchor', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-dual-anchor') const client = await makeClient(pds) const rkey = planArtifactRkey(PROJECT_REQUEST.uri, PROJECT_REQUEST.cid) // Pre-seed a record at the deterministic rkey that carries BOTH project AND a stray goal — the // materializer would resolve it to the goal (goal takes precedence), so adopting it for a project // request would leave that request open forever while the ledger reported it fulfilled. Injected // straight into the store (client.create validates exactly-one-anchor and would reject it). const seededUri = `at://${pds.did}/${COLLECTIONS.artifact}/${rkey}` pds.records.set(seededUri, { uri: seededUri, cid: 'seed-dual-anchor-cid', value: { $type: COLLECTIONS.artifact, request: { uri: PROJECT_REQUEST.uri, cid: PROJECT_REQUEST.cid }, project: PROJECT, goal: { uri: 'at://did:plc:human/com.disnetdev.radial.goal/stray', cid: 'straycid' }, type: 'adr', body: 'the adr', links: {}, createdAt: NOW, }, }) const { server, socketPath } = await startServer(directory, client, { context: { request: PROJECT_REQUEST, mode: { kind: 'artifact', type: 'adr' } }, }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'the adr' }) // Adoption is refused. Two guards enforce this: getOwnRecord validates the fetched record // (exactly-one anchor) and rejects a dual-anchor record on read, and #submitArtifact's adoption // comparison additionally requires the stray anchor to be absent (defense in depth if read-side // validation is ever relaxed). Either way the socket does not adopt and reports failure. assert.equal(response.ok, false) assert.match(response.error, /refusing to adopt tampered record|exactly one of goal or project/) // The stray/seeded record was not handed back as an accepted artifact ref. assert.equal(response.ref, undefined) } finally { await server.close() } }) }) it('project-scoped askQuestion precedence: capture source goal anchor', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-adr-q1') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { request: PROJECT_REQUEST, mode: { kind: 'artifact', type: 'adr' }, questionAnchor: { goal: SOURCE_GOAL } }, }) try { const response = await sendLine(socketPath, { method: 'askQuestion', token: TOKEN, body: 'which format?' }) assert.equal(response.ok, true) const record = [...pds.records.values()].find((r) => r.value.$type === COLLECTIONS.message).value assert.deepEqual(record.goal, SOURCE_GOAL) // anchored to the capture source goal assert.equal(record.artifact, undefined) assert.deepEqual(record.re, { uri: PROJECT_REQUEST.uri, cid: PROJECT_REQUEST.cid }) } finally { await server.close() } }) }) it('project-scoped askQuestion precedence: falls back to the basedOn artifact anchor', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-adr-q2') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { request: PROJECT_REQUEST, mode: { kind: 'artifact', type: 'adr' }, questionAnchor: { artifact: BASED_ON_ARTIFACT } }, }) try { const response = await sendLine(socketPath, { method: 'askQuestion', token: TOKEN, body: 'which format?' }) assert.equal(response.ok, true) const record = [...pds.records.values()].find((r) => r.value.$type === COLLECTIONS.message).value assert.deepEqual(record.artifact, BASED_ON_ARTIFACT) assert.equal(record.goal, undefined) } finally { await server.close() } }) }) it('project-scoped askQuestion with no anchor is rejected with a clear error and writes no record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-adr-q3') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { request: PROJECT_REQUEST, mode: { kind: 'artifact', type: 'adr' } }, // no questionAnchor }) try { const response = await sendLine(socketPath, { method: 'askQuestion', token: TOKEN, body: 'which format?' }) assert.equal(response.ok, false) assert.match(response.error, /project-scoped turn without basedOn cannot ask questions/) assert.equal(pds.records.size, 0) } finally { await server.close() } }) }) it('submitReview: two concurrent submits on one review turn land exactly one record (reservation)', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-review-concurrent') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: REVIEW_MODE }) try { const [a, b] = await Promise.all([ sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [] }), sendLine(socketPath, { method: 'submitReview', token: TOKEN, verdict: 'request_changes', findings: [] }), ]) const oks = [a, b].filter((r) => r.ok) assert.equal(oks.length, 1) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.review) assert.equal(stored.length, 1) } finally { await server.close() } }) }) // ── answer turns ──────────────────────────────────────────────────────────────────────────────── // The reply is the one terminal record a human reads directly, so everything about it is // daemon-owned: the thread anchoring, the backref that closes the request, and who it mentions. const ANSWER_SUBJECT = { uri: 'at://did:plc:human/com.disnetdev.radial.message/q1', cid: 'qcid' } const ANSWER_ROOT = { uri: 'at://did:plc:human/com.disnetdev.radial.message/root', cid: 'rootcid' } const ANSWER_MODE = { mode: { kind: 'answer', subject: ANSWER_SUBJECT, // The subject's OWN anchor, which the dispatcher resolved — not the request's goal. Here they // are the same record; the artifact-anchored case is below. anchor: { goal: REQUEST.goal }, subjectAuthor: 'did:plc:human', }, } it('submitAnswer writes a reply threaded under the subject, at the deterministic answer- rkey', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { ...ANSWER_MODE, model: () => 'anthropic/claude-sonnet-4-8' } }) try { const response = await sendLine(socketPath, { method: 'submitAnswer', token: TOKEN, body: 'It needed a second store to keep in sync.', }) assert.equal(response.ok, true) const expectedRkey = answerRkey(REQUEST.uri, REQUEST.cid) assert.match(expectedRkey, /^answer-[a-z2-7]+$/) assert.equal(response.ref.uri, `at://${pds.did}/${COLLECTIONS.message}/${expectedRkey}`) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.message) assert.equal(stored.length, 1) const record = stored[0].value assert.deepEqual(record.goal, REQUEST.goal) assert.deepEqual(record.parent, ANSWER_SUBJECT) // The subject starts the thread, so it is its own root. assert.deepEqual(record.root, ANSWER_SUBJECT) assert.deepEqual(record.re, { uri: REQUEST.uri, cid: REQUEST.cid }) assert.deepEqual(record.mentions, ['did:plc:human']) assert.equal(record.body, 'It needed a second store to keep in sync.') assert.equal(record.model, 'anthropic/claude-sonnet-4-8') assert.equal(record.createdAt, NOW) } finally { await server.close() } }) }) it('submitAnswer threads under an existing root when the subject is not the root', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer-root') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: { mode: { ...ANSWER_MODE.mode, root: ANSWER_ROOT } }, }) try { await sendLine(socketPath, { method: 'submitAnswer', token: TOKEN, body: 'threaded' }) const record = [...pds.records.values()].find((r) => r.value.$type === COLLECTIONS.message).value assert.deepEqual(record.parent, ANSWER_SUBJECT) assert.deepEqual(record.root, ANSWER_ROOT) } finally { await server.close() } }) }) it('submitAnswer anchors the reply where the subject is, not where the request is', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer-anchor') const client = await makeClient(pds) // A message anchored to an ARTIFACT under the request's goal: the fold files it into that goal's // thread all the same, so it is a message a reader sees and can be asked about. The reply // belongs where it does — `message post --parent` refuses a human's reply anchored anywhere // else, and the daemon must not write a record it would have refused. const artifact = { uri: 'at://did:plc:agent/com.disnetdev.radial.artifact/plan-1', cid: 'acid' } const { server, socketPath } = await startServer(directory, client, { context: { mode: { ...ANSWER_MODE.mode, anchor: { artifact } } }, }) try { const response = await sendLine(socketPath, { method: 'submitAnswer', token: TOKEN, body: 'Because the second read is what makes it safe.', }) assert.equal(response.ok, true) const record = [...pds.records.values()].find((r) => r.value.$type === COLLECTIONS.message).value assert.deepEqual(record.artifact, artifact) assert.equal(record.goal, undefined, 'exactly one anchor, like every other message') assert.deepEqual(record.parent, ANSWER_SUBJECT) assert.deepEqual(record.re, { uri: REQUEST.uri, cid: REQUEST.cid }) } finally { await server.close() } }) }) // --- daemon-brokered pull requests (tangled) -------------------------------- /** A minimal adapter that opens the pull request itself, as `TangledForge` does. */ function brokeringForge(overrides = {}) { const opened = [] return { opened, adapter: { kind: 'tangled', canObserve: true, matchesGitUrl: () => true, ownsPullUrl: () => true, pullBelongsToProject: async () => true, projectIdentity: async () => 'did:plc:owner', getPullRequestState: async () => ({ state: 'open', headRef: '', headRepoFullName: '', baseRef: 'main' }), compareUrl: () => '', async openPullRequest(input) { opened.push(input) return { url: 'at://did:plc:agent/sh.tangled.repo.pull/3ktid1', record: { uri: 'at://did:plc:agent/sh.tangled.repo.pull/3ktid1', cid: 'pullcid' } } }, ...overrides, }, } } const TANGLED_IMPL = (forge, branch) => ({ mode: { kind: 'artifact', type: 'implementation' }, implementation: { gitUrl: 'https://tangled.org/@alice/widget', branch, base: 'main', baseCommit: 'b'.repeat(40), checkoutPath: '/run/checkout', artifactUri: 'at://did:plc:agent/com.disnetdev.radial.artifact/impl-1', pullTitle: 'Ship the thing', forge, }, }) it('a brokered submit needs no --pr: the daemon opens the pull and stamps the link it wrote', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-tangled') const client = await makeClient(pds) const branch = 'radial/impl-123456789abc' const { adapter, opened } = brokeringForge() const { server, socketPath } = await startServer(directory, client, { context: TANGLED_IMPL(adapter, branch) }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'done', branch, commit: 'a'.repeat(40), }) assert.equal(response.ok, true) const record = [...pds.records.values()].find((entry) => entry.value.$type === COLLECTIONS.artifact).value assert.deepEqual(record.links, { branch, commit: 'a'.repeat(40), pr: 'at://did:plc:agent/sh.tangled.repo.pull/3ktid1', }) // Everything the adapter needs to generate a patch and title the pull, from the daemon. assert.equal(opened.length, 1) assert.equal(opened[0].base, 'main') assert.equal(opened[0].baseCommit, 'b'.repeat(40)) assert.equal(opened[0].checkoutPath, '/run/checkout') assert.equal(opened[0].title, 'Ship the thing') // The same provenance a GitHub turn is told to put in its own pull request body. assert.match(opened[0].body, /\[Radial artifact\]\(at:\/\/did:plc:agent\/com\.disnetdev\.radial\.artifact\/impl-1\)/) // And the record it wrote is reported back, so the ledger can remember it across a crash. assert.equal(server.observation.pull.record.cid, 'pullcid') } finally { await server.close() } }) }) it('submitAnswer: adoption compares whichever anchor the reply carries', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer-anchor-adopt') const client = await makeClient(pds) const artifact = { uri: 'at://did:plc:agent/com.disnetdev.radial.artifact/plan-1', cid: 'acid' } const first = await startServer(directory, client, { context: { mode: { ...ANSWER_MODE.mode, anchor: { artifact } } }, }) const firstResponse = await sendLine(first.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'one' }) await first.server.close() assert.equal(firstResponse.ok, true) // The same turn retried adopts its own record… const retry = await startServer(directory, client, { context: { mode: { ...ANSWER_MODE.mode, anchor: { artifact } } }, }) const retried = await sendLine(retry.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'two' }) await retry.server.close() assert.deepEqual(retried.ref, firstResponse.ref) // …while a turn that would anchor it somewhere else refuses to. const moved = await startServer(directory, client, { context: ANSWER_MODE }) const response = await sendLine(moved.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'three' }) await moved.server.close() assert.equal(response.ok, false) assert.match(response.error, /mismatched anchor/) assert.equal([...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.message).length, 1) }) }) it('mutual exclusion: submitAnswer only on an answer turn, and an answer turn submits nothing else', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer-mutex') const client = await makeClient(pds) // Artifact (plan) turn: submitAnswer rejected. const artifactTurn = await startServer(directory, client) const answerOnArtifact = await sendLine(artifactTurn.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'not my turn', }) await artifactTurn.server.close() assert.equal(answerOnArtifact.ok, false) assert.match(answerOnArtifact.error, /not an answer turn/) // Review turn: submitAnswer rejected too. const reviewTurn = await startServer(directory, client, { context: REVIEW_MODE }) const answerOnReview = await sendLine(reviewTurn.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'still not my turn', }) await reviewTurn.server.close() assert.equal(answerOnReview.ok, false) // Answer turn: neither an artifact nor a verdict may be submitted. const answerTurn = await startServer(directory, client, { context: ANSWER_MODE }) const artifactOnAnswer = await sendLine(answerTurn.socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'sneaky artifact', }) const reviewOnAnswer = await sendLine(answerTurn.socketPath, { method: 'submitReview', token: TOKEN, verdict: 'approve', findings: [], }) await answerTurn.server.close() assert.equal(artifactOnAnswer.ok, false) assert.match(artifactOnAnswer.error, /answer turn/) assert.equal(reviewOnAnswer.ok, false) assert.match(reviewOnAnswer.error, /not a review turn/) assert.equal(pds.records.size, 0) }) }) it('submitAnswer: a retried turn adopts its own record rather than double-posting to the thread', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer-adopt') const client = await makeClient(pds) const first = await startServer(directory, client, { context: ANSWER_MODE }) const firstResponse = await sendLine(first.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'first phrasing', }) await first.server.close() assert.equal(firstResponse.ok, true) // A retried answer turn legitimately writes different prose; adoption keys on the anchors only. const second = await startServer(directory, client, { context: ANSWER_MODE }) const secondResponse = await sendLine(second.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'second phrasing, entirely different', }) await second.server.close() assert.equal(secondResponse.ok, true) assert.deepEqual(secondResponse.ref, firstResponse.ref) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.message) assert.equal(stored.length, 1) assert.equal(stored[0].value.body, 'first phrasing') }) }) it('submitAnswer: adoption refuses a record at the answer- rkey with a mismatched anchor', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer-tamper') const client = await makeClient(pds) const first = await startServer(directory, client, { context: ANSWER_MODE }) await sendLine(first.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'original' }) await first.server.close() // A second turn for the same request whose subject is a DIFFERENT message: the rkey collides but // the record at it is not the one this turn would have written. const second = await startServer(directory, client, { context: { mode: { ...ANSWER_MODE.mode, subject: { uri: 'at://did:plc:human/com.disnetdev.radial.message/q2', cid: 'q2cid' }, }, }, }) const response = await sendLine(second.socketPath, { method: 'submitAnswer', token: TOKEN, body: 'about something else', }) await second.server.close() assert.equal(response.ok, false) assert.match(response.error, /mismatched anchor/) const stored = [...pds.records.values()].filter((r) => r.value.$type === COLLECTIONS.message) assert.equal(stored.length, 1) }) }) it('submitAnswer: rejects an oversized reply body and writes no record', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-answer-big') const client = await makeClient(pds) const { server, socketPath } = await startServer(directory, client, { context: ANSWER_MODE, limits: { maxFrameBytes: 2_000_000, frameTimeoutMs: 30_000, artifactBodyChars: 100, messageBodyChars: 10 }, }) try { const response = await sendLine(socketPath, { method: 'submitAnswer', token: TOKEN, body: 'x'.repeat(11), }) assert.equal(response.ok, false) assert.match(response.error, /exceeds 10 characters/) assert.equal(pds.records.size, 0) } finally { await server.close() } }) }) it('a broker failure is a submit error and writes NO artifact, leaving the turn retryable', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-tangled-fail') const client = await makeClient(pds) const branch = 'radial/impl-123456789abc' const { adapter } = brokeringForge({ openPullRequest: async () => { throw new Error('knot unreachable') }, }) const { server, socketPath } = await startServer(directory, client, { context: TANGLED_IMPL(adapter, branch) }) try { const failed = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'done', branch, commit: 'a'.repeat(40) }) assert.equal(failed.ok, false) assert.match(failed.error, /could not open the pull request.*knot unreachable/) assert.equal(pds.records.size, 0) assert.equal(server.observation.artifact, undefined) // The reservation is released, so the same server can still accept a real submit. const { adapter: working } = brokeringForge() server.observation.pull = undefined const retried = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'done', branch, commit: 'a'.repeat(40) }) assert.equal(retried.ok, false, 'the adapter on this server still throws') assert.ok(working.openPullRequest) } finally { await server.close() } }) }) it('a GitHub submit with no --pr is still rejected: only a brokering forge relaxes that', async () => { await withTempDir(async (directory) => { const pds = new LocalPds('did:plc:agent-nopr') const client = await makeClient(pds) const branch = 'radial/impl-123456789abc' const { server, socketPath } = await startServer(directory, client, { context: { mode: { kind: 'artifact', type: 'implementation' }, implementation: { gitUrl: 'https://github.com/acme/widget.git', branch }, }, }) try { const response = await sendLine(socketPath, { method: 'submitArtifact', token: TOKEN, title: 'A short title', body: 'done', branch, commit: 'a'.repeat(40) }) assert.equal(response.ok, false) assert.match(response.error, /requires --branch, --commit, and --pr/) assert.equal(pds.records.size, 0) } finally { await server.close() } }) })