refactor(session): remove summary async facades (#22337)
This commit is contained in:
parent
14ccff4037
commit
dcbf11f41a
8 changed files with 68 additions and 26 deletions
|
|
@ -474,10 +474,14 @@ export const SessionRoutes = lazy(() =>
|
||||||
async (c) => {
|
async (c) => {
|
||||||
const query = c.req.valid("query")
|
const query = c.req.valid("query")
|
||||||
const params = c.req.valid("param")
|
const params = c.req.valid("param")
|
||||||
const result = await SessionSummary.diff({
|
const result = await AppRuntime.runPromise(
|
||||||
sessionID: params.sessionID,
|
SessionSummary.Service.use((summary) =>
|
||||||
messageID: query.messageID,
|
summary.diff({
|
||||||
})
|
sessionID: params.sessionID,
|
||||||
|
messageID: query.messageID,
|
||||||
|
}),
|
||||||
|
),
|
||||||
|
)
|
||||||
return c.json(result)
|
return c.json(result)
|
||||||
},
|
},
|
||||||
)
|
)
|
||||||
|
|
|
||||||
|
|
@ -1,4 +1,4 @@
|
||||||
import { Cause, Deferred, Effect, Layer, Context } from "effect"
|
import { Cause, Deferred, Effect, Layer, Context, Scope } from "effect"
|
||||||
import * as Stream from "effect/Stream"
|
import * as Stream from "effect/Stream"
|
||||||
import { Agent } from "@/agent/agent"
|
import { Agent } from "@/agent/agent"
|
||||||
import { Bus } from "@/bus"
|
import { Bus } from "@/bus"
|
||||||
|
|
@ -89,6 +89,7 @@ export namespace SessionProcessor {
|
||||||
| LLM.Service
|
| LLM.Service
|
||||||
| Permission.Service
|
| Permission.Service
|
||||||
| Plugin.Service
|
| Plugin.Service
|
||||||
|
| SessionSummary.Service
|
||||||
| SessionStatus.Service
|
| SessionStatus.Service
|
||||||
> = Layer.effect(
|
> = Layer.effect(
|
||||||
Service,
|
Service,
|
||||||
|
|
@ -101,6 +102,8 @@ export namespace SessionProcessor {
|
||||||
const llm = yield* LLM.Service
|
const llm = yield* LLM.Service
|
||||||
const permission = yield* Permission.Service
|
const permission = yield* Permission.Service
|
||||||
const plugin = yield* Plugin.Service
|
const plugin = yield* Plugin.Service
|
||||||
|
const summary = yield* SessionSummary.Service
|
||||||
|
const scope = yield* Scope.Scope
|
||||||
const status = yield* SessionStatus.Service
|
const status = yield* SessionStatus.Service
|
||||||
|
|
||||||
const create = Effect.fn("SessionProcessor.create")(function* (input: Input) {
|
const create = Effect.fn("SessionProcessor.create")(function* (input: Input) {
|
||||||
|
|
@ -385,10 +388,12 @@ export namespace SessionProcessor {
|
||||||
}
|
}
|
||||||
ctx.snapshot = undefined
|
ctx.snapshot = undefined
|
||||||
}
|
}
|
||||||
SessionSummary.summarize({
|
yield* summary
|
||||||
sessionID: ctx.sessionID,
|
.summarize({
|
||||||
messageID: ctx.assistantMessage.parentID,
|
sessionID: ctx.sessionID,
|
||||||
})
|
messageID: ctx.assistantMessage.parentID,
|
||||||
|
})
|
||||||
|
.pipe(Effect.ignore, Effect.forkIn(scope))
|
||||||
if (
|
if (
|
||||||
!ctx.assistantMessage.summary &&
|
!ctx.assistantMessage.summary &&
|
||||||
isOverflow({ cfg: yield* config.get(), tokens: usage.tokens, model: ctx.model })
|
isOverflow({ cfg: yield* config.get(), tokens: usage.tokens, model: ctx.model })
|
||||||
|
|
@ -603,6 +608,7 @@ export namespace SessionProcessor {
|
||||||
Layer.provide(LLM.defaultLayer),
|
Layer.provide(LLM.defaultLayer),
|
||||||
Layer.provide(Permission.defaultLayer),
|
Layer.provide(Permission.defaultLayer),
|
||||||
Layer.provide(Plugin.defaultLayer),
|
Layer.provide(Plugin.defaultLayer),
|
||||||
|
Layer.provide(SessionSummary.defaultLayer),
|
||||||
Layer.provide(SessionStatus.defaultLayer),
|
Layer.provide(SessionStatus.defaultLayer),
|
||||||
Layer.provide(Bus.layer),
|
Layer.provide(Bus.layer),
|
||||||
Layer.provide(Config.defaultLayer),
|
Layer.provide(Config.defaultLayer),
|
||||||
|
|
|
||||||
|
|
@ -102,6 +102,7 @@ export namespace SessionPrompt {
|
||||||
const instruction = yield* Instruction.Service
|
const instruction = yield* Instruction.Service
|
||||||
const state = yield* SessionRunState.Service
|
const state = yield* SessionRunState.Service
|
||||||
const revert = yield* SessionRevert.Service
|
const revert = yield* SessionRevert.Service
|
||||||
|
const summary = yield* SessionSummary.Service
|
||||||
const sys = yield* SystemPrompt.Service
|
const sys = yield* SystemPrompt.Service
|
||||||
const llm = yield* LLM.Service
|
const llm = yield* LLM.Service
|
||||||
|
|
||||||
|
|
@ -1444,7 +1445,10 @@ NOTE: At any point in time through this workflow you should feel free to ask the
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
if (step === 1) SessionSummary.summarize({ sessionID, messageID: lastUser.id })
|
if (step === 1)
|
||||||
|
yield* summary
|
||||||
|
.summarize({ sessionID, messageID: lastUser.id })
|
||||||
|
.pipe(Effect.ignore, Effect.forkIn(scope))
|
||||||
|
|
||||||
if (step > 1 && lastFinished) {
|
if (step > 1 && lastFinished) {
|
||||||
for (const m of msgs) {
|
for (const m of msgs) {
|
||||||
|
|
@ -1692,6 +1696,7 @@ NOTE: At any point in time through this workflow you should feel free to ask the
|
||||||
Layer.provide(Plugin.defaultLayer),
|
Layer.provide(Plugin.defaultLayer),
|
||||||
Layer.provide(Session.defaultLayer),
|
Layer.provide(Session.defaultLayer),
|
||||||
Layer.provide(SessionRevert.defaultLayer),
|
Layer.provide(SessionRevert.defaultLayer),
|
||||||
|
Layer.provide(SessionSummary.defaultLayer),
|
||||||
Layer.provide(
|
Layer.provide(
|
||||||
Layer.mergeAll(
|
Layer.mergeAll(
|
||||||
Agent.defaultLayer,
|
Agent.defaultLayer,
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,5 @@
|
||||||
import z from "zod"
|
import z from "zod"
|
||||||
import { Effect, Layer, Context } from "effect"
|
import { Effect, Layer, Context } from "effect"
|
||||||
import { makeRuntime } from "@/effect/run-service"
|
|
||||||
import { Bus } from "@/bus"
|
import { Bus } from "@/bus"
|
||||||
import { Snapshot } from "@/snapshot"
|
import { Snapshot } from "@/snapshot"
|
||||||
import { Storage } from "@/storage/storage"
|
import { Storage } from "@/storage/storage"
|
||||||
|
|
@ -159,17 +158,8 @@ export namespace SessionSummary {
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
const { runPromise } = makeRuntime(Service, defaultLayer)
|
|
||||||
|
|
||||||
export const summarize = (input: { sessionID: SessionID; messageID: MessageID }) =>
|
|
||||||
void runPromise((svc) => svc.summarize(input)).catch(() => {})
|
|
||||||
|
|
||||||
export const DiffInput = z.object({
|
export const DiffInput = z.object({
|
||||||
sessionID: SessionID.zod,
|
sessionID: SessionID.zod,
|
||||||
messageID: MessageID.zod.optional(),
|
messageID: MessageID.zod.optional(),
|
||||||
})
|
})
|
||||||
|
|
||||||
export async function diff(input: z.infer<typeof DiffInput>) {
|
|
||||||
return runPromise((svc) => svc.diff(input))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -18,6 +18,7 @@ import { Session } from "../../src/session"
|
||||||
import { MessageV2 } from "../../src/session/message-v2"
|
import { MessageV2 } from "../../src/session/message-v2"
|
||||||
import { MessageID, PartID, SessionID } from "../../src/session/schema"
|
import { MessageID, PartID, SessionID } from "../../src/session/schema"
|
||||||
import { SessionStatus } from "../../src/session/status"
|
import { SessionStatus } from "../../src/session/status"
|
||||||
|
import { SessionSummary } from "../../src/session/summary"
|
||||||
import { ModelID, ProviderID } from "../../src/provider/schema"
|
import { ModelID, ProviderID } from "../../src/provider/schema"
|
||||||
import type { Provider } from "../../src/provider/provider"
|
import type { Provider } from "../../src/provider/provider"
|
||||||
import * as SessionProcessorModule from "../../src/session/processor"
|
import * as SessionProcessorModule from "../../src/session/processor"
|
||||||
|
|
@ -26,6 +27,15 @@ import { ProviderTest } from "../fake/provider"
|
||||||
|
|
||||||
Log.init({ print: false })
|
Log.init({ print: false })
|
||||||
|
|
||||||
|
const summary = Layer.succeed(
|
||||||
|
SessionSummary.Service,
|
||||||
|
SessionSummary.Service.of({
|
||||||
|
summarize: () => Effect.void,
|
||||||
|
diff: () => Effect.succeed([]),
|
||||||
|
computeDiff: () => Effect.succeed([]),
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
const ref = {
|
const ref = {
|
||||||
providerID: ProviderID.make("test"),
|
providerID: ProviderID.make("test"),
|
||||||
modelID: ModelID.make("test-model"),
|
modelID: ModelID.make("test-model"),
|
||||||
|
|
@ -194,7 +204,7 @@ function llm() {
|
||||||
function liveRuntime(layer: Layer.Layer<LLM.Service>, provider = ProviderTest.fake()) {
|
function liveRuntime(layer: Layer.Layer<LLM.Service>, provider = ProviderTest.fake()) {
|
||||||
const bus = Bus.layer
|
const bus = Bus.layer
|
||||||
const status = SessionStatus.layer.pipe(Layer.provide(bus))
|
const status = SessionStatus.layer.pipe(Layer.provide(bus))
|
||||||
const processor = SessionProcessorModule.SessionProcessor.layer
|
const processor = SessionProcessorModule.SessionProcessor.layer.pipe(Layer.provide(summary))
|
||||||
return ManagedRuntime.make(
|
return ManagedRuntime.make(
|
||||||
Layer.mergeAll(SessionCompaction.layer.pipe(Layer.provide(processor)), processor, bus, status).pipe(
|
Layer.mergeAll(SessionCompaction.layer.pipe(Layer.provide(processor)), processor, bus, status).pipe(
|
||||||
Layer.provide(provider.layer),
|
Layer.provide(provider.layer),
|
||||||
|
|
|
||||||
|
|
@ -16,6 +16,7 @@ import { MessageV2 } from "../../src/session/message-v2"
|
||||||
import { SessionProcessor } from "../../src/session/processor"
|
import { SessionProcessor } from "../../src/session/processor"
|
||||||
import { MessageID, PartID, SessionID } from "../../src/session/schema"
|
import { MessageID, PartID, SessionID } from "../../src/session/schema"
|
||||||
import { SessionStatus } from "../../src/session/status"
|
import { SessionStatus } from "../../src/session/status"
|
||||||
|
import { SessionSummary } from "../../src/session/summary"
|
||||||
import { Snapshot } from "../../src/snapshot"
|
import { Snapshot } from "../../src/snapshot"
|
||||||
import { Log } from "../../src/util/log"
|
import { Log } from "../../src/util/log"
|
||||||
import * as CrossSpawnSpawner from "../../src/effect/cross-spawn-spawner"
|
import * as CrossSpawnSpawner from "../../src/effect/cross-spawn-spawner"
|
||||||
|
|
@ -25,6 +26,15 @@ import { raw, reply, TestLLMServer } from "../lib/llm-server"
|
||||||
|
|
||||||
Log.init({ print: false })
|
Log.init({ print: false })
|
||||||
|
|
||||||
|
const summary = Layer.succeed(
|
||||||
|
SessionSummary.Service,
|
||||||
|
SessionSummary.Service.of({
|
||||||
|
summarize: () => Effect.void,
|
||||||
|
diff: () => Effect.succeed([]),
|
||||||
|
computeDiff: () => Effect.succeed([]),
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
const ref = {
|
const ref = {
|
||||||
providerID: ProviderID.make("test"),
|
providerID: ProviderID.make("test"),
|
||||||
modelID: ModelID.make("test-model"),
|
modelID: ModelID.make("test-model"),
|
||||||
|
|
@ -156,7 +166,10 @@ const deps = Layer.mergeAll(
|
||||||
Provider.defaultLayer,
|
Provider.defaultLayer,
|
||||||
status,
|
status,
|
||||||
).pipe(Layer.provideMerge(infra))
|
).pipe(Layer.provideMerge(infra))
|
||||||
const env = Layer.mergeAll(TestLLMServer.layer, SessionProcessor.layer.pipe(Layer.provideMerge(deps)))
|
const env = Layer.mergeAll(
|
||||||
|
TestLLMServer.layer,
|
||||||
|
SessionProcessor.layer.pipe(Layer.provide(summary), Layer.provideMerge(deps)),
|
||||||
|
)
|
||||||
|
|
||||||
const it = testEffect(env)
|
const it = testEffect(env)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -23,6 +23,7 @@ import { LLM } from "../../src/session/llm"
|
||||||
import { MessageV2 } from "../../src/session/message-v2"
|
import { MessageV2 } from "../../src/session/message-v2"
|
||||||
import { AppFileSystem } from "../../src/filesystem"
|
import { AppFileSystem } from "../../src/filesystem"
|
||||||
import { SessionCompaction } from "../../src/session/compaction"
|
import { SessionCompaction } from "../../src/session/compaction"
|
||||||
|
import { SessionSummary } from "../../src/session/summary"
|
||||||
import { Instruction } from "../../src/session/instruction"
|
import { Instruction } from "../../src/session/instruction"
|
||||||
import { SessionProcessor } from "../../src/session/processor"
|
import { SessionProcessor } from "../../src/session/processor"
|
||||||
import { SessionPrompt } from "../../src/session/prompt"
|
import { SessionPrompt } from "../../src/session/prompt"
|
||||||
|
|
@ -46,6 +47,15 @@ import { reply, TestLLMServer } from "../lib/llm-server"
|
||||||
|
|
||||||
Log.init({ print: false })
|
Log.init({ print: false })
|
||||||
|
|
||||||
|
const summary = Layer.succeed(
|
||||||
|
SessionSummary.Service,
|
||||||
|
SessionSummary.Service.of({
|
||||||
|
summarize: () => Effect.void,
|
||||||
|
diff: () => Effect.succeed([]),
|
||||||
|
computeDiff: () => Effect.succeed([]),
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
const ref = {
|
const ref = {
|
||||||
providerID: ProviderID.make("test"),
|
providerID: ProviderID.make("test"),
|
||||||
modelID: ModelID.make("test-model"),
|
modelID: ModelID.make("test-model"),
|
||||||
|
|
@ -182,12 +192,13 @@ function makeHttp() {
|
||||||
Layer.provideMerge(deps),
|
Layer.provideMerge(deps),
|
||||||
)
|
)
|
||||||
const trunc = Truncate.layer.pipe(Layer.provideMerge(deps))
|
const trunc = Truncate.layer.pipe(Layer.provideMerge(deps))
|
||||||
const proc = SessionProcessor.layer.pipe(Layer.provideMerge(deps))
|
const proc = SessionProcessor.layer.pipe(Layer.provide(summary), Layer.provideMerge(deps))
|
||||||
const compact = SessionCompaction.layer.pipe(Layer.provideMerge(proc), Layer.provideMerge(deps))
|
const compact = SessionCompaction.layer.pipe(Layer.provideMerge(proc), Layer.provideMerge(deps))
|
||||||
return Layer.mergeAll(
|
return Layer.mergeAll(
|
||||||
TestLLMServer.layer,
|
TestLLMServer.layer,
|
||||||
SessionPrompt.layer.pipe(
|
SessionPrompt.layer.pipe(
|
||||||
Layer.provide(SessionRevert.defaultLayer),
|
Layer.provide(SessionRevert.defaultLayer),
|
||||||
|
Layer.provide(summary),
|
||||||
Layer.provideMerge(run),
|
Layer.provideMerge(run),
|
||||||
Layer.provideMerge(compact),
|
Layer.provideMerge(compact),
|
||||||
Layer.provideMerge(proc),
|
Layer.provideMerge(proc),
|
||||||
|
|
|
||||||
|
|
@ -146,12 +146,14 @@ function makeHttp() {
|
||||||
Layer.provideMerge(deps),
|
Layer.provideMerge(deps),
|
||||||
)
|
)
|
||||||
const trunc = Truncate.layer.pipe(Layer.provideMerge(deps))
|
const trunc = Truncate.layer.pipe(Layer.provideMerge(deps))
|
||||||
const proc = SessionProcessor.layer.pipe(Layer.provideMerge(deps))
|
const proc = SessionProcessor.layer.pipe(Layer.provide(SessionSummary.defaultLayer), Layer.provideMerge(deps))
|
||||||
const compact = SessionCompaction.layer.pipe(Layer.provideMerge(proc), Layer.provideMerge(deps))
|
const compact = SessionCompaction.layer.pipe(Layer.provideMerge(proc), Layer.provideMerge(deps))
|
||||||
return Layer.mergeAll(
|
return Layer.mergeAll(
|
||||||
TestLLMServer.layer,
|
TestLLMServer.layer,
|
||||||
|
SessionSummary.defaultLayer,
|
||||||
SessionPrompt.layer.pipe(
|
SessionPrompt.layer.pipe(
|
||||||
Layer.provide(SessionRevert.defaultLayer),
|
Layer.provide(SessionRevert.defaultLayer),
|
||||||
|
Layer.provide(SessionSummary.defaultLayer),
|
||||||
Layer.provideMerge(run),
|
Layer.provideMerge(run),
|
||||||
Layer.provideMerge(compact),
|
Layer.provideMerge(compact),
|
||||||
Layer.provideMerge(proc),
|
Layer.provideMerge(proc),
|
||||||
|
|
@ -200,6 +202,7 @@ it.live("tool execution produces non-empty session diff (snapshot race)", () =>
|
||||||
Effect.fnUntraced(function* ({ dir, llm }) {
|
Effect.fnUntraced(function* ({ dir, llm }) {
|
||||||
const prompt = yield* SessionPrompt.Service
|
const prompt = yield* SessionPrompt.Service
|
||||||
const sessions = yield* Session.Service
|
const sessions = yield* Session.Service
|
||||||
|
const summary = yield* SessionSummary.Service
|
||||||
|
|
||||||
const session = yield* sessions.create({
|
const session = yield* sessions.create({
|
||||||
title: "snapshot race test",
|
title: "snapshot race test",
|
||||||
|
|
@ -244,9 +247,9 @@ it.live("tool execution produces non-empty session diff (snapshot race)", () =>
|
||||||
expect(tool?.state.status).toBe("completed")
|
expect(tool?.state.status).toBe("completed")
|
||||||
|
|
||||||
// Poll for diff — summarize() is fire-and-forget
|
// Poll for diff — summarize() is fire-and-forget
|
||||||
let diff: Awaited<ReturnType<typeof SessionSummary.diff>> = []
|
let diff: Array<{ file: string }> = []
|
||||||
for (let i = 0; i < 50; i++) {
|
for (let i = 0; i < 50; i++) {
|
||||||
diff = yield* Effect.promise(() => SessionSummary.diff({ sessionID: session.id }))
|
diff = yield* summary.diff({ sessionID: session.id })
|
||||||
if (diff.length > 0) break
|
if (diff.length > 0) break
|
||||||
yield* Effect.sleep("100 millis")
|
yield* Effect.sleep("100 millis")
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue