Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
90 changes: 90 additions & 0 deletions apps/server/src/conversation-plan/correction-validation.ts
Original file line number Diff line number Diff line change
@@ -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);
}
82 changes: 82 additions & 0 deletions apps/server/src/conversation-plan/domain-analysis.test.ts
Original file line number Diff line number Diff line change
@@ -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<ConversationPlan.AnalysisRecord, "messageId" | "eventIds"> = {
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();
});
});
97 changes: 97 additions & 0 deletions apps/server/src/conversation-plan/domain-analysis.ts
Original file line number Diff line number Diff line change
@@ -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<ConversationPlan.AnalysisRecord, "messageId" | "eventIds">,
): 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;
}
Loading
Loading