mirror of
https://github.com/anomalyco/opencode.git
synced 2026-08-08 18:30:00 -04:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fd229f7edf | |||
| db664db5f3 | |||
| 59e8dbb066 |
@@ -572,6 +572,7 @@ export type SessionsContextOutput = {
|
|||||||
>
|
>
|
||||||
readonly snapshot?: { readonly start?: string; readonly end?: string; readonly files?: ReadonlyArray<string> }
|
readonly snapshot?: { readonly start?: string; readonly end?: string; readonly files?: ReadonlyArray<string> }
|
||||||
readonly finish?: string
|
readonly finish?: string
|
||||||
|
readonly settlement?: "completed" | "failed" | "interrupted"
|
||||||
readonly cost?: number
|
readonly cost?: number
|
||||||
readonly tokens?: {
|
readonly tokens?: {
|
||||||
readonly input: number
|
readonly input: number
|
||||||
@@ -795,6 +796,14 @@ export type SessionsEventsOutput =
|
|||||||
readonly error: { readonly type: "unknown"; readonly message: string }
|
readonly error: { readonly type: "unknown"; readonly message: string }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
| {
|
||||||
|
readonly id: string
|
||||||
|
readonly metadata?: { readonly [x: string]: unknown }
|
||||||
|
readonly type: "session.next.step.interrupted"
|
||||||
|
readonly durable?: { readonly aggregateID: string; readonly seq: number; readonly version: number }
|
||||||
|
readonly location?: { readonly directory: string; readonly workspaceID?: string }
|
||||||
|
readonly data: { readonly timestamp: number; readonly sessionID: string; readonly assistantMessageID: string }
|
||||||
|
}
|
||||||
| {
|
| {
|
||||||
readonly id: string
|
readonly id: string
|
||||||
readonly metadata?: { readonly [x: string]: unknown }
|
readonly metadata?: { readonly [x: string]: unknown }
|
||||||
@@ -1189,6 +1198,7 @@ export type SessionsMessageOutput = {
|
|||||||
>
|
>
|
||||||
readonly snapshot?: { readonly start?: string; readonly end?: string; readonly files?: ReadonlyArray<string> }
|
readonly snapshot?: { readonly start?: string; readonly end?: string; readonly files?: ReadonlyArray<string> }
|
||||||
readonly finish?: string
|
readonly finish?: string
|
||||||
|
readonly settlement?: "completed" | "failed" | "interrupted"
|
||||||
readonly cost?: number
|
readonly cost?: number
|
||||||
readonly tokens?: {
|
readonly tokens?: {
|
||||||
readonly input: number
|
readonly input: number
|
||||||
|
|||||||
@@ -210,6 +210,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
|
|||||||
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
|
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
|
||||||
draft.time.completed = event.data.timestamp
|
draft.time.completed = event.data.timestamp
|
||||||
draft.finish = event.data.finish
|
draft.finish = event.data.finish
|
||||||
|
draft.settlement = "completed"
|
||||||
draft.cost = event.data.cost
|
draft.cost = event.data.cost
|
||||||
draft.tokens = event.data.tokens
|
draft.tokens = event.data.tokens
|
||||||
if (event.data.snapshot || event.data.files)
|
if (event.data.snapshot || event.data.files)
|
||||||
@@ -224,9 +225,16 @@ export function update(adapter: Adapter, event: SessionEvent.Event) {
|
|||||||
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
|
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
|
||||||
draft.time.completed = event.data.timestamp
|
draft.time.completed = event.data.timestamp
|
||||||
draft.finish = "error"
|
draft.finish = "error"
|
||||||
|
draft.settlement = "failed"
|
||||||
draft.error = event.data.error
|
draft.error = event.data.error
|
||||||
})
|
})
|
||||||
},
|
},
|
||||||
|
"session.next.step.interrupted": (event) => {
|
||||||
|
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
|
||||||
|
draft.time.completed = event.data.timestamp
|
||||||
|
draft.settlement = "interrupted"
|
||||||
|
})
|
||||||
|
},
|
||||||
"session.next.text.started": (event) => {
|
"session.next.text.started": (event) => {
|
||||||
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
|
return updateOwnedAssistant(event.data.assistantMessageID, (draft) => {
|
||||||
draft.content.push(
|
draft.content.push(
|
||||||
|
|||||||
@@ -381,6 +381,7 @@ export const layer = Layer.effectDiscard(
|
|||||||
yield* events.project(SessionEvent.Step.Started, (event) => run(db, event))
|
yield* events.project(SessionEvent.Step.Started, (event) => run(db, event))
|
||||||
yield* events.project(SessionEvent.Step.Ended, (event) => run(db, event))
|
yield* events.project(SessionEvent.Step.Ended, (event) => run(db, event))
|
||||||
yield* events.project(SessionEvent.Step.Failed, (event) => run(db, event))
|
yield* events.project(SessionEvent.Step.Failed, (event) => run(db, event))
|
||||||
|
yield* events.project(SessionEvent.Step.Interrupted, (event) => run(db, event))
|
||||||
yield* events.project(SessionEvent.Text.Started, (event) => run(db, event))
|
yield* events.project(SessionEvent.Text.Started, (event) => run(db, event))
|
||||||
yield* events.project(SessionEvent.Text.Ended, (event) => run(db, event))
|
yield* events.project(SessionEvent.Text.Ended, (event) => run(db, event))
|
||||||
yield* events.project(SessionEvent.Tool.Input.Started, (event) => run(db, event))
|
yield* events.project(SessionEvent.Tool.Input.Started, (event) => run(db, event))
|
||||||
|
|||||||
@@ -290,6 +290,7 @@ export const layer = Layer.effect(
|
|||||||
if (settled._tag === "Failure" && isQuestionRejected(settled.cause)) {
|
if (settled._tag === "Failure" && isQuestionRejected(settled.cause)) {
|
||||||
yield* FiberSet.clear(toolFibers)
|
yield* FiberSet.clear(toolFibers)
|
||||||
yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted"))
|
yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted"))
|
||||||
|
yield* withPublication(publisher.interruptAssistant())
|
||||||
return yield* Effect.interrupt
|
return yield* Effect.interrupt
|
||||||
}
|
}
|
||||||
if (
|
if (
|
||||||
@@ -298,8 +299,7 @@ export const layer = Layer.effect(
|
|||||||
) {
|
) {
|
||||||
yield* FiberSet.clear(toolFibers)
|
yield* FiberSet.clear(toolFibers)
|
||||||
yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted"))
|
yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted"))
|
||||||
if (publisher.hasActiveAssistant())
|
yield* withPublication(publisher.interruptAssistant())
|
||||||
yield* withPublication(publisher.failAssistant("Provider turn interrupted"))
|
|
||||||
}
|
}
|
||||||
if (settled._tag === "Failure" && !Cause.hasInterrupts(settled.cause)) {
|
if (settled._tag === "Failure" && !Cause.hasInterrupts(settled.cause)) {
|
||||||
const failure = Cause.squash(settled.cause)
|
const failure = Cause.squash(settled.cause)
|
||||||
@@ -307,7 +307,7 @@ export const layer = Layer.effect(
|
|||||||
yield* withPublication(publisher.failUnsettledTools(`Tool execution failed: ${message}`))
|
yield* withPublication(publisher.failUnsettledTools(`Tool execution failed: ${message}`))
|
||||||
}
|
}
|
||||||
const stepSettlement = publisher.stepSettlement()
|
const stepSettlement = publisher.stepSettlement()
|
||||||
if (stepSettlement && !publisher.hasProviderError()) {
|
if (stepSettlement && !publisher.hasProviderError() && !publisher.hasAssistantSettled()) {
|
||||||
const endSnapshot = yield* snapshots.capture()
|
const endSnapshot = yield* snapshots.capture()
|
||||||
const files =
|
const files =
|
||||||
startSnapshot && endSnapshot
|
startSnapshot && endSnapshot
|
||||||
|
|||||||
@@ -66,15 +66,13 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
|
|||||||
>()
|
>()
|
||||||
const timestamp = DateTime.now
|
const timestamp = DateTime.now
|
||||||
let assistantMessageID: SessionMessage.ID | undefined
|
let assistantMessageID: SessionMessage.ID | undefined
|
||||||
let assistantActive = false
|
let assistantSettled = false
|
||||||
let assistantFailed = false
|
|
||||||
let providerFailed = false
|
let providerFailed = false
|
||||||
let stepSettlement: { readonly finish: string; readonly tokens: ReturnType<typeof tokens> } | undefined
|
let stepSettlement: { readonly finish: SessionMessage.Finish; readonly tokens: ReturnType<typeof tokens> } | undefined
|
||||||
|
|
||||||
const startAssistant = Effect.fnUntraced(function* () {
|
const startAssistant = Effect.fnUntraced(function* () {
|
||||||
if (assistantMessageID !== undefined) return assistantMessageID
|
if (assistantMessageID !== undefined) return assistantMessageID
|
||||||
assistantMessageID = SessionMessage.ID.create()
|
assistantMessageID = SessionMessage.ID.create()
|
||||||
assistantActive = true
|
|
||||||
yield* events.publish(SessionEvent.Step.Started, {
|
yield* events.publish(SessionEvent.Step.Started, {
|
||||||
...input,
|
...input,
|
||||||
assistantMessageID,
|
assistantMessageID,
|
||||||
@@ -196,20 +194,39 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
|
|||||||
yield* flushFragments()
|
yield* flushFragments()
|
||||||
})
|
})
|
||||||
|
|
||||||
const failAssistant = Effect.fnUntraced(function* (message: string) {
|
const settleAssistant = Effect.fnUntraced(function* (
|
||||||
if (assistantFailed) return
|
publish: (assistantMessageID: SessionMessage.ID) => Effect.Effect<void>,
|
||||||
|
) {
|
||||||
|
if (assistantSettled) return
|
||||||
yield* flush()
|
yield* flush()
|
||||||
const assistantMessageID = yield* startAssistant()
|
const assistantMessageID = yield* startAssistant()
|
||||||
assistantActive = false
|
assistantSettled = true
|
||||||
assistantFailed = true
|
yield* publish(assistantMessageID)
|
||||||
yield* events.publish(SessionEvent.Step.Failed, {
|
|
||||||
sessionID: input.sessionID,
|
|
||||||
timestamp: yield* timestamp,
|
|
||||||
assistantMessageID,
|
|
||||||
error: { type: "unknown", message },
|
|
||||||
})
|
|
||||||
})
|
})
|
||||||
|
|
||||||
|
const failAssistant = (message: string) =>
|
||||||
|
settleAssistant((assistantMessageID) =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
yield* events.publish(SessionEvent.Step.Failed, {
|
||||||
|
sessionID: input.sessionID,
|
||||||
|
timestamp: yield* timestamp,
|
||||||
|
assistantMessageID,
|
||||||
|
error: { type: "unknown", message },
|
||||||
|
})
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
const interruptAssistant = () =>
|
||||||
|
settleAssistant((assistantMessageID) =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
yield* events.publish(SessionEvent.Step.Interrupted, {
|
||||||
|
sessionID: input.sessionID,
|
||||||
|
timestamp: yield* timestamp,
|
||||||
|
assistantMessageID,
|
||||||
|
})
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
const failUnsettledTools = Effect.fn("SessionRunner.failUnsettledTools")(function* (
|
const failUnsettledTools = Effect.fn("SessionRunner.failUnsettledTools")(function* (
|
||||||
message: string,
|
message: string,
|
||||||
hostedOnly = false,
|
hostedOnly = false,
|
||||||
@@ -395,7 +412,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
|
|||||||
}
|
}
|
||||||
case "step-finish":
|
case "step-finish":
|
||||||
yield* flush()
|
yield* flush()
|
||||||
assistantActive = false
|
|
||||||
if (stepSettlement) return yield* Effect.die("Duplicate step finish")
|
if (stepSettlement) return yield* Effect.die("Duplicate step finish")
|
||||||
stepSettlement = { finish: event.reason, tokens: tokens(event.usage) }
|
stepSettlement = { finish: event.reason, tokens: tokens(event.usage) }
|
||||||
return
|
return
|
||||||
@@ -412,9 +428,10 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
|
|||||||
publish,
|
publish,
|
||||||
flush,
|
flush,
|
||||||
failAssistant,
|
failAssistant,
|
||||||
|
interruptAssistant,
|
||||||
failUnsettledTools,
|
failUnsettledTools,
|
||||||
hasActiveAssistant: () => assistantActive,
|
|
||||||
hasAssistantStarted: () => assistantMessageID !== undefined,
|
hasAssistantStarted: () => assistantMessageID !== undefined,
|
||||||
|
hasAssistantSettled: () => assistantSettled,
|
||||||
hasProviderError: () => providerFailed,
|
hasProviderError: () => providerFailed,
|
||||||
stepSettlement: () => stepSettlement,
|
stepSettlement: () => stepSettlement,
|
||||||
startAssistant,
|
startAssistant,
|
||||||
|
|||||||
@@ -70,7 +70,7 @@ const toolResult = (tool: SessionMessage.AssistantTool, providerMetadata: Provid
|
|||||||
const assistant = (message: SessionMessage.Assistant, model: Model) => {
|
const assistant = (message: SessionMessage.Assistant, model: Model) => {
|
||||||
const sameModel =
|
const sameModel =
|
||||||
String(message.model.providerID) === String(model.provider) && String(message.model.id) === String(model.id)
|
String(message.model.providerID) === String(model.provider) && String(message.model.id) === String(model.id)
|
||||||
const reuseProviderMetadata = sameModel && message.error === undefined
|
const reuseProviderMetadata = sameModel && message.error === undefined && message.settlement !== "interrupted"
|
||||||
const content = message.content.flatMap((item): ContentPart[] => {
|
const content = message.content.flatMap((item): ContentPart[] => {
|
||||||
if (item.type === "text") return [{ type: "text", text: item.text }]
|
if (item.type === "text") return [{ type: "text", text: item.text }]
|
||||||
if (item.type === "reasoning")
|
if (item.type === "reasoning")
|
||||||
|
|||||||
@@ -327,7 +327,7 @@ Recent work
|
|||||||
])
|
])
|
||||||
})
|
})
|
||||||
|
|
||||||
test("drops provider-native continuation metadata from failed assistant turns", () => {
|
test("drops provider-native continuation metadata from interrupted assistant turns", () => {
|
||||||
const messages = toLLMMessages(
|
const messages = toLLMMessages(
|
||||||
[
|
[
|
||||||
SessionMessage.Assistant.make({
|
SessionMessage.Assistant.make({
|
||||||
@@ -361,8 +361,7 @@ Recent work
|
|||||||
time: { created, completed: created },
|
time: { created, completed: created },
|
||||||
}),
|
}),
|
||||||
],
|
],
|
||||||
finish: "error",
|
settlement: "interrupted",
|
||||||
error: { type: "unknown", message: "Provider turn interrupted" },
|
|
||||||
time: { created, completed: created },
|
time: { created, completed: created },
|
||||||
}),
|
}),
|
||||||
],
|
],
|
||||||
|
|||||||
@@ -519,6 +519,7 @@ const verifyPartialFlushOnFailure = (kind: FragmentKind) =>
|
|||||||
{
|
{
|
||||||
type: "assistant",
|
type: "assistant",
|
||||||
finish: "error",
|
finish: "error",
|
||||||
|
settlement: "failed",
|
||||||
error: { type: "unknown", message: "Provider unavailable" },
|
error: { type: "unknown", message: "Provider unavailable" },
|
||||||
content: [fixture.expectedContent],
|
content: [fixture.expectedContent],
|
||||||
},
|
},
|
||||||
@@ -542,12 +543,28 @@ const verifyPartialFlushOnInterruption = (kind: FragmentKind) =>
|
|||||||
const fiber = yield* runner.run({ sessionID, force: true }).pipe(Effect.forkChild)
|
const fiber = yield* runner.run({ sessionID, force: true }).pipe(Effect.forkChild)
|
||||||
yield* Deferred.await(streamed)
|
yield* Deferred.await(streamed)
|
||||||
yield* Fiber.interrupt(fiber)
|
yield* Fiber.interrupt(fiber)
|
||||||
|
const { db } = yield* Database.Service
|
||||||
|
const interruptedVersion = SessionEvent.Step.Interrupted.durable?.version
|
||||||
|
expect(interruptedVersion).toBe(2)
|
||||||
|
if (interruptedVersion === undefined) return yield* Effect.die("Step.Interrupted must be durable")
|
||||||
|
const settlements = yield* db
|
||||||
|
.select({ type: EventTable.type })
|
||||||
|
.from(EventTable)
|
||||||
|
.where(eq(EventTable.aggregate_id, sessionID))
|
||||||
|
.all()
|
||||||
|
.pipe(Effect.orDie)
|
||||||
|
expect(
|
||||||
|
settlements.filter(({ type }) =>
|
||||||
|
[SessionEvent.Step.Ended.type, SessionEvent.Step.Failed.type, SessionEvent.Step.Interrupted.type].some((settled) =>
|
||||||
|
type.startsWith(settled),
|
||||||
|
),
|
||||||
|
),
|
||||||
|
).toEqual([{ type: EventV2.versionedType(SessionEvent.Step.Interrupted.type, interruptedVersion) }])
|
||||||
expect(yield* session.context(sessionID)).toMatchObject([
|
expect(yield* session.context(sessionID)).toMatchObject([
|
||||||
{ type: "user", text: prompt },
|
{ type: "user", text: prompt },
|
||||||
{
|
{
|
||||||
type: "assistant",
|
type: "assistant",
|
||||||
finish: "error",
|
settlement: "interrupted",
|
||||||
error: { type: "unknown", message: "Provider turn interrupted" },
|
|
||||||
content: [
|
content: [
|
||||||
kind === "tool input"
|
kind === "tool input"
|
||||||
? { type: "tool", id: fragmentID(kind, "interrupted"), state: { status: "error" } }
|
? { type: "tool", id: fragmentID(kind, "interrupted"), state: { status: "error" } }
|
||||||
@@ -2653,6 +2670,7 @@ describe("SessionRunnerLLM", () => {
|
|||||||
state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
|
state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
|
||||||
},
|
},
|
||||||
],
|
],
|
||||||
|
settlement: "interrupted",
|
||||||
},
|
},
|
||||||
])
|
])
|
||||||
}),
|
}),
|
||||||
@@ -2797,6 +2815,7 @@ describe("SessionRunnerLLM", () => {
|
|||||||
state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
|
state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
|
||||||
},
|
},
|
||||||
],
|
],
|
||||||
|
settlement: "interrupted",
|
||||||
},
|
},
|
||||||
])
|
])
|
||||||
}),
|
}),
|
||||||
|
|||||||
@@ -9,7 +9,7 @@ describe("public event manifest", () => {
|
|||||||
expect(EventManifest.Definitions).toBe(SchemaEventManifest.Definitions)
|
expect(EventManifest.Definitions).toBe(SchemaEventManifest.Definitions)
|
||||||
expect(EventManifest.Latest).toBe(SchemaEventManifest.Latest)
|
expect(EventManifest.Latest).toBe(SchemaEventManifest.Latest)
|
||||||
expect(EventManifest.Durable).toBe(SchemaEventManifest.Durable)
|
expect(EventManifest.Durable).toBe(SchemaEventManifest.Durable)
|
||||||
expect(EventManifest.Latest.size).toBe(88)
|
expect(EventManifest.Latest.size).toBe(89)
|
||||||
expect(EventManifest.Latest.get("session.next.step.ended")).toBe(SessionEvent.Step.Ended)
|
expect(EventManifest.Latest.get("session.next.step.ended")).toBe(SessionEvent.Step.Ended)
|
||||||
expect(EventManifest.Latest.get("todo.updated")).toBe(Todo.Event.Updated)
|
expect(EventManifest.Latest.get("todo.updated")).toBe(Todo.Event.Updated)
|
||||||
expect(EventManifest.Latest.has("ide.installed")).toBe(false)
|
expect(EventManifest.Latest.has("ide.installed")).toBe(false)
|
||||||
|
|||||||
@@ -165,6 +165,7 @@ export namespace Step {
|
|||||||
schema: {
|
schema: {
|
||||||
...Base,
|
...Base,
|
||||||
assistantMessageID: SessionMessage.ID,
|
assistantMessageID: SessionMessage.ID,
|
||||||
|
// Step.Ended v2 was originally persisted with an open string schema.
|
||||||
finish: Schema.String,
|
finish: Schema.String,
|
||||||
cost: Schema.Finite,
|
cost: Schema.Finite,
|
||||||
tokens: Schema.Struct({
|
tokens: Schema.Struct({
|
||||||
@@ -192,6 +193,16 @@ export namespace Step {
|
|||||||
},
|
},
|
||||||
})
|
})
|
||||||
export type Failed = typeof Failed.Type
|
export type Failed = typeof Failed.Type
|
||||||
|
|
||||||
|
export const Interrupted = Event.define({
|
||||||
|
type: "session.next.step.interrupted",
|
||||||
|
...stepSettlementOptions,
|
||||||
|
schema: {
|
||||||
|
...Base,
|
||||||
|
assistantMessageID: SessionMessage.ID,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
export type Interrupted = typeof Interrupted.Type
|
||||||
}
|
}
|
||||||
|
|
||||||
export namespace Text {
|
export namespace Text {
|
||||||
@@ -458,6 +469,7 @@ export const DurableDefinitions = Event.inventory(
|
|||||||
Step.Started,
|
Step.Started,
|
||||||
Step.Ended,
|
Step.Ended,
|
||||||
Step.Failed,
|
Step.Failed,
|
||||||
|
Step.Interrupted,
|
||||||
Text.Started,
|
Text.Started,
|
||||||
Text.Ended,
|
Text.Ended,
|
||||||
Tool.Input.Started,
|
Tool.Input.Started,
|
||||||
@@ -489,6 +501,7 @@ export const Definitions = Event.inventory(
|
|||||||
Step.Started,
|
Step.Started,
|
||||||
Step.Ended,
|
Step.Ended,
|
||||||
Step.Failed,
|
Step.Failed,
|
||||||
|
Step.Interrupted,
|
||||||
Text.Started,
|
Text.Started,
|
||||||
Text.Delta,
|
Text.Delta,
|
||||||
Text.Ended,
|
Text.Ended,
|
||||||
|
|||||||
@@ -21,6 +21,19 @@ export const UnknownError = Schema.Struct({
|
|||||||
message: Schema.String,
|
message: Schema.String,
|
||||||
}).annotate({ identifier: "Session.Error.Unknown" })
|
}).annotate({ identifier: "Session.Error.Unknown" })
|
||||||
|
|
||||||
|
export const Finish = Schema.Literals([
|
||||||
|
"stop",
|
||||||
|
"length",
|
||||||
|
"tool-calls",
|
||||||
|
"content-filter",
|
||||||
|
"error",
|
||||||
|
"unknown",
|
||||||
|
])
|
||||||
|
export type Finish = typeof Finish.Type
|
||||||
|
|
||||||
|
export const Settlement = Schema.Literals(["completed", "failed", "interrupted"])
|
||||||
|
export type Settlement = typeof Settlement.Type
|
||||||
|
|
||||||
const Base = {
|
const Base = {
|
||||||
id: ID,
|
id: ID,
|
||||||
metadata: Schema.Record(Schema.String, Schema.Unknown).pipe(optional),
|
metadata: Schema.Record(Schema.String, Schema.Unknown).pipe(optional),
|
||||||
@@ -169,7 +182,9 @@ export const Assistant = Schema.Struct({
|
|||||||
end: Schema.String.pipe(optional),
|
end: Schema.String.pipe(optional),
|
||||||
files: Schema.Array(RelativePath).pipe(optional),
|
files: Schema.Array(RelativePath).pipe(optional),
|
||||||
}).pipe(optional),
|
}).pipe(optional),
|
||||||
|
// Projected histories predate the typed provider finish model and may contain arbitrary values.
|
||||||
finish: Schema.String.pipe(optional),
|
finish: Schema.String.pipe(optional),
|
||||||
|
settlement: Settlement.pipe(optional),
|
||||||
cost: Schema.Finite.pipe(optional),
|
cost: Schema.Finite.pipe(optional),
|
||||||
tokens: Schema.Struct({
|
tokens: Schema.Struct({
|
||||||
input: Schema.Finite,
|
input: Schema.Finite,
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import { describe, expect, test } from "bun:test"
|
import { describe, expect, test } from "bun:test"
|
||||||
|
import { Schema } from "effect"
|
||||||
import { FileSystem, Integration, Permission, Project, Reference, Session, Workspace } from "../src"
|
import { FileSystem, Integration, Permission, Project, Reference, Session, Workspace } from "../src"
|
||||||
import { EventManifest } from "../src/event-manifest"
|
import { EventManifest } from "../src/event-manifest"
|
||||||
import { IdeEvent } from "../src/ide-event"
|
import { IdeEvent } from "../src/ide-event"
|
||||||
@@ -9,8 +10,8 @@ import { WorkspaceEvent } from "../src/workspace-event"
|
|||||||
|
|
||||||
describe("public event manifest", () => {
|
describe("public event manifest", () => {
|
||||||
test("owns the complete public event surface", () => {
|
test("owns the complete public event surface", () => {
|
||||||
expect(EventManifest.ServerDefinitions.length).toBe(55)
|
expect(EventManifest.ServerDefinitions.length).toBe(59)
|
||||||
expect(EventManifest.Definitions.length).toBe(85)
|
expect(EventManifest.Definitions.length).toBe(89)
|
||||||
expect(SessionV1.Event.Definitions).toEqual([
|
expect(SessionV1.Event.Definitions).toEqual([
|
||||||
SessionV1.Event.Created,
|
SessionV1.Event.Created,
|
||||||
SessionV1.Event.Updated,
|
SessionV1.Event.Updated,
|
||||||
@@ -23,8 +24,8 @@ describe("public event manifest", () => {
|
|||||||
SessionV1.Event.Diff,
|
SessionV1.Event.Diff,
|
||||||
SessionV1.Event.Error,
|
SessionV1.Event.Error,
|
||||||
])
|
])
|
||||||
expect(EventManifest.Latest.size).toBe(85)
|
expect(EventManifest.Latest.size).toBe(89)
|
||||||
expect(EventManifest.Durable.size).toBe(32)
|
expect(EventManifest.Durable.size).toBe(36)
|
||||||
})
|
})
|
||||||
|
|
||||||
test("uses canonical definitions for current public events", () => {
|
test("uses canonical definitions for current public events", () => {
|
||||||
@@ -42,12 +43,28 @@ describe("public event manifest", () => {
|
|||||||
expect(Reference.Event.Definitions).toEqual([Reference.Event.Updated])
|
expect(Reference.Event.Definitions).toEqual([Reference.Event.Updated])
|
||||||
expect(EventManifest.Latest.has("ide.installed")).toBe(false)
|
expect(EventManifest.Latest.has("ide.installed")).toBe(false)
|
||||||
expect(IdeEvent.Definitions).toEqual([IdeEvent.Installed])
|
expect(IdeEvent.Definitions).toEqual([IdeEvent.Installed])
|
||||||
expect(EventManifest.Definitions.slice(40, 43)).toEqual([
|
const partDelta = EventManifest.Definitions.indexOf(SessionV1.Event.PartDelta)
|
||||||
|
expect(partDelta).toBeGreaterThanOrEqual(0)
|
||||||
|
expect(EventManifest.Definitions.slice(partDelta, partDelta + 3)).toEqual([
|
||||||
SessionV1.Event.PartDelta,
|
SessionV1.Event.PartDelta,
|
||||||
SessionV1.Event.Diff,
|
SessionV1.Event.Diff,
|
||||||
SessionV1.Event.Error,
|
SessionV1.Event.Error,
|
||||||
])
|
])
|
||||||
|
expect(EventManifest.Latest.get("session.next.step.interrupted")).toBe(SessionEvent.Step.Interrupted)
|
||||||
expect(EventManifest.Durable.has("session.next.step.ended.1")).toBe(false)
|
expect(EventManifest.Durable.has("session.next.step.ended.1")).toBe(false)
|
||||||
expect(EventManifest.Durable.get("session.next.step.ended.2")).toBe(SessionEvent.Step.Ended)
|
expect(EventManifest.Durable.get("session.next.step.ended.2")).toBe(SessionEvent.Step.Ended)
|
||||||
})
|
})
|
||||||
|
|
||||||
|
test("decodes legacy Step.Ended v2 finish strings", () => {
|
||||||
|
const event = Schema.decodeUnknownSync(SessionEvent.Step.Ended.data)({
|
||||||
|
sessionID: "ses_legacy",
|
||||||
|
timestamp: 0,
|
||||||
|
assistantMessageID: "msg_legacy",
|
||||||
|
finish: "legacy-provider-reason",
|
||||||
|
cost: 0,
|
||||||
|
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
|
||||||
|
})
|
||||||
|
|
||||||
|
expect(event.finish).toBe("legacy-provider-reason")
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -0,0 +1,22 @@
|
|||||||
|
import { expect, test } from "bun:test"
|
||||||
|
import { Schema } from "effect"
|
||||||
|
import { SessionMessage } from "../src/session-message"
|
||||||
|
|
||||||
|
test("does not model interruption as a provider finish reason", () => {
|
||||||
|
expect(() => Schema.decodeUnknownSync(SessionMessage.Finish)("interrupted")).toThrow()
|
||||||
|
expect(Schema.decodeUnknownSync(SessionMessage.Finish)("error")).toBe("error")
|
||||||
|
})
|
||||||
|
|
||||||
|
test("decodes projected assistant histories with arbitrary finish strings", () => {
|
||||||
|
const message = Schema.decodeUnknownSync(SessionMessage.Message)({
|
||||||
|
id: "msg_legacy",
|
||||||
|
type: "assistant",
|
||||||
|
agent: "build",
|
||||||
|
model: { id: "model", providerID: "provider" },
|
||||||
|
content: [],
|
||||||
|
finish: "legacy-provider-reason",
|
||||||
|
time: { created: 0, completed: 1 },
|
||||||
|
})
|
||||||
|
|
||||||
|
expect(message).toMatchObject({ type: "assistant", finish: "legacy-provider-reason" })
|
||||||
|
})
|
||||||
@@ -28,6 +28,7 @@ export type Event =
|
|||||||
| EventSessionNextStepStarted
|
| EventSessionNextStepStarted
|
||||||
| EventSessionNextStepEnded
|
| EventSessionNextStepEnded
|
||||||
| EventSessionNextStepFailed
|
| EventSessionNextStepFailed
|
||||||
|
| EventSessionNextStepInterrupted
|
||||||
| EventSessionNextTextStarted
|
| EventSessionNextTextStarted
|
||||||
| EventSessionNextTextDelta
|
| EventSessionNextTextDelta
|
||||||
| EventSessionNextTextEnded
|
| EventSessionNextTextEnded
|
||||||
@@ -955,6 +956,15 @@ export type GlobalEvent = {
|
|||||||
error: SessionErrorUnknown
|
error: SessionErrorUnknown
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
| {
|
||||||
|
id: string
|
||||||
|
type: "session.next.step.interrupted"
|
||||||
|
properties: {
|
||||||
|
timestamp: number
|
||||||
|
sessionID: string
|
||||||
|
assistantMessageID: string
|
||||||
|
}
|
||||||
|
}
|
||||||
| {
|
| {
|
||||||
id: string
|
id: string
|
||||||
type: "session.next.text.started"
|
type: "session.next.text.started"
|
||||||
@@ -1620,6 +1630,7 @@ export type GlobalEvent = {
|
|||||||
| SyncEventSessionNextStepStarted
|
| SyncEventSessionNextStepStarted
|
||||||
| SyncEventSessionNextStepEnded
|
| SyncEventSessionNextStepEnded
|
||||||
| SyncEventSessionNextStepFailed
|
| SyncEventSessionNextStepFailed
|
||||||
|
| SyncEventSessionNextStepInterrupted
|
||||||
| SyncEventSessionNextTextStarted
|
| SyncEventSessionNextTextStarted
|
||||||
| SyncEventSessionNextTextEnded
|
| SyncEventSessionNextTextEnded
|
||||||
| SyncEventSessionNextReasoningStarted
|
| SyncEventSessionNextReasoningStarted
|
||||||
@@ -2759,6 +2770,7 @@ export type V2Event =
|
|||||||
| V2EventSessionNextStepStarted
|
| V2EventSessionNextStepStarted
|
||||||
| V2EventSessionNextStepEnded
|
| V2EventSessionNextStepEnded
|
||||||
| V2EventSessionNextStepFailed
|
| V2EventSessionNextStepFailed
|
||||||
|
| V2EventSessionNextStepInterrupted
|
||||||
| V2EventSessionNextTextStarted
|
| V2EventSessionNextTextStarted
|
||||||
| V2EventSessionNextTextDelta
|
| V2EventSessionNextTextDelta
|
||||||
| V2EventSessionNextTextEnded
|
| V2EventSessionNextTextEnded
|
||||||
@@ -3400,6 +3412,22 @@ export type SyncEventSessionNextStepFailed = {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export type SyncEventSessionNextStepInterrupted = {
|
||||||
|
type: "sync"
|
||||||
|
id: string
|
||||||
|
syncEvent: {
|
||||||
|
type: "session.next.step.interrupted.2"
|
||||||
|
id: string
|
||||||
|
seq: number
|
||||||
|
aggregateID: string
|
||||||
|
data: {
|
||||||
|
timestamp: number
|
||||||
|
sessionID: string
|
||||||
|
assistantMessageID: string
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export type SyncEventSessionNextTextStarted = {
|
export type SyncEventSessionNextTextStarted = {
|
||||||
type: "sync"
|
type: "sync"
|
||||||
id: string
|
id: string
|
||||||
@@ -3995,6 +4023,7 @@ export type SessionMessageAssistant = {
|
|||||||
files?: Array<string>
|
files?: Array<string>
|
||||||
}
|
}
|
||||||
finish?: string
|
finish?: string
|
||||||
|
settlement?: "completed" | "failed" | "interrupted"
|
||||||
cost?: number
|
cost?: number
|
||||||
tokens?: {
|
tokens?: {
|
||||||
input: number
|
input: number
|
||||||
@@ -4289,6 +4318,25 @@ export type SessionNextStepFailed = {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export type SessionNextStepInterrupted = {
|
||||||
|
id: string
|
||||||
|
metadata?: {
|
||||||
|
[key: string]: unknown
|
||||||
|
}
|
||||||
|
type: "session.next.step.interrupted"
|
||||||
|
durable?: {
|
||||||
|
aggregateID: string
|
||||||
|
seq: number | "NaN" | "Infinity" | "-Infinity"
|
||||||
|
version: number | "NaN" | "Infinity" | "-Infinity"
|
||||||
|
}
|
||||||
|
location?: LocationRef
|
||||||
|
data: {
|
||||||
|
timestamp: number
|
||||||
|
sessionID: string
|
||||||
|
assistantMessageID: string
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export type SessionNextTextStarted = {
|
export type SessionNextTextStarted = {
|
||||||
id: string
|
id: string
|
||||||
metadata?: {
|
metadata?: {
|
||||||
@@ -5344,6 +5392,25 @@ export type V2EventSessionNextStepFailed = {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export type V2EventSessionNextStepInterrupted = {
|
||||||
|
id: string
|
||||||
|
metadata?: {
|
||||||
|
[key: string]: unknown
|
||||||
|
}
|
||||||
|
durable?: {
|
||||||
|
aggregateID: string
|
||||||
|
seq: number
|
||||||
|
version: number
|
||||||
|
}
|
||||||
|
location?: LocationRef
|
||||||
|
type: "session.next.step.interrupted"
|
||||||
|
data: {
|
||||||
|
timestamp: number
|
||||||
|
sessionID: string
|
||||||
|
assistantMessageID: string
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export type V2EventSessionNextTextStarted = {
|
export type V2EventSessionNextTextStarted = {
|
||||||
id: string
|
id: string
|
||||||
metadata?: {
|
metadata?: {
|
||||||
@@ -6931,6 +6998,16 @@ export type EventSessionNextStepFailed = {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export type EventSessionNextStepInterrupted = {
|
||||||
|
id: string
|
||||||
|
type: "session.next.step.interrupted"
|
||||||
|
properties: {
|
||||||
|
timestamp: number
|
||||||
|
sessionID: string
|
||||||
|
assistantMessageID: string
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export type EventSessionNextTextStarted = {
|
export type EventSessionNextTextStarted = {
|
||||||
id: string
|
id: string
|
||||||
type: "session.next.text.started"
|
type: "session.next.text.started"
|
||||||
|
|||||||
@@ -227,6 +227,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
|||||||
if (!currentAssistant) return
|
if (!currentAssistant) return
|
||||||
currentAssistant.time.completed = event.data.timestamp
|
currentAssistant.time.completed = event.data.timestamp
|
||||||
currentAssistant.finish = event.data.finish
|
currentAssistant.finish = event.data.finish
|
||||||
|
currentAssistant.settlement = "completed"
|
||||||
currentAssistant.cost = event.data.cost
|
currentAssistant.cost = event.data.cost
|
||||||
currentAssistant.tokens = event.data.tokens
|
currentAssistant.tokens = event.data.tokens
|
||||||
if (event.data.snapshot)
|
if (event.data.snapshot)
|
||||||
@@ -239,9 +240,18 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
|||||||
if (!currentAssistant) return
|
if (!currentAssistant) return
|
||||||
currentAssistant.time.completed = event.data.timestamp
|
currentAssistant.time.completed = event.data.timestamp
|
||||||
currentAssistant.finish = "error"
|
currentAssistant.finish = "error"
|
||||||
|
currentAssistant.settlement = "failed"
|
||||||
currentAssistant.error = event.data.error
|
currentAssistant.error = event.data.error
|
||||||
})
|
})
|
||||||
break
|
break
|
||||||
|
case "session.next.step.interrupted":
|
||||||
|
message.update(event.data.sessionID, (draft) => {
|
||||||
|
const currentAssistant = message.assistant(draft, event.data.assistantMessageID)
|
||||||
|
if (!currentAssistant) return
|
||||||
|
currentAssistant.time.completed = event.data.timestamp
|
||||||
|
currentAssistant.settlement = "interrupted"
|
||||||
|
})
|
||||||
|
break
|
||||||
case "session.next.text.started":
|
case "session.next.text.started":
|
||||||
message.update(event.data.sessionID, (draft) => {
|
message.update(event.data.sessionID, (draft) => {
|
||||||
message.assistant(draft, event.data.assistantMessageID)?.content.push({
|
message.assistant(draft, event.data.assistantMessageID)?.content.push({
|
||||||
|
|||||||
@@ -370,6 +370,72 @@ test("settles pending tools when a live failure arrives", async () => {
|
|||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
|
test("marks an interrupted assistant inactive without a finish reason", async () => {
|
||||||
|
const events = createEventSource()
|
||||||
|
const calls = createFetch(undefined, events)
|
||||||
|
let sync!: ReturnType<typeof useData>
|
||||||
|
let ready!: () => void
|
||||||
|
const mounted = new Promise<void>((resolve) => {
|
||||||
|
ready = resolve
|
||||||
|
})
|
||||||
|
|
||||||
|
function Probe() {
|
||||||
|
sync = useData()
|
||||||
|
onMount(ready)
|
||||||
|
return <box />
|
||||||
|
}
|
||||||
|
|
||||||
|
const app = await testRender(() => (
|
||||||
|
<TestTuiContexts>
|
||||||
|
<SDKProvider url="http://test" directory={directory} events={events.source} fetch={calls.fetch}>
|
||||||
|
<ProjectProvider>
|
||||||
|
<DataProvider>
|
||||||
|
<Probe />
|
||||||
|
</DataProvider>
|
||||||
|
</ProjectProvider>
|
||||||
|
</SDKProvider>
|
||||||
|
</TestTuiContexts>
|
||||||
|
))
|
||||||
|
|
||||||
|
try {
|
||||||
|
await mounted
|
||||||
|
emitEvent(events, {
|
||||||
|
id: "evt_step_started_interrupted",
|
||||||
|
type: "session.next.step.started",
|
||||||
|
properties: {
|
||||||
|
sessionID: "session-interrupted",
|
||||||
|
assistantMessageID: "msg_interrupted",
|
||||||
|
timestamp: 1,
|
||||||
|
agent: "build",
|
||||||
|
model: { id: "model-1", providerID: "provider-1" },
|
||||||
|
},
|
||||||
|
})
|
||||||
|
emitEvent(events, {
|
||||||
|
id: "evt_step_interrupted",
|
||||||
|
type: "session.next.step.interrupted",
|
||||||
|
properties: {
|
||||||
|
sessionID: "session-interrupted",
|
||||||
|
assistantMessageID: "msg_interrupted",
|
||||||
|
timestamp: 2,
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
|
await wait(() => {
|
||||||
|
const message = sync.session.message.list("session-interrupted")?.[0]
|
||||||
|
return message?.type === "assistant" && message.time.completed === 2
|
||||||
|
})
|
||||||
|
const assistant = sync.session.message.list("session-interrupted")?.[0]
|
||||||
|
expect(assistant).toMatchObject({
|
||||||
|
type: "assistant",
|
||||||
|
settlement: "interrupted",
|
||||||
|
time: { completed: 2 },
|
||||||
|
})
|
||||||
|
expect(assistant).not.toHaveProperty("finish")
|
||||||
|
} finally {
|
||||||
|
app.renderer.destroy()
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
test("renders admitted prompts only after they become model-visible", async () => {
|
test("renders admitted prompts only after they become model-visible", async () => {
|
||||||
const events = createEventSource()
|
const events = createEventSource()
|
||||||
const calls = createFetch(undefined, events)
|
const calls = createFetch(undefined, events)
|
||||||
|
|||||||
Reference in New Issue
Block a user