diff --git a/apps/server/src/auth/RpcAuthorization.ts b/apps/server/src/auth/RpcAuthorization.ts index 5d8e5c0baac..a864f628c5e 100644 --- a/apps/server/src/auth/RpcAuthorization.ts +++ b/apps/server/src/auth/RpcAuthorization.ts @@ -66,6 +66,8 @@ export const RPC_REQUIRED_SCOPES = { [WS_METHODS.shellOpenInEditor]: AuthOrchestrationOperateScope, [WS_METHODS.filesystemBrowse]: AuthOrchestrationReadScope, [WS_METHODS.assetsCreateUrl]: AuthOrchestrationReadScope, + [WS_METHODS.composerDraftUpdate]: AuthOrchestrationOperateScope, + [WS_METHODS.subscribeComposerDraft]: AuthOrchestrationReadScope, [WS_METHODS.subscribeVcsStatus]: AuthOrchestrationReadScope, [WS_METHODS.subscribeResourceTelemetry]: AuthOrchestrationReadScope, [WS_METHODS.vcsRefreshStatus]: AuthOrchestrationReadScope, diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index c697b4bd98f..3175e1f94be 100644 --- a/apps/server/src/environment/ServerEnvironment.ts +++ b/apps/server/src/environment/ServerEnvironment.ts @@ -148,6 +148,7 @@ export const make = Effect.gen(function* () { threadPinning: true, threadPinReorder: true, threadTitleRegeneration: true, + composerDraftSync: true, ...(serverSelfUpdate === null ? {} : { serverSelfUpdate }), ...(serverSelfUpdate === "boot-service" ? { serverSelfUpdateProgress: true } : {}), }, diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts index 8e65295b1ba..58260f9031c 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.test.ts @@ -239,6 +239,50 @@ it.layer(BaseTestLayer)("OrchestrationProjectionPipeline", (it) => { assert.deepEqual(unsettledRows, [{ settledOverride: "active", settledAt: null }]); }), ); + + it.effect("removes the composer draft when its thread is deleted", () => + Effect.gen(function* () { + const projectionPipeline = yield* OrchestrationProjectionPipeline; + const eventStore = yield* OrchestrationEventStore; + const sql = yield* SqlClient.SqlClient; + const threadId = ThreadId.make("thread-with-composer-draft"); + const now = "2026-01-01T00:00:00.000Z"; + const commonJson = + '{"text":"delete me","modelSelection":null,"runtimeMode":null,"interactionMode":null}'; + + yield* sql` + INSERT INTO composer_drafts ( + thread_id, revision, common_json, updated_at, client_mutation_id + ) VALUES ( + ${threadId}, 1, ${commonJson}, ${now}, 'test:delete-composer-draft' + ) + `; + yield* eventStore.append({ + type: "thread.deleted", + eventId: EventId.make("evt-delete-composer-draft"), + aggregateKind: "thread", + aggregateId: threadId, + occurredAt: now, + commandId: CommandId.make("cmd-delete-composer-draft"), + causationEventId: null, + correlationId: CommandId.make("cmd-delete-composer-draft"), + metadata: {}, + payload: { + threadId, + deletedAt: now, + }, + }); + + yield* projectionPipeline.bootstrap; + + const remaining = yield* sql<{ readonly count: number }>` + SELECT COUNT(*) AS count + FROM composer_drafts + WHERE thread_id = ${threadId} + `; + assert.deepEqual(remaining, [{ count: 0 }]); + }), + ); }); it.layer(Layer.fresh(makeProjectionPipelinePrefixedTestLayer("t3-base-")))( diff --git a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts index 38a70240d97..1e64904af8f 100644 --- a/apps/server/src/orchestration/Layers/ProjectionPipeline.ts +++ b/apps/server/src/orchestration/Layers/ProjectionPipeline.ts @@ -828,6 +828,10 @@ const makeOrchestrationProjectionPipeline = Effect.fn("makeOrchestrationProjecti case "thread.deleted": { attachmentSideEffects.deletedThreadIds.add(event.payload.threadId); + yield* sql` + DELETE FROM composer_drafts + WHERE thread_id = ${event.payload.threadId} + `.pipe(Effect.mapError(toPersistenceSqlError("delete thread composer draft"))); const existingRow = yield* projectionThreadRepository.getById({ threadId: event.payload.threadId, }); diff --git a/apps/server/src/persistence/ComposerDrafts.test.ts b/apps/server/src/persistence/ComposerDrafts.test.ts new file mode 100644 index 00000000000..779ef200fe0 --- /dev/null +++ b/apps/server/src/persistence/ComposerDrafts.test.ts @@ -0,0 +1,119 @@ +import { ProviderInstanceId, ThreadId } from "@t3tools/contracts"; +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; + +import * as ComposerDrafts from "./ComposerDrafts.ts"; +import { runMigrations } from "./Migrations.ts"; +import * as NodeSqliteClient from "./NodeSqliteClient.ts"; + +const layer = it.layer( + ComposerDrafts.layer.pipe(Layer.provideMerge(NodeSqliteClient.layerMemory())), +); + +layer("ComposerDraftRepository", (it) => { + it.effect("uses revision compare-and-swap and preserves the winning snapshot", () => + Effect.gen(function* () { + yield* runMigrations({ toMigrationInclusive: 39 }); + const repository = yield* ComposerDrafts.ComposerDraftRepository; + const threadId = ThreadId.make("draft-cas-thread"); + const common = { + text: "hello from device one", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5.4", + }, + runtimeMode: "full-access" as const, + interactionMode: "default" as const, + }; + + const accepted = yield* repository.update({ + threadId, + baseRevision: 0, + common, + clientMutationId: "device-one-1", + }); + assert.equal(accepted._tag, "accepted"); + assert.equal(accepted.snapshot.revision, 1); + + const conflict = yield* repository.update({ + threadId, + baseRevision: 0, + common: { ...common, text: "stale device" }, + clientMutationId: "device-two-1", + }); + assert.equal(conflict._tag, "conflict"); + assert.deepEqual(conflict.snapshot.common, common); + assert.equal(conflict.snapshot.revision, 1); + }), + ); + + it.effect("keeps clears as revisioned tombstones", () => + Effect.gen(function* () { + yield* runMigrations({ toMigrationInclusive: 39 }); + const repository = yield* ComposerDrafts.ComposerDraftRepository; + const threadId = ThreadId.make("draft-tombstone-thread"); + + yield* repository.update({ + threadId, + baseRevision: 0, + common: { + text: "sent later", + modelSelection: null, + runtimeMode: null, + interactionMode: null, + }, + clientMutationId: "write-1", + }); + const cleared = yield* repository.update({ + threadId, + baseRevision: 1, + common: null, + clientMutationId: "clear-2", + }); + + assert.equal(cleared._tag, "accepted"); + assert.equal(cleared.snapshot.revision, 2); + assert.isNull(cleared.snapshot.common); + assert.deepEqual(yield* repository.get({ threadId }), cleared.snapshot); + }), + ); + + it.effect("does not let a delayed send clear a newer device revision", () => + Effect.gen(function* () { + yield* runMigrations({ toMigrationInclusive: 39 }); + const repository = yield* ComposerDrafts.ComposerDraftRepository; + const threadId = ThreadId.make("draft-delayed-send-thread"); + const first = { + text: "message being sent", + modelSelection: null, + runtimeMode: null, + interactionMode: null, + }; + const newer = { ...first, text: "new text from another device" }; + + yield* repository.update({ + threadId, + baseRevision: 0, + common: first, + clientMutationId: "first-device", + }); + yield* repository.update({ + threadId, + baseRevision: 1, + common: newer, + clientMutationId: "second-device", + }); + const delayedClear = yield* repository.update({ + threadId, + baseRevision: 1, + common: null, + clientMutationId: "delayed-send", + }); + + assert.equal(delayedClear._tag, "conflict"); + assert.equal(delayedClear.snapshot.revision, 2); + assert.deepEqual(delayedClear.snapshot.common, newer); + }), + ); +}); diff --git a/apps/server/src/persistence/ComposerDrafts.ts b/apps/server/src/persistence/ComposerDrafts.ts new file mode 100644 index 00000000000..63b84b18126 --- /dev/null +++ b/apps/server/src/persistence/ComposerDrafts.ts @@ -0,0 +1,208 @@ +import { + ComposerDraftCommon, + type ComposerDraftGetInput, + type ComposerDraftSnapshot, + type ComposerDraftUpdateInput, + type ComposerDraftUpdateResult, + ThreadId, +} from "@t3tools/contracts"; +import * as Context from "effect/Context"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; +import * as PubSub from "effect/PubSub"; +import * as Schema from "effect/Schema"; +import * as Stream from "effect/Stream"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; +import * as SqlSchema from "effect/unstable/sql/SqlSchema"; + +export class ComposerDraftPersistenceError extends Schema.TaggedErrorClass()( + "ComposerDraftPersistenceError", + { + operation: Schema.String, + cause: Schema.Defect(), + }, +) {} + +export class ComposerDraftRepository extends Context.Service< + ComposerDraftRepository, + { + readonly get: ( + input: ComposerDraftGetInput, + ) => Effect.Effect; + readonly update: ( + input: ComposerDraftUpdateInput, + ) => Effect.Effect; + readonly subscribe: ( + input: ComposerDraftGetInput, + ) => Stream.Stream; + } +>()("t3/persistence/ComposerDrafts/ComposerDraftRepository") {} + +const DbRow = Schema.Struct({ + threadId: ThreadId, + revision: Schema.Int, + common: Schema.NullOr(Schema.fromJsonString(ComposerDraftCommon)), + updatedAt: Schema.String, + clientMutationId: Schema.String, +}); + +const RawDbRow = Schema.Struct({ + threadId: Schema.Unknown, + revision: Schema.Unknown, + common: Schema.Unknown, + updatedAt: Schema.Unknown, + clientMutationId: Schema.Unknown, +}); + +const WriteRow = Schema.Struct({ + threadId: ThreadId, + baseRevision: Schema.Int, + nextRevision: Schema.Int, + common: Schema.NullOr(Schema.fromJsonString(ComposerDraftCommon)), + updatedAt: Schema.String, + clientMutationId: Schema.String, +}); + +const decodeRow = Schema.decodeUnknownEffect(DbRow); +const currentIsoTimestamp = DateTime.now.pipe(Effect.map(DateTime.formatIso)); + +function emptySnapshot(threadId: ThreadId): ComposerDraftSnapshot { + return { + threadId, + revision: 0, + common: null, + updatedAt: null, + clientMutationId: null, + }; +} + +function toSnapshot(row: typeof DbRow.Type): ComposerDraftSnapshot { + return { + threadId: row.threadId, + revision: row.revision, + common: row.common, + updatedAt: row.updatedAt, + clientMutationId: row.clientMutationId, + }; +} + +function mapPersistenceError(operation: string) { + return (cause: unknown) => new ComposerDraftPersistenceError({ operation, cause }); +} + +export const make = Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const changes = yield* PubSub.unbounded(); + + const getRow = SqlSchema.findOneOption({ + Request: Schema.Struct({ threadId: ThreadId }), + Result: RawDbRow, + execute: ({ threadId }) => sql` + SELECT + thread_id AS "threadId", + revision, + common_json AS "common", + updated_at AS "updatedAt", + client_mutation_id AS "clientMutationId" + FROM composer_drafts + WHERE thread_id = ${threadId} + `, + }); + + const insertRow = SqlSchema.findAll({ + Request: WriteRow, + Result: RawDbRow, + execute: (row) => sql` + INSERT INTO composer_drafts ( + thread_id, revision, common_json, updated_at, client_mutation_id + ) VALUES ( + ${row.threadId}, ${row.nextRevision}, ${row.common}, ${row.updatedAt}, + ${row.clientMutationId} + ) + ON CONFLICT (thread_id) DO NOTHING + RETURNING + thread_id AS "threadId", + revision, + common_json AS "common", + updated_at AS "updatedAt", + client_mutation_id AS "clientMutationId" + `, + }); + + const updateRow = SqlSchema.findAll({ + Request: WriteRow, + Result: RawDbRow, + execute: (row) => sql` + UPDATE composer_drafts + SET + revision = ${row.nextRevision}, + common_json = ${row.common}, + updated_at = ${row.updatedAt}, + client_mutation_id = ${row.clientMutationId} + WHERE thread_id = ${row.threadId} + AND revision = ${row.baseRevision} + RETURNING + thread_id AS "threadId", + revision, + common_json AS "common", + updated_at AS "updatedAt", + client_mutation_id AS "clientMutationId" + `, + }); + + const get: ComposerDraftRepository["Service"]["get"] = Effect.fn("ComposerDraftRepository.get")( + function* (input) { + const row = yield* getRow(input).pipe( + Effect.mapError(mapPersistenceError("ComposerDraftRepository.get:query")), + ); + if (Option.isNone(row)) return emptySnapshot(input.threadId); + const decoded = yield* decodeRow(row.value).pipe( + Effect.mapError(mapPersistenceError("ComposerDraftRepository.get:decode")), + ); + return toSnapshot(decoded); + }, + ); + + const update: ComposerDraftRepository["Service"]["update"] = Effect.fn( + "ComposerDraftRepository.update", + )(function* (input) { + const write = { + ...input, + nextRevision: input.baseRevision + 1, + updatedAt: yield* currentIsoTimestamp, + }; + const rows = yield* (input.baseRevision === 0 ? insertRow(write) : updateRow(write)).pipe( + Effect.mapError(mapPersistenceError("ComposerDraftRepository.update:query")), + ); + const row = rows[0]; + if (row === undefined) { + return { _tag: "conflict", snapshot: yield* get(input) } as const; + } + const decoded = yield* decodeRow(row).pipe( + Effect.mapError(mapPersistenceError("ComposerDraftRepository.update:decode")), + ); + const snapshot = toSnapshot(decoded); + yield* PubSub.publish(changes, snapshot); + return { _tag: "accepted", snapshot } as const; + }); + + const subscribe: ComposerDraftRepository["Service"]["subscribe"] = (input) => + Stream.unwrap( + Effect.gen(function* () { + const subscription = yield* PubSub.subscribe(changes); + const initial = yield* get(input); + return Stream.concat( + Stream.make(initial), + Stream.fromSubscription(subscription).pipe( + Stream.filter((snapshot) => snapshot.threadId === input.threadId), + ), + ); + }), + ); + + return ComposerDraftRepository.of({ get, update, subscribe }); +}); + +export const layer = Layer.effect(ComposerDraftRepository, make); diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index 733c52fab3e..bf8dc2529f7 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -51,6 +51,7 @@ import Migration0035 from "./Migrations/035_ProjectionThreadTitleRegeneration.ts import Migration0036 from "./Migrations/036_ProjectionThreadsPinned.ts"; import Migration0037 from "./Migrations/037_ProjectionTurnsKeysetIndex.ts"; import Migration0038 from "./Migrations/038_ProjectionThreadsPinOrderKey.ts"; +import Migration0039 from "./Migrations/039_ComposerDrafts.ts"; /** * Migration loader with all migrations defined inline. @@ -101,6 +102,7 @@ export const migrationEntries = [ [36, "ProjectionThreadsPinned", Migration0036], [37, "ProjectionTurnsKeysetIndex", Migration0037], [38, "ProjectionThreadsPinOrderKey", Migration0038], + [39, "ComposerDrafts", Migration0039], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/039_ComposerDrafts.ts b/apps/server/src/persistence/Migrations/039_ComposerDrafts.ts new file mode 100644 index 00000000000..3fc47c4728e --- /dev/null +++ b/apps/server/src/persistence/Migrations/039_ComposerDrafts.ts @@ -0,0 +1,16 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +/** Current-state storage for high-churn existing-thread composer drafts. */ +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + yield* sql` + CREATE TABLE IF NOT EXISTS composer_drafts ( + thread_id TEXT PRIMARY KEY NOT NULL, + revision INTEGER NOT NULL CHECK (revision >= 0), + common_json TEXT, + updated_at TEXT NOT NULL, + client_mutation_id TEXT NOT NULL + ) + `; +}); diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index d982c2e192c..19b4b005df4 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -110,6 +110,7 @@ import * as ExternalLauncher from "./process/externalLauncher.ts"; import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import { OrchestrationListenerCallbackError } from "./orchestration/Errors.ts"; import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts"; +import * as ComposerDrafts from "./persistence/ComposerDrafts.ts"; import { SqlitePersistenceMemory } from "./persistence/Layers/Sqlite.ts"; import { PersistenceSqlError } from "./persistence/Errors.ts"; import * as ProviderRegistry from "./provider/Services/ProviderRegistry.ts"; @@ -616,14 +617,17 @@ const buildAppUnderTest = (options?: { }, ).pipe( Layer.provide( - Layer.mock(Keybindings.Keybindings)({ - loadConfigState: Effect.succeed({ - keybindings: [], - issues: [], + Layer.mergeAll( + ComposerDrafts.layer.pipe(Layer.provide(SqlitePersistenceMemory)), + Layer.mock(Keybindings.Keybindings)({ + loadConfigState: Effect.succeed({ + keybindings: [], + issues: [], + }), + streamChanges: Stream.empty, + ...options?.layers?.keybindings, }), - streamChanges: Stream.empty, - ...options?.layers?.keybindings, - }), + ), ), Layer.provide( Layer.mock(ProviderRegistry.ProviderRegistry)({ @@ -4477,6 +4481,105 @@ it.layer(NodeServices.layer)("server router seam", (it) => { ).pipe(Effect.provide(NodeHttpServer.layerTest)), ); + it.effect("shares composer draft updates across websocket sessions", () => + Effect.scoped( + Effect.gen(function* () { + yield* buildAppUnderTest(); + + const wsUrl = yield* getWsServerUrl("/ws"); + const threadId = ThreadId.make("shared-composer-draft"); + const subscribed = yield* Deferred.make(); + const snapshotsFiber = yield* withWsRpcClient(wsUrl, (client) => + client[WS_METHODS.subscribeComposerDraft]({ threadId }).pipe( + Stream.tap(() => Deferred.succeed(subscribed, undefined)), + Stream.take(2), + Stream.runCollect, + ), + ).pipe(Effect.forkScoped); + + yield* Deferred.await(subscribed); + const update = yield* withWsRpcClient(wsUrl, (client) => + client[WS_METHODS.composerDraftUpdate]({ + threadId, + baseRevision: 0, + common: { + text: "shared from another websocket", + modelSelection: null, + runtimeMode: null, + interactionMode: null, + }, + clientMutationId: "test:shared-composer-draft", + }), + ); + const snapshots = Array.from(yield* Fiber.join(snapshotsFiber)); + + assert.equal(update._tag, "accepted"); + assert.equal(snapshots[0]?.revision, 0); + assert.deepEqual(snapshots[1], update.snapshot); + }), + ).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + + it.effect("clears the matching composer draft revision after thread.turn.start", () => + Effect.gen(function* () { + yield* buildAppUnderTest(); + + const wsUrl = yield* getWsServerUrl("/ws"); + const threadId = ThreadId.make("sent-composer-draft"); + const commandId = CommandId.make("cmd-send-composer-draft"); + const createdAt = "2026-01-01T00:00:00.000Z"; + const update = yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + client[WS_METHODS.composerDraftUpdate]({ + threadId, + baseRevision: 0, + common: { + text: "send this draft", + modelSelection: defaultModelSelection, + runtimeMode: "full-access", + interactionMode: "default", + }, + clientMutationId: "test:send-composer-draft", + }), + ), + ); + assert.equal(update._tag, "accepted"); + + yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.turn.start", + commandId, + threadId, + message: { + messageId: MessageId.make("msg-send-composer-draft"), + role: "user", + text: "send this draft", + attachments: [], + }, + modelSelection: defaultModelSelection, + runtimeMode: "full-access", + interactionMode: "default", + composerDraftRevision: update.snapshot.revision, + createdAt, + }), + ), + ); + + const cleared = yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + client[WS_METHODS.subscribeComposerDraft]({ threadId }).pipe( + Stream.runHead, + Effect.map(Option.getOrThrow), + ), + ), + ); + assert.equal(cleared.revision, 2); + assert.isNull(cleared.common); + assert.equal(cleared.clientMutationId, `turn:${commandId}`); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + it.effect("rejects websocket rpc handshake when session authentication is missing", () => Effect.gen(function* () { const fs = yield* FileSystem.FileSystem; diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index a00c5cd0570..28dffdecc94 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -22,6 +22,7 @@ import { fixPath } from "./os-jank.ts"; import { websocketRpcRouteLayer } from "./ws.ts"; import * as ExternalLauncher from "./process/externalLauncher.ts"; import { layerConfig as SqlitePersistenceLayerLive } from "./persistence/Layers/Sqlite.ts"; +import * as ComposerDrafts from "./persistence/ComposerDrafts.ts"; import * as ServerLifecycleEvents from "./serverLifecycleEvents.ts"; import * as AnalyticsService from "./telemetry/AnalyticsService.ts"; import { ProviderSessionDirectoryLive } from "./provider/Layers/ProviderSessionDirectory.ts"; @@ -374,7 +375,7 @@ const RuntimeCoreDependenciesLive = ReactorLayerLive.pipe( Layer.provideMerge(GitLayerLive), Layer.provideMerge(VcsLayerLive), Layer.provideMerge(ProviderRuntimeLayerLive), - Layer.provideMerge(Layer.mergeAll(TerminalLayerLive, PreviewLayerLive)), + Layer.provideMerge(Layer.mergeAll(TerminalLayerLive, PreviewLayerLive, ComposerDrafts.layer)), Layer.provideMerge(PersistenceLayerLive), Layer.provideMerge(Keybindings.layer), Layer.provideMerge(ProviderRegistryLive), diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index fda8f1353e3..d3ea319ab64 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -14,6 +14,7 @@ import { AuthAccessStreamError, type AuthAccessStreamEvent, type AuthEnvironmentScope, + ComposerDraftSyncError, AuthSessionId, CommandId, type DiscoveredLocalServerList, @@ -65,6 +66,7 @@ import { RpcSerialization, RpcServer } from "effect/unstable/rpc"; import * as CheckpointDiffQuery from "./checkpointing/CheckpointDiffQuery.ts"; import * as CommandOutputQuery from "./orchestration/CommandOutputQuery.ts"; +import * as ComposerDrafts from "./persistence/ComposerDrafts.ts"; import * as ServerConfig from "./config.ts"; import * as Keybindings from "./keybindings.ts"; import * as ExternalLauncher from "./process/externalLauncher.ts"; @@ -417,6 +419,7 @@ const makeWsRpcLayer = ( const processResourceMonitor = yield* ProcessResourceMonitor.ProcessResourceMonitor; const resourceTelemetry = yield* ResourceTelemetry.ResourceTelemetry; const usage = yield* UsageService.UsageService; + const composerDrafts = yield* ComposerDrafts.ComposerDraftRepository; const relayClient = yield* RelayClient.RelayClient; const authorizationError = (requiredScope: AuthEnvironmentScope) => new EnvironmentAuthorizationError({ @@ -1034,6 +1037,32 @@ const makeWsRpcLayer = ( .pipe(Effect.ignoreCause({ log: true }), Effect.forkDetach, Effect.asVoid); return WsRpcGroup.of({ + [WS_METHODS.composerDraftUpdate]: (input) => + observeRpcEffect( + WS_METHODS.composerDraftUpdate, + composerDrafts + .update(input) + .pipe( + Effect.mapError( + () => + new ComposerDraftSyncError({ message: "Failed to save the composer draft." }), + ), + ), + { threadId: input.threadId }, + ), + [WS_METHODS.subscribeComposerDraft]: (input) => + observeRpcStream( + WS_METHODS.subscribeComposerDraft, + composerDrafts + .subscribe(input) + .pipe( + Stream.mapError( + () => + new ComposerDraftSyncError({ message: "Composer draft sync was interrupted." }), + ), + ), + { threadId: input.threadId }, + ), [ORCHESTRATION_WS_METHODS.dispatchCommand]: (command) => observeRpcEffect( ORCHESTRATION_WS_METHODS.dispatchCommand, @@ -1055,6 +1084,26 @@ const makeWsRpcLayer = ( ) : false; const result = yield* dispatchNormalizedCommand(normalizedCommand); + if ( + normalizedCommand.type === "thread.turn.start" && + normalizedCommand.composerDraftRevision !== undefined + ) { + yield* composerDrafts + .update({ + threadId: normalizedCommand.threadId, + baseRevision: normalizedCommand.composerDraftRevision, + common: null, + clientMutationId: `turn:${normalizedCommand.commandId}`, + }) + .pipe( + Effect.catch((error) => + Effect.logWarning("failed to clear the sent composer draft", { + threadId: normalizedCommand.threadId, + error: error.message, + }), + ), + ); + } if (normalizedCommand.type === "thread.archive") { if (shouldStopSessionAfterArchive) { yield* Effect.gen(function* () { @@ -2178,6 +2227,7 @@ export const websocketRpcRouteLayer = Layer.unwrap( Effect.gen(function* () { const previewAutomationBroker = yield* PreviewAutomationBroker.PreviewAutomationBroker; const serverSelfUpdate = yield* ServerSelfUpdate.ServerSelfUpdate; + const composerDrafts = yield* ComposerDrafts.ComposerDraftRepository; return HttpRouter.add( "GET", "/ws", @@ -2198,6 +2248,7 @@ export const websocketRpcRouteLayer = Layer.unwrap( }).pipe( Effect.provide( makeWsRpcLayer(session, previewAutomationBroker).pipe( + Layer.provide(Layer.succeed(ComposerDrafts.ComposerDraftRepository, composerDrafts)), Layer.provideMerge(RpcSerialization.layerJson), Layer.provide(ProviderMaintenanceRunner.layer), Layer.provide(Layer.succeed(ServerSelfUpdate.ServerSelfUpdate, serverSelfUpdate)), diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index f141069ff5a..d8bd15a1088 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -235,6 +235,7 @@ import { threadHasOlderTurns, } from "@t3tools/client-runtime/state/threads"; import { vcsEnvironment } from "../state/vcs"; +import { markComposerDraftSent, readComposerDraftRevision } from "../state/composerDrafts"; import { useEnvironments, usePrimaryEnvironment } from "../state/environments"; import { useProject, @@ -5006,6 +5007,9 @@ function ChatViewContent(props: ChatViewProps) { ); const messageIdForSend = newMessageId(); const messageCreatedAt = new Date().toISOString(); + const composerDraftRevision = isServerThread + ? readComposerDraftRevision(routeThreadRef) + : undefined; const outgoingMessageText = formatOutgoingPrompt({ provider: ctxSelectedProvider, model: ctxSelectedModel, @@ -5074,6 +5078,7 @@ function ChatViewContent(props: ChatViewProps) { } promptRef.current = ""; clearComposerDraftContent(composerDraftTarget); + if (isServerThread) markComposerDraftSent(routeThreadRef); composerRef.current?.resetCursorState(); let firstComposerImageName: string | null = null; @@ -5187,6 +5192,7 @@ function ChatViewContent(props: ChatViewProps) { titleSeed: title, runtimeMode, interactionMode, + ...(composerDraftRevision === undefined ? {} : { composerDraftRevision }), ...(bootstrap ? { bootstrap } : {}), createdAt: messageCreatedAt, }, diff --git a/apps/web/src/components/chat/ChatComposer.tsx b/apps/web/src/components/chat/ChatComposer.tsx index 9e7a967dc08..344110c0b57 100644 --- a/apps/web/src/components/chat/ChatComposer.tsx +++ b/apps/web/src/components/chat/ChatComposer.tsx @@ -108,6 +108,7 @@ import { buildExpandedImagePreview, type ExpandedImagePreview } from "./Expanded import { basenameOfPath } from "../../pierre-icons"; import { cn, randomUUID } from "~/lib/utils"; import { Separator } from "../ui/separator"; +import { useServerComposerDraftSync } from "../../state/composerDrafts"; type ComposerCommandMenuPosition = { bottom: number; @@ -677,6 +678,7 @@ export const ChatComposer = memo(function ChatComposer(props: ChatComposerProps) // Store subscriptions (prompt / images / terminal contexts) // ------------------------------------------------------------------ const composerDraft = useComposerThreadDraft(composerDraftTarget); + useServerComposerDraftSync(routeKind === "server" ? routeThreadRef : null); const prompt = composerDraft.prompt; const composerImages = composerDraft.images; const composerTerminalContexts = composerDraft.terminalContexts; diff --git a/apps/web/src/composerDraftStore.test.ts b/apps/web/src/composerDraftStore.test.ts index 4425e10676f..b3c10fc916b 100644 --- a/apps/web/src/composerDraftStore.test.ts +++ b/apps/web/src/composerDraftStore.test.ts @@ -172,6 +172,39 @@ function draftByKey(key: string) { return useComposerDraftStore.getState().draftsByThreadKey[key] ?? undefined; } +describe("composerDraftStore shared draft section", () => { + const threadRef = scopeThreadRef(TEST_ENVIRONMENT_ID, ThreadId.make("thread-synced-common")); + + beforeEach(() => { + resetComposerDraftStore(); + }); + + it("applies remote common state without replacing local attachments", () => { + const image = makeImage({ id: "local-image", previewUrl: "blob:local" }); + const store = useComposerDraftStore.getState(); + store.setPrompt(threadRef, "local"); + store.addImage(threadRef, image); + + store.applySyncedCommon(threadRef, { + text: "remote", + modelSelection: { + instanceId: CODEX_INSTANCE, + model: "gpt-5.4", + }, + runtimeMode: "full-access", + interactionMode: "plan", + }); + + expect(draftByKey(scopedThreadKey(threadRef))).toMatchObject({ + prompt: "remote", + images: [image], + activeProvider: CODEX_INSTANCE, + runtimeMode: "full-access", + interactionMode: "plan", + }); + }); +}); + describe("composerDraftStore addImages", () => { const threadId = ThreadId.make("thread-dedupe"); const threadRef = scopeThreadRef(TEST_ENVIRONMENT_ID, threadId); diff --git a/apps/web/src/composerDraftStore.ts b/apps/web/src/composerDraftStore.ts index 517b642869b..e340ea95d2d 100644 --- a/apps/web/src/composerDraftStore.ts +++ b/apps/web/src/composerDraftStore.ts @@ -1,4 +1,5 @@ import { + type ComposerDraftCommon, DEFAULT_MODEL, DEFAULT_MODEL_BY_PROVIDER, defaultInstanceIdForDriver, @@ -439,6 +440,8 @@ interface ComposerDraftStoreState { clearDraftThread: (threadRef: ComposerThreadTarget) => void; setStickyModelSelection: (modelSelection: ModelSelection | null | undefined) => void; setPrompt: (threadRef: ComposerThreadTarget, prompt: string) => void; + /** Replaces only the cross-device section and preserves local attachments/context. */ + applySyncedCommon: (threadRef: ComposerThreadTarget, common: ComposerDraftCommon | null) => void; setTerminalContexts: (threadRef: ComposerThreadTarget, contexts: TerminalContextDraft[]) => void; setModelSelection: ( threadRef: ComposerThreadTarget, @@ -2725,6 +2728,34 @@ const composerDraftStore = create()( return { draftsByThreadKey: nextDraftsByThreadKey }; }); }, + applySyncedCommon: (threadRef, common) => { + const threadKey = resolveComposerDraftKey(get(), threadRef) ?? ""; + if (threadKey.length === 0) { + return; + } + set((state) => { + const existing = state.draftsByThreadKey[threadKey] ?? createEmptyThreadDraft(); + const nextMap = { ...existing.modelSelectionByProvider }; + if (common?.modelSelection) { + nextMap[common.modelSelection.instanceId] = common.modelSelection; + } + const nextDraft: ComposerThreadDraftState = { + ...existing, + prompt: common?.text ?? "", + modelSelectionByProvider: nextMap, + activeProvider: common?.modelSelection?.instanceId ?? null, + runtimeMode: common?.runtimeMode ?? null, + interactionMode: common?.interactionMode ?? null, + }; + const nextDraftsByThreadKey = { ...state.draftsByThreadKey }; + if (shouldRemoveDraft(nextDraft)) { + delete nextDraftsByThreadKey[threadKey]; + } else { + nextDraftsByThreadKey[threadKey] = nextDraft; + } + return { draftsByThreadKey: nextDraftsByThreadKey }; + }); + }, setTerminalContexts: (threadRef, contexts) => { const threadKey = resolveComposerDraftKey(get(), threadRef); const threadId = resolveComposerThreadId(get(), threadRef); diff --git a/apps/web/src/state/composerDrafts.ts b/apps/web/src/state/composerDrafts.ts new file mode 100644 index 00000000000..8536b01edbc --- /dev/null +++ b/apps/web/src/state/composerDrafts.ts @@ -0,0 +1,176 @@ +import { useAtomValue } from "@effect/atom-react"; +import { + canonicalComposerDraftCommon, + composerDraftCommonEquals, + createComposerDraftEnvironmentAtoms, + createComposerDraftSyncController, + type ComposerDraftSyncController, +} from "@t3tools/client-runtime/state/composer-drafts"; +import { scopedThreadKey } from "@t3tools/client-runtime/environment"; +import type { + ComposerDraftCommon, + ComposerDraftSnapshot, + ScopedThreadRef, +} from "@t3tools/contracts"; +import { AsyncResult, Atom } from "effect/unstable/reactivity"; +import { useEffect, useRef } from "react"; + +import { + type ComposerThreadDraftState, + DraftId, + useComposerThreadDraft, + useComposerDraftStore, +} from "../composerDraftStore"; +import { connectionAtomRuntime } from "../connection/runtime"; +import { appAtomRegistry } from "../rpc/atomRegistry"; +import { serverEnvironment } from "./server"; +import { randomUUID } from "../lib/utils"; + +export const composerDraftEnvironment = createComposerDraftEnvironmentAtoms(connectionAtomRuntime); + +const EMPTY_SYNC_ATOM = Atom.make(null).pipe(Atom.withLabel("composer-draft-sync:disabled")); +const revisions = new Map(); +const suppressedPostSendCommon = new Map(); + +export function readComposerDraftRevision(threadRef: ScopedThreadRef): number | undefined { + return revisions.get(scopedThreadKey(threadRef)); +} + +function commonFromDraft(draft: ComposerThreadDraftState): ComposerDraftCommon | null { + const hasLocalOnlyContext = + draft.images.length > 0 || + draft.persistedAttachments.length > 0 || + draft.terminalContexts.length > 0 || + draft.elementContexts.length > 0 || + draft.issueContexts.length > 0 || + draft.previewAnnotations.length > 0 || + draft.reviewComments.length > 0; + if (hasLocalOnlyContext) return null; + return canonicalComposerDraftCommon({ + text: draft.prompt, + modelSelection: + draft.activeProvider === null + ? null + : (draft.modelSelectionByProvider[draft.activeProvider] ?? null), + runtimeMode: draft.runtimeMode, + interactionMode: draft.interactionMode, + }); +} + +/** Prevents retained model/mode preferences from resurrecting a sent draft. */ +export function markComposerDraftSent(threadRef: ScopedThreadRef): void { + const key = scopedThreadKey(threadRef); + const draft = useComposerDraftStore.getState().getComposerDraft(threadRef); + suppressedPostSendCommon.set(key, draft === null ? null : commonFromDraft(draft)); +} + +function mutationId(): string { + return `web:${randomUUID()}`; +} + +export function useServerComposerDraftSync(threadRef: ScopedThreadRef | null): void { + const serverConfig = useAtomValue( + threadRef === null + ? EMPTY_SYNC_ATOM + : serverEnvironment.configValueAtom(threadRef.environmentId), + ); + const enabled = + threadRef !== null && + serverConfig !== null && + "environment" in serverConfig && + serverConfig.environment.capabilities.composerDraftSync === true; + const streamResult = useAtomValue( + enabled && threadRef !== null + ? composerDraftEnvironment.changes({ + environmentId: threadRef.environmentId, + input: { threadId: threadRef.threadId }, + }) + : EMPTY_SYNC_ATOM, + ); + const draft = useComposerThreadDraft( + threadRef ?? DraftId.make("__composer-draft-sync-disabled__"), + ); + const draftRef = useRef(draft); + draftRef.current = draft; + const controllerRef = useRef(null); + + useEffect(() => { + controllerRef.current?.dispose(); + if (!enabled || threadRef === null) { + controllerRef.current = null; + return; + } + const key = scopedThreadKey(threadRef); + const readLocal = (): ComposerDraftCommon | null => { + const common = commonFromDraft(draftRef.current); + if (!suppressedPostSendCommon.has(key)) return common; + const baseline = suppressedPostSendCommon.get(key) ?? null; + if (composerDraftCommonEquals(common, baseline)) return null; + suppressedPostSendCommon.delete(key); + return common; + }; + const controller = createComposerDraftSyncController({ + threadId: threadRef.threadId, + readLocal, + canApplyRemote: () => { + const current = draftRef.current; + return ( + current.images.length === 0 && + current.persistedAttachments.length === 0 && + current.terminalContexts.length === 0 && + current.elementContexts.length === 0 && + current.issueContexts.length === 0 && + current.previewAnnotations.length === 0 && + current.reviewComments.length === 0 + ); + }, + applyRemote: (common) => + useComposerDraftStore.getState().applySyncedCommon(threadRef, common), + update: async (input) => { + const result = await composerDraftEnvironment.update.run(appAtomRegistry, { + environmentId: threadRef.environmentId, + input, + }); + return AsyncResult.isSuccess(result) ? result.value : null; + }, + createMutationId: mutationId, + scheduleTask: (task, delayMs) => { + const timer = window.setTimeout(task, delayMs); + return () => window.clearTimeout(timer); + }, + onRevisionChange: (snapshot: ComposerDraftSnapshot) => { + revisions.set(key, snapshot.revision); + }, + }); + controllerRef.current = controller; + return () => { + controller.dispose(); + if (controllerRef.current === controller) controllerRef.current = null; + revisions.delete(key); + suppressedPostSendCommon.delete(key); + }; + }, [enabled, threadRef?.environmentId, threadRef?.threadId]); + + useEffect(() => { + if (streamResult !== null && AsyncResult.isSuccess(streamResult)) { + controllerRef.current?.observeSnapshot(streamResult.value); + } + }, [streamResult]); + + useEffect(() => { + controllerRef.current?.observeLocalChange(); + }, [ + draft.activeProvider, + draft.elementContexts, + draft.images, + draft.interactionMode, + draft.issueContexts, + draft.modelSelectionByProvider, + draft.persistedAttachments, + draft.previewAnnotations, + draft.prompt, + draft.reviewComments, + draft.runtimeMode, + draft.terminalContexts, + ]); +} diff --git a/docs/internals/composer-draft-sync.md b/docs/internals/composer-draft-sync.md new file mode 100644 index 00000000000..63dcd8a814c --- /dev/null +++ b/docs/internals/composer-draft-sync.md @@ -0,0 +1,21 @@ +# Composer draft synchronization + +Existing-thread composer drafts are high-churn current state, not orchestration history. The server +stores one JSON `common` section per thread with a monotonic revision and mutation ID. A null +`common` value is a durable tombstone, which prevents a stale revision-zero client from resurrecting +a sent draft. + +Web clients keep their existing local durable cache and subscribe to the server snapshot. Desktop +uses the same web implementation. Updates use revision compare-and-swap. A clean client applies +newer server state; an actively edited client waits for the idle debounce and retries once against a +conflicting revision. On first contact, an existing non-empty local cache is preserved rather than +automatically overwriting an established server draft. Mobile does not participate in synchronization +in this version and retains its existing device-local draft and outbox behavior. + +The shared section contains text, model selection, runtime mode, and interaction mode. Attachment +bytes and surface-specific context never cross this channel. A client projects a draft containing +local-only context to a tombstone and refuses to apply remote state until that context is gone. + +A turn-start command may carry the composer revision captured at send time. After the turn command +is durably accepted, the server conditionally writes a tombstone at that revision. This prevents a +delayed web or desktop send from erasing a newer edit from another client. diff --git a/docs/user/composer-drafts.md b/docs/user/composer-drafts.md new file mode 100644 index 00000000000..308788641e8 --- /dev/null +++ b/docs/user/composer-drafts.md @@ -0,0 +1,17 @@ +# Composer drafts + +Draft text in an existing thread follows you between web and desktop clients connected to the same +T3 Code environment. The selected model, runtime mode, and interaction mode travel with the text. +Changes sync after a short typing pause and converge after an offline client reconnects. Mobile +drafts remain device-local and do not participate in this synchronization. + +Drafts for a new task remain on the device until the thread is created. Images, terminal excerpts, +preview selections, review comments, and other device-specific context also remain local. While an +existing-thread draft contains any of that context, T3 Code withholds the whole draft from other +devices so they cannot send an incomplete version of it. + +Sending clears only the server revision that was visible when Send was pressed. If another device +has already changed the draft, that newer revision is preserved. + +Drafts are stored by the T3 Code server for that environment. They are not shared between separate +servers, even if both servers contain a thread with the same name. diff --git a/packages/client-runtime/package.json b/packages/client-runtime/package.json index d1daa871652..3b35616e184 100644 --- a/packages/client-runtime/package.json +++ b/packages/client-runtime/package.json @@ -51,6 +51,10 @@ "types": "./src/state/connections.ts", "default": "./src/state/connections.ts" }, + "./state/composer-drafts": { + "types": "./src/state/composerDrafts.ts", + "default": "./src/state/composerDrafts.ts" + }, "./state/entities": { "types": "./src/state/entities.ts", "default": "./src/state/entities.ts" diff --git a/packages/client-runtime/src/rpc/client.ts b/packages/client-runtime/src/rpc/client.ts index bfe57a6c0dd..50cd9447d27 100644 --- a/packages/client-runtime/src/rpc/client.ts +++ b/packages/client-runtime/src/rpc/client.ts @@ -52,6 +52,7 @@ export type EnvironmentSubscriptionRpcTag = | typeof WS_METHODS.subscribeResourceTelemetry | typeof WS_METHODS.previewAutomationConnect | typeof WS_METHODS.subscribeVcsStatus + | typeof WS_METHODS.subscribeComposerDraft | typeof WS_METHODS.terminalAttach; export type EnvironmentStreamCommandRpcTag = diff --git a/packages/client-runtime/src/state/composerDrafts.test.ts b/packages/client-runtime/src/state/composerDrafts.test.ts new file mode 100644 index 00000000000..c7d1c6bfb76 --- /dev/null +++ b/packages/client-runtime/src/state/composerDrafts.test.ts @@ -0,0 +1,162 @@ +import { ThreadId, type ComposerDraftCommon, type ComposerDraftSnapshot } from "@t3tools/contracts"; +import { describe, expect, it } from "@effect/vitest"; + +import { createComposerDraftSyncController } from "./composerDrafts.ts"; + +const THREAD_ID = ThreadId.make("thread-1"); +const REMOTE: ComposerDraftCommon = { + text: "from server", + modelSelection: null, + runtimeMode: null, + interactionMode: null, +}; +const LOCAL: ComposerDraftCommon = { + text: "local edit", + modelSelection: null, + runtimeMode: null, + interactionMode: null, +}; + +function snapshot( + revision: number, + common: ComposerDraftCommon | null, + clientMutationId = "server-change", +): ComposerDraftSnapshot { + return { + threadId: THREAD_ID, + revision, + common, + updatedAt: "2026-08-08T12:00:00.000Z", + clientMutationId, + }; +} + +function makeScheduler() { + let scheduled: (() => void) | null = null; + return { + scheduleTask: (task: () => void) => { + scheduled = task; + return () => { + if (scheduled === task) scheduled = null; + }; + }, + run: async () => { + const task = scheduled; + scheduled = null; + task?.(); + await Promise.resolve(); + await Promise.resolve(); + }, + hasTask: () => scheduled !== null, + }; +} + +describe("composer draft sync controller", () => { + it("applies an authoritative remote draft when the local cache is empty", () => { + let local: ComposerDraftCommon | null = null; + const scheduler = makeScheduler(); + const controller = createComposerDraftSyncController({ + threadId: THREAD_ID, + readLocal: () => local, + canApplyRemote: () => true, + applyRemote: (common) => { + local = common; + }, + update: async () => null, + createMutationId: () => "mutation-1", + scheduleTask: scheduler.scheduleTask, + }); + + controller.observeSnapshot(snapshot(1, REMOTE)); + + expect(local).toEqual(REMOTE); + expect(controller.revision()).toBe(1); + expect(scheduler.hasTask()).toBe(false); + }); + + it("does not overwrite a non-empty local cache on first contact", async () => { + let local: ComposerDraftCommon | null = LOCAL; + const writes: Array<{ baseRevision: number; common: ComposerDraftCommon | null }> = []; + const scheduler = makeScheduler(); + const controller = createComposerDraftSyncController({ + threadId: THREAD_ID, + readLocal: () => local, + canApplyRemote: () => true, + applyRemote: (common) => { + local = common; + }, + update: async (input) => { + writes.push({ baseRevision: input.baseRevision, common: input.common }); + return { _tag: "accepted", snapshot: snapshot(2, input.common, input.clientMutationId) }; + }, + createMutationId: () => "mutation-1", + scheduleTask: scheduler.scheduleTask, + }); + + controller.observeSnapshot(snapshot(1, REMOTE)); + expect(local).toEqual(LOCAL); + expect(scheduler.hasTask()).toBe(false); + + local = { ...LOCAL, text: "actively edited" }; + controller.observeLocalChange(); + await scheduler.run(); + + expect(writes).toEqual([{ baseRevision: 1, common: local }]); + }); + + it("retries an active local edit against the revision returned by a conflict", async () => { + let local: ComposerDraftCommon | null = null; + const bases: number[] = []; + const scheduler = makeScheduler(); + const controller = createComposerDraftSyncController({ + threadId: THREAD_ID, + readLocal: () => local, + canApplyRemote: () => true, + applyRemote: (common) => { + local = common; + }, + update: async (input) => { + bases.push(input.baseRevision); + return bases.length === 1 + ? { _tag: "conflict", snapshot: snapshot(2, REMOTE) } + : { _tag: "accepted", snapshot: snapshot(3, input.common, input.clientMutationId) }; + }, + createMutationId: () => `mutation-${bases.length + 1}`, + scheduleTask: scheduler.scheduleTask, + }); + + controller.observeSnapshot(snapshot(0, null)); + local = LOCAL; + controller.observeLocalChange(); + await scheduler.run(); + await scheduler.run(); + + expect(bases).toEqual([0, 2]); + expect(controller.revision()).toBe(3); + expect(local).toEqual(LOCAL); + }); + + it("tombstones the transferable copy while device-only context is present", async () => { + const writes: Array = []; + const scheduler = makeScheduler(); + const controller = createComposerDraftSyncController({ + threadId: THREAD_ID, + readLocal: () => null, + canApplyRemote: () => false, + applyRemote: () => { + throw new Error("remote state must not replace device-only context"); + }, + update: async (input) => { + writes.push(input.common); + return { _tag: "accepted", snapshot: snapshot(2, null, input.clientMutationId) }; + }, + createMutationId: () => "context-tombstone", + scheduleTask: scheduler.scheduleTask, + }); + + controller.observeSnapshot(snapshot(1, REMOTE)); + await scheduler.run(); + + expect(writes).toEqual([null]); + }); +}); diff --git a/packages/client-runtime/src/state/composerDrafts.ts b/packages/client-runtime/src/state/composerDrafts.ts new file mode 100644 index 00000000000..04a9af2fe5c --- /dev/null +++ b/packages/client-runtime/src/state/composerDrafts.ts @@ -0,0 +1,221 @@ +import { + type ComposerDraftCommon, + type ComposerDraftSnapshot, + type ComposerDraftUpdateResult, + type ThreadId, + WS_METHODS, +} from "@t3tools/contracts"; +import { Atom } from "effect/unstable/reactivity"; + +import type { EnvironmentRegistry } from "../connection/registry.ts"; +import { + createAtomCommandScheduler, + createEnvironmentRpcCommand, + createEnvironmentRpcSubscriptionAtomFamily, +} from "./runtime.ts"; + +/** Typed RPC primitives used by the web composer adapter and desktop wrapper. */ +export function createComposerDraftEnvironmentAtoms( + runtime: Atom.AtomRuntime, +) { + const scheduler = createAtomCommandScheduler(); + return { + changes: createEnvironmentRpcSubscriptionAtomFamily(runtime, { + label: "environment-data:composer-draft:changes", + tag: WS_METHODS.subscribeComposerDraft, + idleTtlMs: 1_000, + }), + update: createEnvironmentRpcCommand(runtime, { + label: "environment-data:composer-draft:update", + tag: WS_METHODS.composerDraftUpdate, + scheduler, + concurrency: { + mode: "latest", + key: ({ environmentId, input }) => JSON.stringify([environmentId, input.threadId]), + }, + }), + }; +} + +export function canonicalComposerDraftCommon( + common: ComposerDraftCommon | null, +): ComposerDraftCommon | null { + if ( + common === null || + (common.text.length === 0 && + common.modelSelection === null && + common.runtimeMode === null && + common.interactionMode === null) + ) { + return null; + } + return common; +} + +export function composerDraftCommonEquals( + left: ComposerDraftCommon | null, + right: ComposerDraftCommon | null, +): boolean { + return ( + JSON.stringify(canonicalComposerDraftCommon(left)) === + JSON.stringify(canonicalComposerDraftCommon(right)) + ); +} + +export interface ComposerDraftSyncController { + readonly observeSnapshot: (snapshot: ComposerDraftSnapshot) => void; + readonly observeLocalChange: () => void; + readonly revision: () => number; + readonly dispose: () => void; +} + +/** + * Reconciles one existing thread's local cache with its revisioned server value. + * The surface owns persistence and rendering; this controller only defines the + * conflict and debounce behavior shared by web and desktop. + */ +export function createComposerDraftSyncController(options: { + readonly threadId: ThreadId; + readonly readLocal: () => ComposerDraftCommon | null; + readonly canApplyRemote: () => boolean; + readonly applyRemote: (common: ComposerDraftCommon | null) => void; + readonly update: (input: { + readonly threadId: ThreadId; + readonly baseRevision: number; + readonly common: ComposerDraftCommon | null; + readonly clientMutationId: string; + }) => Promise; + readonly createMutationId: () => string; + readonly scheduleTask: (task: () => void, delayMs: number) => () => void; + readonly onRevisionChange?: (snapshot: ComposerDraftSnapshot) => void; + readonly debounceMs?: number; +}): ComposerDraftSyncController { + const debounceMs = options.debounceMs ?? 1_200; + let disposed = false; + let initialized = false; + let currentRevision = 0; + let lastSynced: ComposerDraftCommon | null = null; + let cancelScheduledTask: (() => void) | null = null; + let inFlight = false; + let pendingAfterFlight = false; + + const cancelTimer = () => { + if (cancelScheduledTask !== null) { + cancelScheduledTask(); + cancelScheduledTask = null; + } + }; + + const schedule = (delay = debounceMs) => { + if (disposed || !initialized) return; + cancelTimer(); + cancelScheduledTask = options.scheduleTask(() => { + cancelScheduledTask = null; + void flush(); + }, delay); + }; + + const acceptSnapshotMetadata = (snapshot: ComposerDraftSnapshot) => { + currentRevision = snapshot.revision; + lastSynced = canonicalComposerDraftCommon(snapshot.common); + options.onRevisionChange?.(snapshot); + }; + + const flush = async () => { + if (disposed || !initialized) return; + if (inFlight) { + pendingAfterFlight = true; + return; + } + const sent = canonicalComposerDraftCommon(options.readLocal()); + if (composerDraftCommonEquals(sent, lastSynced)) return; + + const baseRevision = currentRevision; + const clientMutationId = options.createMutationId(); + inFlight = true; + const result = await options.update({ + threadId: options.threadId, + baseRevision, + common: sent, + clientMutationId, + }); + inFlight = false; + if (disposed) return; + + if (result === null) { + // The subscription/reconnect path will call observeLocalChange again. + return; + } + + acceptSnapshotMetadata(result.snapshot); + const localNow = canonicalComposerDraftCommon(options.readLocal()); + if (result._tag === "conflict") { + // An actively edited local value wins by retrying against the new base. + if (composerDraftCommonEquals(localNow, sent)) schedule(0); + else schedule(); + } else if (!composerDraftCommonEquals(localNow, lastSynced)) { + schedule(); + } + + if (pendingAfterFlight) { + pendingAfterFlight = false; + if (!composerDraftCommonEquals(options.readLocal(), lastSynced)) schedule(); + } + }; + + const observeSnapshot = (snapshot: ComposerDraftSnapshot) => { + if (disposed || snapshot.threadId !== options.threadId || snapshot.revision < currentRevision) { + return; + } + const remote = canonicalComposerDraftCommon(snapshot.common); + const local = canonicalComposerDraftCommon(options.readLocal()); + + if (!initialized) { + initialized = true; + acceptSnapshotMetadata(snapshot); + if (snapshot.revision === 0) { + if (!composerDraftCommonEquals(local, remote)) schedule(); + return; + } + if (!options.canApplyRemote()) { + // Hide a transferable draft as soon as this web/desktop client has + // richer local context; another client must not send an incomplete copy. + if (remote !== null) schedule(); + return; + } + if (local === null || composerDraftCommonEquals(local, remote)) { + if (!composerDraftCommonEquals(local, remote)) options.applyRemote(remote); + } + // Preserve a non-empty local cache on first contact. It is uploaded only + // after the user edits it, avoiding an automatic migration-time overwrite. + return; + } + + if (snapshot.revision === currentRevision && composerDraftCommonEquals(remote, lastSynced)) { + return; + } + + const wasClean = composerDraftCommonEquals(local, lastSynced); + acceptSnapshotMetadata(snapshot); + if (wasClean && options.canApplyRemote()) { + options.applyRemote(remote); + return; + } + // Local typing (or a client-local attachment) wins after the idle window. + schedule(); + }; + + return { + observeSnapshot, + observeLocalChange: () => { + if (!initialized || disposed) return; + if (composerDraftCommonEquals(options.readLocal(), lastSynced)) cancelTimer(); + else schedule(); + }, + revision: () => currentRevision, + dispose: () => { + disposed = true; + cancelTimer(); + }, + }; +} diff --git a/packages/contracts/src/composerDraft.ts b/packages/contracts/src/composerDraft.ts new file mode 100644 index 00000000000..b656d10fa4f --- /dev/null +++ b/packages/contracts/src/composerDraft.ts @@ -0,0 +1,51 @@ +import * as Schema from "effect/Schema"; + +import { IsoDateTime, NonNegativeInt, ThreadId, TrimmedNonEmptyString } from "./baseSchemas.ts"; +import { ModelSelection, ProviderInteractionMode, RuntimeMode } from "./orchestration.ts"; + +/** + * The web/desktop transferable portion of an existing thread's composer draft. + * + * Attachments and surface-specific context intentionally stay out of v1. The + * server replaces this section as one revisioned value so web/desktop clients + * never erase fields owned by local-only composer state. + */ +export const ComposerDraftCommon = Schema.Struct({ + text: Schema.String, + modelSelection: Schema.NullOr(ModelSelection), + runtimeMode: Schema.NullOr(RuntimeMode), + interactionMode: Schema.NullOr(ProviderInteractionMode), +}); +export type ComposerDraftCommon = typeof ComposerDraftCommon.Type; + +/** A null common section is a durable tombstone, not an unknown draft. */ +export const ComposerDraftSnapshot = Schema.Struct({ + threadId: ThreadId, + revision: NonNegativeInt, + common: Schema.NullOr(ComposerDraftCommon), + updatedAt: Schema.NullOr(IsoDateTime), + clientMutationId: Schema.NullOr(TrimmedNonEmptyString), +}); +export type ComposerDraftSnapshot = typeof ComposerDraftSnapshot.Type; + +export const ComposerDraftGetInput = Schema.Struct({ threadId: ThreadId }); +export type ComposerDraftGetInput = typeof ComposerDraftGetInput.Type; + +export const ComposerDraftUpdateInput = Schema.Struct({ + threadId: ThreadId, + baseRevision: NonNegativeInt, + common: Schema.NullOr(ComposerDraftCommon), + clientMutationId: TrimmedNonEmptyString, +}); +export type ComposerDraftUpdateInput = typeof ComposerDraftUpdateInput.Type; + +export const ComposerDraftUpdateResult = Schema.Union([ + Schema.TaggedStruct("accepted", { snapshot: ComposerDraftSnapshot }), + Schema.TaggedStruct("conflict", { snapshot: ComposerDraftSnapshot }), +]); +export type ComposerDraftUpdateResult = typeof ComposerDraftUpdateResult.Type; + +export class ComposerDraftSyncError extends Schema.TaggedErrorClass()( + "ComposerDraftSyncError", + { message: Schema.String }, +) {} diff --git a/packages/contracts/src/environment.ts b/packages/contracts/src/environment.ts index 329ff911503..ec080640528 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -56,6 +56,8 @@ export const ExecutionEnvironmentCapabilities = Schema.Struct({ /** Server understands regenerateTitle on thread.meta.update. Absent on older servers, so clients hide the action instead of sending it. */ threadTitleRegeneration: Schema.optionalKey(Schema.Boolean), + /** Server can persist and stream revisioned web/desktop composer drafts. */ + composerDraftSync: Schema.optionalKey(Schema.Boolean), /** The update path clients should offer for this server. Absent on servers that must be relaunched manually (dev checkouts, Windows foreground runs, pre-update servers). */ diff --git a/packages/contracts/src/index.ts b/packages/contracts/src/index.ts index 6181391eca3..eb03012336f 100644 --- a/packages/contracts/src/index.ts +++ b/packages/contracts/src/index.ts @@ -1,5 +1,6 @@ export * from "./baseSchemas.ts"; export * from "./background.ts"; +export * from "./composerDraft.ts"; export * from "./auth.ts"; export * from "./environment.ts"; export * from "./environmentHttp.ts"; diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 3b1f66109e8..1efbbdca2af 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -810,6 +810,8 @@ export const ThreadTurnStartCommand = Schema.Struct({ ), bootstrap: Schema.optional(ThreadTurnStartBootstrap), sourceProposedPlan: Schema.optional(SourceProposedPlanReference), + /** Clears this composer revision only if it is still current after dispatch. */ + composerDraftRevision: Schema.optional(NonNegativeInt), createdAt: IsoDateTime, }); @@ -829,6 +831,7 @@ const ClientThreadTurnStartCommand = Schema.Struct({ interactionMode: ProviderInteractionMode, bootstrap: Schema.optional(ThreadTurnStartBootstrap), sourceProposedPlan: Schema.optional(SourceProposedPlanReference), + composerDraftRevision: Schema.optional(NonNegativeInt), createdAt: IsoDateTime, }); diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index 801052a177f..c5b4abf0eea 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -172,6 +172,13 @@ import { SourceControlRepositoryLookupInput, } from "./sourceControl.ts"; import { VcsError } from "./vcs.ts"; +import { + ComposerDraftGetInput, + ComposerDraftSnapshot, + ComposerDraftSyncError, + ComposerDraftUpdateInput, + ComposerDraftUpdateResult, +} from "./composerDraft.ts"; export const WS_METHODS = { // Project registry methods @@ -191,6 +198,9 @@ export const WS_METHODS = { filesystemBrowse: "filesystem.browse", assetsCreateUrl: "assets.createUrl", + // Existing-thread composer draft synchronization + composerDraftUpdate: "composerDraft.update", + // VCS methods vcsPull: "vcs.pull", vcsRefreshStatus: "vcs.refreshStatus", @@ -276,8 +286,22 @@ export const WS_METHODS = { subscribeAuthAccess: "subscribeAuthAccess", subscribeBackgroundPolicy: "subscribeBackgroundPolicy", subscribeResourceTelemetry: "subscribeResourceTelemetry", + subscribeComposerDraft: "subscribeComposerDraft", } as const; +export const WsComposerDraftUpdateRpc = Rpc.make(WS_METHODS.composerDraftUpdate, { + payload: ComposerDraftUpdateInput, + success: ComposerDraftUpdateResult, + error: Schema.Union([ComposerDraftSyncError, EnvironmentAuthorizationError]), +}); + +export const WsSubscribeComposerDraftRpc = Rpc.make(WS_METHODS.subscribeComposerDraft, { + payload: ComposerDraftGetInput, + success: ComposerDraftSnapshot, + error: Schema.Union([ComposerDraftSyncError, EnvironmentAuthorizationError]), + stream: true, +}); + export const WsServerUpsertKeybindingRpc = Rpc.make(WS_METHODS.serverUpsertKeybinding, { payload: ServerUpsertKeybindingInput, success: ServerUpsertKeybindingResult, @@ -877,6 +901,7 @@ export const WsRpcGroup = RpcGroup.make( WsShellOpenInEditorRpc, WsFilesystemBrowseRpc, WsAssetsCreateUrlRpc, + WsComposerDraftUpdateRpc, WsSubscribeVcsStatusRpc, WsVcsPullRpc, WsVcsRefreshStatusRpc, @@ -917,6 +942,7 @@ export const WsRpcGroup = RpcGroup.make( WsSubscribeAuthAccessRpc, WsSubscribeBackgroundPolicyRpc, WsSubscribeResourceTelemetryRpc, + WsSubscribeComposerDraftRpc, WsOrchestrationDispatchCommandRpc, WsOrchestrationGetWorkflowScriptRpc, WsOrchestrationGetCommandOutputRpc,