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: 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/reasoning_timeout.py b/bot/reasoning_timeout.py new file mode 100644 index 000000000000..b7c74734d99c --- /dev/null +++ b/bot/reasoning_timeout.py @@ -0,0 +1,76 @@ +""" +Отслеживание зависания модели после 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 — запускаем (или перезапускаем) таймер.""" + 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) + + 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/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 b3bd38fc2df6..af31576306c7 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 @@ -120,13 +121,12 @@ 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, ...}} + # 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 +165,12 @@ 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) 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" @@ -194,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 "" @@ -227,6 +237,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 @@ -239,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""" @@ -768,6 +762,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 +1148,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 +1515,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 +1541,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() 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/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(), 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,