From 51812bf61d293b8e671ed5c5fe2a5a00956ce4bc Mon Sep 17 00:00:00 2001 From: grishberg Date: Fri, 17 Jul 2026 15:54:45 +0300 Subject: [PATCH 1/9] feat(bot): nudge model after 10min reasoning silence --- bot/reasoning_timeout.py | 73 ++++++++++++++++++++++++++++++++++++++++ bot/vk_longpoll.py | 47 ++++++++++++++++++++++++++ 2 files changed, 120 insertions(+) create mode 100644 bot/reasoning_timeout.py diff --git a/bot/reasoning_timeout.py b/bot/reasoning_timeout.py new file mode 100644 index 000000000000..4140c96ee60f --- /dev/null +++ b/bot/reasoning_timeout.py @@ -0,0 +1,73 @@ +""" +Отслеживание зависания модели после reasoning. + +Если reasoning.ended сработал, но за ним SILENCE_TIMEOUT минут не пришло +ни text.ended, ни step.ended, ни tool.called — отправляет промпт-напоминание. +""" +import asyncio +from typing import Optional, Callable, Awaitable + +from logging_config import logger + + +SILENCE_TIMEOUT = 600 # 10 минут +NUDGE_PROMPT = "На чем остановился?" + + +class ReasoningTimeout: + """Watchdog: если после reasoning тишина > SILENCE_TIMEOUT — нуджить модель.""" + + def __init__(self, send_prompt: Callable[[str, str], Awaitable[bool]]): + """ + Args: + send_prompt: функция (session_id, text) -> bool для отправки промпта. + """ + self.send_prompt = send_prompt + + # session_id -> состояние + self._active: dict[str, dict] = {} + self._timers: dict[str, asyncio.Task] = {} + + def reasoning_ended(self, session_id: str, user_id: int): + """reasoning.ended — запускаем таймер.""" + if session_id in self._active: + return # уже есть активный таймер + + self._active[session_id] = {"user_id": user_id} + self._start_timer(session_id) + + def activity(self, session_id: str): + """Любая активность (text.ended, step.ended, tool.called) — сбрасываем таймер.""" + self._cancel_timer(session_id) + self._active.pop(session_id, None) + + def _start_timer(self, session_id: str): + self._cancel_timer(session_id) + task = asyncio.create_task(self._wait_and_nudge(session_id)) + task.set_name(f"reasoning_timeout:{session_id[:8]}") + self._timers[session_id] = task + + def _cancel_timer(self, session_id: str): + task = self._timers.pop(session_id, None) + if task and not task.done(): + task.cancel() + + async def _wait_and_nudge(self, session_id: str): + try: + await asyncio.sleep(SILENCE_TIMEOUT) + except asyncio.CancelledError: + return + + state = self._active.pop(session_id, None) + self._timers.pop(session_id, None) + if not state: + return + + user_id = state["user_id"] + logger.warning( + f"Reasoning timeout for session {session_id} (user {user_id}), sending nudge" + ) + try: + await self.send_prompt(session_id, NUDGE_PROMPT) + except Exception as e: + logger.exception(f"Failed to send nudge prompt: {e}") diff --git a/bot/vk_longpoll.py b/bot/vk_longpoll.py index b3bd38fc2df6..a304a79aa0e8 100644 --- a/bot/vk_longpoll.py +++ b/bot/vk_longpoll.py @@ -36,6 +36,7 @@ from opencode_client import OpenCodeClient from opencode_process import OpenCodeProcess from session_manager import SessionManager +from reasoning_timeout import ReasoningTimeout from sse_listener import SSEEventListener as SSEListener from vk_client import VKClient @@ -127,6 +128,9 @@ def __init__( # Child session (subagent/subtask) tracking self.parent_child_map: Dict[str, Dict[str, dict]] = {} # parent_id -> {child_id: {title, ...}} + # Watchdog: нуджить модель если после reasoning тишина > 10 мин + self.reasoning_timeout: Optional[ReasoningTimeout] = None + # Временные хранилища self.waiting_for_answer: Dict = {} # user_id -> question_id или (peer_id, child_id) -> question_id self.pending_permissions: Dict[str, Tuple[str, int, int]] = {} @@ -165,6 +169,8 @@ async def _on_text_ended(self, event_type: str, data: dict): text = data.get("text", "") if not text.strip(): return + if self.reasoning_timeout: + self.reasoning_timeout.activity(session_id) self._pending_texts.setdefault(session_id, []).append(text) async def _on_reasoning_ended(self, event_type: str, data: dict): @@ -178,6 +184,8 @@ async def _on_reasoning_ended(self, event_type: str, data: dict): return target = THINKING_PEER_ID if THINKING_PEER_ID else user_id await self.vk.send_message(target, f"🧠:\n{text}") + if self.reasoning_timeout: + self.reasoning_timeout.reasoning_ended(session_id, user_id) async def _on_tool_event(self, event_type: str, data: dict): """Обрабатывает события инструментов с детальным описанием""" @@ -186,6 +194,9 @@ async def _on_tool_event(self, event_type: str, data: dict): if not user_id: return + if self.reasoning_timeout: + self.reasoning_timeout.activity(session_id) + tool_name = data.get("tool", "") or data.get("name", "") target = THINKING_PEER_ID if THINKING_PEER_ID else user_id action = event_type.split(".")[-1] if "." in event_type else "event" @@ -227,6 +238,9 @@ async def _on_step_event(self, event_type: str, data: dict): if not user_id: return + if self.reasoning_timeout: + self.reasoning_timeout.activity(session_id) + step_type = event_type.split(".")[-1] if "." in event_type else "step" target = THINKING_PEER_ID if THINKING_PEER_ID else user_id @@ -768,6 +782,10 @@ async def _handle_message_new(self, event: list): await self._handle_sessions_command(user_id) return + if cmd == "/compact": + await self._handle_compact_command(user_id) + return + if cmd.startswith("/logs"): await self._handle_logs_command(user_id) return @@ -1150,6 +1168,32 @@ async def _handle_sessions_command(self, user_id: int): sessions_text += f"• `{sid}` (user={uid}) {marker}\n" await self.vk.send_message(user_id, f"📋 **Список сессий**:\n\n{sessions_text}") + async def _handle_compact_command(self, user_id: int): + """Обрабатывает команду /compact — принудительное сжатие контекста текущей сессии.""" + session_id = self.session_mgr.sessions.get(user_id) + if not session_id: + await self.vk.send_message(user_id, "❌ Нет активной сессии") + return + + # provider/model берём из текущей модели бота + model_alias = bot_config.DEFAULT_MODEL + model_info = bot_config.MODELS.get(model_alias, {}) + model_id = model_info.get("model", model_alias) + provider_id = model_info.get("provider", "cli") + + await self.vk.send_message( + user_id, + f"🔄 Сжимаю контекст сессии {session_id}...\nМожет занять до минуты.", + ) + ok = await self.opencode_client.summarize_session(session_id, provider_id, model_id) + if ok: + await self.vk.send_message(user_id, "✅ Контекст сжат") + else: + await self.vk.send_message( + user_id, + "❌ Не удалось сжать контекст (см. логи). Возможно модель не поддерживает сжатие.", + ) + async def _handle_logs_command(self, user_id: int): """Обрабатывает команду /logs""" await self.vk.send_file( @@ -1491,6 +1535,7 @@ async def _send_help(self, user_id: int): /sysmon - Показать статус системы (GPU, RAM, процессы) /logs - Отправить файл логов /sessions - Показать список всех сессий +/compact - Принудительно сжать (суммаризировать) контекст текущей сессии /newsession [path] - Создать новую сессию (очищает старые) /n [path] - То же что /newsession /newsession /path/to/project - Смена рабочей директории opencode serve @@ -1516,6 +1561,8 @@ async def run(self): self.opencode_client = OpenCodeClient() await self.opencode_client.__aenter__() + self.reasoning_timeout = ReasoningTimeout(self.opencode_client.send_prompt) + self.sse_listener = SSEListener(OPENCODE_URL) self._register_sse_callbacks() await self.sse_listener.start() From 9baf031058795e598bab8093cfed4882c1e03ae3 Mon Sep 17 00:00:00 2001 From: grishberg Date: Fri, 17 Jul 2026 16:19:04 +0300 Subject: [PATCH 2/9] fix(bot): filter processor error from text delivery to chat --- bot/sse_listener.py | 4 ++++ bot/vk_longpoll.py | 4 ++++ 2 files changed, 8 insertions(+) diff --git a/bot/sse_listener.py b/bot/sse_listener.py index 003b4a872a78..1f0c73f49256 100644 --- a/bot/sse_listener.py +++ b/bot/sse_listener.py @@ -100,6 +100,10 @@ async def _dispatch(self, raw: str): except Exception as e: logger.exception(f"SSE callback error for {event_type}: {e}") + # Log full event data for ended/failed events (text, reasoning, step) + if event_type.endswith(".ended") or event_type.endswith(".failed"): + logger.debug(f"SSE {event_type} data: {json.dumps(event_data, ensure_ascii=False)[:2000]}") + for cb in self._callbacks.get("*", []): try: await cb(event_type, event_data) diff --git a/bot/vk_longpoll.py b/bot/vk_longpoll.py index a304a79aa0e8..0e58ed973eca 100644 --- a/bot/vk_longpoll.py +++ b/bot/vk_longpoll.py @@ -169,6 +169,10 @@ async def _on_text_ended(self, event_type: str, data: dict): text = data.get("text", "") if not text.strip(): return + # Suppress internal processor error messages leaked to SSE + if text.startswith("[ERROR]"): + logger.debug(f"Suppressing processor error in text: {text[:100]}") + return if self.reasoning_timeout: self.reasoning_timeout.activity(session_id) self._pending_texts.setdefault(session_id, []).append(text) From 8f80994f6692aee39861019f3860687175340524 Mon Sep 17 00:00:00 2001 From: grishberg Date: Mon, 20 Jul 2026 07:33:53 +0300 Subject: [PATCH 3/9] tune processor.ts --- bot/opencode_client.py | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/bot/opencode_client.py b/bot/opencode_client.py index 201166c77d9a..47d55f3236a7 100644 --- a/bot/opencode_client.py +++ b/bot/opencode_client.py @@ -171,6 +171,27 @@ async def get_child_sessions(self, parent_id: str) -> List[dict]: all_sessions = await self.get_all_sessions() return [s for s in all_sessions if s.get("parentID") == parent_id] + async def summarize_session(self, session_id: str, provider_id: str, model_id: str) -> bool: + """Принудительно сжимает (суммаризирует) контекст сессии. + + Использует POST /session/:sessionID/summarize. Операция может выполняться + долго (реальное AI-сжатие), поэтому используется увеличенный таймаут. + """ + url = f"{self.base_url}/session/{session_id}/summarize" + data = {"providerID": provider_id, "modelID": model_id} + long_timeout = ClientTimeout(total=600) + try: + async with self.session.post(url, json=data, timeout=long_timeout) as resp: + if resp.status in (200, 204): + logger.info(f"Session {session_id} summarized") + return True + text = await resp.text() + logger.error(f"Failed to summarize session: {resp.status} - {text}") + return False + except Exception as e: + logger.warning(f"Failed to summarize session {session_id}: {e}") + return False + async def delete_session(self, session_id: str) -> bool: """Удаляет сессию по ID.""" try: From bf2f04e378d12933228d7c71f144e4f0f8728e3a Mon Sep 17 00:00:00 2001 From: grishberg Date: Tue, 21 Jul 2026 14:05:35 +0300 Subject: [PATCH 4/9] fix(bot): restart reasoning timeout on repeated reasoning events --- bot/reasoning_timeout.py | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/bot/reasoning_timeout.py b/bot/reasoning_timeout.py index 4140c96ee60f..b7c74734d99c 100644 --- a/bot/reasoning_timeout.py +++ b/bot/reasoning_timeout.py @@ -29,15 +29,18 @@ def __init__(self, send_prompt: Callable[[str, str], Awaitable[bool]]): self._timers: dict[str, asyncio.Task] = {} def reasoning_ended(self, session_id: str, user_id: int): - """reasoning.ended — запускаем таймер.""" - if session_id in self._active: - return # уже есть активный таймер - + """reasoning.ended — запускаем (или перезапускаем) таймер.""" + was_active = session_id in self._active self._active[session_id] = {"user_id": user_id} + if was_active: + logger.debug(f"Reasoning timer restarted for session {session_id}") + else: + logger.debug(f"Reasoning timer started for session {session_id}") self._start_timer(session_id) def activity(self, session_id: str): """Любая активность (text.ended, step.ended, tool.called) — сбрасываем таймер.""" + logger.debug(f"Reasoning timer cancelled for session {session_id}") self._cancel_timer(session_id) self._active.pop(session_id, None) From 22f48fa578cc5da5c4a3b06354aa05585fb5c6b7 Mon Sep 17 00:00:00 2001 From: grishberg Date: Tue, 21 Jul 2026 15:54:15 +0300 Subject: [PATCH 5/9] refactor: remove V2 session event system (session.next.*), keep only V1 --- bot/opencode_process.py | 2 - bot/vk_longpoll.py | 141 ------- packages/opencode/src/effect/runtime-flags.ts | 1 - packages/opencode/src/session/compaction.ts | 22 -- packages/opencode/src/session/processor.ts | 357 +----------------- packages/opencode/src/session/prompt.ts | 71 ---- 6 files changed, 8 insertions(+), 586 deletions(-) diff --git a/bot/opencode_process.py b/bot/opencode_process.py index d725e7f44595..6f64872dbbd6 100644 --- a/bot/opencode_process.py +++ b/bot/opencode_process.py @@ -122,8 +122,6 @@ def _build_args(self) -> list[str]: def _build_env(self) -> dict[str, str]: """Строит environment для opencode процесса.""" env = os.environ.copy() - # Включаем v2 события (session.next.text.ended и т.д.) - env["OPENCODE_EXPERIMENTAL_EVENT_SYSTEM"] = "true" # Передаём путь к конфигу с MCP серверами env["OPENCODE_CONFIG"] = str(Path(OPENCODE_CONFIG_PATH).expanduser()) return env diff --git a/bot/vk_longpoll.py b/bot/vk_longpoll.py index 0e58ed973eca..d7388a2ff372 100644 --- a/bot/vk_longpoll.py +++ b/bot/vk_longpoll.py @@ -121,10 +121,6 @@ def __init__( # session_id -> user_id для маршрутизации SSE событий self.session_to_user: Dict[str, int] = {} # session_id -> user_id - # Buffers for debouncing text.ended alongside tool calls - self._pending_texts: Dict[str, list[str]] = {} # session_id -> [text] - self._step_has_tools: Dict[str, bool] = {} # session_id -> bool - # Child session (subagent/subtask) tracking self.parent_child_map: Dict[str, Dict[str, dict]] = {} # parent_id -> {child_id: {title, ...}} @@ -142,14 +138,6 @@ def __init__( def _register_sse_callbacks(self): """Регистрирует SSE колбэки для обработки событий OpenCode""" - self.sse_listener.on("session.next.text.ended", self._on_text_ended) - self.sse_listener.on("session.next.reasoning.ended", self._on_reasoning_ended) - self.sse_listener.on("session.next.tool.called", self._on_tool_event) - self.sse_listener.on("session.next.tool.success", self._on_tool_event) - self.sse_listener.on("session.next.tool.failed", self._on_tool_event) - self.sse_listener.on("session.next.step.started", self._on_step_event) - self.sse_listener.on("session.next.step.ended", self._on_step_event) - self.sse_listener.on("session.next.step.failed", self._on_step_event) self.sse_listener.on("permission.asked", self._on_permission) self.sse_listener.on("permission.v2.asked", self._on_permission) # v2 API fallback self.sse_listener.on("question.asked", self._on_question) @@ -160,135 +148,6 @@ def _register_sse_callbacks(self): self.sse_listener.on("todo.updated", self._on_todo_updated) self.sse_listener.on_any(self._on_any_event) - async def _on_text_ended(self, event_type: str, data: dict): - """Буферизует текст — доставит только если шаг без tool calls.""" - session_id = data.get("sessionID", "") - user_id = self.session_to_user.get(session_id) - if not user_id: - return - text = data.get("text", "") - if not text.strip(): - return - # Suppress internal processor error messages leaked to SSE - if text.startswith("[ERROR]"): - logger.debug(f"Suppressing processor error in text: {text[:100]}") - return - if self.reasoning_timeout: - self.reasoning_timeout.activity(session_id) - self._pending_texts.setdefault(session_id, []).append(text) - - async def _on_reasoning_ended(self, event_type: str, data: dict): - """Обрабатывает завершение рассуждения""" - session_id = data.get("sessionID", "") - user_id = self.session_to_user.get(session_id) - if not user_id: - return - text = data.get("text", "") - if not text.strip(): - return - target = THINKING_PEER_ID if THINKING_PEER_ID else user_id - await self.vk.send_message(target, f"🧠:\n{text}") - if self.reasoning_timeout: - self.reasoning_timeout.reasoning_ended(session_id, user_id) - - async def _on_tool_event(self, event_type: str, data: dict): - """Обрабатывает события инструментов с детальным описанием""" - session_id = data.get("sessionID", "") - user_id = self.session_to_user.get(session_id) - if not user_id: - return - - if self.reasoning_timeout: - self.reasoning_timeout.activity(session_id) - - tool_name = data.get("tool", "") or data.get("name", "") - target = THINKING_PEER_ID if THINKING_PEER_ID else user_id - action = event_type.split(".")[-1] if "." in event_type else "event" - - if action == "called": - self._step_has_tools[session_id] = True - tool_input = data.get("input", {}) - if isinstance(tool_input, dict): - # Показываем ключевые параметры (не весь JSON чтобы не спамить) - keys = list(tool_input.keys())[:5] - params = ", ".join(f"{k}={repr(tool_input[k])[:50]}" for k in keys) - desc = f"({params})" if params else "" - else: - desc = "" - await self.vk.send_message(target, f"🔧 Вызов: {tool_name}{desc}") - - elif action == "success": - result = data.get("structured", data.get("result", {})) - if isinstance(result, dict): - keys = list(result.keys())[:3] - preview = ", ".join(f"{k}" for k in keys) - desc = f" → [{preview}]" if preview else "" - else: - desc = "" - await self.vk.send_message(target, f"🔧 Готово: {tool_name}{desc}") - - elif action == "failed": - error = data.get("error", {}) - msg = error.get("message", "?") if isinstance(error, dict) else str(error) - await self.vk.send_message(target, f"❌ Ошибка: {tool_name} — {msg}") - - else: - await self.vk.send_message(target, f"🔧 Tool {tool_name} ({action})") - - async def _on_step_event(self, event_type: str, data: dict): - """Обрабатывает события шагов с детальным описанием""" - session_id = data.get("sessionID", "") - user_id = self.session_to_user.get(session_id) - if not user_id: - return - - if self.reasoning_timeout: - self.reasoning_timeout.activity(session_id) - - step_type = event_type.split(".")[-1] if "." in event_type else "step" - target = THINKING_PEER_ID if THINKING_PEER_ID else user_id - - if step_type == "started": - agent = data.get("agent", "?") - model = data.get("model", {}) - model_id = model.get("id", "?") if isinstance(model, dict) else str(model) - await self.vk.send_message(target, f"🧠: Шаг начат — агент={agent}, модель={model_id}") - # Reset tool flag for the new step - self._step_has_tools[session_id] = False - - elif step_type == "ended": - # Deliver buffered text only if this step had no tool calls - if not self._step_has_tools.get(session_id, False): - texts = self._pending_texts.pop(session_id, []) - if texts: - await self.vk.send_message(user_id, "\n\n".join(texts)) - else: - self._pending_texts.pop(session_id, None) - self._step_has_tools.pop(session_id, None) - - finish = data.get("finish", "?") - tokens = data.get("tokens", {}) - input_tok = tokens.get("input", 0) - output_tok = tokens.get("output", 0) - await self.vk.send_message( - target, - f"🧠: Шаг завершён — причина={finish}, токены in={input_tok} out={output_tok}", - ) - - elif step_type == "failed": - # On failure, still deliver any buffered text - texts = self._pending_texts.pop(session_id, []) - if texts: - await self.vk.send_message(user_id, "\n\n".join(texts)) - self._step_has_tools.pop(session_id, None) - - error = data.get("error", {}) - msg = error.get("message", "?") if isinstance(error, dict) else str(error) - await self.vk.send_message(target, f"❌: Шаг провалился — {msg}") - - else: - await self.vk.send_message(target, f"🧠: Step {step_type}") - async def _on_permission(self, event_type: str, data: dict): """Обрабатывает запрос разрешения через SSE""" session_id = data.get("sessionID", "") diff --git a/packages/opencode/src/effect/runtime-flags.ts b/packages/opencode/src/effect/runtime-flags.ts index 58dc50d0278c..ede1d705d213 100644 --- a/packages/opencode/src/effect/runtime-flags.ts +++ b/packages/opencode/src/effect/runtime-flags.ts @@ -45,7 +45,6 @@ export class Service extends ConfigService.Service()("@opencode/Runtime experimentalLspTool: enabledByExperimental("OPENCODE_EXPERIMENTAL_LSP_TOOL"), experimentalOxfmt: enabledByExperimental("OPENCODE_EXPERIMENTAL_OXFMT"), experimentalPlanMode: enabledByExperimental("OPENCODE_EXPERIMENTAL_PLAN_MODE"), - experimentalEventSystem: enabledByExperimental("OPENCODE_EXPERIMENTAL_EVENT_SYSTEM"), experimentalWorkspaces: enabledByExperimental("OPENCODE_EXPERIMENTAL_WORKSPACES"), experimentalIconDiscovery: enabledByExperimental("OPENCODE_EXPERIMENTAL_ICON_DISCOVERY"), outputTokenMax: positiveInteger("OPENCODE_EXPERIMENTAL_OUTPUT_TOKEN_MAX"), diff --git a/packages/opencode/src/session/compaction.ts b/packages/opencode/src/session/compaction.ts index a12348c92778..6a655d97b327 100644 --- a/packages/opencode/src/session/compaction.ts +++ b/packages/opencode/src/session/compaction.ts @@ -13,14 +13,11 @@ import { Config } from "@/config/config" import { NotFoundError } from "@/storage/storage" import { Effect, Layer, Context } from "effect" -import * as DateTime from "effect/DateTime" import { InstanceState } from "@/effect/instance-state" import { isOverflow as overflow, usable } from "./overflow" import { serviceUse } from "@opencode-ai/core/effect/service-use" import { RuntimeFlags } from "@/effect/runtime-flags" import { EventV2Bridge } from "@/event-v2-bridge" -import { SessionEvent } from "@opencode-ai/core/session/event" -import { SessionMessage } from "@opencode-ai/core/session/message" import { ProviderV2 } from "@opencode-ai/core/provider" import { ModelV2 } from "@opencode-ai/core/model" import { EventV2 } from "@opencode-ai/core/event" @@ -535,17 +532,6 @@ export const layer = Layer.effect( parts: [], }, ) - if (flags.experimentalEventSystem) { - if (summary) - yield* events.publish(SessionEvent.Compaction.Ended, { - sessionID: input.sessionID, - messageID: SessionMessage.ID.make(input.parentID), - timestamp: DateTime.makeUnsafe(Date.now()), - reason: input.auto ? "auto" : "manual", - text: summary ?? "", - recent, - }) - } yield* events.publish(Event.Compacted, { sessionID: input.sessionID }) } return result @@ -574,14 +560,6 @@ export const layer = Layer.effect( auto: input.auto, overflow: input.overflow, }) - if (flags.experimentalEventSystem) { - yield* events.publish(SessionEvent.Compaction.Started, { - sessionID: input.sessionID, - messageID: SessionMessage.ID.make(msg.id), - timestamp: DateTime.makeUnsafe(Date.now()), - reason: input.auto ? "auto" : "manual", - }) - } }) return Service.of({ diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 05398a4411d6..26c3cd5edb69 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -25,13 +25,7 @@ import { errorMessage } from "@/util/error" import { isRecord } from "@/util/record" import { EventV2Bridge } from "@/event-v2-bridge" import { Database } from "@opencode-ai/core/database/database" -import { SessionEvent } from "@opencode-ai/core/session/event" -import { SessionMessage } from "@opencode-ai/core/session/message" import { findToolCallName } from "@opencode-ai/core/session/runner/publish-llm-event" -import { ModelV2 } from "@opencode-ai/core/model" -import { ProviderV2 } from "@opencode-ai/core/provider" -import * as DateTime from "effect/DateTime" -import { RuntimeFlags } from "@/effect/runtime-flags" import { ToolOutput, Usage, type LLMEvent } from "@opencode-ai/llm" const DOOM_LOOP_THRESHOLD = 3 @@ -66,7 +60,6 @@ export interface Interface { } type ToolCall = { - assistantMessageID?: SessionMessage.ID partID: SessionV1.ToolPart["id"] messageID: SessionV1.ToolPart["messageID"] sessionID: SessionV1.ToolPart["sessionID"] @@ -84,7 +77,6 @@ interface ProcessorContext extends Input { currentText: SessionV1.TextPart | undefined currentTextID: string | undefined reasoningMap: Record - v2AssistantMessageID: SessionMessage.ID | undefined textBasedToolCall: boolean textBasedToolCallCount: number } @@ -108,7 +100,6 @@ export const layer = Layer.effect( const status = yield* SessionStatus.Service const image = yield* Image.Service const events = yield* EventV2Bridge.Service - const flags = yield* RuntimeFlags.Service const database = yield* Database.Service const create = Effect.fn("SessionProcessor.create")(function* (input: Input) { @@ -128,12 +119,10 @@ export const layer = Layer.effect( currentText: undefined, currentTextID: undefined, reasoningMap: {}, - v2AssistantMessageID: undefined, textBasedToolCall: false, textBasedToolCallCount: 0, } const TEXT_BASED_TOOL_CALL_MAX = 3 - const mirrorAssistant = flags.experimentalEventSystem && !input.assistantMessage.summary let aborted = false const parse = (e: unknown) => @@ -148,34 +137,6 @@ export const layer = Layer.effect( if (done) yield* Deferred.succeed(done, undefined).pipe(Effect.ignore) }) - const ensureV2AssistantMessage = Effect.fn("SessionProcessor.ensureV2AssistantMessage")(function* () { - if (ctx.v2AssistantMessageID) return ctx.v2AssistantMessageID - ctx.v2AssistantMessageID = SessionMessage.ID.create() - yield* events.publish(SessionEvent.Step.Started, { - sessionID: ctx.sessionID, - assistantMessageID: ctx.v2AssistantMessageID, - agent: input.assistantMessage.agent, - model: { - id: ModelV2.ID.make(ctx.model.id), - providerID: ProviderV2.ID.make(ctx.model.providerID), - variant: ModelV2.VariantID.make(input.assistantMessage.variant ?? "default"), - }, - snapshot: ctx.snapshot, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - return ctx.v2AssistantMessageID - }) - - const requireV2AssistantMessage = (toolCall?: ToolCall) => - toolCall?.assistantMessageID === undefined - ? Effect.die("V2 tool settlement has no owning assistant message") - : Effect.succeed(toolCall.assistantMessageID) - - const currentV2AssistantMessage = () => - ctx.v2AssistantMessageID === undefined - ? Effect.die("V2 step settlement has no owning assistant message") - : Effect.succeed(ctx.v2AssistantMessageID) - const readToolCall = Effect.fn("SessionProcessor.readToolCall")(function* (toolCallID: string) { const call = ctx.toolcalls[toolCallID] if (!call) return undefined @@ -264,17 +225,6 @@ export const layer = Layer.effect( ctx.reasoningMap[reasoningID].text = `[ERROR] You outputted a tool call as text ("${reasoningToolCall}") in your reasoning instead of using the structured tool_calls field in the API response. The tool was NOT executed. You MUST retry using the proper tool_calls mechanism. (Attempt ${ctx.textBasedToolCallCount}/${TEXT_BASED_TOOL_CALL_MAX})` } - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Reasoning.Ended, { - sessionID: ctx.sessionID, - assistantMessageID: yield* currentV2AssistantMessage(), - reasoningID, - text: ctx.reasoningMap[reasoningID].text, - providerMetadata: ctx.reasoningMap[reasoningID].metadata, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } // oxlint-disable-next-line no-self-assign -- reactivity trigger ctx.reasoningMap[reasoningID].text = ctx.reasoningMap[reasoningID].text ctx.reasoningMap[reasoningID].time = { ...ctx.reasoningMap[reasoningID].time, end: Date.now() } @@ -282,33 +232,6 @@ export const layer = Layer.effect( delete ctx.reasoningMap[reasoningID] }) - const flushV2Fragments = Effect.fn("SessionProcessor.flushV2Fragments")(function* () { - if (!mirrorAssistant) return - if (!ctx.assistantMessage.summary && ctx.currentText && ctx.currentTextID) { - yield* events.publish(SessionEvent.Text.Ended, { - sessionID: ctx.sessionID, - assistantMessageID: yield* currentV2AssistantMessage(), - textID: ctx.currentTextID, - text: ctx.currentText.text, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } - yield* Effect.forEach(Object.entries(ctx.reasoningMap), ([reasoningID, part]) => - currentV2AssistantMessage().pipe( - Effect.flatMap((assistantMessageID) => - events.publish(SessionEvent.Reasoning.Ended, { - sessionID: ctx.sessionID, - assistantMessageID, - reasoningID, - text: part.text, - providerMetadata: part.metadata, - timestamp: DateTime.makeUnsafe(Date.now()), - }), - ), - ), - ) - }) - const ensureToolCall = Effect.fn("SessionProcessor.ensureToolCall")(function* (input: { id: string name: string @@ -329,17 +252,6 @@ export const layer = Layer.effect( } return { call: ctx.toolcalls[input.id], part } } - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - const assistantMessageID = mirrorAssistant ? yield* ensureV2AssistantMessage() : undefined - if (assistantMessageID) { - yield* events.publish(SessionEvent.Tool.Input.Started, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: input.id, - name: input.name, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } const part = yield* session.updatePart({ id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -351,7 +263,6 @@ export const layer = Layer.effect( metadata: input.providerExecuted ? { providerExecuted: true } : undefined, } satisfies SessionV1.ToolPart) ctx.toolcalls[input.id] = { - assistantMessageID, done: yield* Deferred.make(), partID: part.id, messageID: part.messageID, @@ -389,16 +300,6 @@ export const layer = Layer.effect( switch (value.type) { case "reasoning-start": if (value.id in ctx.reasoningMap) return - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Reasoning.Started, { - sessionID: ctx.sessionID, - assistantMessageID: yield* ensureV2AssistantMessage(), - reasoningID: value.id, - providerMetadata: value.providerMetadata, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } ctx.reasoningMap[value.id] = { id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -416,15 +317,6 @@ export const layer = Layer.effect( if (!(value.id in ctx.reasoningMap)) return ctx.reasoningMap[value.id].text += value.text if (value.providerMetadata) ctx.reasoningMap[value.id].metadata = value.providerMetadata - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Reasoning.Delta, { - sessionID: ctx.sessionID, - assistantMessageID: yield* currentV2AssistantMessage(), - reasoningID: value.id, - delta: value.text, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } yield* session.updatePartDelta({ sessionID: ctx.reasoningMap[value.id].sessionID, messageID: ctx.reasoningMap[value.id].messageID, @@ -451,33 +343,12 @@ export const layer = Layer.effect( case "tool-input-delta": { const toolCall = yield* ensureToolCall(value) - const assistantMessageID = mirrorAssistant ? yield* requireV2AssistantMessage(toolCall.call) : undefined - if (assistantMessageID) { - yield* events.publish(SessionEvent.Tool.Input.Delta, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - delta: value.text, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } ctx.toolcalls[value.id] = { ...toolCall.call, raw: toolCall.call.raw + value.text } } return case "tool-input-end": { const toolCall = yield* ensureToolCall(value) - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - const assistantMessageID = yield* requireV2AssistantMessage(toolCall.call) - yield* events.publish(SessionEvent.Tool.Input.Ended, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - text: toolCall.call.raw, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } ctx.toolcalls[value.id] = { ...toolCall.call, inputEnded: true } return } @@ -488,35 +359,6 @@ export const layer = Layer.effect( } const toolCall = yield* ensureToolCall(value) const input = isRecord(value.input) ? value.input : { value: value.input } - if (!toolCall.call.inputEnded) { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - const assistantMessageID = yield* requireV2AssistantMessage(toolCall.call) - yield* events.publish(SessionEvent.Tool.Input.Ended, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - text: toolCall.call.raw, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } - } - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - const assistantMessageID = yield* requireV2AssistantMessage(toolCall.call) - yield* events.publish(SessionEvent.Tool.Called, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - tool: value.name, - input, - provider: { - executed: toolCall.part.metadata?.providerExecuted === true, - ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), - }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } yield* updateToolCall(value.id, (match) => ({ ...match, tool: value.name, @@ -567,22 +409,6 @@ export const layer = Layer.effect( const toolCall = yield* readToolCall(value.id) if (!toolCall && value.result.type === "error") return if (value.result.type === "error") { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - const assistantMessageID = yield* requireV2AssistantMessage(toolCall?.call) - yield* events.publish(SessionEvent.Tool.Failed, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - error: { type: "unknown", message: errorMessage(value.result.value) }, - result: value.result, - provider: { - executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, - ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), - }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } yield* failToolCall(value.id, value.result.value) return } @@ -608,81 +434,12 @@ export const layer = Layer.effect( : `${rawOutput.output}\n\n[${omitted} image${omitted === 1 ? "" : "s"} omitted: could not be resized below the image size limit.]`, attachments: attachments.length ? attachments : undefined, } - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - const assistantMessageID = yield* requireV2AssistantMessage(toolCall?.call) - const content = [ - { type: "text" as const, text: output.output }, - ...(output.attachments?.map( - (item: SessionV1.FilePart) => - ({ - type: "file", - uri: item.url, - mime: item.mime, - name: item.filename, - }) as const, - ) ?? []), - ] - const unsupported = content.find((item) => item.type === "file" && !item.uri.startsWith("data:")) - if (unsupported?.type === "file") { - const error = new Error( - `Tool attachment URI "${unsupported.uri}" must be materialized before durable V2 settlement`, - ) - yield* events.publish(SessionEvent.Tool.Failed, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - error: { - type: "unknown", - message: error.message, - }, - provider: { - executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, - ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), - }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - yield* failToolCall(value.id, error) - return - } else - yield* events.publish(SessionEvent.Tool.Success, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - structured: output.metadata, - content, - result: value.result, - provider: { - executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, - ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), - }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } yield* completeToolCall(value.id, output) return } case "tool-error": { const toolCall = yield* readToolCall(value.id) - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - const assistantMessageID = yield* requireV2AssistantMessage(toolCall?.call) - yield* events.publish(SessionEvent.Tool.Failed, { - sessionID: ctx.sessionID, - assistantMessageID, - callID: value.id, - error: { - type: "unknown", - message: value.message, - }, - provider: { - executed: toolCall?.part.metadata?.providerExecuted === true, - ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), - }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } yield* failToolCall(value.id, value.error ?? new Error(value.message)) return } @@ -692,12 +449,6 @@ export const layer = Layer.effect( case "step-start": if (!ctx.snapshot) ctx.snapshot = yield* snapshot.track() - if (!ctx.assistantMessage.summary) { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - yield* ensureV2AssistantMessage() - } - } yield* session.updatePart({ id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -715,21 +466,6 @@ export const layer = Layer.effect( usage: value.usage ?? new Usage({}), metadata: value.providerMetadata, }) - if (!ctx.assistantMessage.summary) { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Step.Ended, { - sessionID: ctx.sessionID, - assistantMessageID: yield* currentV2AssistantMessage(), - finish: value.reason, - cost: usage.cost, - tokens: usage.tokens, - snapshot: completedSnapshot, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - ctx.v2AssistantMessageID = undefined - } - } ctx.assistantMessage.finish = value.reason // A text-based tool call was detected and rewritten as an error. Keep the loop // running (as if tools were called) so the model retries with a proper call @@ -778,17 +514,6 @@ export const layer = Layer.effect( } case "text-start": - if (!ctx.assistantMessage.summary) { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Text.Started, { - sessionID: ctx.sessionID, - assistantMessageID: yield* ensureV2AssistantMessage(), - timestamp: DateTime.makeUnsafe(Date.now()), - textID: value.id, - }) - } - } ctx.currentText = { id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -806,15 +531,6 @@ export const layer = Layer.effect( if (!ctx.currentText) return ctx.currentText.text += value.text if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Text.Delta, { - sessionID: ctx.sessionID, - assistantMessageID: yield* currentV2AssistantMessage(), - textID: value.id, - delta: value.text, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } yield* session.updatePartDelta({ sessionID: ctx.currentText.sessionID, messageID: ctx.currentText.messageID, @@ -849,18 +565,6 @@ export const layer = Layer.effect( `[ERROR] You outputted a tool call as text ("${textToolCall}") instead of using the structured tool_calls field in the API response. The tool was NOT executed. You MUST retry using the proper tool_calls mechanism. (Attempt ${ctx.textBasedToolCallCount}/${TEXT_BASED_TOOL_CALL_MAX})` } } - if (!ctx.assistantMessage.summary) { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Text.Ended, { - sessionID: ctx.sessionID, - assistantMessageID: yield* currentV2AssistantMessage(), - text: ctx.currentText.text, - timestamp: DateTime.makeUnsafe(Date.now()), - textID: value.id, - }) - } - } { const end = Date.now() ctx.currentText.time = { start: ctx.currentText.time?.start ?? end, end } @@ -919,16 +623,6 @@ export const layer = Layer.effect( const match = yield* readToolCall(toolCallID) if (!match) continue const part = match.part - if (mirrorAssistant && match.call.assistantMessageID) { - yield* events.publish(SessionEvent.Tool.Failed, { - sessionID: ctx.sessionID, - assistantMessageID: match.call.assistantMessageID, - callID: toolCallID, - error: { type: "unknown", message: "Tool execution aborted" }, - provider: { executed: part.metadata?.providerExecuted === true }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } const end = Date.now() const metadata = "metadata" in part.state && isRecord(part.state.metadata) ? part.state.metadata : {} yield* session.updatePart({ @@ -955,7 +649,6 @@ export const layer = Layer.effect( stack: e instanceof Error ? e.stack : undefined, }) const error = parse(e) - yield* flushV2Fragments() if (SessionV1.ContextOverflowError.isInstance(error)) { if ((yield* config.get()).compaction?.auto === false && !ctx.assistantMessage.summary) { ctx.assistantMessage.error = error @@ -968,20 +661,6 @@ export const layer = Layer.effect( yield* events.publish(Session.Event.Error, { sessionID: ctx.sessionID, error }) return } - if (!ctx.assistantMessage.summary) { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (mirrorAssistant) { - yield* events.publish(SessionEvent.Step.Failed, { - sessionID: ctx.sessionID, - assistantMessageID: yield* ensureV2AssistantMessage(), - error: { - type: "unknown", - message: errorMessage(e), - }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - } - } ctx.assistantMessage.error = error yield* events.publish(Session.Event.Error, { sessionID: ctx.assistantMessage.sessionID, @@ -1029,32 +708,14 @@ export const layer = Layer.effect( SessionRetry.policy({ provider: input.model.providerID, parse, - set: (info) => { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - const event = mirrorAssistant - ? events.publish(SessionEvent.Retried, { - sessionID: ctx.sessionID, - attempt: info.attempt, - error: { - message: info.message, - isRetryable: true, - }, - timestamp: DateTime.makeUnsafe(Date.now()), - }) - : Effect.void - return flushV2Fragments().pipe( - Effect.andThen(event), - Effect.andThen( - status.set(ctx.sessionID, { - type: "retry", - attempt: info.attempt, - message: info.message, - action: info.action, - next: info.next, - }), - ), - ) - }, + set: (info) => + status.set(ctx.sessionID, { + type: "retry", + attempt: info.attempt, + message: info.message, + action: info.action, + next: info.next, + }), }), ), Effect.catch(halt), @@ -1105,7 +766,6 @@ export const defaultLayer = Layer.suspend(() => Layer.provide(SessionStatus.defaultLayer), Layer.provide(Image.defaultLayer), Layer.provide(Config.defaultLayer), - Layer.provide(RuntimeFlags.defaultLayer), Layer.provide(Database.defaultLayer), Layer.provide(EventV2Bridge.defaultLayer), ), @@ -1123,7 +783,6 @@ export const node = LayerNode.make(layer, [ SessionStatus.node, Image.node, EventV2Bridge.node, - RuntimeFlags.node, Database.node, ]) diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index f2bcab8e5a5d..6b565e029c5b 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -49,12 +49,9 @@ import { SessionRunState } from "./run-state" import { RuntimeFlags } from "@/effect/runtime-flags" import { EventV2Bridge } from "@/event-v2-bridge" import { Database } from "@opencode-ai/core/database/database" -import { SessionEvent } from "@opencode-ai/core/session/event" -import { SessionMessage } from "@opencode-ai/core/session/message" import { ModelV2 } from "@opencode-ai/core/model" import { ProviderV2 } from "@opencode-ai/core/provider" import { AgentAttachment, FileAttachment, Prompt, Source } from "@opencode-ai/core/session/prompt" -import * as DateTime from "effect/DateTime" import { eq } from "drizzle-orm" import { SessionTable } from "@opencode-ai/core/session/sql" import { SessionReminders } from "./reminders" @@ -500,15 +497,6 @@ export const layer = Layer.effect( }, } yield* sessions.updatePart(part) - if (flags.experimentalEventSystem) { - yield* events.publish(SessionEvent.Shell.Started, { - sessionID: input.sessionID, - messageID: SessionMessage.ID.create(), - timestamp: DateTime.makeUnsafe(started), - callID: part.callID, - command: input.command, - }) - } return { msg, part, cwd: ctx.directory } }).pipe(Effect.ensuring(markReady)) @@ -524,14 +512,6 @@ export const layer = Layer.effect( output += "\n\n" + ["", "User aborted the command", ""].join("\n") } const completed = Date.now() - if (flags.experimentalEventSystem) { - yield* events.publish(SessionEvent.Shell.Ended, { - sessionID: input.sessionID, - timestamp: DateTime.makeUnsafe(completed), - callID: part.callID, - output, - }) - } if (!msg.time.completed) { msg.time.completed = completed yield* sessions.updateMessage(msg) @@ -676,31 +656,6 @@ export const layer = Layer.effect( format: input.format, } - if (current?.agent !== info.agent) { - yield* events.publish(SessionEvent.AgentSwitched, { - sessionID: input.sessionID, - messageID: SessionMessage.ID.create(), - timestamp: DateTime.makeUnsafe(info.time.created), - agent: info.agent, - }) - } - if ( - current?.model?.providerID !== info.model.providerID || - current.model.id !== info.model.modelID || - (current.model.variant === "default" ? undefined : current.model.variant) !== info.model.variant - ) { - yield* events.publish(SessionEvent.ModelSwitched, { - sessionID: input.sessionID, - messageID: SessionMessage.ID.create(), - timestamp: DateTime.makeUnsafe(info.time.created), - model: { - id: ModelV2.ID.make(info.model.modelID), - providerID: ProviderV2.ID.make(info.model.providerID), - variant: ModelV2.VariantID.make(info.model.variant ?? "default"), - }, - }) - } - yield* Effect.addFinalizer(() => instruction.clear(info.id)) type Draft = T extends SessionV1.Part ? Omit & { id?: string } : never @@ -1073,32 +1028,6 @@ export const layer = Layer.effect( synthetic: [] as string[], }, ) - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (flags.experimentalEventSystem) { - yield* events.publish(SessionEvent.Prompted, { - sessionID: input.sessionID, - messageID: SessionMessage.ID.create(), - timestamp: DateTime.makeUnsafe(info.time.created), - delivery: "steer", - prompt: new Prompt({ - text: nextPrompt.text.join("\n"), - files: nextPrompt.files, - agents: nextPrompt.agents, - }), - }) - } - for (const text of nextPrompt.synthetic) { - // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (flags.experimentalEventSystem) { - yield* events.publish(SessionEvent.Synthetic, { - sessionID: input.sessionID, - messageID: SessionMessage.ID.create(), - timestamp: DateTime.makeUnsafe(info.time.created), - text, - }) - } - } - return { info, parts } }, Effect.scoped) From 4b68139518c8176f672f8303f504814b5d94e0af Mon Sep 17 00:00:00 2001 From: grishberg Date: Tue, 21 Jul 2026 16:16:49 +0300 Subject: [PATCH 6/9] fix: clean up remaining experimentalEventSystem references --- packages/opencode/src/plugin/tui/internal.ts | 7 ++----- packages/opencode/src/plugin/tui/runtime.ts | 2 +- packages/opencode/test/effect/runtime-flags.test.ts | 1 - packages/opencode/test/session/compaction.test.ts | 6 +++--- packages/opencode/test/session/processor-effect.test.ts | 2 +- packages/opencode/test/session/prompt.test.ts | 8 ++++---- packages/opencode/test/session/snapshot-tool-race.test.ts | 2 +- packages/tui/src/feature-plugins/builtins.ts | 2 +- 8 files changed, 13 insertions(+), 17 deletions(-) diff --git a/packages/opencode/src/plugin/tui/internal.ts b/packages/opencode/src/plugin/tui/internal.ts index cc33b15beeef..718ede4ad556 100644 --- a/packages/opencode/src/plugin/tui/internal.ts +++ b/packages/opencode/src/plugin/tui/internal.ts @@ -1,10 +1,7 @@ import { createBuiltinPlugins, type BuiltinTuiPlugin } from "@opencode-ai/tui/builtins" -import type { RuntimeFlags } from "@/effect/runtime-flags" export type InternalTuiPlugin = BuiltinTuiPlugin -export function internalTuiPlugins(flags: Pick): InternalTuiPlugin[] { - return createBuiltinPlugins({ - experimentalEventSystem: flags.experimentalEventSystem, - }) +export function internalTuiPlugins(): InternalTuiPlugin[] { + return createBuiltinPlugins({}) } diff --git a/packages/opencode/src/plugin/tui/runtime.ts b/packages/opencode/src/plugin/tui/runtime.ts index 4673805cf3bc..1e6cc696f295 100644 --- a/packages/opencode/src/plugin/tui/runtime.ts +++ b/packages/opencode/src/plugin/tui/runtime.ts @@ -1089,7 +1089,7 @@ async function load(input: { if (Flag.OPENCODE_PURE && pluginOrigins.length) { } - for (const item of internalTuiPlugins(flags)) { + for (const item of internalTuiPlugins()) { const entry = loadInternalPlugin(item) const meta = createMeta(entry.source, entry.spec, entry.target, undefined, entry.id) addPluginEntry(next, { diff --git a/packages/opencode/test/effect/runtime-flags.test.ts b/packages/opencode/test/effect/runtime-flags.test.ts index 2e1226b38ba4..f87f74c58081 100644 --- a/packages/opencode/test/effect/runtime-flags.test.ts +++ b/packages/opencode/test/effect/runtime-flags.test.ts @@ -55,7 +55,6 @@ describe("RuntimeFlags", () => { expect(flags.experimentalLspTool).toBe(true) expect(flags.experimentalOxfmt).toBe(true) expect(flags.experimentalPlanMode).toBe(true) - expect(flags.experimentalEventSystem).toBe(true) expect(flags.experimentalWorkspaces).toBe(true) expect(flags.experimentalIconDiscovery).toBe(true) expect(flags.experimentalNativeLlm).toBe(false) diff --git a/packages/opencode/test/session/compaction.test.ts b/packages/opencode/test/session/compaction.test.ts index 63276bfe197d..bcb9913a9205 100644 --- a/packages/opencode/test/session/compaction.test.ts +++ b/packages/opencode/test/session/compaction.test.ts @@ -232,7 +232,7 @@ const deps = Layer.mergeAll( Plugin.defaultLayer, EventV2Bridge.defaultLayer, Config.defaultLayer, - RuntimeFlags.layer({ experimentalEventSystem: true }), + RuntimeFlags.layer({}), Database.defaultLayer, EventV2Bridge.defaultLayer, ) @@ -274,7 +274,7 @@ function compactionProcessLayer(options?: CompactionProcessOptions) { ? SessionProcessorModule.SessionProcessor.layer.pipe( Layer.provide(summary), Layer.provide(Image.defaultLayer), - Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })), + Layer.provide(RuntimeFlags.layer({})), Layer.provide(status), ) : layer(options?.result ?? "continue") @@ -289,7 +289,7 @@ function compactionProcessLayer(options?: CompactionProcessOptions) { Layer.provide(status), Layer.provide(events), Layer.provide(options?.config ?? Config.defaultLayer), - Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })), + Layer.provide(RuntimeFlags.layer({})), Layer.provide(EventV2Bridge.defaultLayer), ) } diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index c8f40d0de1f9..259a20a8aa97 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -178,7 +178,7 @@ const root = LayerNode.group([ ]) const replacements = [ LayerNode.replace(SessionSummary.node, summary), - LayerNode.replace(RuntimeFlags.node, RuntimeFlags.layer({ experimentalEventSystem: true })), + LayerNode.replace(RuntimeFlags.node, RuntimeFlags.layer({})), ] const env = LayerNode.buildLayer(LayerNode.group([root, LayerNode.make(TestLLMServer.layer, [])]), { replacements }) diff --git a/packages/opencode/test/session/prompt.test.ts b/packages/opencode/test/session/prompt.test.ts index 08828018a4a3..8aea62edc69b 100644 --- a/packages/opencode/test/session/prompt.test.ts +++ b/packages/opencode/test/session/prompt.test.ts @@ -192,7 +192,7 @@ function makePrompt(input?: { processor?: "blocking" }) { Layer.provide(Git.defaultLayer), Layer.provide(Ripgrep.defaultLayer), Layer.provide(Format.defaultLayer), - Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })), + Layer.provide(RuntimeFlags.layer({})), Layer.provideMerge(todo), Layer.provideMerge(question), Layer.provideMerge(deps), @@ -204,11 +204,11 @@ function makePrompt(input?: { processor?: "blocking" }) { : SessionProcessor.layer.pipe( Layer.provide(summary), Layer.provide(Image.defaultLayer), - Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })), + Layer.provide(RuntimeFlags.layer({})), Layer.provideMerge(deps), ) const compact = SessionCompaction.layer.pipe( - Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })), + Layer.provide(RuntimeFlags.layer({})), Layer.provideMerge(proc), Layer.provideMerge(deps), ) @@ -223,7 +223,7 @@ function makePrompt(input?: { processor?: "blocking" }) { Layer.provideMerge(trunc), Layer.provide(Instruction.defaultLayer), Layer.provide(SystemPrompt.defaultLayer), - Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })), + Layer.provide(RuntimeFlags.layer({})), Layer.provideMerge(deps), Layer.provide(summary), ) diff --git a/packages/opencode/test/session/snapshot-tool-race.test.ts b/packages/opencode/test/session/snapshot-tool-race.test.ts index 8a3701e12518..8ef85988555d 100644 --- a/packages/opencode/test/session/snapshot-tool-race.test.ts +++ b/packages/opencode/test/session/snapshot-tool-race.test.ts @@ -89,7 +89,7 @@ const it = testEffect( replacements: [ LayerNode.replace(MCP.node, mcp), LayerNode.replace(LSP.node, lsp), - LayerNode.replace(RuntimeFlags.node, RuntimeFlags.layer({ experimentalEventSystem: true })), + LayerNode.replace(RuntimeFlags.node, RuntimeFlags.layer({})), ], }), ) diff --git a/packages/tui/src/feature-plugins/builtins.ts b/packages/tui/src/feature-plugins/builtins.ts index b67923f3c5d1..a4126d5ab131 100644 --- a/packages/tui/src/feature-plugins/builtins.ts +++ b/packages/tui/src/feature-plugins/builtins.ts @@ -18,7 +18,7 @@ export type BuiltinTuiPlugin = Omit & { enabled?: boolean } -export function createBuiltinPlugins(options: { experimentalEventSystem: boolean }): BuiltinTuiPlugin[] { +export function createBuiltinPlugins(_options?: {}): BuiltinTuiPlugin[] { return [ HomeFooter, HomeTips, From e479ffb064449b7cbe213607179abd34cfdb50c4 Mon Sep 17 00:00:00 2001 From: grishberg Date: Tue, 21 Jul 2026 16:38:33 +0300 Subject: [PATCH 7/9] Revert "refactor: remove V2 session event system (session.next.*), keep only V1" This reverts commit 22f48fa578cc5da5c4a3b06354aa05585fb5c6b7. --- bot/opencode_process.py | 2 + bot/vk_longpoll.py | 141 +++++++ packages/opencode/src/effect/runtime-flags.ts | 1 + packages/opencode/src/session/compaction.ts | 22 ++ packages/opencode/src/session/processor.ts | 357 +++++++++++++++++- packages/opencode/src/session/prompt.ts | 71 ++++ 6 files changed, 586 insertions(+), 8 deletions(-) diff --git a/bot/opencode_process.py b/bot/opencode_process.py index 6f64872dbbd6..d725e7f44595 100644 --- a/bot/opencode_process.py +++ b/bot/opencode_process.py @@ -122,6 +122,8 @@ def _build_args(self) -> list[str]: def _build_env(self) -> dict[str, str]: """Строит environment для opencode процесса.""" env = os.environ.copy() + # Включаем v2 события (session.next.text.ended и т.д.) + env["OPENCODE_EXPERIMENTAL_EVENT_SYSTEM"] = "true" # Передаём путь к конфигу с MCP серверами env["OPENCODE_CONFIG"] = str(Path(OPENCODE_CONFIG_PATH).expanduser()) return env diff --git a/bot/vk_longpoll.py b/bot/vk_longpoll.py index d7388a2ff372..0e58ed973eca 100644 --- a/bot/vk_longpoll.py +++ b/bot/vk_longpoll.py @@ -121,6 +121,10 @@ def __init__( # session_id -> user_id для маршрутизации SSE событий self.session_to_user: Dict[str, int] = {} # session_id -> user_id + # Buffers for debouncing text.ended alongside tool calls + self._pending_texts: Dict[str, list[str]] = {} # session_id -> [text] + self._step_has_tools: Dict[str, bool] = {} # session_id -> bool + # Child session (subagent/subtask) tracking self.parent_child_map: Dict[str, Dict[str, dict]] = {} # parent_id -> {child_id: {title, ...}} @@ -138,6 +142,14 @@ def __init__( def _register_sse_callbacks(self): """Регистрирует SSE колбэки для обработки событий OpenCode""" + self.sse_listener.on("session.next.text.ended", self._on_text_ended) + self.sse_listener.on("session.next.reasoning.ended", self._on_reasoning_ended) + self.sse_listener.on("session.next.tool.called", self._on_tool_event) + self.sse_listener.on("session.next.tool.success", self._on_tool_event) + self.sse_listener.on("session.next.tool.failed", self._on_tool_event) + self.sse_listener.on("session.next.step.started", self._on_step_event) + self.sse_listener.on("session.next.step.ended", self._on_step_event) + self.sse_listener.on("session.next.step.failed", self._on_step_event) self.sse_listener.on("permission.asked", self._on_permission) self.sse_listener.on("permission.v2.asked", self._on_permission) # v2 API fallback self.sse_listener.on("question.asked", self._on_question) @@ -148,6 +160,135 @@ def _register_sse_callbacks(self): self.sse_listener.on("todo.updated", self._on_todo_updated) self.sse_listener.on_any(self._on_any_event) + async def _on_text_ended(self, event_type: str, data: dict): + """Буферизует текст — доставит только если шаг без tool calls.""" + session_id = data.get("sessionID", "") + user_id = self.session_to_user.get(session_id) + if not user_id: + return + text = data.get("text", "") + if not text.strip(): + return + # Suppress internal processor error messages leaked to SSE + if text.startswith("[ERROR]"): + logger.debug(f"Suppressing processor error in text: {text[:100]}") + return + if self.reasoning_timeout: + self.reasoning_timeout.activity(session_id) + self._pending_texts.setdefault(session_id, []).append(text) + + async def _on_reasoning_ended(self, event_type: str, data: dict): + """Обрабатывает завершение рассуждения""" + session_id = data.get("sessionID", "") + user_id = self.session_to_user.get(session_id) + if not user_id: + return + text = data.get("text", "") + if not text.strip(): + return + target = THINKING_PEER_ID if THINKING_PEER_ID else user_id + await self.vk.send_message(target, f"🧠:\n{text}") + if self.reasoning_timeout: + self.reasoning_timeout.reasoning_ended(session_id, user_id) + + async def _on_tool_event(self, event_type: str, data: dict): + """Обрабатывает события инструментов с детальным описанием""" + session_id = data.get("sessionID", "") + user_id = self.session_to_user.get(session_id) + if not user_id: + return + + if self.reasoning_timeout: + self.reasoning_timeout.activity(session_id) + + tool_name = data.get("tool", "") or data.get("name", "") + target = THINKING_PEER_ID if THINKING_PEER_ID else user_id + action = event_type.split(".")[-1] if "." in event_type else "event" + + if action == "called": + self._step_has_tools[session_id] = True + tool_input = data.get("input", {}) + if isinstance(tool_input, dict): + # Показываем ключевые параметры (не весь JSON чтобы не спамить) + keys = list(tool_input.keys())[:5] + params = ", ".join(f"{k}={repr(tool_input[k])[:50]}" for k in keys) + desc = f"({params})" if params else "" + else: + desc = "" + await self.vk.send_message(target, f"🔧 Вызов: {tool_name}{desc}") + + elif action == "success": + result = data.get("structured", data.get("result", {})) + if isinstance(result, dict): + keys = list(result.keys())[:3] + preview = ", ".join(f"{k}" for k in keys) + desc = f" → [{preview}]" if preview else "" + else: + desc = "" + await self.vk.send_message(target, f"🔧 Готово: {tool_name}{desc}") + + elif action == "failed": + error = data.get("error", {}) + msg = error.get("message", "?") if isinstance(error, dict) else str(error) + await self.vk.send_message(target, f"❌ Ошибка: {tool_name} — {msg}") + + else: + await self.vk.send_message(target, f"🔧 Tool {tool_name} ({action})") + + async def _on_step_event(self, event_type: str, data: dict): + """Обрабатывает события шагов с детальным описанием""" + session_id = data.get("sessionID", "") + user_id = self.session_to_user.get(session_id) + if not user_id: + return + + if self.reasoning_timeout: + self.reasoning_timeout.activity(session_id) + + step_type = event_type.split(".")[-1] if "." in event_type else "step" + target = THINKING_PEER_ID if THINKING_PEER_ID else user_id + + if step_type == "started": + agent = data.get("agent", "?") + model = data.get("model", {}) + model_id = model.get("id", "?") if isinstance(model, dict) else str(model) + await self.vk.send_message(target, f"🧠: Шаг начат — агент={agent}, модель={model_id}") + # Reset tool flag for the new step + self._step_has_tools[session_id] = False + + elif step_type == "ended": + # Deliver buffered text only if this step had no tool calls + if not self._step_has_tools.get(session_id, False): + texts = self._pending_texts.pop(session_id, []) + if texts: + await self.vk.send_message(user_id, "\n\n".join(texts)) + else: + self._pending_texts.pop(session_id, None) + self._step_has_tools.pop(session_id, None) + + finish = data.get("finish", "?") + tokens = data.get("tokens", {}) + input_tok = tokens.get("input", 0) + output_tok = tokens.get("output", 0) + await self.vk.send_message( + target, + f"🧠: Шаг завершён — причина={finish}, токены in={input_tok} out={output_tok}", + ) + + elif step_type == "failed": + # On failure, still deliver any buffered text + texts = self._pending_texts.pop(session_id, []) + if texts: + await self.vk.send_message(user_id, "\n\n".join(texts)) + self._step_has_tools.pop(session_id, None) + + error = data.get("error", {}) + msg = error.get("message", "?") if isinstance(error, dict) else str(error) + await self.vk.send_message(target, f"❌: Шаг провалился — {msg}") + + else: + await self.vk.send_message(target, f"🧠: Step {step_type}") + async def _on_permission(self, event_type: str, data: dict): """Обрабатывает запрос разрешения через SSE""" session_id = data.get("sessionID", "") diff --git a/packages/opencode/src/effect/runtime-flags.ts b/packages/opencode/src/effect/runtime-flags.ts index ede1d705d213..58dc50d0278c 100644 --- a/packages/opencode/src/effect/runtime-flags.ts +++ b/packages/opencode/src/effect/runtime-flags.ts @@ -45,6 +45,7 @@ export class Service extends ConfigService.Service()("@opencode/Runtime experimentalLspTool: enabledByExperimental("OPENCODE_EXPERIMENTAL_LSP_TOOL"), experimentalOxfmt: enabledByExperimental("OPENCODE_EXPERIMENTAL_OXFMT"), experimentalPlanMode: enabledByExperimental("OPENCODE_EXPERIMENTAL_PLAN_MODE"), + experimentalEventSystem: enabledByExperimental("OPENCODE_EXPERIMENTAL_EVENT_SYSTEM"), experimentalWorkspaces: enabledByExperimental("OPENCODE_EXPERIMENTAL_WORKSPACES"), experimentalIconDiscovery: enabledByExperimental("OPENCODE_EXPERIMENTAL_ICON_DISCOVERY"), outputTokenMax: positiveInteger("OPENCODE_EXPERIMENTAL_OUTPUT_TOKEN_MAX"), diff --git a/packages/opencode/src/session/compaction.ts b/packages/opencode/src/session/compaction.ts index 6a655d97b327..a12348c92778 100644 --- a/packages/opencode/src/session/compaction.ts +++ b/packages/opencode/src/session/compaction.ts @@ -13,11 +13,14 @@ import { Config } from "@/config/config" import { NotFoundError } from "@/storage/storage" import { Effect, Layer, Context } from "effect" +import * as DateTime from "effect/DateTime" import { InstanceState } from "@/effect/instance-state" import { isOverflow as overflow, usable } from "./overflow" import { serviceUse } from "@opencode-ai/core/effect/service-use" import { RuntimeFlags } from "@/effect/runtime-flags" import { EventV2Bridge } from "@/event-v2-bridge" +import { SessionEvent } from "@opencode-ai/core/session/event" +import { SessionMessage } from "@opencode-ai/core/session/message" import { ProviderV2 } from "@opencode-ai/core/provider" import { ModelV2 } from "@opencode-ai/core/model" import { EventV2 } from "@opencode-ai/core/event" @@ -532,6 +535,17 @@ export const layer = Layer.effect( parts: [], }, ) + if (flags.experimentalEventSystem) { + if (summary) + yield* events.publish(SessionEvent.Compaction.Ended, { + sessionID: input.sessionID, + messageID: SessionMessage.ID.make(input.parentID), + timestamp: DateTime.makeUnsafe(Date.now()), + reason: input.auto ? "auto" : "manual", + text: summary ?? "", + recent, + }) + } yield* events.publish(Event.Compacted, { sessionID: input.sessionID }) } return result @@ -560,6 +574,14 @@ export const layer = Layer.effect( auto: input.auto, overflow: input.overflow, }) + if (flags.experimentalEventSystem) { + yield* events.publish(SessionEvent.Compaction.Started, { + sessionID: input.sessionID, + messageID: SessionMessage.ID.make(msg.id), + timestamp: DateTime.makeUnsafe(Date.now()), + reason: input.auto ? "auto" : "manual", + }) + } }) return Service.of({ diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 26c3cd5edb69..05398a4411d6 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -25,7 +25,13 @@ import { errorMessage } from "@/util/error" import { isRecord } from "@/util/record" import { EventV2Bridge } from "@/event-v2-bridge" import { Database } from "@opencode-ai/core/database/database" +import { SessionEvent } from "@opencode-ai/core/session/event" +import { SessionMessage } from "@opencode-ai/core/session/message" import { findToolCallName } from "@opencode-ai/core/session/runner/publish-llm-event" +import { ModelV2 } from "@opencode-ai/core/model" +import { ProviderV2 } from "@opencode-ai/core/provider" +import * as DateTime from "effect/DateTime" +import { RuntimeFlags } from "@/effect/runtime-flags" import { ToolOutput, Usage, type LLMEvent } from "@opencode-ai/llm" const DOOM_LOOP_THRESHOLD = 3 @@ -60,6 +66,7 @@ export interface Interface { } type ToolCall = { + assistantMessageID?: SessionMessage.ID partID: SessionV1.ToolPart["id"] messageID: SessionV1.ToolPart["messageID"] sessionID: SessionV1.ToolPart["sessionID"] @@ -77,6 +84,7 @@ interface ProcessorContext extends Input { currentText: SessionV1.TextPart | undefined currentTextID: string | undefined reasoningMap: Record + v2AssistantMessageID: SessionMessage.ID | undefined textBasedToolCall: boolean textBasedToolCallCount: number } @@ -100,6 +108,7 @@ export const layer = Layer.effect( const status = yield* SessionStatus.Service const image = yield* Image.Service const events = yield* EventV2Bridge.Service + const flags = yield* RuntimeFlags.Service const database = yield* Database.Service const create = Effect.fn("SessionProcessor.create")(function* (input: Input) { @@ -119,10 +128,12 @@ export const layer = Layer.effect( currentText: undefined, currentTextID: undefined, reasoningMap: {}, + v2AssistantMessageID: undefined, textBasedToolCall: false, textBasedToolCallCount: 0, } const TEXT_BASED_TOOL_CALL_MAX = 3 + const mirrorAssistant = flags.experimentalEventSystem && !input.assistantMessage.summary let aborted = false const parse = (e: unknown) => @@ -137,6 +148,34 @@ export const layer = Layer.effect( if (done) yield* Deferred.succeed(done, undefined).pipe(Effect.ignore) }) + const ensureV2AssistantMessage = Effect.fn("SessionProcessor.ensureV2AssistantMessage")(function* () { + if (ctx.v2AssistantMessageID) return ctx.v2AssistantMessageID + ctx.v2AssistantMessageID = SessionMessage.ID.create() + yield* events.publish(SessionEvent.Step.Started, { + sessionID: ctx.sessionID, + assistantMessageID: ctx.v2AssistantMessageID, + agent: input.assistantMessage.agent, + model: { + id: ModelV2.ID.make(ctx.model.id), + providerID: ProviderV2.ID.make(ctx.model.providerID), + variant: ModelV2.VariantID.make(input.assistantMessage.variant ?? "default"), + }, + snapshot: ctx.snapshot, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + return ctx.v2AssistantMessageID + }) + + const requireV2AssistantMessage = (toolCall?: ToolCall) => + toolCall?.assistantMessageID === undefined + ? Effect.die("V2 tool settlement has no owning assistant message") + : Effect.succeed(toolCall.assistantMessageID) + + const currentV2AssistantMessage = () => + ctx.v2AssistantMessageID === undefined + ? Effect.die("V2 step settlement has no owning assistant message") + : Effect.succeed(ctx.v2AssistantMessageID) + const readToolCall = Effect.fn("SessionProcessor.readToolCall")(function* (toolCallID: string) { const call = ctx.toolcalls[toolCallID] if (!call) return undefined @@ -225,6 +264,17 @@ export const layer = Layer.effect( ctx.reasoningMap[reasoningID].text = `[ERROR] You outputted a tool call as text ("${reasoningToolCall}") in your reasoning instead of using the structured tool_calls field in the API response. The tool was NOT executed. You MUST retry using the proper tool_calls mechanism. (Attempt ${ctx.textBasedToolCallCount}/${TEXT_BASED_TOOL_CALL_MAX})` } + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Reasoning.Ended, { + sessionID: ctx.sessionID, + assistantMessageID: yield* currentV2AssistantMessage(), + reasoningID, + text: ctx.reasoningMap[reasoningID].text, + providerMetadata: ctx.reasoningMap[reasoningID].metadata, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } // oxlint-disable-next-line no-self-assign -- reactivity trigger ctx.reasoningMap[reasoningID].text = ctx.reasoningMap[reasoningID].text ctx.reasoningMap[reasoningID].time = { ...ctx.reasoningMap[reasoningID].time, end: Date.now() } @@ -232,6 +282,33 @@ export const layer = Layer.effect( delete ctx.reasoningMap[reasoningID] }) + const flushV2Fragments = Effect.fn("SessionProcessor.flushV2Fragments")(function* () { + if (!mirrorAssistant) return + if (!ctx.assistantMessage.summary && ctx.currentText && ctx.currentTextID) { + yield* events.publish(SessionEvent.Text.Ended, { + sessionID: ctx.sessionID, + assistantMessageID: yield* currentV2AssistantMessage(), + textID: ctx.currentTextID, + text: ctx.currentText.text, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + yield* Effect.forEach(Object.entries(ctx.reasoningMap), ([reasoningID, part]) => + currentV2AssistantMessage().pipe( + Effect.flatMap((assistantMessageID) => + events.publish(SessionEvent.Reasoning.Ended, { + sessionID: ctx.sessionID, + assistantMessageID, + reasoningID, + text: part.text, + providerMetadata: part.metadata, + timestamp: DateTime.makeUnsafe(Date.now()), + }), + ), + ), + ) + }) + const ensureToolCall = Effect.fn("SessionProcessor.ensureToolCall")(function* (input: { id: string name: string @@ -252,6 +329,17 @@ export const layer = Layer.effect( } return { call: ctx.toolcalls[input.id], part } } + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + const assistantMessageID = mirrorAssistant ? yield* ensureV2AssistantMessage() : undefined + if (assistantMessageID) { + yield* events.publish(SessionEvent.Tool.Input.Started, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: input.id, + name: input.name, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } const part = yield* session.updatePart({ id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -263,6 +351,7 @@ export const layer = Layer.effect( metadata: input.providerExecuted ? { providerExecuted: true } : undefined, } satisfies SessionV1.ToolPart) ctx.toolcalls[input.id] = { + assistantMessageID, done: yield* Deferred.make(), partID: part.id, messageID: part.messageID, @@ -300,6 +389,16 @@ export const layer = Layer.effect( switch (value.type) { case "reasoning-start": if (value.id in ctx.reasoningMap) return + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Reasoning.Started, { + sessionID: ctx.sessionID, + assistantMessageID: yield* ensureV2AssistantMessage(), + reasoningID: value.id, + providerMetadata: value.providerMetadata, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } ctx.reasoningMap[value.id] = { id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -317,6 +416,15 @@ export const layer = Layer.effect( if (!(value.id in ctx.reasoningMap)) return ctx.reasoningMap[value.id].text += value.text if (value.providerMetadata) ctx.reasoningMap[value.id].metadata = value.providerMetadata + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Reasoning.Delta, { + sessionID: ctx.sessionID, + assistantMessageID: yield* currentV2AssistantMessage(), + reasoningID: value.id, + delta: value.text, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } yield* session.updatePartDelta({ sessionID: ctx.reasoningMap[value.id].sessionID, messageID: ctx.reasoningMap[value.id].messageID, @@ -343,12 +451,33 @@ export const layer = Layer.effect( case "tool-input-delta": { const toolCall = yield* ensureToolCall(value) + const assistantMessageID = mirrorAssistant ? yield* requireV2AssistantMessage(toolCall.call) : undefined + if (assistantMessageID) { + yield* events.publish(SessionEvent.Tool.Input.Delta, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + delta: value.text, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } ctx.toolcalls[value.id] = { ...toolCall.call, raw: toolCall.call.raw + value.text } } return case "tool-input-end": { const toolCall = yield* ensureToolCall(value) + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + const assistantMessageID = yield* requireV2AssistantMessage(toolCall.call) + yield* events.publish(SessionEvent.Tool.Input.Ended, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + text: toolCall.call.raw, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } ctx.toolcalls[value.id] = { ...toolCall.call, inputEnded: true } return } @@ -359,6 +488,35 @@ export const layer = Layer.effect( } const toolCall = yield* ensureToolCall(value) const input = isRecord(value.input) ? value.input : { value: value.input } + if (!toolCall.call.inputEnded) { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + const assistantMessageID = yield* requireV2AssistantMessage(toolCall.call) + yield* events.publish(SessionEvent.Tool.Input.Ended, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + text: toolCall.call.raw, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + } + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + const assistantMessageID = yield* requireV2AssistantMessage(toolCall.call) + yield* events.publish(SessionEvent.Tool.Called, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + tool: value.name, + input, + provider: { + executed: toolCall.part.metadata?.providerExecuted === true, + ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), + }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } yield* updateToolCall(value.id, (match) => ({ ...match, tool: value.name, @@ -409,6 +567,22 @@ export const layer = Layer.effect( const toolCall = yield* readToolCall(value.id) if (!toolCall && value.result.type === "error") return if (value.result.type === "error") { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + const assistantMessageID = yield* requireV2AssistantMessage(toolCall?.call) + yield* events.publish(SessionEvent.Tool.Failed, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + error: { type: "unknown", message: errorMessage(value.result.value) }, + result: value.result, + provider: { + executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, + ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), + }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } yield* failToolCall(value.id, value.result.value) return } @@ -434,12 +608,81 @@ export const layer = Layer.effect( : `${rawOutput.output}\n\n[${omitted} image${omitted === 1 ? "" : "s"} omitted: could not be resized below the image size limit.]`, attachments: attachments.length ? attachments : undefined, } + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + const assistantMessageID = yield* requireV2AssistantMessage(toolCall?.call) + const content = [ + { type: "text" as const, text: output.output }, + ...(output.attachments?.map( + (item: SessionV1.FilePart) => + ({ + type: "file", + uri: item.url, + mime: item.mime, + name: item.filename, + }) as const, + ) ?? []), + ] + const unsupported = content.find((item) => item.type === "file" && !item.uri.startsWith("data:")) + if (unsupported?.type === "file") { + const error = new Error( + `Tool attachment URI "${unsupported.uri}" must be materialized before durable V2 settlement`, + ) + yield* events.publish(SessionEvent.Tool.Failed, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + error: { + type: "unknown", + message: error.message, + }, + provider: { + executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, + ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), + }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + yield* failToolCall(value.id, error) + return + } else + yield* events.publish(SessionEvent.Tool.Success, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + structured: output.metadata, + content, + result: value.result, + provider: { + executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, + ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), + }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } yield* completeToolCall(value.id, output) return } case "tool-error": { const toolCall = yield* readToolCall(value.id) + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + const assistantMessageID = yield* requireV2AssistantMessage(toolCall?.call) + yield* events.publish(SessionEvent.Tool.Failed, { + sessionID: ctx.sessionID, + assistantMessageID, + callID: value.id, + error: { + type: "unknown", + message: value.message, + }, + provider: { + executed: toolCall?.part.metadata?.providerExecuted === true, + ...(value.providerMetadata ? { metadata: value.providerMetadata } : {}), + }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } yield* failToolCall(value.id, value.error ?? new Error(value.message)) return } @@ -449,6 +692,12 @@ export const layer = Layer.effect( case "step-start": if (!ctx.snapshot) ctx.snapshot = yield* snapshot.track() + if (!ctx.assistantMessage.summary) { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + yield* ensureV2AssistantMessage() + } + } yield* session.updatePart({ id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -466,6 +715,21 @@ export const layer = Layer.effect( usage: value.usage ?? new Usage({}), metadata: value.providerMetadata, }) + if (!ctx.assistantMessage.summary) { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Step.Ended, { + sessionID: ctx.sessionID, + assistantMessageID: yield* currentV2AssistantMessage(), + finish: value.reason, + cost: usage.cost, + tokens: usage.tokens, + snapshot: completedSnapshot, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + ctx.v2AssistantMessageID = undefined + } + } ctx.assistantMessage.finish = value.reason // A text-based tool call was detected and rewritten as an error. Keep the loop // running (as if tools were called) so the model retries with a proper call @@ -514,6 +778,17 @@ export const layer = Layer.effect( } case "text-start": + if (!ctx.assistantMessage.summary) { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Text.Started, { + sessionID: ctx.sessionID, + assistantMessageID: yield* ensureV2AssistantMessage(), + timestamp: DateTime.makeUnsafe(Date.now()), + textID: value.id, + }) + } + } ctx.currentText = { id: PartID.ascending(), messageID: ctx.assistantMessage.id, @@ -531,6 +806,15 @@ export const layer = Layer.effect( if (!ctx.currentText) return ctx.currentText.text += value.text if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Text.Delta, { + sessionID: ctx.sessionID, + assistantMessageID: yield* currentV2AssistantMessage(), + textID: value.id, + delta: value.text, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } yield* session.updatePartDelta({ sessionID: ctx.currentText.sessionID, messageID: ctx.currentText.messageID, @@ -565,6 +849,18 @@ export const layer = Layer.effect( `[ERROR] You outputted a tool call as text ("${textToolCall}") instead of using the structured tool_calls field in the API response. The tool was NOT executed. You MUST retry using the proper tool_calls mechanism. (Attempt ${ctx.textBasedToolCallCount}/${TEXT_BASED_TOOL_CALL_MAX})` } } + if (!ctx.assistantMessage.summary) { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Text.Ended, { + sessionID: ctx.sessionID, + assistantMessageID: yield* currentV2AssistantMessage(), + text: ctx.currentText.text, + timestamp: DateTime.makeUnsafe(Date.now()), + textID: value.id, + }) + } + } { const end = Date.now() ctx.currentText.time = { start: ctx.currentText.time?.start ?? end, end } @@ -623,6 +919,16 @@ export const layer = Layer.effect( const match = yield* readToolCall(toolCallID) if (!match) continue const part = match.part + if (mirrorAssistant && match.call.assistantMessageID) { + yield* events.publish(SessionEvent.Tool.Failed, { + sessionID: ctx.sessionID, + assistantMessageID: match.call.assistantMessageID, + callID: toolCallID, + error: { type: "unknown", message: "Tool execution aborted" }, + provider: { executed: part.metadata?.providerExecuted === true }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } const end = Date.now() const metadata = "metadata" in part.state && isRecord(part.state.metadata) ? part.state.metadata : {} yield* session.updatePart({ @@ -649,6 +955,7 @@ export const layer = Layer.effect( stack: e instanceof Error ? e.stack : undefined, }) const error = parse(e) + yield* flushV2Fragments() if (SessionV1.ContextOverflowError.isInstance(error)) { if ((yield* config.get()).compaction?.auto === false && !ctx.assistantMessage.summary) { ctx.assistantMessage.error = error @@ -661,6 +968,20 @@ export const layer = Layer.effect( yield* events.publish(Session.Event.Error, { sessionID: ctx.sessionID, error }) return } + if (!ctx.assistantMessage.summary) { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (mirrorAssistant) { + yield* events.publish(SessionEvent.Step.Failed, { + sessionID: ctx.sessionID, + assistantMessageID: yield* ensureV2AssistantMessage(), + error: { + type: "unknown", + message: errorMessage(e), + }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + } + } ctx.assistantMessage.error = error yield* events.publish(Session.Event.Error, { sessionID: ctx.assistantMessage.sessionID, @@ -708,14 +1029,32 @@ export const layer = Layer.effect( SessionRetry.policy({ provider: input.model.providerID, parse, - set: (info) => - status.set(ctx.sessionID, { - type: "retry", - attempt: info.attempt, - message: info.message, - action: info.action, - next: info.next, - }), + set: (info) => { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + const event = mirrorAssistant + ? events.publish(SessionEvent.Retried, { + sessionID: ctx.sessionID, + attempt: info.attempt, + error: { + message: info.message, + isRetryable: true, + }, + timestamp: DateTime.makeUnsafe(Date.now()), + }) + : Effect.void + return flushV2Fragments().pipe( + Effect.andThen(event), + Effect.andThen( + status.set(ctx.sessionID, { + type: "retry", + attempt: info.attempt, + message: info.message, + action: info.action, + next: info.next, + }), + ), + ) + }, }), ), Effect.catch(halt), @@ -766,6 +1105,7 @@ export const defaultLayer = Layer.suspend(() => Layer.provide(SessionStatus.defaultLayer), Layer.provide(Image.defaultLayer), Layer.provide(Config.defaultLayer), + Layer.provide(RuntimeFlags.defaultLayer), Layer.provide(Database.defaultLayer), Layer.provide(EventV2Bridge.defaultLayer), ), @@ -783,6 +1123,7 @@ export const node = LayerNode.make(layer, [ SessionStatus.node, Image.node, EventV2Bridge.node, + RuntimeFlags.node, Database.node, ]) diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index 6b565e029c5b..f2bcab8e5a5d 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -49,9 +49,12 @@ import { SessionRunState } from "./run-state" import { RuntimeFlags } from "@/effect/runtime-flags" import { EventV2Bridge } from "@/event-v2-bridge" import { Database } from "@opencode-ai/core/database/database" +import { SessionEvent } from "@opencode-ai/core/session/event" +import { SessionMessage } from "@opencode-ai/core/session/message" import { ModelV2 } from "@opencode-ai/core/model" import { ProviderV2 } from "@opencode-ai/core/provider" import { AgentAttachment, FileAttachment, Prompt, Source } from "@opencode-ai/core/session/prompt" +import * as DateTime from "effect/DateTime" import { eq } from "drizzle-orm" import { SessionTable } from "@opencode-ai/core/session/sql" import { SessionReminders } from "./reminders" @@ -497,6 +500,15 @@ export const layer = Layer.effect( }, } yield* sessions.updatePart(part) + if (flags.experimentalEventSystem) { + yield* events.publish(SessionEvent.Shell.Started, { + sessionID: input.sessionID, + messageID: SessionMessage.ID.create(), + timestamp: DateTime.makeUnsafe(started), + callID: part.callID, + command: input.command, + }) + } return { msg, part, cwd: ctx.directory } }).pipe(Effect.ensuring(markReady)) @@ -512,6 +524,14 @@ export const layer = Layer.effect( output += "\n\n" + ["", "User aborted the command", ""].join("\n") } const completed = Date.now() + if (flags.experimentalEventSystem) { + yield* events.publish(SessionEvent.Shell.Ended, { + sessionID: input.sessionID, + timestamp: DateTime.makeUnsafe(completed), + callID: part.callID, + output, + }) + } if (!msg.time.completed) { msg.time.completed = completed yield* sessions.updateMessage(msg) @@ -656,6 +676,31 @@ export const layer = Layer.effect( format: input.format, } + if (current?.agent !== info.agent) { + yield* events.publish(SessionEvent.AgentSwitched, { + sessionID: input.sessionID, + messageID: SessionMessage.ID.create(), + timestamp: DateTime.makeUnsafe(info.time.created), + agent: info.agent, + }) + } + if ( + current?.model?.providerID !== info.model.providerID || + current.model.id !== info.model.modelID || + (current.model.variant === "default" ? undefined : current.model.variant) !== info.model.variant + ) { + yield* events.publish(SessionEvent.ModelSwitched, { + sessionID: input.sessionID, + messageID: SessionMessage.ID.create(), + timestamp: DateTime.makeUnsafe(info.time.created), + model: { + id: ModelV2.ID.make(info.model.modelID), + providerID: ProviderV2.ID.make(info.model.providerID), + variant: ModelV2.VariantID.make(info.model.variant ?? "default"), + }, + }) + } + yield* Effect.addFinalizer(() => instruction.clear(info.id)) type Draft = T extends SessionV1.Part ? Omit & { id?: string } : never @@ -1028,6 +1073,32 @@ export const layer = Layer.effect( synthetic: [] as string[], }, ) + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (flags.experimentalEventSystem) { + yield* events.publish(SessionEvent.Prompted, { + sessionID: input.sessionID, + messageID: SessionMessage.ID.create(), + timestamp: DateTime.makeUnsafe(info.time.created), + delivery: "steer", + prompt: new Prompt({ + text: nextPrompt.text.join("\n"), + files: nextPrompt.files, + agents: nextPrompt.agents, + }), + }) + } + for (const text of nextPrompt.synthetic) { + // TODO(v2): Temporary dual-write while migrating session messages to v2 events. + if (flags.experimentalEventSystem) { + yield* events.publish(SessionEvent.Synthetic, { + sessionID: input.sessionID, + messageID: SessionMessage.ID.create(), + timestamp: DateTime.makeUnsafe(info.time.created), + text, + }) + } + } + return { info, parts } }, Effect.scoped) From ae6d83e4a583253b9d6aa0e0b8202494004296ae Mon Sep 17 00:00:00 2001 From: grishberg Date: Tue, 21 Jul 2026 16:40:33 +0300 Subject: [PATCH 8/9] fix: restore V2 session events, make unconditional (remove experimental flag) --- bot/vk_longpoll.py | 40 +++++---------------- packages/opencode/src/session/compaction.ts | 4 +-- packages/opencode/src/session/processor.ts | 2 +- packages/opencode/src/session/prompt.ts | 8 ++--- 4 files changed, 15 insertions(+), 39 deletions(-) diff --git a/bot/vk_longpoll.py b/bot/vk_longpoll.py index 0e58ed973eca..af31576306c7 100644 --- a/bot/vk_longpoll.py +++ b/bot/vk_longpoll.py @@ -121,10 +121,6 @@ def __init__( # session_id -> user_id для маршрутизации SSE событий self.session_to_user: Dict[str, int] = {} # session_id -> user_id - # Buffers for debouncing text.ended alongside tool calls - self._pending_texts: Dict[str, list[str]] = {} # session_id -> [text] - self._step_has_tools: Dict[str, bool] = {} # session_id -> bool - # Child session (subagent/subtask) tracking self.parent_child_map: Dict[str, Dict[str, dict]] = {} # parent_id -> {child_id: {title, ...}} @@ -209,7 +205,6 @@ async def _on_tool_event(self, event_type: str, data: dict): self._step_has_tools[session_id] = True tool_input = data.get("input", {}) if isinstance(tool_input, dict): - # Показываем ключевые параметры (не весь JSON чтобы не спамить) keys = list(tool_input.keys())[:5] params = ", ".join(f"{k}={repr(tool_input[k])[:50]}" for k in keys) desc = f"({params})" if params else "" @@ -257,37 +252,18 @@ async def _on_step_event(self, event_type: str, data: dict): self._step_has_tools[session_id] = False elif step_type == "ended": - # Deliver buffered text only if this step had no tool calls - if not self._step_has_tools.get(session_id, False): - texts = self._pending_texts.pop(session_id, []) - if texts: - await self.vk.send_message(user_id, "\n\n".join(texts)) - else: - self._pending_texts.pop(session_id, None) - self._step_has_tools.pop(session_id, None) - - finish = data.get("finish", "?") - tokens = data.get("tokens", {}) - input_tok = tokens.get("input", 0) - output_tok = tokens.get("output", 0) - await self.vk.send_message( - target, - f"🧠: Шаг завершён — причина={finish}, токены in={input_tok} out={output_tok}", - ) - - elif step_type == "failed": - # On failure, still deliver any buffered text + # Step ended — deliver buffered text if no tools were used texts = self._pending_texts.pop(session_id, []) - if texts: - await self.vk.send_message(user_id, "\n\n".join(texts)) - self._step_has_tools.pop(session_id, None) + has_tools = self._step_has_tools.pop(session_id, False) + if not has_tools and texts: + combined = "\n\n".join(texts) + if combined.strip(): + await self.vk.send_message(target, combined) + elif step_type == "failed": error = data.get("error", {}) msg = error.get("message", "?") if isinstance(error, dict) else str(error) - await self.vk.send_message(target, f"❌: Шаг провалился — {msg}") - - else: - await self.vk.send_message(target, f"🧠: Step {step_type}") + await self.vk.send_message(target, f"❌ Шаг упал — {msg}") async def _on_permission(self, event_type: str, data: dict): """Обрабатывает запрос разрешения через SSE""" diff --git a/packages/opencode/src/session/compaction.ts b/packages/opencode/src/session/compaction.ts index a12348c92778..b3918f45429a 100644 --- a/packages/opencode/src/session/compaction.ts +++ b/packages/opencode/src/session/compaction.ts @@ -535,7 +535,7 @@ export const layer = Layer.effect( parts: [], }, ) - if (flags.experimentalEventSystem) { + { if (summary) yield* events.publish(SessionEvent.Compaction.Ended, { sessionID: input.sessionID, @@ -574,7 +574,7 @@ export const layer = Layer.effect( auto: input.auto, overflow: input.overflow, }) - if (flags.experimentalEventSystem) { + { yield* events.publish(SessionEvent.Compaction.Started, { sessionID: input.sessionID, messageID: SessionMessage.ID.make(msg.id), diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 05398a4411d6..2ef490dcfb7e 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -133,7 +133,7 @@ export const layer = Layer.effect( textBasedToolCallCount: 0, } const TEXT_BASED_TOOL_CALL_MAX = 3 - const mirrorAssistant = flags.experimentalEventSystem && !input.assistantMessage.summary + const mirrorAssistant = !input.assistantMessage.summary let aborted = false const parse = (e: unknown) => diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index f2bcab8e5a5d..44f39fc6ddd9 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -500,7 +500,7 @@ export const layer = Layer.effect( }, } yield* sessions.updatePart(part) - if (flags.experimentalEventSystem) { + { yield* events.publish(SessionEvent.Shell.Started, { sessionID: input.sessionID, messageID: SessionMessage.ID.create(), @@ -524,7 +524,7 @@ export const layer = Layer.effect( output += "\n\n" + ["", "User aborted the command", ""].join("\n") } const completed = Date.now() - if (flags.experimentalEventSystem) { + { yield* events.publish(SessionEvent.Shell.Ended, { sessionID: input.sessionID, timestamp: DateTime.makeUnsafe(completed), @@ -1074,7 +1074,7 @@ export const layer = Layer.effect( }, ) // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (flags.experimentalEventSystem) { + { yield* events.publish(SessionEvent.Prompted, { sessionID: input.sessionID, messageID: SessionMessage.ID.create(), @@ -1089,7 +1089,7 @@ export const layer = Layer.effect( } for (const text of nextPrompt.synthetic) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. - if (flags.experimentalEventSystem) { + { yield* events.publish(SessionEvent.Synthetic, { sessionID: input.sessionID, messageID: SessionMessage.ID.create(), From f9f261b31139b31945fe2268a813bfa4cf1740f5 Mon Sep 17 00:00:00 2001 From: grishberg Date: Tue, 21 Jul 2026 16:40:55 +0300 Subject: [PATCH 9/9] fix: remove obsolete OPENCODE_EXPERIMENTAL_EVENT_SYSTEM env var --- bot/opencode_process.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/bot/opencode_process.py b/bot/opencode_process.py index d725e7f44595..6f64872dbbd6 100644 --- a/bot/opencode_process.py +++ b/bot/opencode_process.py @@ -122,8 +122,6 @@ def _build_args(self) -> list[str]: def _build_env(self) -> dict[str, str]: """Строит environment для opencode процесса.""" env = os.environ.copy() - # Включаем v2 события (session.next.text.ended и т.д.) - env["OPENCODE_EXPERIMENTAL_EVENT_SYSTEM"] = "true" # Передаём путь к конфигу с MCP серверами env["OPENCODE_CONFIG"] = str(Path(OPENCODE_CONFIG_PATH).expanduser()) return env