Compare commits

...

3 Commits

Author SHA1 Message Date
Dax Raad 1036e82c6b fix(cli): defer update check until service resolves 2026-08-12 12:14:01 -04:00
Dax Raad fcc8856f36 fix(tui): narrow pending delivery updates 2026-08-12 09:30:27 -04:00
Dax Raad 984e62502b feat(core): defer active session moves 2026-08-12 09:27:31 -04:00
22 changed files with 472 additions and 75 deletions
@@ -15,8 +15,6 @@ export default Runtime.handler(Commands, (input) =>
Effect.gen(function* () { Effect.gen(function* () {
const requestedDirectory = Option.getOrUndefined(input.directory) const requestedDirectory = Option.getOrUndefined(input.directory)
if (requestedDirectory !== undefined) process.chdir(requestedDirectory) if (requestedDirectory !== undefined) process.chdir(requestedDirectory)
const updater = yield* Updater.Service
yield* updater.check().pipe(Effect.forkScoped)
const preflight = UpdatePreflight.make() const preflight = UpdatePreflight.make()
yield* Effect.addFinalizer(() => Effect.promise(() => preflight.close())) yield* Effect.addFinalizer(() => Effect.promise(() => preflight.close()))
const server = yield* ServerConnection.resolve({ const server = yield* ServerConnection.resolve({
@@ -36,6 +34,8 @@ export default Runtime.handler(Commands, (input) =>
Effect.promise(() => preflight.fail("OpenCode update could not start the new background service")), Effect.promise(() => preflight.fail("OpenCode update could not start the new background service")),
), ),
) )
const updater = yield* Updater.Service
yield* updater.check().pipe(Effect.forkScoped)
preflight.loading() preflight.loading()
const config = yield* Config.Service const config = yield* Config.Service
const npm = yield* Npm.Service const npm = yield* Npm.Service
+10
View File
@@ -358,6 +358,15 @@ export type Endpoint5_31Output =
readonly previous?: Model.Ref | undefined readonly previous?: Model.Ref | undefined
} }
} }
| {
readonly id: Event.ID
readonly created: DateTime.Utc
readonly metadata?: { readonly [x: string]: unknown } | undefined
readonly type: "session.move.admitted"
readonly durable: { readonly aggregateID: string; readonly seq: Event.Seq; readonly version: Event.Version }
readonly location?: Location.Ref | undefined
readonly data: { readonly sessionID: Session.ID; readonly move: SessionPending.MoveData }
}
| { | {
readonly id: Event.ID readonly id: Event.ID
readonly created: DateTime.Utc readonly created: DateTime.Utc
@@ -367,6 +376,7 @@ export type Endpoint5_31Output =
readonly location?: Location.Ref | undefined readonly location?: Location.Ref | undefined
readonly data: { readonly data: {
readonly sessionID: Session.ID readonly sessionID: Session.ID
readonly moveID?: Event.ID | undefined
readonly location: Location.Ref readonly location: Location.Ref
readonly projectID?: Project.ID | undefined readonly projectID?: Project.ID | undefined
readonly subpath?: RelativePath | undefined readonly subpath?: RelativePath | undefined
+28 -2
View File
@@ -420,6 +420,8 @@ export type SessionMessageLocationSwitched = {
previous?: { location: LocationRef; projectID?: string; subpath?: string } previous?: { location: LocationRef; projectID?: string; subpath?: string }
} }
export type SessionPendingMoveData = { location: LocationRef; projectID: string; subpath?: string }
export type SessionCreated = { export type SessionCreated = {
id: string id: string
created: number created: number
@@ -468,7 +470,7 @@ export type SessionMoved = {
type: "session.moved" type: "session.moved"
durable: { aggregateID: string; seq: number; version: 1 } durable: { aggregateID: string; seq: number; version: 1 }
location?: LocationRef location?: LocationRef
data: { sessionID: string; location: LocationRef; projectID?: string; subpath?: string } data: { sessionID: string; moveID?: string; location: LocationRef; projectID?: string; subpath?: string }
} }
export type SessionRenamed = { export type SessionRenamed = {
@@ -1527,6 +1529,24 @@ export type VcsInfo = { branch: VcsBranch }
export type PermissionRuleset = Array<PermissionRule> export type PermissionRuleset = Array<PermissionRule>
export type SessionPendingMove = {
id: string
sessionID: string
timeCreated: number
type: "move"
data: SessionPendingMoveData
}
export type SessionMoveAdmitted = {
id: string
created: number
metadata?: { [x: string]: any }
type: "session.move.admitted"
durable: { aggregateID: string; seq: number; version: 1 }
location?: LocationRef
data: { sessionID: string; move: SessionPendingMoveData }
}
export type SessionInfo = { export type SessionInfo = {
id: string id: string
parentID?: string parentID?: string
@@ -1914,7 +1934,11 @@ export type FormFields = [FormField, ...Array<FormField>]
export type FormFields3 = [FormField1, ...Array<FormField1>] export type FormFields3 = [FormField1, ...Array<FormField1>]
export type SessionPendingInfo = SessionPendingUser | SessionPendingSynthetic | SessionPendingCompaction export type SessionPendingInfo =
| SessionPendingUser
| SessionPendingSynthetic
| SessionPendingCompaction
| SessionPendingMove
export type SessionPendingMessage = SessionPendingUserMessage | SessionPendingSyntheticMessage export type SessionPendingMessage = SessionPendingUserMessage | SessionPendingSyntheticMessage
@@ -1983,6 +2007,7 @@ export type SessionEventDurable =
| SessionCreated | SessionCreated
| SessionAgentSelected | SessionAgentSelected
| SessionModelSelected | SessionModelSelected
| SessionMoveAdmitted
| SessionMoved | SessionMoved
| SessionRenamed | SessionRenamed
| SessionDeleted | SessionDeleted
@@ -2046,6 +2071,7 @@ export type V2Event =
| SessionCreated | SessionCreated
| SessionAgentSelected | SessionAgentSelected
| SessionModelSelected | SessionModelSelected
| SessionMoveAdmitted
| SessionMoved | SessionMoved
| SessionRenamed | SessionRenamed
| SessionUsageUpdated | SessionUsageUpdated
+11 -13
View File
@@ -185,8 +185,8 @@ export interface Interface {
) => Effect.Effect<SessionMessage.Info[], NotFoundError | MessageDecodeError> ) => Effect.Effect<SessionMessage.Info[], NotFoundError | MessageDecodeError>
/** /**
* Durable admitted session work not yet visible in projected history, * Durable admitted session work not yet visible in projected history,
* ordered by admission. Includes unpromoted user and synthetic inputs and * ordered by admission. Includes unpromoted user and synthetic inputs,
* unhandled compaction barriers. * unhandled compaction barriers, and deferred moves.
*/ */
readonly pending: (sessionID: SessionSchema.ID) => Effect.Effect<SessionPending.Info[], NotFoundError> readonly pending: (sessionID: SessionSchema.ID) => Effect.Effect<SessionPending.Info[], NotFoundError>
readonly cancelPending: (input: PendingInputRef) => Effect.Effect<void, NotFoundError | PendingInputConflictError> readonly cancelPending: (input: PendingInputRef) => Effect.Effect<void, NotFoundError | PendingInputConflictError>
@@ -738,23 +738,21 @@ const layer = Layer.effect(
const info = yield* fs.stat(directory).pipe(Effect.catch(() => Effect.succeed(undefined))) const info = yield* fs.stat(directory).pipe(Effect.catch(() => Effect.succeed(undefined)))
if (!info) return yield* new DestinationNotFoundError({ directory }) if (!info) return yield* new DestinationNotFoundError({ directory })
if (info.type !== "Directory") return yield* new DestinationNotDirectoryError({ directory }) if (info.type !== "Directory") return yield* new DestinationNotDirectoryError({ directory })
if (current.location.directory === directory && current.location.workspaceID === input.workspaceID) return const pending = yield* SessionPending.move(db, input.sessionID)
if (!pending && current.location.directory === directory && current.location.workspaceID === input.workspaceID)
return
const project = yield* projects.resolve(directory) const project = yield* projects.resolve(directory)
yield* persistProject(project) yield* persistProject(project)
if ((yield* execution.active).has(input.sessionID)) { yield* SessionPending.admitMove(db, bus, {
yield* execution.interrupt(input.sessionID) sessionID: input.sessionID,
yield* execution.awaitIdle(input.sessionID) source: current.location,
} data: {
yield* bus.publish(
SessionEvent.Moved,
{
sessionID: input.sessionID,
location: Location.Ref.make({ directory, workspaceID: input.workspaceID }), location: Location.Ref.make({ directory, workspaceID: input.workspaceID }),
projectID: project.id, projectID: project.id,
subpath: RelativePath.make(path.relative(project.directory, directory).replaceAll("\\", "/")), subpath: RelativePath.make(path.relative(project.directory, directory).replaceAll("\\", "/")),
}, },
{ location: current.location }, })
) yield* execution.wake(input.sessionID)
}), }),
compact: Effect.fn("Session.compact")(function* (input) { compact: Effect.fn("Session.compact")(function* (input) {
yield* result.get(input.sessionID) yield* result.get(input.sessionID)
+26 -15
View File
@@ -11,6 +11,8 @@ import { SessionSchema } from "./schema.js"
import { SessionStore } from "./store.js" import { SessionStore } from "./store.js"
import { toSessionError } from "./to-session-error.js" import { toSessionError } from "./to-session-error.js"
import { UserInterruptedError } from "./error.js" import { UserInterruptedError } from "./error.js"
import { Database } from "../database/database.js"
import { SessionPending } from "./pending.js"
export interface Interface { export interface Interface {
/** Snapshots active execution owned by this process. */ /** Snapshots active execution owned by this process. */
@@ -45,6 +47,7 @@ export const layer = Layer.effect(
const store = yield* SessionStore.Service const store = yield* SessionStore.Service
const locations = yield* LocationServiceMap.Service const locations = yield* LocationServiceMap.Service
const bus = yield* Bus.Service const bus = yield* Bus.Service
const db = (yield* Database.Service).db
const reportLifecycle = <A>(sessionID: SessionSchema.ID, effect: Effect.Effect<A>) => const reportLifecycle = <A>(sessionID: SessionSchema.ID, effect: Effect.Effect<A>) =>
effect.pipe( effect.pipe(
Effect.tapCause((cause) => Effect.tapCause((cause) =>
@@ -54,7 +57,6 @@ export const layer = Layer.effect(
Effect.annotateLogs({ sessionID }), Effect.annotateLogs({ sessionID }),
), ),
), ),
Effect.asVoid,
) )
// Write-ahead claim: starting records the durable intent that a turn is in flight, in the same // Write-ahead claim: starting records the durable intent that a turn is in flight, in the same
// transaction as the started event. Terminals release it — except shutdown interruption, which // transaction as the started event. Terminals release it — except shutdown interruption, which
@@ -72,7 +74,7 @@ export const layer = Layer.effect(
reportLifecycle( reportLifecycle(
sessionID, sessionID,
bus.publish(SessionEvent.Execution.Started, { sessionID }, claimOnCommit(sessionID)), bus.publish(SessionEvent.Execution.Started, { sessionID }, claimOnCommit(sessionID)),
), ).pipe(Effect.asVoid),
drain: Effect.fnUntraced(function* (sessionID: SessionSchema.ID, force) { drain: Effect.fnUntraced(function* (sessionID: SessionSchema.ID, force) {
const session = yield* store.get(sessionID) const session = yield* store.get(sessionID)
if (!session) return yield* Effect.die(new Error(`Session not found: ${sessionID}`)) if (!session) return yield* Effect.die(new Error(`Session not found: ${sessionID}`))
@@ -91,11 +93,9 @@ export const layer = Layer.effect(
sessionID, sessionID,
Effect.gen(function* () { Effect.gen(function* () {
const outcome = terminal(exit, reason) const outcome = terminal(exit, reason)
if (outcome.type === "succeeded") { if (outcome.type === "succeeded")
yield* bus.publish(SessionEvent.Execution.Succeeded, { sessionID }, releaseOnCommit(sessionID)) yield* bus.publish(SessionEvent.Execution.Succeeded, { sessionID }, releaseOnCommit(sessionID))
return if (outcome.type === "interrupted")
}
if (outcome.type === "interrupted") {
// A user cancel (or a superseding execution) releases the claim: the turn must not // A user cancel (or a superseding execution) releases the claim: the turn must not
// resurrect at the next boot. Shutdown interruption keeps it for restart continuity. // resurrect at the next boot. Shutdown interruption keeps it for restart continuity.
yield* bus.publish( yield* bus.publish(
@@ -103,16 +103,27 @@ export const layer = Layer.effect(
{ sessionID, reason: outcome.reason }, { sessionID, reason: outcome.reason },
outcome.reason === "shutdown" ? undefined : releaseOnCommit(sessionID), outcome.reason === "shutdown" ? undefined : releaseOnCommit(sessionID),
) )
return if (outcome.type === "failed")
} yield* bus.publish(
SessionEvent.Execution.Failed,
{
sessionID,
error: outcome.error,
},
releaseOnCommit(sessionID),
)
if (outcome.type === "interrupted" && outcome.reason === "shutdown") return false
const pending = yield* SessionPending.move(db, sessionID)
if (!pending) return false
const session = yield* store.get(sessionID)
if (!session) return yield* Effect.die(new Error(`Session not found: ${sessionID}`))
yield* bus.publish( yield* bus.publish(
SessionEvent.Execution.Failed, SessionEvent.Moved,
{ { sessionID, moveID: pending.id, ...pending.data },
sessionID, { location: session.location },
error: outcome.error,
},
releaseOnCommit(sessionID),
) )
return yield* SessionPending.has(db, sessionID, "any")
}), }),
), ),
}) })
@@ -130,7 +141,7 @@ export const layer = Layer.effect(
export const node = makeGlobalNode({ export const node = makeGlobalNode({
service: Service, service: Service,
layer, layer,
deps: [SessionStore.node, LocationServiceMap.node, Bus.node], deps: [SessionStore.node, LocationServiceMap.node, Bus.node, Database.node],
}) })
/** Low-level compatibility layer for callers that only need durable Session recording. */ /** Low-level compatibility layer for callers that only need durable Session recording. */
+12 -1
View File
@@ -7,6 +7,8 @@ import { SessionEvent } from "../event.js"
import { SessionExecution } from "../execution.js" import { SessionExecution } from "../execution.js"
import { SessionSchema } from "../schema.js" import { SessionSchema } from "../schema.js"
import { SessionStore } from "../store.js" import { SessionStore } from "../store.js"
import { Database } from "../../database/database.js"
import { SessionPending } from "../pending.js"
const CONTINUE_AFTER_SERVER_RESTART = const CONTINUE_AFTER_SERVER_RESTART =
"The server restarted while you were working. Continue from where you left off without repeating completed work." "The server restarted while you were working. Continue from where you left off without repeating completed work."
@@ -62,6 +64,7 @@ export const layer = (options?: Options) =>
const store = yield* SessionStore.Service const store = yield* SessionStore.Service
const execution = yield* SessionExecution.Service const execution = yield* SessionExecution.Service
const bus = yield* Bus.Service const bus = yield* Bus.Service
const db = (yield* Database.Service).db
const scope = yield* Effect.scope const scope = yield* Effect.scope
const maxAttempts = options?.maxAttempts ?? DEFAULT_MAX_ATTEMPTS const maxAttempts = options?.maxAttempts ?? DEFAULT_MAX_ATTEMPTS
@@ -103,6 +106,14 @@ export const layer = (options?: Options) =>
// them would only inject a stray continuation into a live turn. // them would only inject a stray continuation into a live turn.
const orphaned = (yield* store.listSuspended()).filter((sessionID) => !active.has(sessionID)) const orphaned = (yield* store.listSuspended()).filter((sessionID) => !active.has(sessionID))
yield* Effect.forEach(orphaned, resumeOne, { concurrency: "unbounded", discard: true }) yield* Effect.forEach(orphaned, resumeOne, { concurrency: "unbounded", discard: true })
const claimed = new Set(orphaned)
yield* Effect.forEach(
(yield* SessionPending.moveSessions(db)).filter(
(sessionID) => !active.has(sessionID) && !claimed.has(sessionID),
),
execution.wake,
{ concurrency: "unbounded", discard: true },
)
}), }),
}) })
}), }),
@@ -111,5 +122,5 @@ export const layer = (options?: Options) =>
export const node = makeGlobalNode({ export const node = makeGlobalNode({
service: Service, service: Service,
layer: layer(), layer: layer(),
deps: [SessionStore.node, SessionExecution.node, Bus.node], deps: [SessionStore.node, SessionExecution.node, Bus.node, Database.node],
}) })
@@ -106,6 +106,7 @@ export function update(adapter: Adapter, event: SessionEvent.DurableEvent) {
) )
}) })
}, },
"session.move.admitted": () => Effect.void,
"session.renamed": () => Effect.void, "session.renamed": () => Effect.void,
"session.deleted": () => Effect.void, "session.deleted": () => Effect.void,
"session.forked": () => Effect.void, "session.forked": () => Effect.void,
+118 -12
View File
@@ -7,6 +7,8 @@ import {
Delivery, Delivery,
Info, Info,
Message, Message,
Move,
MoveData,
Synthetic, Synthetic,
SyntheticData, SyntheticData,
User, User,
@@ -19,10 +21,11 @@ import { SessionEvent } from "./event.js"
import { SessionMessage } from "./message.js" import { SessionMessage } from "./message.js"
import { SessionSchema } from "./schema.js" import { SessionSchema } from "./schema.js"
import { SessionMessageTable, SessionPendingTable } from "./sql.js" import { SessionMessageTable, SessionPendingTable } from "./sql.js"
import { Event } from "@opencode-ai/schema/event"
type DatabaseService = Database.Interface["db"] type DatabaseService = Database.Interface["db"]
export { Compaction, Delivery, Info, Message, Synthetic, SyntheticData, User, UserData } export { Compaction, Delivery, Info, Message, Move, MoveData, Synthetic, SyntheticData, User, UserData }
/** /**
* Which pending input `promote` may consume: "steer" promotes steers only (a step * Which pending input `promote` may consume: "steer" promotes steers only (a step
@@ -35,6 +38,8 @@ const decodeUser = Schema.decodeUnknownSync(UserData)
const encodeUser = Schema.encodeSync(UserData) const encodeUser = Schema.encodeSync(UserData)
const decodeSynthetic = Schema.decodeUnknownSync(SyntheticData) const decodeSynthetic = Schema.decodeUnknownSync(SyntheticData)
const encodeSynthetic = Schema.encodeSync(SyntheticData) const encodeSynthetic = Schema.encodeSync(SyntheticData)
const decodeMove = Schema.decodeUnknownSync(MoveData)
const encodeMove = Schema.encodeSync(MoveData)
const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info) const decodeMessage = Schema.decodeUnknownSync(SessionMessage.Info)
const inboxLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>() const inboxLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>()
type PendingRef = { readonly id: SessionMessage.ID; readonly sessionID: SessionSchema.ID } type PendingRef = { readonly id: SessionMessage.ID; readonly sessionID: SessionSchema.ID }
@@ -42,21 +47,24 @@ type PendingRef = { readonly id: SessionMessage.ID; readonly sessionID: SessionS
export class LifecycleConflict extends Schema.TaggedErrorClass<LifecycleConflict>()( export class LifecycleConflict extends Schema.TaggedErrorClass<LifecycleConflict>()(
"SessionPending.LifecycleConflict", "SessionPending.LifecycleConflict",
{ {
id: SessionMessage.ID, id: Schema.Union([SessionMessage.ID, Event.ID]),
}, },
) {} ) {}
const fromRow = (row: typeof SessionPendingTable.$inferSelect): Info => { const fromRow = (row: typeof SessionPendingTable.$inferSelect): Info => {
const base = { const base = {
id: SessionMessage.ID.make(row.id),
sessionID: SessionSchema.ID.make(row.session_id), sessionID: SessionSchema.ID.make(row.session_id),
timeCreated: DateTime.makeUnsafe(row.time_created), timeCreated: DateTime.makeUnsafe(row.time_created),
} }
if (row.type === "compaction") return Compaction.make({ ...base, type: "compaction" }) if (row.type === "move")
if (!row.delivery) throw new LifecycleConflict({ id: base.id }) return Move.make({ ...base, id: Event.ID.make(row.id), type: "move", data: decodeMove(row.data) })
const id = SessionMessage.ID.make(row.id)
if (row.type === "compaction") return Compaction.make({ ...base, id, type: "compaction" })
if (!row.delivery) throw new LifecycleConflict({ id })
if (row.type === "user") if (row.type === "user")
return User.make({ return User.make({
...base, ...base,
id,
type: "user", type: "user",
data: decodeUser(row.data), data: decodeUser(row.data),
delivery: row.delivery, delivery: row.delivery,
@@ -64,11 +72,12 @@ const fromRow = (row: typeof SessionPendingTable.$inferSelect): Info => {
if (row.type === "synthetic") if (row.type === "synthetic")
return Synthetic.make({ return Synthetic.make({
...base, ...base,
id,
type: "synthetic", type: "synthetic",
data: decodeSynthetic(row.data), data: decodeSynthetic(row.data),
delivery: row.delivery, delivery: row.delivery,
}) })
throw new LifecycleConflict({ id: base.id }) throw new LifecycleConflict({ id })
} }
export const find = Effect.fn("SessionPending.find")(function* (db: DatabaseService, id: SessionMessage.ID) { export const find = Effect.fn("SessionPending.find")(function* (db: DatabaseService, id: SessionMessage.ID) {
@@ -98,6 +107,44 @@ export const compaction = Effect.fn("SessionPending.compaction")(function* (
return entry.type === "compaction" ? entry : undefined return entry.type === "compaction" ? entry : undefined
}) })
export const move = Effect.fn("SessionPending.move")(function* (db: DatabaseService, sessionID: SessionSchema.ID) {
const row = yield* db
.select()
.from(SessionPendingTable)
.where(and(eq(SessionPendingTable.session_id, sessionID), eq(SessionPendingTable.type, "move")))
.orderBy(asc(SessionPendingTable.admitted_seq))
.limit(1)
.get()
.pipe(Effect.orDie)
if (!row) return
const entry = fromRow(row)
return entry.type === "move" ? entry : undefined
})
export const admitMove = Effect.fn("SessionPending.admitMove")(function* (
db: DatabaseService,
bus: Bus.Interface,
input: { readonly sessionID: SessionSchema.ID; readonly data: MoveData; readonly source: MoveData["location"] },
) {
return yield* inboxLocks.withLock(input.sessionID)(
Effect.gen(function* () {
const pending = yield* move(db, input.sessionID)
if (pending && JSON.stringify(encodeMove(pending.data)) === JSON.stringify(encodeMove(input.data))) return pending
const event = yield* bus.publish(
SessionEvent.MoveAdmitted,
{
sessionID: input.sessionID,
move: input.data,
},
{ location: input.source },
)
const stored = yield* move(db, input.sessionID)
if (stored) return stored
return yield* Effect.die(new LifecycleConflict({ id: event.id }))
}),
)
})
const promotedFromMessage = Effect.fn("SessionPending.promotedFromMessage")(function* ( const promotedFromMessage = Effect.fn("SessionPending.promotedFromMessage")(function* (
db: DatabaseService, db: DatabaseService,
sessionID: SessionSchema.ID, sessionID: SessionSchema.ID,
@@ -288,6 +335,35 @@ export const projectCompactionAdmitted = Effect.fn("SessionPending.projectCompac
return yield* Effect.die(new LifecycleConflict({ id: input.id })) return yield* Effect.die(new LifecycleConflict({ id: input.id }))
}) })
export const projectMoveAdmitted = Effect.fn("SessionPending.projectMoveAdmitted")(function* (
db: DatabaseService,
input: {
readonly admittedSeq: number
readonly id: Event.ID
readonly sessionID: SessionSchema.ID
readonly data: MoveData
readonly timeCreated: DateTime.Utc
},
) {
yield* db
.delete(SessionPendingTable)
.where(and(eq(SessionPendingTable.session_id, input.sessionID), eq(SessionPendingTable.type, "move")))
.run()
.pipe(Effect.orDie)
yield* db
.insert(SessionPendingTable)
.values({
id: input.id,
session_id: input.sessionID,
type: "move",
data: input.data,
admitted_seq: input.admittedSeq,
time_created: DateTime.toEpochMillis(input.timeCreated),
})
.run()
.pipe(Effect.orDie)
})
/** /**
* Consume one pending row at promotion. The row's content feeds the projected * Consume one pending row at promotion. The row's content feeds the projected
* message insert inside the same event transaction; the deleted row is what * message insert inside the same event transaction; the deleted row is what
@@ -297,7 +373,8 @@ export const projectPromoted = Effect.fn("SessionPending.projectPromoted")(funct
db: DatabaseService, db: DatabaseService,
input: PendingRef, input: PendingRef,
) { ) {
if (yield* compaction(db, input.sessionID)) return yield* Effect.die(new LifecycleConflict({ id: input.id })) if ((yield* compaction(db, input.sessionID)) || (yield* move(db, input.sessionID)))
return yield* Effect.die(new LifecycleConflict({ id: input.id }))
const deleted = yield* db const deleted = yield* db
.delete(SessionPendingTable) .delete(SessionPendingTable)
.where(and(eq(SessionPendingTable.id, input.id), eq(SessionPendingTable.session_id, input.sessionID))) .where(and(eq(SessionPendingTable.id, input.id), eq(SessionPendingTable.session_id, input.sessionID)))
@@ -306,7 +383,8 @@ export const projectPromoted = Effect.fn("SessionPending.projectPromoted")(funct
.pipe(Effect.orDie) .pipe(Effect.orDie)
if (!deleted) return yield* Effect.die(new LifecycleConflict({ id: input.id })) if (!deleted) return yield* Effect.die(new LifecycleConflict({ id: input.id }))
const stored = fromRow(deleted) const stored = fromRow(deleted)
if (stored.type === "compaction") return yield* Effect.die(new LifecycleConflict({ id: input.id })) if (stored.type === "compaction" || stored.type === "move")
return yield* Effect.die(new LifecycleConflict({ id: input.id }))
return stored return stored
}) })
@@ -374,6 +452,33 @@ export const settleCompaction = Effect.fn("SessionPending.settleCompaction")(fun
return undefined return undefined
}) })
export const settleMove = Effect.fn("SessionPending.settleMove")(function* (
db: DatabaseService,
input: { readonly sessionID: SessionSchema.ID; readonly id: Event.ID },
) {
yield* db
.delete(SessionPendingTable)
.where(
and(
eq(SessionPendingTable.id, input.id),
eq(SessionPendingTable.session_id, input.sessionID),
eq(SessionPendingTable.type, "move"),
),
)
.run()
.pipe(Effect.orDie)
})
export const moveSessions = Effect.fn("SessionPending.moveSessions")(function* (db: DatabaseService) {
const rows = yield* db
.select({ sessionID: SessionPendingTable.session_id })
.from(SessionPendingTable)
.where(eq(SessionPendingTable.type, "move"))
.all()
.pipe(Effect.orDie)
return [...new Set(rows.map((row) => row.sessionID))]
})
export const list = Effect.fn("SessionPending.list")(function* (db: DatabaseService, sessionID: SessionSchema.ID) { export const list = Effect.fn("SessionPending.list")(function* (db: DatabaseService, sessionID: SessionSchema.ID) {
const rows = yield* db const rows = yield* db
.select() .select()
@@ -397,7 +502,7 @@ export const has = Effect.fn("SessionPending.has")(function* (
sessionID: SessionSchema.ID, sessionID: SessionSchema.ID,
scope: Scope, scope: Scope,
) { ) {
if (scope !== "any" && (yield* compaction(db, sessionID))) return false if (scope !== "any" && ((yield* compaction(db, sessionID)) || (yield* move(db, sessionID)))) return false
const row = yield* db const row = yield* db
.select({ id: SessionPendingTable.id }) .select({ id: SessionPendingTable.id })
.from(SessionPendingTable) .from(SessionPendingTable)
@@ -473,12 +578,13 @@ const publish = Effect.fn("SessionPending.publish")(function* (
sessionID: SessionSchema.ID, sessionID: SessionSchema.ID,
rows: ReadonlyArray<typeof SessionPendingTable.$inferSelect>, rows: ReadonlyArray<typeof SessionPendingTable.$inferSelect>,
) { ) {
if (yield* compaction(db, sessionID)) return 0 if ((yield* compaction(db, sessionID)) || (yield* move(db, sessionID))) return 0
yield* Effect.forEach( yield* Effect.forEach(
rows, rows,
(row) => { (row) => {
const entry = fromRow(row) const entry = fromRow(row)
if (entry.type === "compaction") return Effect.die(new LifecycleConflict({ id: entry.id })) if (entry.type === "compaction" || entry.type === "move")
return Effect.die(new LifecycleConflict({ id: entry.id }))
return bus return bus
.publish(SessionEvent.InputPromoted, { .publish(SessionEvent.InputPromoted, {
sessionID, sessionID,
@@ -512,7 +618,7 @@ export const promote = Effect.fn("SessionPending.promote")(function* (
) { ) {
return yield* inboxLocks.withLock(sessionID)( return yield* inboxLocks.withLock(sessionID)(
Effect.gen(function* () { Effect.gen(function* () {
if (yield* compaction(db, sessionID)) return 0 if ((yield* compaction(db, sessionID)) || (yield* move(db, sessionID))) return 0
const steers = yield* db const steers = yield* db
.select() .select()
.from(SessionPendingTable) .from(SessionPendingTable)
+15
View File
@@ -433,6 +433,8 @@ const layer = Layer.effectDiscard(
.run() .run()
.pipe(Effect.orDie) .pipe(Effect.orDie)
yield* InstructionState.reset(db, event.data.sessionID) yield* InstructionState.reset(db, event.data.sessionID)
if (event.data.moveID)
yield* SessionPending.settleMove(db, { sessionID: event.data.sessionID, id: event.data.moveID })
}), }),
) )
yield* bus.project(SessionEvent.Deleted, (event) => yield* bus.project(SessionEvent.Deleted, (event) =>
@@ -522,6 +524,19 @@ const layer = Layer.effectDiscard(
.pipe(Effect.orDie) .pipe(Effect.orDie)
}), }),
) )
yield* bus.project(SessionEvent.MoveAdmitted, (event) =>
Effect.gen(function* () {
if (event.durable === undefined)
return yield* Effect.die(new Error("Durable Session event is missing aggregate sequence"))
yield* SessionPending.projectMoveAdmitted(db, {
admittedSeq: event.durable.seq,
id: event.id,
sessionID: event.data.sessionID,
data: event.data.move,
timeCreated: event.created,
})
}),
)
yield* bus.project(SessionEvent.InputCancelled, (event) => yield* bus.project(SessionEvent.InputCancelled, (event) =>
SessionPending.projectCancelled(db, { SessionPending.projectCancelled(db, {
id: event.data.inputID, id: event.data.inputID,
+24 -3
View File
@@ -1,6 +1,6 @@
export * as SessionRunCoordinator from "./run-coordinator.js" export * as SessionRunCoordinator from "./run-coordinator.js"
import { Deferred, Effect, Exit, Fiber, FiberSet, Scope } from "effect" import { Cause, Deferred, Effect, Exit, Fiber, FiberSet, Scope } from "effect"
/** Serializes execution for each key while allowing different keys to run concurrently. */ /** Serializes execution for each key while allowing different keys to run concurrently. */
export interface Coordinator<Key, E, Reason = never> { export interface Coordinator<Key, E, Reason = never> {
@@ -50,11 +50,17 @@ export const make = <Key, E, Reason = never>(options: {
* Runs in the execution fiber for every exit, including interruption, after the final * Runs in the execution fiber for every exit, including interruption, after the final
* drain and before the execution settles (waiters resolve after it completes). * drain and before the execution settles (waiters resolve after it completes).
*/ */
readonly settled?: (key: Key, exit: Exit.Exit<void, E>, reason?: Reason) => Effect.Effect<void> readonly settled?: (key: Key, exit: Exit.Exit<void, E>, reason?: Reason) => Effect.Effect<boolean | void>
}): Effect.Effect<Coordinator<Key, E, Reason>, never, Scope.Scope> => }): Effect.Effect<Coordinator<Key, E, Reason>, never, Scope.Scope> =>
Effect.gen(function* () { Effect.gen(function* () {
const executions = new Map<Key, Execution<E, Reason>>() const executions = new Map<Key, Execution<E, Reason>>()
const fork = yield* FiberSet.makeRuntime<never, void, never>() const fork = yield* FiberSet.makeRuntime<never, void, never>()
const closing = { value: false }
yield* Effect.addFinalizer(() =>
Effect.sync(() => {
closing.value = true
}),
)
const loop = (key: Key, execution: Execution<E, Reason>, force: boolean): Effect.Effect<void, E> => const loop = (key: Key, execution: Execution<E, Reason>, force: boolean): Effect.Effect<void, E> =>
Effect.suspend(() => options.drain(key, force)).pipe( Effect.suspend(() => options.drain(key, force)).pipe(
@@ -85,7 +91,22 @@ export const make = <Key, E, Reason = never>(options: {
Effect.onExit((exit) => Effect.onExit((exit) =>
Effect.sync(() => { Effect.sync(() => {
execution.owner = undefined execution.owner = undefined
}).pipe(Effect.andThen(options.settled?.(key, exit, execution.interruptionReason) ?? Effect.void)), if (closing.value && Exit.isFailure(exit) && Cause.hasInterrupts(exit.cause)) {
execution.stopping = true
execution.pendingWake = false
}
}).pipe(
Effect.andThen(options.settled?.(key, exit, execution.interruptionReason) ?? Effect.void),
Effect.map(Boolean),
Effect.tap((restart) =>
restart && !execution.stopping
? Effect.sync(() => {
execution.pendingWake = true
})
: Effect.void,
),
Effect.asVoid,
),
), ),
Effect.onExit((exit) => Effect.sync(() => settle(key, execution, exit))), Effect.onExit((exit) => Effect.sync(() => settle(key, execution, exit))),
Effect.exit, Effect.exit,
+4
View File
@@ -143,6 +143,7 @@ const layer = Layer.effect(
let promotable: SessionPending.Promotable = "input" let promotable: SessionPending.Promotable = "input"
let step = 1 let step = 1
while (true) { while (true) {
if (yield* SessionPending.move(db, sessionID)) return
const result = yield* runStep(sessionID, promotable, step) const result = yield* runStep(sessionID, promotable, step)
if (step === 1) yield* startTitle(sessionID) if (step === 1) yield* startTitle(sessionID)
yield* runPendingCompaction(sessionID) yield* runPendingCompaction(sessionID)
@@ -236,6 +237,8 @@ const layer = Layer.effect(
// a blocked first step leaves pending inputs untouched. // a blocked first step leaves pending inputs untouched.
yield* InstructionState.prepare(db, bus, selected.instructions, selected.session.id) yield* InstructionState.prepare(db, bus, selected.instructions, selected.session.id)
const promoted = promotable ? yield* SessionPending.promote(db, bus, selected.session.id, promotable) : 0 const promoted = promotable ? yield* SessionPending.promote(db, bus, selected.session.id, promotable) : 0
if (promotable && promoted === 0 && (yield* SessionPending.move(db, sessionID)))
return CallOutcome.Completed({ needsContinuation: false, step })
// Promoted input opens a fresh step allowance. // Promoted input opens a fresh step allowance.
const currentStep = promoted > 0 ? 1 : step const currentStep = promoted > 0 ? 1 : step
const loaded = yield* context.load(selected) const loaded = yield* context.load(selected)
@@ -482,6 +485,7 @@ const layer = Layer.effect(
const runPendingCompaction = Effect.fn("SessionRunner.runPendingCompaction")(function* ( const runPendingCompaction = Effect.fn("SessionRunner.runPendingCompaction")(function* (
sessionID: SessionSchema.ID, sessionID: SessionSchema.ID,
) { ) {
if (yield* SessionPending.move(db, sessionID)) return
const pending = yield* SessionPending.compaction(db, sessionID) const pending = yield* SessionPending.compaction(db, sessionID)
if (!pending) return if (!pending) return
const session = yield* getSession(sessionID) const session = yield* getSession(sessionID)
+4 -2
View File
@@ -96,13 +96,15 @@ export const SessionMessageTable = sqliteTable(
export const SessionPendingTable = sqliteTable( export const SessionPendingTable = sqliteTable(
"session_pending", "session_pending",
{ {
id: text().$type<SessionMessage.ID>().primaryKey(), id: text().$type<SessionPending.Info["id"]>().primaryKey(),
session_id: text() session_id: text()
.$type<SessionSchema.ID>() .$type<SessionSchema.ID>()
.notNull() .notNull()
.references(() => SessionTable.id, { onDelete: "cascade" }), .references(() => SessionTable.id, { onDelete: "cascade" }),
type: text().$type<SessionPending.Info["type"]>().notNull(), type: text().$type<SessionPending.Info["type"]>().notNull(),
data: text({ mode: "json" }).$type<UserData | SyntheticData | Record<string, never>>().notNull(), data: text({ mode: "json" })
.$type<UserData | SyntheticData | SessionPending.MoveData | Record<string, never>>()
.notNull(),
delivery: text().$type<SessionPending.Delivery>(), delivery: text().$type<SessionPending.Delivery>(),
admitted_seq: integer().notNull(), admitted_seq: integer().notNull(),
time_created: integer() time_created: integer()
+105 -2
View File
@@ -8,7 +8,7 @@ import { LocationServiceMap } from "@opencode-ai/core/location-service-map"
import type { LocationServices } from "@opencode-ai/core/location-services" import type { LocationServices } from "@opencode-ai/core/location-services"
import { Project } from "@opencode-ai/core/project" import { Project } from "@opencode-ai/core/project"
import { ProjectTable } from "@opencode-ai/core/project/sql" import { ProjectTable } from "@opencode-ai/core/project/sql"
import { AbsolutePath } from "@opencode-ai/core/schema" import { AbsolutePath, RelativePath } from "@opencode-ai/core/schema"
import { Session } from "@opencode-ai/core/session" import { Session } from "@opencode-ai/core/session"
import { SessionExecution } from "@opencode-ai/core/session/execution" import { SessionExecution } from "@opencode-ai/core/session/execution"
import { SessionRestart } from "@opencode-ai/core/session/execution/restart" import { SessionRestart } from "@opencode-ai/core/session/execution/restart"
@@ -17,11 +17,16 @@ import { SessionEvent } from "@opencode-ai/core/session/event"
import { SessionRunner } from "@opencode-ai/core/session/runner" import { SessionRunner } from "@opencode-ai/core/session/runner"
import { SessionTable } from "@opencode-ai/core/session/sql" import { SessionTable } from "@opencode-ai/core/session/sql"
import { SessionStore } from "@opencode-ai/core/session/store" import { SessionStore } from "@opencode-ai/core/session/store"
import { SessionProjector } from "@opencode-ai/core/session/projector"
import { SessionPending } from "@opencode-ai/core/session/pending"
import { Location } from "@opencode-ai/core/location"
import { Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from "effect" import { Context, Deferred, Effect, Exit, Fiber, Layer, LayerMap, Scope } from "effect"
import { eq } from "drizzle-orm" import { eq } from "drizzle-orm"
import { testEffect } from "./lib/effect" import { testEffect } from "./lib/effect"
const it = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionStore.node]))) const it = testEffect(
AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node])),
)
describe("SessionExecution lifecycle", () => { describe("SessionExecution lifecycle", () => {
test("classifies success and typed failure terminals", () => { test("classifies success and typed failure terminals", () => {
@@ -133,6 +138,104 @@ describe("SessionExecution lifecycle", () => {
}), }),
) )
it.effect("applies a deferred move only after the active execution settles", () =>
Effect.gen(function* () {
const database = yield* Database.Service
const bus = yield* Bus.Service
const store = yield* SessionStore.Service
const sessionID = Session.ID.make("ses_deferred_move")
yield* seedSessions(database, [sessionID])
const draining = yield* Deferred.make<void>()
const release = yield* Deferred.make<void>()
const scope = yield* Scope.make()
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
const context = yield* buildExecution(scope, () =>
Deferred.succeed(draining, undefined).pipe(Effect.andThen(Deferred.await(release))),
)
const execution = Context.get(context, SessionExecution.Service)
yield* execution.resume(sessionID).pipe(Effect.forkIn(scope))
yield* Deferred.await(draining)
yield* bus.publish(SessionEvent.MoveAdmitted, {
sessionID,
move: {
location: Location.Ref.make({ directory: AbsolutePath.make("/destination") }),
projectID: Project.ID.global,
subpath: RelativePath.make(""),
},
})
expect((yield* store.get(sessionID))?.location.directory).toBe(AbsolutePath.make("/project"))
expect((yield* SessionPending.move(database.db, sessionID))?.data.location.directory).toBe(
AbsolutePath.make("/destination"),
)
yield* Deferred.succeed(release, undefined)
yield* execution.awaitIdle(sessionID)
expect((yield* store.get(sessionID))?.location.directory).toBe(AbsolutePath.make("/destination"))
expect(yield* SessionPending.move(database.db, sessionID)).toBeUndefined()
}),
)
it.effect("settling one move preserves a newer admitted destination", () =>
Effect.gen(function* () {
const database = yield* Database.Service
const bus = yield* Bus.Service
const store = yield* SessionStore.Service
const sessionID = Session.ID.make("ses_move_replacement")
yield* seedSessions(database, [sessionID])
const first = {
location: Location.Ref.make({ directory: AbsolutePath.make("/first") }),
projectID: Project.ID.global,
subpath: RelativePath.make("first"),
}
const second = {
location: Location.Ref.make({ directory: AbsolutePath.make("/second") }),
projectID: Project.ID.global,
subpath: RelativePath.make("second"),
}
const admittedFirst = yield* bus.publish(SessionEvent.MoveAdmitted, { sessionID, move: first })
const admittedSecond = yield* bus.publish(SessionEvent.MoveAdmitted, { sessionID, move: second })
yield* bus.publish(SessionEvent.Moved, { sessionID, moveID: admittedFirst.id, ...first })
expect((yield* store.get(sessionID))?.location.directory).toBe(first.location.directory)
expect((yield* SessionPending.move(database.db, sessionID))?.id).toBe(admittedSecond.id)
yield* bus.publish(SessionEvent.Moved, { sessionID, moveID: admittedSecond.id, ...second })
expect((yield* store.get(sessionID))?.location.directory).toBe(second.location.directory)
expect(yield* SessionPending.move(database.db, sessionID)).toBeUndefined()
}),
)
it.effect("recovers an unclaimed deferred move on startup", () =>
Effect.gen(function* () {
const database = yield* Database.Service
const bus = yield* Bus.Service
const store = yield* SessionStore.Service
const sessionID = Session.ID.make("ses_move_recovery")
yield* seedSessions(database, [sessionID])
yield* bus.publish(SessionEvent.MoveAdmitted, {
sessionID,
move: {
location: Location.Ref.make({ directory: AbsolutePath.make("/recovered") }),
projectID: Project.ID.global,
},
})
const scope = yield* Scope.make()
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void))
const context = yield* buildExecution(scope, () => Effect.void)
const execution = Context.get(context, SessionExecution.Service)
yield* Context.get(context, SessionRestart.Service).resumeSuspendedSessions
yield* execution.awaitIdle(sessionID)
expect((yield* store.get(sessionID))?.location.directory).toBe(AbsolutePath.make("/recovered"))
expect(yield* SessionPending.move(database.db, sessionID)).toBeUndefined()
}),
)
it.effect("starts every claimed execution without waiting for earlier drains to finish", () => it.effect("starts every claimed execution without waiting for earlier drains to finish", () =>
Effect.gen(function* () { Effect.gen(function* () {
const database = yield* Database.Service const database = yield* Database.Service
+20 -17
View File
@@ -1,5 +1,6 @@
import { describe, expect } from "bun:test" import { describe, expect } from "bun:test"
import path from "path" import path from "path"
import fs from "fs/promises"
import { Effect, Layer } from "effect" import { Effect, Layer } from "effect"
import { Bus } from "@opencode-ai/core/bus" import { Bus } from "@opencode-ai/core/bus"
import { Database } from "@opencode-ai/core/database/database" import { Database } from "@opencode-ai/core/database/database"
@@ -34,7 +35,7 @@ const it = testEffect(
) )
describe("Session.move", () => { describe("Session.move", () => {
it.effect("moves a session whose source directory no longer exists", () => it.effect("durably admits a move when the source directory no longer exists", () =>
Effect.acquireRelease( Effect.acquireRelease(
Effect.promise(() => tmpdir()), Effect.promise(() => tmpdir()),
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()), (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
@@ -49,24 +50,26 @@ describe("Session.move", () => {
yield* session.move({ sessionID: created.id, directory: destination }) yield* session.move({ sessionID: created.id, directory: destination })
expect((yield* session.get(created.id)).location.directory).toBe(destination) expect((yield* session.get(created.id)).location.directory).toBe(
const messages = yield* session.messages({ sessionID: created.id, order: "asc" }) AbsolutePath.make(path.join(tmp.path, "deleted")),
expect(messages).toEqual([ )
expect.objectContaining({ expect(yield* session.pending(created.id)).toMatchObject([
type: "location-switched", {
location: { directory: destination }, type: "move",
projectID: Project.ID.global, data: { location: { directory: destination }, projectID: Project.ID.global },
previous: { },
location: { directory: path.join(tmp.path, "deleted") },
projectID: Project.ID.global,
subpath: "",
},
subpath: "",
}),
]) ])
yield* session.move({ sessionID: created.id, directory: destination }) const replacement = AbsolutePath.make(path.join(tmp.path, "replacement"))
expect(yield* session.messages({ sessionID: created.id, order: "asc" })).toEqual(messages) yield* Effect.promise(() => fs.mkdir(replacement))
yield* session.move({ sessionID: created.id, directory: replacement })
expect(yield* session.pending(created.id)).toMatchObject([
{
type: "move",
data: { location: { directory: replacement }, projectID: Project.ID.global },
},
])
}), }),
), ),
), ),
@@ -143,19 +143,26 @@ describe("SessionRunCoordinator", () => {
it.effect("cleans active executions when its scope closes", () => it.effect("cleans active executions when its scope closes", () =>
Effect.gen(function* () { Effect.gen(function* () {
const started = yield* Deferred.make<void>() const started = yield* Deferred.make<void>()
let runs = 0
const coordinator = yield* Effect.scoped( const coordinator = yield* Effect.scoped(
Effect.gen(function* () { Effect.gen(function* () {
const coordinator = yield* SessionRunCoordinator.make({ const coordinator = yield* SessionRunCoordinator.make({
drain: () => Deferred.succeed(started, undefined).pipe(Effect.andThen(Effect.never)), drain: () =>
Effect.sync(() => runs++).pipe(
Effect.andThen(Deferred.succeed(started, undefined)),
Effect.andThen(Effect.never),
),
}) })
yield* coordinator.wake("session") yield* coordinator.wake("session")
yield* Deferred.await(started) yield* Deferred.await(started)
yield* coordinator.wake("session")
expect(Array.from(yield* coordinator.active)).toEqual(["session"]) expect(Array.from(yield* coordinator.active)).toEqual(["session"])
return coordinator return coordinator
}), }),
) )
expect(Array.from(yield* coordinator.active)).toEqual([]) expect(Array.from(yield* coordinator.active)).toEqual([])
expect(runs).toBe(1)
}), }),
) )
@@ -517,6 +524,31 @@ describe("SessionRunCoordinator", () => {
), ),
) )
it.effect("starts one successor when settlement requests it", () =>
Effect.scoped(
Effect.gen(function* () {
const successor = yield* Deferred.make<void>()
let drains = 0
let settlements = 0
const coordinator = yield* SessionRunCoordinator.make<string, never>({
drain: () =>
Effect.sync(() => {
drains++
if (drains === 2) Deferred.doneUnsafe(successor, Effect.void)
}),
settled: () => Effect.sync(() => ++settlements === 1),
})
yield* coordinator.wake("session")
yield* Deferred.await(successor)
yield* coordinator.awaitIdle("session")
expect(drains).toBe(2)
expect(settlements).toBe(2)
}),
),
)
it.effect("trampolines synchronous self-waking execution", () => it.effect("trampolines synchronous self-waking execution", () =>
Effect.scoped( Effect.scoped(
Effect.gen(function* () { Effect.gen(function* () {
+1 -1
View File
@@ -308,7 +308,7 @@ export const makeSessionGroup = <I extends HttpApiMiddleware.AnyId, S>(sessionLo
OpenApi.annotations({ OpenApi.annotations({
identifier: "v2.session.move", identifier: "v2.session.move",
summary: "Move session", summary: "Move session",
description: "Move a session to another project directory, optionally transferring local changes.", description: "Move a session to another project directory after any active execution settles.",
}), }),
), ),
) )
+12
View File
@@ -90,6 +90,7 @@ export const Moved = Event.durable({
...options, ...options,
schema: { schema: {
...Base, ...Base,
moveID: Event.ID.pipe(optional),
location: Location.Ref, location: Location.Ref,
projectID: Project.ID.pipe(optional), projectID: Project.ID.pipe(optional),
subpath: RelativePath.pipe(optional), subpath: RelativePath.pipe(optional),
@@ -97,6 +98,16 @@ export const Moved = Event.durable({
}) })
export type Moved = typeof Moved.Type export type Moved = typeof Moved.Type
export const MoveAdmitted = Event.durable({
type: "session.move.admitted",
...options,
schema: {
...Base,
move: SessionPending.MoveData,
},
})
export type MoveAdmitted = typeof MoveAdmitted.Type
export const Renamed = Event.durable({ export const Renamed = Event.durable({
type: "session.renamed", type: "session.renamed",
...options, ...options,
@@ -597,6 +608,7 @@ export const Definitions = Event.inventory(
Created, Created,
AgentSelected, AgentSelected,
ModelSelected, ModelSelected,
MoveAdmitted,
Moved, Moved,
Renamed, Renamed,
UsageUpdated, UsageUpdated,
+21 -1
View File
@@ -7,6 +7,10 @@ import { DateTimeUtcFromMillis } from "./schema.js"
import { SessionDelivery } from "./session-delivery.js" import { SessionDelivery } from "./session-delivery.js"
import { SessionID } from "./session-id.js" import { SessionID } from "./session-id.js"
import { SessionMessage } from "./session-message.js" import { SessionMessage } from "./session-message.js"
import { Event } from "./event.js"
import { Location } from "./location.js"
import { Project } from "./project.js"
import { RelativePath } from "./schema.js"
export const Delivery = SessionDelivery.Delivery export const Delivery = SessionDelivery.Delivery
export type Delivery = SessionDelivery.Delivery export type Delivery = SessionDelivery.Delivery
@@ -68,7 +72,23 @@ export const Compaction = Schema.Struct({
type: Schema.tag("compaction"), type: Schema.tag("compaction"),
}).annotate({ identifier: "SessionPending.Compaction" }) }).annotate({ identifier: "SessionPending.Compaction" })
export const Info = Schema.Union([User, Synthetic, Compaction]).pipe( export interface MoveData extends Schema.Schema.Type<typeof MoveData> {}
export const MoveData = Schema.Struct({
location: Location.Ref,
projectID: Project.ID,
subpath: RelativePath.pipe(optional),
}).annotate({ identifier: "SessionPending.MoveData" })
export interface Move extends Schema.Schema.Type<typeof Move> {}
export const Move = Schema.Struct({
id: Event.ID,
sessionID: SessionID,
timeCreated: DateTimeUtcFromMillis,
type: Schema.tag("move"),
data: MoveData,
}).annotate({ identifier: "SessionPending.Move" })
export const Info = Schema.Union([User, Synthetic, Compaction, Move]).pipe(
Schema.toTaggedUnion("type"), Schema.toTaggedUnion("type"),
Schema.annotate({ identifier: "SessionPending.Info" }), Schema.annotate({ identifier: "SessionPending.Info" }),
) )
@@ -90,6 +90,16 @@ describe("contract hygiene", () => {
}) })
}) })
test("pending moves omit absent placement details", () => {
expect(
Schema.encodeSync(SessionPending.MoveData)({
location: { directory: AbsolutePath.make("/project"), workspaceID: undefined },
projectID: Project.ID.global,
subpath: undefined,
}),
).toEqual({ location: { directory: "/project" }, projectID: "global" })
})
test("forms require at least one field", () => { test("forms require at least one field", () => {
expect(() => expect(() =>
Schema.decodeUnknownSync(Form.Info)({ Schema.decodeUnknownSync(Form.Info)({
@@ -78,6 +78,7 @@ describe("public event manifest", () => {
"session.deleted.2", "session.deleted.2",
"session.agent.selected.1", "session.agent.selected.1",
"session.model.selected.1", "session.model.selected.1",
"session.move.admitted.1",
"session.moved.1", "session.moved.1",
"session.renamed.1", "session.renamed.1",
"session.usage.recorded.1", "session.usage.recorded.1",
+2 -1
View File
@@ -201,7 +201,8 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
function updatePending(sessionID: string, inputID: string, delivery: SessionPending.Delivery) { function updatePending(sessionID: string, inputID: string, delivery: SessionPending.Delivery) {
const index = store.session.pending[sessionID]?.findIndex((item) => item.id === inputID) ?? -1 const index = store.session.pending[sessionID]?.findIndex((item) => item.id === inputID) ?? -1
const item = store.session.pending[sessionID]?.[index] const item = store.session.pending[sessionID]?.[index]
if (index < 0 || !item || item.type === "compaction" || item.delivery === delivery) return if (index < 0 || !item || (item.type !== "user" && item.type !== "synthetic") || item.delivery === delivery)
return
setStore("session", "pending", sessionID, index, { ...item, delivery }) setStore("session", "pending", sessionID, index, { ...item, delivery })
} }
+12 -2
View File
@@ -982,7 +982,12 @@ test("updates and removes queued inputs from durable lifecycle events", async ()
await wait(() => await wait(() =>
data.session.pending data.session.pending
.list(sessionID) .list(sessionID)
.some((item) => item.id === "message-queued" && item.type !== "compaction" && item.delivery === "steer"), .some(
(item) =>
item.id === "message-queued" &&
(item.type === "user" || item.type === "synthetic") &&
item.delivery === "steer",
),
) )
expect(rows).toContainEqual({ type: "message", messageID: "message-queued" }) expect(rows).toContainEqual({ type: "message", messageID: "message-queued" })
@@ -996,7 +1001,12 @@ test("updates and removes queued inputs from durable lifecycle events", async ()
await wait(() => await wait(() =>
data.session.pending data.session.pending
.list(sessionID) .list(sessionID)
.some((item) => item.id === "message-queued" && item.type !== "compaction" && item.delivery === "queue"), .some(
(item) =>
item.id === "message-queued" &&
(item.type === "user" || item.type === "synthetic") &&
item.delivery === "queue",
),
) )
expect(rows).not.toContainEqual({ type: "message", messageID: "message-queued" }) expect(rows).not.toContainEqual({ type: "message", messageID: "message-queued" })