Compare commits
2 commits
dev
...
nxl/fix-au
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8ce4c507a4 | ||
|
|
4095839858 |
2 changed files with 158 additions and 12 deletions
|
|
@ -82,6 +82,28 @@ const STRUCTURED_OUTPUT_SYSTEM_PROMPT = `IMPORTANT: The user has requested struc
|
||||||
|
|
||||||
const log = Log.create({ service: "session.prompt" })
|
const log = Log.create({ service: "session.prompt" })
|
||||||
const elog = EffectLogger.create({ service: "session.prompt" })
|
const elog = EffectLogger.create({ service: "session.prompt" })
|
||||||
|
const COMPACTION_WITHOUT_TAIL = "none"
|
||||||
|
|
||||||
|
function latestAutoCompactionBoundary(messages: MessageV2.WithParts[]) {
|
||||||
|
const users = new Map(
|
||||||
|
messages.flatMap((message) => {
|
||||||
|
if (message.info.role !== "user") return []
|
||||||
|
const part = message.parts.find((item): item is MessageV2.CompactionPart => item.type === "compaction" && item.auto)
|
||||||
|
return part ? [[message.info.id, part] as const] : []
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
const latest = messages
|
||||||
|
.flatMap((message) => {
|
||||||
|
if (message.info.role !== "assistant") return []
|
||||||
|
if (!message.info.summary || !message.info.finish || message.info.error) return []
|
||||||
|
if (!users.has(message.info.parentID)) return []
|
||||||
|
return [message.info]
|
||||||
|
})
|
||||||
|
.sort((a, b) => a.id.localeCompare(b.id))
|
||||||
|
.at(-1)
|
||||||
|
if (!latest) return
|
||||||
|
return users.get(latest.parentID)?.tail_start_id ?? COMPACTION_WITHOUT_TAIL
|
||||||
|
}
|
||||||
|
|
||||||
type ReferencePromptMetadata = {
|
type ReferencePromptMetadata = {
|
||||||
name: string
|
name: string
|
||||||
|
|
@ -1646,8 +1668,66 @@ NOTE: At any point in time through this workflow you should feel free to ask the
|
||||||
const slog = elog.with({ sessionID })
|
const slog = elog.with({ sessionID })
|
||||||
let structured: unknown
|
let structured: unknown
|
||||||
let step = 0
|
let step = 0
|
||||||
|
let pendingAutoCompactionBoundary = latestAutoCompactionBoundary(yield* MessageV2.filterCompactedEffect(sessionID))
|
||||||
|
let awaitingAutoCompactionProgress = false
|
||||||
const session = yield* sessions.get(sessionID).pipe(Effect.orDie)
|
const session = yield* sessions.get(sessionID).pipe(Effect.orDie)
|
||||||
|
|
||||||
|
const requestAutoCompaction = Effect.fn("SessionPrompt.requestAutoCompaction")(function* (input: {
|
||||||
|
lastUser: MessageV2.User
|
||||||
|
model: Provider.Model
|
||||||
|
messages: MessageV2.WithParts[]
|
||||||
|
overflow?: boolean
|
||||||
|
assistant?: MessageV2.Assistant
|
||||||
|
}) {
|
||||||
|
const boundary = latestAutoCompactionBoundary(input.messages)
|
||||||
|
// A completed auto-compaction only counts as progress when it advances
|
||||||
|
// the retained tail boundary. Otherwise another auto-trigger would loop.
|
||||||
|
if (!awaitingAutoCompactionProgress || boundary !== pendingAutoCompactionBoundary) {
|
||||||
|
pendingAutoCompactionBoundary = boundary
|
||||||
|
awaitingAutoCompactionProgress = true
|
||||||
|
yield* compaction.create({
|
||||||
|
sessionID,
|
||||||
|
agent: input.lastUser.agent,
|
||||||
|
model: input.lastUser.model,
|
||||||
|
auto: true,
|
||||||
|
overflow: input.overflow,
|
||||||
|
})
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
const error = new MessageV2.ContextOverflowError({
|
||||||
|
message:
|
||||||
|
"Automatic compaction could not reduce this session enough to continue. Please start a new session.",
|
||||||
|
}).toObject()
|
||||||
|
yield* slog.info("auto compaction blocked", { boundary })
|
||||||
|
if (input.assistant) {
|
||||||
|
input.assistant.error = error
|
||||||
|
input.assistant.finish = "error"
|
||||||
|
input.assistant.time.completed ??= Date.now()
|
||||||
|
yield* sessions.updateMessage(input.assistant)
|
||||||
|
} else {
|
||||||
|
yield* sessions.updateMessage({
|
||||||
|
id: MessageID.ascending(),
|
||||||
|
parentID: input.lastUser.id,
|
||||||
|
role: "assistant",
|
||||||
|
mode: input.lastUser.agent,
|
||||||
|
agent: input.lastUser.agent,
|
||||||
|
variant: input.lastUser.model.variant,
|
||||||
|
path: { cwd: ctx.directory, root: ctx.worktree },
|
||||||
|
cost: 0,
|
||||||
|
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
|
||||||
|
modelID: input.model.id,
|
||||||
|
providerID: input.model.providerID,
|
||||||
|
time: { created: Date.now(), completed: Date.now() },
|
||||||
|
sessionID,
|
||||||
|
finish: "error",
|
||||||
|
error,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
yield* bus.publish(Session.Event.Error, { sessionID, error })
|
||||||
|
return false
|
||||||
|
})
|
||||||
|
|
||||||
while (true) {
|
while (true) {
|
||||||
yield* status.set(sessionID, { type: "busy" })
|
yield* status.set(sessionID, { type: "busy" })
|
||||||
yield* slog.info("loop", { step })
|
yield* slog.info("loop", { step })
|
||||||
|
|
@ -1712,8 +1792,8 @@ NOTE: At any point in time through this workflow you should feel free to ask the
|
||||||
lastFinished.summary !== true &&
|
lastFinished.summary !== true &&
|
||||||
(yield* compaction.isOverflow({ tokens: lastFinished.tokens, model }))
|
(yield* compaction.isOverflow({ tokens: lastFinished.tokens, model }))
|
||||||
) {
|
) {
|
||||||
yield* compaction.create({ sessionID, agent: lastUser.agent, model: lastUser.model, auto: true })
|
if (yield* requestAutoCompaction({ lastUser, model, messages: msgs })) continue
|
||||||
continue
|
break
|
||||||
}
|
}
|
||||||
|
|
||||||
const agent = yield* agents.get(lastUser.agent)
|
const agent = yield* agents.get(lastUser.agent)
|
||||||
|
|
@ -1852,13 +1932,14 @@ NOTE: At any point in time through this workflow you should feel free to ask the
|
||||||
|
|
||||||
if (result === "stop") return "break" as const
|
if (result === "stop") return "break" as const
|
||||||
if (result === "compact") {
|
if (result === "compact") {
|
||||||
yield* compaction.create({
|
const created = yield* requestAutoCompaction({
|
||||||
sessionID,
|
lastUser,
|
||||||
agent: lastUser.agent,
|
model,
|
||||||
model: lastUser.model,
|
messages: msgs,
|
||||||
auto: true,
|
assistant: handle.message,
|
||||||
overflow: !handle.message.finish,
|
overflow: !handle.message.finish,
|
||||||
})
|
})
|
||||||
|
if (!created) return "break" as const
|
||||||
}
|
}
|
||||||
return "continue" as const
|
return "continue" as const
|
||||||
}).pipe(
|
}).pipe(
|
||||||
|
|
|
||||||
|
|
@ -163,7 +163,37 @@ const blockingProcessor = Layer.succeed(
|
||||||
}),
|
}),
|
||||||
)
|
)
|
||||||
|
|
||||||
function makeHttp(input?: { processor?: "blocking" }) {
|
const compactLoopProcessor = Layer.effect(
|
||||||
|
SessionProcessor.Service,
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const sessions = yield* Session.Service
|
||||||
|
return SessionProcessor.Service.of({
|
||||||
|
create: (input) =>
|
||||||
|
Effect.succeed({
|
||||||
|
message: input.assistantMessage,
|
||||||
|
updateToolCall: () => Effect.succeed(undefined),
|
||||||
|
completeToolCall: () => Effect.void,
|
||||||
|
process: () =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
if (!input.assistantMessage.summary) return "compact" as const
|
||||||
|
input.assistantMessage.finish = "stop"
|
||||||
|
input.assistantMessage.time.completed = Date.now()
|
||||||
|
yield* sessions.updateMessage(input.assistantMessage)
|
||||||
|
yield* sessions.updatePart({
|
||||||
|
id: PartID.ascending(),
|
||||||
|
messageID: input.assistantMessage.id,
|
||||||
|
sessionID: input.assistantMessage.sessionID,
|
||||||
|
type: "text",
|
||||||
|
text: "summary",
|
||||||
|
})
|
||||||
|
return "continue" as const
|
||||||
|
}),
|
||||||
|
}),
|
||||||
|
})
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
|
function makeHttp(input?: { processor?: "blocking" | "compact-loop" }) {
|
||||||
const deps = Layer.mergeAll(
|
const deps = Layer.mergeAll(
|
||||||
Session.defaultLayer,
|
Session.defaultLayer,
|
||||||
Snapshot.defaultLayer,
|
Snapshot.defaultLayer,
|
||||||
|
|
@ -199,15 +229,21 @@ function makeHttp(input?: { processor?: "blocking" }) {
|
||||||
Layer.provideMerge(deps),
|
Layer.provideMerge(deps),
|
||||||
)
|
)
|
||||||
const trunc = Truncate.layer.pipe(Layer.provideMerge(deps))
|
const trunc = Truncate.layer.pipe(Layer.provideMerge(deps))
|
||||||
const proc =
|
const proc = (() => {
|
||||||
input?.processor === "blocking"
|
switch (input?.processor) {
|
||||||
? blockingProcessor
|
case "blocking":
|
||||||
: SessionProcessor.layer.pipe(
|
return blockingProcessor
|
||||||
|
case "compact-loop":
|
||||||
|
return compactLoopProcessor.pipe(Layer.provideMerge(deps))
|
||||||
|
default:
|
||||||
|
return SessionProcessor.layer.pipe(
|
||||||
Layer.provide(summary),
|
Layer.provide(summary),
|
||||||
Layer.provide(Image.defaultLayer),
|
Layer.provide(Image.defaultLayer),
|
||||||
Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })),
|
Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })),
|
||||||
Layer.provideMerge(deps),
|
Layer.provideMerge(deps),
|
||||||
)
|
)
|
||||||
|
}
|
||||||
|
})()
|
||||||
const compact = SessionCompaction.layer.pipe(
|
const compact = SessionCompaction.layer.pipe(
|
||||||
Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })),
|
Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })),
|
||||||
Layer.provideMerge(proc),
|
Layer.provideMerge(proc),
|
||||||
|
|
@ -235,6 +271,7 @@ function makeHttp(input?: { processor?: "blocking" }) {
|
||||||
|
|
||||||
const it = testEffect(makeHttp())
|
const it = testEffect(makeHttp())
|
||||||
const race = testEffect(makeHttp({ processor: "blocking" }))
|
const race = testEffect(makeHttp({ processor: "blocking" }))
|
||||||
|
const compactLoop = testEffect(makeHttp({ processor: "compact-loop" }))
|
||||||
const unix = process.platform !== "win32" ? it.instance : it.instance.skip
|
const unix = process.platform !== "win32" ? it.instance : it.instance.skip
|
||||||
|
|
||||||
// Config that registers a custom "test" provider with a "test-model" model
|
// Config that registers a custom "test" provider with a "test-model" model
|
||||||
|
|
@ -477,6 +514,34 @@ it.instance(
|
||||||
{ git: true },
|
{ git: true },
|
||||||
)
|
)
|
||||||
|
|
||||||
|
compactLoop.instance(
|
||||||
|
"loop stops when auto-compaction does not make progress",
|
||||||
|
() =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
yield* useServerConfig(providerCfg)
|
||||||
|
const prompt = yield* SessionPrompt.Service
|
||||||
|
const sessions = yield* Session.Service
|
||||||
|
const chat = yield* sessions.create({ title: "Pinned" })
|
||||||
|
yield* prompt.prompt({
|
||||||
|
sessionID: chat.id,
|
||||||
|
agent: "build",
|
||||||
|
noReply: true,
|
||||||
|
parts: [{ type: "text", text: "keep going until compaction fails" }],
|
||||||
|
})
|
||||||
|
|
||||||
|
const result = yield* prompt.loop({ sessionID: chat.id })
|
||||||
|
expect(result.info.role).toBe("assistant")
|
||||||
|
if (result.info.role === "assistant") {
|
||||||
|
expect(result.info.finish).toBe("error")
|
||||||
|
expect(result.info.error?.name).toBe("ContextOverflowError")
|
||||||
|
}
|
||||||
|
|
||||||
|
const messages = yield* sessions.messages({ sessionID: chat.id })
|
||||||
|
expect(messages.flatMap((message) => message.parts).filter((part) => part.type === "compaction")).toHaveLength(2)
|
||||||
|
}),
|
||||||
|
{ git: true },
|
||||||
|
)
|
||||||
|
|
||||||
it.instance(
|
it.instance(
|
||||||
"prompt emits v2 prompted and synthetic events",
|
"prompt emits v2 prompted and synthetic events",
|
||||||
() =>
|
() =>
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue