From c2fabdee6f00e9cf9522930f701e54360bb51f40 Mon Sep 17 00:00:00 2001 From: Ecko95 <8972676+Ecko95@users.noreply.github.com> Date: Sat, 20 Jun 2026 15:13:51 +0200 Subject: [PATCH 01/10] feat(visual-plan): contracts + orchestration wiring for native visual plans Add a GITS-native visual plan model (block registry, content patches, anchored comments) and wire a thread.visual-plan-upserted event through the orchestration engine end-to-end: decider, projector, SQL projection (migration 031), and the snapshot read-model hydration paths. This is the data plane for the builder.io visual-plan parity feature (Phase 1). Co-Authored-By: Claude Opus 4.8 (1M context) --- .../Layers/OrchestrationEngine.test.ts | 2 + .../Layers/ProjectionPipeline.ts | 32 ++ .../Layers/ProjectionSnapshotQuery.test.ts | 1 + .../Layers/ProjectionSnapshotQuery.ts | 118 ++++- apps/server/src/orchestration/Schemas.ts | 2 + .../orchestration/commandInvariants.test.ts | 2 + apps/server/src/orchestration/decider.ts | 21 + apps/server/src/orchestration/projector.ts | 33 ++ .../Layers/ProjectionThreadVisualPlans.ts | 125 ++++++ apps/server/src/persistence/Migrations.ts | 2 + .../031_ProjectionThreadVisualPlans.ts | 23 + .../Services/ProjectionThreadVisualPlans.ts | 65 +++ apps/server/src/server.test.ts | 2 + apps/server/src/ws.ts | 2 + packages/contracts/src/index.ts | 1 + packages/contracts/src/orchestration.ts | 37 ++ packages/contracts/src/visualPlan.ts | 407 ++++++++++++++++++ 17 files changed, 874 insertions(+), 1 deletion(-) create mode 100644 apps/server/src/persistence/Layers/ProjectionThreadVisualPlans.ts create mode 100644 apps/server/src/persistence/Migrations/031_ProjectionThreadVisualPlans.ts create mode 100644 apps/server/src/persistence/Services/ProjectionThreadVisualPlans.ts create mode 100644 packages/contracts/src/visualPlan.ts diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index e11720d72ca..2bce74b42ef 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -149,6 +149,7 @@ describe("OrchestrationEngine", () => { deletedAt: null, messages: [], proposedPlans: [], + visualPlans: [], activities: [], checkpoints: [], session: null, @@ -161,6 +162,7 @@ describe("OrchestrationEngine", () => { ...thread, messages: [], proposedPlans: [], + visualPlans: [], activities: [], checkpoints: [], })), diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index c201e1d9f9e..c85616fdd7f 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -27,6 +27,7 @@ import { type ProjectionThreadProposedPlan, ProjectionThreadProposedPlanRepository, } from "../../persistence/Services/ProjectionThreadProposedPlans.ts"; +import { ProjectionThreadVisualPlanRepository } from "../../persistence/Services/ProjectionThreadVisualPlans.ts"; import { ProjectionThreadSessionRepository } from "../../persistence/Services/ProjectionThreadSessions.ts"; import { type ProjectionTurn, @@ -39,6 +40,7 @@ import { ProjectionStateRepositoryLive } from "../../persistence/Layers/Projecti import { ProjectionThreadActivityRepositoryLive } from "../../persistence/Layers/ProjectionThreadActivities.ts"; import { ProjectionThreadMessageRepositoryLive } from "../../persistence/Layers/ProjectionThreadMessages.ts"; import { ProjectionThreadProposedPlanRepositoryLive } from "../../persistence/Layers/ProjectionThreadProposedPlans.ts"; +import { ProjectionThreadVisualPlanRepositoryLive } from "../../persistence/Layers/ProjectionThreadVisualPlans.ts"; import { ProjectionThreadSessionRepositoryLive } from "../../persistence/Layers/ProjectionThreadSessions.ts"; import { ProjectionTurnRepositoryLive } from "../../persistence/Layers/ProjectionTurns.ts"; import { ProjectionThreadRepositoryLive } from "../../persistence/Layers/ProjectionThreads.ts"; @@ -59,6 +61,7 @@ export const ORCHESTRATION_PROJECTOR_NAMES = { threads: "projection.threads", threadMessages: "projection.thread-messages", threadProposedPlans: "projection.thread-proposed-plans", + threadVisualPlans: "projection.thread-visual-plans", threadActivities: "projection.thread-activities", threadSessions: "projection.thread-sessions", threadTurns: "projection.thread-turns", @@ -452,6 +455,7 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti const projectionThreadRepository = yield* ProjectionThreadRepository; const projectionThreadMessageRepository = yield* ProjectionThreadMessageRepository; const projectionThreadProposedPlanRepository = yield* ProjectionThreadProposedPlanRepository; + const projectionThreadVisualPlanRepository = yield* ProjectionThreadVisualPlanRepository; const projectionThreadActivityRepository = yield* ProjectionThreadActivityRepository; const projectionThreadSessionRepository = yield* ProjectionThreadSessionRepository; const projectionTurnRepository = yield* ProjectionTurnRepository; @@ -913,6 +917,29 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti } }); + const applyThreadVisualPlansProjection: ProjectorDefinition["apply"] = Effect.fn( + "applyThreadVisualPlansProjection", + )(function* (event, _attachmentSideEffects) { + switch (event.type) { + case "thread.visual-plan-upserted": + yield* projectionThreadVisualPlanRepository.upsert({ + planId: event.payload.visualPlan.id, + threadId: event.payload.threadId, + turnId: event.payload.visualPlan.turnId, + content: event.payload.visualPlan.content, + comments: event.payload.visualPlan.comments, + createdAt: event.payload.visualPlan.createdAt, + updatedAt: event.payload.visualPlan.updatedAt, + }); + return; + + // Visual plans are review documents rather than per-turn outputs, so a + // checkpoint revert intentionally leaves them intact. + default: + return; + } + }); + const applyThreadActivitiesProjection: ProjectorDefinition["apply"] = Effect.fn( "applyThreadActivitiesProjection", )(function* (event, _attachmentSideEffects) { @@ -1378,6 +1405,10 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti name: ORCHESTRATION_PROJECTOR_NAMES.threadProposedPlans, apply: applyThreadProposedPlansProjection, }, + { + name: ORCHESTRATION_PROJECTOR_NAMES.threadVisualPlans, + apply: applyThreadVisualPlansProjection, + }, { name: ORCHESTRATION_PROJECTOR_NAMES.threadActivities, apply: applyThreadActivitiesProjection, @@ -1500,6 +1531,7 @@ export const OrchestrationProjectionPipelineLive = Layer.effect( Layer.provideMerge(ProjectionThreadRepositoryLive), Layer.provideMerge(ProjectionThreadMessageRepositoryLive), Layer.provideMerge(ProjectionThreadProposedPlanRepositoryLive), + Layer.provideMerge(ProjectionThreadVisualPlanRepositoryLive), Layer.provideMerge(ProjectionThreadActivityRepositoryLive), Layer.provideMerge(ProjectionThreadSessionRepositoryLive), Layer.provideMerge(ProjectionTurnRepositoryLive), diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index 7db2a23e5ec..058a3555e80 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -332,6 +332,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { updatedAt: "2026-02-24T00:00:05.500Z", }, ], + visualPlans: [], activities: [ { id: asEventId("activity-1"), diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index e629d1604b3..60f97b717ca 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -20,6 +20,9 @@ import { type OrchestrationSession, type OrchestrationThreadActivity, type OrchestrationThreadShell, + type OrchestrationVisualPlan, + PlanComment, + PlanContent, ModelSelection, ProjectId, ThreadId, @@ -46,6 +49,7 @@ import { ProjectionState } from "../../persistence/Services/ProjectionState.ts"; import { ProjectionThreadActivity } from "../../persistence/Services/ProjectionThreadActivities.ts"; import { ProjectionThreadMessage } from "../../persistence/Services/ProjectionThreadMessages.ts"; import { ProjectionThreadProposedPlan } from "../../persistence/Services/ProjectionThreadProposedPlans.ts"; +import { ProjectionThreadVisualPlan } from "../../persistence/Services/ProjectionThreadVisualPlans.ts"; import { ProjectionThreadSession } from "../../persistence/Services/ProjectionThreadSessions.ts"; import { ProjectionThread } from "../../persistence/Services/ProjectionThreads.ts"; import { RepositoryIdentityResolver } from "../../project/Services/RepositoryIdentityResolver.ts"; @@ -74,6 +78,12 @@ const ProjectionThreadMessageDbRowSchema = ProjectionThreadMessage.mapFields( }), ); const ProjectionThreadProposedPlanDbRowSchema = ProjectionThreadProposedPlan; +const ProjectionThreadVisualPlanDbRowSchema = ProjectionThreadVisualPlan.mapFields( + Struct.assign({ + content: Schema.fromJsonString(PlanContent), + comments: Schema.fromJsonString(Schema.Array(PlanComment)), + }), +); const ProjectionThreadDbRowSchema = ProjectionThread.mapFields( Struct.assign({ modelSelection: Schema.fromJsonString(ModelSelection), @@ -253,6 +263,19 @@ function mapProposedPlanRow( }; } +function mapVisualPlanRow( + row: Schema.Schema.Type, +): OrchestrationVisualPlan { + return { + id: row.planId, + turnId: row.turnId, + content: row.content, + comments: row.comments, + createdAt: row.createdAt, + updatedAt: row.updatedAt, + }; +} + function toPersistenceSqlOrDecodeError(sqlOperation: string, decodeOperation: string) { return (cause: unknown): ProjectionRepositoryError => Schema.isSchemaError(cause) @@ -442,6 +465,24 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { `, }); + const listThreadVisualPlanRows = SqlSchema.findAll({ + Request: Schema.Void, + Result: ProjectionThreadVisualPlanDbRowSchema, + execute: () => + sql` + SELECT + plan_id AS "planId", + thread_id AS "threadId", + turn_id AS "turnId", + content_json AS "content", + comments_json AS "comments", + created_at AS "createdAt", + updated_at AS "updatedAt" + FROM projection_thread_visual_plans + ORDER BY thread_id ASC, created_at ASC, plan_id ASC + `, + }); + const listThreadActivityRows = SqlSchema.findAll({ Request: Schema.Void, Result: ProjectionThreadActivityDbRowSchema, @@ -807,6 +848,25 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { `, }); + const listThreadVisualPlanRowsByThread = SqlSchema.findAll({ + Request: ThreadIdLookupInput, + Result: ProjectionThreadVisualPlanDbRowSchema, + execute: ({ threadId }) => + sql` + SELECT + plan_id AS "planId", + thread_id AS "threadId", + turn_id AS "turnId", + content_json AS "content", + comments_json AS "comments", + created_at AS "createdAt", + updated_at AS "updatedAt" + FROM projection_thread_visual_plans + WHERE thread_id = ${threadId} + ORDER BY created_at ASC, plan_id ASC + `, + }); + const listThreadActivityRowsByThread = SqlSchema.findAll({ Request: ThreadIdLookupInput, Result: ProjectionThreadActivityDbRowSchema, @@ -966,6 +1026,14 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { ), ), ), + listThreadVisualPlanRows(undefined).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionSnapshotQuery.getSnapshot:listThreadVisualPlans:query", + "ProjectionSnapshotQuery.getSnapshot:listThreadVisualPlans:decodeRows", + ), + ), + ), listThreadActivityRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( @@ -1015,6 +1083,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { threadRows, messageRows, proposedPlanRows, + visualPlanRows, activityRows, sessionRows, checkpointRows, @@ -1024,6 +1093,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { Effect.gen(function* () { const messagesByThread = new Map>(); const proposedPlansByThread = new Map>(); + const visualPlansByThread = new Map>(); const activitiesByThread = new Map>(); const checkpointsByThread = new Map>(); const sessionsByThread = new Map(); @@ -1072,6 +1142,13 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { proposedPlansByThread.set(row.threadId, threadProposedPlans); } + for (const row of visualPlanRows) { + updatedAt = maxIso(updatedAt, row.updatedAt); + const threadVisualPlans = visualPlansByThread.get(row.threadId) ?? []; + threadVisualPlans.push(mapVisualPlanRow(row)); + visualPlansByThread.set(row.threadId, threadVisualPlans); + } + for (const row of activityRows) { updatedAt = maxIso(updatedAt, row.createdAt); const threadActivities = activitiesByThread.get(row.threadId) ?? []; @@ -1188,6 +1265,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { deletedAt: row.deletedAt, messages: messagesByThread.get(row.threadId) ?? [], proposedPlans: proposedPlansByThread.get(row.threadId) ?? [], + visualPlans: visualPlansByThread.get(row.threadId) ?? [], activities: activitiesByThread.get(row.threadId) ?? [], checkpoints: checkpointsByThread.get(row.threadId) ?? [], session: sessionsByThread.get(row.threadId) ?? null, @@ -1243,6 +1321,14 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { ), ), ), + listThreadVisualPlanRows(undefined).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionSnapshotQuery.getCommandReadModel:listThreadVisualPlans:query", + "ProjectionSnapshotQuery.getCommandReadModel:listThreadVisualPlans:decodeRows", + ), + ), + ), listThreadSessionRows(undefined).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( @@ -1271,7 +1357,15 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { ) .pipe( Effect.flatMap( - ([projectRows, threadRows, proposedPlanRows, sessionRows, latestTurnRows, stateRows]) => + ([ + projectRows, + threadRows, + proposedPlanRows, + visualPlanRows, + sessionRows, + latestTurnRows, + stateRows, + ]) => Effect.sync(() => { let updatedAt: string | null = null; const projects: OrchestrationProject[] = []; @@ -1345,6 +1439,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { latestTurnByThread.set(row.threadId, mapLatestTurn(row)); } const proposedPlansByThread = new Map>(); + const visualPlansByThread = new Map>(); const sessionByThread = new Map(); for (let index = 0; index < sessionRows.length; index += 1) { @@ -1365,6 +1460,16 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { proposedPlansByThread.set(row.threadId, threadProposedPlans); } + for (let index = 0; index < visualPlanRows.length; index += 1) { + const row = visualPlanRows[index]; + if (!row) { + continue; + } + const threadVisualPlans = visualPlansByThread.get(row.threadId) ?? []; + threadVisualPlans.push(mapVisualPlanRow(row)); + visualPlansByThread.set(row.threadId, threadVisualPlans); + } + for (let index = 0; index < threadRows.length; index += 1) { const row = threadRows[index]; if (!row) { @@ -1386,6 +1491,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { deletedAt: row.deletedAt, messages: [], proposedPlans: proposedPlansByThread.get(row.threadId) ?? [], + visualPlans: visualPlansByThread.get(row.threadId) ?? [], activities: [], checkpoints: [], session: sessionByThread.get(row.threadId) ?? null, @@ -1900,6 +2006,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { threadRow, messageRows, proposedPlanRows, + visualPlanRows, activityRows, checkpointRows, latestTurnRow, @@ -1929,6 +2036,14 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { ), ), ), + listThreadVisualPlanRowsByThread({ threadId }).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionSnapshotQuery.getThreadDetailById:listVisualPlans:query", + "ProjectionSnapshotQuery.getThreadDetailById:listVisualPlans:decodeRows", + ), + ), + ), listThreadActivityRowsByThread({ threadId }).pipe( Effect.mapError( toPersistenceSqlOrDecodeError( @@ -1997,6 +2112,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { return message; }), proposedPlans: proposedPlanRows.map(mapProposedPlanRow), + visualPlans: visualPlanRows.map(mapVisualPlanRow), activities: activityRows.map((row) => { const activity = { id: row.activityId, diff --git a/apps/server/src/orchestration/Schemas.ts b/apps/server/src/orchestration/Schemas.ts index f7ebf693440..b0b2f0cd1ad 100644 --- a/apps/server/src/orchestration/Schemas.ts +++ b/apps/server/src/orchestration/Schemas.ts @@ -11,6 +11,7 @@ import { ThreadUnarchivedPayload as ContractsThreadUnarchivedPayloadSchema, ThreadMessageSentPayload as ContractsThreadMessageSentPayloadSchema, ThreadProposedPlanUpsertedPayload as ContractsThreadProposedPlanUpsertedPayloadSchema, + ThreadVisualPlanUpsertedPayload as ContractsThreadVisualPlanUpsertedPayloadSchema, ThreadSessionSetPayload as ContractsThreadSessionSetPayloadSchema, ThreadTurnDiffCompletedPayload as ContractsThreadTurnDiffCompletedPayloadSchema, ThreadRevertedPayload as ContractsThreadRevertedPayloadSchema, @@ -37,6 +38,7 @@ export const ThreadUnarchivedPayload = ContractsThreadUnarchivedPayloadSchema; export const MessageSentPayloadSchema = ContractsThreadMessageSentPayloadSchema; export const ThreadProposedPlanUpsertedPayload = ContractsThreadProposedPlanUpsertedPayloadSchema; +export const ThreadVisualPlanUpsertedPayload = ContractsThreadVisualPlanUpsertedPayloadSchema; export const ThreadSessionSetPayload = ContractsThreadSessionSetPayloadSchema; export const ThreadTurnDiffCompletedPayload = ContractsThreadTurnDiffCompletedPayloadSchema; export const ThreadRevertedPayload = ContractsThreadRevertedPayloadSchema; diff --git a/apps/server/src/orchestration/commandInvariants.test.ts b/apps/server/src/orchestration/commandInvariants.test.ts index d52f0535fbb..ab05370e451 100644 --- a/apps/server/src/orchestration/commandInvariants.test.ts +++ b/apps/server/src/orchestration/commandInvariants.test.ts @@ -73,6 +73,7 @@ const readModel: OrchestrationReadModel = { session: null, activities: [], proposedPlans: [], + visualPlans: [], checkpoints: [], deletedAt: null, }, @@ -96,6 +97,7 @@ const readModel: OrchestrationReadModel = { session: null, activities: [], proposedPlans: [], + visualPlans: [], checkpoints: [], deletedAt: null, }, diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 0d4af771ca8..76dbde24063 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -675,6 +675,27 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" }; } + case "thread.visual-plan.upsert": { + yield* requireThread({ + readModel, + command, + threadId: command.threadId, + }); + return { + ...(yield* withEventBase({ + aggregateKind: "thread", + aggregateId: command.threadId, + occurredAt: command.createdAt, + commandId: command.commandId, + })), + type: "thread.visual-plan-upserted", + payload: { + threadId: command.threadId, + visualPlan: command.visualPlan, + }, + }; + } + case "thread.turn.diff.complete": { yield* requireThread({ readModel, diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index 0c92f965433..3b7bf642a3c 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -21,6 +21,7 @@ import { ThreadInteractionModeSetPayload, ThreadMetaUpdatedPayload, ThreadProposedPlanUpsertedPayload, + ThreadVisualPlanUpsertedPayload, ThreadRuntimeModeSetPayload, ThreadUnarchivedPayload, ThreadRevertedPayload, @@ -499,6 +500,38 @@ export function projectEvent( }; }); + case "thread.visual-plan-upserted": + return Effect.gen(function* () { + const payload = yield* decodeForEvent( + ThreadVisualPlanUpsertedPayload, + event.payload, + event.type, + "payload", + ); + const thread = nextBase.threads.find((entry) => entry.id === payload.threadId); + if (!thread) { + return nextBase; + } + + const visualPlans = [ + ...thread.visualPlans.filter((entry) => entry.id !== payload.visualPlan.id), + payload.visualPlan, + ] + .toSorted( + (left, right) => + left.createdAt.localeCompare(right.createdAt) || left.id.localeCompare(right.id), + ) + .slice(-50); + + return { + ...nextBase, + threads: updateThread(nextBase.threads, payload.threadId, { + visualPlans, + updatedAt: event.occurredAt, + }), + }; + }); + case "thread.turn-diff-completed": return Effect.gen(function* () { const payload = yield* decodeForEvent( diff --git a/apps/server/src/persistence/Layers/ProjectionThreadVisualPlans.ts b/apps/server/src/persistence/Layers/ProjectionThreadVisualPlans.ts new file mode 100644 index 00000000000..1d2c34560f5 --- /dev/null +++ b/apps/server/src/persistence/Layers/ProjectionThreadVisualPlans.ts @@ -0,0 +1,125 @@ +import * as SqlClient from "effect/unstable/sql/SqlClient"; +import * as SqlSchema from "effect/unstable/sql/SqlSchema"; +import { PlanComment, PlanContent } from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Schema from "effect/Schema"; +import * as Struct from "effect/Struct"; + +import { toPersistenceDecodeError, toPersistenceSqlError } from "../Errors.ts"; +import { + DeleteProjectionThreadVisualPlansInput, + ListProjectionThreadVisualPlansInput, + ProjectionThreadVisualPlan, + ProjectionThreadVisualPlanRepository, + type ProjectionThreadVisualPlanRepositoryShape, +} from "../Services/ProjectionThreadVisualPlans.ts"; + +const ProjectionThreadVisualPlanDbRowSchema = ProjectionThreadVisualPlan.mapFields( + Struct.assign({ + content: Schema.fromJsonString(PlanContent), + comments: Schema.fromJsonString(Schema.Array(PlanComment)), + }), +); + +function toPersistenceSqlOrDecodeError(sqlOperation: string, decodeOperation: string) { + return (cause: unknown) => + Schema.isSchemaError(cause) + ? toPersistenceDecodeError(decodeOperation)(cause) + : toPersistenceSqlError(sqlOperation)(cause); +} + +const makeProjectionThreadVisualPlanRepository = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + const upsertProjectionThreadVisualPlanRow = SqlSchema.void({ + Request: ProjectionThreadVisualPlan, + execute: (row) => sql` + INSERT INTO projection_thread_visual_plans ( + plan_id, + thread_id, + turn_id, + content_json, + comments_json, + created_at, + updated_at + ) + VALUES ( + ${row.planId}, + ${row.threadId}, + ${row.turnId}, + ${JSON.stringify(row.content)}, + ${JSON.stringify(row.comments)}, + ${row.createdAt}, + ${row.updatedAt} + ) + ON CONFLICT (plan_id) + DO UPDATE SET + thread_id = excluded.thread_id, + turn_id = excluded.turn_id, + content_json = excluded.content_json, + comments_json = excluded.comments_json, + created_at = excluded.created_at, + updated_at = excluded.updated_at + `, + }); + + const listProjectionThreadVisualPlanRows = SqlSchema.findAll({ + Request: ListProjectionThreadVisualPlansInput, + Result: ProjectionThreadVisualPlanDbRowSchema, + execute: ({ threadId }) => sql` + SELECT + plan_id AS "planId", + thread_id AS "threadId", + turn_id AS "turnId", + content_json AS "content", + comments_json AS "comments", + created_at AS "createdAt", + updated_at AS "updatedAt" + FROM projection_thread_visual_plans + WHERE thread_id = ${threadId} + ORDER BY created_at ASC, plan_id ASC + `, + }); + + const deleteProjectionThreadVisualPlanRows = SqlSchema.void({ + Request: DeleteProjectionThreadVisualPlansInput, + execute: ({ threadId }) => sql` + DELETE FROM projection_thread_visual_plans + WHERE thread_id = ${threadId} + `, + }); + + const upsert: ProjectionThreadVisualPlanRepositoryShape["upsert"] = (row) => + upsertProjectionThreadVisualPlanRow(row).pipe( + Effect.mapError(toPersistenceSqlError("ProjectionThreadVisualPlanRepository.upsert:query")), + ); + + const listByThreadId: ProjectionThreadVisualPlanRepositoryShape["listByThreadId"] = (input) => + listProjectionThreadVisualPlanRows(input).pipe( + Effect.mapError( + toPersistenceSqlOrDecodeError( + "ProjectionThreadVisualPlanRepository.listByThreadId:query", + "ProjectionThreadVisualPlanRepository.listByThreadId:decode", + ), + ), + ); + + const deleteByThreadId: ProjectionThreadVisualPlanRepositoryShape["deleteByThreadId"] = (input) => + deleteProjectionThreadVisualPlanRows(input).pipe( + Effect.mapError( + toPersistenceSqlError("ProjectionThreadVisualPlanRepository.deleteByThreadId:query"), + ), + ); + + return { + upsert, + listByThreadId, + deleteByThreadId, + } satisfies ProjectionThreadVisualPlanRepositoryShape; +}); + +export const ProjectionThreadVisualPlanRepositoryLive = Layer.effect( + ProjectionThreadVisualPlanRepository, + makeProjectionThreadVisualPlanRepository, +); diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index cc5024d5f51..8f0e1ad345a 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -43,6 +43,7 @@ import Migration0027 from "./Migrations/027_ProviderSessionRuntimeInstanceId.ts" import Migration0028 from "./Migrations/028_ProjectionThreadSessionInstanceId.ts"; import Migration0029 from "./Migrations/029_ProjectionThreadDetailOrderingIndexes.ts"; import Migration0030 from "./Migrations/030_ProjectionThreadShellArchiveIndexes.ts"; +import Migration0031 from "./Migrations/031_ProjectionThreadVisualPlans.ts"; /** * Migration loader with all migrations defined inline. @@ -85,6 +86,7 @@ export const migrationEntries = [ [28, "ProjectionThreadSessionInstanceId", Migration0028], [29, "ProjectionThreadDetailOrderingIndexes", Migration0029], [30, "ProjectionThreadShellArchiveIndexes", Migration0030], + [31, "ProjectionThreadVisualPlans", Migration0031], ] as const; export const makeMigrationLoader = (throughId?: number) => diff --git a/apps/server/src/persistence/Migrations/031_ProjectionThreadVisualPlans.ts b/apps/server/src/persistence/Migrations/031_ProjectionThreadVisualPlans.ts new file mode 100644 index 00000000000..0bb874ef8ad --- /dev/null +++ b/apps/server/src/persistence/Migrations/031_ProjectionThreadVisualPlans.ts @@ -0,0 +1,23 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + yield* sql` + CREATE TABLE IF NOT EXISTS projection_thread_visual_plans ( + plan_id TEXT PRIMARY KEY, + thread_id TEXT NOT NULL, + turn_id TEXT, + content_json TEXT NOT NULL, + comments_json TEXT NOT NULL, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL + ) + `; + + yield* sql` + CREATE INDEX IF NOT EXISTS idx_projection_thread_visual_plans_thread_created + ON projection_thread_visual_plans(thread_id, created_at) + `; +}); diff --git a/apps/server/src/persistence/Services/ProjectionThreadVisualPlans.ts b/apps/server/src/persistence/Services/ProjectionThreadVisualPlans.ts new file mode 100644 index 00000000000..70f00d64cb8 --- /dev/null +++ b/apps/server/src/persistence/Services/ProjectionThreadVisualPlans.ts @@ -0,0 +1,65 @@ +/** + * ProjectionThreadVisualPlanRepository - Projection repository interface for + * thread visual plans (the GITS-native builder.io-style visual plan documents). + * + * Owns persistence for the per-thread `PlanContent` document + comments + * projected from `thread.visual-plan-upserted` orchestration events. Content + * and comments are stored as JSON columns. + * + * @module ProjectionThreadVisualPlanRepository + */ +import { + IsoDateTime, + OrchestrationVisualPlanId, + PlanComment, + PlanContent, + ThreadId, + TurnId, +} from "@t3tools/contracts"; +import * as Schema from "effect/Schema"; +import * as Context from "effect/Context"; +import type * as Effect from "effect/Effect"; + +import type { ProjectionRepositoryError } from "../Errors.ts"; + +export const ProjectionThreadVisualPlan = Schema.Struct({ + planId: OrchestrationVisualPlanId, + threadId: ThreadId, + turnId: Schema.NullOr(TurnId), + content: PlanContent, + comments: Schema.Array(PlanComment), + createdAt: IsoDateTime, + updatedAt: IsoDateTime, +}); +export type ProjectionThreadVisualPlan = typeof ProjectionThreadVisualPlan.Type; + +export const ListProjectionThreadVisualPlansInput = Schema.Struct({ + threadId: ThreadId, +}); +export type ListProjectionThreadVisualPlansInput = + typeof ListProjectionThreadVisualPlansInput.Type; + +export const DeleteProjectionThreadVisualPlansInput = Schema.Struct({ + threadId: ThreadId, +}); +export type DeleteProjectionThreadVisualPlansInput = + typeof DeleteProjectionThreadVisualPlansInput.Type; + +export interface ProjectionThreadVisualPlanRepositoryShape { + readonly upsert: ( + visualPlan: ProjectionThreadVisualPlan, + ) => Effect.Effect; + readonly listByThreadId: ( + input: ListProjectionThreadVisualPlansInput, + ) => Effect.Effect, ProjectionRepositoryError>; + readonly deleteByThreadId: ( + input: DeleteProjectionThreadVisualPlansInput, + ) => Effect.Effect; +} + +export class ProjectionThreadVisualPlanRepository extends Context.Service< + ProjectionThreadVisualPlanRepository, + ProjectionThreadVisualPlanRepositoryShape +>()( + "t3/persistence/Services/ProjectionThreadVisualPlans/ProjectionThreadVisualPlanRepository", +) {} diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 4ed78d1998d..f67ea001e39 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -582,6 +582,7 @@ const makeDefaultOrchestrationReadModel = () => { session: null, activities: [], proposedPlans: [], + visualPlans: [], checkpoints: [], deletedAt: null, }, @@ -4691,6 +4692,7 @@ it.layer(NodeServices.layer)("server router seam", (it) => { session: null, activities: [], proposedPlans: [], + visualPlans: [], checkpoints: [], deletedAt: null, }, diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 79d1d3866f9..189a43f8622 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -133,6 +133,7 @@ function isThreadDetailEvent(event: OrchestrationEvent): event is Extract< type: | "thread.message-sent" | "thread.proposed-plan-upserted" + | "thread.visual-plan-upserted" | "thread.activity-appended" | "thread.turn-diff-completed" | "thread.reverted" @@ -142,6 +143,7 @@ function isThreadDetailEvent(event: OrchestrationEvent): event is Extract< return ( event.type === "thread.message-sent" || event.type === "thread.proposed-plan-upserted" || + event.type === "thread.visual-plan-upserted" || event.type === "thread.activity-appended" || event.type === "thread.turn-diff-completed" || event.type === "thread.reverted" || diff --git a/packages/contracts/src/index.ts b/packages/contracts/src/index.ts index 7f8b51aeaf6..a5f173bc362 100644 --- a/packages/contracts/src/index.ts +++ b/packages/contracts/src/index.ts @@ -15,6 +15,7 @@ export * from "./settings.ts"; export * from "./git.ts"; export * from "./vcs.ts"; export * from "./sourceControl.ts"; +export * from "./visualPlan.ts"; export * from "./orchestration.ts"; export * from "./editor.ts"; export * from "./project.ts"; diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 401928171c8..4071baa5cc0 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -21,6 +21,7 @@ import { TurnId, } from "./baseSchemas.ts"; import { ProviderInstanceId } from "./providerInstance.ts"; +import { PlanComment, PlanContent } from "./visualPlan.ts"; export const ORCHESTRATION_WS_METHODS = { dispatchCommand: "orchestration.dispatchCommand", @@ -246,6 +247,19 @@ const SourceProposedPlanReference = Schema.Struct({ planId: OrchestrationProposedPlanId, }); +export const OrchestrationVisualPlanId = TrimmedNonEmptyString; +export type OrchestrationVisualPlanId = typeof OrchestrationVisualPlanId.Type; + +export const OrchestrationVisualPlan = Schema.Struct({ + id: OrchestrationVisualPlanId, + turnId: Schema.NullOr(TurnId), + content: PlanContent, + comments: Schema.Array(PlanComment).pipe(Schema.withDecodingDefault(Effect.succeed([]))), + createdAt: IsoDateTime, + updatedAt: IsoDateTime, +}); +export type OrchestrationVisualPlan = typeof OrchestrationVisualPlan.Type; + export const OrchestrationSessionStatus = Schema.Literals([ "idle", "starting", @@ -350,6 +364,9 @@ export const OrchestrationThread = Schema.Struct({ proposedPlans: Schema.Array(OrchestrationProposedPlan).pipe( Schema.withDecodingDefault(Effect.succeed([])), ), + visualPlans: Schema.Array(OrchestrationVisualPlan).pipe( + Schema.withDecodingDefault(Effect.succeed([])), + ), activities: Schema.Array(OrchestrationThreadActivity), checkpoints: Schema.Array(OrchestrationCheckpointSummary), session: Schema.NullOr(OrchestrationSession), @@ -721,6 +738,14 @@ const ThreadProposedPlanUpsertCommand = Schema.Struct({ createdAt: IsoDateTime, }); +const ThreadVisualPlanUpsertCommand = Schema.Struct({ + type: Schema.Literal("thread.visual-plan.upsert"), + commandId: CommandId, + threadId: ThreadId, + visualPlan: OrchestrationVisualPlan, + createdAt: IsoDateTime, +}); + const ThreadTurnDiffCompleteCommand = Schema.Struct({ type: Schema.Literal("thread.turn.diff.complete"), commandId: CommandId, @@ -756,6 +781,7 @@ const InternalOrchestrationCommand = Schema.Union([ ThreadMessageAssistantDeltaCommand, ThreadMessageAssistantCompleteCommand, ThreadProposedPlanUpsertCommand, + ThreadVisualPlanUpsertCommand, ThreadTurnDiffCompleteCommand, ThreadActivityAppendCommand, ThreadRevertCompleteCommand, @@ -789,6 +815,7 @@ export const OrchestrationEventType = Schema.Literals([ "thread.session-stop-requested", "thread.session-set", "thread.proposed-plan-upserted", + "thread.visual-plan-upserted", "thread.turn-diff-completed", "thread.activity-appended", ]); @@ -949,6 +976,11 @@ export const ThreadProposedPlanUpsertedPayload = Schema.Struct({ proposedPlan: OrchestrationProposedPlan, }); +export const ThreadVisualPlanUpsertedPayload = Schema.Struct({ + threadId: ThreadId, + visualPlan: OrchestrationVisualPlan, +}); + export const ThreadTurnDiffCompletedPayload = Schema.Struct({ threadId: ThreadId, turnId: TurnId, @@ -1087,6 +1119,11 @@ export const OrchestrationEvent = Schema.Union([ type: Schema.Literal("thread.proposed-plan-upserted"), payload: ThreadProposedPlanUpsertedPayload, }), + Schema.Struct({ + ...EventBaseFields, + type: Schema.Literal("thread.visual-plan-upserted"), + payload: ThreadVisualPlanUpsertedPayload, + }), Schema.Struct({ ...EventBaseFields, type: Schema.Literal("thread.turn-diff-completed"), diff --git a/packages/contracts/src/visualPlan.ts b/packages/contracts/src/visualPlan.ts new file mode 100644 index 00000000000..c9f0824caf6 --- /dev/null +++ b/packages/contracts/src/visualPlan.ts @@ -0,0 +1,407 @@ +import * as Effect from "effect/Effect"; +import * as Schema from "effect/Schema"; +import { IsoDateTime, TrimmedNonEmptyString } from "./baseSchemas.ts"; + +/** + * Visual plan content model — a GITS-native reimplementation of the + * builder.io `visual-plan` block registry. The agent authors a `PlanContent` + * document through the visual-plan MCP tools; the GITS web app renders it as + * an interactive side panel. + * + * v1 renders the document blocks below. `diagram`/`custom-html` are defined + * here but rendered in Phase 2; `wireframe`/`prototype`/`canvas` arrive in + * Phase 3 and are intentionally not part of the v1 union yet. + */ + +export const PlanBlockId = TrimmedNonEmptyString; +export type PlanBlockId = typeof PlanBlockId.Type; + +const blockBase = { + id: PlanBlockId, + title: Schema.optional(Schema.String), + summary: Schema.optional(Schema.String), + editable: Schema.optional(Schema.Boolean), +}; + +export const PlanRichTextBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("rich-text"), + data: Schema.Struct({ markdown: Schema.String }), +}); + +export const PlanCalloutTone = Schema.Literals(["info", "decision", "risk", "warning", "success"]); +export type PlanCalloutTone = typeof PlanCalloutTone.Type; +export const PlanCalloutBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("callout"), + data: Schema.Struct({ + tone: Schema.optional(PlanCalloutTone), + body: Schema.String, + }), +}); + +export const PlanChecklistItem = Schema.Struct({ + id: PlanBlockId, + label: Schema.String, + checked: Schema.optional(Schema.Boolean), + note: Schema.optional(Schema.String), +}); +export const PlanChecklistBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("checklist"), + data: Schema.Struct({ items: Schema.Array(PlanChecklistItem) }), +}); + +export const PlanTableDensity = Schema.Literals(["compact", "normal", "relaxed"]); +export const PlanTableBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("table"), + data: Schema.Struct({ + columns: Schema.Array(Schema.String), + rows: Schema.Array(Schema.Array(Schema.String)), + density: Schema.optional(PlanTableDensity), + }), +}); + +export const PlanCodeBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("code"), + data: Schema.Struct({ + language: Schema.optional(Schema.String), + code: Schema.String, + filename: Schema.optional(Schema.String), + caption: Schema.optional(Schema.String), + }), +}); + +export const PlanCodeAnnotation = Schema.Struct({ + lines: Schema.String, + label: Schema.optional(Schema.String), + note: Schema.String, +}); +export const PlanAnnotatedCodeBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("annotated-code"), + data: Schema.Struct({ + filename: Schema.optional(Schema.String), + language: Schema.optional(Schema.String), + code: Schema.String, + annotations: Schema.optional(Schema.Array(PlanCodeAnnotation)), + }), +}); + +export const PlanChangeKind = Schema.Literals(["added", "modified", "removed", "renamed"]); +export const PlanFileTreeEntry = Schema.Struct({ + path: Schema.String, + change: Schema.optional(PlanChangeKind), + note: Schema.optional(Schema.String), +}); +export const PlanFileTreeBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("file-tree"), + data: Schema.Struct({ entries: Schema.Array(PlanFileTreeEntry) }), +}); + +export const PlanImplementationMapFile = Schema.Struct({ + path: Schema.String, + title: Schema.optional(Schema.String), + note: Schema.String, + language: Schema.optional(Schema.String), + snippet: Schema.optional(Schema.String), +}); +export const PlanImplementationMapBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("implementation-map"), + data: Schema.Struct({ files: Schema.Array(PlanImplementationMapFile) }), +}); + +export const PlanApiMethod = Schema.Literals([ + "GET", + "POST", + "PUT", + "PATCH", + "DELETE", + "HEAD", + "OPTIONS", +]); +export const PlanApiParam = Schema.Struct({ + name: Schema.String, + in: Schema.Literals(["path", "query", "header", "body"]), + type: Schema.optional(Schema.String), + required: Schema.optional(Schema.Boolean), + description: Schema.optional(Schema.String), +}); +export const PlanApiResponse = Schema.Struct({ + status: Schema.String, + description: Schema.optional(Schema.String), + example: Schema.optional(Schema.String), +}); +export const PlanApiEndpointBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("api-endpoint"), + data: Schema.Struct({ + method: PlanApiMethod, + path: Schema.String, + summary: Schema.optional(Schema.String), + description: Schema.optional(Schema.String), + params: Schema.optional(Schema.Array(PlanApiParam)), + responses: Schema.optional(Schema.Array(PlanApiResponse)), + }), +}); + +export const PlanDataModelField = Schema.Struct({ + name: Schema.String, + type: Schema.optional(Schema.String), + pk: Schema.optional(Schema.Boolean), + fk: Schema.optional(Schema.Boolean), + nullable: Schema.optional(Schema.Boolean), + note: Schema.optional(Schema.String), +}); +export const PlanDataModelEntity = Schema.Struct({ + id: PlanBlockId, + name: Schema.String, + note: Schema.optional(Schema.String), + fields: Schema.Array(PlanDataModelField), +}); +export const PlanDataModelRelation = Schema.Struct({ + from: Schema.String, + to: Schema.String, + kind: Schema.optional(Schema.Literals(["1-1", "1-n", "n-n"])), + label: Schema.optional(Schema.String), +}); +export const PlanDataModelBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("data-model"), + data: Schema.Struct({ + entities: Schema.Array(PlanDataModelEntity), + relations: Schema.optional(Schema.Array(PlanDataModelRelation)), + }), +}); + +export const PlanQuestionMode = Schema.Literals(["single", "multi", "freeform"]); +export const PlanQuestionOption = Schema.Struct({ + id: PlanBlockId, + label: Schema.String, + recommended: Schema.optional(Schema.Boolean), +}); +export const PlanQuestion = Schema.Struct({ + id: PlanBlockId, + title: Schema.String, + subtitle: Schema.optional(Schema.String), + mode: PlanQuestionMode, + options: Schema.optional(Schema.Array(PlanQuestionOption)), + allowOther: Schema.optional(Schema.Boolean), + placeholder: Schema.optional(Schema.String), + required: Schema.optional(Schema.Boolean), +}); +export type PlanQuestion = typeof PlanQuestion.Type; +export const PlanQuestionFormBlock = Schema.Struct({ + ...blockBase, + type: Schema.Literal("question-form"), + data: Schema.Struct({ + questions: Schema.Array(PlanQuestion), + submitLabel: Schema.optional(Schema.String), + }), +}); + +// Inert scoped HTML/CSS — rendered in Phase 2 (no