From bcc68842f9468fd285472305b39f8b2b31435cc3 Mon Sep 17 00:00:00 2001 From: Lex Date: Mon, 28 Sep 2026 08:48:50 +0800 Subject: [PATCH 1/2] fix: persist opt-in reasoning distillation after each turn --- .github/releases/v1.0.51.md | 11 + docs/reasoning-distillation-design.md | 53 ++-- packages/client/src/generated/types.ts | 1 + packages/core/src/config.ts | 2 +- .../core/src/config/reasoning-distillation.ts | 6 +- packages/core/src/event.ts | 20 +- packages/core/src/session/message-updater.ts | 1 + .../session/reasoning-distillation/adopt.ts | 95 +++++++ .../reasoning-distillation/adoption.ts | 31 ++ .../reasoning-distillation/schedule.ts | 38 +++ packages/core/src/session/runner/llm.ts | 110 +++++--- .../session/runner/reasoning-distillation.ts | 29 +- .../core/src/session/runner/to-llm-message.ts | 3 + .../config/reasoning-distillation.test.ts | 4 +- packages/core/test/session-runner.test.ts | 151 +++++++++- .../test/session/reasoning-adoption.test.ts | 145 ++++++++++ .../test/session/reasoning-schedule.test.ts | 141 ++++++++++ packages/opencode/src/session/llm.ts | 239 +++++++++------- packages/opencode/src/session/processor.ts | 3 + packages/opencode/src/session/prompt.ts | 70 ++++- .../src/session/reasoning-adoption.ts | 142 ++++++++++ .../src/session/reasoning-distillation.ts | 13 +- packages/opencode/src/session/session.ts | 1 + packages/opencode/test/config/config.test.ts | 10 + .../opencode/test/session/compaction.test.ts | 1 + packages/opencode/test/session/llm.test.ts | 265 ++++++++++-------- .../test/session/processor-effect.test.ts | 7 +- packages/opencode/test/session/prompt.test.ts | 157 ++++++++++- .../session/reasoning-distillation.test.ts | 8 +- .../opencode/test/session/session.test.ts | 179 +++++++++++- packages/schema/src/reasoning-distillation.ts | 8 + packages/schema/src/session-event.ts | 2 + packages/schema/src/session-message.ts | 2 + packages/schema/src/session-v1.ts | 3 + packages/sdk/js/src/v2/gen/types.gen.ts | 15 + .../cmd/tui/sync-reasoning-adoption.test.tsx | 77 +++++ 36 files changed, 1737 insertions(+), 306 deletions(-) create mode 100644 .github/releases/v1.0.51.md create mode 100644 packages/core/src/session/reasoning-distillation/adopt.ts create mode 100644 packages/core/src/session/reasoning-distillation/adoption.ts create mode 100644 packages/core/src/session/reasoning-distillation/schedule.ts create mode 100644 packages/core/test/session/reasoning-adoption.test.ts create mode 100644 packages/core/test/session/reasoning-schedule.test.ts create mode 100644 packages/opencode/src/session/reasoning-adoption.ts create mode 100644 packages/schema/src/reasoning-distillation.ts create mode 100644 packages/tui/test/cli/cmd/tui/sync-reasoning-adoption.test.tsx diff --git a/.github/releases/v1.0.51.md b/.github/releases/v1.0.51.md new file mode 100644 index 0000000000..54aa365d57 --- /dev/null +++ b/.github/releases/v1.0.51.md @@ -0,0 +1,11 @@ +# GraphAgent v1.0.51 + +Reasoning distillation is now opt-in and updates conversation history as well as future model context. + +- Configure `reasoningDistillation.enabled` in `opencode.json`. The default is `false`. +- When enabled, the first completed user turn waits for distillation; later completed turns run distillation in the background. Tool steps remain part of their user turn. +- Validated results replace the original reasoning part in place. The TUI updates through its normal message events, reloaded sessions keep the adopted text, and subsequent requests reuse it. +- Failed, cancelled or stale work retains the current reasoning. Disabling the feature prevents new work and pending adoption. Existing compatibility and fidelity checks still protect unsupported, signed and encrypted reasoning. +- Host AI SDK, native adapter and Core runner share the same persisted replacement semantics; Core history caches are invalidated after adoption. + +Existing conversations are not automatically reprocessed when the feature is enabled. Background work requires the application to remain running. diff --git a/docs/reasoning-distillation-design.md b/docs/reasoning-distillation-design.md index 713c9f76d0..5286d7908d 100644 --- a/docs/reasoning-distillation-design.md +++ b/docs/reasoning-distillation-design.md @@ -1,5 +1,26 @@ # 内化推理蒸馏:开发设计与验收规格 +> **2026-09-28 product contract:** `opencode.json` controls `reasoningDistillation.enabled`, default `false`. +> Once enabled, the first completed user turn is distilled synchronously before the turn becomes idle. +> Each subsequent completed user turn is queued in the background. Tool steps do not count as turns. +> Successful adoption replaces the persisted reasoning part and emits its normal update event; the TUI, +> reopened history and later model requests use that same text. In-flight requests retain their snapshot. +> Failed, incompatible, stale or cancelled work retains the original. Disabling blocks new work and adoption. +> This contract supersedes the historical request-time, default-on and projection-only behavior below. + +```json +{ + "reasoningDistillation": { + "enabled": false + } +} +``` + +Set `enabled` to `true` to opt in. Existing exact-provider compatibility evidence and fidelity checks still apply; +turn scheduling alone does not authorize rewriting signed, encrypted or unsupported reasoning. +The original source is retained as host-owned provenance, excluded from provider request conversion. +Background work belongs to the running application scope; shutdown cancels unfinished work without changing history. + - 日期:2026-09-22,第 4 版。 - 状态:提案。产品实现未开始;文档检查通过不代表产品、模型质量或上游兼容性通过验收。 - 调研基线:`cedcfb3647e50d9a25633b65e4f3e599adc7fb6c`;`origin/dev` = `f3f4e50a0164c3b11b3b0126b2cd9c1ebe25dc37`。实施前重新核对开发分支。 @@ -74,7 +95,7 @@ W1/W2/W3 描述可定位的候选形态,**并不自动授予改写资格**。 上一轮读取的本地配置中,small_model 为 `local-proxy-compatible/deepseek`;多个 compatible 模型配置了 `interleaved.field = reasoning_content`,另有 Anthropic 通道和外部 think MCP。这是当时的配置观察,不是本版重新实测上游的结果,也不是可发布的兼容白名单。 -该 compatible 配置说明 W1 值得优先验证,但尚不能得出“完全可改写”或“全部历史思绪均被上游计费”的结论。成本需要用实际发送载荷和 usage 证明。默认开启功能时,这些槽位仍受 P5 门控。 +该 compatible 配置说明 W1 值得优先验证,但尚不能得出“完全可改写”或“全部历史思绪均被上游计费”的结论。成本需要用实际发送载荷和 usage 证明。显式开启功能时,这些槽位仍受 P5 门控。 ## 3. think 内化的含义与边界 @@ -92,20 +113,20 @@ W1/W2/W3 描述可定位的候选形态,**并不自动授予改写资格**。 ## 4. 已确认目标与形式化边界 -| 编号 | 已确认目标 | 本版解释 | -| ---- | ---------------------------- | ------------------------------------------------------------------ | -| D01 | 内化,无外部 MCP/插件依赖 | 共享核心与宿主适配承担 | -| D02 | 默认开启 | 开关默认开启,兼容授权仍默认保护 | -| D03 | 仅在压缩/动态压缩触发时执行 | 无后台任务;每请求只做有界开关、用途和预算判定,不做无条件模型整理 | -| D04 | 替换而非附加回传思绪 | conversation 中原槽位替换;不旁路添加额外消息 | -| D05 | 协议保护优先 | 无兼容证据不改写;收益不能覆盖保护失败 | -| D06 | 不写回原始历史 | 缓存、验证、审计和投影均为派生数据 | -| D07 | 小模型 → agent 模型 → 主模型 | 分级解析与实际尝试区分,所有实际调用共用预算 | -| D08 | 信息守恒审计 | 按具体声明核对证据;违反规则不等于证明模型具有欺骗意图 | -| D09 | 结构化、可回溯 | 每条 claim 有原文锚点,支持证据与来源锚点分开 | -| D10 | 回传信息 think 化 | 正反意见与理由保留,不以工具式 think 替代 | -| D11 | 安静、可关闭 | 诊断不含正文;关闭后不应用已有候选缓存 | -| D12 | 与折叠、全文压缩独立验收 | 各自计收益和回归,验证组合次序 | +| 编号 | 已确认目标 | 本版解释 | +| ---- | ---------------------------- | ------------------------------------------------------------ | +| D01 | 内化,无外部 MCP/插件依赖 | 共享核心与宿主适配承担 | +| D02 | 默认关闭、显式开启 | opencode.json 中 reasoningDistillation.enabled,缺省为 false | +| D03 | 完整回合结束时执行 | 首轮同步,此后每轮后台异步;工具步骤不单独计轮 | +| D04 | 替换而非附加回传思绪 | conversation 中原槽位替换;不旁路添加额外消息 | +| D05 | 协议保护优先 | 无兼容证据不改写;收益不能覆盖保护失败 | +| D06 | 保存已采用的思考 | 事务更新原思考槽位,保留来源信息,同步显示与后续回传 | +| D07 | 小模型 → agent 模型 → 主模型 | 分级解析与实际尝试区分,所有实际调用共用预算 | +| D08 | 信息守恒审计 | 按具体声明核对证据;违反规则不等于证明模型具有欺骗意图 | +| D09 | 结构化、可回溯 | 每条 claim 有原文锚点,支持证据与来源锚点分开 | +| D10 | 回传信息 think 化 | 正反意见与理由保留,不以工具式 think 替代 | +| D11 | 安静、可关闭 | 诊断不含正文;关闭后不应用已有候选缓存 | +| D12 | 与折叠、全文压缩独立验收 | 各自计收益和回归,验证组合次序 | ### 4.1 决策等价下的最小表示 @@ -531,7 +552,7 @@ claim store 由宿主按 InstanceState 管理,以 location/session 隔离, ## 7. 待决事项与剩余风险 1. **触发判据 — 已解决(2026-09-22,预算制)**:owner 确认采用 `ContextFoldingBudget.overBudget === true`(见 §5.1),不以 duplicatePlan 非空触发;实施前在 dev 分支复核预算估算路径与软阈值取值。 -2. **实现治理 — 已豁免(2026-09-22)**:项目所有者明确授权开发本功能,覆盖 `AGENTS.md` v1 focused-maintenance 的“禁止新增平台特性”约束。豁免仅限推理蒸馏本身及其直接依赖的模块/测试/配置,不扩展到其它无关平台特性或基础重构;是否同步修订 `AGENTS.md` 治理措辞属独立决定。默认开启仍为已确认产品目标。 +2. **实现治理 — 已豁免(2026-09-22)**:项目所有者明确授权开发本功能,覆盖 `AGENTS.md` v1 focused-maintenance 的“禁止新增平台特性”约束。豁免仅限推理蒸馏本身及其直接依赖的模块/测试/配置,不扩展到其它无关平台特性或基础重构;是否同步修订 `AGENTS.md` 治理措辞属独立决定。默认开启目标已由 2026-09-28 的默认关闭产品规则取代。 3. **兼容授权,启用阻塞**:当前配置只证明 W1 形态,不能生成上游白名单。每个实际启用组合必须完成 §2.1;签名/加密通道收益不对等,不能靠破坏保护补收益。 4. **语义能力上限**:一般命题等价和支持依赖 judged 证据,仍可能误判。需要跨度级证据、留出基准与明确归因;未知不放行,不能宣称模型审阅等同确定性证明。 5. **证据完整性上限**:历史和调用清单可能缺失,导致只能 unverifiable。结果正文清理不抹除调用状态;若需要新增持久化执行账本,那是另一个需批准的范围,不在首版偷偷加入。 diff --git a/packages/client/src/generated/types.ts b/packages/client/src/generated/types.ts index ac888719f4..b4f2abf82b 100644 --- a/packages/client/src/generated/types.ts +++ b/packages/client/src/generated/types.ts @@ -529,6 +529,7 @@ export type SessionsContextOutput = { readonly type: "reasoning" readonly id: string readonly text: string + readonly distillation?: { readonly originalText: string; readonly sourceFingerprint: string } | null readonly providerMetadata?: { readonly [x: string]: { readonly [x: string]: JsonValue } } | null } | { diff --git a/packages/core/src/config.ts b/packages/core/src/config.ts index 35f78240d1..1767f912d6 100644 --- a/packages/core/src/config.ts +++ b/packages/core/src/config.ts @@ -93,7 +93,7 @@ export class Info extends Schema.Class("Config.Info")({ }), reasoningDistillation: ConfigReasoningDistillation.Info.pipe(Schema.optional).annotate({ description: - "Reasoning distillation of returned model thoughts (default-on; any rewrite is gated by per-provider compatibility evidence)", + "Reasoning distillation of returned model thoughts (default-off; enable in opencode.json; any rewrite is gated by per-provider compatibility evidence)", }), skills: Schema.String.pipe(Schema.Array, Schema.optional).annotate({ description: "Additional paths or URLs to discover skills from", diff --git a/packages/core/src/config/reasoning-distillation.ts b/packages/core/src/config/reasoning-distillation.ts index f8ec887556..ca09353daf 100644 --- a/packages/core/src/config/reasoning-distillation.ts +++ b/packages/core/src/config/reasoning-distillation.ts @@ -14,7 +14,7 @@ export class Compatibility extends Schema.Class("ConfigV2.Reasoni }) {} /** - * Reasoning-distillation switch (design D02: default-on, but compatibility authorization still defaults to protected). + * Reasoning-distillation switch (default-off; explicit opt-in still requires compatibility authorization). * The feature only ever rewrites a reasoning slot that carries a §2.1 dual-evidence compatibility record; absent that * record every slot is P5-protected, so enabling the switch is safe and never rewrites an unproven provider. */ @@ -37,9 +37,9 @@ export type EnableInput = Readonly<{ enabled?: boolean }> -/** Resolve the effective switch without mutating persisted config. Default-on per D02 unless disabled. */ +/** Resolve the effective switch without mutating persisted config. Disabled unless explicitly enabled. */ export function resolveEnabled(input: EnableInput): EnableResolution { if (input.disabledByEnvironment) return { enabled: false, source: "environment" } if (input.enabled !== undefined) return { enabled: input.enabled, source: "config" } - return { enabled: true, source: "default" } + return { enabled: false, source: "default" } } diff --git a/packages/core/src/event.ts b/packages/core/src/event.ts index d2e23c627d..8df487eadd 100644 --- a/packages/core/src/event.ts +++ b/packages/core/src/event.ts @@ -75,7 +75,7 @@ export interface Interface { /** Batch durable publish: one transaction for the whole batch, contiguous seq per aggregate, projectors run in entry order, single durable wake after commit. */ readonly publishMany: ( events: ReadonlyArray, - options?: { readonly location?: Location.Ref }, + options?: { readonly location?: Location.Ref; readonly validate?: Effect.Effect }, ) => Effect.Effect> readonly subscribe: (definition: D) => Stream.Stream> readonly all: () => Stream.Stream @@ -139,11 +139,9 @@ export const layerWith = (options?: LayerOptions) => ) const wakeDurable = (aggregateID: string) => - Effect.forEach( - pubsub.durable.get(aggregateID) ?? [], - (wake) => PubSub.publish(wake, undefined), - { discard: true }, - ) + Effect.forEach(pubsub.durable.get(aggregateID) ?? [], (wake) => PubSub.publish(wake, undefined), { + discard: true, + }) /** Transaction-scoped single durable event commit: seq allocation, owner checks, projectors, UPSERT + INSERT. */ function commitDurableEventInner( @@ -416,9 +414,7 @@ export const layerWith = (options?: LayerOptions) => function notify(event: Payload) { return Effect.gen(function* () { const snapshot = Array.from(listeners) - forkListeners( - Effect.forEach(snapshot, (listener) => observe(event, listener), { discard: true }), - ) + forkListeners(Effect.forEach(snapshot, (listener) => observe(event, listener), { discard: true })) const typed = pubsub.typed.get(event.type) if (typed) yield* PubSub.publish(typed, event) yield* PubSub.publish(pubsub.all, event) @@ -447,7 +443,10 @@ export const layerWith = (options?: LayerOptions) => }) } - function publishMany(events: ReadonlyArray, options?: { readonly location?: Location.Ref }) { + function publishMany( + events: ReadonlyArray, + options?: { readonly location?: Location.Ref; readonly validate?: Effect.Effect }, + ) { return Effect.gen(function* () { const serviceLocation = Option.getOrUndefined(yield* Effect.serviceOption(Location.Service)) const location = @@ -507,6 +506,7 @@ export const layerWith = (options?: LayerOptions) => // Aligned with entries by index: a deduped entry yields // undefined so the payload pairing below stays positional. const results = new Array<{ aggregateID: string; seq: number } | undefined>() + if (options?.validate) yield* options.validate for (const entry of entries) { // No replay input: seq is allocated contiguously from the latest sequence inside the transaction. const result = yield* commitDurableEventInner( diff --git a/packages/core/src/session/message-updater.ts b/packages/core/src/session/message-updater.ts index a97afb58b6..45dd5eb3ec 100644 --- a/packages/core/src/session/message-updater.ts +++ b/packages/core/src/session/message-updater.ts @@ -365,6 +365,7 @@ export function update(adapter: Adapter, event: SessionEvent.Event) { const match = latestReasoning(draft, event.data.reasoningID) if (match) { match.text = event.data.text + if (event.data.distillation !== undefined) match.distillation = event.data.distillation if (event.data.providerMetadata !== undefined) match.providerMetadata = event.data.providerMetadata } }) diff --git a/packages/core/src/session/reasoning-distillation/adopt.ts b/packages/core/src/session/reasoning-distillation/adopt.ts new file mode 100644 index 0000000000..c59a846ff6 --- /dev/null +++ b/packages/core/src/session/reasoning-distillation/adopt.ts @@ -0,0 +1,95 @@ +import { Cause, DateTime, Effect, Schema } from "effect" +import { isDeepStrictEqual } from "node:util" +import { eq } from "drizzle-orm" +import { Database } from "../../database/database" +import { EventV2 } from "../../event" +import { SessionEvent } from "../event" +import { SessionMessage } from "../message" +import { SessionSchema } from "../schema" +import { SessionMessageTable } from "../sql" +import { adoptionProvenance, type ReasoningReplacement } from "./adoption" + +export const adoptReasoning = Effect.fn("CoreReasoningDistillation.adopt")(function* ( + events: EventV2.Interface, + db: Database.Interface["db"], + sessionID: SessionSchema.ID, + messages: readonly SessionMessage.Message[], + replacements: readonly ReasoningReplacement[], + canAdopt: Effect.Effect = Effect.succeed(true), +) { + if (replacements.length === 0) return false + const entries: EventV2.BatchEvent[] = [] + const checks: Array<{ messageID: SessionMessage.ID; part: SessionMessage.AssistantReasoning }> = [] + const seen = new Set() + for (const replacement of replacements) { + const message = messages.find((message) => message.id === replacement.messageID) + if (!message || message.type !== "assistant" || message.time.completed === undefined) return false + const part = message.content.find((part) => part.id === replacement.partID) + const key = JSON.stringify([message.id, replacement.partID]) + if ( + !part || + part.type !== "reasoning" || + part.distillation || + part.text !== replacement.before || + !replacement.after.trim() || + seen.has(key) + ) + return false + seen.add(key) + checks.push({ messageID: message.id, part }) + entries.push({ + definition: SessionEvent.Reasoning.Ended, + data: { + sessionID, + assistantMessageID: message.id, + reasoningID: part.id, + text: replacement.after, + providerMetadata: part.providerMetadata, + distillation: adoptionProvenance(part.text), + timestamp: DateTime.makeUnsafe(Date.now()), + }, + }) + } + const validate = Effect.gen(function* () { + if (!(yield* canAdopt)) yield* Effect.die("reasoning distillation was disabled") + for (const user of messages.filter((message) => message.type === "user")) { + const row = yield* db + .select() + .from(SessionMessageTable) + .where(eq(SessionMessageTable.id, user.id)) + .get() + .pipe(Effect.orDie) + const current = row + ? yield* Schema.decodeUnknownEffect(SessionMessage.Message)({ ...row.data, id: row.id, type: row.type }).pipe( + Effect.orDie, + ) + : undefined + if (row?.session_id !== sessionID || !isDeepStrictEqual(current, user)) + yield* Effect.die("reasoning source user was changed") + } + for (const check of checks) { + const row = yield* db + .select() + .from(SessionMessageTable) + .where(eq(SessionMessageTable.id, check.messageID)) + .get() + .pipe(Effect.orDie) + const message = row + ? yield* Schema.decodeUnknownEffect(SessionMessage.Message)({ ...row.data, id: row.id, type: row.type }).pipe( + Effect.orDie, + ) + : undefined + const current = + message?.type === "assistant" ? message.content.find((part) => part.id === check.part.id) : undefined + if (row?.session_id !== sessionID || !isDeepStrictEqual(current, check.part)) + yield* Effect.die("stale reasoning adoption") + } + }) + return yield* events.publishMany(entries, { validate }).pipe( + Effect.as(true), + Effect.catchCauseIf( + (cause) => !Cause.hasInterrupts(cause), + () => Effect.logWarning("reasoning adoption failed; retaining current history").pipe(Effect.as(false)), + ), + ) +}) diff --git a/packages/core/src/session/reasoning-distillation/adoption.ts b/packages/core/src/session/reasoning-distillation/adoption.ts new file mode 100644 index 0000000000..01a2afec4b --- /dev/null +++ b/packages/core/src/session/reasoning-distillation/adoption.ts @@ -0,0 +1,31 @@ +import { readWirePath } from "../context-folding/wire-value" +import { Hash } from "../../util/hash" + +export type ReasoningReplacement = Readonly<{ + messageID: string + partID: string + before: string + after: string +}> + +/** Collect every applied slot from the final request, including earlier slots in a multi-slot cycle. */ +export function reasoningReplacements( + request: unknown, + slots: readonly Readonly<{ + messageID: string + partID: string + text: string + bodyPath: readonly (string | number)[] + }>[], +): ReasoningReplacement[] { + return slots.flatMap((slot) => { + const value = readWirePath(request, slot.bodyPath) + if (!value.ok || typeof value.value !== "string" || !value.value.trim() || value.value === slot.text) return [] + return [{ messageID: slot.messageID, partID: slot.partID, before: slot.text, after: value.value }] + }) +} + +export const adoptionProvenance = (originalText: string) => ({ + originalText, + sourceFingerprint: Hash.sha256(originalText), +}) diff --git a/packages/core/src/session/reasoning-distillation/schedule.ts b/packages/core/src/session/reasoning-distillation/schedule.ts new file mode 100644 index 0000000000..5c82661cee --- /dev/null +++ b/packages/core/src/session/reasoning-distillation/schedule.ts @@ -0,0 +1,38 @@ +import { Cause, Effect, Scope, Semaphore } from "effect" + +/** One submission per completed user turn. First submission blocks; later ones belong to the service scope. */ +export const makeTurnScheduler = (scope: Scope.Scope) => { + const sessions = new Map; lock: Semaphore.Semaphore }>() + return Effect.fn("ReasoningDistillation.schedule")(function* (input: { + sessionID: string + turnID: string + previouslyDistilled?: boolean + enabled: Effect.Effect + work: Effect.Effect + }) { + if (!(yield* input.enabled)) return + let state = sessions.get(input.sessionID) + const first = !state && !input.previouslyDistilled + if (!state) { + state = { turns: new Set(), lock: Semaphore.makeUnsafe(1) } + sessions.set(input.sessionID, state) + } + if (state.turns.has(input.turnID)) return + state.turns.add(input.turnID) + const work = state.lock + .withPermit( + Effect.gen(function* () { + if (yield* input.enabled) yield* input.work + }), + ) + .pipe( + Effect.catchCause((cause) => + Cause.hasInterrupts(cause) + ? Effect.interrupt + : Effect.logWarning("reasoning distillation failed; retaining current history"), + ), + ) + if (first) yield* work + else yield* work.pipe(Effect.forkIn(scope)) + }) +} diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index b8037c9046..58a6d4045e 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -1,3 +1,5 @@ +import { makeTurnScheduler } from "../reasoning-distillation/schedule" +import { adoptReasoning } from "../reasoning-distillation/adopt" import { LLM, LLMClient, @@ -9,7 +11,7 @@ import { isContextOverflowFailure, type ProviderErrorEvent, } from "@opencode-ai/llm" -import { Cause, DateTime, Deferred, Duration, Effect, FiberSet, Layer, Option, Semaphore, Stream } from "effect" +import { Cause, DateTime, Deferred, Duration, Effect, FiberSet, Layer, Option, Scope, Semaphore, Stream } from "effect" import { AgentV2 } from "../../agent" import { Config } from "../../config" import { ConfigCompaction } from "../../config/compaction" @@ -127,6 +129,8 @@ const toolWaitTimeoutError = () => export const layer = Layer.effect( Service, Effect.gen(function* () { + const scope = yield* Scope.Scope + const scheduleDistillation = makeTurnScheduler(scope) const events = yield* EventV2.Service const llm = yield* LLMClient.Service const agents = yield* AgentV2.Service @@ -394,45 +398,7 @@ export const layer = Layer.effect( projectionPlan: folding.plan, }), ) - let request = folding.request - const reasoningConfig = Config.latest(yield* config.entries(), "reasoningDistillation") - const reasoningResolution = ConfigReasoningDistillation.resolveEnabled({ - disabledByEnvironment: Flag.OPENCODE_DISABLE_REASONING_DISTILLATION, - enabled: reasoningConfig?.enabled, - }) - if ( - reasoningResolution.enabled && - (reasoningConfig?.compatibility?.length ?? 0) > 0 && - conversion.reasoningBindings.length > 0 - ) { - const prepared = Option.getOrUndefined(yield* llm.prepare(request).pipe(Effect.option)) - if (prepared) { - const distilled = yield* reasoningDistillation.distill({ - sessionID: session.id, - ...(session.model?.variant === undefined ? {} : { variant: session.model.variant }), - request, - prepared, - sourceMessages: context, - bindings: conversion.reasoningBindings, - config: reasoningConfig ?? new ConfigReasoningDistillation.Info({}), - }) - request = distilled.request - yield* Effect.logInfo("reasoning distillation", { - "reasoning_distillation.runtime": "core-runner", - "reasoning_distillation.enabled": true, - "reasoning_distillation.source": reasoningResolution.source, - "reasoning_distillation.attempted": distilled.attempted, - "reasoning_distillation.applied": distilled.applied, - "reasoning_distillation.skip_reason": distilled.skipReason ?? "none", - "reasoning_distillation.slot_count": conversion.reasoningBindings.length, - "reasoning_distillation.aux_reserved_tokens": distilled.usage?.reservedTokens ?? 0, - "reasoning_distillation.aux_actual_tokens": distilled.usage?.actualTokens ?? 0, - "reasoning_distillation.aux_unknown_usage_calls": distilled.usage?.unknownUsageCalls ?? 0, - "reasoning_distillation.aux_latency_ms": distilled.usage?.latencyMs ?? 0, - "reasoning_distillation.paid_admission_paused": distilled.usage?.paidAdmissionPaused ?? false, - }) - } - } + const request = folding.request if (yield* compaction.compactIfNeeded({ sessionID: session.id, entries, model, request })) return yield* Effect.die(continueAfterCompaction(currentStep)) const startSnapshot = yield* Snapshot.captureDeduped(history.snapshots, snapshots.capture) @@ -593,6 +559,70 @@ export const layer = Layer.effect( yield* withPublication(batch.flush()) if (stream._tag === "Failure") return yield* Effect.failCause(stream.cause) if (settled._tag === "Failure") return yield* Effect.failCause(settled.cause) + if (!publisher.hasProviderError() && !needsContinuation) { + const completed = yield* getContext(session.id) + const userIndex = completed.findLastIndex((message) => message.type === "user") + const turnID = completed[userIndex]?.id + if (turnID) { + const turnIDs = new Set(completed.slice(userIndex + 1).map((message) => message.id)) + const latestConversion = toLLMMessagesWithBindings(completed, model) + const enabled = config.entries().pipe( + Effect.map( + (entries) => + ConfigReasoningDistillation.resolveEnabled({ + disabledByEnvironment: Flag.OPENCODE_DISABLE_REASONING_DISTILLATION, + enabled: Config.latest(entries, "reasoningDistillation")?.enabled, + }).enabled, + ), + ) + yield* restore( + scheduleDistillation({ + sessionID: session.id, + turnID, + previouslyDistilled: completed + .slice(0, userIndex) + .some( + (message) => + message.type === "assistant" && + message.content.some((part) => part.type === "reasoning" && part.distillation !== undefined), + ), + enabled, + work: Effect.gen(function* () { + const reasoningConfig = Config.latest(yield* config.entries(), "reasoningDistillation") + if (!reasoningConfig?.compatibility?.length) return + const completedRequest = LLM.updateRequest(request, { messages: latestConversion.messages }) + const prepared = yield* llm.prepare(completedRequest) + const distilled = yield* reasoningDistillation.distill({ + sessionID: session.id, + variant: session.model?.variant, + request: completedRequest, + prepared, + sourceMessages: completed, + bindings: latestConversion.reasoningBindings.filter((binding) => + turnIDs.has(SessionMessage.ID.make(binding.ref.messageID)), + ), + config: reasoningConfig, + }) + if (distilled.applied && (yield* enabled)) + yield* adoptReasoning( + events, + db, + session.id, + completed, + distilled.replacements ?? [], + enabled, + ).pipe( + Effect.tap((adopted) => + Effect.sync(() => { + if (adopted) cursors.delete(session.id) + }), + ), + ) + }).pipe(Effect.asVoid), + }), + ) + } + } return { needsContinuation: !publisher.hasProviderError() && needsContinuation, step: currentStep } }), ) diff --git a/packages/core/src/session/runner/reasoning-distillation.ts b/packages/core/src/session/runner/reasoning-distillation.ts index 97bae835ca..a405a127f5 100644 --- a/packages/core/src/session/runner/reasoning-distillation.ts +++ b/packages/core/src/session/runner/reasoning-distillation.ts @@ -1,3 +1,4 @@ +import { reasoningReplacements, type ReasoningReplacement } from "../reasoning-distillation/adoption" import { LLM, LLMResponse, @@ -93,6 +94,7 @@ const emptyState: LifecycleState = { type Slot = Readonly export type Result = Readonly<{ + replacements?: readonly ReasoningReplacement[] request: LLMRequest attempted: "none" | "propose" | "judge" applied: boolean @@ -476,7 +478,7 @@ const plan = (input: Input, slot: Slot, state: LifecycleState, messages: unknown const latestIndex = input.sourceMessages.findIndex((message) => message.id === latestSource) const turn = input.sourceMessages.slice(0, latestIndex + 1).filter((message) => message.type === "user").length const current = cached && isCertificateCurrent(cached, inventory.fingerprint) - const trigger = turn === 0 ? undefined : current ? "replay" : (turn - 1) % 3 === 0 ? "scheduled" : "idle" + const trigger = turn === 0 ? undefined : current ? "replay" : "scheduled" const preparedBudget = budget(input.request, plainValue(input.prepared.body) ?? null) const evidence = { spans: [sourceSpan(slot)], @@ -601,7 +603,7 @@ export const make = (llm: LLMClientShape) => { return yield* lock.withPermits(1)( Effect.gen(function* () { const slots: Slot[] = input.bindings.map((binding) => ({ ...binding, structureRewritable: true })) - const eligible = slots.filter((item) => item.settled && !item.signed && !item.encrypted) + const eligible = slots.filter((item) => item.settled && !item.signed && !item.encrypted && !item.distilled) if (eligible.length === 0) return unchanged(input.request, "no-rewritable-slot") let state = states.get(input.sessionID) ?? emptyState let request = input.request @@ -746,7 +748,28 @@ export const make = (llm: LLMClientShape) => { }) const distill = Effect.fn("CoreReasoningDistillation.distill")(function* (input: Input) { - const result = yield* distillCycle(input) + let cycle: Result = unchanged(input.request) + if (!input.sourceMessages.some((message) => message.type === "user")) { + cycle = yield* distillCycle(input) + } else { + for (const binding of input.bindings) { + const next = yield* distillCycle({ ...input, request: cycle.request, bindings: [binding] }) + cycle = { ...next, applied: cycle.applied || next.applied } + } + } + const result = { + ...cycle, + replacements: cycle.applied + ? reasoningReplacements( + cycle.request, + input.bindings.map((binding) => ({ + ...binding, + messageID: binding.ref.messageID, + partID: binding.ref.partID, + })), + ) + : [], + } if (!result.applied || !input.sourceMessages.some((message) => message.type === "user")) return result // Recheck the fully lowered provider body, including tool schemas and protocol fields. const prepared = yield* llm.prepare(result.request).pipe(Effect.option) diff --git a/packages/core/src/session/runner/to-llm-message.ts b/packages/core/src/session/runner/to-llm-message.ts index 63d501793e..c92b18df6e 100644 --- a/packages/core/src/session/runner/to-llm-message.ts +++ b/packages/core/src/session/runner/to-llm-message.ts @@ -85,6 +85,7 @@ export type ToolMessageBinding = Readonly<{ }> export type ReasoningMessageBinding = Readonly<{ + distilled?: boolean ref: Readonly<{ messageID: string; partID: string }> /** Exact path in the canonical LLM request, resolved without matching on text. */ bodyPath: readonly (string | number)[] @@ -153,6 +154,7 @@ const assistant = (message: SessionMessage.Assistant, model: Model): MessageConv new Set(["encrypted_content", "encryptedContent", "reasoningEncryptedContent"]), ), settled: message.time.completed !== undefined, + distilled: item.distillation !== undefined, }) canonicalContentIndex++ } @@ -278,6 +280,7 @@ export const toLLMMessagesWithBindings = ( signed: binding.signed, encrypted: binding.encrypted, settled: binding.settled, + distilled: binding.distilled, })), ) messageOffset += item.messages.length diff --git a/packages/core/test/config/reasoning-distillation.test.ts b/packages/core/test/config/reasoning-distillation.test.ts index e729c3fb7b..848f87b446 100644 --- a/packages/core/test/config/reasoning-distillation.test.ts +++ b/packages/core/test/config/reasoning-distillation.test.ts @@ -3,9 +3,9 @@ import { Schema } from "effect" import { ConfigReasoningDistillation } from "../../src/config/reasoning-distillation" describe("ConfigReasoningDistillation.resolveEnabled (D02)", () => { - test("defaults on when nothing is set", () => { + test("defaults off when nothing is set", () => { expect(ConfigReasoningDistillation.resolveEnabled({ disabledByEnvironment: false })).toEqual({ - enabled: true, + enabled: false, source: "default", }) }) diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 6df50ac575..68b6b4eaf5 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -1,5 +1,10 @@ +import { ConfigReasoningDistillation } from "../src/config/reasoning-distillation" +import { adapterVersion } from "../src/session/runner/reasoning-distillation" +import { Hash } from "../src/util/hash" import { describe, expect } from "bun:test" import { + LLMResponse, + PreparedRequest, LLMClient, LLMError, LLMEvent, @@ -59,6 +64,24 @@ import { testEffect } from "./lib/effect" const questions = QuestionV2.layer.pipe(Layer.provide(EventV2.defaultLayer)) const requests: LLMRequest[] = [] +let reasoningConfig: ConfigReasoningDistillation.Info | undefined +let auxiliary: LLMClientShape["generate"] | undefined +const prepareDistillation = (request: LLMRequest) => + request.model.route.body + .from(request) + .pipe( + Effect.map( + (body) => + new PreparedRequest({ + id: "distillation-test", + route: request.model.route.id, + protocol: request.model.route.protocol, + model: request.model, + body, + }), + ), + ) + let response: LLMEvent[] = [] let responses: LLMEvent[][] | undefined let responseStream: Stream.Stream | undefined @@ -73,7 +96,8 @@ let maxActiveToolExecutions = 0 const client = Layer.succeed( LLMClient.Service, LLMClient.Service.of({ - prepare: () => Effect.die("unused"), + prepare: ((request: LLMRequest) => + reasoningConfig ? prepareDistillation(request) : Effect.die("unused")) as LLMClientShape["prepare"], stream: ((request: LLMRequest) => { requests.push(request) if (responseStream) { @@ -92,7 +116,7 @@ const client = Layer.succeed( ), ) }) as unknown as LLMClientShape["stream"], - generate: () => Effect.die("unused"), + generate: (request) => (auxiliary ? auxiliary(request) : Effect.die("unused")), }), ) const model = Model.make({ id: "fake-model", provider: "fake", route: OpenAIChat.route }) @@ -222,6 +246,7 @@ const config = Layer.succeed( new Config.Document({ type: "document", info: new Config.Info({ + reasoningDistillation: reasoningConfig, compaction: new ConfigCompaction.Info({ buffer: 3_000, keep: new ConfigCompaction.Keep({ tokens: 1_000 }), @@ -315,6 +340,8 @@ const insertSession = (id: SessionV2.ID) => const setup = Effect.gen(function* () { const { db } = yield* Database.Service response = [] + reasoningConfig = undefined + auxiliary = undefined systemBaseline = "Initial context" systemRemoved = false systemUnavailable = false @@ -3214,3 +3241,123 @@ describe("SessionRunnerLLM", () => { }), ) }) + +it.effect("Core runner adopts every slot at turn completion and asynchronously replaces the next turn", () => + Effect.gen(function* () { + yield* setup + currentModel = Model.make({ + id: "distillation", + provider: "fake", + route: OpenAIChat.route.with({ limits: { context: 100_000, output: 4096 } }), + }) + reasoningConfig = new ConfigReasoningDistillation.Info({ + enabled: true, + compatibility: [ + new ConfigReasoningDistillation.Compatibility({ + runtime: "core-runner", + protocol: currentModel.route.protocol, + providerModelVariant: "fake/distillation/default", + endpointIdentity: Hash.sha256(currentModel.route.id), + adapterVersion, + optionsFingerprint: Hash.sha256(JSON.stringify({ openai: { promptCacheKey: sessionID } })), + transportVerified: true, + upstreamVerified: true, + }), + ], + }) + const started = yield* Deferred.make() + const release = yield* Deferred.make() + let calls = 0 + auxiliary = (request) => + Effect.gen(function* () { + if (calls++ === 4) { + yield* Deferred.succeed(started, undefined) + yield* Deferred.await(release) + } + const prompt = request.messages + .flatMap((message) => message.content.flatMap((part) => (part.type === "text" ? [part.text] : []))) + .join("\n") + const ref = /messageID=([^,]+),partID=([^,]+),长度=(\d+)/.exec(prompt) + const output = ref + ? { + claims: [ + { + id: "claim", + kind: "decision", + text: `采用-${ref[2]}`, + scope: "当前回合", + sources: [{ messageID: ref[1], partID: ref[2], start: 0, end: Number(ref[3]) }], + evidence: [], + status: "unverified", + }, + ], + preserved: [], + coverage: [ + { + source: { messageID: ref[1], partID: ref[2], start: 0, end: Number(ref[3]) }, + action: "keep", + claimID: "claim", + }, + ], + } + : { + retention: { verdict: "supported" }, + support: [{ claimID: "claim", verdict: "supported", method: "judged" }], + } + const result = LLMResponse.fromEvents([ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.textStart({ id: "aux" }), + LLMEvent.textDelta({ id: "aux", text: JSON.stringify(output) }), + LLMEvent.textEnd({ id: "aux" }), + LLMEvent.stepFinish({ + index: 0, + reason: "stop", + usage: { inputTokens: 10, outputTokens: 10, totalTokens: 20 }, + }), + LLMEvent.finish({ reason: "stop" }), + ]) + if (!result) throw new Error("missing auxiliary response") + return result + }) + const session = yield* SessionV2.Service + const events = yield* EventV2.Service + const answer = (turn: number) => [ + LLMEvent.stepStart({ index: 0 }), + ...[0, 1].flatMap((slot) => [ + LLMEvent.reasoningStart({ id: `r-${turn}-${slot}` }), + LLMEvent.reasoningDelta({ id: `r-${turn}-${slot}`, text: `original-${turn}-${slot}。`.repeat(1000) }), + LLMEvent.reasoningEnd({ id: `r-${turn}-${slot}` }), + ]), + LLMEvent.stepFinish({ index: 0, reason: "stop" }), + LLMEvent.finish({ reason: "stop" }), + ] + yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "first" }), resume: false }) + response = answer(1) + yield* session.resume(sessionID) + expect(calls).toBe(4) + expect(JSON.stringify(yield* session.context(sessionID))).toContain("采用-r-1-0") + expect(JSON.stringify(yield* session.context(sessionID))).toContain("采用-r-1-1") + + const updates = yield* events.subscribe(SessionEvent.Reasoning.Ended).pipe( + Stream.filter((event) => event.data.distillation !== undefined), + Stream.take(2), + Stream.runDrain, + Effect.forkChild, + ) + yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "second" }), resume: false }) + response = answer(2) + yield* session.resume(sessionID) + yield* Deferred.await(started) + expect(JSON.stringify(requests.at(-1)?.messages)).toContain("采用-r-1-0") + expect(JSON.stringify(requests.at(-1)?.messages)).not.toContain("original-1") + expect(JSON.stringify(yield* session.context(sessionID))).not.toContain("采用-r-2-0") + yield* Deferred.succeed(release, undefined) + yield* Fiber.join(updates) + expect(calls).toBe(8) + yield* session.prompt({ sessionID, prompt: Prompt.make({ text: "third" }), resume: false }) + response = [] + yield* session.resume(sessionID) + expect(JSON.stringify(requests.at(-1)?.messages)).toContain("采用-r-2-0") + expect(JSON.stringify(requests.at(-1)?.messages)).not.toContain("original-2") + }), +) diff --git a/packages/core/test/session/reasoning-adoption.test.ts b/packages/core/test/session/reasoning-adoption.test.ts new file mode 100644 index 0000000000..3b91cc57eb --- /dev/null +++ b/packages/core/test/session/reasoning-adoption.test.ts @@ -0,0 +1,145 @@ +import { describe, expect } from "bun:test" +import { DateTime, Effect, Layer, Schema } from "effect" +import { eq } from "drizzle-orm" +import { Database } from "../../src/database/database" +import { EventV2 } from "../../src/event" +import { EventTable } from "../../src/event/sql" +import { Project } from "../../src/project" +import { ProjectTable } from "../../src/project/sql" +import { AbsolutePath } from "../../src/schema" +import { ModelV2 } from "../../src/model" +import { ProviderV2 } from "../../src/provider" +import { SessionSchema } from "../../src/session/schema" +import { SessionMessage } from "../../src/session/message" +import { SessionProjector } from "../../src/session/projector" +import { SessionStore } from "../../src/session/store" +import { SessionMessageTable, SessionTable } from "../../src/session/sql" +import { adoptReasoning } from "../../src/session/reasoning-distillation/adopt" +import { reasoningReplacements } from "../../src/session/reasoning-distillation/adoption" +import { toLLMMessagesWithBindings } from "../../src/session/runner/to-llm-message" +import { Model } from "@opencode-ai/llm" +import * as OpenAIChat from "@opencode-ai/llm/protocols/openai-chat" +import { testEffect } from "../lib/effect" + +const it = testEffect( + Layer.mergeAll(Database.defaultLayer, EventV2.defaultLayer, SessionProjector.defaultLayer, SessionStore.defaultLayer), +) +const model = Model.make({ id: "model", provider: "provider", route: OpenAIChat.route }) + +const seed = Effect.fnUntraced(function* () { + const { db } = yield* Database.Service + const events = yield* EventV2.Service + const store = yield* SessionStore.Service + const sessionID = SessionSchema.ID.make("ses_adoption") + 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: "adoption", + directory: "/project", + title: "adoption", + version: "test", + }) + .run() + .pipe(Effect.orDie) + const message: SessionMessage.Assistant = { + id: SessionMessage.ID.make("msg_adoption"), + type: "assistant", + agent: "build", + model: { id: ModelV2.ID.make("model"), providerID: ProviderV2.ID.make("provider") }, + time: { created: DateTime.makeUnsafe(1), completed: DateTime.makeUnsafe(2) }, + content: [ + { type: "reasoning", id: "r1", text: "first original" }, + { type: "text", id: "t1", text: "answer unchanged" }, + { type: "reasoning", id: "r2", text: "second original" }, + ], + } + const { id: _id, type, ...data } = Schema.encodeSync(SessionMessage.Assistant)(message) + yield* db + .insert(SessionMessageTable) + .values({ id: message.id, type, data, session_id: sessionID, seq: 0, time_created: 1 }) + .run() + .pipe(Effect.orDie) + const replacements = [ + { messageID: message.id, partID: "r1", before: "first original", after: "first adopted" }, + { messageID: message.id, partID: "r2", before: "second original", after: "second adopted" }, + ] + return { db, events, store, sessionID, message, replacements } +}) + +describe("durable reasoning adoption", () => { + it.effect("persists all accepted slots and reloads the same text into model context without provenance leakage", () => + Effect.gen(function* () { + const { db, events, store, sessionID, message, replacements } = yield* seed() + expect(yield* adoptReasoning(events, db, sessionID, [message], replacements)).toBe(true) + const stored = yield* store.message(message.id) + expect(stored?.message).toMatchObject({ + content: [ + { type: "reasoning", id: "r1", text: "first adopted", distillation: { originalText: "first original" } }, + { type: "text", id: "t1", text: "answer unchanged" }, + { type: "reasoning", id: "r2", text: "second adopted", distillation: { originalText: "second original" } }, + ], + }) + if (!stored) throw new Error("missing stored message") + const conversion = toLLMMessagesWithBindings([stored.message], model) + expect(conversion.reasoningBindings.every((part) => part.distilled)).toBe(true) + expect(JSON.stringify(conversion.messages)).not.toContain("original") + expect(JSON.stringify(conversion.messages)).toContain("first adopted") + const switched = toLLMMessagesWithBindings( + [stored.message], + Model.make({ id: "other", provider: model.provider, route: model.route }), + ) + expect(JSON.stringify(switched.messages)).toContain("first adopted") + expect(JSON.stringify(switched.messages)).not.toContain("original") + expect(yield* adoptReasoning(events, db, sessionID, [stored.message], replacements)).toBe(false) + }), + ) + + it.effect("rejects the entire batch before emitting events when any persisted source has changed", () => + Effect.gen(function* () { + const { db, events, store, sessionID, message, replacements } = yield* seed() + const changed = { + ...message, + content: message.content.map((part) => (part.id === "r2" ? { ...part, text: "newer second" } : part)), + } + const { id: _id, type: _type, ...data } = Schema.encodeSync(SessionMessage.Assistant)(changed) + yield* db + .update(SessionMessageTable) + .set({ data }) + .where(eq(SessionMessageTable.id, message.id)) + .run() + .pipe(Effect.orDie) + expect(yield* adoptReasoning(events, db, sessionID, [message], replacements)).toBe(false) + expect((yield* store.message(message.id))?.message).toEqual(changed) + expect(yield* db.select().from(EventTable).all().pipe(Effect.orDie)).toEqual([]) + }), + ) + + it.effect("extracts every changed wire slot rather than only the last cycle plan", () => + Effect.sync(() => { + expect( + reasoningReplacements({ messages: [{ reasoning: "first adopted" }, { reasoning: "second adopted" }] }, [ + { messageID: "m1", partID: "r1", text: "first original", bodyPath: ["messages", 0, "reasoning"] }, + { messageID: "m2", partID: "r2", text: "second original", bodyPath: ["messages", 1, "reasoning"] }, + ]), + ).toEqual([ + { messageID: "m1", partID: "r1", before: "first original", after: "first adopted" }, + { messageID: "m2", partID: "r2", before: "second original", after: "second adopted" }, + ]) + }), + ) + it.effect("checks the switch again inside the adoption transaction", () => + Effect.gen(function* () { + const { db, events, store, sessionID, message, replacements } = yield* seed() + expect(yield* adoptReasoning(events, db, sessionID, [message], replacements, Effect.succeed(false))).toBe(false) + expect((yield* store.message(message.id))?.message).toEqual(message) + expect(yield* db.select().from(EventTable).all().pipe(Effect.orDie)).toEqual([]) + }), + ) +}) diff --git a/packages/core/test/session/reasoning-schedule.test.ts b/packages/core/test/session/reasoning-schedule.test.ts new file mode 100644 index 0000000000..a6b34e38b4 --- /dev/null +++ b/packages/core/test/session/reasoning-schedule.test.ts @@ -0,0 +1,141 @@ +import { expect, test } from "bun:test" +import { Deferred, Effect, Fiber, Scope } from "effect" +import { makeTurnScheduler } from "../../src/session/reasoning-distillation/schedule" + +test("default-off admission, first completion blocks, later turns run in order without blocking", async () => { + await Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const schedule = makeTurnScheduler(yield* Scope.Scope) + let enabled = false + const calls: string[] = [] + const first = yield* Deferred.make() + const second = yield* Deferred.make() + const started = yield* Deferred.make() + const background = yield* Deferred.make() + const input = (turnID: string, work: Effect.Effect) => ({ + sessionID: "s", + turnID, + enabled: Effect.sync(() => enabled), + work, + }) + yield* schedule( + input( + "off", + Effect.sync(() => { + calls.push("off") + }), + ), + ) + expect(calls).toEqual([]) + enabled = true + const sync = yield* schedule( + input( + "1", + Effect.gen(function* () { + calls.push("1-start") + yield* Deferred.succeed(started, undefined) + yield* Deferred.await(first) + calls.push("1-end") + }), + ), + ).pipe(Effect.forkChild) + yield* Deferred.await(started) + expect(calls).toEqual(["1-start"]) + expect(sync.pollUnsafe()).toBeUndefined() + yield* Deferred.succeed(first, undefined) + yield* Fiber.join(sync) + yield* schedule( + input( + "2", + Effect.gen(function* () { + calls.push("2-start") + yield* Deferred.succeed(background, undefined) + yield* Deferred.await(second) + calls.push("2-end") + }), + ), + ) + yield* Deferred.await(background) + expect(calls).toEqual(["1-start", "1-end", "2-start"]) + const done = yield* Deferred.make() + yield* schedule( + input( + "3", + Effect.gen(function* () { + calls.push("3") + yield* Deferred.succeed(done, undefined) + }), + ), + ) + yield* schedule( + input( + "3", + Effect.sync(() => { + calls.push("duplicate") + }), + ), + ) + yield* Deferred.succeed(second, undefined) + yield* Deferred.await(done) + expect(calls).toEqual(["1-start", "1-end", "2-start", "2-end", "3"]) + }), + ), + ) +}) + +test("turn failure retains progress and disabled queued turns do no work", async () => { + await Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const schedule = makeTurnScheduler(yield* Scope.Scope) + let enabled = true + const input = (turnID: string, work: Effect.Effect) => ({ + sessionID: "s", + turnID, + enabled: Effect.sync(() => enabled), + work, + }) + yield* schedule(input("1", Effect.fail("auxiliary failed"))) + const busy = yield* Deferred.make() + const release = yield* Deferred.make() + yield* schedule(input("2", Deferred.succeed(busy, undefined).pipe(Effect.andThen(Deferred.await(release))))) + yield* Deferred.await(busy) + let ran = false + yield* schedule( + input( + "3", + Effect.sync(() => { + ran = true + }), + ), + ) + enabled = false + yield* Deferred.succeed(release, undefined) + yield* Effect.yieldNow + expect(ran).toBe(false) + }), + ), + ) +}) + +test("reopened history with an adopted prior turn resumes background scheduling", async () => { + await Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const schedule = makeTurnScheduler(yield* Scope.Scope) + const started = yield* Deferred.make() + const release = yield* Deferred.make() + yield* schedule({ + sessionID: "restored", + turnID: "next", + previouslyDistilled: true, + enabled: Effect.succeed(true), + work: Deferred.succeed(started, undefined).pipe(Effect.andThen(Deferred.await(release))), + }) + yield* Deferred.await(started) + yield* Deferred.succeed(release, undefined) + }), + ), + ) +}) diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index 605b9ab749..f739381cbe 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -1,10 +1,14 @@ +import { + reasoningReplacements, + type ReasoningReplacement, +} from "@opencode-ai/core/session/reasoning-distillation/adoption" import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { llmClient } from "@opencode-ai/core/effect/layer-node-platform" import { PermissionV1 } from "@opencode-ai/core/v1/permission" import { Provider } from "@/provider/provider" import { SessionV1 } from "@opencode-ai/core/v1/session" import { serviceUse } from "@opencode-ai/core/effect/service-use" -import { Context, Effect, Layer } from "effect" +import { Context, Effect, Layer, Semaphore } from "effect" import * as Stream from "effect/Stream" import { asSchema, generateText, streamText, wrapLanguageModel, type ModelMessage, type Tool } from "ai" import { LLMRequest, type LLMEvent } from "@opencode-ai/llm" @@ -134,6 +138,7 @@ export type StreamInput = { purpose?: RequestPurpose contextFolding?: ContextFoldingHistorySnapshot reasoningDistillation?: ReasoningDistillationHistorySnapshot + adoptReasoning?: (replacements: readonly ReasoningReplacement[]) => Effect.Effect } export type StreamRequest = StreamInput & { @@ -141,6 +146,7 @@ export type StreamRequest = StreamInput & { } export interface Interface { + readonly distill: (input: StreamInput) => Effect.Effect readonly stream: (input: StreamInput) => Stream.Stream } @@ -171,13 +177,19 @@ const live: Layer.Layer< const llmClient = yield* LLMClient.Service const flags = yield* RuntimeFlags.Service const distillationState = yield* InstanceState.make(() => - Effect.succeed({ - current: ReasoningDistillation.emptyLifecycleState, - busy: false, - }), + Effect.succeed( + new Map< + string, + { + current: typeof ReasoningDistillation.emptyLifecycleState + lock: Semaphore.Semaphore + generation: number + } + >(), + ), ) - const run = Effect.fn("LLM.run")(function* (input: StreamRequest) { + const run = Effect.fn("LLM.run")(function* (input: StreamRequest, distillationOnly = false) { yield* Effect.logInfo("stream", { providerID: input.model.providerID, modelID: input.model.id, @@ -231,7 +243,7 @@ const live: Layer.Layer< ? input.model.capabilities.interleaved.field : undefined const organizer = - distillationResolution.enabled && interleavedField + distillationOnly && distillationResolution.enabled && interleavedField ? ((yield* provider.getSmallModel(input.model.providerID)) ?? input.model) : undefined const organizerResolution = organizer @@ -374,25 +386,18 @@ const live: Layer.Layer< media: hasMedia(args.sourceMessages) ? ("unknown" as const) : ("none" as const), } const estimatedBudget = estimateContextFoldingBudget(budget) - return InstanceState.useEffect(distillationState, (state) => { - // Non-blocking admission: while one auxiliary call is in flight, another request's cycle is skipped - // instead of queueing behind it. The old semaphore stalled unrelated sessions for up to the whole aux - // timeout; skipping only forgoes an opportunistic distillation cycle, and the next ordinary request - // runs its own. - if (state.busy) - return Effect.logInfo("reasoning distillation skipped: an auxiliary call is already in flight", { - "reasoning_distillation.runtime": args.runtime, - "reasoning_distillation.enabled": true, - "reasoning_distillation.source": distillationResolution.source, - "reasoning_distillation.attempted": "none", - "reasoning_distillation.applied": false, - "reasoning_distillation.skip_reason": "aux-busy", - }).pipe(Effect.as(args.request)) - state.busy = true - return Effect.tryPromise({ - try: async () => { - try { - return await ReasoningDistillation.runDistillationCycle(state.current, { + return InstanceState.useEffect(distillationState, (states) => { + let state = states.get(input.sessionID) + if (!state) { + state = { current: ReasoningDistillation.emptyLifecycleState, lock: Semaphore.makeUnsafe(1), generation: 0 } + states.set(input.sessionID, state) + } + const current = state + return current.lock.withPermit( + Effect.tryPromise({ + try: async () => { + const generation = ++current.generation + return await ReasoningDistillation.runDistillationCycle(current.current, { request: args.request, identity: { adapter: args.adapterVersion, @@ -419,39 +424,69 @@ const live: Layer.Layer< originalTokens: Math.ceil(selected.reduce((total, slot) => total + slot.text.length, 0) / 4), callPropose: callAuxiliary, callJudge: callAuxiliary, - commitState: (next) => void (state.current = next), + commitState: (next) => { + if (!input.abort.aborted && current.generation === generation) current.current = next + }, }) - } finally { - state.busy = false - } - }, - catch: (cause) => cause, - }).pipe( - Effect.tap((result) => Effect.sync(() => void (state.current = result.state))), - Effect.tap((cycle) => { - const usage = cycle.state.usageBySession[input.sessionID] - return Effect.logInfo("reasoning distillation", { - "session.id": input.sessionID, - "reasoning_distillation.turn": input.reasoningDistillation?.reasoningTurn ?? 0, - "reasoning_distillation.capability": capability, - "reasoning_distillation.runtime": args.runtime, - "reasoning_distillation.enabled": true, - "reasoning_distillation.source": distillationResolution.source, - "reasoning_distillation.attempted": cycle.attempted, - "reasoning_distillation.applied": cycle.projection.applied, - "reasoning_distillation.skip_reason": cycle.projection.skipReason ?? "none", - "reasoning_distillation.slot_count": args.slots.length, - "reasoning_distillation.estimated_input_tokens": estimatedBudget.estimatedInputTokens ?? "unknown", - "reasoning_distillation.target_tokens": estimatedBudget.targetTokens ?? "unknown", - "reasoning_distillation.budget_skip_reason": estimatedBudget.skipReason ?? "none", - "reasoning_distillation.aux_reserved_tokens": usage?.reservedTokens ?? 0, - "reasoning_distillation.aux_actual_tokens": usage?.actualTokens ?? 0, - "reasoning_distillation.aux_unknown_usage_calls": usage?.unknownUsageCalls ?? 0, - "reasoning_distillation.aux_latency_ms": usage?.latencyMs ?? 0, - "reasoning_distillation.paid_admission_paused": usage?.paidAdmissionPaused ?? false, - }) - }), - Effect.map((cycle) => cycle.projection.request), + }, + catch: (cause) => cause, + }).pipe( + Effect.tap((result) => Effect.sync(() => void (current.current = result.state))), + Effect.flatMap((cycle) => { + if (!cycle.projection.applied || !input.adoptReasoning) return Effect.succeed(cycle) + const replacements = reasoningReplacements(cycle.projection.request, selected) + return Effect.gen(function* () { + const latest = yield* config.get() + if ( + input.abort.aborted || + !ConfigReasoningDistillation.resolveEnabled({ + disabledByEnvironment: Flag.OPENCODE_DISABLE_REASONING_DISTILLATION, + enabled: latest.reasoningDistillation?.enabled, + }).enabled + ) + return false + return yield* input.adoptReasoning!(replacements) + }).pipe( + Effect.map((adopted) => + adopted + ? cycle + : { + ...cycle, + projection: { + ...cycle.projection, + request: args.request, + applied: false, + skipReason: "projection-failed" as const, + }, + }, + ), + ) + }), + Effect.tap((cycle) => { + const usage = cycle.state.usageBySession[input.sessionID] + return Effect.logInfo("reasoning distillation", { + "session.id": input.sessionID, + "reasoning_distillation.turn": input.reasoningDistillation?.reasoningTurn ?? 0, + "reasoning_distillation.capability": capability, + "reasoning_distillation.runtime": args.runtime, + "reasoning_distillation.enabled": true, + "reasoning_distillation.source": distillationResolution.source, + "reasoning_distillation.attempted": cycle.attempted, + "reasoning_distillation.applied": cycle.projection.applied, + "reasoning_distillation.skip_reason": cycle.projection.skipReason ?? "none", + "reasoning_distillation.slot_count": args.slots.length, + "reasoning_distillation.estimated_input_tokens": estimatedBudget.estimatedInputTokens ?? "unknown", + "reasoning_distillation.target_tokens": estimatedBudget.targetTokens ?? "unknown", + "reasoning_distillation.budget_skip_reason": estimatedBudget.skipReason ?? "none", + "reasoning_distillation.aux_reserved_tokens": usage?.reservedTokens ?? 0, + "reasoning_distillation.aux_actual_tokens": usage?.actualTokens ?? 0, + "reasoning_distillation.aux_unknown_usage_calls": usage?.unknownUsageCalls ?? 0, + "reasoning_distillation.aux_latency_ms": usage?.latencyMs ?? 0, + "reasoning_distillation.paid_admission_paused": usage?.paidAdmissionPaused ?? false, + }) + }), + Effect.map((cycle) => cycle.projection.request), + ), ) }).pipe( Effect.catch((cause) => @@ -591,7 +626,11 @@ const live: Layer.Layer< abort: input.abort, contextFolding: folding, reasoningDistillation: - distillationResolution.enabled && interleavedField && organizerResolution && input.reasoningDistillation + distillationOnly && + distillationResolution.enabled && + interleavedField && + organizerResolution && + input.reasoningDistillation ? ({ request, sourceMessages, transformedMessages }) => { const messages = plainWireMessages(request.messages) if (!messages) return Effect.succeed(request) @@ -661,6 +700,42 @@ const live: Layer.Layer< }) } + if (distillationOnly) { + if (distillationResolution.enabled && interleavedField && input.reasoningDistillation) { + const sourceMessages = ContextFolding.copyModelMessages(prepared.messages) + if (sourceMessages) { + const transformed = ProviderTransform.message( + prepared.messages, + input.model, + prepared.messageTransformOptions, + ) + const messages = plainWireMessages(transformed) + if (messages) { + const lineage = ReasoningDistillation.bindInterleavedReasoningLineage( + sourceMessages, + transformed, + interleavedField, + input.reasoningDistillation, + ) + yield* distillRequest({ + runtime: "opencode-ai-sdk", + adapterVersion: REASONING_DISTILLATION_ADAPTER_VERSION, + request: { messages }, + messages, + sourceMessages, + slots: ReasoningDistillation.extractInterleavedReasoningSlots( + messages, + interleavedField, + ["messages"], + lineage, + ), + }) + } + } + } + return { type: "native" as const, stream: Stream.empty } + } + yield* Effect.logInfo("llm runtime selected", { "llm.runtime": "ai-sdk", "llm.provider": input.model.providerID, @@ -753,40 +828,6 @@ const live: Layer.Layer< projectionPlan = projected.plan if (projected.applied) outbound = projected.request.messages } - if ( - distillationResolution.enabled && - interleavedField && - organizerResolution && - sourceMessages && - input.reasoningDistillation - ) { - const messages = plainWireMessages(outbound) - if (messages) { - const lineage = ReasoningDistillation.bindInterleavedReasoningLineage( - sourceMessages, - transformed, - interleavedField, - input.reasoningDistillation, - ) - const observed = ReasoningDistillation.extractInterleavedReasoningSlots( - messages, - interleavedField, - ["messages"], - lineage, - ) - const projected = await bridge.promise( - distillRequest({ - runtime: "opencode-ai-sdk", - adapterVersion: REASONING_DISTILLATION_ADAPTER_VERSION, - request: { messages }, - messages, - sourceMessages, - slots: observed, - }), - ) - outbound = projected.messages as ModelMessage[] - } - } await bridge.promise( Effect.logInfo( "context folding", @@ -846,7 +887,17 @@ const live: Layer.Layer< ), ) - return Service.of({ stream }) + const distill: Interface["distill"] = (input) => + Effect.scoped( + Effect.gen(function* () { + const ctrl = yield* Effect.acquireRelease( + Effect.sync(() => new AbortController()), + (ctrl) => Effect.sync(() => ctrl.abort()), + ) + yield* run({ ...input, abort: ctrl.signal }, true) + }), + ) + return Service.of({ stream, distill }) }), ) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 94d11089e6..b7f22451f4 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -388,6 +388,9 @@ export const layer = Layer.effect( sessionID: ctx.assistantMessage.sessionID, type: "reasoning", text: "", + ...(ctx.v2AssistantMessageID + ? { v2: { messageID: ctx.v2AssistantMessageID, reasoningID: value.id } } + : {}), time: { start: Date.now() }, metadata: value.providerMetadata, } diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index 402bf2577b..f2398ba8bd 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -1,3 +1,7 @@ +import { Flag } from "@opencode-ai/core/flag/flag" +import { makeTurnScheduler } from "@opencode-ai/core/session/reasoning-distillation/schedule" +import { ConfigReasoningDistillation } from "@opencode-ai/core/config/reasoning-distillation" +import { adoptReasoning } from "./reasoning-adoption" import { withHookFeedback } from "@/hook/trigger-result" import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { PermissionV1 } from "@opencode-ai/core/v1/permission" @@ -224,6 +228,7 @@ export const layer = Layer.effect( const todoSvc = yield* Todo.Service const spawner = yield* ChildProcessSpawner.ChildProcessSpawner const scope = yield* Scope.Scope + const scheduleDistillation = makeTurnScheduler(scope) const instruction = yield* Instruction.Service const state = yield* SessionRunState.Service const revert = yield* SessionRevert.Service @@ -1767,6 +1772,7 @@ export const layer = Layer.effect( const runLoop: (sessionID: SessionID) => Effect.Effect = Effect.fn("SessionPrompt.run")( function* (sessionID: SessionID) { const ctx = yield* InstanceState.context + let distillationInput: LLM.StreamInput | undefined let structured: unknown let step = 0 // Error carried by the assistant message when the turn loop breaks — drives @@ -1923,6 +1929,61 @@ export const layer = Layer.effect( }) } } + if ( + !turnError && + distillationInput && + ConfigReasoningDistillation.resolveEnabled({ + disabledByEnvironment: Flag.OPENCODE_DISABLE_REASONING_DISTILLATION, + enabled: (yield* config.get()).reasoningDistillation?.enabled, + }).enabled + ) { + const sources = yield* sessions.messages({ sessionID }).pipe(Effect.orDie) + const current = sources.filter( + (message) => message.info.role === "assistant" && message.info.parentID === lastUser.id, + ) + const history = ReasoningDistillation.reasoningHistory(sources) + const ids = new Set(current.map((message) => message.info.id)) + const snapshot = { + ...history, + groups: history.groups.map((group) => + ids.has(group.messageID) + ? group + : { ...group, parts: group.parts.map((part) => ({ ...part, distilled: true })) }, + ), + } + const modelMessages = yield* MessageV2.toModelMessagesEffect(sources, distillationInput.model) + const enabled = config.get().pipe( + Effect.map( + (cfg) => + ConfigReasoningDistillation.resolveEnabled({ + disabledByEnvironment: Flag.OPENCODE_DISABLE_REASONING_DISTILLATION, + enabled: cfg.reasoningDistillation?.enabled, + }).enabled, + ), + ) + yield* scheduleDistillation({ + sessionID, + turnID: lastUser.id, + previouslyDistilled: sources.some( + (message) => + !ids.has(message.info.id) && + message.parts.some((part) => part.type === "reasoning" && part.distillation !== undefined), + ), + enabled, + work: llm.distill({ + ...distillationInput, + messages: modelMessages, + contextFolding: undefined, + reasoningDistillation: snapshot, + adoptReasoning: (replacements) => + adoptReasoning({ sessionID, sources, replacements, canAdopt: enabled }).pipe( + Effect.provideService(Database.Service, database), + Effect.provideService(Session.Service, sessions), + Effect.provideService(EventV2Bridge.Service, events), + ), + }), + }) + } return yield* getLastAssistant(sessionID) } } @@ -2136,7 +2197,6 @@ export const layer = Layer.effect( memoryDocs, modelMsgs, contextFoldingHistory, - reasoningDistillationHistory, ] = yield* Effect.all( [ sys.skills(agent), @@ -2148,7 +2208,6 @@ export const layer = Layer.effect( sys.memory({ sessionID, messages: msgs, main: !session.parentID }), MessageV2.toModelMessagesEffect(msgs, model), MessageV2.contextFoldingHistory({ messages: msgs, ledger: toolSources }), - Effect.sync(() => ReasoningDistillation.reasoningHistory(msgs)), ], { concurrency: "unbounded" }, ) @@ -2163,7 +2222,7 @@ export const layer = Layer.effect( ] const format = lastUser.format ?? { type: "text" as const } if (format.type === "json_schema") system.push(STRUCTURED_OUTPUT_SYSTEM_PROMPT) - const result = yield* handle.process({ + const processInput: LLM.StreamInput = { user: lastUser, agent, permission: session.permission, @@ -2179,8 +2238,9 @@ export const layer = Layer.effect( toolChoice: format.type === "json_schema" ? "required" : undefined, purpose: "conversation", contextFolding: contextFoldingHistory, - reasoningDistillation: reasoningDistillationHistory, - }) + } + distillationInput = processInput + const result = yield* handle.process(processInput) if (structured !== undefined) { handle.message.structured = structured diff --git a/packages/opencode/src/session/reasoning-adoption.ts b/packages/opencode/src/session/reasoning-adoption.ts new file mode 100644 index 0000000000..d7f64f8a38 --- /dev/null +++ b/packages/opencode/src/session/reasoning-adoption.ts @@ -0,0 +1,142 @@ +import { Database } from "@opencode-ai/core/database/database" +import { SessionMessageTable } from "@opencode-ai/core/session/sql" +import { eq } from "drizzle-orm" +import { Cause, DateTime, Effect, Schema } from "effect" +import { isDeepStrictEqual } from "node:util" +import { EventV2 } from "@opencode-ai/core/event" +import { SessionV1 } from "@opencode-ai/core/v1/session" +import { SessionEvent } from "@opencode-ai/core/session/event" +import { SessionMessage } from "@opencode-ai/core/session/message" +import { + adoptionProvenance, + type ReasoningReplacement, +} from "@opencode-ai/core/session/reasoning-distillation/adoption" +import { EventV2Bridge } from "@/event-v2-bridge" +import { Session } from "./session" +import type { SessionID } from "./schema" + +/** Commit the same accepted text consumed by the outgoing request and by existing live/reload UI paths. */ +export const adoptReasoning = Effect.fn("Session.adoptReasoning")(function* (input: { + sessionID: SessionID + sources: readonly SessionV1.WithParts[] + replacements: readonly ReasoningReplacement[] + canAdopt?: Effect.Effect +}) { + if (input.replacements.length === 0) return false + const session = yield* Session.Service + const events = yield* EventV2Bridge.Service + const { db } = yield* Database.Service + const entries: EventV2.BatchEvent[] = [] + const sources: SessionV1.ReasoningPart[] = [] + const seen = new Set() + for (const replacement of input.replacements) { + const part = input.sources + .find((message) => message.info.id === replacement.messageID && message.info.role === "assistant") + ?.parts.find((part) => part.id === replacement.partID) + if ( + !part || + part.type !== "reasoning" || + part.sessionID !== input.sessionID || + part.text !== replacement.before || + part.time.end === undefined || + part.distillation || + !replacement.after.trim() || + seen.has(part.id) + ) + return false + seen.add(part.id) + sources.push(part) + const distillation = adoptionProvenance(part.text) + entries.push({ + definition: SessionV1.Event.PartUpdated, + data: { + sessionID: input.sessionID, + part: { ...part, text: replacement.after, distillation }, + time: Date.now(), + }, + }) + if (part.v2) + entries.push({ + definition: SessionEvent.Reasoning.Ended, + data: { + sessionID: input.sessionID, + assistantMessageID: SessionMessage.ID.make(part.v2.messageID), + reasoningID: part.v2.reasoningID, + text: replacement.after, + providerMetadata: part.metadata, + distillation, + timestamp: DateTime.makeUnsafe(Date.now()), + }, + }) + } + // Validate before any projector runs, inside the same transaction as the entire event batch. + const validate = Effect.gen(function* () { + if (input.canAdopt && !(yield* input.canAdopt)) yield* Effect.die("reasoning distillation was disabled") + const currentSession = yield* session.get(input.sessionID).pipe(Effect.orDie) + if (currentSession.revert) yield* Effect.die("reasoning source was reverted") + const currentMessages = yield* session.messages({ sessionID: input.sessionID }).pipe(Effect.orDie) + for (const user of input.sources.filter((message) => message.info.role === "user")) { + if ( + !isDeepStrictEqual( + currentMessages.find((message) => message.info.id === user.info.id), + user, + ) + ) + yield* Effect.die("reasoning source user was changed") + } + const sourceIDs = new Set(input.sources.map((message) => message.info.id)) + const parentIDs = new Set( + input.sources + .filter( + (message) => message.info.role === "assistant" && sources.some((part) => part.messageID === message.info.id), + ) + .map((message) => (message.info.role === "assistant" ? message.info.parentID : undefined)), + ) + if ( + currentMessages.some( + (message) => + message.info.role === "assistant" && parentIDs.has(message.info.parentID) && !sourceIDs.has(message.info.id), + ) + ) + yield* Effect.die("reasoning turn was retried") + for (const source of sources) { + const current = yield* session.getPart({ + sessionID: input.sessionID, + messageID: source.messageID, + partID: source.id, + }) + if (!isDeepStrictEqual(current, source)) yield* Effect.die("stale reasoning adoption") + if (source.v2) { + const row = yield* db + .select() + .from(SessionMessageTable) + .where(eq(SessionMessageTable.id, SessionMessage.ID.make(source.v2.messageID))) + .get() + .pipe(Effect.orDie) + const message = row + ? yield* Schema.decodeUnknownEffect(SessionMessage.Message)({ ...row.data, id: row.id, type: row.type }).pipe( + Effect.orDie, + ) + : undefined + const content = message?.type === "assistant" ? message.content : [] + const matches = content.filter((part) => part.id === source.v2?.reasoningID) + if ( + row?.session_id !== input.sessionID || + matches.length !== 1 || + matches[0].type !== "reasoning" || + matches[0].text !== source.text || + matches[0].distillation || + !isDeepStrictEqual(matches[0].providerMetadata, source.metadata) + ) + yield* Effect.die("stale mirrored reasoning adoption") + } + } + }) + return yield* events.publishMany(entries, { validate }).pipe( + Effect.as(true), + Effect.catchCauseIf( + (cause) => !Cause.hasInterrupts(cause), + () => Effect.logWarning("reasoning adoption failed; retaining current history").pipe(Effect.as(false)), + ), + ) +}) diff --git a/packages/opencode/src/session/reasoning-distillation.ts b/packages/opencode/src/session/reasoning-distillation.ts index 141878f03c..2ea7001eaf 100644 --- a/packages/opencode/src/session/reasoning-distillation.ts +++ b/packages/opencode/src/session/reasoning-distillation.ts @@ -540,7 +540,7 @@ const runPreparedCycle = async ( } /** - * Replay every exact reasoning slot in wire order and prepare at most one new slot per request. + * Replay exact reasoning slots in wire order. Completed-turn preparation processes every new slot. * Synchronous preparation can perform a proposal followed by its independent review. * Reapply all current certificates before advancing a pending judge or a fresh proposal. A paid call never prevents * another validated slot from appearing in this request rebuilt from persisted original history. @@ -587,7 +587,7 @@ export const runDistillationCycle = async ( request = cycle.projection.request applied ||= cycle.projection.applied last = cycle - if (cycle.attempted !== "none") break + if (!input.synchronous && cycle.attempted !== "none") break } if (!last) return runSingleDistillationCycle(state, input) return { @@ -911,6 +911,7 @@ export const parseSupport = (raw: unknown): ClaimSupport[] | undefined => { * ambiguous, or multi-part lineage remains P4-protected until ordered multi-part evidence mapping is implemented. */ export type InterleavedSourcePart = Readonly<{ + distilled?: boolean messageID: string partID: string text: string @@ -962,7 +963,7 @@ export const extractInterleavedReasoningSlots = ( signed: known ? source.signed : false, encrypted: known ? source.encrypted : false, settled: known ? source.settled : false, - structureRewritable: known, + structureRewritable: known && !source.distilled, }) } return slots @@ -1015,7 +1016,7 @@ export const extractNativeInterleavedReasoningSlots = ( signed: known ? source.signed : false, encrypted: known ? source.encrypted : false, settled: known ? source.settled : false, - structureRewritable: known, + structureRewritable: known && !source.distilled, }) } return slots @@ -1195,8 +1196,7 @@ export type ScopedReasoningEvidence = Readonly<{ evidenceReferences: readonly EvidenceRef[] }> -export const isDistillationTurn = (turn: number): boolean => - Number.isSafeInteger(turn) && turn > 0 && (turn - 1) % 3 === 0 +export const isDistillationTurn = (turn: number): boolean => Number.isSafeInteger(turn) && turn > 0 export type ReasoningHistorySnapshot = Readonly<{ /** User turn owning the newest settled reasoning; tool steps do not increment it. */ @@ -1295,6 +1295,7 @@ export const reasoningHistory = (messages: readonly SessionV1.WithParts[]): Reas new Set(["encrypted_content", "encryptedContent", "reasoningEncryptedContent"]), ), settled: part.time.end !== undefined, + distilled: part.distillation !== undefined, }) } if (part.type !== "tool") continue diff --git a/packages/opencode/src/session/session.ts b/packages/opencode/src/session/session.ts index 7316550c3f..47892173d7 100644 --- a/packages/opencode/src/session/session.ts +++ b/packages/opencode/src/session/session.ts @@ -875,6 +875,7 @@ export const layer: Layer.Layer< messageID: cloned.id, sessionID: session.id, } + if (p.type === "reasoning") delete p.v2 if (p.type === "compaction" && p.tail_start_id) { p.tail_start_id = idMap.get(p.tail_start_id) } diff --git a/packages/opencode/test/config/config.test.ts b/packages/opencode/test/config/config.test.ts index 1f0f4c9ed3..b633a0526b 100644 --- a/packages/opencode/test/config/config.test.ts +++ b/packages/opencode/test/config/config.test.ts @@ -2144,3 +2144,13 @@ test("parseManagedPlist handles empty config", async () => { ) expect(config.$schema).toBe("https://opencode.ai/config.json") }) + +for (const enabled of [true, false]) { + it.instance(`reads reasoning distillation opt-in ${enabled} from opencode.json`, () => + Effect.gen(function* () { + const test = yield* TestInstance + yield* writeConfigEffect(test.directory, { reasoningDistillation: { enabled } }) + expect((yield* Config.use.get()).reasoningDistillation?.enabled).toBe(enabled) + }), + ) +} diff --git a/packages/opencode/test/session/compaction.test.ts b/packages/opencode/test/session/compaction.test.ts index cfd04d6d3f..61215bd8c7 100644 --- a/packages/opencode/test/session/compaction.test.ts +++ b/packages/opencode/test/session/compaction.test.ts @@ -328,6 +328,7 @@ function llm() { layer: Layer.succeed( LLM.Service, LLM.Service.of({ + distill: () => Effect.void, stream: (input) => { const item = queue.shift() ?? Stream.empty const stream = typeof item === "function" ? item(input) : item diff --git a/packages/opencode/test/session/llm.test.ts b/packages/opencode/test/session/llm.test.ts index 1dd074539d..4ad6dc02b6 100644 --- a/packages/opencode/test/session/llm.test.ts +++ b/packages/opencode/test/session/llm.test.ts @@ -80,15 +80,25 @@ const drainWith = (layer: Layer.Layer, input: LLM.StreamInput) => ) }) -const drainSequenceWith = (layer: Layer.Layer, inputs: readonly LLM.StreamInput[]) => +const drainSequenceWith = (layer: Layer.Layer, inputs: readonly LLM.StreamInput[], distill = false) => Effect.gen(function* () { const ctx = yield* InstanceRef if (!ctx) return yield* Effect.die("InstanceRef not provided") return yield* Effect.promise(() => Effect.runPromise( - Effect.forEach(inputs, (input) => LLM.Service.use((svc) => svc.stream(input).pipe(Stream.runDrain)), { - discard: true, - }).pipe(Effect.provide(layer), Effect.provideService(InstanceRef, ctx)), + Effect.forEach( + inputs, + (input) => + LLM.Service.use((svc) => + Effect.gen(function* () { + if (distill) yield* svc.distill(input) + yield* svc.stream(input).pipe(Stream.runDrain) + }), + ), + { + discard: true, + }, + ).pipe(Effect.provide(layer), Effect.provideService(InstanceRef, ctx)), ), ) }) @@ -1842,127 +1852,150 @@ describe("session.llm.stream", () => { ) for (const runtime of ["opencode-ai-sdk", "opencode-native"] as const) { - it.instance( - `organizes before the first ${runtime} send below threshold and replays on an idle turn`, - () => - Effect.gen(function* () { - const propose = waitRequest( - "/chat/completions", - auxiliaryResponse({ - claims: [ - { - id: "c1", - kind: "decision", - text: "已确认方案A。", - scope: "当前会话", - sources: [ - { + for (const adopted of [true, false]) { + it.instance( + `${runtime} distills independently and only replays persisted replacements; persistence accepted=${adopted}`, + () => + Effect.gen(function* () { + const propose = waitRequest( + "/chat/completions", + auxiliaryResponse({ + claims: [ + { + id: "c1", + kind: "decision", + text: "已确认方案A。", + scope: "当前会话", + sources: [ + { + messageID: "msg-reasoning-source", + partID: "prt-reasoning-source", + start: 0, + end: distillationBody.length, + }, + ], + evidence: [], + status: "unverified", + }, + ], + preserved: [], + coverage: [ + { + source: { messageID: "msg-reasoning-source", partID: "prt-reasoning-source", start: 0, end: distillationBody.length, }, - ], - evidence: [], - status: "unverified", - }, - ], - preserved: [], - coverage: [ - { - source: { - messageID: "msg-reasoning-source", - partID: "prt-reasoning-source", - start: 0, - end: distillationBody.length, + action: "keep", + claimID: "c1", }, - action: "keep", - claimID: "c1", - }, - ], - }), - ) - const judge = waitRequest( - "/chat/completions", - auxiliaryResponse({ - retention: { verdict: "supported" }, - support: [{ claimID: "c1", verdict: "supported", method: "judged" }], - }), - ) - const first = waitRequest( - "/chat/completions", - new Response(createChatStream("first"), { headers: { "Content-Type": "text/event-stream" } }), - ) - const second = waitRequest( - "/chat/completions", - new Response(createChatStream("second"), { headers: { "Content-Type": "text/event-stream" } }), - ) + ], + }), + ) + const judge = waitRequest( + "/chat/completions", + auxiliaryResponse({ + retention: { verdict: "supported" }, + support: [{ claimID: "c1", verdict: "supported", method: "judged" }], + }), + ) + const first = waitRequest( + "/chat/completions", + new Response(createChatStream("first"), { headers: { "Content-Type": "text/event-stream" } }), + ) + const second = waitRequest( + "/chat/completions", + new Response(createChatStream("second"), { headers: { "Content-Type": "text/event-stream" } }), + ) - const resolved = yield* Provider.use.getModel( - ProviderV2.ID.make("custom-provider"), - ModelV2.ID.make("deepseek-test-r1"), - ) - expect(resolved.limit).toMatchObject({ context: 65_536, output: 4_096 }) - const sessionID = SessionID.make(`session-test-distillation-${runtime}`) - const agent = { - name: "test", - mode: "primary", - options: {}, - permission: [{ permission: "*", pattern: "*", action: "allow" }], - } satisfies Agent.Info - const messages = distillationMessages() - const before = JSON.stringify(messages) - const input = (id: string): LLM.StreamInput => ({ - user: { - id: MessageID.make(id), + const resolved = yield* Provider.use.getModel( + ProviderV2.ID.make("custom-provider"), + ModelV2.ID.make("deepseek-test-r1"), + ) + expect(resolved.limit).toMatchObject({ context: 65_536, output: 4_096 }) + const sessionID = SessionID.make(`session-test-distillation-${runtime}`) + const agent = { + name: "test", + mode: "primary", + options: {}, + permission: [{ permission: "*", pattern: "*", action: "allow" }], + } satisfies Agent.Info + const messages = distillationMessages() + const before = JSON.stringify(messages) + const applied: Array = [] + const input = (id: string): LLM.StreamInput => ({ + user: { + id: MessageID.make(id), + sessionID, + role: "user", + time: { created: Date.now() }, + agent: agent.name, + model: { providerID: ProviderV2.ID.make("custom-provider"), modelID: resolved.id }, + } satisfies SessionV1.User, sessionID, - role: "user", - time: { created: Date.now() }, - agent: agent.name, - model: { providerID: ProviderV2.ID.make("custom-provider"), modelID: resolved.id }, - } satisfies SessionV1.User, - sessionID, - model: resolved, - agent, - system: [], - messages, - tools: {}, - purpose: "conversation", - reasoningDistillation: distillationHistory(), - }) - - yield* drainSequenceWith( - llmLayerWithExecutor(RequestExecutor.defaultLayer, { - experimentalNativeLlm: runtime === "opencode-native", - outputTokenMax: 4_096, - }), - [ - input("msg_user-distillation-first"), - { - ...input("msg_user-distillation-second"), - reasoningDistillation: { ...distillationHistory()!, reasoningTurn: 2 }, - }, - ], - ) + model: resolved, + agent, + system: [], + messages, + tools: {}, + purpose: "conversation", + reasoningDistillation: distillationHistory(), + adoptReasoning: (replacements) => + Effect.sync(() => { + applied.push(replacements) + if (adopted) { + const assistant = messages.find((message) => message.role === "assistant") + if (assistant && Array.isArray(assistant.content)) { + for (const part of assistant.content) + if (part.type === "reasoning") part.text = replacements[0].after + } + } + return adopted + }), + }) + + yield* drainSequenceWith( + llmLayerWithExecutor(RequestExecutor.defaultLayer, { + experimentalNativeLlm: runtime === "opencode-native", + outputTokenMax: 4_096, + }), + [ + input("msg_user-distillation-first"), + { + ...input("msg_user-distillation-second"), + reasoningDistillation: { ...distillationHistory()!, reasoningTurn: 2 }, + }, + ], + true, + ) - const [proposeCapture, firstCapture, judgeCapture, secondCapture] = yield* Effect.promise(() => - Promise.all([propose, first, judge, second]), - ) - expect(proposeCapture.body.stream).not.toBe(true) - expect(judgeCapture.body.stream).not.toBe(true) - expect(JSON.stringify(judgeCapture.body.messages)).toContain("最终发送文本") - expect(JSON.stringify(judgeCapture.body.messages)).toContain("retention") - const reasoning = (capture: Capture) => - (capture.body.messages as Array> | undefined)?.find( - (message) => message.role === "assistant", - )?.reasoning_content - expect(reasoning(firstCapture)).toContain("已确认方案A。") - expect(reasoning(secondCapture)).toContain("已确认方案A。") - expect(reasoning(secondCapture)).not.toBe(distillationBody) - expect(JSON.stringify(messages)).toBe(before) - }), - { config: () => distillationConfig(runtime) }, - ) + const [proposeCapture, firstCapture, judgeCapture, secondCapture] = yield* Effect.promise(() => + Promise.all([propose, first, judge, second]), + ) + expect(proposeCapture.body.stream).not.toBe(true) + expect(judgeCapture.body.stream).not.toBe(true) + expect(JSON.stringify(judgeCapture.body.messages)).toContain("最终发送文本") + expect(JSON.stringify(judgeCapture.body.messages)).toContain("retention") + const reasoning = (capture: Capture) => + (capture.body.messages as Array> | undefined)?.find( + (message) => message.role === "assistant", + )?.reasoning_content + expect(applied).toHaveLength(adopted ? 1 : 2) + expect(applied[0][0].before).toBe(distillationBody) + expect(applied[0][0].after).toContain("已确认方案A。") + if (adopted) { + expect(reasoning(firstCapture)).toBe(applied[0][0].after) + expect(reasoning(secondCapture)).toBe(applied[0][0].after) + } else { + expect(reasoning(firstCapture)).toBe(distillationBody) + expect(reasoning(secondCapture)).toBe(distillationBody) + } + if (!adopted) expect(JSON.stringify(messages)).toBe(before) + }), + { config: () => distillationConfig(runtime) }, + ) + } } it.instance( diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index fb23f8cff1..7b94b27d9f 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -199,6 +199,7 @@ const native = testEffect( const providerErrorLLM = Layer.succeed( LLM.Service, LLM.Service.of({ + distill: () => Effect.void, stream: () => Stream.make( LLMEvent.stepStart({ index: 0 }), @@ -224,6 +225,7 @@ const itProviderError = testEffect(providerErrorEnv) const fragmentFailureLLM = Layer.succeed( LLM.Service, LLM.Service.of({ + distill: () => Effect.void, stream: () => Stream.make( LLMEvent.stepStart({ index: 0 }), @@ -243,10 +245,9 @@ const itFragmentFailure = testEffect(fragmentFailureEnv) const classifiedProviderErrorLLM = Layer.succeed( LLM.Service, LLM.Service.of({ + distill: () => Effect.void, stream: () => - Stream.make( - LLMEvent.providerError({ message: "request entity too large", classification: "context-overflow" }), - ), + Stream.make(LLMEvent.providerError({ message: "request entity too large", classification: "context-overflow" })), }), ) const classifiedProviderErrorEnv = LayerNode.buildLayer( diff --git a/packages/opencode/test/session/prompt.test.ts b/packages/opencode/test/session/prompt.test.ts index 3bf4686bdb..c79ccc2cc8 100644 --- a/packages/opencode/test/session/prompt.test.ts +++ b/packages/opencode/test/session/prompt.test.ts @@ -264,6 +264,7 @@ const snapshotFenceAgentLayer: Layer.Layer = Layer.succeed( ) type PromptLayerOptions = { + distill?: LLM.Interface["distill"] mcpInstructions?: MCP.ServerInstructions[] mcpTools?: Record processor?: "blocking" @@ -295,7 +296,7 @@ function makePrompt(input?: PromptLayerOptions) { statusReason: () => Effect.succeed(undefined), status: () => Effect.succeed("Memory on"), }) - const llmLayer = input?.native + const baseLLM = input?.native ? LLM.layer.pipe( Layer.provide(Auth.defaultLayer), Layer.provide(Config.defaultLayer), @@ -307,6 +308,15 @@ function makePrompt(input?: PromptLayerOptions) { Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true, experimentalNativeLlm: true })), ) : LLM.defaultLayer + const llmLayer = input?.distill + ? Layer.effect( + LLM.Service, + Effect.gen(function* () { + const service = yield* LLM.Service + return LLM.Service.of({ ...service, distill: input.distill! }) + }), + ).pipe(Layer.provide(baseLLM)) + : baseLLM const deps = Layer.mergeAll( hookRecorderLayer, memoryLayer, @@ -5042,3 +5052,148 @@ for (const dynamic of [false, true]) { }), ) } + +let controlledDistill: LLM.Interface["distill"] = () => Effect.die("distillation fixture not configured") +const distillationIt = testEffect(makeHttp({ distill: (input) => Effect.suspend(() => controlledDistill(input)) })) + +distillationIt.instance( + "first full turn waits for adoption; later turns return before background replacement and replay adopted history", + () => + Effect.gen(function* () { + const { llm } = yield* useServerConfig((url) => ({ + ...providerCfg(url), + reasoningDistillation: { enabled: true }, + })) + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + const chat = yield* sessions.create({ + title: "Pinned", + permission: [{ permission: "*", pattern: "*", action: "allow" }], + }) + const starts = [yield* Deferred.make(), yield* Deferred.make()] + const releases = [yield* Deferred.make(), yield* Deferred.make()] + const finishes = [yield* Deferred.make(), yield* Deferred.make()] + let calls = 0 + controlledDistill = (input) => + Effect.gen(function* () { + const index = calls++ + if (index > 1) return + const slots = input + .reasoningDistillation!.groups.flatMap((group) => group.parts) + .filter((part) => !part.distilled) + expect(slots).toHaveLength(1) + expect(slots[0].text).toBe(`original-${index}`) + yield* Deferred.succeed(starts[index], undefined) + yield* Deferred.await(releases[index]) + expect( + yield* input.adoptReasoning!( + slots.map((slot) => ({ + messageID: slot.messageID, + partID: slot.partID, + before: slot.text, + after: `distilled-${index}`, + })), + ), + ).toBe(true) + yield* Deferred.succeed(finishes[index], undefined) + }) + const send = (text: string) => + prompt.prompt({ sessionID: chat.id, agent: "build", noReply: true, parts: [{ type: "text", text }] }) + const thinking = (messages: SessionV1.WithParts[]) => + messages + .flatMap((message) => message.parts) + .filter((part) => part.type === "reasoning") + .map((part) => part.text) + yield* send("first") + yield* llm.push(reply().reason("original-0").text("first answer").stop()) + const first = yield* prompt.loop({ sessionID: chat.id }).pipe(Effect.forkChild) + yield* awaitWithTimeout(Deferred.await(starts[0]), "first distillation did not start") + expect(first.pollUnsafe()).toBeUndefined() + expect(thinking(yield* sessions.messages({ sessionID: chat.id }))).toEqual(["original-0"]) + yield* Deferred.succeed(releases[0], undefined) + const completed = yield* Fiber.join(first) + expect(completed.parts.some((part) => part.type === "reasoning" && part.text === "distilled-0")).toBe(true) + + yield* send("second") + yield* llm.push(reply().reason("original-1").text("second answer").stop()) + yield* prompt.loop({ sessionID: chat.id }) + yield* awaitWithTimeout(Deferred.await(starts[1]), "background distillation did not start") + expect(thinking(yield* sessions.messages({ sessionID: chat.id }))).toEqual(["distilled-0", "original-1"]) + yield* send("third while distillation is pending") + yield* llm.text("third answer") + yield* prompt.loop({ sessionID: chat.id }) + const pendingPayload = JSON.stringify((yield* llm.hits).at(-1)?.body) + expect(pendingPayload).toContain("original-1") + expect(pendingPayload).not.toContain("distilled-1") + yield* Deferred.succeed(releases[1], undefined) + yield* awaitWithTimeout(Deferred.await(finishes[1]), "background distillation did not finish") + expect(thinking(yield* sessions.messages({ sessionID: chat.id }))).toEqual(["distilled-0", "distilled-1"]) + + yield* send("fourth after replacement") + yield* llm.text("fourth answer") + yield* prompt.loop({ sessionID: chat.id }) + const payload = JSON.stringify((yield* llm.hits).at(-1)?.body) + expect(payload).toContain("distilled-0") + expect(payload).toContain("distilled-1") + expect(payload).not.toContain("original-0") + expect(payload).not.toContain("original-1") + }), +) + +distillationIt.instance("tool steps form one distillation turn and include all newly completed reasoning slots", () => + Effect.gen(function* () { + const { llm, dir } = yield* useServerConfig((url) => ({ + ...providerCfg(url), + reasoningDistillation: { enabled: true }, + })) + yield* writeText(path.join(dir, "distillation-source.txt"), "source data") + const prompt = yield* SessionPrompt.Service + const sessions = yield* Session.Service + const chat = yield* sessions.create({ + title: "Pinned", + permission: [{ permission: "*", pattern: "*", action: "allow" }], + }) + let calls = 0 + controlledDistill = (input) => + Effect.gen(function* () { + calls++ + const slots = input + .reasoningDistillation!.groups.flatMap((group) => group.parts) + .filter((part) => !part.distilled) + expect(slots.map((slot) => slot.text)).toEqual(["tool thinking", "final thinking"]) + expect( + yield* input.adoptReasoning!( + slots.map((slot) => ({ + messageID: slot.messageID, + partID: slot.partID, + before: slot.text, + after: `adopted ${slot.text}`, + })), + ), + ).toBe(true) + }) + yield* prompt.prompt({ + sessionID: chat.id, + agent: "build", + noReply: true, + parts: [{ type: "text", text: "read and explain" }], + }) + yield* llm.push( + reply() + .reason("tool thinking") + .tool("read", { filePath: path.join(dir, "distillation-source.txt") }, "read-source"), + reply().reason("final thinking").text("answer").stop(), + ) + yield* prompt.loop({ sessionID: chat.id }) + expect(calls).toBe(1) + const parts = (yield* sessions.messages({ sessionID: chat.id })).flatMap((message) => message.parts) + expect(parts.filter((part) => part.type === "reasoning").map((part) => part.text)).toEqual([ + "adopted tool thinking", + "adopted final thinking", + ]) + expect(parts.find((part) => part.type === "tool")).toMatchObject({ + callID: "read-source", + state: { status: "completed" }, + }) + }), +) diff --git a/packages/opencode/test/session/reasoning-distillation.test.ts b/packages/opencode/test/session/reasoning-distillation.test.ts index fa3d907a99..bd5329bc58 100644 --- a/packages/opencode/test/session/reasoning-distillation.test.ts +++ b/packages/opencode/test/session/reasoning-distillation.test.ts @@ -137,7 +137,7 @@ describe("buildReasoningEvidence", () => { }) describe("buildSlotMappings (§2.1 gating)", () => { - test("without a compatibility record every slot is P5-protected (default-on still a no-op)", () => { + test("without a compatibility record every slot is P5-protected (explicit opt-in still a no-op)", () => { const mappings = buildSlotMappings([slot()], capability(), []) expect(mappings[0].eligibility).toEqual({ allowed: false, protection: "P5" }) expect(mappings[0].sourceFingerprint).toBe(Hash.sha256("原始冗长思绪")) @@ -1479,8 +1479,10 @@ describe("review regressions: conservation and source binding", () => { }) }) -test("reasoning preparation cadence is first turn then every three user turns", () => { - expect(Array.from({ length: 11 }, (_, turn) => turn).filter(isDistillationTurn)).toEqual([1, 4, 7, 10]) +test("reasoning preparation cadence includes every completed user turn", () => { + expect(Array.from({ length: 11 }, (_, turn) => turn).filter(isDistillationTurn)).toEqual([ + 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, + ]) expect(isDistillationTurn(NaN)).toBe(false) expect(isDistillationTurn(1.5)).toBe(false) }) diff --git a/packages/opencode/test/session/session.test.ts b/packages/opencode/test/session/session.test.ts index 4323703b07..546ba13235 100644 --- a/packages/opencode/test/session/session.test.ts +++ b/packages/opencode/test/session/session.test.ts @@ -1,9 +1,23 @@ +import { SessionEvent } from "@opencode-ai/core/session/event" +import { SessionMessage } from "@opencode-ai/core/session/message" +import { SessionMessageTable } from "@opencode-ai/core/session/sql" +import { eq } from "drizzle-orm" +import { adoptReasoning } from "@/session/reasoning-adoption" +import { + reasoningHistory, + bindInterleavedReasoningLineage, + extractInterleavedReasoningSlots, +} from "@/session/reasoning-distillation" +import { ProviderTest } from "../fake/provider" +import { ProviderTransform } from "@/provider/transform" +import { ProviderV2 } from "@opencode-ai/core/provider" +import { ModelV2 } from "@opencode-ai/core/model" import { describe, expect } from "bun:test" import { SessionV1 } from "@opencode-ai/core/v1/session" import { Database } from "@opencode-ai/core/database/database" import { EventV2 } from "@opencode-ai/core/event" import { SessionProjector } from "@opencode-ai/core/session/projector" -import { Deferred, Effect, Exit, Layer } from "effect" +import { Schema, DateTime, Deferred, Effect, Exit, Layer } from "effect" import { Session as SessionNs } from "@/session/session" import { Goal } from "@/goal/goal" import { SessionAutomationLease } from "@/session/automation-lease" @@ -32,6 +46,7 @@ const it = testEffect( Layer.provide(SessionAutomationLease.defaultLayer), Layer.provide(Dag.defaultLayer), ), + Database.defaultLayer, CrossSpawnSpawner.defaultLayer, testInstanceStoreLayer, ), @@ -252,3 +267,165 @@ describe("Session", () => { }), ) }) + +describe("adopted reasoning", () => { + const seed = Effect.fnUntraced(function* () { + const session = yield* SessionNs.Service + const chat = yield* session.create({}) + const info: SessionV1.Assistant = { + id: MessageID.ascending(), + sessionID: chat.id, + parentID: MessageID.ascending(), + role: "assistant", + agent: "build", + modelID: ModelV2.ID.make("model"), + providerID: ProviderV2.ID.make("provider"), + mode: "build", + path: { cwd: chat.directory, root: chat.directory }, + time: { created: 1, completed: 2 }, + cost: 0, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + } + yield* session.updateMessage(info) + const reasoning: SessionV1.ReasoningPart = { + id: PartID.ascending(), + sessionID: chat.id, + messageID: info.id, + type: "reasoning", + text: "original private reasoning", + time: { start: 1, end: 2 }, + } + const events = yield* EventV2Bridge.Service + const assistantMessageID = SessionMessage.ID.create() + yield* events.publish(SessionEvent.Step.Started, { + sessionID: chat.id, + assistantMessageID, + agent: "build", + model: { id: info.modelID, providerID: info.providerID }, + timestamp: DateTime.makeUnsafe(1), + }) + yield* events.publish(SessionEvent.Reasoning.Started, { + sessionID: chat.id, + assistantMessageID, + reasoningID: "r1", + timestamp: DateTime.makeUnsafe(1), + }) + yield* events.publish(SessionEvent.Reasoning.Ended, { + sessionID: chat.id, + assistantMessageID, + reasoningID: "r1", + text: reasoning.text, + timestamp: DateTime.makeUnsafe(2), + }) + reasoning.v2 = { messageID: assistantMessageID, reasoningID: "r1" } + yield* session.updatePart(reasoning) + const sources: SessionV1.WithParts[] = [{ info, parts: [reasoning] }] + const replacements = [{ messageID: info.id, partID: reasoning.id, before: reasoning.text, after: "采用后的思考" }] + return { session, chat, info, reasoning, sources, replacements } + }) + + it.instance("updates subscribers, reloaded history and model context to the same adopted text", () => + Effect.gen(function* () { + const { session, chat, info, reasoning, sources, replacements } = yield* seed() + const events = yield* EventV2Bridge.Service + const received = yield* Deferred.make() + const unsub = yield* events.listen((event) => { + if (event.type === SessionV1.Event.PartUpdated.type) { + const part = Schema.decodeUnknownSync(SessionV1.Event.PartUpdated.data)(event.data).part + if (part.type === "reasoning" && part.distillation) + Deferred.doneUnsafe(received, Effect.succeed(part as SessionV1.Part)) + } + return Effect.void + }) + yield* Effect.addFinalizer(() => unsub) + expect(yield* adoptReasoning({ sessionID: chat.id, sources, replacements })).toBe(true) + expect(yield* awaitDeferred(received, "adoption event missing")).toMatchObject({ + id: reasoning.id, + text: replacements[0].after, + }) + const reloaded = yield* session.messages({ sessionID: chat.id }) + const part = reloaded + .find((message) => message.info.id === info.id) + ?.parts.find((part) => part.id === reasoning.id) + expect(part).toMatchObject({ text: replacements[0].after, distillation: { originalText: reasoning.text } }) + const { db } = yield* Database.Service + const mirror = yield* db + .select() + .from(SessionMessageTable) + .where(eq(SessionMessageTable.id, SessionMessage.ID.make(reasoning.v2!.messageID))) + .get() + .pipe(Effect.orDie) + expect(JSON.stringify(mirror?.data)).toContain(replacements[0].after) + const model = ProviderTest.model({ id: info.modelID, providerID: info.providerID }) + const plain = yield* MessageV2.toModelMessagesEffect(reloaded, model) + expect(JSON.stringify(plain)).toContain(replacements[0].after) + expect(JSON.stringify(plain)).not.toContain(reasoning.text) + expect(JSON.stringify(plain)).not.toContain("sourceFingerprint") + const interleaved = { + ...model, + capabilities: { ...model.capabilities, interleaved: { field: "reasoning_content" as const } }, + } + const transformed = ProviderTransform.message(structuredClone(plain), interleaved, {}) + const snapshot = reasoningHistory(reloaded) + const lineage = bindInterleavedReasoningLineage(plain, transformed, "reasoning_content", snapshot) + const slots = extractInterleavedReasoningSlots(transformed, "reasoning_content", ["messages"], lineage) + expect(slots).toHaveLength(1) + expect(slots[0]).toMatchObject({ text: replacements[0].after, structureRewritable: false }) + expect(yield* adoptReasoning({ sessionID: chat.id, sources, replacements })).toBe(false) + yield* session.remove(chat.id) + }), + ) + + it.instance("rejects stale or streaming sources without replacing history", () => + Effect.gen(function* () { + const { session, chat, reasoning, sources, replacements } = yield* seed() + const changed = { ...reasoning, text: "a newer source" } + yield* session.updatePart(changed) + expect(yield* adoptReasoning({ sessionID: chat.id, sources, replacements })).toBe(false) + expect( + yield* session.getPart({ sessionID: chat.id, messageID: reasoning.messageID, partID: reasoning.id }), + ).toMatchObject({ text: changed.text }) + const pending = { ...reasoning, time: { start: 1 } } + yield* session.updatePart(pending) + expect( + yield* adoptReasoning({ sessionID: chat.id, sources: [{ ...sources[0], parts: [pending] }], replacements }), + ).toBe(false) + yield* session.remove(chat.id) + }), + ) + + it.instance("rolls back every slot if a later source changed", () => + Effect.gen(function* () { + const { session, chat, info, reasoning, sources, replacements } = yield* seed() + const second = { ...reasoning, id: PartID.ascending(), text: "second source" } + yield* session.updatePart(second) + sources[0].parts.push(second) + yield* session.updatePart({ ...second, metadata: { provider: { signature: "new-signature" } } }) + expect( + yield* adoptReasoning({ + sessionID: chat.id, + sources, + replacements: [ + ...replacements, + { messageID: info.id, partID: second.id, before: second.text, after: "second replacement" }, + ], + }), + ).toBe(false) + expect(yield* session.getPart({ sessionID: chat.id, messageID: info.id, partID: reasoning.id })).toEqual( + reasoning, + ) + yield* session.remove(chat.id) + }), + ) + it.instance("a disabled switch blocks an already prepared adoption", () => + Effect.gen(function* () { + const { session, chat, reasoning, sources, replacements } = yield* seed() + expect( + yield* adoptReasoning({ sessionID: chat.id, sources, replacements, canAdopt: Effect.succeed(false) }), + ).toBe(false) + expect( + yield* session.getPart({ sessionID: chat.id, messageID: reasoning.messageID, partID: reasoning.id }), + ).toEqual(reasoning) + }), + ) +}) diff --git a/packages/schema/src/reasoning-distillation.ts b/packages/schema/src/reasoning-distillation.ts new file mode 100644 index 0000000000..28966e9e95 --- /dev/null +++ b/packages/schema/src/reasoning-distillation.ts @@ -0,0 +1,8 @@ +import { Schema } from "effect" + +/** Host-owned adoption provenance, separate from provider metadata and model context. */ +export const ReasoningDistillation = Schema.Struct({ + originalText: Schema.String, + sourceFingerprint: Schema.String, +}).annotate({ identifier: "ReasoningDistillation" }) +export type ReasoningDistillation = typeof ReasoningDistillation.Type diff --git a/packages/schema/src/session-event.ts b/packages/schema/src/session-event.ts index 2e71667e37..cc1fabc6ff 100644 --- a/packages/schema/src/session-event.ts +++ b/packages/schema/src/session-event.ts @@ -1,3 +1,4 @@ +import { ReasoningDistillation } from "./reasoning-distillation" export * as SessionEvent from "./session-event" import { Schema } from "effect" @@ -264,6 +265,7 @@ export namespace Reasoning { assistantMessageID: SessionMessageID.ID, reasoningID: Schema.String, text: Schema.String, + distillation: ReasoningDistillation.pipe(Schema.optional), providerMetadata: ProviderMetadata.pipe(Schema.optional), }, }) diff --git a/packages/schema/src/session-message.ts b/packages/schema/src/session-message.ts index 63cc25b43d..f79e5bc7dd 100644 --- a/packages/schema/src/session-message.ts +++ b/packages/schema/src/session-message.ts @@ -1,3 +1,4 @@ +import { ReasoningDistillation } from "./reasoning-distillation" export * as SessionMessage from "./session-message" import { Schema } from "effect" @@ -145,6 +146,7 @@ export const AssistantReasoning = Schema.Struct({ type: Schema.Literal("reasoning"), id: Schema.String, text: Schema.String, + distillation: ReasoningDistillation.pipe(Schema.optional), providerMetadata: ProviderMetadata.pipe(Schema.optional), }).annotate({ identifier: "Session.Message.Assistant.Reasoning" }) diff --git a/packages/schema/src/session-v1.ts b/packages/schema/src/session-v1.ts index 0af94d73ab..82d21229e6 100644 --- a/packages/schema/src/session-v1.ts +++ b/packages/schema/src/session-v1.ts @@ -1,3 +1,4 @@ +import { ReasoningDistillation } from "./reasoning-distillation" export * as SessionV1 from "./session-v1" import { Effect, Schema, Types } from "effect" @@ -119,6 +120,8 @@ export const ReasoningPart = Schema.Struct({ ...partBase, type: Schema.Literal("reasoning"), text: Schema.String, + distillation: ReasoningDistillation.pipe(Schema.optional), + v2: Schema.Struct({ messageID: Schema.String, reasoningID: Schema.String }).pipe(Schema.optional), metadata: Schema.optional(Schema.Record(Schema.String, Schema.Any)), time: Schema.Struct({ start: NonNegativeInt, diff --git a/packages/sdk/js/src/v2/gen/types.gen.ts b/packages/sdk/js/src/v2/gen/types.gen.ts index 58b90147ac..c5798d612c 100644 --- a/packages/sdk/js/src/v2/gen/types.gen.ts +++ b/packages/sdk/js/src/v2/gen/types.gen.ts @@ -416,12 +416,22 @@ export type SubtaskPart = { command?: string } +export type ReasoningDistillation = { + originalText: string + sourceFingerprint: string +} + export type ReasoningPart = { id: string sessionID: string messageID: string type: "reasoning" text: string + distillation?: ReasoningDistillation + v2?: { + messageID: string + reasoningID: string + } metadata?: { [key: string]: unknown } @@ -1076,6 +1086,7 @@ export type GlobalEvent = { assistantMessageID: string reasoningID: string text: string + distillation?: ReasoningDistillation providerMetadata?: { [key: string]: { [key: string]: unknown @@ -3628,6 +3639,7 @@ export type SyncEventSessionNextReasoningEnded = { assistantMessageID: string reasoningID: string text: string + distillation?: ReasoningDistillation providerMetadata?: { [key: string]: { [key: string]: unknown @@ -4154,6 +4166,7 @@ export type SessionMessageAssistantReasoning = { type: "reasoning" id: string text: string + distillation?: ReasoningDistillation providerMetadata?: { [key: string]: { [key: string]: unknown @@ -5118,6 +5131,7 @@ export type V2EventSessionNextReasoningEnded = { assistantMessageID: string reasoningID: string text: string + distillation?: ReasoningDistillation providerMetadata?: { [key: string]: { [key: string]: unknown @@ -6832,6 +6846,7 @@ export type EventSessionNextReasoningEnded = { assistantMessageID: string reasoningID: string text: string + distillation?: ReasoningDistillation providerMetadata?: { [key: string]: { [key: string]: unknown diff --git a/packages/tui/test/cli/cmd/tui/sync-reasoning-adoption.test.tsx b/packages/tui/test/cli/cmd/tui/sync-reasoning-adoption.test.tsx new file mode 100644 index 0000000000..634a289ea7 --- /dev/null +++ b/packages/tui/test/cli/cmd/tui/sync-reasoning-adoption.test.tsx @@ -0,0 +1,77 @@ +/** @jsxImportSource @opentui/solid */ +import { expect, test } from "bun:test" +import type { AssistantMessage, ReasoningPart } from "@opencode-ai/sdk/v2" +import { tmpdir } from "../../../fixture/fixture" +import { directory, json, mount, wait } from "./sync-fixture" + +test("adopted thinking replaces the existing displayed part and survives a fresh session sync", async () => { + await using tmp = await tmpdir() + await Bun.write(`${tmp.path}/kv.json`, "{}") + const sessionID = "ses_reasoning_adoption" + const messageID = "msg_reasoning_adoption" + const partID = "prt_reasoning_adoption" + const session = { + id: sessionID, + title: "reasoning", + slug: "reasoning", + projectID: "proj_test", + time: { created: 1, updated: 1 }, + version: "test", + directory, + } + const info: AssistantMessage = { + id: messageID, + sessionID, + parentID: "msg_user", + role: "assistant", + agent: "build", + mode: "build", + modelID: "model", + providerID: "provider", + path: { cwd: directory, root: directory }, + time: { created: 1, completed: 2 }, + cost: 0, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + } + let part: ReasoningPart = { + id: partID, + messageID, + sessionID, + type: "reasoning", + text: "original thinking", + time: { start: 1, end: 2 }, + } + const serve = (url: URL) => { + if (url.pathname === "/session") return json([session]) + if (url.pathname === `/session/${sessionID}`) return json(session) + if (url.pathname === `/session/${sessionID}/message`) return json([{ info, parts: [part] }]) + if (url.pathname === `/session/${sessionID}/todo` || url.pathname === `/session/${sessionID}/diff`) return json([]) + return undefined + } + const first = await mount(serve, tmp.path) + try { + await first.sync.session.sync(sessionID) + expect(first.sync.data.part[messageID][0].type).toBe("reasoning") + expect(first.sync.data.part[messageID][0]).toMatchObject({ text: "original thinking" }) + part = { ...part, text: "蒸馏后的思考", distillation: { originalText: part.text, sourceFingerprint: "source" } } + first.emit({ + directory, + payload: { id: "evt_adoption", type: "message.part.updated", properties: { sessionID, part, time: 3 } }, + }) + await wait(() => + first.sync.data.part[messageID].some((item) => item.type === "reasoning" && item.text === part.text), + ) + expect(first.sync.data.part[messageID]).toHaveLength(1) + expect(first.sync.data.part[messageID][0]).toMatchObject({ id: partID, text: part.text }) + } finally { + first.app.renderer.destroy() + } + const restored = await mount(serve, tmp.path) + try { + await restored.sync.session.sync(sessionID) + expect(restored.sync.data.part[messageID]).toHaveLength(1) + expect(restored.sync.data.part[messageID][0]).toMatchObject({ id: partID, text: "蒸馏后的思考" }) + } finally { + restored.app.renderer.destroy() + } +}) From bee301af534b1b0f823ab66749534c75ebd2ed60 Mon Sep 17 00:00:00 2001 From: Lex Date: Mon, 28 Sep 2026 08:56:02 +0800 Subject: [PATCH 2/2] fix: skip disabled distillation history work and retain diagnostics --- packages/core/src/session/runner/llm.ts | 35 ++++++++++++++++--------- 1 file changed, 22 insertions(+), 13 deletions(-) diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index 58a6d4045e..f29734ae44 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -559,22 +559,22 @@ export const layer = Layer.effect( yield* withPublication(batch.flush()) if (stream._tag === "Failure") return yield* Effect.failCause(stream.cause) if (settled._tag === "Failure") return yield* Effect.failCause(settled.cause) - if (!publisher.hasProviderError() && !needsContinuation) { + const enabled = config.entries().pipe( + Effect.map( + (entries) => + ConfigReasoningDistillation.resolveEnabled({ + disabledByEnvironment: Flag.OPENCODE_DISABLE_REASONING_DISTILLATION, + enabled: Config.latest(entries, "reasoningDistillation")?.enabled, + }).enabled, + ), + ) + if (!publisher.hasProviderError() && !needsContinuation && (yield* enabled)) { const completed = yield* getContext(session.id) const userIndex = completed.findLastIndex((message) => message.type === "user") const turnID = completed[userIndex]?.id if (turnID) { const turnIDs = new Set(completed.slice(userIndex + 1).map((message) => message.id)) const latestConversion = toLLMMessagesWithBindings(completed, model) - const enabled = config.entries().pipe( - Effect.map( - (entries) => - ConfigReasoningDistillation.resolveEnabled({ - disabledByEnvironment: Flag.OPENCODE_DISABLE_REASONING_DISTILLATION, - enabled: Config.latest(entries, "reasoningDistillation")?.enabled, - }).enabled, - ), - ) yield* restore( scheduleDistillation({ sessionID: session.id, @@ -603,8 +603,10 @@ export const layer = Layer.effect( ), config: reasoningConfig, }) - if (distilled.applied && (yield* enabled)) - yield* adoptReasoning( + const adopted = + distilled.applied && + (yield* enabled) && + (yield* adoptReasoning( events, db, session.id, @@ -617,7 +619,14 @@ export const layer = Layer.effect( if (adopted) cursors.delete(session.id) }), ), - ) + )) + yield* Effect.logInfo("reasoning distillation", { + "reasoning_distillation.runtime": "core-runner", + "reasoning_distillation.applied": adopted, + "reasoning_distillation.attempted": distilled.attempted, + "reasoning_distillation.skip_reason": distilled.skipReason ?? "none", + "reasoning_distillation.aux_actual_tokens": distilled.usage?.actualTokens ?? 0, + }) }).pipe(Effect.asVoid), }), )