From 25614f45067623edc2537af59b05267a38d426f0 Mon Sep 17 00:00:00 2001 From: Maggie Appleton <5599295+MaggieAppleton@users.noreply.github.com> Date: Sat, 3 Oct 2026 07:09:30 +0100 Subject: [PATCH 1/2] Add pure Jev event replay and validation --- .../correction-validation.ts | 90 +++++ .../conversation-plan/domain-analysis.test.ts | 82 +++++ .../src/conversation-plan/domain-analysis.ts | 97 +++++ .../domain-corrections.test.ts | 182 ++++++++++ .../conversation-plan/domain-corrections.ts | 135 +++++++ .../src/conversation-plan/domain-initial.ts | 15 + .../src/conversation-plan/domain-research.ts | 88 +++++ .../conversation-plan/domain.test-fixtures.ts | 60 +++ .../src/conversation-plan/domain.test.ts | 213 +++++++++++ apps/server/src/conversation-plan/domain.ts | 6 + .../conversation-plan/event-contributions.ts | 120 ++++++ .../conversation-plan/event-corrections.ts | 111 ++++++ .../src/conversation-plan/event-lifecycle.ts | 129 +++++++ .../src/conversation-plan/event-settlement.ts | 160 ++++++++ .../conversation-plan/event-support.test.ts | 159 ++++++++ .../src/conversation-plan/event-support.ts | 69 ++++ .../src/conversation-plan/event-validation.ts | 304 ++++++++++++++++ .../src/conversation-plan/events.test.ts | 280 ++++++++++++++ apps/server/src/conversation-plan/events.ts | 110 ++++++ .../src/conversation-plan/preference.ts | 94 +++++ .../conversation-plan/research-snapshots.ts | 91 +++++ .../research-validation.test.ts | 153 ++++++++ .../conversation-plan/research-validation.ts | 174 +++++++++ .../src/conversation-plan/state-provenance.ts | 53 +++ .../conversation-plan/state-restoration.ts | 72 ++++ .../state-validation.test.ts | 237 ++++++++++++ .../src/conversation-plan/state-validation.ts | 265 ++++++++++++++ .../conversation-plan/validation-fields.ts | 54 +++ .../src/conversation-plan/validation.test.ts | 342 ++++++++++++++++++ .../src/conversation-plan/validation.ts | 17 + 30 files changed, 3962 insertions(+) create mode 100644 apps/server/src/conversation-plan/correction-validation.ts create mode 100644 apps/server/src/conversation-plan/domain-analysis.test.ts create mode 100644 apps/server/src/conversation-plan/domain-analysis.ts create mode 100644 apps/server/src/conversation-plan/domain-corrections.test.ts create mode 100644 apps/server/src/conversation-plan/domain-corrections.ts create mode 100644 apps/server/src/conversation-plan/domain-initial.ts create mode 100644 apps/server/src/conversation-plan/domain-research.ts create mode 100644 apps/server/src/conversation-plan/domain.test-fixtures.ts create mode 100644 apps/server/src/conversation-plan/domain.test.ts create mode 100644 apps/server/src/conversation-plan/domain.ts create mode 100644 apps/server/src/conversation-plan/event-contributions.ts create mode 100644 apps/server/src/conversation-plan/event-corrections.ts create mode 100644 apps/server/src/conversation-plan/event-lifecycle.ts create mode 100644 apps/server/src/conversation-plan/event-settlement.ts create mode 100644 apps/server/src/conversation-plan/event-support.test.ts create mode 100644 apps/server/src/conversation-plan/event-support.ts create mode 100644 apps/server/src/conversation-plan/event-validation.ts create mode 100644 apps/server/src/conversation-plan/events.test.ts create mode 100644 apps/server/src/conversation-plan/events.ts create mode 100644 apps/server/src/conversation-plan/preference.ts create mode 100644 apps/server/src/conversation-plan/research-snapshots.ts create mode 100644 apps/server/src/conversation-plan/research-validation.test.ts create mode 100644 apps/server/src/conversation-plan/research-validation.ts create mode 100644 apps/server/src/conversation-plan/state-provenance.ts create mode 100644 apps/server/src/conversation-plan/state-restoration.ts create mode 100644 apps/server/src/conversation-plan/state-validation.test.ts create mode 100644 apps/server/src/conversation-plan/state-validation.ts create mode 100644 apps/server/src/conversation-plan/validation-fields.ts create mode 100644 apps/server/src/conversation-plan/validation.test.ts create mode 100644 apps/server/src/conversation-plan/validation.ts diff --git a/apps/server/src/conversation-plan/correction-validation.ts b/apps/server/src/conversation-plan/correction-validation.ts new file mode 100644 index 00000000..12eb5d5f --- /dev/null +++ b/apps/server/src/conversation-plan/correction-validation.ts @@ -0,0 +1,90 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { id, knownKeys, record, text, version } from "./validation-fields"; + +export function assertCorrectionChange( + value: unknown, +): asserts value is ConversationPlan.CorrectionChange { + let change = record(value); + switch (change.kind) { + case "add-excerpt": + knownKeys(change, [ + "kind", + "messageId", + "start", + "end", + "contributionKind", + "targetOptionId", + ]); + id(change.messageId); + if ( + !Number.isSafeInteger(change.start) || !Number.isSafeInteger(change.end) + || (change.start as number) < 0 || (change.end as number) <= (change.start as number) + || (change.end as number) - (change.start as number) > 500 + ) throw new Error("invalid excerpt range"); + if (!["option", "reason", "constraint"].includes(change.contributionKind as string)) { + throw new Error("invalid excerpt contribution kind"); + } + if (change.targetOptionId !== undefined) id(change.targetOptionId); + break; + case "edit": + knownKeys(change, ["kind", "field", "contributionId", "text"]); + if (!["question", "contribution", "decision"].includes(change.field as string)) { + throw new Error("invalid edit field"); + } + text(change.text); + if (change.field === "contribution") id(change.contributionId); + else if ("contributionId" in change) { + throw new Error("unexpected contribution ID on question or decision edit"); + } + break; + case "move": + knownKeys(change, ["kind", "contributionId", "targetThreadId", "targetVersion"]); + id(change.contributionId); + id(change.targetThreadId); + version(change.targetVersion); + break; + case "set-status": + knownKeys(change, ["kind", "status"]); + if (!["exploring", "leaning", "reopened"].includes(change.status as string)) { + throw new Error("invalid status correction"); + } + break; + case "retarget-stance": + knownKeys(change, ["kind", "stanceId", "optionId"]); + id(change.stanceId); + if (change.optionId !== undefined) id(change.optionId); + break; + case "dismiss-stance": + knownKeys(change, ["kind", "stanceId"]); + id(change.stanceId); + break; + case "retarget-contribution": + knownKeys(change, ["kind", "contributionId", "targetId"]); + id(change.contributionId); + id(change.targetId); + break; + case "record-decision": + knownKeys(change, ["kind", "text", "optionId"]); + text(change.text); + if (change.optionId !== undefined) id(change.optionId); + break; + case "confirm-candidate": + case "reject-candidate": + knownKeys(change, ["kind", "candidateId"]); + id(change.candidateId); + break; + default: + throw new Error("invalid correction action"); + } +} + +export function assertCorrectionAction( + value: unknown, +): asserts value is ConversationPlan.CorrectionAction { + let action = record(value); + knownKeys(action, ["actionId", "threadId", "expectedVersion", "change"]); + id(action.actionId); + id(action.threadId); + version(action.expectedVersion); + assertCorrectionChange(action.change); +} diff --git a/apps/server/src/conversation-plan/domain-analysis.test.ts b/apps/server/src/conversation-plan/domain-analysis.test.ts new file mode 100644 index 00000000..542f68b6 --- /dev/null +++ b/apps/server/src/conversation-plan/domain-analysis.test.ts @@ -0,0 +1,82 @@ +import { describe, expect, test } from "bun:test"; +import type { ConversationPlan } from "@chopin/protocol"; +import { applyInference, completeAnalysis, enqueue, initialState, restoreState } from "./domain"; +import { alice, bob, opened } from "./domain.test-fixtures"; + +describe("conversation plan events", () => { + test("message completion applies events and bounded debug with queue removal", () => { + let queued = enqueue(initialState(), "m1"); + let finished = completeAnalysis(queued, "m1", [opened()], alice, { + questionSetVersion: "v1", + modelVersion: "jev-test", + status: "applied", + passes: [{ stage: "triage", answers: { question: { type: "noul", noul: 0.9 } } }], + policyGate: "clear question", + latencyMs: 12, + }); + expect(finished.queue).toEqual([]); + expect(finished.analysis[0].eventIds).toEqual(["m1:0:thread.opened:v1"]); + expect(finished.threads).toHaveLength(1); + let bad = { ...opened(), id: "bad", observedThreadVersion: 3 }; + expect(() => + completeAnalysis(queued, "m1", [bad], alice, { + questionSetVersion: "v1", + modelVersion: "jev-test", + status: "applied", + passes: [], + }) + ).toThrow(); + expect(queued.threads).toHaveLength(0); + }); + + test("completed debug data is detached from the caller's mutable response", () => { + let debug: Omit = { + questionSetVersion: "v1", + modelVersion: "jev-test", + status: "unlinked", + passes: [], + }; + let finished = completeAnalysis(enqueue(initialState(), "m1"), "m1", [], alice, debug); + debug.passes.push({ stage: "triage", answers: {} }); + expect(finished.analysis[0].passes).toEqual([]); + }); + + test("absent state is initial, malformed present snapshots fail closed", () => { + expect(restoreState(undefined)).toEqual(initialState()); + expect(() => + restoreState({ + schemaVersion: 1, + revision: 0, + events: [], + threads: [{}], + queue: [], + analysis: [], + }) + ).toThrow(); + let state = applyInference(initialState(), opened(), alice); + expect(() => restoreState(state, [bob])).toThrow(/missing/i); + }); + + test("queue deduplicates message IDs across retry and restore", () => { + let state = enqueue(initialState(), "m1"); + expect(enqueue(state, "m1")).toEqual(state); + expect(restoreState(structuredClone(state)).queue).toEqual([{ + messageId: "m1", + status: "pending", + attempts: 0, + }]); + }); + + test("restore rejects unbounded probability labels in debug data", () => { + let state = initialState(); + state.analysis.push({ + messageId: "m1", + questionSetVersion: "v1", + modelVersion: "jev-test", + status: "unlinked", + passes: [{ stage: "triage", answers: { ["x".repeat(1000)]: { type: "noul", noul: 1 } } }], + eventIds: [], + }); + expect(() => restoreState(state)).toThrow(); + }); +}); diff --git a/apps/server/src/conversation-plan/domain-analysis.ts b/apps/server/src/conversation-plan/domain-analysis.ts new file mode 100644 index 00000000..a1178480 --- /dev/null +++ b/apps/server/src/conversation-plan/domain-analysis.ts @@ -0,0 +1,97 @@ +import type { Chat, ConversationPlan } from "@chopin/protocol"; +import { initialState } from "./domain-initial"; +import { applyEvent } from "./events"; +import { validateSource } from "./sources"; +import { assertStateShape, MAX_ANALYSIS, MAX_QUEUE } from "./validation"; + +type State = ConversationPlan.State; +type Event = ConversationPlan.Event; + +export function enqueue(state: State, messageId: string): State { + if (!messageId || messageId.length > 200) throw new Error("invalid message ID"); + if (state.queue.some((item) => item.messageId === messageId)) return state; + if (state.queue.length >= MAX_QUEUE) throw new Error("conversation analysis queue is full"); + return { + ...state, + revision: state.revision + 1, + queue: [...state.queue, { messageId, status: "pending", attempts: 0 }], + }; +} + +export function retryMessage(state: State, messageId: string): State { + let index = state.queue.findIndex((item) => item.messageId === messageId); + if (index < 0 || state.queue[index].status !== "failed") throw new Error("message is not failed"); + let queue = state.queue.map((item, i) => + i === index + ? { messageId, status: "pending" as const, attempts: item.attempts + 1 } + : item + ); + return { ...state, revision: state.revision + 1, queue }; +} + +export function replay(events: readonly Event[]): State { + let state = initialState(); + for (let event of events) state = applyEvent(state, event); + return state; +} + +export function applyInference(state: State, event: Event, message: Chat.Entry): State { + if (state.events.some((accepted) => accepted.id === event.id)) return state; + if (event.origin === "human") throw new Error("human events require authenticated correction"); + if ( + event.type === "decision.recorded" || event.type === "decision.reopened" + || (event.type === "candidate.proposed" && event.candidate.kind === "resolution") + ) { + throw new Error("decisions are recorded on the card"); + } + if (!("source" in event) || !event.source) throw new Error("inference needs a source"); + validateSource(event.source, message); + return applyEvent(state, event); +} + +export function completeAnalysis( + state: State, + messageId: string, + events: readonly Event[], + message: Chat.Entry, + analysis: Omit, +): State { + let queued = state.queue.find((item) => item.messageId === messageId); + if (message.id !== messageId || !queued || queued.status === "failed") { + throw new Error("conversation analysis message is not queued"); + } + if (events.length > 12) throw new Error("too many events from one message"); + if (!["applied", "unlinked", "failed"].includes(analysis.status)) { + throw new Error("analysis is not terminal"); + } + if (analysis.status !== "applied" && events.length > 0) { + throw new Error("non-applied analysis has events"); + } + let next = state; + for (let event of events) { + if (!("source" in event) || event.source?.messageId !== messageId) { + throw new Error("analysis event belongs to another message"); + } + next = applyInference(next, event, message); + } + let eventIds = events.filter((event) => next.events.some((accepted) => accepted.id === event.id)) + .map((event) => event.id); + let record: ConversationPlan.AnalysisRecord = { + ...structuredClone(analysis), + messageId, + eventIds, + }; + let history = [...next.analysis.filter((item) => item.messageId !== messageId), record].slice( + -MAX_ANALYSIS, + ); + let queue = analysis.status === "failed" + ? next.queue.map((item) => + item.messageId === messageId + ? { ...item, status: "failed" as const, error: analysis.error?.slice(0, 300) } + : item + ) + : next.queue.filter((item) => item.messageId !== messageId); + let completed = { ...next, revision: next.revision + 1, queue, analysis: history }; + assertStateShape(completed); + return completed; +} diff --git a/apps/server/src/conversation-plan/domain-corrections.test.ts b/apps/server/src/conversation-plan/domain-corrections.test.ts new file mode 100644 index 00000000..afddaac7 --- /dev/null +++ b/apps/server/src/conversation-plan/domain-corrections.test.ts @@ -0,0 +1,182 @@ +import { describe, expect, test } from "bun:test"; +import type { ConversationPlan } from "@chopin/protocol"; +import { applyCorrection, applyInference, initialState, replay, restoreState } from "./domain"; +import { applyEvent } from "./events"; +import { alice, bob, opened, option, source } from "./domain.test-fixtures"; + +describe("conversation plan events", () => { + test("restores JSON-persisted optional fields after decide, reopen and status correction", () => { + let state = applyInference(initialState(), opened(), alice); + state = applyEvent(state, { + id: "m7:0:decision.recorded:v1", + type: "decision.recorded", + threadId: "t1", + observedThreadVersion: 1, + origin: "human", + actor: { kind: "member", handle: "alice" }, + at: 7, + text: "Use an outline", + explicit: true, + }); + state = applyEvent(state, { + id: "m8:0:decision.reopened:v1", + type: "decision.reopened", + threadId: "t1", + observedThreadVersion: 2, + origin: "human", + actor: { kind: "member", handle: "bob" }, + at: 8, + explicit: true, + }); + state = applyCorrection( + state, + { + actionId: "status-exploring", + threadId: "t1", + expectedVersion: 3, + change: { kind: "set-status", status: "exploring" }, + }, + { kind: "member", handle: "alice" }, + 9, + ); + let saved = JSON.parse(JSON.stringify(state)); + let restored = restoreState(saved, [alice]); + expect(restored).toEqual(saved); + }); + + test("human wording survives later classifier work and a wrong target can move", () => { + let state = applyInference(initialState(), opened(), alice); + state = applyInference(state, option(state), bob); + let secondMessage = { ...alice, id: "m5", text: "What about blank canvas?" }; + let second: ConversationPlan.Event = { + ...opened(), + id: "m5:0:thread.opened:v1", + threadId: "t2", + at: 5, + source: source(secondMessage, "question"), + question: "What about blank canvas?", + }; + state = applyInference(state, second, secondMessage); + state = applyCorrection( + state, + { + actionId: "edit-o1", + threadId: "t1", + expectedVersion: 2, + change: { + kind: "edit", + field: "contribution", + contributionId: "o1", + text: "Optional starter outline", + }, + }, + { kind: "member", handle: "alice" }, + 6, + ); + expect(state.threads[0].contributions[0].authoring).toBe("human-edited"); + let stale = option(state); + stale.id = "m2:retry"; + stale.observedThreadVersion = 2; + expect(() => applyInference(state, stale, bob)).toThrow(/stale/i); + state = applyCorrection( + state, + { + actionId: "move-o1", + threadId: "t1", + expectedVersion: state.threads[0].version, + change: { kind: "move", contributionId: "o1", targetThreadId: "t2", targetVersion: 1 }, + }, + { kind: "member", handle: "alice" }, + 7, + ); + expect(state.threads[0].contributions).toHaveLength(0); + expect(state.threads[1].contributions[0].text).toBe("Optional starter outline"); + expect(state.threads[1].contributions[0].sources[0].messageId).toBe("m2"); + expect(replay(state.events).threads).toEqual(state.threads); + }); + + test("legacy resolution candidates replay, while new confirmations are refused", () => { + let state = applyInference(initialState(), opened(), alice); + let candidateMessage = { ...bob, id: "m6", text: "Maybe we agreed on an outline?" }; + let candidate: ConversationPlan.Event = { + id: "m6:0:candidate.proposed:v1", + type: "candidate.proposed", + threadId: "t1", + observedThreadVersion: 1, + origin: "classifier", + actor: { kind: "classifier" }, + at: 6, + source: source(candidateMessage, "resolution"), + candidate: { id: "c1", kind: "resolution", text: "Use an outline" }, + }; + state = applyEvent(state, candidate); + state = applyCorrection( + state, + { + actionId: "reject-c1", + threadId: "t1", + expectedVersion: 2, + change: { kind: "reject-candidate", candidateId: "c1" }, + }, + { kind: "member", handle: "alice" }, + 7, + ); + expect(state.threads[0].candidates[0].status).toBe("rejected"); + expect(state.threads[0].status).toBe("exploring"); + let second: ConversationPlan.Event = { + ...candidate, + id: "m6:1:candidate.proposed:v1", + observedThreadVersion: state.threads[0].version, + candidate: { id: "c2", kind: "resolution", text: "Use an outline" }, + }; + state = applyEvent(state, second); + expect(() => + applyCorrection( + state, + { + actionId: "confirm-c2", + threadId: "t1", + expectedVersion: state.threads[0].version, + change: { kind: "confirm-candidate", candidateId: "c2" }, + }, + { kind: "member", handle: "alice" }, + 8, + ) + ) + .toThrow("decisions are recorded on the card"); + state = applyEvent(state, { + id: "historic-confirm-c2", + type: "candidate.confirmed", + threadId: "t1", + observedThreadVersion: state.threads[0].version, + origin: "human", + actor: { kind: "member", handle: "alice" }, + at: 8, + candidateId: "c2", + }); + expect(state.threads[0].status).toBe("decided"); + expect(state.threads[0].decision?.actor).toEqual({ kind: "member", handle: "alice" }); + expect(state.threads[0].decision?.sources[0].messageId).toBe("m6"); + expect(replay(state.events).threads).toEqual(state.threads); + expect(restoreState(JSON.parse(JSON.stringify(state))).threads).toEqual(state.threads); + }); + + test("stale corrections fail and stable action IDs make retries idempotent", () => { + let state = applyInference(initialState(), opened(), alice); + let action: ConversationPlan.CorrectionAction = { + actionId: "edit-question", + threadId: "t1", + expectedVersion: 1, + change: { kind: "edit", field: "question", text: "Should we offer an outline?" }, + }; + let changed = applyCorrection(state, action, { kind: "member", handle: "bob" }, 3); + expect(changed.threads[0].question).toBe("Should we offer an outline?"); + expect(applyCorrection(changed, action, { kind: "member", handle: "bob" }, 4)).toEqual(changed); + expect(() => + applyCorrection(changed, { ...action, actionId: "another-action" }, { + kind: "member", + handle: "bob", + }, 4) + ).toThrow(/stale/i); + }); +}); diff --git a/apps/server/src/conversation-plan/domain-corrections.ts b/apps/server/src/conversation-plan/domain-corrections.ts new file mode 100644 index 00000000..904a495f --- /dev/null +++ b/apps/server/src/conversation-plan/domain-corrections.ts @@ -0,0 +1,135 @@ +import type { Chat, ConversationPlan } from "@chopin/protocol"; +import { isDeepStrictEqual } from "node:util"; +import { ulid } from "@chopin/dialect"; +import { applyEvent } from "./events"; +import { validateSource } from "./sources"; +import { assertCorrectionAction } from "./validation"; + +type State = ConversationPlan.State; +type Event = ConversationPlan.Event; +type Member = Extract; + +export function applyCorrection( + state: State, + request: ConversationPlan.CorrectionAction, + actor: Member, + at: number, + messages?: ReadonlyMap | readonly Chat.Entry[], +): State { + assertCorrectionAction(request); + if (actor.kind !== "member" || !actor.handle) { + throw new Error("correction requires a human member"); + } + let id = `human:${actor.handle}:${request.actionId}`; + let base: ConversationPlan.EventBase = { + id, + threadId: request.threadId, + observedThreadVersion: request.expectedVersion, + origin: "human", + actor, + at, + }; + let change = request.change; + let existing = state.events.find((accepted) => accepted.id === id); + let event: Event; + switch (change.kind) { + case "add-excerpt": { + let saved = Array.isArray(messages) + ? messages.find(item => item.id === change.messageId) + : (messages as ReadonlyMap | undefined)?.get(change.messageId); + if (!saved) throw new Error("saved excerpt message is missing"); + let quote = saved.text.slice(change.start, change.end); + let source: ConversationPlan.SourceRef = { + messageId: change.messageId, + author: saved.author as ConversationPlan.SourceAuthor, + quote, + start: change.start, + end: change.end, + role: change.contributionKind, + }; + validateSource(source, saved); + if (!quote.trim() || quote.length > 500) throw new Error("invalid excerpt quote"); + let outcome = state.analysis.find(item => item.messageId === change.messageId) + ?.outcomes?.find(item => + ["review", "ignored"].includes(item.status) + && item.start <= change.start && item.end >= change.end + ); + if (!outcome && !existing) throw new Error("excerpt has no held analysis outcome"); + let target = state.threads.find(item => item.id === request.threadId); + if (!target || ["decided", "discarded"].includes(target.status) && !existing) { + throw new Error("excerpt target thread is not open"); + } + if (target.version !== request.expectedVersion && !existing) { + throw new Error("stale conversation thread version"); + } + if ( + !existing && change.targetOptionId && ( + change.contributionKind === "option" + || !target.contributions.some(item => + item.kind === "option" && item.id === change.targetOptionId + ) + ) + ) throw new Error("invalid excerpt target option"); + if ( + !existing + && state.events.some(item => + item.threadId === request.threadId && "source" in item && item.source + && item.source.messageId === change.messageId + && item.source.start === change.start && item.source.end === change.end + ) + ) throw new Error("duplicate excerpt source"); + let contributionId = existing && "contribution" in existing + ? existing.contribution.id + : ulid(); + event = { + ...base, + type: `${change.contributionKind}.added`, + source, + contribution: { + id: contributionId, + text: quote, + authoring: "quoted", + targetId: change.targetOptionId ?? request.threadId, + }, + }; + break; + } + case "record-decision": + event = { + ...base, + type: "decision.recorded", + text: change.text, + optionId: change.optionId, + explicit: true, + }; + break; + case "confirm-candidate": + event = { ...base, type: "candidate.confirmed", candidateId: change.candidateId }; + break; + case "reject-candidate": + event = { ...base, type: "candidate.rejected", candidateId: change.candidateId }; + break; + default: + event = { ...base, type: "card.corrected", change }; + } + if (existing) { + let saved = JSON.parse(JSON.stringify({ ...existing, at: 0 })); + let submitted = JSON.parse(JSON.stringify({ ...event, at: 0 })); + if (!isDeepStrictEqual(saved, submitted)) throw new Error("human action ID collision"); + return state; + } + let thread = state.threads.find((item) => item.id === request.threadId); + if ( + change.kind === "record-decision" + || (change.kind === "confirm-candidate" && ( + thread?.questionnaireId + || thread?.candidates.some((item) => + item.id === change.candidateId && item.kind === "resolution" + ) + )) + || (thread?.questionnaireId && ( + change.kind === "set-status" || (change.kind === "edit" && change.field === "decision") + )) + ) throw new Error("decisions are recorded on the card"); + return applyEvent(state, event); +} diff --git a/apps/server/src/conversation-plan/domain-initial.ts b/apps/server/src/conversation-plan/domain-initial.ts new file mode 100644 index 00000000..649bec33 --- /dev/null +++ b/apps/server/src/conversation-plan/domain-initial.ts @@ -0,0 +1,15 @@ +import type { ConversationPlan } from "@chopin/protocol"; + +type State = ConversationPlan.State; + +export function initialState(): State { + return { + schemaVersion: 1, + revision: 0, + events: [], + threads: [], + queue: [], + analysis: [], + researchOffers: [], + }; +} diff --git a/apps/server/src/conversation-plan/domain-research.ts b/apps/server/src/conversation-plan/domain-research.ts new file mode 100644 index 00000000..9dbdc9db --- /dev/null +++ b/apps/server/src/conversation-plan/domain-research.ts @@ -0,0 +1,88 @@ +import type { Chat, ConversationPlan } from "@chopin/protocol"; +import { isDeepStrictEqual } from "node:util"; +import { validateSource } from "./sources"; +import { assertResearchOfferShape, assertStateShape, MAX_RESEARCH_OFFERS } from "./validation"; +import { namesStaleResearchOption, taskMatches } from "./research-snapshots"; + +type State = ConversationPlan.State; + +export function offerResearch( + state: State, + proposal: Omit, + message: Chat.Entry, +): State { + if ( + !proposal || typeof proposal !== "object" || Array.isArray(proposal) + || Object.keys(proposal).some((key) => + !["id", "needId", "contextId", "source", "brief", "threadId", "task"].includes(key) + ) + ) throw new Error("invalid research offer proposal"); + let offer: ConversationPlan.ResearchOffer = { + ...structuredClone(proposal), + status: "offered", + }; + assertResearchOfferShape(offer); + validateSource({ ...offer.source, role: "support" }, message); + let current = state.researchOffers ?? []; + let existing = current.find((item) => item.id === offer.id); + if (existing) { + let { status: _status, action: _action, ...identity } = existing; + if (isDeepStrictEqual(identity, proposal)) return state; + throw new Error("research offer ID already has different content"); + } + if (offer.threadId && !state.threads.some((thread) => thread.id === offer.threadId)) { + throw new Error("research offer thread is missing"); + } + if (offer.task && !taskMatches(state, offer.task)) { + throw new Error("research task options do not match the current thread"); + } + if (offer.task && namesStaleResearchOption(state, offer.task, offer.source.quote)) { + throw new Error("research task source names a stale option"); + } + if (offer.task && offer.task.observedEventCount !== state.events.length) { + throw new Error("research task event prefix is stale"); + } + if (current.length >= MAX_RESEARCH_OFFERS) throw new Error("research offers are full"); + let next = { + ...state, + revision: state.revision + 1, + researchOffers: [...current, offer], + }; + assertStateShape(next); + return next; +} + +export function actOnResearchOffer( + state: State, + offerId: string, + action: ConversationPlan.ResearchAction, +): State { + let current = state.researchOffers ?? []; + let index = current.findIndex((item) => item.id === offerId); + if (index < 0) throw new Error("research offer is missing"); + let offer = current[index]; + let status: ConversationPlan.ResearchOffer["status"] = action.kind === "research" + ? "accepted" + : "dismissed"; + assertResearchOfferShape({ ...offer, status, action }); + if (offer.status !== "offered") { + if ( + offer.action?.id === action.id && offer.action.kind === action.kind + && isDeepStrictEqual(offer.action.actor, action.actor) + && offer.action.principalId === action.principalId + ) return state; + throw new Error("research offer already has a terminal action"); + } + if (current.some((item) => item.action?.id === action.id)) { + throw new Error("research action ID already used"); + } + let next: State = { + ...state, + revision: state.revision + 1, + researchOffers: current.map((item, position) => + position === index ? { ...item, status, action: structuredClone(action) } : item + ), + }; + assertStateShape(next); + return next; +} diff --git a/apps/server/src/conversation-plan/domain.test-fixtures.ts b/apps/server/src/conversation-plan/domain.test-fixtures.ts new file mode 100644 index 00000000..9a4492f7 --- /dev/null +++ b/apps/server/src/conversation-plan/domain.test-fixtures.ts @@ -0,0 +1,60 @@ +import type { Chat, ConversationPlan } from "@chopin/protocol"; + +export let alice: Chat.Entry = { + id: "m1", + author: { kind: "member", handle: "alice" }, + text: "Should we start with an outline?", + ts: 1, +}; +export let bob: Chat.Entry = { + id: "m2", + author: { kind: "member", handle: "bob" }, + text: "I prefer an optional outline.", + ts: 2, +}; + +export function source( + message: Chat.Entry, + role: ConversationPlan.SourceRole, +): ConversationPlan.SourceRef { + return { + messageId: message.id, + author: message.author as ConversationPlan.SourceAuthor, + quote: message.text, + start: 0, + end: message.text.length, + role, + }; +} + +export function opened(): Extract { + return { + id: "m1:0:thread.opened:v1", + type: "thread.opened", + threadId: "t1", + observedThreadVersion: 0, + origin: "classifier", + actor: { kind: "classifier" }, + at: 1, + source: source(alice, "question"), + question: "Should we start with an outline?", + }; +} + +export function option(state: ConversationPlan.State): ConversationPlan.Event { + return { + id: "m2:0:option.added:v1", + type: "option.added", + threadId: "t1", + observedThreadVersion: state.threads[0].version, + origin: "classifier", + actor: { kind: "classifier" }, + at: 2, + source: source(bob, "option"), + contribution: { + id: "o1", + text: "Optional outline", + authoring: "scribe", + }, + }; +} diff --git a/apps/server/src/conversation-plan/domain.test.ts b/apps/server/src/conversation-plan/domain.test.ts new file mode 100644 index 00000000..69c7ab7e --- /dev/null +++ b/apps/server/src/conversation-plan/domain.test.ts @@ -0,0 +1,213 @@ +import { describe, expect, test } from "bun:test"; +import type { Chat, ConversationPlan } from "@chopin/protocol"; +import { applyInference, initialState, replay, restoreState } from "./domain"; +import { applyEvent } from "./events"; +import { alice, bob, opened, option, source } from "./domain.test-fixtures"; + +describe("conversation plan events", () => { + test("relabeling preserves quoted evidence and replays with the new display label", () => { + let state = applyInference(initialState(), opened(), alice); + state = applyInference(state, option(state), bob); + state = applyEvent(state, { + id: "linked", + type: "card.linked", + threadId: "t1", + observedThreadVersion: state.threads[0]!.version, + origin: "classifier", + actor: { kind: "classifier" }, + at: 3, + questionnaireId: "card-1", + }); + let before = structuredClone(state.threads[0]!.contributions[0]!); + let event: ConversationPlan.Event = { + id: "relabel-1", + type: "option.relabeled", + threadId: "t1", + observedThreadVersion: state.threads[0]!.version, + origin: "planner", + actor: { kind: "agent" }, + at: 4, + optionId: "o1", + label: "Outline when useful", + observedCardRevision: 0, + }; + let changed = applyEvent(state, event); + expect(changed.threads[0]!.contributions[0]).toMatchObject({ + ...before, + displayLabel: "Outline when useful", + }); + expect(state.threads[0]!.contributions[0]).toEqual(before); + expect(restoreState(changed, [alice, bob])).toEqual(changed); + expect(replay(changed.events).threads).toEqual(changed.threads); + expect(applyEvent(changed, event)).toBe(changed); + expect(() => applyEvent(state, { ...event, label: "A or B" })).toThrow(); + }); + test("a Planner-scribed question opens without a fabricated chat source and replays", () => { + let event: ConversationPlan.Event = { + id: "planner-question", + type: "thread.opened", + threadId: "planner-thread", + observedThreadVersion: 0, + origin: "planner", + actor: { kind: "agent" }, + at: 1, + question: "Which approach should we take?", + }; + let state = applyEvent(initialState(), event); + expect(state.threads[0]).toMatchObject({ + questionSources: [], + questionAuthoring: "scribe", + }); + expect(restoreState(state, [])).toEqual(state); + expect(replay(state.events).threads).toEqual(state.threads); + expect(() => + applyEvent(initialState(), { + ...event, + origin: "classifier", + actor: { kind: "classifier" }, + }) + ).toThrow("opening requires a question source"); + }); + + test("accepted events replay to the same threads and revision", () => { + let state = applyInference(initialState(), opened(), alice); + state = applyInference(state, option(state), bob); + let restored = restoreState(structuredClone(state)); + let rebuilt = replay(state.events); + expect(restored).toEqual(state); + expect(rebuilt.threads).toEqual(state.threads); + expect(rebuilt.revision).toBe(state.revision); + }); + + test("rejects forged quote offsets and source authors", () => { + let forged = opened(); + forged.source!.quote = "a different question"; + expect(() => applyInference(initialState(), forged, alice)).toThrow(); + forged = opened(); + forged.source!.author = { kind: "member", handle: "mallory" }; + expect(() => applyInference(initialState(), forged, alice)).toThrow(); + }); + + test("deduplicates accepted IDs before checking a stale version", () => { + let first = opened(); + let state = applyInference(initialState(), first, alice); + expect(applyInference(state, first, alice)).toEqual(state); + let stale = option(state); + stale.observedThreadVersion = 0; + expect(() => applyInference(state, stale, bob)).toThrow(/stale/i); + }); + + test("accepted event data cannot be changed by mutating the candidate later", () => { + let candidate = opened(); + let state = applyInference(initialState(), candidate, alice); + candidate.source!.quote = "forged later"; + let accepted = state.events[0]; + expect("source" in accepted && accepted.source?.quote).toBe("Should we start with an outline?"); + expect(state.threads[0].questionSources[0].quote).toBe("Should we start with an outline?"); + }); + + test("shows the latest participant stance and retains the reversal", () => { + let state = applyInference(initialState(), opened(), alice); + state = applyInference(state, option(state), bob); + let first: ConversationPlan.Event = { + id: "m3:0:stance.changed:v1", + type: "stance.changed", + threadId: "t1", + observedThreadVersion: state.threads[0].version, + origin: "classifier", + actor: { kind: "classifier" }, + at: 3, + source: source({ ...alice, id: "m3", text: "I support this" }, "support"), + optionId: "o1", + position: "support", + }; + let supportMessage = { ...alice, id: "m3", text: "I support this" }; + state = applyInference(state, first, supportMessage); + let reversalMessage = { ...alice, id: "m4", text: "Actually I oppose this" }; + let reversal: ConversationPlan.Event = { + ...first, + id: "m4:0:stance.changed:v1", + observedThreadVersion: state.threads[0].version, + at: 4, + source: source(reversalMessage, "objection"), + position: "oppose", + }; + state = applyInference(state, reversal, reversalMessage); + expect(state.threads[0].stances).toMatchObject([ + { participant: "alice", optionId: "o1", position: "oppose" }, + ]); + expect(state.threads[0].stanceHistory.map((stance) => stance.position)).toEqual([ + "support", + "oppose", + ]); + }); + + test("preserves a decision after an explicit human reopening", () => { + let state = applyInference(initialState(), opened(), alice); + let decision: ConversationPlan.Event = { + id: "m3:0:decision.recorded:v1", + type: "decision.recorded", + threadId: "t1", + observedThreadVersion: 1, + origin: "human", + actor: { kind: "member", handle: "alice" }, + at: 3, + text: "Use an outline", + explicit: true, + }; + state = applyEvent(state, decision); + expect(state.threads[0].status).toBe("decided"); + let reopen: ConversationPlan.Event = { + id: "m4:0:decision.reopened:v1", + type: "decision.reopened", + threadId: "t1", + observedThreadVersion: state.threads[0].version, + origin: "human", + actor: { kind: "member", handle: "bob" }, + at: 4, + explicit: true, + }; + state = applyEvent(state, reopen); + expect(state.threads[0].status).toBe("reopened"); + expect(state.threads[0].decision).toBeUndefined(); + expect(state.threads[0].decisionHistory).toHaveLength(1); + expect(replay(state.events).threads).toEqual(state.threads); + }); + + test("an agent cannot decide or cast a participant stance", () => { + let state = applyInference(initialState(), opened(), alice); + let agent: Chat.Entry = { + id: "agent1", + author: { kind: "agent" }, + text: "We decided to use an outline", + ts: 4, + }; + let decision: ConversationPlan.Event = { + id: "agent1:0:decision.recorded:v1", + type: "decision.recorded", + threadId: "t1", + observedThreadVersion: 1, + origin: "planner", + actor: { kind: "agent" }, + at: 4, + source: source(agent, "resolution"), + text: "Use an outline", + explicit: true, + }; + expect(() => applyInference(state, decision, agent)).toThrow( + "decisions are recorded on the card", + ); + let stance: ConversationPlan.Event = { + id: "agent1:0:stance.changed:v1", + type: "stance.changed", + threadId: "t1", + observedThreadVersion: 1, + origin: "planner", + actor: { kind: "agent" }, + at: 4, + source: source(agent, "support"), + position: "support", + }; + expect(() => applyInference(state, stance, agent)).toThrow(/human|member/i); + }); +}); diff --git a/apps/server/src/conversation-plan/domain.ts b/apps/server/src/conversation-plan/domain.ts new file mode 100644 index 00000000..5619a230 --- /dev/null +++ b/apps/server/src/conversation-plan/domain.ts @@ -0,0 +1,6 @@ +export { applyInference, completeAnalysis, enqueue, replay, retryMessage } from "./domain-analysis"; +export { applyCorrection } from "./domain-corrections"; +export { initialState } from "./domain-initial"; +export { actOnResearchOffer, offerResearch } from "./domain-research"; +export { validateStateSources } from "./state-provenance"; +export { restoreState } from "./state-restoration"; diff --git a/apps/server/src/conversation-plan/event-contributions.ts b/apps/server/src/conversation-plan/event-contributions.ts new file mode 100644 index 00000000..c50e3393 --- /dev/null +++ b/apps/server/src/conversation-plan/event-contributions.ts @@ -0,0 +1,120 @@ +import { MAX_CONTRIBUTIONS } from "./validation-fields"; +import type { ConversationPlan } from "@chopin/protocol"; +import { activeScopedSupport, currentScopedProposal, targetsScopedProposal } from "./event-support"; + +type Event = Extract< + ConversationPlan.Event, + { + type: + | "option.added" + | "reason.added" + | "constraint.added" + | "stance.changed" + | "option.relabeled"; + } +>; + +export function applyContributionEvent( + next: ConversationPlan.State, + thread: ConversationPlan.Thread, + event: Event, +): void { + switch (event.type) { + case "option.added": + case "reason.added": + case "constraint.added": { + let kind = event.type.split(".")[0] as ConversationPlan.Contribution["kind"]; + if (event.source && event.source.role !== kind) { + throw new Error("contribution source role disagrees"); + } + if (event.origin !== "human" && event.contribution.authoring === "human-edited") { + throw new Error("inference cannot claim human wording"); + } + if ( + event.source && event.contribution.authoring === "quoted" + && event.contribution.text !== event.source.quote + ) throw new Error("quoted wording must match source"); + if ( + next.threads.some((item) => + item.contributions.some((value) => value.id === event.contribution.id) + ) + ) throw new Error("duplicate contribution ID"); + if (thread.contributions.length >= MAX_CONTRIBUTIONS) { + throw new Error("conversation thread contribution limit reached"); + } + if ( + event.contribution.targetId && event.contribution.targetId !== thread.id + && !thread.contributions.some((value) => value.id === event.contribution.targetId) + ) throw new Error("unknown contribution target"); + thread.contributions.push({ + ...event.contribution, + kind, + sources: event.source ? [event.source] : [], + actor: event.actor, + }); + break; + } + case "stance.changed": { + if (event.source.author.kind !== "member") { + throw new Error("human member required for stance"); + } + if ( + event.optionId + && !thread.contributions.some((item) => + item.id === event.optionId && item.kind === "option" + ) + ) throw new Error("unknown stance option"); + if ( + (event.position === "support" && event.source.role !== "support") + || (event.position === "oppose" && event.source.role !== "objection") + || (event.position === "neutral" && event.source.role !== "withdrawal") + ) throw new Error("stance source role disagrees"); + let stance: ConversationPlan.Stance = { + id: event.id, + participant: event.source.author.handle, + optionId: event.optionId, + position: event.position, + sources: [event.source], + at: event.at, + }; + thread.stanceHistory.push(stance); + thread.stances = thread.stances.filter((item) => + item.participant !== stance.participant || item.optionId !== stance.optionId + ); + thread.stances.push(stance); + if (thread.pendingScopedChoice) { + let proposal = currentScopedProposal(thread, next.events); + if ( + proposal && targetsScopedProposal(event, proposal.id, proposal.optionId) + && activeScopedSupport([...next.events, event], proposal).length === 0 + ) thread.pendingScopedChoice = undefined; + } + break; + } + case "option.relabeled": { + if (!thread.questionnaireId || thread.status === "decided" || thread.status === "discarded") { + throw new Error("option's decision is not open and linked"); + } + let option = thread.contributions.find(item => + item.kind === "option" && item.id === event.optionId + ); + if (!option) throw new Error("unknown relabel option"); + if ( + option.displayLabel === event.label || !option.displayLabel && option.text === event.label + ) { + throw new Error("option label is unchanged"); + } + let folded = event.label.toLocaleLowerCase(); + if ( + thread.contributions.some(item => + item.kind === "option" && item.id !== option.id + && (item.displayLabel ?? item.text).trim().toLocaleLowerCase() === folded + ) + ) { + throw new Error("duplicate option label"); + } + option.displayLabel = event.label; + break; + } + } +} diff --git a/apps/server/src/conversation-plan/event-corrections.ts b/apps/server/src/conversation-plan/event-corrections.ts new file mode 100644 index 00000000..7bd8b982 --- /dev/null +++ b/apps/server/src/conversation-plan/event-corrections.ts @@ -0,0 +1,111 @@ +import { MAX_CONTRIBUTIONS } from "./validation-fields"; +import type { Chat, ConversationPlan } from "@chopin/protocol"; + +type Event = Extract; +type Member = Extract; + +export function applyCorrectionEvent( + next: ConversationPlan.State, + thread: ConversationPlan.Thread, + event: Event, +): void { + switch (event.type) { + case "card.corrected": { + let change = event.change; + if (change.kind === "edit" && change.field === "question") { + thread.question = change.text; + thread.questionAuthoring = "human-edited"; + thread.questionEditedBy = (event.actor as Member).handle; + } else if (change.kind === "edit" && change.field === "contribution") { + let contribution = thread.contributions.find((item) => item.id === change.contributionId); + if (!contribution) throw new Error("unknown contribution"); + contribution.text = change.text; + contribution.authoring = "human-edited"; + contribution.editedBy = (event.actor as Member).handle; + } else if (change.kind === "edit" && change.field === "decision") { + if (!thread.decision) throw new Error("no current decision"); + thread.decision.text = change.text; + thread.decision.editedBy = (event.actor as Member).handle; + let historical = thread.decisionHistory.find((item) => item.id === thread.decision?.id); + if (historical) { + historical.text = change.text; + historical.editedBy = (event.actor as Member).handle; + } + } else if (change.kind === "set-status") { + if (change.status === "reopened" && (thread.status !== "decided" || !thread.decision)) { + throw new Error("only a decided thread can reopen"); + } + thread.status = change.status; + thread.decision = undefined; + } else if (change.kind === "move") { + let target = next.threads.find((item) => item.id === change.targetThreadId); + if (!target || target.id === thread.id || target.version !== change.targetVersion) { + throw new Error("stale target thread version"); + } + if (target.contributions.length >= MAX_CONTRIBUTIONS) { + throw new Error("conversation thread contribution limit reached"); + } + let index = thread.contributions.findIndex((item) => item.id === change.contributionId); + if (index < 0) throw new Error("unknown contribution"); + let contribution = thread.contributions[index]; + if ( + contribution.kind === "option" + && (thread.stanceHistory.some((item) => item.optionId === contribution.id) + || thread.contributions.some((item) => item.targetId === contribution.id) + || thread.decisionHistory.some((item) => item.optionId === contribution.id)) + ) throw new Error("referenced option cannot move alone"); + thread.contributions.splice(index, 1); + contribution.targetId = target.id; + target.contributions.push(contribution); + target.version++; + } else if (change.kind === "retarget-stance" || change.kind === "dismiss-stance") { + let index = thread.stances.findIndex((item) => item.id === change.stanceId); + if (index < 0) throw new Error("stance is not current"); + let stance = thread.stances[index]; + if (change.kind === "retarget-stance") { + if ( + change.optionId !== undefined + && !thread.contributions.some((item) => + item.id === change.optionId && item.kind === "option" + ) + ) throw new Error("unknown stance option"); + if (stance.optionId === change.optionId) throw new Error("stance target is unchanged"); + if ( + thread.stances.some((item) => + item.id !== stance.id && item.participant === stance.participant + && item.optionId === change.optionId + ) + ) throw new Error("participant already has a current stance on target"); + let corrected: ConversationPlan.Stance = { + ...stance, + id: event.id, + optionId: change.optionId, + corrects: stance.id, + correctedBy: (event.actor as Member).handle, + }; + thread.stances.splice(index, 1, corrected); + thread.stanceHistory.push(corrected); + } else { + thread.stances.splice(index, 1); + } + } else if (change.kind === "retarget-contribution") { + let contribution = thread.contributions.find((item) => item.id === change.contributionId); + if (!contribution || contribution.kind === "option") { + throw new Error("target correction requires a reason or constraint"); + } + if ( + change.targetId !== thread.id + && !thread.contributions.some((item) => + item.id === change.targetId && item.kind === "option" + ) + ) throw new Error("unknown contribution target"); + if (contribution.targetId === change.targetId) { + throw new Error("contribution target is unchanged"); + } + contribution.targetId = change.targetId; + contribution.targetEditedBy = (event.actor as Member).handle; + } + break; + } + } +} diff --git a/apps/server/src/conversation-plan/event-lifecycle.ts b/apps/server/src/conversation-plan/event-lifecycle.ts new file mode 100644 index 00000000..19571045 --- /dev/null +++ b/apps/server/src/conversation-plan/event-lifecycle.ts @@ -0,0 +1,129 @@ +import type { Chat, ConversationPlan } from "@chopin/protocol"; + +type Event = Extract< + ConversationPlan.Event, + { + type: + | "thread.leaning" + | "decision.recorded" + | "decision.reopened" + | "candidate.proposed" + | "candidate.confirmed" + | "candidate.rejected" + | "card.linked" + | "thread.discarded"; + } +>; +type Member = Extract; + +export function applyLifecycleEvent(thread: ConversationPlan.Thread, event: Event): void { + switch (event.type) { + case "thread.leaning": { + if (thread.status === "decided" || thread.status === "discarded") { + throw new Error("cannot infer leaning over a closed thread"); + } + let supporters = new Set( + thread.stances.filter((item) => + item.position === "support" && item.optionId === event.optionId + ).map((item) => item.participant), + ); + if (supporters.size < 2) throw new Error("leaning needs several participant stances"); + thread.status = "leaning"; + break; + } + case "decision.recorded": { + if (thread.status === "discarded") throw new Error("thread is discarded"); + if (thread.status === "decided") throw new Error("decision must be reopened first"); + if ( + event.origin !== "human" + && (event.source?.author.kind !== "member" || event.source.role !== "resolution") + ) throw new Error("decision requires explicit human resolution"); + if ( + event.optionId + && !thread.contributions.some((item) => + item.id === event.optionId && item.kind === "option" + ) + ) throw new Error("unknown decision option"); + if ( + event.origin !== "human" + && thread.stances.some((item) => + item.position === "oppose" && (!item.optionId || item.optionId === event.optionId) + ) + ) throw new Error("contradictory stance needs human resolution"); + let decision: ConversationPlan.Decision = { + id: event.id, + text: event.text, + optionId: event.optionId, + sources: event.source ? [event.source] : [], + actor: event.origin === "human" ? event.actor as Member : event.source!.author as Member, + at: event.at, + }; + thread.decision = decision; + thread.decisionHistory.push(decision); + thread.status = "decided"; + thread.pendingSettle = undefined; + thread.pendingScopedChoice = undefined; + break; + } + case "decision.reopened": + if (thread.status !== "decided") throw new Error("only a decided thread can reopen"); + if ( + event.origin !== "human" + && (event.source?.author.kind !== "member" || event.source.role !== "reopening") + ) throw new Error("reopening requires explicit human statement"); + thread.decision = undefined; + thread.status = "reopened"; + thread.pendingSettle = undefined; + break; + case "candidate.proposed": + if (event.source.role !== event.candidate.kind) { + throw new Error("candidate source role disagrees"); + } + if (thread.candidates.some((candidate) => candidate.id === event.candidate.id)) { + throw new Error("duplicate candidate ID"); + } + if ((event.candidate.kind === "reopening") !== (thread.status === "decided")) { + throw new Error("candidate does not match thread status"); + } + thread.candidates.push({ ...event.candidate, sources: [event.source], status: "pending" }); + break; + case "candidate.confirmed": + case "candidate.rejected": { + let candidate = thread.candidates.find((item) => item.id === event.candidateId); + if (!candidate || candidate.status !== "pending") throw new Error("candidate is not pending"); + if ( + event.type === "candidate.confirmed" + && ((candidate.kind === "reopening") !== (thread.status === "decided")) + ) throw new Error("candidate no longer matches thread status"); + candidate.status = event.type === "candidate.confirmed" ? "confirmed" : "rejected"; + candidate.actedBy = (event.actor as Member).handle; + if (event.type === "candidate.confirmed") { + if (candidate.kind === "resolution") { + let decision: ConversationPlan.Decision = { + id: event.id, + text: candidate.text, + sources: candidate.sources, + actor: event.actor as Member, + at: event.at, + }; + thread.decision = decision; + thread.decisionHistory.push(decision); + thread.status = "decided"; + } else { + thread.decision = undefined; + thread.status = "reopened"; + } + } + break; + } + case "card.linked": + if (thread.questionnaireId) throw new Error("thread already has a card"); + thread.questionnaireId = event.questionnaireId; + break; + case "thread.discarded": + thread.status = "discarded"; + thread.pendingSettle = undefined; + thread.pendingScopedChoice = undefined; + break; + } +} diff --git a/apps/server/src/conversation-plan/event-settlement.ts b/apps/server/src/conversation-plan/event-settlement.ts new file mode 100644 index 00000000..375db462 --- /dev/null +++ b/apps/server/src/conversation-plan/event-settlement.ts @@ -0,0 +1,160 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { isDeepStrictEqual } from "node:util"; +import { activeSettleDeferral } from "./preference"; +import { activeScopedSupport, currentScopedProposal } from "./event-support"; + +type Event = Extract< + ConversationPlan.Event, + { + type: + | "settle.suggested" + | "settle.agreed" + | "settle.deferred" + | "settle.resumed" + | "scoped-choice.proposed" + | "scoped-choice.agreed" + | "scoped-choice.saved"; + } +>; + +export function applySettlementEvent( + next: ConversationPlan.State, + thread: ConversationPlan.Thread, + event: Event, +): void { + switch (event.type) { + case "settle.suggested": + if (event.source.role !== "resolution") throw new Error("settle source role disagrees"); + if (event.source.author.kind !== "member") throw new Error("human member required to settle"); + if (thread.status === "decided" || thread.status === "discarded") { + throw new Error("thread is not open for settling"); + } + if ( + !thread.contributions.some((item) => item.id === event.optionId && item.kind === "option") + ) { + throw new Error("unknown settle option"); + } + if (!activeSettleDeferral(thread, next.events)) { + thread.pendingSettle = { + optionId: event.optionId, + proposer: event.source.author.handle, + messageId: event.source.messageId, + }; + } + break; + case "settle.agreed": + if (event.source.role !== "support") throw new Error("agreement source role disagrees"); + if (!thread.pendingSettle) throw new Error("no pending proposal to settle"); + if (thread.pendingSettle.optionId !== event.optionId) { + throw new Error("agreement names another option"); + } + if ( + event.source.author.kind !== "member" + || event.source.author.handle === thread.pendingSettle.proposer + ) throw new Error("agreement needs another member"); + break; + case "settle.deferred": { + let pending = thread.pendingSettle; + let proposal = next.events.find(item => + item.id === event.proposalId && item.type === "settle.suggested" + && item.threadId === thread.id + ); + if ( + !pending || !proposal || proposal.type !== "settle.suggested" + || proposal.optionId !== pending.optionId + || proposal.source.messageId !== pending.messageId + || proposal.source.author.kind !== "member" + || proposal.source.author.handle !== pending.proposer + || event.source.author.kind !== "member" + || event.source.author.handle !== pending.proposer + || !next.events.some(item => + item.type === "stance.changed" && item.threadId === thread.id + && item.source.messageId === event.source.messageId + && item.source.author.kind === "member" + && item.source.author.handle === pending.proposer + && item.optionId === pending.optionId && item.position === "neutral" + ) + || activeSettleDeferral(thread, next.events) + ) throw new Error("deferral does not target the active proposal"); + break; + } + case "settle.resumed": { + let deferred = activeSettleDeferral(thread, next.events); + if ( + !deferred || deferred.id !== event.deferredEventId + || deferred.proposalId !== event.proposalId + || !thread.pendingSettle || event.source.author.kind !== "member" + ) throw new Error("verification does not target the active deferral"); + break; + } + case "scoped-choice.proposed": + if ( + thread.status === "decided" || thread.status === "discarded" + || thread.questionnaireId !== event.cardId + || event.source.author.kind !== "member" + ) throw new Error("scoped choice needs an open linked card"); + thread.pendingScopedChoice = { + proposalId: event.id, + cardId: event.cardId, + optionId: event.optionId, + label: event.label, + scope: "spike", + proposer: event.source.author.handle, + messageId: event.source.messageId, + }; + break; + case "scoped-choice.agreed": { + let pending = thread.pendingScopedChoice; + let proposal = currentScopedProposal(thread, next.events); + if ( + thread.status === "decided" || thread.status === "discarded" + || !pending || !proposal || proposal.id !== event.proposalId + || pending.proposalId !== undefined && pending.proposalId !== event.proposalId + || thread.questionnaireId !== event.cardId + || pending.cardId !== event.cardId || pending.optionId !== event.optionId + || pending.label !== event.label || pending.scope !== event.scope + || event.source.author.kind !== "member" + || event.source.author.handle === pending.proposer + ) throw new Error("scoped choice agreement does not match the current proposal"); + break; + } + case "scoped-choice.saved": { + let pending = thread.pendingScopedChoice; + let proposal = currentScopedProposal(thread, next.events); + let agreement = event.agreementId + ? next.events.find(item => item.id === event.agreementId) + : undefined; + let active = proposal ? activeScopedSupport(next.events, proposal) : []; + let lineageMatches = event.supportEventIds !== undefined + ? active.length > 0 + && isDeepStrictEqual(event.supportEventIds, active.map(item => item.id)) + && isDeepStrictEqual(event.sources, active.map(item => item.source)) + : active.length > 0 + && isDeepStrictEqual(active.map(item => item.id), [ + proposal?.id, + ...(agreement && agreement.type === "scoped-choice.agreed" + ? [agreement.id] + : []), + ]) + && isDeepStrictEqual(event.sources, active.map(item => item.source)); + if ( + thread.status === "decided" || thread.status === "discarded" + || !pending || !proposal || proposal.id !== event.proposalId + || thread.questionnaireId !== event.cardId + || pending.cardId !== event.cardId || pending.optionId !== event.optionId + || pending.label !== event.label || pending.scope !== event.scope + || next.events.some(item => + item.type === "scoped-choice.saved" && item.proposalId === event.proposalId + ) + || !lineageMatches + || event.agreementId !== undefined && ( + agreement?.type !== "scoped-choice.agreed" + || agreement.proposalId !== proposal.id || agreement.threadId !== thread.id + || agreement.cardId !== event.cardId || agreement.optionId !== event.optionId + || agreement.label !== event.label || agreement.scope !== event.scope + ) + ) throw new Error("saved scoped choice does not match the current proposal"); + break; + } + } +} diff --git a/apps/server/src/conversation-plan/event-support.test.ts b/apps/server/src/conversation-plan/event-support.test.ts new file mode 100644 index 00000000..3594dc43 --- /dev/null +++ b/apps/server/src/conversation-plan/event-support.test.ts @@ -0,0 +1,159 @@ +import { expect, test } from "bun:test"; +import type { ConversationPlan } from "@chopin/protocol"; +import { + activeScopedSupport, + applyEvent, + currentScopedProposal, + targetsScopedProposal, +} from "./events"; + +type Event = ConversationPlan.Event; +function source( + role: ConversationPlan.SourceRole, + quote: string, + handle = "maggie", + messageId = "message-1", +): ConversationPlan.SourceRef { + return { + messageId, + author: { kind: "member", handle }, + quote, + start: 0, + end: quote.length, + role, + }; +} +function base(id: string, version = 0): ConversationPlan.EventBase { + return { + id, + threadId: "thread-1", + observedThreadVersion: version, + origin: "classifier", + actor: { kind: "classifier" }, + at: 1, + }; +} +function scoped() { + let state: ConversationPlan.State = { + schemaVersion: 1, + revision: 0, + events: [], + threads: [], + queue: [], + analysis: [], + }; + state = applyEvent(state, { + ...base("open"), + type: "thread.opened", + source: source("question", "Which?"), + question: "Which?", + }); + state = applyEvent(state, { ...base("link", 1), type: "card.linked", questionnaireId: "card-1" }); + let proposal: Extract = { + ...base("proposal", 2), + type: "scoped-choice.proposed", + source: source("support", "I'd pick Bun for the spike."), + cardId: "card-1", + optionId: "option-1", + label: "Bun", + scope: "spike", + }; + state = applyEvent(state, proposal); + return { state, proposal }; +} +function agreement( + version: number, + id = "agreement", +): Extract { + return { + ...base(id, version), + type: "scoped-choice.agreed", + source: source("support", "yep, Bun for the spike.", "alex", id), + proposalId: "proposal", + cardId: "card-1", + optionId: "option-1", + label: "Bun", + scope: "spike", + }; +} + +test("support helpers resolve the current proposal and latest member evidence in event order", () => { + let { state, proposal } = scoped(); + let first = agreement(3); + state = applyEvent(state, first); + let latest = agreement(4, "latest-agreement"); + state = applyEvent(state, latest); + expect(currentScopedProposal(state.threads[0]!, state.events)?.id).toBe(proposal.id); + expect(activeScopedSupport(state.events, proposal).map(event => event.id)).toEqual([ + "proposal", + "latest-agreement", + ]); + let legacy = structuredClone(state.threads[0]!); + delete legacy.pendingScopedChoice!.proposalId; + expect(currentScopedProposal(legacy, state.events)?.id).toBe(proposal.id); + expect(currentScopedProposal(legacy, [...state.events, { ...proposal, id: "ambiguous" }])) + .toBeUndefined(); +}); + +test("withdrawal targets retain legacy absence while explicit null does not retract scoped support", () => { + let { state, proposal } = scoped(); + let withdrawal: Extract = { + ...base("withdraw", 3), + type: "stance.changed", + source: source("withdrawal", "I withdraw"), + position: "neutral", + }; + expect(targetsScopedProposal(withdrawal, proposal.id, proposal.optionId)).toBe(true); + expect( + targetsScopedProposal( + { ...withdrawal, scopedProposalId: null }, + proposal.id, + proposal.optionId, + ), + ).toBe(false); + expect( + targetsScopedProposal( + { ...withdrawal, scopedProposalId: "other" }, + proposal.id, + proposal.optionId, + ), + ).toBe(false); + expect( + activeScopedSupport([...state.events, { ...withdrawal, scopedProposalId: null }], proposal).map( + event => event.id, + ), + ).toEqual(["proposal"]); + state = applyEvent(state, withdrawal); + expect(activeScopedSupport(state.events, proposal)).toEqual([]); + expect(state.threads[0]!.pendingScopedChoice).toBeUndefined(); +}); + +test("scoped saves bind exact ordered active sources and preserve the provisional thread", () => { + let { state, proposal } = scoped(); + let agreed = agreement(3); + state = applyEvent(state, agreed); + let save: Extract = { + ...base("save", 4), + origin: "human", + actor: { kind: "member", handle: "maggie" }, + type: "scoped-choice.saved", + proposalId: proposal.id, + supportEventIds: [proposal.id, agreed.id], + cardId: "card-1", + optionId: "option-1", + label: "Bun", + scope: "spike", + sources: [proposal.source, agreed.source], + expectedGeneration: 0, + }; + expect(() => applyEvent(state, { ...save, supportEventIds: [agreed.id, proposal.id] })).toThrow( + "saved scoped choice does not match the current proposal", + ); + let saved = applyEvent(state, save); + expect(saved.threads[0]!.status).toBe("exploring"); + expect(saved.threads[0]!.decision).toBeUndefined(); + expect(saved.threads[0]!.pendingScopedChoice?.proposalId).toBe(proposal.id); + expect(() => applyEvent(saved, { ...save, id: "save-again", observedThreadVersion: 5 })).toThrow( + "saved scoped choice does not match the current proposal", + ); +}); diff --git a/apps/server/src/conversation-plan/event-support.ts b/apps/server/src/conversation-plan/event-support.ts new file mode 100644 index 00000000..0c7814d8 --- /dev/null +++ b/apps/server/src/conversation-plan/event-support.ts @@ -0,0 +1,69 @@ +import type { ConversationPlan } from "@chopin/protocol"; + +type Event = ConversationPlan.Event; + +export function currentScopedProposal( + thread: ConversationPlan.Thread, + events: readonly Event[], +): Extract | undefined { + let pending = thread.pendingScopedChoice; + if (!pending) return; + let matches = events.filter(( + event, + ): event is Extract => + event.type === "scoped-choice.proposed" && event.threadId === thread.id + && event.cardId === pending.cardId && event.optionId === pending.optionId + && event.label === pending.label && event.scope === pending.scope + && event.source.messageId === pending.messageId + && event.source.author.kind === "member" + && event.source.author.handle === pending.proposer + ); + return pending.proposalId + ? matches.find(event => event.id === pending.proposalId) + : matches.length === 1 + ? matches[0] + : undefined; +} + +/** An absent link preserves previously persisted optionless stance behavior. */ +export function targetsScopedProposal( + event: Extract, + proposalId: string | undefined, + optionId: string, +): boolean { + if (event.position === "support" || event.scopedProposalId === null) return false; + if (event.optionId !== undefined && event.optionId !== optionId) return false; + return event.scopedProposalId === undefined || event.scopedProposalId === proposalId; +} + +type ScopedSupportEvent = Extract; + +/** Accepted support in event order, with each member's latest active source. */ +export function activeScopedSupport( + events: readonly Event[], + proposal: Extract, +): ScopedSupportEvent[] { + let start = events.findIndex(item => item.id === proposal.id); + if (start < 0 || proposal.source.author.kind !== "member") return []; + let active = new Map(); + for (let item of events.slice(start)) { + if (item.id === proposal.id) { + active.set(proposal.source.author.handle, proposal); + } else if ( + item.type === "scoped-choice.agreed" && item.threadId === proposal.threadId + && item.proposalId === proposal.id && item.cardId === proposal.cardId + && item.optionId === proposal.optionId && item.label === proposal.label + && item.scope === proposal.scope && item.source.author.kind === "member" + ) { + active.delete(item.source.author.handle); + active.set(item.source.author.handle, item); + } else if ( + item.type === "stance.changed" && item.threadId === proposal.threadId + && item.source.author.kind === "member" + && targetsScopedProposal(item, proposal.id, proposal.optionId) + ) active.delete(item.source.author.handle); + } + return [...active.values()]; +} diff --git a/apps/server/src/conversation-plan/event-validation.ts b/apps/server/src/conversation-plan/event-validation.ts new file mode 100644 index 00000000..cf1a35d4 --- /dev/null +++ b/apps/server/src/conversation-plan/event-validation.ts @@ -0,0 +1,304 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { ULID } from "@chopin/dialect"; +import { limits as questionLimits } from "@chopin/question"; +import { assertSourceShape } from "./sources"; +import { assertCorrectionChange } from "./correction-validation"; +import { + actor, + id, + knownKeys, + MAX_EVENTS, + record, + spikeAgreementLabel, + spikePreferenceLabel, + text, + version, +} from "./validation-fields"; + +export function assertEventShape(value: unknown): asserts value is ConversationPlan.Event { + let event = record(value); + let base = ["id", "type", "threadId", "observedThreadVersion", "origin", "actor", "at"]; + id(event.id); + id(event.threadId); + version(event.observedThreadVersion); + version(event.at); + actor(event.actor); + if (!["classifier", "planner", "human"].includes(event.origin as string)) { + throw new Error("invalid event origin"); + } + let kind = (event.actor as { kind: string }).kind; + if ( + (event.origin === "classifier" && kind !== "classifier") + || (event.origin === "planner" && kind !== "agent") + || (event.origin === "human" && kind !== "member") + ) throw new Error("event origin and actor disagree"); + if (event.source !== undefined) assertSourceShape(event.source); + let humanExcerpt = event.origin === "human" + && ["option.added", "reason.added", "constraint.added"].includes(event.type as string) + && event.source !== undefined + && (event.contribution as { authoring?: unknown } | undefined)?.authoring === "quoted"; + if ( + event.origin === "human" && (!humanExcerpt && event.source !== undefined || ![ + "decision.recorded", + "decision.reopened", + "thread.discarded", + "scoped-choice.saved", + "option.added", + "reason.added", + "constraint.added", + "candidate.confirmed", + "candidate.rejected", + "card.corrected", + ].includes(event.type as string)) + ) throw new Error("human action cannot claim a message quote"); + switch (event.type) { + case "thread.opened": + knownKeys(event, [...base, "source", "question"]); + if (event.source === undefined) { + if (event.origin !== "planner") throw new Error("opening requires a question source"); + } else assertSourceShape(event.source); + text(event.question, event.source === undefined ? 1_000 : 500); + break; + case "option.added": + case "reason.added": + case "constraint.added": { + knownKeys(event, [...base, "source", "contribution"]); + let cardOption = event.type === "option.added" && event.origin === "planner" + && typeof event.id === "string" + && /^card:[0-9A-HJKMNP-TV-Z]{26}:option:[0-9A-HJKMNP-TV-Z]{26}$/.test(event.id); + let inferred = event.origin === "classifier" + || (event.origin === "planner" && event.source !== undefined + && /^classifier:[0-9a-f]{24}$/.test(event.id as string)); + if (inferred) { + assertSourceShape(event.source); + if (event.origin === "planner" && event.source.author.kind !== "agent") { + throw new Error("Planner quote must come from an agent message"); + } + } else if ( + !humanExcerpt && (event.type !== "option.added" + || event.source !== undefined && !cardOption) + ) { + throw new Error("only card options can omit a source"); + } + let contribution = record(event.contribution); + knownKeys(contribution, ["id", "text", "authoring", "targetId", "relation"]); + id(contribution.id); + if (cardOption && !String(event.id).endsWith(`:option:${contribution.id}`)) { + throw new Error("card option event ID disagrees with contribution"); + } + if (!inferred && !ULID.test(contribution.id as string)) { + throw new Error("card option IDs must be ULIDs"); + } + text(contribution.text); + if (!["quoted", "scribe", "human-edited"].includes(contribution.authoring as string)) { + throw new Error("invalid contribution authoring"); + } + if ( + humanExcerpt && contribution.authoring !== "quoted" + || !inferred && !humanExcerpt && ( + (event.origin === "human" && contribution.authoring !== "human-edited") + || (event.origin === "planner" && contribution.authoring !== "scribe") + ) + ) throw new Error("card option authoring disagrees with origin"); + if (contribution.targetId !== undefined) id(contribution.targetId); + if ( + contribution.relation !== undefined + && !["supports", "challenges", "qualifies"].includes(contribution.relation as string) + ) throw new Error("invalid contribution relation"); + break; + } + case "stance.changed": + knownKeys(event, [...base, "source", "optionId", "scopedProposalId", "position"]); + assertSourceShape(event.source); + if (event.optionId !== undefined) id(event.optionId); + if (event.scopedProposalId !== undefined && event.scopedProposalId !== null) { + id(event.scopedProposalId); + } + if (!["support", "oppose", "neutral"].includes(event.position as string)) { + throw new Error("invalid stance"); + } + break; + case "option.relabeled": + knownKeys(event, [...base, "optionId", "label", "observedCardRevision"]); + id(event.optionId); + version(event.observedCardRevision); + text(event.label, questionLimits.MAX_LABEL); + let label = event.label as string; + if ( + event.origin !== "planner" || label !== label.trim() + || /[\r\n]/.test(label) + || /\b(?:and|or|versus|vs\.?)\b/i.test(label) + || /,[^,]+,|,[^,]+\band\b/i.test(label) + || /(?<=[a-z0-9])\s*\/\s*(?=[a-z0-9])/i.test(label) + ) throw new Error("invalid atomic option label"); + break; + case "thread.leaning": + knownKeys(event, [...base, "source", "optionId"]); + assertSourceShape(event.source); + if (event.optionId !== undefined) id(event.optionId); + break; + case "card.linked": + knownKeys(event, [...base, "questionnaireId"]); + id(event.questionnaireId); + if (event.origin !== "classifier") throw new Error("card link is a system event"); + break; + case "settle.suggested": + knownKeys(event, [...base, "source", "optionId"]); + assertSourceShape(event.source); + id(event.optionId); + if (event.origin !== "classifier") throw new Error("settle suggestion is inferred"); + break; + case "settle.agreed": + knownKeys(event, [...base, "source", "optionId"]); + assertSourceShape(event.source); + id(event.optionId); + if (event.origin !== "classifier") throw new Error("agreement is inferred"); + break; + case "settle.deferred": + knownKeys(event, [...base, "source", "proposalId"]); + assertSourceShape(event.source); + id(event.proposalId); + if (event.origin !== "classifier" || event.source.role !== "constraint") { + throw new Error("deferral requires a sourced constraint"); + } + break; + case "settle.resumed": + knownKeys(event, [...base, "source", "proposalId", "deferredEventId"]); + assertSourceShape(event.source); + id(event.proposalId); + id(event.deferredEventId); + if (event.origin !== "classifier" || event.source.role !== "verification") { + throw new Error("resume requires sourced verification"); + } + break; + case "scoped-choice.proposed": + knownKeys(event, [...base, "source", "cardId", "optionId", "label", "scope"]); + assertSourceShape(event.source); + id(event.cardId); + id(event.optionId); + text(event.label, questionLimits.MAX_LABEL); + if ( + event.origin !== "classifier" || event.scope !== "spike" + || event.source.role !== "support" || event.source.author.kind !== "member" + || spikePreferenceLabel(event.source.quote)?.toLocaleLowerCase() + !== (event.label as string).toLocaleLowerCase() + ) throw new Error("invalid scoped choice proposal"); + break; + case "scoped-choice.agreed": + knownKeys(event, [ + ...base, + "source", + "proposalId", + "cardId", + "optionId", + "label", + "scope", + ]); + assertSourceShape(event.source); + id(event.proposalId); + id(event.cardId); + id(event.optionId); + text(event.label, questionLimits.MAX_LABEL); + if ( + event.origin !== "classifier" || event.scope !== "spike" + || event.source.role !== "support" || event.source.author.kind !== "member" + || spikeAgreementLabel(event.source.quote)?.toLocaleLowerCase() + !== (event.label as string).toLocaleLowerCase() + ) throw new Error("invalid scoped choice agreement"); + break; + case "scoped-choice.saved": { + knownKeys(event, [ + ...base, + "proposalId", + "supportEventIds", + "agreementId", + "cardId", + "optionId", + "label", + "scope", + "sources", + "expectedGeneration", + "expectedLabel", + ]); + id(event.proposalId); + if (event.supportEventIds !== undefined) { + if ( + !Array.isArray(event.supportEventIds) || event.supportEventIds.length < 1 + || event.supportEventIds.length > MAX_EVENTS + || new Set(event.supportEventIds).size !== event.supportEventIds.length + ) throw new Error("invalid saved scoped choice source IDs"); + for (let sourceId of event.supportEventIds) id(sourceId); + } + if (event.agreementId !== undefined) id(event.agreementId); + id(event.cardId); + id(event.optionId); + text(event.label, questionLimits.MAX_LABEL); + version(event.expectedGeneration); + if (event.expectedLabel !== undefined) text(event.expectedLabel, questionLimits.MAX_LABEL); + if ( + event.origin !== "human" || event.scope !== "spike" + || !Array.isArray(event.sources) + || (event.supportEventIds === undefined + ? event.sources.length !== (event.agreementId === undefined ? 1 : 2) + : event.agreementId !== undefined + || event.sources.length !== event.supportEventIds.length) + ) throw new Error("invalid saved scoped choice"); + let members = new Set(); + for (let source of event.sources) { + assertSourceShape(source); + if (source.role !== "support" || source.author.kind !== "member") { + throw new Error("invalid saved scoped choice source"); + } + if (event.supportEventIds !== undefined) { + if (members.has(source.author.handle)) throw new Error("duplicate scoped supporter"); + members.add(source.author.handle); + } + } + break; + } + case "thread.discarded": + knownKeys(event, base); + if (event.origin !== "human") throw new Error("only a person discards"); + break; + case "decision.recorded": + knownKeys(event, [...base, "source", "text", "optionId", "explicit"]); + text(event.text); + if (event.optionId !== undefined) id(event.optionId); + if (event.explicit !== true) throw new Error("decision requires explicit resolution"); + if (event.origin !== "human") assertSourceShape(event.source); + break; + case "decision.reopened": + knownKeys(event, [...base, "source", "explicit"]); + if (event.explicit !== true) throw new Error("reopening requires explicit statement"); + if (event.origin !== "human") assertSourceShape(event.source); + break; + case "candidate.proposed": { + knownKeys(event, [...base, "source", "candidate"]); + assertSourceShape(event.source); + let candidate = record(event.candidate); + knownKeys(candidate, ["id", "kind", "text"]); + id(candidate.id); + text(candidate.text); + if (!["resolution", "reopening"].includes(candidate.kind as string)) { + throw new Error("invalid candidate kind"); + } + break; + } + case "candidate.confirmed": + case "candidate.rejected": + knownKeys(event, [...base, "candidateId"]); + id(event.candidateId); + if (event.origin !== "human") throw new Error("candidate action requires human"); + break; + case "card.corrected": + knownKeys(event, [...base, "change"]); + if (event.origin !== "human") throw new Error("card correction requires human"); + assertCorrectionChange(event.change); + if ( + ["record-decision", "confirm-candidate", "reject-candidate"].includes(event.change.kind) + ) throw new Error("invalid card correction event"); + break; + default: + throw new Error("invalid conversation plan event type"); + } +} diff --git a/apps/server/src/conversation-plan/events.test.ts b/apps/server/src/conversation-plan/events.test.ts new file mode 100644 index 00000000..a2c3d425 --- /dev/null +++ b/apps/server/src/conversation-plan/events.test.ts @@ -0,0 +1,280 @@ +import { expect, test } from "bun:test"; +import type { ConversationPlan } from "@chopin/protocol"; +import { applyEvent } from "./events"; +import { effectivePending, effectivePreference } from "./preference"; + +function empty(): ConversationPlan.State { + return { schemaVersion: 1, revision: 0, events: [], threads: [], queue: [], analysis: [] }; +} +function opening( + threadId = "thread-1", + id = "open-1", +): Extract { + return { + id, + threadId, + observedThreadVersion: 0, + origin: "planner", + actor: { kind: "agent" }, + at: 1, + type: "thread.opened", + question: "Which runtime?", + }; +} + +test("replay clones an accepted event and advances the initial thread and revision", () => { + let before = empty(); + let event = opening(); + let after = applyEvent(before, event); + expect(before).toEqual(empty()); + expect(after.revision).toBe(1); + expect(after.threads[0]!.version).toBe(1); + expect(after.threads[0]!.question).toBe("Which runtime?"); + event.question = "Mutated caller input"; + expect(after.threads[0]!.question).toBe("Which runtime?"); + expect(after.events[0]).not.toBe(event); +}); + +function citation( + role: ConversationPlan.SourceRole, + quote = "Bun", + handle = "maggie", + messageId = "message-1", +): ConversationPlan.SourceRef { + return { + messageId, + author: { kind: "member", handle }, + quote, + start: 0, + end: quote.length, + role, + }; +} +function accept( + state: ConversationPlan.State, + body: Record, + threadId = "thread-1", + human = false, +): ConversationPlan.State { + return applyEvent(state, { + id: `event-${state.events.length}`, + threadId, + observedThreadVersion: state.threads.find(thread => thread.id === threadId)!.version, + origin: human ? "human" : "classifier", + actor: human ? { kind: "member", handle: "maggie" } : { kind: "classifier" }, + at: state.events.length, + ...body, + } as ConversationPlan.Event); +} +function withOption(): ConversationPlan.State { + return accept(applyEvent(empty(), opening()), { + type: "option.added", + source: citation("option"), + contribution: { id: "option-1", text: "Bun", authoring: "quoted" }, + }); +} + +test("event ID deduplication preserves the inherited return-before-validation behavior", () => { + let current = applyEvent(empty(), opening()); + let duplicate = { ...opening(), question: "Different body", observedThreadVersion: 99 }; + expect(applyEvent(current, duplicate)).toBe(current); + expect(applyEvent(current, { ...duplicate, question: "" })).toBe(current); + expect(() => + applyEvent(current, { + id: "stale", + type: "thread.discarded", + threadId: "thread-1", + observedThreadVersion: 0, + at: 1, + origin: "human", + actor: { kind: "member", handle: "maggie" }, + }) + ).toThrow("stale conversation thread version"); + expect(() => applyEvent(current, { ...opening(), id: "different" })).toThrow( + "stale thread opening", + ); + expect(() => + applyEvent(current, { + id: "missing", + type: "thread.discarded", + threadId: "missing", + observedThreadVersion: 0, + at: 1, + origin: "human", + actor: { kind: "member", handle: "maggie" }, + }) + ).toThrow("stale conversation thread version"); +}); + +test("contribution, stance, link, and relabel replay retain exact source wording", () => { + let original = withOption(); + let before = structuredClone(original); + let current = accept(original, { type: "card.linked", questionnaireId: "card-1" }); + current = accept(current, { + type: "stance.changed", + source: citation("support"), + optionId: "option-1", + position: "support", + }); + current = accept(current, { + type: "stance.changed", + source: citation("support", "Bun", "alex", "message-2"), + optionId: "option-1", + position: "support", + }); + current = accept(current, { + type: "thread.leaning", + source: citation("support"), + optionId: "option-1", + }); + current = accept(current, { + type: "option.relabeled", + optionId: "option-1", + label: "Bun runtime", + observedCardRevision: 0, + origin: "planner", + actor: { kind: "agent" }, + }); + expect(original).toEqual(before); + expect(current.threads[0]!.status).toBe("leaning"); + expect(current.threads[0]!.contributions[0]!.text).toBe("Bun"); + expect(current.threads[0]!.contributions[0]!.displayLabel).toBe("Bun runtime"); + expect(current.threads[0]!.stanceHistory).toHaveLength(2); + expect(() => + accept(current, { + type: "stance.changed", + source: citation("support"), + optionId: "unknown", + position: "support", + }) + ).toThrow("unknown stance option"); + expect(() => + accept(current, { + type: "reason.added", + source: citation("reason"), + contribution: { id: "reason-1", text: "Different", authoring: "quoted" }, + }) + ).toThrow("quoted wording must match source"); +}); + +test("candidate confirmation, explicit reopening, and discard preserve decision history", () => { + let current = withOption(); + current = accept(current, { + type: "candidate.proposed", + source: citation("resolution", "Use Bun"), + candidate: { id: "candidate-1", kind: "resolution", text: "Use Bun" }, + }); + current = accept( + current, + { type: "candidate.confirmed", candidateId: "candidate-1" }, + "thread-1", + true, + ); + expect(current.threads[0]!.status).toBe("decided"); + expect(current.threads[0]!.decision!.text).toBe("Use Bun"); + current = accept(current, { type: "decision.reopened", explicit: true }, "thread-1", true); + expect(current.threads[0]!.decision).toBeUndefined(); + expect(current.threads[0]!.decisionHistory).toHaveLength(1); + current = accept(current, { type: "thread.discarded" }, "thread-1", true); + expect(current.threads[0]!.status).toBe("discarded"); +}); + +test("settle deferral blocks effective preference until its exact verification resumes it", () => { + let current = withOption(); + current = accept(current, { + type: "settle.suggested", + source: citation("resolution", "Use Bun", "maggie", "proposal-message"), + optionId: "option-1", + }); + let proposalId = current.events.at(-1)!.id; + current = accept(current, { + type: "stance.changed", + source: citation("withdrawal", "Wait", "maggie", "defer-message"), + optionId: "option-1", + position: "neutral", + }); + current = accept(current, { + type: "settle.deferred", + source: citation("constraint", "Wait", "maggie", "defer-message"), + proposalId, + }); + let deferredEventId = current.events.at(-1)!.id; + expect(effectivePending(current.threads[0]!, current.events)).toBeUndefined(); + expect(effectivePreference(current.threads[0]!, current.events)).toBeUndefined(); + expect(() => + accept(current, { + type: "settle.resumed", + source: citation("verification", "Verified"), + proposalId, + deferredEventId: "wrong", + }) + ).toThrow("verification does not target the active deferral"); + current = accept(current, { + type: "settle.resumed", + source: citation("verification", "Verified"), + proposalId, + deferredEventId, + }); + expect(effectivePending(current.threads[0]!, current.events)?.optionId).toBe("option-1"); + expect(effectivePreference(current.threads[0]!, current.events)).toEqual({ + optionId: "option-1", + messageIds: ["proposal-message"], + }); +}); + +test("moving a contribution advances both threads once and rejects stale targets", () => { + let current = withOption(); + current = applyEvent(current, opening("thread-2", "open-2")); + let before = structuredClone(current); + let fromVersion = current.threads[0]!.version; + let targetVersion = current.threads[1]!.version; + expect(() => + accept( + current, + { + type: "card.corrected", + change: { + kind: "move", + contributionId: "option-1", + targetThreadId: "thread-2", + targetVersion: targetVersion + 1, + }, + }, + "thread-1", + true, + ) + ).toThrow("stale target thread version"); + let next = accept( + current, + { + type: "card.corrected", + change: { + kind: "move", + contributionId: "option-1", + targetThreadId: "thread-2", + targetVersion, + }, + }, + "thread-1", + true, + ); + expect(current).toEqual(before); + expect(next.threads[0]!.contributions).toHaveLength(0); + expect(next.threads[1]!.contributions[0]!.targetId).toBe("thread-2"); + expect(next.threads[0]!.version).toBe(fromVersion + 1); + expect(next.threads[1]!.version).toBe(targetVersion + 1); + expect(next.revision).toBe(current.revision + 1); + expect(next.events).toHaveLength(current.events.length + 1); + let later = accept( + next, + { + type: "card.corrected", + change: { kind: "edit", field: "question", text: "Which editor?" }, + }, + "thread-2", + true, + ); + expect(later.threads[1]!.question).toBe("Which editor?"); + expect(next.threads[1]!.question).toBe("Which runtime?"); + expect(current).toEqual(before); +}); diff --git a/apps/server/src/conversation-plan/events.ts b/apps/server/src/conversation-plan/events.ts new file mode 100644 index 00000000..9f8f71ec --- /dev/null +++ b/apps/server/src/conversation-plan/events.ts @@ -0,0 +1,110 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { assertEventShape, MAX_EVENTS, MAX_THREADS } from "./validation"; +import { applyContributionEvent } from "./event-contributions"; +import { applyLifecycleEvent } from "./event-lifecycle"; +import { applySettlementEvent } from "./event-settlement"; +import { applyCorrectionEvent } from "./event-corrections"; + +export { activeScopedSupport, currentScopedProposal, targetsScopedProposal } from "./event-support"; + +type State = ConversationPlan.State; +type Event = ConversationPlan.Event; + +export class ConversationCapacityError extends Error { + constructor() { + super("Conversation history is full"); + } +} + +/** Durable card actions reserve the event slots their FIFO mirror will need. */ +export function assertEventCapacity(state: State, reserved: number): void { + if (state.events.length + reserved > MAX_EVENTS) throw new ConversationCapacityError(); +} + +/** One accepted event changes a cloned view; callers persist the returned state. */ +export function applyEvent(state: State, event: Event): State { + if (state.events.some((accepted) => accepted.id === event.id)) return state; + assertEventShape(event); + event = structuredClone(event); + assertEventCapacity(state, 1); + let existing = state.threads.find((thread) => thread.id === event.threadId); + if (event.type === "thread.opened") { + if (existing || event.observedThreadVersion !== 0) throw new Error("stale thread opening"); + if (state.threads.length >= MAX_THREADS) throw new Error("conversation thread limit reached"); + if (event.source && event.source.role !== "question") { + throw new Error("opening requires a question source"); + } + let thread: ConversationPlan.Thread = { + id: event.threadId, + question: event.question, + questionSources: event.source ? [event.source] : [], + questionAuthoring: event.source && event.question === event.source.quote + ? "quoted" + : "scribe", + status: "exploring", + contributions: [], + stances: [], + stanceHistory: [], + decisionHistory: [], + candidates: [], + version: 1, + }; + return { + ...state, + revision: state.revision + 1, + events: [...state.events, event], + threads: [...state.threads, thread], + }; + } + if (!existing || existing.version !== event.observedThreadVersion) { + throw new Error("stale conversation thread version"); + } + let threads = [...state.threads]; + let threadIndex = threads.indexOf(existing); + threads[threadIndex] = structuredClone(existing); + // A move also changes its destination. Other threads and sidecars are read-only here. + if (event.type === "card.corrected" && event.change.kind === "move") { + let targetThreadId = event.change.targetThreadId; + let targetIndex = threads.findIndex((item) => item.id === targetThreadId); + if (targetIndex >= 0 && targetIndex !== threadIndex) { + threads[targetIndex] = structuredClone(threads[targetIndex]); + } + } + let next: State = { ...state, threads, events: [...state.events] }; + let thread = threads[threadIndex]; + switch (event.type) { + case "option.added": + case "reason.added": + case "constraint.added": + case "stance.changed": + case "option.relabeled": + applyContributionEvent(next, thread, event); + break; + case "thread.leaning": + case "decision.recorded": + case "decision.reopened": + case "candidate.proposed": + case "candidate.confirmed": + case "candidate.rejected": + case "card.linked": + case "thread.discarded": + applyLifecycleEvent(thread, event); + break; + case "settle.suggested": + case "settle.agreed": + case "settle.deferred": + case "settle.resumed": + case "scoped-choice.proposed": + case "scoped-choice.agreed": + case "scoped-choice.saved": + applySettlementEvent(next, thread, event); + break; + case "card.corrected": + applyCorrectionEvent(next, thread, event); + break; + } + thread.version++; + next.revision++; + next.events.push(event); + return next; +} diff --git a/apps/server/src/conversation-plan/preference.ts b/apps/server/src/conversation-plan/preference.ts new file mode 100644 index 00000000..76b966b3 --- /dev/null +++ b/apps/server/src/conversation-plan/preference.ts @@ -0,0 +1,94 @@ +import type { ConversationPlan } from "@chopin/protocol"; + +type Event = ConversationPlan.Event; +type Thread = ConversationPlan.Thread; + +/** A later suggestion cannot override an unmet verification gate. */ +export function activeSettleDeferral( + thread: Thread, + events: readonly Event[], +): Extract | undefined { + if (thread.status === "decided" || thread.status === "discarded") return; + let deferred = events.findLast((event): event is Extract => + event.type === "settle.deferred" && event.threadId === thread.id + ); + if (!deferred) return; + return events.some(event => + event.type === "settle.resumed" && event.threadId === thread.id + && event.proposalId === deferred.proposalId && event.deferredEventId === deferred.id + ) + ? undefined + : deferred; +} + +/** A historical proposal remains on the thread for replay, even after its source withdraws. */ +export function effectivePending( + thread: Thread, + events: readonly Event[], +): Thread["pendingSettle"] { + let pending = thread.pendingSettle; + if (!pending || thread.status === "decided" || thread.status === "discarded") return; + if (activeSettleDeferral(thread, events)) return; + let proposedAt = events.findLastIndex(event => + event.threadId === thread.id && event.type === "settle.suggested" + && event.optionId === pending.optionId + && event.source.messageId === pending.messageId + && event.source.author.kind === "member" + && event.source.author.handle === pending.proposer + ); + if (proposedAt < 0) return pending; + let resumedAt = events.findLastIndex(event => + event.type === "settle.resumed" && event.threadId === thread.id + ); + let withdrawn = events.slice(proposedAt + 1).some(event => + event.threadId === thread.id && event.type === "stance.changed" + && event.position !== "support" && event.optionId === pending.optionId + && event.source.author.kind === "member" + && event.source.author.handle === pending.proposer + && events.indexOf(event) > resumedAt + ); + return withdrawn ? undefined : pending; +} + +/** Keep only independently valid sources for the latest pending settle proposal. */ +export function effectivePreference( + thread: Thread, + events: readonly Event[], +): { optionId: string; messageIds: string[] } | undefined { + let pending = thread.pendingSettle; + if (!pending || thread.status === "decided" || thread.status === "discarded") return; + if (activeSettleDeferral(thread, events)) return; + let proposedAt = events.findLastIndex(event => + event.threadId === thread.id && event.type === "settle.suggested" + && event.optionId === pending.optionId + && event.source.messageId === pending.messageId + && event.source.author.kind === "member" + && event.source.author.handle === pending.proposer + ); + if (proposedAt < 0) return { optionId: pending.optionId, messageIds: [pending.messageId] }; + let resumedAt = events.findLastIndex(event => + event.type === "settle.resumed" && event.threadId === thread.id + ); + let sources = events.slice(proposedAt).filter(event => + event.threadId === thread.id + && (event.type === "settle.suggested" && event === events[proposedAt] + || event.type === "settle.agreed" && event.optionId === pending.optionId) + ); + let messageIds = sources.flatMap(source => { + if (source.type !== "settle.suggested" && source.type !== "settle.agreed") return []; + if (source.source.author.kind !== "member") return []; + let actor = source.source.author.handle; + let sourceAt = events.indexOf(source); + let withdrawn = events.slice(sourceAt + 1).some(event => + event.threadId === thread.id && event.type === "stance.changed" + && event.position !== "support" && event.optionId === pending.optionId + && event.source.author.kind === "member" && event.source.author.handle === actor + && events.indexOf(event) > resumedAt + ); + return withdrawn ? [] : [source.source.messageId]; + }); + messageIds = [...new Set(messageIds)]; + return messageIds.length + ? { optionId: pending.optionId, messageIds: messageIds.slice(-8) } + : undefined; +} diff --git a/apps/server/src/conversation-plan/research-snapshots.ts b/apps/server/src/conversation-plan/research-snapshots.ts new file mode 100644 index 00000000..ac13b94b --- /dev/null +++ b/apps/server/src/conversation-plan/research-snapshots.ts @@ -0,0 +1,91 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { isDeepStrictEqual } from "node:util"; +import { initialState } from "./domain-initial"; +import { applyEvent } from "./events"; +import { researchNamedOptionIds } from "./validation"; + +type State = ConversationPlan.State; +type Event = ConversationPlan.Event; + +export function taskOptionLabels(state: State, task: ConversationPlan.ResearchTask) { + let thread = state.threads.find(item => item.id === task.threadId); + return task.options.map(option => { + let contribution = thread?.contributions.find(item => + item.kind === "option" && item.id === option.id + ); + return contribution ? contribution.displayLabel ?? contribution.text : undefined; + }); +} + +export function taskSnapshot(state: State, task: ConversationPlan.ResearchTask) { + let labels = taskOptionLabels(state, task); + if (task.kind === "current-cost-comparison") return labels; + let thread = state.threads.find(item => item.id === task.threadId); + return { + labels, + currentOptionIds: thread?.contributions.filter(item => item.kind === "option") + .map(item => item.id), + }; +} + +export function taskMatches(state: State, task: ConversationPlan.ResearchTask): boolean { + let thread = state.threads.find(item => item.id === task.threadId); + return !!thread && thread.version === task.observedThreadVersion + && (task.kind === "current-cost-comparison" + || isDeepStrictEqual( + thread.contributions.filter(item => item.kind === "option").map(item => item.id), + task.options.map(item => item.id), + )) + && taskOptionLabels(state, task).every((label, index) => + label === task.options[index].labelAtOffer + ); +} + +export function namesStaleResearchOption( + state: State, + task: ConversationPlan.ResearchTask, + quote: string, +): boolean { + if (task.kind !== "current-cost-concern") return false; + let currentNames = new Set(researchNamedOptionIds(quote, task.options)); + let currentLabels = new Map(task.options.map(option => [option.id, option.labelAtOffer])); + let displayLabels = new Set(); + let oldName = (id: string, label: string) => + label !== currentLabels.get(id) && !currentNames.has(id) + && researchNamedOptionIds(quote, [{ id, labelAtOffer: label }]).length > 0; + for (let event of state.events.slice(0, task.observedEventCount)) { + if (event.type === "option.added" && currentLabels.has(event.contribution.id)) { + if (oldName(event.contribution.id, event.contribution.text)) return true; + } else if (event.type === "option.relabeled" && currentLabels.has(event.optionId)) { + displayLabels.add(event.optionId); + if (oldName(event.optionId, event.label)) return true; + } else if ( + event.type === "card.corrected" && event.change.kind === "edit" + && event.change.field === "contribution" + && currentLabels.has(event.change.contributionId) + && !displayLabels.has(event.change.contributionId) + ) { + if (oldName(event.change.contributionId, event.change.text)) return true; + } + } + return false; +} + +export function taskLabelChanges( + events: readonly Event[], + tasks: readonly ConversationPlan.ResearchOffer[], +) { + let changes = new Map(); + if (tasks.length === 0) return changes; + let state = initialState(); + for (let [eventIndex, event] of events.entries()) { + let before = tasks.map(offer => taskSnapshot(state, offer.task!)); + state = applyEvent(state, event); + for (let [index, offer] of tasks.entries()) { + if (!isDeepStrictEqual(before[index], taskSnapshot(state, offer.task!))) { + changes.set(offer.id, [...(changes.get(offer.id) ?? []), eventIndex]); + } + } + } + return changes; +} diff --git a/apps/server/src/conversation-plan/research-validation.test.ts b/apps/server/src/conversation-plan/research-validation.test.ts new file mode 100644 index 00000000..59e2c935 --- /dev/null +++ b/apps/server/src/conversation-plan/research-validation.test.ts @@ -0,0 +1,153 @@ +import { expect, test } from "bun:test"; +import type { ConversationPlan } from "@chopin/protocol"; +import { assertResearchOfferShape, renderResearchTask, researchNamedOptionIds } from "./validation"; + +function legacy(): ConversationPlan.ResearchOffer { + return { + id: "offer-1", + needId: "need-1", + contextId: "context-1", + source: { + messageId: "message-1", + author: { kind: "member", handle: "maggie" }, + quote: "Compare costs", + start: 0, + end: 13, + }, + brief: "Compare costs", + status: "offered", + }; +} + +test("legacy research offers keep the exact source quote as their brief", () => { + expect(() => assertResearchOfferShape(legacy())).not.toThrow(); + expect(() => assertResearchOfferShape({ ...legacy(), brief: "Compare prices" })) + .toThrow("research brief must equal source quote"); +}); + +function comparison(): ConversationPlan.ResearchTask { + return { + kind: "current-cost-comparison", + threadId: "thread-1", + observedEventCount: 2, + observedThreadVersion: 1, + options: [{ id: "r2", labelAtOffer: "Cloudflare R2" }, { id: "s3", labelAtOffer: "Amazon S3" }], + }; +} +function taskOffer( + task: ConversationPlan.ResearchTask, + quote = "Compare costs", +): ConversationPlan.ResearchOffer { + let source = { ...legacy().source, quote, end: quote.length }; + return { + ...legacy(), + source, + threadId: "thread-1", + task, + brief: renderResearchTask(task, source), + }; +} + +test("comparison and concern briefs use frozen labels and exact source context", () => { + let pair = taskOffer(comparison()); + expect(pair.brief).toBe( + "Compare current costs for Cloudflare R2 and Amazon S3. Discussion context: “Compare costs”", + ); + expect(() => assertResearchOfferShape(pair)).not.toThrow(); + let concern: ConversationPlan.ResearchTask = { + kind: "current-cost-concern", + threadId: "thread-1", + observedEventCount: 2, + observedThreadVersion: 1, + options: [ + { id: "r2", labelAtOffer: "Cloudflare R2" }, + { id: "s3", labelAtOffer: "Amazon S3" }, + { id: "b2", labelAtOffer: "Backblaze B2" }, + ], + }; + let across = taskOffer(concern); + expect(across.brief).toBe( + "Investigate current costs across Cloudflare R2, Amazon S3, Backblaze B2. Discussion context: “Compare costs”", + ); + expect(() => assertResearchOfferShape(across)).not.toThrow(); + let focused = taskOffer({ ...concern, focusOptionId: "r2" }, "R2 costs"); + expect(focused.brief).toBe( + "Investigate current costs for Cloudflare R2; consider Amazon S3, Backblaze B2 as context. Discussion context: “R2 costs”", + ); + expect(() => assertResearchOfferShape(focused)).not.toThrow(); + expect(() => assertResearchOfferShape({ ...focused, brief: focused.brief + " changed" })).toThrow( + /brief must match/, + ); + expect(() => assertResearchOfferShape(taskOffer({ ...concern, focusOptionId: "r2" }))).toThrow( + /not grounded/, + ); +}); + +test("research provider codes resolve only unique bounded mentions", () => { + let options = [{ id: "r2", labelAtOffer: "Cloudflare R2" }, { + id: "s3", + labelAtOffer: "Amazon S3", + }]; + expect(researchNamedOptionIds("R2 costs", options)).toEqual(["r2"]); + expect(researchNamedOptionIds("AR2B costs", options)).toEqual([]); + expect(researchNamedOptionIds("R2 and S3 costs", options)).toEqual(["r2", "s3"]); + expect(researchNamedOptionIds("R2", [...options, { id: "other", labelAtOffer: "Other R2" }])) + .toEqual([]); + expect(researchNamedOptionIds("R2", [{ id: "mixed", labelAtOffer: "R2 S3" }])).toEqual([]); +}); + +test("frozen research task fields, option identities, labels, and focus remain bounded", () => { + let offer = taskOffer(comparison()); + for ( + let task of [ + { ...comparison(), threadId: "other" }, + { ...comparison(), observedEventCount: -1 }, + { ...comparison(), focusOptionId: "r2" }, + { ...comparison(), options: [{ id: "r2", labelAtOffer: "R2" }] }, + { + ...comparison(), + options: [{ id: "r2", labelAtOffer: "R2" }, { id: "r2", labelAtOffer: "S3" }], + }, + { + ...comparison(), + options: [{ id: "r2", labelAtOffer: " R2" }, { id: "s3", labelAtOffer: "S3" }], + }, + { + ...comparison(), + options: [{ id: "r2", labelAtOffer: "x".repeat(161) }, { id: "s3", labelAtOffer: "S3" }], + }, + { ...comparison(), extra: true }, + ] + ) expect(() => assertResearchOfferShape({ ...offer, task })).toThrow(); +}); + +test("research status and action actors agree while unknown fields are rejected", () => { + let action = { + id: "action-1", + kind: "research", + actor: { kind: "member", handle: "maggie" }, + principalId: "principal-1", + at: 1, + }; + expect(() => assertResearchOfferShape({ ...legacy(), status: "accepted", action })).not.toThrow(); + expect(() => + assertResearchOfferShape({ + ...legacy(), + status: "dismissed", + action: { ...action, kind: "dismiss" }, + }) + ).not.toThrow(); + for ( + let offer of [ + { ...legacy(), action }, + { ...legacy(), status: "unknown" }, + { ...legacy(), status: "accepted" }, + { ...legacy(), status: "dismissed", action }, + { ...legacy(), status: "accepted", action: { ...action, actor: { kind: "agent" } } }, + { ...legacy(), status: "accepted", action: { ...action, at: -1 } }, + { ...legacy(), source: { ...legacy().source, author: { kind: "agent" } } }, + { ...legacy(), source: { ...legacy().source, role: "support" } }, + { ...legacy(), extra: true }, + ] + ) expect(() => assertResearchOfferShape(offer)).toThrow(); +}); diff --git a/apps/server/src/conversation-plan/research-validation.ts b/apps/server/src/conversation-plan/research-validation.ts new file mode 100644 index 00000000..1ec4fc31 --- /dev/null +++ b/apps/server/src/conversation-plan/research-validation.ts @@ -0,0 +1,174 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { assertSourceShape } from "./sources"; +import { actor, id, knownKeys, record, text, version } from "./validation-fields"; + +const MAX_RESEARCH_OPTION_LABEL = 160; + +/** Both the Chat card and worker receive this exact, code-owned wording. */ +export function renderResearchTask( + task: ConversationPlan.ResearchTask, + source: ConversationPlan.ResearchSource, +): string { + if (task.kind === "current-cost-comparison") { + return `Compare current costs for ${task.options[0].labelAtOffer} and ${ + task.options[1].labelAtOffer + }. Discussion context: “${source.quote}”`; + } + let focus = task.options.find(item => item.id === task.focusOptionId); + let labels = task.options.map(item => item.labelAtOffer); + return focus + ? `Investigate current costs for ${focus.labelAtOffer}; consider ${ + task.options.filter( + item => item.id !== focus.id, + ).map(item => item.labelAtOffer).join(", ") + } as context. Discussion context: “${source.quote}”` + : `Investigate current costs across ${ + labels.join(", ") + }. Discussion context: “${source.quote}”`; +} + +/** A short provider code is usable only when it uniquely resolves within this frozen task. */ +export function researchNamedOptionIds( + quote: string, + options: readonly { id: string; labelAtOffer: string }[], +): string[] { + let codes = options.map( + option => [ + ...new Set( + option.labelAtOffer.match(/(? { + let escaped = name.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); + return new RegExp(`(^|[^\\p{L}\\p{N}])${escaped}(?=$|[^\\p{L}\\p{N}])`, "iu") + .test(quote); + }; + return options.filter((option, index) => { + let label = option.labelAtOffer; + let optionCodes = codes[index]!; + if (optionCodes.length > 1) return false; + let code = optionCodes[0]; + return mentions(label) || !!code + && codes.filter(item => item.includes(code)).length === 1 + && mentions(code); + }).map(option => option.id); +} + +export function assertResearchOfferShape( + value: unknown, +): asserts value is ConversationPlan.ResearchOffer { + let offer = record(value); + knownKeys(offer, [ + "id", + "needId", + "contextId", + "source", + "brief", + "threadId", + "task", + "status", + "action", + ]); + id(offer.id); + id(offer.needId); + id(offer.contextId); + let source = record(offer.source); + knownKeys(source, ["messageId", "author", "quote", "start", "end"]); + assertSourceShape({ ...source, role: "support" }); + if ((source.author as { kind: string }).kind !== "member") { + throw new Error("research offer source requires a member message"); + } + actor(source.author); + text(offer.brief, 2048); + if (offer.threadId !== undefined) id(offer.threadId); + if (offer.task === undefined) { + if (offer.brief !== source.quote) throw new Error("research brief must equal source quote"); + } else { + let task = record(offer.task); + knownKeys(task, [ + "kind", + "threadId", + "observedEventCount", + "observedThreadVersion", + "options", + "focusOptionId", + ]); + if (task.kind !== "current-cost-comparison" && task.kind !== "current-cost-concern") { + throw new Error("invalid research task kind"); + } + id(task.threadId); + version(task.observedEventCount); + version(task.observedThreadVersion); + if (offer.threadId !== task.threadId) throw new Error("research task thread disagrees"); + let count = task.kind === "current-cost-comparison" ? 2 : 3; + if ( + !Array.isArray(task.options) + || task.options.length < count + || task.options.length > (task.kind === "current-cost-comparison" ? 2 : 4) + ) { + throw new Error( + task.kind === "current-cost-comparison" + ? "research task requires two options" + : "research concern requires three or four options", + ); + } + let optionIds = new Set(); + for (let value of task.options) { + let option = record(value); + knownKeys(option, ["id", "labelAtOffer"]); + id(option.id); + text(option.labelAtOffer, MAX_RESEARCH_OPTION_LABEL); + if ( + option.labelAtOffer !== (option.labelAtOffer as string).trim() + || /[\r\n]/.test(option.labelAtOffer as string) + ) { + throw new Error("invalid research option label"); + } + optionIds.add(option.id as string); + } + if (optionIds.size !== task.options.length) throw new Error("duplicate research task option"); + if (task.kind === "current-cost-comparison") { + if (task.focusOptionId !== undefined) throw new Error("invalid research task focus"); + } else { + if (task.focusOptionId !== undefined) { + id(task.focusOptionId); + if (!optionIds.has(task.focusOptionId as string)) { + throw new Error("invalid research task focus"); + } + } + let named = researchNamedOptionIds( + source.quote as string, + task.options as Array<{ id: string; labelAtOffer: string }>, + ); + if (named.length > 1 || task.focusOptionId !== named[0]) { + throw new Error("research task focus is not grounded in its source"); + } + } + if ( + offer.brief !== renderResearchTask( + task as ConversationPlan.ResearchTask, + source as ConversationPlan.ResearchSource, + ) + ) throw new Error("research brief must match task and source"); + } + if (!["offered", "dismissed", "accepted"].includes(offer.status as string)) { + throw new Error("invalid research offer status"); + } + if (offer.status === "offered") { + if (offer.action !== undefined) throw new Error("unacted research offer has an action"); + return; + } + let action = record(offer.action); + knownKeys(action, ["id", "kind", "actor", "principalId", "at"]); + id(action.id); + id(action.principalId); + if (action.kind !== (offer.status === "accepted" ? "research" : "dismiss")) { + throw new Error("research offer action disagrees with status"); + } + actor(action.actor); + if ((action.actor as { kind: string }).kind !== "member") { + throw new Error("research offer action requires a member"); + } + version(action.at); +} diff --git a/apps/server/src/conversation-plan/state-provenance.ts b/apps/server/src/conversation-plan/state-provenance.ts new file mode 100644 index 00000000..e4e5e742 --- /dev/null +++ b/apps/server/src/conversation-plan/state-provenance.ts @@ -0,0 +1,53 @@ +import type { Chat, ConversationPlan } from "@chopin/protocol"; +import { validateSource } from "./sources"; +import { taskLabelChanges } from "./research-snapshots"; + +type State = ConversationPlan.State; + +function eventSecond(at: number): number { + return Math.floor(at >= 100_000_000_000 ? at / 1_000 : at); +} + +export function validateStateSourcesWithChanges( + state: State, + messages: ReadonlyMap | readonly Chat.Entry[], + labelChanges: ReadonlyMap, +): void { + let lookup = Array.isArray(messages) + ? new Map(messages.map((message) => [message.id, message])) + : messages as ReadonlyMap; + for (let event of state.events) { + if ("source" in event && event.source) { + let message = lookup.get(event.source.messageId); + if (!message) throw new Error("conversation source message missing"); + validateSource(event.source, message); + } + if (event.type === "scoped-choice.saved") { + for (let source of event.sources) { + let message = lookup.get(source.messageId); + if (!message) throw new Error("conversation source message missing"); + validateSource(source, message); + } + } + } + for (let offer of state.researchOffers ?? []) { + let message = lookup.get(offer.source.messageId); + if (!message) throw new Error("research offer source message missing"); + validateSource({ ...offer.source, role: "support" }, message); + if ( + offer.task + && (labelChanges.get(offer.id) ?? []).some(index => + index >= offer.task!.observedEventCount + && eventSecond(state.events[index].at) <= message.ts + ) + ) throw new Error("research task source postdates its captured event prefix"); + } +} + +export function validateStateSources( + state: State, + messages: ReadonlyMap | readonly Chat.Entry[], +): void { + let taskOffers = (state.researchOffers ?? []).filter(offer => offer.task); + validateStateSourcesWithChanges(state, messages, taskLabelChanges(state.events, taskOffers)); +} diff --git a/apps/server/src/conversation-plan/state-restoration.ts b/apps/server/src/conversation-plan/state-restoration.ts new file mode 100644 index 00000000..dc0278f0 --- /dev/null +++ b/apps/server/src/conversation-plan/state-restoration.ts @@ -0,0 +1,72 @@ +import type { Chat, ConversationPlan } from "@chopin/protocol"; +import { isDeepStrictEqual } from "node:util"; +import { initialState } from "./domain-initial"; +import { applyEvent, currentScopedProposal } from "./events"; +import { assertStateShape } from "./validation"; +import { namesStaleResearchOption, taskMatches, taskSnapshot } from "./research-snapshots"; +import { validateStateSourcesWithChanges } from "./state-provenance"; + +type State = ConversationPlan.State; + +export function restoreState( + raw: unknown, + messages?: ReadonlyMap | readonly Chat.Entry[], +): State { + if (raw === undefined) return initialState(); + assertStateShape(raw); + let tasks = (raw.researchOffers ?? []).filter(offer => offer.task); + let captures = new Map(); + for (let offer of tasks) { + let count = offer.task!.observedEventCount; + captures.set(count, [...(captures.get(count) ?? []), offer]); + } + let matched = new Set(); + let rebuilt = initialState(); + let changes = new Map(); + let checkCapture = (count: number) => { + for (let offer of captures.get(count) ?? []) { + if (offer.task && taskMatches(rebuilt, offer.task)) matched.add(offer.id); + } + }; + checkCapture(0); + for (let [eventIndex, event] of raw.events.entries()) { + let before = tasks.map(offer => taskSnapshot(rebuilt, offer.task!)); + rebuilt = applyEvent(rebuilt, event); + for (let [index, offer] of tasks.entries()) { + if (!isDeepStrictEqual(before[index], taskSnapshot(rebuilt, offer.task!))) { + changes.set(offer.id, [...(changes.get(offer.id) ?? []), eventIndex]); + } + } + checkCapture(rebuilt.events.length); + } + let restored = structuredClone(raw); + for (let thread of restored.threads) { + let pending = thread.pendingScopedChoice; + if (!pending || pending.proposalId !== undefined) continue; + let proposal = currentScopedProposal(thread, restored.events); + if (!proposal) throw new Error("scoped choice proposal is ambiguous or missing"); + pending.proposalId = proposal.id; + } + if ( + !isDeepStrictEqual( + JSON.parse(JSON.stringify(rebuilt.threads)), + JSON.parse(JSON.stringify(restored.threads)), + ) + || rebuilt.events.length !== raw.events.length || raw.revision < rebuilt.revision + ) { + throw new Error("conversation plan snapshot does not match its events"); + } + if (matched.size !== tasks.length) { + throw new Error("research task capture does not match history"); + } + if ( + tasks.some(offer => + offer.task + && namesStaleResearchOption(raw, offer.task, offer.source.quote) + ) + ) { + throw new Error("research task source names a stale option"); + } + if (messages) validateStateSourcesWithChanges(restored, messages, changes); + return restored; +} diff --git a/apps/server/src/conversation-plan/state-validation.test.ts b/apps/server/src/conversation-plan/state-validation.test.ts new file mode 100644 index 00000000..e49169cc --- /dev/null +++ b/apps/server/src/conversation-plan/state-validation.test.ts @@ -0,0 +1,237 @@ +import { expect, test } from "bun:test"; +import type { ConversationPlan } from "@chopin/protocol"; +import { + assertStateShape, + MAX_ANALYSIS, + MAX_EVENTS, + MAX_QUEUE, + MAX_RESEARCH_OFFERS, + MAX_THREADS, +} from "./validation"; + +function state(): ConversationPlan.State { + return { schemaVersion: 1, revision: 0, events: [], threads: [], queue: [], analysis: [] }; +} + +test("the empty version-one snapshot validates without adding stored fields", () => { + let snapshot = state(); + let before = structuredClone(snapshot); + expect(() => assertStateShape(snapshot)).not.toThrow(); + expect(snapshot).toEqual(before); + expect(() => assertStateShape({ ...snapshot, schemaVersion: 2 })).toThrow( + "unsupported conversation plan schema", + ); +}); + +function analysis(): ConversationPlan.AnalysisRecord { + return { + messageId: "message-1", + questionSetVersion: "questions-1", + modelVersion: "fixture", + status: "applied", + passes: [], + eventIds: [], + }; +} +function withAnswer(answer: unknown) { + return { + ...state(), + analysis: [{ ...analysis(), passes: [{ stage: "triage", answers: { flag: answer } }] }], + }; +} + +test("snapshot array bounds and revision reject malformed saved containers", () => { + for ( + let [field, maximum] of [ + ["events", MAX_EVENTS], + ["threads", MAX_THREADS], + ["queue", MAX_QUEUE], + ["analysis", MAX_ANALYSIS], + ] as const + ) { + expect(() => + assertStateShape({ ...state(), [field]: Array.from({ length: maximum + 1 }, () => null) }) + ) + .toThrow("invalid conversation plan snapshot arrays"); + expect(() => assertStateShape({ ...state(), [field]: {} })).toThrow( + "invalid conversation plan snapshot arrays", + ); + } + expect(() => + assertStateShape({ + ...state(), + researchOffers: Array.from({ length: MAX_RESEARCH_OFFERS + 1 }, () => null), + }) + ) + .toThrow("invalid research offers array"); + expect(() => assertStateShape({ ...state(), revision: -1 })).toThrow(/version/); + expect(() => assertStateShape({ ...state(), extra: true })).toThrow(/unknown/); +}); + +test("queued messages retain unique IDs and bounded attempts and errors", () => { + let queued = { messageId: "message-1", status: "pending", attempts: 0 }; + expect(() => assertStateShape({ ...state(), queue: [queued] })).not.toThrow(); + expect(() => assertStateShape({ ...state(), queue: [queued, queued] })).toThrow( + "duplicate queued message", + ); + for ( + let fields of [{ status: "unknown" }, { attempts: -1 }, { error: "x".repeat(301) }, { + extra: true, + }] + ) { + expect(() => assertStateShape({ ...state(), queue: [{ ...queued, ...fields }] })).toThrow(); + } +}); + +test("raw analysis probabilities, winning choices, and weighted scores validate", () => { + let choice = { + type: "choice", + choice: "save", + confidence: 0.8, + probabilities: { save: 0.8, skip: 0.2 }, + }; + let score = { + type: "score", + score: 0.75, + confidence: 0.9, + legend: { "0": "Low", "1": "High" }, + probabilities: { "0": 0.25, "1": 0.75 }, + }; + for (let answer of [{ type: "noul", noul: 0.5 }, choice, score]) { + expect(() => assertStateShape(withAnswer(answer))).not.toThrow(); + } + for ( + let answer of [ + { type: "noul", noul: Number.NaN }, + { type: "noul", noul: 1.1 }, + { ...choice, choice: "skip" }, + { ...choice, probabilities: { save: 0.1, skip: 0.2 } }, + { ...choice, probabilities: { save: 1 } }, + { ...choice, confidence: -1 }, + { ...score, score: 0.1 }, + { ...score, legend: { "1": "Low", "2": "High" } }, + { ...score, probabilities: { "0": 0.25, "2": 0.75 } }, + { type: "noul", noul: 0.5, extra: true }, + ] + ) expect(() => assertStateShape(withAnswer(answer))).toThrow(); +}); + +test("analysis pass order, question counts, and aggregate distributions stay bounded", () => { + let passes = [ + { stage: "triage", answers: {} }, + { stage: "targeting", answers: {} }, + { stage: "clarification", version: "bare-editor-clarification-1", answers: {} }, + ]; + expect(() => assertStateShape({ ...state(), analysis: [{ ...analysis(), passes }] })).not + .toThrow(); + expect(() => assertStateShape({ ...state(), analysis: [{ ...analysis(), passes: [passes[2]] }] })) + .toThrow(/clarification/); + let answers = Object.fromEntries( + Array.from({ length: 46 }, (_, index) => [`q${index}`, { type: "noul", noul: 0.5 }]), + ); + expect(() => + assertStateShape({ + ...state(), + analysis: [{ ...analysis(), passes: [{ stage: "triage", answers }] }], + }) + ).toThrow("too many analysis questions"); + let probabilities = Object.fromEntries( + Array.from({ length: 200 }, (_, index) => [`option${index}`, 0.005]), + ); + let manyAnswers = Object.fromEntries( + Array.from( + { length: 3 }, + ( + _, + index, + ) => [`q${index}`, { type: "choice", choice: "option0", confidence: 0.005, probabilities }], + ), + ); + expect(() => + assertStateShape({ + ...state(), + analysis: [{ ...analysis(), passes: [{ stage: "triage", answers: manyAnswers }] }], + }) + ).toThrow("too many analysis answer entries"); +}); + +test("candidate outcome event IDs must belong to the analysis record", () => { + let record = { + ...analysis(), + eventIds: ["event-1"], + outcomes: [{ start: 0, end: 3, status: "accepted", gate: "fixture", eventIds: ["event-1"] }], + }; + expect(() => assertStateShape({ ...state(), analysis: [record] })).not.toThrow(); + expect(() => assertStateShape({ ...state(), analysis: [{ ...record, eventIds: [] }] })) + .toThrow("candidate outcome event is not accepted"); + expect(() => + assertStateShape({ + ...state(), + analysis: [{ ...record, outcomes: [{ ...record.outcomes[0], end: 0 }] }], + }) + ) + .toThrow("invalid candidate outcome"); + expect(() => + assertStateShape({ + ...state(), + analysis: [{ ...analysis(), quoteValidation: [{ start: 0, end: 3, valid: "yes" }] }], + }) + ) + .toThrow("invalid quote validation debug"); +}); + +function offer(index: number): ConversationPlan.ResearchOffer { + return { + id: `offer-${index}`, + needId: `need-${index}`, + contextId: "context-1", + source: { + messageId: `message-${index}`, + author: { kind: "member", handle: "maggie" }, + quote: "Costs", + start: 0, + end: 5, + }, + brief: "Costs", + status: "offered", + }; +} + +test("snapshot research identity guards duplicate offers, sources, needs, and actions", () => { + let first = offer(1); + let second = offer(2); + expect(() => assertStateShape({ ...state(), researchOffers: [first, second] })).not.toThrow(); + for ( + let changed of [ + { ...second, id: first.id }, + { ...second, source: first.source }, + { ...second, needId: first.needId }, + ] + ) { + expect(() => assertStateShape({ ...state(), researchOffers: [first, changed] })).toThrow( + "duplicate research offer identity", + ); + } + let action: ConversationPlan.ResearchAction = { + id: "action-1", + kind: "research", + actor: { kind: "member", handle: "maggie" }, + principalId: "principal-1", + at: 1, + }; + expect(() => + assertStateShape({ + ...state(), + researchOffers: [{ ...first, status: "accepted", action }, { + ...second, + status: "accepted", + action, + }], + }) + ) + .toThrow("duplicate research offer identity"); + expect(() => + assertStateShape({ ...state(), researchOffers: [{ ...first, threadId: "missing" }] }) + ) + .toThrow("research offer thread is missing"); +}); diff --git a/apps/server/src/conversation-plan/state-validation.ts b/apps/server/src/conversation-plan/state-validation.ts new file mode 100644 index 00000000..61619286 --- /dev/null +++ b/apps/server/src/conversation-plan/state-validation.ts @@ -0,0 +1,265 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { assertEventShape } from "./event-validation"; +import { assertResearchOfferShape } from "./research-validation"; +import { + id, + knownKeys, + MAX_ANALYSIS, + MAX_EVENTS, + MAX_QUEUE, + MAX_RESEARCH_OFFERS, + MAX_THREADS, + record, + version, +} from "./validation-fields"; + +function unitProbability(value: unknown): void { + if (typeof value !== "number" || !Number.isFinite(value) || value < 0 || value > 1) { + throw new Error("invalid analysis probability"); + } +} + +function analysisDistribution(value: unknown, max: number): Record { + let distribution = record(value); + let keys = Object.keys(distribution); + if (keys.length < 2 || keys.length > max) throw new Error("invalid analysis distribution"); + for (let key of keys) { + id(key); + unitProbability(distribution[key]); + } + let total = Object.values(distribution).reduce((sum: number, item) => sum + (item as number), 0); + if (Math.abs(total - 1) > 0.08) throw new Error("invalid analysis distribution"); + return distribution as Record; +} + +export function assertStateShape(value: unknown): asserts value is ConversationPlan.State { + let state = record(value); + knownKeys(state, [ + "schemaVersion", + "revision", + "events", + "threads", + "queue", + "analysis", + "researchOffers", + ]); + if (state.schemaVersion !== 1) throw new Error("unsupported conversation plan schema"); + version(state.revision); + if ( + !Array.isArray(state.events) || state.events.length > MAX_EVENTS + || !Array.isArray(state.threads) || state.threads.length > MAX_THREADS + || !Array.isArray(state.queue) || state.queue.length > MAX_QUEUE + || !Array.isArray(state.analysis) || state.analysis.length > MAX_ANALYSIS + ) throw new Error("invalid conversation plan snapshot arrays"); + if (state.researchOffers !== undefined) { + if (!Array.isArray(state.researchOffers) || state.researchOffers.length > MAX_RESEARCH_OFFERS) { + throw new Error("invalid research offers array"); + } + let offerIds = new Set(); + let sources = new Set(); + let needs = new Set(); + let actionIds = new Set(); + let threadIds = new Set(state.threads.map((item: ConversationPlan.Thread) => item.id)); + for (let offer of state.researchOffers) { + assertResearchOfferShape(offer); + let need = JSON.stringify([offer.needId, offer.contextId]); + if ( + offerIds.has(offer.id) || sources.has(offer.source.messageId) || needs.has(need) + || offer.action && actionIds.has(offer.action.id) + ) throw new Error("duplicate research offer identity"); + if (offer.threadId !== undefined && !threadIds.has(offer.threadId)) { + throw new Error("research offer thread is missing"); + } + offerIds.add(offer.id); + sources.add(offer.source.messageId); + needs.add(need); + if (offer.action) actionIds.add(offer.action.id); + } + } + for (let event of state.events) assertEventShape(event); + for (let item of state.queue) { + let queued = record(item); + knownKeys(queued, ["messageId", "status", "attempts", "error"]); + id(queued.messageId); + if (!["pending", "processing", "failed"].includes(queued.status as string)) { + throw new Error("invalid queue status"); + } + version(queued.attempts); + if ( + queued.error !== undefined && (typeof queued.error !== "string" || queued.error.length > 300) + ) throw new Error("invalid queue error"); + } + if (new Set(state.queue.map((item) => item.messageId)).size !== state.queue.length) { + throw new Error("duplicate queued message"); + } + for (let item of state.analysis) { + let analysis = record(item); + knownKeys(analysis, [ + "messageId", + "questionSetVersion", + "modelVersion", + "status", + "passes", + "selectedTarget", + "quoteValidation", + "outcomes", + "policyGate", + "candidates", + "eventIds", + "latencyMs", + "error", + ]); + id(analysis.messageId); + id(analysis.questionSetVersion); + id(analysis.modelVersion); + if ( + !["queued", "running", "applied", "unlinked", "failed"].includes(analysis.status as string) + ) throw new Error("invalid analysis status"); + if (!Array.isArray(analysis.passes) || analysis.passes.length > 3) { + throw new Error("invalid analysis passes"); + } + for (let [index, pass] of analysis.passes.entries()) { + let p = record(pass); + knownKeys(p, ["stage", "version", "answers"]); + if (!["triage", "targeting", "clarification"].includes(p.stage as string)) { + throw new Error("invalid analysis pass stage"); + } + if (p.stage === "clarification") { + if ( + index !== 2 || analysis.passes[0]?.stage !== "triage" + || analysis.passes[1]?.stage !== "targeting" + || p.version !== "bare-editor-clarification-1" + ) throw new Error("invalid clarification pass version or order"); + } else if (p.version !== undefined) { + throw new Error("invalid analysis pass version"); + } + let answers = record(p.answers); + if (Object.keys(answers).length > 45) throw new Error("too many analysis questions"); + let answerEntries = 0; + for (let [question, value] of Object.entries(answers)) { + id(question); + let answer = record(value); + if (answer.type === "noul") { + knownKeys(answer, ["type", "noul"]); + unitProbability(answer.noul); + continue; + } + if (answer.type === "choice") { + knownKeys(answer, ["type", "choice", "confidence", "probabilities"]); + id(answer.choice); + unitProbability(answer.confidence); + let distribution = analysisDistribution(answer.probabilities, 255); + answerEntries += Object.keys(distribution).length; + if ( + !(answer.choice as string in distribution) + || distribution[answer.choice as string] + 0.01 + < Math.max(...Object.values(distribution)) + ) { + throw new Error("invalid analysis choice"); + } + continue; + } + if (answer.type === "score") { + knownKeys(answer, ["type", "score", "confidence", "legend", "probabilities"]); + unitProbability(answer.confidence); + let legend = record(answer.legend); + let levels = Object.keys(legend); + if ( + levels.length < 2 || levels.length > 10 + || levels.some((level, index) => + level !== String(index) || typeof legend[level] !== "string" + || !(legend[level] as string).trim() || (legend[level] as string).length > 500 + ) + ) throw new Error("invalid analysis score legend"); + let distribution = analysisDistribution(answer.probabilities, 10); + answerEntries += Object.keys(distribution).length; + if ( + Object.keys(distribution).length !== levels.length + || levels.some((level) => !(level in distribution)) + ) { + throw new Error("invalid analysis score distribution"); + } + let weighted = levels.reduce((sum, level, index) => sum + index * distribution[level], 0); + if ( + typeof answer.score !== "number" || !Number.isFinite(answer.score) + || answer.score < 0 || answer.score > levels.length - 1 + || Math.abs(answer.score - weighted) > 0.12 + ) { + throw new Error("invalid analysis score"); + } + continue; + } + throw new Error("invalid analysis answer type"); + } + if (answerEntries > 512) throw new Error("too many analysis answer entries"); + } + if (!Array.isArray(analysis.eventIds) || analysis.eventIds.length > 12) { + throw new Error("invalid analysis event IDs"); + } + for (let eventId of analysis.eventIds) id(eventId); + if (analysis.outcomes !== undefined) { + if (!Array.isArray(analysis.outcomes) || analysis.outcomes.length > 4) { + throw new Error("invalid candidate outcomes"); + } + for (let outcome of analysis.outcomes) { + let item = record(outcome); + knownKeys(item, ["start", "end", "status", "gate", "targetId", "eventIds"]); + version(item.start); + version(item.end); + if ( + (item.end as number) <= (item.start as number) + || !["accepted", "review", "ignored"].includes(item.status as string) + || typeof item.gate !== "string" || !item.gate || item.gate.length > 300 + || !Array.isArray(item.eventIds) || item.eventIds.length > 12 + ) { + throw new Error("invalid candidate outcome"); + } + if (item.targetId !== undefined) id(item.targetId); + for (let eventId of item.eventIds) { + id(eventId); + if (!analysis.eventIds.includes(eventId)) { + throw new Error("candidate outcome event is not accepted"); + } + } + } + } + if ( + analysis.error !== undefined + && (typeof analysis.error !== "string" || analysis.error.length > 300) + ) throw new Error("invalid analysis error"); + if (analysis.latencyMs !== undefined) version(analysis.latencyMs); + if (analysis.selectedTarget !== undefined) id(analysis.selectedTarget); + if ( + analysis.policyGate !== undefined + && (typeof analysis.policyGate !== "string" || analysis.policyGate.length > 300) + ) throw new Error("invalid analysis policy gate"); + if (analysis.candidates !== undefined) { + if (!Array.isArray(analysis.candidates) || analysis.candidates.length > 12) { + throw new Error("invalid analysis candidates"); + } + for (let candidate of analysis.candidates) { + let c = record(candidate); + knownKeys(c, ["id", "kind", "targetId"]); + id(c.id); + if (!["resolution", "reopening"].includes(c.kind as string)) { + throw new Error("invalid analysis candidate kind"); + } + if (c.targetId !== undefined) id(c.targetId); + } + } + if (analysis.quoteValidation !== undefined) { + if (!Array.isArray(analysis.quoteValidation) || analysis.quoteValidation.length > 12) { + throw new Error("invalid quote validation debug"); + } + for (let quote of analysis.quoteValidation) { + let q = record(quote); + knownKeys(q, ["start", "end", "valid"]); + version(q.start); + version(q.end); + if ((q.end as number) <= (q.start as number) || typeof q.valid !== "boolean") { + throw new Error("invalid quote validation debug"); + } + } + } + } +} diff --git a/apps/server/src/conversation-plan/validation-fields.ts b/apps/server/src/conversation-plan/validation-fields.ts new file mode 100644 index 00000000..a5422a59 --- /dev/null +++ b/apps/server/src/conversation-plan/validation-fields.ts @@ -0,0 +1,54 @@ +export const MAX_CONTRIBUTIONS = 64; +export const MAX_THREADS = 20; +export const MAX_EVENTS = 4096; +export const MAX_ANALYSIS = 64; +export const MAX_QUEUE = 128; +export const MAX_RESEARCH_OFFERS = 64; + +export function record(value: unknown): Record { + if (!value || typeof value !== "object" || Array.isArray(value)) { + throw new Error("invalid conversation plan record"); + } + return value as Record; +} + +export function knownKeys(value: Record, allowed: readonly string[]): void { + if (Object.keys(value).some((key) => !allowed.includes(key))) { + throw new Error("unknown conversation plan field"); + } +} + +export function id(value: unknown): void { + if (typeof value !== "string" || !value || value.length > 200) { + throw new Error("invalid conversation plan ID"); + } +} + +export function text(value: unknown, max = 500): void { + if (typeof value !== "string" || !value.trim() || value.length > max) { + throw new Error("invalid conversation plan text"); + } +} + +/** Only an owned, direct preference with explicit provisional scope can propose a spike choice. */ +export function spikePreferenceLabel(quote: string): string | undefined { + return /^I(?:['’]d| would) pick (.+?) for (?:the|a) spike[;.!]?$/iu.exec(quote.trim())?.[1]; +} + +export function spikeAgreementLabel(quote: string): string | undefined { + return /^yep,\s+(.+?) for (?:the|a) spike[.!]$/iu.exec(quote.trim())?.[1]; +} + +export function version(value: unknown): void { + if (!Number.isSafeInteger(value) || (value as number) < 0) { + throw new Error("invalid conversation plan version"); + } +} + +export function actor(value: unknown): void { + let a = record(value); + knownKeys(a, a.kind === "member" ? ["kind", "handle"] : ["kind"]); + if (a.kind === "classifier" || a.kind === "agent") return; + if (a.kind === "member") return id(a.handle); + throw new Error("invalid conversation plan actor"); +} diff --git a/apps/server/src/conversation-plan/validation.test.ts b/apps/server/src/conversation-plan/validation.test.ts new file mode 100644 index 00000000..b9a4a4da --- /dev/null +++ b/apps/server/src/conversation-plan/validation.test.ts @@ -0,0 +1,342 @@ +import { expect, test } from "bun:test"; +import type { ConversationPlan } from "@chopin/protocol"; +import { + assertCorrectionAction, + assertEventShape, + MAX_ANALYSIS, + MAX_EVENTS, + MAX_QUEUE, + MAX_RESEARCH_OFFERS, + MAX_THREADS, + spikeAgreementLabel, + spikePreferenceLabel, +} from "./validation"; + +function base(): ConversationPlan.EventBase { + return { + id: "event-1", + threadId: "thread-1", + observedThreadVersion: 0, + origin: "classifier", + actor: { kind: "classifier" }, + at: 1, + }; +} +function human(): ConversationPlan.EventBase { + return { ...base(), origin: "human", actor: { kind: "member", handle: "maggie" } }; +} +function source( + role: ConversationPlan.SourceRole = "question", + quote = "Which?", +): ConversationPlan.SourceRef { + return { + messageId: "message-1", + author: { kind: "member", handle: "maggie" }, + quote, + start: 0, + end: quote.length, + role, + }; +} + +test("explicit human decisions validate without claiming a chat quote", () => { + expect(() => + assertEventShape({ ...human(), type: "decision.recorded", text: "Use Bun", explicit: true }) + ) + .not.toThrow(); + expect(() => + assertEventShape({ ...human(), type: "decision.recorded", text: "Use Bun", explicit: false }) + ) + .toThrow("decision requires explicit resolution"); + expect(() => + assertEventShape({ + ...human(), + type: "decision.recorded", + text: "Use Bun", + explicit: true, + source: source(), + }) + ) + .toThrow("human action cannot claim a message quote"); +}); + +let optionId = "01ARZ3NDEKTSV4RRFFQ69G5FAV"; +function planner(): ConversationPlan.EventBase { + return { ...base(), origin: "planner", actor: { kind: "agent" } }; +} + +test("all archived event kinds retain their field-validation path", () => { + let events = [ + { ...base(), type: "thread.opened", source: source(), question: "Which?" }, + { ...planner(), type: "thread.opened", question: "Which?" }, + ...["option.added", "reason.added", "constraint.added"].map(type => ({ + ...base(), + type, + source: source("option"), + contribution: { id: "contribution-1", text: "Bun", authoring: "quoted" }, + })), + { + ...human(), + type: "option.added", + contribution: { id: optionId, text: "Bun", authoring: "human-edited" }, + }, + { + ...human(), + type: "reason.added", + source: source("reason"), + contribution: { id: optionId, text: "Which?", authoring: "quoted" }, + }, + { + ...base(), + type: "stance.changed", + source: source("support"), + position: "support", + scopedProposalId: null, + }, + { ...planner(), type: "option.relabeled", optionId, label: "Bun", observedCardRevision: 0 }, + { ...base(), type: "thread.leaning", source: source("support") }, + { ...base(), type: "card.linked", questionnaireId: "card-1" }, + ...["settle.suggested", "settle.agreed"].map(type => ({ + ...base(), + type, + source: source("support"), + optionId, + })), + { ...base(), type: "settle.deferred", source: source("constraint"), proposalId: "proposal-1" }, + { + ...base(), + type: "settle.resumed", + source: source("verification"), + proposalId: "proposal-1", + deferredEventId: "deferred-1", + }, + { + ...base(), + type: "scoped-choice.proposed", + source: source("support", "I'd pick Bun for the spike."), + cardId: "card-1", + optionId, + label: "Bun", + scope: "spike", + }, + { + ...base(), + type: "scoped-choice.agreed", + source: source("support", "yep, Bun for the spike."), + proposalId: "proposal-1", + cardId: "card-1", + optionId, + label: "Bun", + scope: "spike", + }, + { + ...human(), + type: "scoped-choice.saved", + proposalId: "proposal-1", + cardId: "card-1", + optionId, + label: "Bun", + scope: "spike", + sources: [source("support")], + expectedGeneration: 0, + }, + { ...human(), type: "thread.discarded" }, + { + ...base(), + type: "decision.recorded", + source: source("resolution"), + text: "Use Bun", + explicit: true, + }, + { ...human(), type: "decision.reopened", explicit: true }, + { ...base(), type: "decision.reopened", source: source("reopening"), explicit: true }, + { + ...base(), + type: "candidate.proposed", + source: source("resolution"), + candidate: { id: "candidate-1", kind: "resolution", text: "Use Bun" }, + }, + ...["candidate.confirmed", "candidate.rejected"].map(type => ({ + ...human(), + type, + candidateId: "candidate-1", + })), + { + ...human(), + type: "card.corrected", + change: { kind: "edit", field: "question", text: "Which runtime?" }, + }, + ]; + for (let event of events) expect(() => assertEventShape(event)).not.toThrow(); +}); + +test("event IDs, versions, unknown fields, and actor origin must validate", () => { + let event = { ...base(), type: "thread.opened", source: source(), question: "Which?" }; + for ( + let override of [ + { id: "" }, + { id: "x".repeat(201) }, + { threadId: "" }, + { observedThreadVersion: -1 }, + { observedThreadVersion: 0.5 }, + { at: Number.MAX_SAFE_INTEGER + 1 }, + { actor: { kind: "member", handle: "maggie" } }, + { origin: "other" }, + { actor: { kind: "classifier", extra: true } }, + { extra: true }, + { source: { ...source(), quote: "Wrong" } }, + { question: "x".repeat(501) }, + { type: "unknown" }, + ] + ) expect(() => assertEventShape({ ...event, ...override })).toThrow(); + expect(() => assertEventShape({ ...event, id: "x".repeat(200), question: "x".repeat(500) })).not + .toThrow(); + expect(() => + assertEventShape({ ...planner(), type: "thread.opened", question: "x".repeat(1000) }) + ).not.toThrow(); + expect(() => + assertEventShape({ ...planner(), type: "thread.opened", question: "x".repeat(1001) }) + ).toThrow(); +}); + +test("human event restrictions and explicit decisions remain distinct from correction actions", () => { + for ( + let event of [ + { ...base(), type: "decision.recorded", text: "Use Bun", explicit: true }, + { ...human(), type: "decision.reopened", explicit: false }, + { ...base(), type: "thread.discarded" }, + { ...base(), type: "candidate.confirmed", candidateId: "candidate-1" }, + { + ...base(), + type: "card.corrected", + change: { kind: "edit", field: "question", text: "Which?" }, + }, + { ...human(), type: "card.corrected", change: { kind: "record-decision", text: "Use Bun" } }, + { + ...human(), + type: "option.added", + contribution: { id: "not-a-ulid", text: "Bun", authoring: "human-edited" }, + }, + ] + ) expect(() => assertEventShape(event)).toThrow(); + expect(() => + assertCorrectionAction({ + actionId: "action-1", + threadId: "thread-1", + expectedVersion: 0, + change: { kind: "record-decision", text: "Use Bun" }, + }) + ).not.toThrow(); +}); + +test("scoped saves validate bounded unique support IDs and matching member sources", () => { + let event = { + ...human(), + type: "scoped-choice.saved", + proposalId: "proposal-1", + cardId: "card-1", + optionId, + label: "Bun", + scope: "spike", + expectedGeneration: 0, + }; + expect(() => + assertEventShape({ ...event, supportEventIds: ["support-1"], sources: [source("support")] }) + ).not.toThrow(); + for ( + let fields of [ + { supportEventIds: [], sources: [] }, + { + supportEventIds: ["support-1", "support-1"], + sources: [source("support"), source("support")], + }, + { + supportEventIds: Array.from({ length: MAX_EVENTS + 1 }, (_, index) => `support-${index}`), + sources: [], + }, + { supportEventIds: ["support-1"], sources: [] }, + { supportEventIds: ["support-1"], agreementId: "agreement-1", sources: [source("support")] }, + { + supportEventIds: ["support-1", "support-2"], + sources: [source("support"), source("support")], + }, + { sources: [source("question")] }, + ] + ) expect(() => assertEventShape({ ...event, ...fields })).toThrow(); +}); + +test("correction action variants validate scoped IDs and fields without applying changes", () => { + let changes = [ + { kind: "add-excerpt", messageId: "message-1", start: 0, end: 500, contributionKind: "option" }, + { kind: "edit", field: "question", text: "Which?" }, + { kind: "edit", field: "decision", text: "Use Bun" }, + { kind: "edit", field: "contribution", contributionId: "contribution-1", text: "Bun" }, + { + kind: "move", + contributionId: "contribution-1", + targetThreadId: "thread-2", + targetVersion: 0, + }, + { kind: "set-status", status: "reopened" }, + { kind: "retarget-stance", stanceId: "stance-1", optionId }, + { kind: "dismiss-stance", stanceId: "stance-1" }, + { kind: "retarget-contribution", contributionId: "contribution-1", targetId: optionId }, + { kind: "record-decision", text: "Use Bun", optionId }, + { kind: "confirm-candidate", candidateId: "candidate-1" }, + { kind: "reject-candidate", candidateId: "candidate-1" }, + ]; + for (let change of changes) { + expect(() => + assertCorrectionAction({ + actionId: "action-1", + threadId: "thread-1", + expectedVersion: 0, + change, + }) + ).not.toThrow(); + } + let action = { + actionId: "action-1", + threadId: "thread-1", + expectedVersion: 0, + change: { kind: "edit", field: "question", text: "Which?" }, + }; + for ( + let fields of [ + { actionId: "" }, + { threadId: "x".repeat(201) }, + { expectedVersion: -1 }, + { extra: true }, + { change: { ...action.change, extra: true } }, + { change: { ...action.change, contributionId: "unexpected" } }, + { change: { kind: "edit", field: "contribution", text: "Bun" } }, + { change: { kind: "set-status", status: "decided" } }, + { + change: { + kind: "add-excerpt", + messageId: "message-1", + start: 0, + end: 501, + contributionKind: "option", + }, + }, + { + change: { + kind: "move", + contributionId: "contribution-1", + targetThreadId: "thread-2", + targetVersion: 0.5, + }, + }, + ] + ) expect(() => assertCorrectionAction({ ...action, ...fields })).toThrow(); +}); + +test("spike wording and exported snapshot bounds retain the archived contract", () => { + expect(spikePreferenceLabel("I'd pick Bun for the spike.")).toBe("Bun"); + expect(spikeAgreementLabel("yep, Bun for a spike!")).toBe("Bun"); + expect(spikePreferenceLabel("We should use Bun.")).toBeUndefined(); + expect(spikeAgreementLabel("yep, Bun for production.")).toBeUndefined(); + expect([MAX_THREADS, MAX_EVENTS, MAX_ANALYSIS, MAX_QUEUE, MAX_RESEARCH_OFFERS]) + .toEqual([20, 4096, 64, 128, 64]); +}); diff --git a/apps/server/src/conversation-plan/validation.ts b/apps/server/src/conversation-plan/validation.ts new file mode 100644 index 00000000..4499b68c --- /dev/null +++ b/apps/server/src/conversation-plan/validation.ts @@ -0,0 +1,17 @@ +export { assertCorrectionAction } from "./correction-validation"; +export { assertEventShape } from "./event-validation"; +export { + assertResearchOfferShape, + renderResearchTask, + researchNamedOptionIds, +} from "./research-validation"; +export { assertStateShape } from "./state-validation"; +export { + MAX_ANALYSIS, + MAX_EVENTS, + MAX_QUEUE, + MAX_RESEARCH_OFFERS, + MAX_THREADS, + spikeAgreementLabel, + spikePreferenceLabel, +} from "./validation-fields"; From 095d6021aa79179e8a1caaf99d788aadbb2e5e0e Mon Sep 17 00:00:00 2001 From: Maggie Appleton <5599295+MaggieAppleton@users.noreply.github.com> Date: Sat, 3 Oct 2026 11:46:19 +0100 Subject: [PATCH 2/2] Clear pending choices when a resolution closes a thread --- .../src/conversation-plan/event-lifecycle.ts | 3 +++ .../src/conversation-plan/events.test.ts | 18 ++++++++++++++++++ 2 files changed, 21 insertions(+) diff --git a/apps/server/src/conversation-plan/event-lifecycle.ts b/apps/server/src/conversation-plan/event-lifecycle.ts index 19571045..64dcac9e 100644 --- a/apps/server/src/conversation-plan/event-lifecycle.ts +++ b/apps/server/src/conversation-plan/event-lifecycle.ts @@ -74,6 +74,7 @@ export function applyLifecycleEvent(thread: ConversationPlan.Thread, event: Even thread.decision = undefined; thread.status = "reopened"; thread.pendingSettle = undefined; + thread.pendingScopedChoice = undefined; break; case "candidate.proposed": if (event.source.role !== event.candidate.kind) { @@ -113,6 +114,8 @@ export function applyLifecycleEvent(thread: ConversationPlan.Thread, event: Even thread.decision = undefined; thread.status = "reopened"; } + thread.pendingSettle = undefined; + thread.pendingScopedChoice = undefined; } break; } diff --git a/apps/server/src/conversation-plan/events.test.ts b/apps/server/src/conversation-plan/events.test.ts index a2c3d425..448bbe2f 100644 --- a/apps/server/src/conversation-plan/events.test.ts +++ b/apps/server/src/conversation-plan/events.test.ts @@ -159,6 +159,22 @@ test("contribution, stance, link, and relabel replay retain exact source wording test("candidate confirmation, explicit reopening, and discard preserve decision history", () => { let current = withOption(); + current = accept(current, { type: "card.linked", questionnaireId: "card-1" }); + current = accept(current, { + type: "settle.suggested", + source: citation("resolution", "Use Bun"), + optionId: "option-1", + }); + current = accept(current, { + type: "scoped-choice.proposed", + source: citation("support", "I'd pick Bun for the spike."), + cardId: "card-1", + optionId: "option-1", + label: "Bun", + scope: "spike", + }); + expect(current.threads[0]!.pendingSettle).toBeDefined(); + expect(current.threads[0]!.pendingScopedChoice).toBeDefined(); current = accept(current, { type: "candidate.proposed", source: citation("resolution", "Use Bun"), @@ -172,6 +188,8 @@ test("candidate confirmation, explicit reopening, and discard preserve decision ); expect(current.threads[0]!.status).toBe("decided"); expect(current.threads[0]!.decision!.text).toBe("Use Bun"); + expect(current.threads[0]!.pendingSettle).toBeUndefined(); + expect(current.threads[0]!.pendingScopedChoice).toBeUndefined(); current = accept(current, { type: "decision.reopened", explicit: true }, "thread-1", true); expect(current.threads[0]!.decision).toBeUndefined(); expect(current.threads[0]!.decisionHistory).toHaveLength(1);