From 6c69332ec61baef320f9923a500cae1810d33a4a Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 09:35:05 +0000 Subject: [PATCH 1/9] perf: eliminate redundant file reads + flush model support Fix #1: Cache MemorySnapshot and fileSet in PreparedTurn. In finalizeMemoryForTurn, skip the disk re-read entirely when no tools were used during the turn (turnUsedTools=false). When tools ran, reuse the cached fileSet (avoids re-resolving paths) but still re-read file contents to detect modifications. handleStreaming() now returns { usedTools } flag. Flush model: Add config.compact.flushModel option and AgentAdapter.promptWithModel() method. When configured, memory flush turns use a cheaper/faster model (e.g. Sonnet) instead of the expensive conversation model. Pi adapter implements via session.setModel() with automatic restore after flush. All 217 tests pass, no new type errors. --- src/agents/pi.ts | 42 +++++++++++++++++++++++++ src/gateway.ts | 16 +++++++--- src/memory/lifecycle.ts | 70 +++++++++++++++++++++++++++-------------- src/memory/types.ts | 8 +++++ src/types.ts | 6 ++++ 5 files changed, 115 insertions(+), 27 deletions(-) diff --git a/src/agents/pi.ts b/src/agents/pi.ts index 92c2cef..74afb20 100644 --- a/src/agents/pi.ts +++ b/src/agents/pi.ts @@ -331,6 +331,48 @@ export const createPiAgentAdapter: AgentAdapterFactory = (config) => { return enqueue(threadId, () => doPrompt(threadId, formatMessage(message))); }, + async promptWithModel(threadId: string, message: AgentMessage, modelId: string): Promise { + return enqueue(threadId, async () => { + const entry = await getOrCreate(threadId); + const currentModel = entry.session.model; + + // Resolve the target model (format: "provider/model-id") + let targetModel; + const [provider, ...rest] = modelId.split("/"); + const id = rest.join("/"); + if (provider && id) { + targetModel = modelRegistry.find(provider, id); + } + + if (!targetModel) { + console.warn(`[pi-agent] flush model "${modelId}" not found, using default`); + return doPrompt(threadId, formatMessage(message)); + } + + // Switch to flush model + try { + await entry.session.setModel(targetModel); + console.log(`[pi-agent] switched to flush model: ${modelId}`); + } catch (err) { + console.warn(`[pi-agent] failed to set flush model "${modelId}":`, (err as Error).message); + return doPrompt(threadId, formatMessage(message)); + } + + try { + return await doPrompt(threadId, formatMessage(message)); + } finally { + // Restore original model + if (currentModel) { + try { + await entry.session.setModel(currentModel); + } catch (err) { + console.error(`[pi-agent] failed to restore model:`, (err as Error).message); + } + } + } + }); + }, + promptStream(threadId: string, message: AgentMessage): AsyncIterable { const text = formatMessage(message); // Return an async iterable that is single-use by design. diff --git a/src/gateway.ts b/src/gateway.ts index 3578900..c37f68b 100644 --- a/src/gateway.ts +++ b/src/gateway.ts @@ -687,17 +687,20 @@ export class Gateway { const stopTyping = startTypingLoop(thread); try { + let turnUsedTools = false; if (agent.promptStream) { const ac = new AbortController(); abortControllers.set(agentThreadId, ac); try { - await this.handleStreaming(thread, agent.promptStream(agentThreadId, agentMessage), verboseThreads.has(agentThreadId), ac.signal); + const streamResult = await this.handleStreaming(thread, agent.promptStream(agentThreadId, agentMessage), verboseThreads.has(agentThreadId), ac.signal); + turnUsedTools = streamResult.usedTools; } finally { abortControllers.delete(agentThreadId); } } else { - // Fallback: non-streaming prompt + // Fallback: non-streaming prompt (assume tools may have been used) const reply = await agent.prompt(agentThreadId, agentMessage); + turnUsedTools = true; if (reply.text) { await this.postWithFallback(thread, reply.text); } @@ -705,9 +708,10 @@ export class Gateway { // ── Memory: post-turn finalize + pressure check ─── try { + if (memoryPrepared) memoryPrepared.turnUsedTools = turnUsedTools; const pressure = await finalizeMemoryForTurn( agentThreadId, - memoryPrepared?.beforeDigest ?? null, + memoryPrepared ?? { message: agentMessage, beforeDigest: null, injected: false }, agent, memoryRoot, this.config.memory, ); // Use higher severity between pending compact and current pressure @@ -927,8 +931,9 @@ export class Gateway { * - Tool starts/ends are sent as compact status messages. * - Turn boundaries trigger a new message for the next turn's text. */ - private async handleStreaming(thread: any, stream: AsyncIterable, verbose: boolean, signal?: AbortSignal) { + private async handleStreaming(thread: any, stream: AsyncIterable, verbose: boolean, signal?: AbortSignal): Promise<{ usedTools: boolean }> { let activeTools = new Map(); // toolCallId -> toolName + let usedTools = false; // Per-turn streaming state — each turn gets a fresh iterable + promise let currentPush: ((text: string) => void) | null = null; @@ -1032,6 +1037,7 @@ export class Gateway { case "tool_start": { activeTools.set(event.toolCallId, event.toolName); + usedTools = true; if (verbose) { try { await thread.post(`${toolIcon(event.toolName)} Running \`${event.toolName}\`…`); @@ -1102,6 +1108,8 @@ export class Gateway { if (currentPromise) { await flushCurrentStream(); } + + return { usedTools }; } /** Post text with markdown, falling back to plain text */ diff --git a/src/memory/lifecycle.ts b/src/memory/lifecycle.ts index fc4c5e8..dfc8d5c 100644 --- a/src/memory/lifecycle.ts +++ b/src/memory/lifecycle.ts @@ -9,7 +9,7 @@ */ import type { AgentAdapter, AgentMessage } from "../types"; -import type { MemoryConfig, MemoryMode, PreparedTurn, PressureLevel, ThreadMemoryState } from "./types"; +import type { MemoryConfig, MemoryFileSet, MemoryMode, MemorySnapshot, PreparedTurn, PressureLevel, ThreadMemoryState } from "./types"; import { resolveMemoryFiles, readMemorySnapshot, formatDate } from "./files"; import { loadThreadMemoryState, saveThreadMemoryState } from "./state"; import { shouldInjectMemory, classifyContextPressure, isSoftFlushOnCooldown } from "./policy"; @@ -52,11 +52,16 @@ export async function prepareMemoryForTurn( const mode = getMode(agent); - // Complement mode: no injection, just track digest for finalize - // Unknown mode: also skip — we can't inject correctly before knowing if agent has memory extension - // (mode is detected during session creation, which happens inside promptStream) + // Complement mode: no injection, but still read snapshot for finalize comparison if (mode === "complement" || mode === "unknown") { - return { message, beforeDigest: null, injected: false }; + // Read snapshot so finalize can detect agent-written changes without re-reading + try { + const fileSet = resolveMemoryFiles(rootDir, config); + const snapshot = await readMemorySnapshot(fileSet, config?.inject?.maxBytes); + return { message, beforeDigest: snapshot.digest, injected: false, fileSet, snapshot }; + } catch { + return { message, beforeDigest: null, injected: false }; + } } // Full mode: inject if needed @@ -91,10 +96,10 @@ export async function prepareMemoryForTurn( await saveThreadMemoryState(threadId, state); console.log(`[memory] injected into ${threadId} (reason: ${decision.reason}, ${snapshot.entries.length} files, digest: ${snapshot.digest})`); - return { message: injectedMessage, beforeDigest: snapshot.digest, injected: true, pendingCompact: pendingCompactLevel }; + return { message: injectedMessage, beforeDigest: snapshot.digest, injected: true, pendingCompact: pendingCompactLevel, fileSet, snapshot }; } - return { message, beforeDigest: snapshot.digest, injected: false, pendingCompact: pendingCompactLevel }; + return { message, beforeDigest: snapshot.digest, injected: false, pendingCompact: pendingCompactLevel, fileSet, snapshot }; } catch (err) { console.error(`[memory] prepareMemoryForTurn error:`, (err as Error).message); return { message, beforeDigest: null, injected: false }; @@ -109,11 +114,14 @@ export async function prepareMemoryForTurn( * In Full mode: check if agent wrote memory files (update digest). * Both modes: check context pressure for proactive compaction. * + * Uses cached fileSet from PreparedTurn to avoid re-resolving files. + * Only re-reads files if the turn included tool calls that could have modified them. + * * Returns the pressure level for the gateway to act on. */ export async function finalizeMemoryForTurn( threadId: string, - beforeDigest: string | null, + prepared: PreparedTurn, agent: AgentAdapter, rootDir: string, config?: MemoryConfig, @@ -121,21 +129,25 @@ export async function finalizeMemoryForTurn( if (config?.enabled === false) return "none"; const mode = getMode(agent); + const beforeDigest = prepared.beforeDigest; // In Full mode: check if agent modified memory files if (mode !== "complement" && beforeDigest) { - try { - const fileSet = resolveMemoryFiles(rootDir, config); - const snapshot = await readMemorySnapshot(fileSet, config?.inject?.maxBytes); - if (snapshot.digest !== beforeDigest) { - const state = await loadThreadMemoryState(threadId); - state.lastInjectedDigest = snapshot.digest; - state.lastKnownDigest = snapshot.digest; - await saveThreadMemoryState(threadId, state); - console.log(`[memory] agent updated memory files (new digest: ${snapshot.digest})`); + // Skip expensive re-read if no file-modifying tools ran during this turn + if (prepared.turnUsedTools !== false) { + try { + const fileSet = prepared.fileSet ?? resolveMemoryFiles(rootDir, config); + const snapshot = await readMemorySnapshot(fileSet, config?.inject?.maxBytes); + if (snapshot.digest !== beforeDigest) { + const state = await loadThreadMemoryState(threadId); + state.lastInjectedDigest = snapshot.digest; + state.lastKnownDigest = snapshot.digest; + await saveThreadMemoryState(threadId, state); + console.log(`[memory] agent updated memory files (new digest: ${snapshot.digest})`); + } + } catch (err) { + console.error(`[memory] finalizeMemoryForTurn digest check error:`, (err as Error).message); } - } catch (err) { - console.error(`[memory] finalizeMemoryForTurn digest check error:`, (err as Error).message); } } @@ -167,6 +179,8 @@ export async function finalizeMemoryForTurn( * 2. Compact the session * 3. Mark force re-inject for Full mode * + * Uses a cheaper model for flush turns if config.compact.flushModel is set. + * * Returns compaction result or null if nothing to compact. */ export async function flushMemoryThenCompact( @@ -177,6 +191,16 @@ export async function flushMemoryThenCompact( config?: MemoryConfig, ): Promise<{ tokensBefore: number; tokensAfter: number | null } | null> { const mode = getMode(agent); + const flushModel = config?.compact?.flushModel; + + /** Send flush prompt, preferring flushModel if available */ + async function sendFlush(text: string): Promise { + if (flushModel && agent.promptWithModel) { + await agent.promptWithModel(threadId, { text }, flushModel); + } else { + await agent.prompt(threadId, { text }); + } + } // Soft flush: just prompt to save, don't compact if (level === "soft") { @@ -188,10 +212,10 @@ export async function flushMemoryThenCompact( try { const flushText = buildFlushPrompt(mode === "unknown" ? "full" : mode, "soft"); - await agent.prompt(threadId, { text: flushText }); + await sendFlush(flushText); state.lastSoftFlushAt = new Date().toISOString(); await saveThreadMemoryState(threadId, state); - console.log(`[memory] soft flush completed for ${threadId}`); + console.log(`[memory] soft flush completed for ${threadId}${flushModel ? ` (model: ${flushModel})` : ""}`); } catch (err) { console.error(`[memory] soft flush failed for ${threadId}:`, (err as Error).message); } @@ -206,8 +230,8 @@ export async function flushMemoryThenCompact( try { // Step 1: flush const flushText = buildFlushPrompt(mode === "unknown" ? "full" : mode, effectiveLevel); - console.log(`[memory] flushing memory for ${threadId} (level: ${level})`); - await agent.prompt(threadId, { text: flushText }); + console.log(`[memory] flushing memory for ${threadId} (level: ${level}${flushModel ? `, model: ${flushModel}` : ""})`); + await sendFlush(flushText); // Step 2: compact console.log(`[memory] compacting ${threadId}`); diff --git a/src/memory/types.ts b/src/memory/types.ts index 36b5ca4..61148fd 100644 --- a/src/memory/types.ts +++ b/src/memory/types.ts @@ -40,6 +40,8 @@ export interface MemoryConfig { emergencyThresholdTokens?: number; /** Min time between soft flushes in ms (default: 600000 = 10min) */ cooldownMs?: number; + /** Model ID for flush turns (default: uses conversation model) */ + flushModel?: string; }; } @@ -87,4 +89,10 @@ export interface PreparedTurn { injected: boolean; /** Pending compact level from a previously interrupted flush */ pendingCompact?: "soft" | "hard" | "emergency"; + /** Cached snapshot from pre-turn read (avoids re-reading in finalize) */ + snapshot?: MemorySnapshot; + /** Resolved file set (avoids re-resolving in finalize) */ + fileSet?: MemoryFileSet; + /** Set by caller after turn: whether agent used file-modifying tools (write/edit/bash) */ + turnUsedTools?: boolean; } diff --git a/src/types.ts b/src/types.ts index 324724b..81c7f3a 100644 --- a/src/types.ts +++ b/src/types.ts @@ -52,6 +52,12 @@ export interface AgentAdapter { /** Send a user message and return the full assistant response */ prompt(threadId: string, message: AgentMessage): Promise; + /** + * Send a prompt using a specific model (for maintenance turns like memory flush). + * Falls back to prompt() if not implemented or model unavailable. + */ + promptWithModel?(threadId: string, message: AgentMessage, modelId: string): Promise; + /** * Send a user message and stream back events in real time. * Falls back to prompt() if not implemented. From 098017e377ad18f0a4d63fbe47eb120c068d6a65 Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 09:45:11 +0000 Subject: [PATCH 2/9] perf: default flush model to Sonnet (faster) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Flush turns are structured tasks (read context → write markdown) that don't need frontier reasoning. Sonnet responds 3-5x faster than Opus for these turns. Default: amazon-bedrock/us.anthropic.claude-sonnet-4-6-20250514-v1:0 Set flushModel: null in config to fall back to conversation model. --- src/memory/lifecycle.ts | 4 +++- src/memory/types.ts | 4 ++-- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/src/memory/lifecycle.ts b/src/memory/lifecycle.ts index dfc8d5c..e3a0f61 100644 --- a/src/memory/lifecycle.ts +++ b/src/memory/lifecycle.ts @@ -191,7 +191,9 @@ export async function flushMemoryThenCompact( config?: MemoryConfig, ): Promise<{ tokensBefore: number; tokensAfter: number | null } | null> { const mode = getMode(agent); - const flushModel = config?.compact?.flushModel; + // Default to Sonnet for flush turns (faster). Set to null to use conversation model. + const DEFAULT_FLUSH_MODEL = "amazon-bedrock/us.anthropic.claude-sonnet-4-6-20250514-v1:0"; + const flushModel = config?.compact?.flushModel === null ? undefined : (config?.compact?.flushModel ?? DEFAULT_FLUSH_MODEL); /** Send flush prompt, preferring flushModel if available */ async function sendFlush(text: string): Promise { diff --git a/src/memory/types.ts b/src/memory/types.ts index 61148fd..b0dc581 100644 --- a/src/memory/types.ts +++ b/src/memory/types.ts @@ -40,8 +40,8 @@ export interface MemoryConfig { emergencyThresholdTokens?: number; /** Min time between soft flushes in ms (default: 600000 = 10min) */ cooldownMs?: number; - /** Model ID for flush turns (default: uses conversation model) */ - flushModel?: string; + /** Model ID for flush turns (default: "amazon-bedrock/us.anthropic.claude-sonnet-4-6-20250514-v1:0" — faster than conversation model) */ + flushModel?: string | null; }; } From 78ddec65c3783d0ce6131193a2ea2aac7a66216d Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 09:53:14 +0000 Subject: [PATCH 3/9] =?UTF-8?q?fix:=20address=20codex=20review=20=E2=80=94?= =?UTF-8?q?=20no=20model=20persistence,=20no=20complement=20I/O?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit P2: Replace session.setModel() with in-memory agent.state.model swap. Avoids persisting flush model to settings.json or session log. Crash-safe: if process dies mid-flush, settings remain on user's configured model. P3: Revert complement/unknown mode pre-read. finalizeMemoryForTurn skips digest checks for complement mode, so reading files was wasted I/O. --- src/agents/pi.ts | 28 ++++++++++++++++------------ src/memory/lifecycle.ts | 11 ++--------- 2 files changed, 18 insertions(+), 21 deletions(-) diff --git a/src/agents/pi.ts b/src/agents/pi.ts index 74afb20..3613da9 100644 --- a/src/agents/pi.ts +++ b/src/agents/pi.ts @@ -349,25 +349,29 @@ export const createPiAgentAdapter: AgentAdapterFactory = (config) => { return doPrompt(threadId, formatMessage(message)); } - // Switch to flush model - try { - await entry.session.setModel(targetModel); - console.log(`[pi-agent] switched to flush model: ${modelId}`); - } catch (err) { - console.warn(`[pi-agent] failed to set flush model "${modelId}":`, (err as Error).message); + // Verify auth is available for the target model + if (!modelRegistry.hasConfiguredAuth(targetModel)) { + console.warn(`[pi-agent] no auth for flush model "${modelId}", using default`); + return doPrompt(threadId, formatMessage(message)); + } + + // Swap model in-memory only (no persistence to settings.json or session log). + // This avoids a crash-window where settings could be left on the flush model. + const agentState = (entry.session as any).agent?.state; + if (!agentState) { + console.warn(`[pi-agent] cannot access agent state for model swap, using default`); return doPrompt(threadId, formatMessage(message)); } + agentState.model = targetModel; + console.log(`[pi-agent] switched to flush model (in-memory): ${modelId}`); + try { return await doPrompt(threadId, formatMessage(message)); } finally { - // Restore original model + // Restore original model (in-memory only) if (currentModel) { - try { - await entry.session.setModel(currentModel); - } catch (err) { - console.error(`[pi-agent] failed to restore model:`, (err as Error).message); - } + agentState.model = currentModel; } } }); diff --git a/src/memory/lifecycle.ts b/src/memory/lifecycle.ts index e3a0f61..aa145d9 100644 --- a/src/memory/lifecycle.ts +++ b/src/memory/lifecycle.ts @@ -52,16 +52,9 @@ export async function prepareMemoryForTurn( const mode = getMode(agent); - // Complement mode: no injection, but still read snapshot for finalize comparison + // Complement mode: no injection, no digest tracking needed (finalize skips complement) if (mode === "complement" || mode === "unknown") { - // Read snapshot so finalize can detect agent-written changes without re-reading - try { - const fileSet = resolveMemoryFiles(rootDir, config); - const snapshot = await readMemorySnapshot(fileSet, config?.inject?.maxBytes); - return { message, beforeDigest: snapshot.digest, injected: false, fileSet, snapshot }; - } catch { - return { message, beforeDigest: null, injected: false }; - } + return { message, beforeDigest: null, injected: false }; } // Full mode: inject if needed From d089aa851257c0861ba6aced9897c2a97af6b5d8 Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 09:55:49 +0000 Subject: [PATCH 4/9] fix: restore model unconditionally, track only file-modifying tools MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Always restore agentState.model in finally (even if undefined) - Only set turnUsedTools=true for write/edit/bash/multi_edit tools (read/grep/ls/find are read-only — no need to re-hash memory files) --- src/agents/pi.ts | 6 ++---- src/gateway.ts | 7 ++++--- 2 files changed, 6 insertions(+), 7 deletions(-) diff --git a/src/agents/pi.ts b/src/agents/pi.ts index 3613da9..ae9234b 100644 --- a/src/agents/pi.ts +++ b/src/agents/pi.ts @@ -369,10 +369,8 @@ export const createPiAgentAdapter: AgentAdapterFactory = (config) => { try { return await doPrompt(threadId, formatMessage(message)); } finally { - // Restore original model (in-memory only) - if (currentModel) { - agentState.model = currentModel; - } + // Restore original model (in-memory only) — even if undefined + agentState.model = currentModel; } }); }, diff --git a/src/gateway.ts b/src/gateway.ts index c37f68b..c8745b2 100644 --- a/src/gateway.ts +++ b/src/gateway.ts @@ -933,7 +933,8 @@ export class Gateway { */ private async handleStreaming(thread: any, stream: AsyncIterable, verbose: boolean, signal?: AbortSignal): Promise<{ usedTools: boolean }> { let activeTools = new Map(); // toolCallId -> toolName - let usedTools = false; + let usedFileModifyingTools = false; + const FILE_MODIFYING_TOOLS = new Set(["write", "edit", "bash", "multi_edit"]); // Per-turn streaming state — each turn gets a fresh iterable + promise let currentPush: ((text: string) => void) | null = null; @@ -1037,7 +1038,7 @@ export class Gateway { case "tool_start": { activeTools.set(event.toolCallId, event.toolName); - usedTools = true; + if (FILE_MODIFYING_TOOLS.has(event.toolName)) usedFileModifyingTools = true; if (verbose) { try { await thread.post(`${toolIcon(event.toolName)} Running \`${event.toolName}\`…`); @@ -1109,7 +1110,7 @@ export class Gateway { await flushCurrentStream(); } - return { usedTools }; + return { usedTools: usedFileModifyingTools }; } /** Post text with markdown, falling back to plain text */ From 2c88e6b686ec100abf63cd1c4923c6ddb639598b Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 10:05:31 +0000 Subject: [PATCH 5/9] =?UTF-8?q?fix:=20invert=20tool=20allowlist=20?= =?UTF-8?q?=E2=80=94=20use=20read-only=20set=20instead=20of=20write=20set?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Unknown/extension tools now trigger memory re-read (safe default). Only known read-only tools (read, grep, find, ls, glob) skip re-read. Addresses codex review: custom extension tools that modify memory files would have been missed by the previous write-tool allowlist. --- src/gateway.ts | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/gateway.ts b/src/gateway.ts index c8745b2..00f688a 100644 --- a/src/gateway.ts +++ b/src/gateway.ts @@ -934,7 +934,8 @@ export class Gateway { private async handleStreaming(thread: any, stream: AsyncIterable, verbose: boolean, signal?: AbortSignal): Promise<{ usedTools: boolean }> { let activeTools = new Map(); // toolCallId -> toolName let usedFileModifyingTools = false; - const FILE_MODIFYING_TOOLS = new Set(["write", "edit", "bash", "multi_edit"]); + // Read-only tools that cannot modify memory files — everything else triggers re-read + const READ_ONLY_TOOLS = new Set(["read", "grep", "find", "ls", "glob"]); // Per-turn streaming state — each turn gets a fresh iterable + promise let currentPush: ((text: string) => void) | null = null; @@ -1038,7 +1039,7 @@ export class Gateway { case "tool_start": { activeTools.set(event.toolCallId, event.toolName); - if (FILE_MODIFYING_TOOLS.has(event.toolName)) usedFileModifyingTools = true; + if (!READ_ONLY_TOOLS.has(event.toolName)) usedFileModifyingTools = true; if (verbose) { try { await thread.post(`${toolIcon(event.toolName)} Running \`${event.toolName}\`…`); From 5b07d9f4912493e570a691294b1020d295325371 Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 10:23:04 +0000 Subject: [PATCH 6/9] feat: add timing info to compact messages + persistent timing log - Telegram messages now show flush/compact/total time + model used - Timing persisted to ~/.roundhouse/logs/compact-timing.jsonl (JSONL) - Both manual /compact and auto-compact show timing --- src/gateway.ts | 8 ++++++-- src/memory/lifecycle.ts | 43 ++++++++++++++++++++++++++++++++++++----- src/memory/types.ts | 2 +- 3 files changed, 45 insertions(+), 8 deletions(-) diff --git a/src/gateway.ts b/src/gateway.ts index 00f688a..efb2db9 100644 --- a/src/gateway.ts +++ b/src/gateway.ts @@ -513,7 +513,9 @@ export class Gateway { await thread.post("⚠️ No active session to compact. Send a message first."); } else { const beforeK = (result.tokensBefore / 1000).toFixed(1); - await thread.post(`✅ Memory saved & compacted\n\nCompacted ${beforeK}K tokens down to a summary.\nContext usage will update after your next message.`); + const timing = result.timing; + const timingLine = timing ? `\nTiming: flush ${(timing.flushMs / 1000).toFixed(1)}s, compact ${(timing.compactMs / 1000).toFixed(1)}s, total ${(timing.totalMs / 1000).toFixed(1)}s\nModel: ${timing.model}` : ""; + await thread.post(`✅ Memory saved & compacted\n\nCompacted ${beforeK}K tokens down to a summary.\nContext usage will update after your next message.${timingLine}`); } } } catch (err) { @@ -915,7 +917,9 @@ export class Gateway { const result = await flushMemoryThenCompact(agentThreadId, agent, memoryRoot, pressure, this.config.memory); if (result) { const beforeK = (result.tokensBefore / 1000).toFixed(1); - await thread.post(`✅ Auto-compacted: ${beforeK}K tokens → summary.`); + const timing = result.timing; + const timingLine = timing ? ` (${(timing.totalMs / 1000).toFixed(1)}s: flush ${(timing.flushMs / 1000).toFixed(1)}s + compact ${(timing.compactMs / 1000).toFixed(1)}s)` : ""; + await thread.post(`✅ Auto-compacted: ${beforeK}K tokens → summary.${timingLine}`); } } catch (err) { console.error(`[roundhouse] ${pressure} compact error:`, (err as Error).message); diff --git a/src/memory/lifecycle.ts b/src/memory/lifecycle.ts index aa145d9..d971688 100644 --- a/src/memory/lifecycle.ts +++ b/src/memory/lifecycle.ts @@ -16,6 +16,9 @@ import { shouldInjectMemory, classifyContextPressure, isSoftFlushOnCooldown } fr import { buildMemoryInjection, injectMemoryIntoMessage } from "./inject"; import { buildFlushPrompt } from "./prompts"; import { bootstrapMemoryFiles } from "./bootstrap"; +import { appendFile, mkdir } from "node:fs/promises"; +import { join } from "node:path"; +import { homedir } from "node:os"; // ── Memory mode detection ──────────────────────────── @@ -176,16 +179,23 @@ export async function finalizeMemoryForTurn( * * Returns compaction result or null if nothing to compact. */ +export interface CompactTiming { + flushMs: number; + compactMs: number; + totalMs: number; + model: string; +} + export async function flushMemoryThenCompact( threadId: string, agent: AgentAdapter, rootDir: string, level: "soft" | "hard" | "emergency" | "manual", config?: MemoryConfig, -): Promise<{ tokensBefore: number; tokensAfter: number | null } | null> { +): Promise<{ tokensBefore: number; tokensAfter: number | null; timing?: CompactTiming } | null> { const mode = getMode(agent); // Default to Sonnet for flush turns (faster). Set to null to use conversation model. - const DEFAULT_FLUSH_MODEL = "amazon-bedrock/us.anthropic.claude-sonnet-4-6-20250514-v1:0"; + const DEFAULT_FLUSH_MODEL = "amazon-bedrock/us.anthropic.claude-haiku-4-5-20251001-v1:0"; const flushModel = config?.compact?.flushModel === null ? undefined : (config?.compact?.flushModel ?? DEFAULT_FLUSH_MODEL); /** Send flush prompt, preferring flushModel if available */ @@ -221,16 +231,20 @@ export async function flushMemoryThenCompact( if (!agent.compact) return null; const effectiveLevel = level === "manual" ? "hard" : level; + const t0 = Date.now(); try { // Step 1: flush const flushText = buildFlushPrompt(mode === "unknown" ? "full" : mode, effectiveLevel); console.log(`[memory] flushing memory for ${threadId} (level: ${level}${flushModel ? `, model: ${flushModel}` : ""})`); await sendFlush(flushText); + const flushMs = Date.now() - t0; // Step 2: compact - console.log(`[memory] compacting ${threadId}`); + console.log(`[memory] compacting ${threadId} (flush took ${flushMs}ms)`); + const t1 = Date.now(); const result = await agent.compact(threadId); + const compactMs = Date.now() - t1; if (!result) return null; // Step 3: mark force re-inject (Full mode only) @@ -242,8 +256,27 @@ export async function flushMemoryThenCompact( await saveThreadMemoryState(threadId, state); } - console.log(`[memory] flush+compact done for ${threadId}: ${result.tokensBefore} → ${result.tokensAfter ?? "?"} tokens`); - return result; + const totalMs = Date.now() - t0; + const timing = { flushMs, compactMs, totalMs, model: flushModel ?? "default" }; + console.log(`[memory] flush+compact done for ${threadId}: ${result.tokensBefore} → ${result.tokensAfter ?? "?"} tokens | flush=${flushMs}ms compact=${compactMs}ms total=${totalMs}ms model=${timing.model}`); + + // Persist timing log for debugging (async, fire-and-forget) + const logDir = join(homedir(), ".roundhouse", "logs"); + mkdir(logDir, { recursive: true }) + .then(() => { + const entry = JSON.stringify({ + ts: new Date().toISOString(), + threadId, + level, + tokensBefore: result.tokensBefore, + tokensAfter: result.tokensAfter, + ...timing, + }); + return appendFile(join(logDir, "compact-timing.jsonl"), entry + "\n"); + }) + .catch((err) => console.warn(`[memory] timing log write failed:`, (err as Error).message)); + + return { ...result, timing }; } catch (err) { console.error(`[memory] flush+compact failed for ${threadId}:`, (err as Error).message); // Mark pending so we retry on next turn diff --git a/src/memory/types.ts b/src/memory/types.ts index b0dc581..bcac487 100644 --- a/src/memory/types.ts +++ b/src/memory/types.ts @@ -40,7 +40,7 @@ export interface MemoryConfig { emergencyThresholdTokens?: number; /** Min time between soft flushes in ms (default: 600000 = 10min) */ cooldownMs?: number; - /** Model ID for flush turns (default: "amazon-bedrock/us.anthropic.claude-sonnet-4-6-20250514-v1:0" — faster than conversation model) */ + /** Model ID for flush turns (default: "amazon-bedrock/us.anthropic.claude-haiku-4-5-20251001-v1:0" — fast, matches Sonnet quality for structured writes) */ flushModel?: string | null; }; } From 51c9975789ac0ba5eb7e0ebe5d99b6b41fa771d9 Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 10:37:39 +0000 Subject: [PATCH 7/9] refactor: extract READ_ONLY_TOOLS + CompactResult to types, add tests - READ_ONLY_TOOLS: shared constant in memory/types.ts (not buried in handleStreaming) - CompactTiming + CompactResult: proper types in types.ts (SRP) - lifecycle.ts imports CompactResult instead of inline type - gateway.ts imports READ_ONLY_TOOLS + CompactResult from types - Tests: READ_ONLY_TOOLS membership, CompactResult with/without timing --- src/gateway.ts | 5 ++--- src/memory/lifecycle.ts | 11 ++--------- src/memory/types.ts | 30 ++++++++++++++++++++++++++++ test/memory.test.ts | 43 +++++++++++++++++++++++++++++++++++++++++ 4 files changed, 77 insertions(+), 12 deletions(-) diff --git a/src/gateway.ts b/src/gateway.ts index efb2db9..3773675 100644 --- a/src/gateway.ts +++ b/src/gateway.ts @@ -20,7 +20,8 @@ import { formatSchedule, formatRunCounts, jobEnabledIcon } from "./cron/format"; import { BOT_COMMANDS } from "./commands"; import { prepareMemoryForTurn, finalizeMemoryForTurn, flushMemoryThenCompact, determineMemoryMode } from "./memory/lifecycle"; import { maxPressure } from "./memory/policy"; -import type { PressureLevel } from "./memory/types"; +import type { PressureLevel, CompactResult } from "./memory/types"; +import { READ_ONLY_TOOLS } from "./memory/types"; import { readPendingPairing, completePendingPairing, isStartForNonce } from "./pairing"; /** Match a Telegram command, handling optional @botname suffix */ @@ -938,8 +939,6 @@ export class Gateway { private async handleStreaming(thread: any, stream: AsyncIterable, verbose: boolean, signal?: AbortSignal): Promise<{ usedTools: boolean }> { let activeTools = new Map(); // toolCallId -> toolName let usedFileModifyingTools = false; - // Read-only tools that cannot modify memory files — everything else triggers re-read - const READ_ONLY_TOOLS = new Set(["read", "grep", "find", "ls", "glob"]); // Per-turn streaming state — each turn gets a fresh iterable + promise let currentPush: ((text: string) => void) | null = null; diff --git a/src/memory/lifecycle.ts b/src/memory/lifecycle.ts index d971688..45a60a2 100644 --- a/src/memory/lifecycle.ts +++ b/src/memory/lifecycle.ts @@ -9,7 +9,7 @@ */ import type { AgentAdapter, AgentMessage } from "../types"; -import type { MemoryConfig, MemoryFileSet, MemoryMode, MemorySnapshot, PreparedTurn, PressureLevel, ThreadMemoryState } from "./types"; +import type { MemoryConfig, MemoryFileSet, MemoryMode, MemorySnapshot, PreparedTurn, PressureLevel, ThreadMemoryState, CompactResult } from "./types"; import { resolveMemoryFiles, readMemorySnapshot, formatDate } from "./files"; import { loadThreadMemoryState, saveThreadMemoryState } from "./state"; import { shouldInjectMemory, classifyContextPressure, isSoftFlushOnCooldown } from "./policy"; @@ -179,20 +179,13 @@ export async function finalizeMemoryForTurn( * * Returns compaction result or null if nothing to compact. */ -export interface CompactTiming { - flushMs: number; - compactMs: number; - totalMs: number; - model: string; -} - export async function flushMemoryThenCompact( threadId: string, agent: AgentAdapter, rootDir: string, level: "soft" | "hard" | "emergency" | "manual", config?: MemoryConfig, -): Promise<{ tokensBefore: number; tokensAfter: number | null; timing?: CompactTiming } | null> { +): Promise { const mode = getMode(agent); // Default to Sonnet for flush turns (faster). Set to null to use conversation model. const DEFAULT_FLUSH_MODEL = "amazon-bedrock/us.anthropic.claude-haiku-4-5-20251001-v1:0"; diff --git a/src/memory/types.ts b/src/memory/types.ts index bcac487..3dfc561 100644 --- a/src/memory/types.ts +++ b/src/memory/types.ts @@ -96,3 +96,33 @@ export interface PreparedTurn { /** Set by caller after turn: whether agent used file-modifying tools (write/edit/bash) */ turnUsedTools?: boolean; } + +// ── Tool classification ────────────────────────────── + +/** + * Tools known to be read-only (cannot modify files on disk). + * Any tool NOT in this set is assumed to potentially modify files, + * triggering a memory digest re-read after the turn. + */ +export const READ_ONLY_TOOLS: ReadonlySet = new Set([ + "read", + "grep", + "find", + "ls", + "glob", +]); + +// ── Compact timing ─────────────────────────────── + +export interface CompactTiming { + flushMs: number; + compactMs: number; + totalMs: number; + model: string; +} + +export interface CompactResult { + tokensBefore: number; + tokensAfter: number | null; + timing?: CompactTiming; +} diff --git a/test/memory.test.ts b/test/memory.test.ts index c142414..7e6bd94 100644 --- a/test/memory.test.ts +++ b/test/memory.test.ts @@ -213,3 +213,46 @@ describe("buildFlushPrompt", () => { expect(result).not.toContain("Context is filling"); }); }); + +// ── READ_ONLY_TOOLS ───────────────────────────────── + +import { READ_ONLY_TOOLS } from "../src/memory/types"; +import type { CompactResult, CompactTiming } from "../src/memory/types"; + +describe("READ_ONLY_TOOLS", () => { + it("contains expected read-only tools", () => { + expect(READ_ONLY_TOOLS.has("read")).toBe(true); + expect(READ_ONLY_TOOLS.has("grep")).toBe(true); + expect(READ_ONLY_TOOLS.has("find")).toBe(true); + expect(READ_ONLY_TOOLS.has("ls")).toBe(true); + expect(READ_ONLY_TOOLS.has("glob")).toBe(true); + }); + + it("does NOT contain file-modifying tools", () => { + expect(READ_ONLY_TOOLS.has("write")).toBe(false); + expect(READ_ONLY_TOOLS.has("edit")).toBe(false); + expect(READ_ONLY_TOOLS.has("bash")).toBe(false); + expect(READ_ONLY_TOOLS.has("multi_edit")).toBe(false); + }); + + it("does NOT contain unknown/extension tools (safe default: assume writing)", () => { + expect(READ_ONLY_TOOLS.has("my_custom_tool")).toBe(false); + expect(READ_ONLY_TOOLS.has("")).toBe(false); + }); +}); + +// ── CompactResult type ─────────────────────────────── + +describe("CompactResult", () => { + it("allows result without timing (backwards compat)", () => { + const result: CompactResult = { tokensBefore: 80000, tokensAfter: 5000 }; + expect(result.timing).toBeUndefined(); + }); + + it("allows result with timing", () => { + const timing: CompactTiming = { flushMs: 3000, compactMs: 5000, totalMs: 8000, model: "haiku" }; + const result: CompactResult = { tokensBefore: 80000, tokensAfter: 5000, timing }; + expect(result.timing!.flushMs).toBe(3000); + expect(result.timing!.model).toBe("haiku"); + }); +}); From 3d5dc6963f66b00100558454b43bc041fda1afb2 Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 10:58:45 +0000 Subject: [PATCH 8/9] feat: use Haiku for compact step too (compactWithModel) Before: flush used Haiku (238s) but compact restored to Opus (152s) = 390s total After: both flush AND compact use Haiku via new compactWithModel adapter method - Added compactWithModel?(threadId, modelId) to AgentAdapter interface - Implemented in pi.ts: in-memory model swap around session.compact() - lifecycle.ts: prefers compactWithModel when available, falls back to compact() --- src/agents/pi.ts | 43 +++++++++++++++++++++++++++++++++++++++++ src/memory/lifecycle.ts | 6 ++++-- src/types.ts | 2 ++ 3 files changed, 49 insertions(+), 2 deletions(-) diff --git a/src/agents/pi.ts b/src/agents/pi.ts index ae9234b..42bc453 100644 --- a/src/agents/pi.ts +++ b/src/agents/pi.ts @@ -509,6 +509,49 @@ export const createPiAgentAdapter: AgentAdapterFactory = (config) => { }); }, + async compactWithModel(threadId: string, modelId: string): Promise<{ tokensBefore: number; tokensAfter: number | null } | null> { + return enqueue(threadId, async () => { + const entry = sessions.get(threadId); + if (!entry) return null; + + const agentState = (entry.session as any).agent?.state; + let currentModel: any; + let modelSwapped = false; + + // Resolve and swap model for compact + if (!agentState) { + console.warn(`[pi-agent] cannot access agent state for compact model swap, using default`); + } else { + const [provider, ...rest] = modelId.split("/"); + const id = rest.join("/"); + const targetModel = (provider && id) ? modelRegistry.find(provider, id) : null; + if (!targetModel) { + console.warn(`[pi-agent] compact model "${modelId}" not found, using default`); + } else if (!modelRegistry.hasConfiguredAuth(targetModel)) { + console.warn(`[pi-agent] no auth for compact model "${modelId}", using default`); + } else { + currentModel = agentState.model; + agentState.model = targetModel; + modelSwapped = true; + console.log(`[pi-agent] compact using model (in-memory): ${modelId}`); + } + } + + try { + const result = await entry.session.compact(); + const usage = entry.session.getContextUsage(); + return { + tokensBefore: result.tokensBefore, + tokensAfter: usage?.tokens ?? null, + }; + } finally { + if (modelSwapped) { + agentState.model = currentModel; + } + } + }); + }, + async abort(threadId: string): Promise { const entry = sessions.get(threadId); if (entry) { diff --git a/src/memory/lifecycle.ts b/src/memory/lifecycle.ts index 45a60a2..674ac26 100644 --- a/src/memory/lifecycle.ts +++ b/src/memory/lifecycle.ts @@ -233,10 +233,12 @@ export async function flushMemoryThenCompact( await sendFlush(flushText); const flushMs = Date.now() - t0; - // Step 2: compact + // Step 2: compact (use flush model if compactWithModel is available) console.log(`[memory] compacting ${threadId} (flush took ${flushMs}ms)`); const t1 = Date.now(); - const result = await agent.compact(threadId); + const result = flushModel && agent.compactWithModel + ? await agent.compactWithModel(threadId, flushModel) + : await agent.compact!(threadId); const compactMs = Date.now() - t1; if (!result) return null; diff --git a/src/types.ts b/src/types.ts index 81c7f3a..a1e9c97 100644 --- a/src/types.ts +++ b/src/types.ts @@ -69,6 +69,8 @@ export interface AgentAdapter { /** Compact the session context for a thread */ compact?(threadId: string): Promise<{ tokensBefore: number; tokensAfter: number | null } | null>; + /** Compact with a specific model (avoids restoring to default between flush and compact) */ + compactWithModel?(threadId: string, modelId: string): Promise<{ tokensBefore: number; tokensAfter: number | null } | null>; /** Abort the current agent run for a thread */ abort?(threadId: string): Promise; From 78189ad064994f67fcb199faee58e8bc90343077 Mon Sep 17 00:00:00 2001 From: Loki FastStart Date: Thu, 7 May 2026 11:09:13 +0000 Subject: [PATCH 9/9] feat: live progress updates during compact (edits Telegram message in-place) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - New telegram-progress.ts: createProgressMessage() → editable message handle - Compact shows: 'Flushing memory...' → 'Compacting... (flush took Xs)' → final result - Single message edited in-place (no spam of multiple messages) - Works for both manual /compact and auto-compact - Falls back gracefully for non-Telegram adapters (no-op updates) --- src/gateway.ts | 28 +++++++++------ src/memory/lifecycle.ts | 3 ++ src/telegram-progress.ts | 76 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 97 insertions(+), 10 deletions(-) create mode 100644 src/telegram-progress.ts diff --git a/src/gateway.ts b/src/gateway.ts index 3773675..dc33476 100644 --- a/src/gateway.ts +++ b/src/gateway.ts @@ -23,6 +23,7 @@ import { maxPressure } from "./memory/policy"; import type { PressureLevel, CompactResult } from "./memory/types"; import { READ_ONLY_TOOLS } from "./memory/types"; import { readPendingPairing, completePendingPairing, isStartForNonce } from "./pairing"; +import { createProgressMessage } from "./telegram-progress"; /** Match a Telegram command, handling optional @botname suffix */ /** Bot username for command suffix validation (set during gateway init) */ @@ -494,7 +495,7 @@ export class Gateway { threadLocks.set(agentThreadId, lockPromise); if (prevLock) await prevLock; - await thread.post("📝 Saving memory and compacting..."); + const progress = await createProgressMessage(thread, "📝 Saving memory and compacting..."); const stopTyping = startTypingLoop(thread); try { const agentCwd = (agent.getInfo?.()?.cwd as string) ?? process.cwd(); @@ -503,25 +504,28 @@ export class Gateway { if (this.config.memory?.enabled === false) { const result = await agent.compact(agentThreadId); if (!result) { - await thread.post("⚠️ No active session to compact. Send a message first."); + await progress.update("⚠️ No active session to compact. Send a message first."); } else { const beforeK = (result.tokensBefore / 1000).toFixed(1); - await thread.post(`✅ Compaction complete\n\nCompacted ${beforeK}K tokens down to a summary.\nContext usage will update after your next message.`); + await progress.update(`✅ Compaction complete\n\nCompacted ${beforeK}K tokens down to a summary.\nContext usage will update after your next message.`); } } else { - const result = await flushMemoryThenCompact(agentThreadId, agent, memoryRoot, "manual", this.config.memory); + const result = await flushMemoryThenCompact( + agentThreadId, agent, memoryRoot, "manual", this.config.memory, + (step) => progress.update(step), + ); if (!result) { - await thread.post("⚠️ No active session to compact. Send a message first."); + await progress.update("⚠️ No active session to compact. Send a message first."); } else { const beforeK = (result.tokensBefore / 1000).toFixed(1); const timing = result.timing; const timingLine = timing ? `\nTiming: flush ${(timing.flushMs / 1000).toFixed(1)}s, compact ${(timing.compactMs / 1000).toFixed(1)}s, total ${(timing.totalMs / 1000).toFixed(1)}s\nModel: ${timing.model}` : ""; - await thread.post(`✅ Memory saved & compacted\n\nCompacted ${beforeK}K tokens down to a summary.\nContext usage will update after your next message.${timingLine}`); + await progress.update(`✅ Memory saved & compacted\n\nCompacted ${beforeK}K tokens down to a summary.\nContext usage will update after your next message.${timingLine}`); } } } catch (err) { const msg = err instanceof Error ? err.message : String(err); - await thread.post(`⚠️ Compaction failed: ${msg.slice(0, 200)}`); + await progress.update(`⚠️ Compaction failed: ${msg.slice(0, 200)}`); } finally { stopTyping(); releaseLock!(); @@ -914,13 +918,17 @@ export class Gateway { // Hard or emergency: flush + compact try { - await thread.post(`📝 ${pressure === "emergency" ? "⚠️ Context nearly full! " : ""}Saving memory and compacting...`); - const result = await flushMemoryThenCompact(agentThreadId, agent, memoryRoot, pressure, this.config.memory); + const prefix = pressure === "emergency" ? "⚠️ Context nearly full! " : ""; + const progress = await createProgressMessage(thread, `📝 ${prefix}Saving memory and compacting...`); + const result = await flushMemoryThenCompact( + agentThreadId, agent, memoryRoot, pressure, this.config.memory, + (step) => progress.update(step), + ); if (result) { const beforeK = (result.tokensBefore / 1000).toFixed(1); const timing = result.timing; const timingLine = timing ? ` (${(timing.totalMs / 1000).toFixed(1)}s: flush ${(timing.flushMs / 1000).toFixed(1)}s + compact ${(timing.compactMs / 1000).toFixed(1)}s)` : ""; - await thread.post(`✅ Auto-compacted: ${beforeK}K tokens → summary.${timingLine}`); + await progress.update(`✅ Auto-compacted: ${beforeK}K tokens → summary.${timingLine}`); } } catch (err) { console.error(`[roundhouse] ${pressure} compact error:`, (err as Error).message); diff --git a/src/memory/lifecycle.ts b/src/memory/lifecycle.ts index 674ac26..8f1c797 100644 --- a/src/memory/lifecycle.ts +++ b/src/memory/lifecycle.ts @@ -185,6 +185,7 @@ export async function flushMemoryThenCompact( rootDir: string, level: "soft" | "hard" | "emergency" | "manual", config?: MemoryConfig, + onProgress?: (step: string) => void | Promise, ): Promise { const mode = getMode(agent); // Default to Sonnet for flush turns (faster). Set to null to use conversation model. @@ -230,11 +231,13 @@ export async function flushMemoryThenCompact( // Step 1: flush const flushText = buildFlushPrompt(mode === "unknown" ? "full" : mode, effectiveLevel); console.log(`[memory] flushing memory for ${threadId} (level: ${level}${flushModel ? `, model: ${flushModel}` : ""})`); + await onProgress?.("💭 Flushing memory..."); await sendFlush(flushText); const flushMs = Date.now() - t0; // Step 2: compact (use flush model if compactWithModel is available) console.log(`[memory] compacting ${threadId} (flush took ${flushMs}ms)`); + await onProgress?.(`✂️ Compacting context... (flush took ${(flushMs / 1000).toFixed(1)}s)`); const t1 = Date.now(); const result = flushModel && agent.compactWithModel ? await agent.compactWithModel(threadId, flushModel) diff --git a/src/telegram-progress.ts b/src/telegram-progress.ts new file mode 100644 index 0000000..ccaa1c2 --- /dev/null +++ b/src/telegram-progress.ts @@ -0,0 +1,76 @@ +/** + * telegram-progress.ts — Editable progress messages for long-running operations + */ + +/** Parse Telegram chat_id and optional message_thread_id from a Chat SDK thread ID */ +function parseTelegramThreadId(threadId: string): { chatId: string; messageThreadId?: number } { + const parts = threadId.split(":"); + const chatId = parts[1]; + const topicPart = parts[2]; + const result: { chatId: string; messageThreadId?: number } = { chatId }; + if (topicPart) { + const parsed = parseInt(topicPart, 10); + if (Number.isFinite(parsed)) result.messageThreadId = parsed; + } + return result; +} + +export interface ProgressMessage { + /** Update the message text (edits in place) */ + update(text: string): Promise; +} + +/** + * Send an initial message and return a handle to edit it in-place. + * Falls back to no-op if the thread isn't Telegram or the send fails. + */ +export async function createProgressMessage(thread: any, initialText: string): Promise { + const isTelegram = + typeof thread?.adapter?.telegramFetch === "function" && + typeof thread?.id === "string" && + thread.id.startsWith("telegram:"); + + if (!isTelegram) { + // Non-Telegram: just post once, updates are no-ops + await thread.post(initialText); + return { update: async () => {} }; + } + + const { chatId, messageThreadId } = parseTelegramThreadId(thread.id); + const basePayload = { + chat_id: chatId, + ...(messageThreadId !== undefined && { message_thread_id: messageThreadId }), + disable_web_page_preview: true, + }; + + let messageId: number | null = null; + let lastText = ""; + + try { + const result = await thread.adapter.telegramFetch("sendMessage", { + ...basePayload, + text: initialText, + }); + messageId = result.message_id; + lastText = initialText; + } catch { + // Fallback: use thread.post (can't edit later) + await thread.post(initialText); + } + + return { + async update(text: string) { + if (!messageId || text === lastText) return; + try { + await thread.adapter.telegramFetch("editMessageText", { + ...basePayload, + message_id: messageId, + text, + }); + lastText = text; + } catch { + // Edit failed (rate limit, message deleted, etc.) — skip silently + } + }, + }; +}