Skip to content

Commit fb43c15

Browse files
authored
refactor(core): simplify event model (anomalyco#33238)
1 parent ca006a2 commit fb43c15

15 files changed

Lines changed: 231 additions & 569 deletions

File tree

packages/core/src/event.ts

Lines changed: 135 additions & 174 deletions
Large diffs are not rendered by default.

packages/core/src/public/session.ts

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
export * as Session from "./session"
22

33
import { Effect, Schema, Stream } from "effect"
4-
import { EventV2 } from "../event"
54
import { ModelV2 } from "../model"
65
import { SessionV2 } from "../session"
76
import { MessageDecodeError } from "../session/error"
@@ -34,9 +33,7 @@ export type Delivery = SessionInput.Delivery
3433
export const ListInput = SessionV2.ListInput
3534
export type ListInput = SessionV2.ListInput
3635

37-
export const EventCursor = EventV2.Cursor
38-
export type EventCursor = EventV2.Cursor
39-
export type Event = EventV2.CursorEvent<SessionEvent.DurableEvent>
36+
export type Event = SessionEvent.DurableEvent
4037

4138
export const NotFoundError = SessionV2.NotFoundError
4239
export type NotFoundError = SessionV2.NotFoundError
@@ -99,7 +96,7 @@ export interface MessageInput {
9996

10097
export interface EventsInput {
10198
readonly sessionID: ID
102-
readonly after?: EventCursor
99+
readonly after?: number
103100
}
104101

105102
export interface Interface {

packages/core/src/session.ts

Lines changed: 6 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -124,8 +124,8 @@ export interface Interface {
124124
) => Effect.Effect<SessionMessage.Message[], NotFoundError | MessageDecodeError>
125125
readonly events: (input: {
126126
sessionID: SessionSchema.ID
127-
after?: EventV2.Cursor
128-
}) => Stream.Stream<EventV2.CursorEvent<SessionEvent.DurableEvent>, NotFoundError>
127+
after?: number
128+
}) => Stream.Stream<SessionEvent.DurableEvent, NotFoundError>
129129
readonly switchAgent: (input: {
130130
sessionID: SessionSchema.ID
131131
agent: string
@@ -339,11 +339,9 @@ export const layer = Layer.effect(
339339
Stream.unwrap(
340340
result
341341
.get(input.sessionID)
342-
.pipe(Effect.as(events.aggregateEvents({ aggregateID: input.sessionID, after: input.after }))),
342+
.pipe(Effect.as(events.durable({ aggregateID: input.sessionID, after: input.after }))),
343343
).pipe(
344-
Stream.filter((event): event is EventV2.CursorEvent<SessionEvent.DurableEvent> =>
345-
isDurableSessionEvent(event.event),
346-
),
344+
Stream.filter((event): event is SessionEvent.DurableEvent => isDurableSessionEvent(event)),
347345
),
348346
prompt: Effect.fn("V2Session.prompt")((input) =>
349347
Effect.uninterruptible(
@@ -413,9 +411,9 @@ export const layer = Layer.effect(
413411
sessionID,
414412
timestamp: yield* DateTime.now,
415413
})
416-
if (event.seq === undefined)
414+
if (event.durable === undefined)
417415
return yield* Effect.die("Interrupt request event is missing aggregate sequence")
418-
yield* execution.interrupt(sessionID, event.seq)
416+
yield* execution.interrupt(sessionID, event.durable.seq)
419417
}),
420418
),
421419
),

packages/core/src/session/event.ts

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,13 @@ const Base = {
2727
}
2828

2929
const options = {
30-
sync: {
30+
durable: {
3131
aggregate: "sessionID",
3232
version: 1,
3333
},
3434
} as const
3535
const stepSettlementOptions = {
36-
sync: {
36+
durable: {
3737
aggregate: "sessionID",
3838
version: 2,
3939
},
@@ -456,7 +456,7 @@ export namespace Compaction {
456456

457457
export const Ended = EventV2.define({
458458
type: "session.next.compaction.ended",
459-
sync: { aggregate: "sessionID", version: 2 },
459+
durable: { aggregate: "sessionID", version: 2 },
460460
schema: {
461461
...Base,
462462
messageID: SessionMessageID.ID,

packages/core/src/session/input.ts

Lines changed: 2 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -74,11 +74,11 @@ export const admit = Effect.fn("SessionInput.admit")(function* (
7474
})
7575
.pipe(
7676
Effect.flatMap((event) =>
77-
event.seq === undefined
77+
event.durable === undefined
7878
? Effect.die("Prompt admission event is missing aggregate sequence")
7979
: Effect.succeed(
8080
new Admitted({
81-
admittedSeq: event.seq,
81+
admittedSeq: event.durable.seq,
8282
id: input.id,
8383
sessionID: input.sessionID,
8484
prompt: input.prompt,
@@ -117,13 +117,6 @@ export const projectAdmitted = Effect.fn("SessionInput.projectAdmitted")(functio
117117
readonly timeCreated: DateTime.Utc
118118
},
119119
) {
120-
const message = yield* db
121-
.select({ id: SessionMessageTable.id })
122-
.from(SessionMessageTable)
123-
.where(eq(SessionMessageTable.id, input.id))
124-
.get()
125-
.pipe(Effect.orDie)
126-
if (message) return yield* Effect.die(new LifecycleConflict({ id: input.id }))
127120
const stored = yield* db
128121
.insert(SessionInputTable)
129122
.values({
@@ -208,37 +201,6 @@ const matchesPrompt = (input: Admitted, expected: { readonly sessionID: SessionS
208201
input.sessionID === expected.sessionID &&
209202
JSON.stringify(encodePrompt(input.prompt)) === JSON.stringify(encodePrompt(expected.prompt))
210203

211-
export const guardReservedID = Effect.fn("SessionInput.guardReservedID")(function* (
212-
db: DatabaseService,
213-
event: EventV2.Payload,
214-
) {
215-
if (
216-
Schema.is(SessionEvent.PromptLifecycle.Admitted)(event) ||
217-
Schema.is(SessionEvent.PromptLifecycle.Promoted)(event)
218-
)
219-
return
220-
const id = reservedID(event)
221-
if (id === undefined) return
222-
const admitted = yield* db
223-
.select({ id: SessionInputTable.id })
224-
.from(SessionInputTable)
225-
.where(eq(SessionInputTable.id, id))
226-
.get()
227-
.pipe(Effect.orDie)
228-
if (admitted === undefined) return
229-
return yield* Effect.die(new LifecycleConflict({ id }))
230-
})
231-
232-
const reservedID = (event: EventV2.Payload) => {
233-
if (Schema.is(SessionEvent.Step.Started)(event)) return event.data.assistantMessageID
234-
if (Schema.is(SessionEvent.AgentSwitched)(event)) return event.data.messageID
235-
if (Schema.is(SessionEvent.ModelSwitched)(event)) return event.data.messageID
236-
if (Schema.is(SessionEvent.Prompted)(event)) return event.data.messageID
237-
if (Schema.is(SessionEvent.Synthetic)(event)) return event.data.messageID
238-
if (Schema.is(SessionEvent.Shell.Started)(event)) return event.data.messageID
239-
if (Schema.is(SessionEvent.Compaction.Started)(event)) return event.data.messageID
240-
}
241-
242204
export const projectLegacyPrompted = Effect.fn("SessionInput.projectLegacyPrompted")(function* (
243205
db: DatabaseService,
244206
input: {

packages/core/src/session/projector.ts

Lines changed: 22 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -115,7 +115,7 @@ function run(db: DatabaseService, event: SessionEvent.Event) {
115115
const decodeRow = (row: typeof SessionMessageTable.$inferSelect) =>
116116
decodeMessage({ ...row.data, id: row.id, type: row.type })
117117
const updateMessage = (message: SessionMessage.Message) => {
118-
if (event.seq === undefined) return Effect.die("Synchronized Session event is missing aggregate sequence")
118+
if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence")
119119
const encoded = encodeMessage(message)
120120
const { id, type, ...data } = encoded
121121
return db
@@ -192,7 +192,7 @@ function run(db: DatabaseService, event: SessionEvent.Event) {
192192
}
193193

194194
function insertMessage(db: DatabaseService, event: SessionEvent.Event, message: SessionMessage.Message) {
195-
if (event.seq === undefined) return Effect.die("Synchronized Session event is missing aggregate sequence")
195+
if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence")
196196
const encoded = encodeMessage(message)
197197
const { id, type, ...data } = encoded
198198
return db
@@ -201,7 +201,7 @@ function insertMessage(db: DatabaseService, event: SessionEvent.Event, message:
201201
id: SessionMessage.ID.make(id),
202202
session_id: event.data.sessionID,
203203
type,
204-
seq: event.seq,
204+
seq: event.durable.seq,
205205
time_created: DateTime.toEpochMillis(message.time.created),
206206
data,
207207
})
@@ -213,7 +213,6 @@ export const layer = Layer.effectDiscard(
213213
Effect.gen(function* () {
214214
const events = yield* EventV2.Service
215215
const { db } = yield* Database.Service
216-
yield* events.beforeCommit((event) => SessionInput.guardReservedID(db, event))
217216
yield* events.project(SessionV1.Event.Created, (event) =>
218217
Effect.gen(function* () {
219218
const stored = yield* db
@@ -331,7 +330,7 @@ export const layer = Layer.effectDiscard(
331330
}),
332331
)
333332
yield* events.project(SessionEvent.AgentSwitched, (event) => {
334-
if (event.seq === undefined) return Effect.die("Synchronized Session event is missing aggregate sequence")
333+
if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence")
335334
return db
336335
.update(SessionTable)
337336
.set({ agent: event.data.agent, time_updated: DateTime.toEpochMillis(event.data.timestamp) })
@@ -340,7 +339,7 @@ export const layer = Layer.effectDiscard(
340339
.pipe(
341340
Effect.orDie,
342341
Effect.andThen(run(db, event)),
343-
Effect.andThen(SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.seq)),
342+
Effect.andThen(SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.durable.seq)),
344343
)
345344
})
346345
yield* events.project(SessionEvent.ModelSwitched, (event) =>
@@ -352,9 +351,9 @@ export const layer = Layer.effectDiscard(
352351
.run()
353352
.pipe(Effect.orDie)
354353
yield* run(db, event)
355-
if (event.seq === undefined)
356-
return yield* Effect.die("Synchronized Session event is missing aggregate sequence")
357-
yield* SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.seq)
354+
if (event.durable === undefined)
355+
return yield* Effect.die("Durable Session event is missing aggregate sequence")
356+
yield* SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.durable.seq)
358357
}),
359358
)
360359
yield* events.project(SessionEvent.Prompted, (event) =>
@@ -368,24 +367,24 @@ export const layer = Layer.effectDiscard(
368367
.pipe(Effect.orDie)
369368
if (existing) return yield* Effect.die(new PromptAlreadyProjected())
370369
yield* run(db, event)
371-
if (event.seq === undefined)
372-
return yield* Effect.die("Synchronized Session event is missing aggregate sequence")
370+
if (event.durable === undefined)
371+
return yield* Effect.die("Durable Session event is missing aggregate sequence")
373372
yield* SessionInput.projectLegacyPrompted(db, {
374373
id: messageID,
375374
sessionID: event.data.sessionID,
376375
prompt: event.data.prompt,
377376
delivery: event.data.delivery,
378377
timeCreated: event.data.timestamp,
379-
promotedSeq: event.seq,
378+
promotedSeq: event.durable.seq,
380379
})
381380
}),
382381
)
383382
yield* events.project(SessionEvent.PromptLifecycle.Admitted, (event) =>
384383
Effect.gen(function* () {
385-
if (event.seq === undefined)
386-
return yield* Effect.die("Synchronized Session event is missing aggregate sequence")
384+
if (event.durable === undefined)
385+
return yield* Effect.die("Durable Session event is missing aggregate sequence")
387386
yield* SessionInput.projectAdmitted(db, {
388-
admittedSeq: event.seq,
387+
admittedSeq: event.durable.seq,
389388
id: event.data.messageID,
390389
sessionID: event.data.sessionID,
391390
prompt: event.data.prompt,
@@ -396,8 +395,8 @@ export const layer = Layer.effectDiscard(
396395
)
397396
yield* events.project(SessionEvent.PromptLifecycle.Promoted, (event) =>
398397
Effect.gen(function* () {
399-
if (event.seq === undefined)
400-
return yield* Effect.die("Synchronized Session event is missing aggregate sequence")
398+
if (event.durable === undefined)
399+
return yield* Effect.die("Durable Session event is missing aggregate sequence")
401400
yield* insertMessage(
402401
db,
403402
event,
@@ -406,18 +405,14 @@ export const layer = Layer.effectDiscard(
406405
sessionID: event.data.sessionID,
407406
prompt: event.data.prompt,
408407
timeCreated: event.data.timeCreated,
409-
promotedSeq: event.seq,
408+
promotedSeq: event.durable.seq,
410409
}),
411410
)
412411
}),
413412
)
414413
yield* events.project(SessionEvent.InterruptRequested, () => Effect.void)
415-
yield* events.project(SessionEvent.ContextUpdated, (event) => {
416-
if (!event.replay || event.seq === undefined) return run(db, event)
417-
return run(db, event).pipe(
418-
Effect.andThen(SessionContextEpoch.requestReplacement(db, event.data.sessionID, event.seq)),
419-
)
420-
})
414+
// TODO: Reconstruct context epoch replacement state during replay without adding replay state to every EventV2 payload.
415+
yield* events.project(SessionEvent.ContextUpdated, (event) => run(db, event))
421416
yield* events.project(SessionEvent.Synthetic, (event) => run(db, event))
422417
yield* events.project(SessionEvent.Shell.Started, (event) => run(db, event))
423418
yield* events.project(SessionEvent.Shell.Ended, (event) => run(db, event))
@@ -436,9 +431,9 @@ export const layer = Layer.effectDiscard(
436431
yield* events.project(SessionEvent.Reasoning.Ended, (event) => run(db, event))
437432
// yield* events.project(SessionEvent.Retried, (event) => run(db, event))
438433
yield* events.project(SessionEvent.Compaction.Ended, (event) => {
439-
if (event.version === 1) return Effect.void
440-
const seq = event.seq
441-
if (seq === undefined) return Effect.die("Synchronized Session event is missing aggregate sequence")
434+
if (event.durable === undefined) return Effect.die("Durable Session event is missing aggregate sequence")
435+
if (event.durable.version === 1) return Effect.void
436+
const seq = event.durable.seq
442437
return Effect.gen(function* () {
443438
yield* run(db, event)
444439
yield* SessionContextEpoch.requestReplacement(db, event.data.sessionID, seq)

packages/core/src/v1/session.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -502,7 +502,7 @@ export type WithParts = {
502502
}
503503

504504
const options = {
505-
sync: {
505+
durable: {
506506
aggregate: "sessionID",
507507
version: 1,
508508
},

0 commit comments

Comments
 (0)