fix(tui): undo pending session input
This commit is contained in:
parent
9abd9594de
commit
330ab9acae
24 changed files with 493 additions and 71 deletions
|
|
@ -209,6 +209,10 @@ export interface Interface {
|
|||
* unhandled compaction barriers.
|
||||
*/
|
||||
readonly pending: (sessionID: SessionSchema.ID) => Effect.Effect<SessionPending.Info[], NotFoundError>
|
||||
readonly withdraw: (input: {
|
||||
sessionID: SessionSchema.ID
|
||||
inputID: SessionMessage.ID
|
||||
}) => Effect.Effect<boolean, NotFoundError>
|
||||
/**
|
||||
* Durable, ordered session log read. Replays durable session bus after
|
||||
* the exclusive `after` cursor, emits a `Synced` marker at the captured
|
||||
|
|
@ -540,6 +544,10 @@ const layer = Layer.effect(
|
|||
yield* result.get(sessionID)
|
||||
return yield* SessionPending.list(db, sessionID)
|
||||
}),
|
||||
withdraw: Effect.fn("Session.withdraw")(function* (input) {
|
||||
yield* result.get(input.sessionID)
|
||||
return yield* SessionPending.withdraw(db, bus, input)
|
||||
}),
|
||||
log: (input) =>
|
||||
Stream.unwrap(
|
||||
result
|
||||
|
|
|
|||
|
|
@ -175,6 +175,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
|
|||
"session.forked": () => Effect.void,
|
||||
"session.input.promoted": () => Effect.void,
|
||||
"session.input.admitted": () => Effect.void,
|
||||
"session.input.withdrawn": () => Effect.void,
|
||||
"session.execution.started": () => Effect.void,
|
||||
"session.execution.succeeded": () => clearCurrentRetry,
|
||||
"session.execution.failed": () => clearCurrentRetry,
|
||||
|
|
|
|||
|
|
@ -38,10 +38,15 @@ const encodeUser = Schema.encodeSync(UserData)
|
|||
const decodeSynthetic = Schema.decodeUnknownSync(SyntheticData)
|
||||
const encodeSynthetic = Schema.encodeSync(SyntheticData)
|
||||
const decodeAdmittedEvent = Schema.decodeUnknownOption(SessionEvent.InputAdmitted.data)
|
||||
const decodeWithdrawnEvent = Schema.decodeUnknownOption(SessionEvent.InputWithdrawn.data)
|
||||
const admittedEventType = Bus.versionedType(
|
||||
SessionEvent.InputAdmitted.type,
|
||||
SessionEvent.InputAdmitted.durable.version,
|
||||
)
|
||||
const withdrawnEventType = Bus.versionedType(
|
||||
SessionEvent.InputWithdrawn.type,
|
||||
SessionEvent.InputWithdrawn.durable.version,
|
||||
)
|
||||
const inboxLocks = KeyedMutex.makeUnsafe<SessionSchema.ID>()
|
||||
|
||||
export class LifecycleConflict extends Schema.TaggedErrorClass<LifecycleConflict>()(
|
||||
|
|
@ -103,26 +108,11 @@ export const compaction = Effect.fn("SessionPending.compaction")(function* (
|
|||
return entry.type === "compaction" ? entry : undefined
|
||||
})
|
||||
|
||||
/**
|
||||
* Reconstruct the admitted record for a pending row that was already consumed
|
||||
* by promotion. The projected `session_message` row proves promotion happened;
|
||||
* the durable `session.input.admitted` event retains the exact admitted
|
||||
* message, including delivery.
|
||||
*/
|
||||
const promotedFromHistory = Effect.fn("SessionPending.promotedFromHistory")(function* (
|
||||
const admittedFromHistory = Effect.fn("SessionPending.admittedFromHistory")(function* (
|
||||
db: DatabaseService,
|
||||
sessionID: SessionSchema.ID,
|
||||
id: SessionMessage.ID,
|
||||
) {
|
||||
const message = yield* db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(eq(SessionMessageTable.id, id))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (message === undefined) return undefined
|
||||
if (message.session_id !== sessionID || (message.type !== "user" && message.type !== "synthetic"))
|
||||
return yield* Effect.die(new LifecycleConflict({ id }))
|
||||
const rows = yield* db
|
||||
.select()
|
||||
.from(EventTable)
|
||||
|
|
@ -146,6 +136,46 @@ const promotedFromHistory = Effect.fn("SessionPending.promotedFromHistory")(func
|
|||
return yield* Effect.die(new LifecycleConflict({ id }))
|
||||
})
|
||||
|
||||
/**
|
||||
* Reconstruct the admitted record for a pending row that was already consumed
|
||||
* by promotion. The projected `session_message` row proves promotion happened;
|
||||
* the durable `session.input.admitted` event retains the exact admitted
|
||||
* message, including delivery.
|
||||
*/
|
||||
const promotedFromHistory = Effect.fn("SessionPending.promotedFromHistory")(function* (
|
||||
db: DatabaseService,
|
||||
sessionID: SessionSchema.ID,
|
||||
id: SessionMessage.ID,
|
||||
) {
|
||||
const message = yield* db
|
||||
.select()
|
||||
.from(SessionMessageTable)
|
||||
.where(eq(SessionMessageTable.id, id))
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (message === undefined) return undefined
|
||||
if (message.session_id !== sessionID || (message.type !== "user" && message.type !== "synthetic"))
|
||||
return yield* Effect.die(new LifecycleConflict({ id }))
|
||||
return yield* admittedFromHistory(db, sessionID, id)
|
||||
})
|
||||
|
||||
const wasWithdrawn = Effect.fn("SessionPending.wasWithdrawn")(function* (
|
||||
db: DatabaseService,
|
||||
sessionID: SessionSchema.ID,
|
||||
id: SessionMessage.ID,
|
||||
) {
|
||||
const rows = yield* db
|
||||
.select({ data: EventTable.data })
|
||||
.from(EventTable)
|
||||
.where(and(eq(EventTable.aggregate_id, sessionID), eq(EventTable.type, withdrawnEventType)))
|
||||
.all()
|
||||
.pipe(Effect.orDie)
|
||||
return rows.some((row) => {
|
||||
const decoded = decodeWithdrawnEvent(row.data)
|
||||
return decoded._tag === "Some" && decoded.value.inputID === id
|
||||
})
|
||||
})
|
||||
|
||||
export const admit = Effect.fn("SessionPending.admit")(function* (
|
||||
db: DatabaseService,
|
||||
bus: Bus.Interface,
|
||||
|
|
@ -162,6 +192,8 @@ export const admit = Effect.fn("SessionPending.admit")(function* (
|
|||
}
|
||||
const promoted = yield* promotedFromHistory(db, request.sessionID, request.id)
|
||||
if (promoted !== undefined) return promoted
|
||||
if (yield* wasWithdrawn(db, request.sessionID, request.id))
|
||||
return yield* admittedFromHistory(db, request.sessionID, request.id)
|
||||
return yield* bus
|
||||
.publish(SessionEvent.InputAdmitted, {
|
||||
inputID: request.id,
|
||||
|
|
@ -329,6 +361,25 @@ export const projectPromoted = Effect.fn("SessionPending.projectPromoted")(funct
|
|||
return stored
|
||||
})
|
||||
|
||||
export const projectWithdrawn = Effect.fn("SessionPending.projectWithdrawn")(function* (
|
||||
db: DatabaseService,
|
||||
input: {
|
||||
readonly id: SessionMessage.ID
|
||||
readonly sessionID: SessionSchema.ID
|
||||
},
|
||||
) {
|
||||
const deleted = yield* db
|
||||
.delete(SessionPendingTable)
|
||||
.where(and(eq(SessionPendingTable.id, input.id), eq(SessionPendingTable.session_id, input.sessionID)))
|
||||
.returning()
|
||||
.get()
|
||||
.pipe(Effect.orDie)
|
||||
if (!deleted) return yield* Effect.die(new LifecycleConflict({ id: input.id }))
|
||||
const stored = fromRow(deleted)
|
||||
if (stored.type === "compaction") return yield* Effect.die(new LifecycleConflict({ id: input.id }))
|
||||
return stored
|
||||
})
|
||||
|
||||
export const settleCompaction = Effect.fn("SessionPending.settleCompaction")(function* (
|
||||
db: DatabaseService,
|
||||
input: { readonly sessionID: SessionSchema.ID },
|
||||
|
|
@ -406,6 +457,30 @@ export const equivalent = (
|
|||
return false
|
||||
}
|
||||
|
||||
export const withdraw = Effect.fn("SessionPending.withdraw")(function* (
|
||||
db: DatabaseService,
|
||||
bus: Bus.Interface,
|
||||
input: { readonly sessionID: SessionSchema.ID; readonly inputID: SessionMessage.ID },
|
||||
) {
|
||||
return yield* inboxLocks.withLock(input.sessionID)(
|
||||
Effect.gen(function* () {
|
||||
const pending = yield* find(db, input.inputID)
|
||||
if (!pending) return yield* wasWithdrawn(db, input.sessionID, input.inputID)
|
||||
if (pending.sessionID !== input.sessionID || pending.type === "compaction") return false
|
||||
yield* bus
|
||||
.publish(SessionEvent.InputWithdrawn, input)
|
||||
.pipe(
|
||||
Effect.catchDefect((defect) =>
|
||||
wasWithdrawn(db, input.sessionID, input.inputID).pipe(
|
||||
Effect.flatMap((withdrawn) => (withdrawn ? Effect.void : Effect.die(defect))),
|
||||
),
|
||||
),
|
||||
)
|
||||
return true
|
||||
}),
|
||||
)
|
||||
})
|
||||
|
||||
const publish = Effect.fn("SessionPending.publish")(function* (
|
||||
db: DatabaseService,
|
||||
bus: Bus.Interface,
|
||||
|
|
|
|||
|
|
@ -665,6 +665,12 @@ const layer = Layer.effectDiscard(
|
|||
.pipe(Effect.orDie)
|
||||
}),
|
||||
)
|
||||
yield* bus.project(SessionEvent.InputWithdrawn, (event) =>
|
||||
SessionPending.projectWithdrawn(db, {
|
||||
id: event.data.inputID,
|
||||
sessionID: event.data.sessionID,
|
||||
}),
|
||||
)
|
||||
yield* bus.project(SessionEvent.Compaction.Admitted, (event) =>
|
||||
Effect.gen(function* () {
|
||||
if (event.durable === undefined)
|
||||
|
|
|
|||
|
|
@ -996,6 +996,52 @@ describe("Session.pending", () => {
|
|||
}),
|
||||
)
|
||||
|
||||
it.effect("withdraws an interrupted input before promotion without resurrecting exact retries", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
const session = yield* Session.Service
|
||||
|
||||
const admitted = yield* session.prompt({
|
||||
id: SessionMessage.ID.make("msg_withdrawn"),
|
||||
sessionID,
|
||||
text: "Withdraw me",
|
||||
resume: false,
|
||||
})
|
||||
|
||||
yield* session.interrupt(sessionID)
|
||||
expect(yield* session.withdraw({ sessionID, inputID: admitted.id })).toBe(true)
|
||||
expect(yield* session.pending(sessionID)).toEqual([])
|
||||
expect(yield* session.messages({ sessionID })).toEqual([])
|
||||
expect(yield* session.withdraw({ sessionID, inputID: admitted.id })).toBe(true)
|
||||
|
||||
const retried = yield* session.prompt({
|
||||
id: admitted.id,
|
||||
sessionID,
|
||||
text: "Withdraw me",
|
||||
resume: false,
|
||||
})
|
||||
expect(retried.id).toBe(admitted.id)
|
||||
expect(yield* session.pending(sessionID)).toEqual([])
|
||||
expect(yield* eventCount(Bus.versionedType(SessionEvent.InputAdmitted.type, 1))).toBe(1)
|
||||
expect(yield* eventCount(Bus.versionedType(SessionEvent.InputWithdrawn.type, 1))).toBe(1)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("leaves promoted input for revert when withdrawal loses the race", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
const session = yield* Session.Service
|
||||
const bus = yield* Bus.Service
|
||||
const { db } = yield* Database.Service
|
||||
const admitted = yield* session.prompt({ sessionID, text: "Promote me", resume: false })
|
||||
|
||||
yield* SessionPending.promote(db, bus, sessionID, "steer")
|
||||
|
||||
expect(yield* session.withdraw({ sessionID, inputID: admitted.id })).toBe(false)
|
||||
expect(yield* session.messages({ sessionID })).toMatchObject([{ id: admitted.id, type: "user" }])
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("lists an unhandled compaction barrier until it settles", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue