Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
57 commits
Select commit Hold shift + click to select a range
6f9db24
feat(core): add embedded session prompt admission
kitlangton Jun 2, 2026
ecf75b5
feat(core): add idempotent session prompt admission
kitlangton Jun 2, 2026
6b4343e
feat(core): add idempotent session creation
kitlangton Jun 2, 2026
f3ba685
feat(core): add local session runtime boundary
kitlangton Jun 2, 2026
2cf2f5c
feat(core): add session runtime resume boundary
kitlangton Jun 2, 2026
d0975e4
feat(core): add durable local session runner
kitlangton Jun 2, 2026
21d9416
feat(core): add permission-checked read tool
kitlangton Jun 2, 2026
1c8b42c
feat(core): add bounded session continuation
kitlangton Jun 2, 2026
96cdba4
fix(core): harden local tool settlement
kitlangton Jun 2, 2026
abd3627
feat(core): start sessions after prompt admission
kitlangton Jun 2, 2026
193358e
refactor(llm): use namespaced tool constructor
kitlangton Jun 2, 2026
78b01d3
refactor(core): route session execution by location
kitlangton Jun 2, 2026
fcb09d7
refactor(core): clarify session projection retries
kitlangton Jun 2, 2026
5ed36f5
feat(core): eagerly settle durable local tools
kitlangton Jun 3, 2026
ff2a351
refactor(core): inline session execution test layer
kitlangton Jun 3, 2026
961d35b
perf(core): index session projection lookups
kitlangton Jun 3, 2026
a822ecb
refactor(llm): remove in-memory tool orchestration
kitlangton Jun 3, 2026
befbbc3
perf(core): keep text deltas ephemeral
kitlangton Jun 3, 2026
2ed6e6d
perf(core): keep reasoning deltas ephemeral
kitlangton Jun 3, 2026
6ef2a30
perf(core): keep tool input deltas ephemeral
kitlangton Jun 3, 2026
9480434
feat(core): add replayable session output cursors
kitlangton Jun 3, 2026
746cd98
feat(core): derive ids from external keys
kitlangton Jun 3, 2026
49ea8f6
refactor(core): tighten replayable session contracts
kitlangton Jun 3, 2026
5e2c797
feat(core): add durable session steering
kitlangton Jun 3, 2026
786c17d
test(core): cover replay publication fencing
kitlangton Jun 3, 2026
54ea6b8
feat(core): add bounded list tool
kitlangton Jun 3, 2026
49c5ce2
fix(core): track session projection order migration
kitlangton Jun 3, 2026
d230b29
chore: regenerate v2 session artifacts
kitlangton Jun 3, 2026
7291176
fix(core): harden bounded filesystem tools
kitlangton Jun 3, 2026
7a02ca8
fix(opencode): share hardened path containment
kitlangton Jun 3, 2026
6583b4b
fix(opencode): run V2 HTTP prompts locally
kitlangton Jun 3, 2026
43add7e
feat(core): add durable session input inbox
kitlangton Jun 3, 2026
2d8e417
chore: regenerate session inbox artifacts
kitlangton Jun 3, 2026
b6f7744
feat(core): harden durable session inbox delivery
kitlangton Jun 3, 2026
23d6314
feat(core): add v2 tool foundation
kitlangton Jun 3, 2026
4e4ccda
fix(core): adapt v2 tools to upstream filesystem
kitlangton Jun 3, 2026
cec6e7b
test(core): update migration count after rebase
kitlangton Jun 3, 2026
3ce6db9
fix(core): reconcile rebased migration snapshots
kitlangton Jun 3, 2026
ae9a94e
fix(core): harden v2 tool authority boundaries
kitlangton Jun 3, 2026
339c2ad
fix(core): stabilize v2 embedded runtime
kitlangton Jun 3, 2026
b30e764
fix(core): bound v2 provider streams
kitlangton Jun 3, 2026
1913a33
perf(core): coalesce durable tail wakeups
kitlangton Jun 3, 2026
f37c008
fix(core): retain durable replay handoff wakes
kitlangton Jun 3, 2026
a9da6b2
docs(v2): defer launch hardening cleanup
kitlangton Jun 3, 2026
7a0ba4d
feat(core): add sequential v2 apply patch
kitlangton Jun 3, 2026
67df366
feat(core): add v2 skill tool and restart recovery
kitlangton Jun 3, 2026
4981241
test(llm): typecheck dynamic tool output regression
kitlangton Jun 3, 2026
32cff84
feat(core): add system context algebra
kitlangton Jun 3, 2026
586144a
fix(core): close v2 pre-pr safety gaps
kitlangton Jun 3, 2026
52cb0db
refactor(core): simplify v2 mutation commits
kitlangton Jun 4, 2026
2e0ba3e
fix(test): stabilize v2 cross-platform fixtures
kitlangton Jun 4, 2026
fd271e2
fix(core): stabilize embedded v2 runtime
kitlangton Jun 4, 2026
8d69375
fix(core): break config permission import cycle
kitlangton Jun 4, 2026
ab2e84f
fix(test): await core webfetch server shutdown
kitlangton Jun 4, 2026
d5d0502
fix(core): reuse embedded sqlite native handles
kitlangton Jun 4, 2026
ed78e06
fix(test): normalize core windows fixtures
kitlangton Jun 4, 2026
63cb0f7
fix(test): normalize windows mutation expectations
kitlangton Jun 4, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
perf(core): keep text deltas ephemeral
  • Loading branch information
kitlangton committed Jun 4, 2026
commit befbbc3627e991fae5023853146c08c788c161fd
9 changes: 6 additions & 3 deletions packages/core/src/event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -256,10 +256,13 @@ export const layer = Layer.effect(

function publishEvent<D extends Definition>(event: Payload<D>) {
return Effect.gen(function* () {
for (const sync of syncHandlers) {
yield* sync(event as Payload)
const durable = registry.get(event.type)?.sync !== undefined
if (durable) {
for (const sync of syncHandlers) {
yield* sync(event as Payload)
}
yield* commitSyncEvent(event as Payload)
}
yield* commitSyncEvent(event as Payload)
for (const listener of listeners) {
yield* listener(event as Payload)
}
Expand Down
2 changes: 1 addition & 1 deletion packages/core/src/session/event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -159,9 +159,9 @@ export namespace Text {
})
export type Started = typeof Started.Type

// Stream fragments are live-only; Text.Ended is the replayable full-value boundary.
export const Delta = EventV2.define({
type: "session.next.text.delta",
...options,
schema: {
...Base,
textID: Schema.String,
Expand Down
1 change: 0 additions & 1 deletion packages/core/src/session/projector.ts
Original file line number Diff line number Diff line change
Expand Up @@ -332,7 +332,6 @@ export const layer = Layer.effectDiscard(
yield* events.project(SessionEvent.Step.Ended, (event) => run(db, event))
yield* events.project(SessionEvent.Step.Failed, (event) => run(db, event))
yield* events.project(SessionEvent.Text.Started, (event) => run(db, event))
yield* events.project(SessionEvent.Text.Delta, (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.Delta, (event) => run(db, event))
Expand Down
9 changes: 5 additions & 4 deletions packages/core/src/session/runner/llm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,7 @@ export const layer = Layer.effect(
let settledLocalTool = false
const context = yield* getContext(session.id)
const request = LLM.request({ model, messages: toLLMMessages(context), tools: yield* tools.definitions() })
const publishLLMEvent = createLLMEventPublisher(events, {
const publisher = createLLMEventPublisher(events, {
sessionID: session.id,
agent: session.agent ?? "build",
model: {
Expand All @@ -110,10 +110,10 @@ export const layer = Layer.effect(
},
})
const publication = Semaphore.makeUnsafe(1)
const publish = (event: LLMEvent) => publication.withPermit(publishLLMEvent(event))
const publish = (event: LLMEvent) => publication.withPermit(publisher.publish(event))

yield* llm.stream(request).pipe(
Stream.runForEach((event) =>
yield* llm.stream(request).pipe(
Stream.runForEach((event) =>
Effect.gen(function* () {
yield* publish(event)
if (event.type !== "tool-call" || event.providerExecuted) return
Expand All @@ -127,6 +127,7 @@ export const layer = Layer.effect(
)
}),
),
Effect.ensuring(publication.withPermit(publisher.flushText())),
)
yield* FiberSet.awaitEmpty(settlements)
return settledLocalTool
Expand Down
47 changes: 31 additions & 16 deletions packages/core/src/session/runner/publish-llm-event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ const output = (result: ToolResultValue): ToolOutput => {

/** Persist one provider turn without executing tools or starting a continuation turn. */
export const createLLMEventPublisher = (events: EventV2.Interface, input: Input) => {
const text = new Map<string, string>()
const text = new Map<string, string[]>()
const reasoning = new Map<string, string>()
const tools = new Map<
string,
Expand All @@ -78,6 +78,22 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
? Effect.die("Tool event before assistant step start")
: Effect.succeed(assistantMessageID)

const endText = Effect.fnUntraced(function* (textID: string) {
const chunks = text.get(textID)
if (!chunks) return yield* Effect.die(`Text end before start: ${textID}`)
yield* events.publish(SessionEvent.Text.Ended, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
textID,
text: chunks.join(""),
})
text.delete(textID)
})

const flushText = Effect.fn("SessionRunner.flushText")(function* () {
for (const textID of text.keys()) yield* endText(textID)
})

const startToolInput = Effect.fnUntraced(function* (event: { readonly id: string; readonly name: string }) {
if (tools.has(event.id)) return yield* Effect.die(`Duplicate tool input start: ${event.id}`)
const assistantMessageID = yield* currentAssistantMessageID()
Expand Down Expand Up @@ -106,37 +122,32 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
})
})

return Effect.fn("SessionRunner.publishLLMEvent")(function* (event: LLMEvent) {
const publish = Effect.fn("SessionRunner.publishLLMEvent")(function* (event: LLMEvent) {
switch (event.type) {
case "step-start":
assistantMessageID = (yield* events.publish(SessionEvent.Step.Started, { ...input, timestamp: yield* timestamp })).id
return
case "text-start":
text.set(event.id, "")
if (text.has(event.id)) return yield* Effect.die(`Duplicate text start: ${event.id}`)
text.set(event.id, [])
yield* events.publish(SessionEvent.Text.Started, { sessionID: input.sessionID, timestamp: yield* timestamp, textID: event.id })
return
case "text-delta":
if (!text.has(event.id)) return yield* Effect.die(`Text delta before start: ${event.id}`)
text.set(event.id, `${text.get(event.id)}${event.text}`)
{
const chunks = text.get(event.id)
if (!chunks) return yield* Effect.die(`Text delta before start: ${event.id}`)
chunks.push(event.text)
}
yield* events.publish(SessionEvent.Text.Delta, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
textID: event.id,
delta: event.text,
})
return
case "text-end": {
const value = text.get(event.id)
if (value === undefined) return yield* Effect.die(`Text end before start: ${event.id}`)
text.delete(event.id)
yield* events.publish(SessionEvent.Text.Ended, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
textID: event.id,
text: value,
})
case "text-end":
yield* endText(event.id)
return
}
case "reasoning-start":
if (reasoning.has(event.id)) return yield* Effect.die(`Duplicate reasoning start: ${event.id}`)
reasoning.set(event.id, "")
Expand Down Expand Up @@ -266,6 +277,7 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
return
}
case "step-finish":
yield* flushText()
yield* events.publish(SessionEvent.Step.Ended, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
Expand All @@ -277,6 +289,7 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
case "finish":
return
case "provider-error":
yield* flushText()
yield* events.publish(SessionEvent.Step.Failed, {
sessionID: input.sessionID,
timestamp: yield* timestamp,
Expand All @@ -285,4 +298,6 @@ export const createLLMEventPublisher = (events: EventV2.Interface, input: Input)
return
}
})

return { publish, flushText }
}
18 changes: 18 additions & 0 deletions packages/core/test/event.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,24 @@ describe("EventV2", () => {
}),
)

it.effect("does not synchronize live-only events", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
const synchronized = new Array<string>()
const unsubscribe = yield* events.sync((event) =>
Effect.sync(() => {
synchronized.push(event.type)
}),
)
yield* Effect.addFinalizer(() => unsubscribe)

yield* events.publish(Message, { text: "live only" })
yield* events.publish(SyncMessage, { id: "one", text: "durable" })

expect(synchronized).toEqual([SyncMessage.type])
}),
)

it.effect("inserts sync event rows on publish", () =>
Effect.gen(function* () {
const events = yield* EventV2.Service
Expand Down
2 changes: 0 additions & 2 deletions packages/core/test/session-runner-recorded.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -91,8 +91,6 @@ describe("SessionRunnerLLM recorded", () => {
"session.next.prompted.1",
"session.next.step.started.1",
"session.next.text.started.1",
"session.next.text.delta.1",
"session.next.text.delta.1",
"session.next.text.ended.1",
"session.next.step.ended.1",
])
Expand Down
148 changes: 148 additions & 0 deletions packages/core/test/session-runner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { Project } from "@opencode-ai/core/project"
import { ProjectTable } from "@opencode-ai/core/project/sql"
import { AbsolutePath } from "@opencode-ai/core/schema"
import { SessionV2 } from "@opencode-ai/core/session"
import { SessionEvent } from "@opencode-ai/core/session/event"
import { Prompt } from "@opencode-ai/core/session/prompt"
import { SessionProjector } from "@opencode-ai/core/session/projector"
import { SessionExecution } from "@opencode-ai/core/session/execution"
Expand Down Expand Up @@ -749,6 +750,153 @@ describe("SessionRunnerLLM", () => {
}),
)

it.effect("broadcasts provider text deltas without storing projection rewrites", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Stream many chunks" }), resume: false })
const live: SessionEvent.Text.Delta[] = []
const unsubscribe = yield* EventV2.Service.pipe(
Effect.flatMap((events) => events.listen((event) => Effect.sync(() => {
if (event.type === SessionEvent.Text.Delta.type) live.push(event as SessionEvent.Text.Delta)
}))),
)
yield* Effect.addFinalizer(() => unsubscribe)

responses = undefined
streamGate = undefined
streamStarted = undefined
response = [
LLMEvent.stepStart({ index: 0 }),
LLMEvent.textStart({ id: "text-many" }),
...Array.from({ length: 32 }, (_, index) => LLMEvent.textDelta({ id: "text-many", text: `${index},` })),
LLMEvent.textEnd({ id: "text-many" }),
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
LLMEvent.finish({ reason: "stop" }),
]

yield* session.resume(sessionID)

const { db } = yield* Database.Service
const events = yield* EventV2.Service
const deltas = yield* db
.select({ type: EventTable.type })
.from(EventTable)
.where(eq(EventTable.type, EventV2.versionedType(SessionEvent.Text.Delta.type, 1)))
.all()
.pipe(Effect.orDie)
expect(live).toHaveLength(32)
expect(deltas).toHaveLength(0)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Stream many chunks" },
{
type: "assistant",
finish: "stop",
content: [{ type: "text", id: "text-many", text: Array.from({ length: 32 }, (_, index) => `${index},`).join("") }],
},
])
const recorded = yield* db
.select()
.from(EventTable)
.where(eq(EventTable.aggregate_id, sessionID))
.orderBy(asc(EventTable.seq))
.all()
.pipe(Effect.orDie)
yield* events.remove(sessionID)
yield* db.delete(SessionMessageTable).where(eq(SessionMessageTable.session_id, sessionID)).run().pipe(Effect.orDie)
yield* events.replayAll(
recorded.map((event) => ({
id: event.id,
aggregateID: event.aggregate_id,
seq: event.seq,
type: event.type,
data: event.data,
})),
)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Stream many chunks" },
{
type: "assistant",
finish: "stop",
content: [{ type: "text", id: "text-many", text: Array.from({ length: 32 }, (_, index) => `${index},`).join("") }],
},
])
}),
)

it.effect("durably closes partial text when the provider stream fails", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fail after text" }), resume: false })
const failure = new LLMError({
module: "test",
method: "stream",
reason: new TransportReason({ message: "Provider unavailable" }),
})

responses = undefined
streamGate = undefined
streamStarted = undefined
responseStream = Stream.concat(
Stream.fromIterable([
LLMEvent.stepStart({ index: 0 }),
LLMEvent.textStart({ id: "text-partial" }),
LLMEvent.textDelta({ id: "text-partial", text: "Partial" }),
]),
Stream.fail(failure),
)

expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Fail after text" },
{ type: "assistant", content: [{ type: "text", id: "text-partial", text: "Partial" }] },
])
}),
)

it.effect("durably closes partial text when the provider stream is interrupted", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Interrupt after text" }), resume: false })
const streamed = yield* Deferred.make<void>()

responses = undefined
streamGate = undefined
streamStarted = undefined
responseStream = Stream.concat(
Stream.fromIterable([
LLMEvent.stepStart({ index: 0 }),
LLMEvent.textStart({ id: "text-interrupted" }),
LLMEvent.textDelta({ id: "text-interrupted", text: "Partial" }),
]),
Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)),
)

const fiber = yield* session.resume(sessionID).pipe(Effect.forkChild)
yield* Deferred.await(streamed)
yield* Fiber.interrupt(fiber)
expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Interrupt after text" },
{ type: "assistant", content: [{ type: "text", id: "text-interrupted", text: "Partial" }] },
])
}),
)

it.effect("rejects duplicate streamed text starts", () =>
Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
responses = undefined
streamGate = undefined
streamStarted = undefined
response = [LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })]

expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe("Duplicate text start: text-1")
}),
)

it.effect("rejects malformed streamed tool input ordering", () =>
Effect.gen(function* () {
yield* setup
Expand Down