Compare commits
1 commit
dev
...
normalize-
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
179a8748cb |
6 changed files with 60 additions and 16 deletions
|
|
@ -32,12 +32,6 @@ const options = {
|
||||||
version: 1,
|
version: 1,
|
||||||
},
|
},
|
||||||
} as const
|
} as const
|
||||||
const stepSettlementOptions = {
|
|
||||||
durable: {
|
|
||||||
aggregate: "sessionID",
|
|
||||||
version: 2,
|
|
||||||
},
|
|
||||||
} as const
|
|
||||||
|
|
||||||
export const UnknownError = Schema.Struct({
|
export const UnknownError = Schema.Struct({
|
||||||
type: Schema.Literal("unknown"),
|
type: Schema.Literal("unknown"),
|
||||||
|
|
@ -181,7 +175,7 @@ export namespace Step {
|
||||||
|
|
||||||
export const Ended = EventV2.define({
|
export const Ended = EventV2.define({
|
||||||
type: "session.next.step.ended",
|
type: "session.next.step.ended",
|
||||||
...stepSettlementOptions,
|
...options,
|
||||||
schema: {
|
schema: {
|
||||||
...Base,
|
...Base,
|
||||||
assistantMessageID: SessionMessageID.ID,
|
assistantMessageID: SessionMessageID.ID,
|
||||||
|
|
@ -203,7 +197,7 @@ export namespace Step {
|
||||||
|
|
||||||
export const Failed = EventV2.define({
|
export const Failed = EventV2.define({
|
||||||
type: "session.next.step.failed",
|
type: "session.next.step.failed",
|
||||||
...stepSettlementOptions,
|
...options,
|
||||||
schema: {
|
schema: {
|
||||||
...Base,
|
...Base,
|
||||||
assistantMessageID: SessionMessageID.ID,
|
assistantMessageID: SessionMessageID.ID,
|
||||||
|
|
|
||||||
|
|
@ -41,6 +41,55 @@ const assistantRow = (
|
||||||
}
|
}
|
||||||
|
|
||||||
describe("SessionProjector", () => {
|
describe("SessionProjector", () => {
|
||||||
|
it.effect("stores step settlement events as version 1", () =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const { db } = yield* Database.Service
|
||||||
|
yield* db
|
||||||
|
.insert(ProjectTable)
|
||||||
|
.values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
|
||||||
|
.run()
|
||||||
|
.pipe(Effect.orDie)
|
||||||
|
yield* db
|
||||||
|
.insert(SessionTable)
|
||||||
|
.values({
|
||||||
|
id: sessionID,
|
||||||
|
project_id: Project.ID.global,
|
||||||
|
slug: "test",
|
||||||
|
directory: "/project",
|
||||||
|
title: "test",
|
||||||
|
version: "test",
|
||||||
|
})
|
||||||
|
.run()
|
||||||
|
.pipe(Effect.orDie)
|
||||||
|
|
||||||
|
const events = yield* EventV2.Service
|
||||||
|
const assistantMessageID = SessionMessage.ID.make("msg_assistant")
|
||||||
|
yield* events.publish(SessionEvent.Step.Ended, {
|
||||||
|
sessionID,
|
||||||
|
assistantMessageID,
|
||||||
|
timestamp: created,
|
||||||
|
finish: "stop",
|
||||||
|
cost: 0,
|
||||||
|
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
|
||||||
|
})
|
||||||
|
yield* events.publish(SessionEvent.Step.Failed, {
|
||||||
|
sessionID,
|
||||||
|
assistantMessageID,
|
||||||
|
timestamp: created,
|
||||||
|
error: { type: "unknown", message: "failed" },
|
||||||
|
})
|
||||||
|
|
||||||
|
expect(
|
||||||
|
(yield* db
|
||||||
|
.select({ type: EventTable.type })
|
||||||
|
.from(EventTable)
|
||||||
|
.where(eq(EventTable.aggregate_id, sessionID))
|
||||||
|
.orderBy(EventTable.seq)
|
||||||
|
.all()).map((event) => event.type),
|
||||||
|
).toEqual(["session.next.step.ended.1", "session.next.step.failed.1"])
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
it.effect("orders projected messages and context by durable aggregate sequence", () =>
|
it.effect("orders projected messages and context by durable aggregate sequence", () =>
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const { db } = yield* Database.Service
|
const { db } = yield* Database.Service
|
||||||
|
|
|
||||||
|
|
@ -180,7 +180,7 @@ describe("SessionRunnerLLM recorded", () => {
|
||||||
"session.next.step.started.1",
|
"session.next.step.started.1",
|
||||||
"session.next.text.started.1",
|
"session.next.text.started.1",
|
||||||
"session.next.text.ended.1",
|
"session.next.text.ended.1",
|
||||||
"session.next.step.ended.2",
|
"session.next.step.ended.1",
|
||||||
])
|
])
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -3334,7 +3334,7 @@ export type SyncEventSessionNextStepEnded = {
|
||||||
type: "sync"
|
type: "sync"
|
||||||
id: string
|
id: string
|
||||||
syncEvent: {
|
syncEvent: {
|
||||||
type: "session.next.step.ended.2"
|
type: "session.next.step.ended.1"
|
||||||
id: string
|
id: string
|
||||||
seq: number
|
seq: number
|
||||||
aggregateID: string
|
aggregateID: string
|
||||||
|
|
@ -3362,7 +3362,7 @@ export type SyncEventSessionNextStepFailed = {
|
||||||
type: "sync"
|
type: "sync"
|
||||||
id: string
|
id: string
|
||||||
syncEvent: {
|
syncEvent: {
|
||||||
type: "session.next.step.failed.2"
|
type: "session.next.step.failed.1"
|
||||||
id: string
|
id: string
|
||||||
seq: number
|
seq: number
|
||||||
aggregateID: string
|
aggregateID: string
|
||||||
|
|
|
||||||
|
|
@ -24782,7 +24782,7 @@
|
||||||
"properties": {
|
"properties": {
|
||||||
"type": {
|
"type": {
|
||||||
"type": "string",
|
"type": "string",
|
||||||
"enum": ["session.next.step.ended.2"]
|
"enum": ["session.next.step.ended.1"]
|
||||||
},
|
},
|
||||||
"id": {
|
"id": {
|
||||||
"type": "string",
|
"type": "string",
|
||||||
|
|
@ -24874,7 +24874,7 @@
|
||||||
"properties": {
|
"properties": {
|
||||||
"type": {
|
"type": {
|
||||||
"type": "string",
|
"type": "string",
|
||||||
"enum": ["session.next.step.failed.2"]
|
"enum": ["session.next.step.failed.1"]
|
||||||
},
|
},
|
||||||
"id": {
|
"id": {
|
||||||
"type": "string",
|
"type": "string",
|
||||||
|
|
|
||||||
|
|
@ -3,6 +3,7 @@
|
||||||
## 2026-06-22: Reset Unpublished Compaction Event
|
## 2026-06-22: Reset Unpublished Compaction Event
|
||||||
|
|
||||||
- Replace the unpublished `session.next.compaction.ended.1` payload with the current checkpoint payload and remove its legacy decoder.
|
- Replace the unpublished `session.next.compaction.ended.1` payload with the current checkpoint payload and remove its legacy decoder.
|
||||||
|
- Reset the current `session.next.step.ended` and `session.next.step.failed` payloads to version `1` now that incompatible experimental history has been discarded.
|
||||||
- Reset experimental events, sequences, Session inputs, projected Session messages, Context Epochs, synchronized workspace rows, and Session workspace links.
|
- Reset experimental events, sequences, Session inputs, projected Session messages, Context Epochs, synchronized workspace rows, and Session workspace links.
|
||||||
- Preserve canonical V1 `session`, `message`, and `part` rows.
|
- Preserve canonical V1 `session`, `message`, and `part` rows.
|
||||||
|
|
||||||
|
|
@ -96,11 +97,11 @@ Compatibility:
|
||||||
- Existing generated Session and Event IDs retain their current prefixes and generation behavior.
|
- Existing generated Session and Event IDs retain their current prefixes and generation behavior.
|
||||||
- Deterministic constructors are additive internal helpers; public ID schemas remain strings with their existing prefixes.
|
- Deterministic constructors are additive internal helpers; public ID schemas remain strings with their existing prefixes.
|
||||||
|
|
||||||
### Durable Step Settlement Ownership
|
### Earlier Durable Step Settlement Ownership
|
||||||
|
|
||||||
Affected schema:
|
Affected schema:
|
||||||
|
|
||||||
- `session.next.step.ended` and `session.next.step.failed` synchronized event version `2`.
|
- Earlier experimental `session.next.step.ended` and `session.next.step.failed` synchronized event version `2`.
|
||||||
|
|
||||||
Change:
|
Change:
|
||||||
|
|
||||||
|
|
@ -112,7 +113,7 @@ Reason:
|
||||||
|
|
||||||
Compatibility:
|
Compatibility:
|
||||||
|
|
||||||
- Step settlement uses synchronized event version `2` because the durable payload changed.
|
- The payload originally used synchronized event version `2`; the June 22 reset later made the current payload version `1`.
|
||||||
|
|
||||||
### Durable Session Input Inbox
|
### Durable Session Input Inbox
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue