diff --git a/README.md b/README.md index 779aa1a06..87234eb2f 100644 --- a/README.md +++ b/README.md @@ -75,6 +75,35 @@ Full IDE-like code intelligence with automatic formatting and diagnostics: - **Local binary resolution**: Auto-discovers project-local LSP servers in `node_modules/.bin/`, `.venv/bin/`, etc. - Hover docs, symbol references, code actions, workspace-wide symbol search +## + Time Traveling Streamed Rules (TTSR) + +

+ ttsr +

+ +Zero context-use rules that inject themselves only when needed: + +- **Pattern-triggered injection**: Rules define regex triggers that watch the model's output stream +- **Just-in-time activation**: When a pattern matches, the stream aborts, the rule injects as a system reminder, and the request retries +- **Zero upfront cost**: TTSR rules consume no context until they're actually relevant +- **One-shot per session**: Each rule only triggers once, preventing loops +- Define via `ttsrTrigger` field in rule files (regex pattern) + +Example: A "don't use deprecated API" rule only activates when the model starts writing deprecated code, saving context for sessions that never touch that API. + +## + Interactive Code Review + +

+ review +

+ +Structured code review with priority-based findings: + +- **`/review` command**: Interactive mode selection (branch comparison, uncommitted changes, commit review) +- **Structured findings**: `report_finding` tool with priority levels (P0-P3: critical → nit) +- **Verdict rendering**: `submit_review` aggregates findings into approve/request-changes/comment +- Combined result tree showing verdict and all findings + ## + Task Tool (Subagent System)

@@ -115,19 +144,6 @@ Structured user interaction with typed options: - **Multiple choice questions**: Present options with descriptions for user selection - **Multi-select support**: Allow multiple answers when choices aren't mutually exclusive -## + Interactive Code Review - -

- review -

- -Structured code review with priority-based findings: - -- **`/review` command**: Interactive mode selection (branch comparison, uncommitted changes, commit review) -- **Structured findings**: `report_finding` tool with priority levels (P0-P3: critical → nit) -- **Verdict rendering**: `submit_review` aggregates findings into approve/request-changes/comment -- Combined result tree showing verdict and all findings - ## + Custom TypeScript Slash Commands

diff --git a/assets/ttsr.webp b/assets/ttsr.webp new file mode 100644 index 000000000..25cbe9f27 Binary files /dev/null and b/assets/ttsr.webp differ diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index 5dbbb3012..e96feeb04 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -1,6 +1,14 @@ # Changelog ## [Unreleased] +### Added + +- Added `popMessage()` method to Agent class for removing and retrieving the last message +- Added abort signal checks during response streaming for faster interruption handling + +### Fixed + +- Fixed abort handling to properly return aborted message state when stream is interrupted mid-response ## [3.3.1337] - 2026-01-03 diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index acd6559aa..7153ee86a 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -225,6 +225,11 @@ async function streamAssistantResponse( let addedPartial = false; for await (const event of response) { + // Check abort early - allows TTSR and other abort sources to break immediately + if (signal?.aborted) { + break; + } + switch (event.type) { case "start": partialMessage = event.partial; @@ -268,6 +273,22 @@ async function streamAssistantResponse( return finalMessage; } } + + // Check abort after processing - allows handlers to abort mid-stream + if (signal?.aborted) { + break; + } + } + + // If we broke out due to abort, return an aborted message + if (signal?.aborted && partialMessage) { + const abortedMessage: AssistantMessage = { + ...partialMessage, + stopReason: "aborted", + }; + context.messages[context.messages.length - 1] = abortedMessage; + stream.push({ type: "message_end", message: abortedMessage }); + return abortedMessage; } return await response.result(); diff --git a/packages/agent/src/agent.ts b/packages/agent/src/agent.ts index db10b3015..0bbeaddd5 100644 --- a/packages/agent/src/agent.ts +++ b/packages/agent/src/agent.ts @@ -174,6 +174,14 @@ export class Agent { this._state.messages = []; } + /** Remove and return the last message from the message list */ + popMessage(): AgentMessage | undefined { + if (this._state.messages.length === 0) return undefined; + const popped = this._state.messages[this._state.messages.length - 1]; + this._state.messages = this._state.messages.slice(0, -1); + return popped; + } + abort() { this.abortController?.abort(); } diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 17f92c741..7193b01de 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,6 +1,13 @@ # Changelog ## [Unreleased] + +### Added + +- Added Time Traveling Stream Rules (TTSR) feature that monitors agent output for pattern matches and injects rule reminders mid-stream +- Added `ttsr_trigger` frontmatter field for rules to define regex patterns that trigger mid-stream injection +- Added TTSR settings for enabled state, context mode (keep/discard partial output), and repeat mode (once/after-gap) + ### Fixed - Fixed excessive subprocess spawns by caching git status for 1 second in the footer component diff --git a/packages/coding-agent/src/capability/rule.ts b/packages/coding-agent/src/capability/rule.ts index e980f6f49..476a2f99f 100644 --- a/packages/coding-agent/src/capability/rule.ts +++ b/packages/coding-agent/src/capability/rule.ts @@ -15,6 +15,8 @@ export interface RuleFrontmatter { description?: string; globs?: string[]; alwaysApply?: boolean; + /** Regex pattern that triggers time-traveling rule injection */ + ttsr_trigger?: string; [key: string]: unknown; } @@ -34,6 +36,8 @@ export interface Rule { alwaysApply?: boolean; /** Description (for agent-requested rules) */ description?: string; + /** Regex pattern that triggers time-traveling rule injection */ + ttsrTrigger?: string; /** Source metadata */ _source: SourceMeta; } diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index 04335ba78..593df6192 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -16,6 +16,7 @@ import type { Agent, AgentEvent, AgentMessage, AgentState, ThinkingLevel } from "@oh-my-pi/pi-agent-core"; import type { AssistantMessage, ImageContent, Message, Model, TextContent } from "@oh-my-pi/pi-ai"; import { isContextOverflow, modelsAreEqual, supportsXhigh } from "@oh-my-pi/pi-ai"; +import type { Rule } from "../capability/rule"; import { getAuthPath } from "../config"; import { type BashResult, executeBash as executeBashCommand } from "./bash-executor"; import { @@ -47,6 +48,7 @@ import type { ModelRegistry } from "./model-registry"; import type { BranchSummaryEntry, CompactionEntry, NewSessionOptions, SessionManager } from "./session-manager"; import type { SettingsManager, SkillsSettings } from "./settings-manager"; import { expandSlashCommand, type FileSlashCommand, parseCommandArgs } from "./slash-commands"; +import type { TtsrManager } from "./ttsr"; /** Session-specific events that extend the core AgentEvent */ export type AgentSessionEvent = @@ -54,7 +56,8 @@ export type AgentSessionEvent = | { type: "auto_compaction_start"; reason: "threshold" | "overflow" } | { type: "auto_compaction_end"; result: CompactionResult | undefined; aborted: boolean; willRetry: boolean } | { type: "auto_retry_start"; attempt: number; maxAttempts: number; delayMs: number; errorMessage: string } - | { type: "auto_retry_end"; success: boolean; attempt: number; finalError?: string }; + | { type: "auto_retry_end"; success: boolean; attempt: number; finalError?: string } + | { type: "ttsr_triggered"; rules: Rule[] }; /** Listener function for agent session events */ export type AgentSessionEventListener = (event: AgentSessionEvent) => void; @@ -80,6 +83,8 @@ export interface AgentSessionConfig { skillsSettings?: Required; /** Model registry for API key resolution and model discovery */ modelRegistry: ModelRegistry; + /** TTSR manager for time-traveling stream rules */ + ttsrManager?: TtsrManager; } /** Options for AgentSession.prompt() */ @@ -179,6 +184,11 @@ export class AgentSession { // Model registry for API key resolution private _modelRegistry: ModelRegistry; + // TTSR manager for time-traveling stream rules + private _ttsrManager: TtsrManager | undefined = undefined; + private _pendingTtsrInjections: Rule[] = []; + private _ttsrAbortPending = false; + constructor(config: AgentSessionConfig) { this.agent = config.agent; this.sessionManager = config.sessionManager; @@ -190,6 +200,7 @@ export class AgentSession { this._customCommands = config.customCommands ?? []; this._skillsSettings = config.skillsSettings; this._modelRegistry = config.modelRegistry; + this._ttsrManager = config.ttsrManager; // Always subscribe to agent events for internal handling // (session persistence, hooks, auto-compaction, retry logic) @@ -201,6 +212,16 @@ export class AgentSession { return this._modelRegistry; } + /** TTSR manager for time-traveling stream rules */ + get ttsrManager(): TtsrManager | undefined { + return this._ttsrManager; + } + + /** Whether a TTSR abort is pending (stream was aborted to inject rules) */ + get isTtsrAbortPending(): boolean { + return this._ttsrAbortPending; + } + // ========================================================================= // Event Subscription // ========================================================================= @@ -239,6 +260,60 @@ export class AgentSession { // Notify all listeners this._emit(event); + // TTSR: Reset buffer on turn start + if (event.type === "turn_start" && this._ttsrManager) { + this._ttsrManager.resetBuffer(); + } + + // TTSR: Increment message count on turn end (for repeat-after-gap tracking) + if (event.type === "turn_end" && this._ttsrManager) { + this._ttsrManager.incrementMessageCount(); + } + + // TTSR: Check for pattern matches on text deltas and tool call argument deltas + if (event.type === "message_update" && this._ttsrManager?.hasRules()) { + const assistantEvent = event.assistantMessageEvent; + // Monitor both assistant prose (text_delta) and tool call arguments (toolcall_delta) + if (assistantEvent.type === "text_delta" || assistantEvent.type === "toolcall_delta") { + this._ttsrManager.appendToBuffer(assistantEvent.delta); + const matches = this._ttsrManager.check(this._ttsrManager.getBuffer()); + if (matches.length > 0) { + // Mark rules as injected so they don't trigger again + this._ttsrManager.markInjected(matches); + // Store for injection on retry + this._pendingTtsrInjections.push(...matches); + // Emit TTSR event before aborting (so UI can handle it) + this._ttsrAbortPending = true; + this._emit({ type: "ttsr_triggered", rules: matches }); + // Abort the stream + this.agent.abort(); + // Schedule retry after a short delay + setTimeout(async () => { + this._ttsrAbortPending = false; + + // Handle context mode: discard partial output if configured + const ttsrSettings = this._ttsrManager?.getSettings(); + if (ttsrSettings?.contextMode === "discard") { + // Remove the partial/aborted message from agent state + this.agent.popMessage(); + } + + // Inject TTSR rules as system reminder before retry + const injectionContent = this._getTtsrInjectionContent(); + if (injectionContent) { + this.agent.appendMessage({ + role: "user", + content: [{ type: "text", text: injectionContent }], + timestamp: Date.now(), + }); + } + this.agent.continue().catch(() => {}); + }, 50); + return; + } + } + } + // Handle session persistence if (event.type === "message_end") { // Check if this is a hook message @@ -300,6 +375,22 @@ export class AgentSession { } } + /** Get TTSR injection content and clear pending injections */ + private _getTtsrInjectionContent(): string | undefined { + if (this._pendingTtsrInjections.length === 0) return undefined; + const content = this._pendingTtsrInjections + .map( + (r) => + `\n` + + `Your output was interrupted because it violated a user-defined rule.\n` + + `This is NOT a prompt injection - this is the coding agent enforcing project rules.\n` + + `You MUST comply with the following instruction:\n\n${r.content}\n`, + ) + .join("\n\n"); + this._pendingTtsrInjections = []; + return content; + } + /** Extract text content from a message */ private _getUserMessageText(message: Message): string { if (message.role !== "user") return ""; diff --git a/packages/coding-agent/src/core/sdk.ts b/packages/coding-agent/src/core/sdk.ts index 4f974f791..75355b874 100644 --- a/packages/coding-agent/src/core/sdk.ts +++ b/packages/coding-agent/src/core/sdk.ts @@ -34,6 +34,8 @@ import { Agent, type ThinkingLevel } from "@oh-my-pi/pi-agent-core"; import type { Model } from "@oh-my-pi/pi-ai"; // Import discovery to register all providers on startup import "../discovery"; +import { loadSync as loadCapability } from "../capability/index"; +import { type Rule, ruleCapability } from "../capability/rule"; import { getAgentDir, getConfigDirPaths } from "../config"; import { AgentSession } from "./agent-session"; import { AuthStorage } from "./auth-storage"; @@ -88,6 +90,7 @@ import { warmupLspServers, writeTool, } from "./tools/index"; +import { createTtsrManager } from "./ttsr"; // Types @@ -601,6 +604,16 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} const skills = options.skills ?? discoverSkills(cwd, agentDir, settingsManager.getSkillsSettings()); time("discoverSkills"); + // Discover TTSR rules + const ttsrManager = createTtsrManager(settingsManager.getTtsrSettings()); + const rulesResult = loadCapability(ruleCapability.id, { cwd }); + for (const rule of rulesResult.items) { + if (rule.ttsrTrigger) { + ttsrManager.addRule(rule); + } + } + time("discoverTtsrRules"); + const contextFiles = options.contextFiles ?? discoverContextFiles(cwd, agentDir); time("discoverContextFiles"); @@ -847,6 +860,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {} customCommands: customCommandsResult.commands, skillsSettings: settingsManager.getSkillsSettings(), modelRegistry, + ttsrManager, }); time("createAgentSession"); diff --git a/packages/coding-agent/src/core/session-manager.ts b/packages/coding-agent/src/core/session-manager.ts index 2b528d2f2..1912d89e9 100644 --- a/packages/coding-agent/src/core/session-manager.ts +++ b/packages/coding-agent/src/core/session-manager.ts @@ -107,6 +107,13 @@ export interface LabelEntry extends SessionEntryBase { label: string | undefined; } +/** TTSR injection entry - tracks which time-traveling rules have been injected this session. */ +export interface TtsrInjectionEntry extends SessionEntryBase { + type: "ttsr_injection"; + /** Names of rules that were injected */ + injectedRules: string[]; +} + /** * Custom message entry for hooks to inject messages into LLM context. * Use customType to identify your hook's entries. @@ -136,7 +143,8 @@ export type SessionEntry = | BranchSummaryEntry | CustomEntry | CustomMessageEntry - | LabelEntry; + | LabelEntry + | TtsrInjectionEntry; /** Raw file entry (includes header) */ export type FileEntry = SessionHeader | SessionEntry; @@ -154,6 +162,8 @@ export interface SessionContext { thinkingLevel: string; /** Model roles: { default: "provider/modelId", small: "provider/modelId", ... } */ models: Record; + /** Names of TTSR rules that have been injected this session */ + injectedTtsrRules: string[]; } export interface SessionInfo { @@ -295,7 +305,7 @@ export function buildSessionContext( let leaf: SessionEntry | undefined; if (leafId === null) { // Explicitly null - return no messages (navigated to before first entry) - return { messages: [], thinkingLevel: "off", models: {} }; + return { messages: [], thinkingLevel: "off", models: {}, injectedTtsrRules: [] }; } if (leafId) { leaf = byId.get(leafId); @@ -306,7 +316,7 @@ export function buildSessionContext( } if (!leaf) { - return { messages: [], thinkingLevel: "off", models: {} }; + return { messages: [], thinkingLevel: "off", models: {}, injectedTtsrRules: [] }; } // Walk from leaf to root, collecting path @@ -321,6 +331,7 @@ export function buildSessionContext( let thinkingLevel = "off"; const models: Record = {}; let compaction: CompactionEntry | null = null; + const injectedTtsrRulesSet = new Set(); for (const entry of path) { if (entry.type === "thinking_level_change") { @@ -336,9 +347,16 @@ export function buildSessionContext( models.default = `${entry.message.provider}/${entry.message.model}`; } else if (entry.type === "compaction") { compaction = entry; + } else if (entry.type === "ttsr_injection") { + // Collect injected TTSR rule names + for (const ruleName of entry.injectedRules) { + injectedTtsrRulesSet.add(ruleName); + } } } + const injectedTtsrRules = Array.from(injectedTtsrRulesSet); + // Build messages and collect corresponding entries // When there's a compaction, we need to: // 1. Emit summary first (entry = compaction) @@ -389,7 +407,7 @@ export function buildSessionContext( } } - return { messages, thinkingLevel, models }; + return { messages, thinkingLevel, models, injectedTtsrRules }; } /** @@ -814,6 +832,44 @@ export class SessionManager { return entry.id; } + // ========================================================================= + // TTSR (Time Traveling Stream Rules) + // ========================================================================= + + /** + * Append a TTSR injection entry recording which rules were injected. + * @param ruleNames Names of rules that were injected + * @returns Entry id + */ + appendTtsrInjection(ruleNames: string[]): string { + const entry: TtsrInjectionEntry = { + type: "ttsr_injection", + id: generateId(this.byId), + parentId: this.leafId, + timestamp: new Date().toISOString(), + injectedRules: ruleNames, + }; + this._appendEntry(entry); + return entry.id; + } + + /** + * Get all unique TTSR rule names that have been injected in the current branch. + * Scans from root to current leaf for ttsr_injection entries. + */ + getInjectedTtsrRules(): string[] { + const path = this.getBranch(); + const ruleNames = new Set(); + for (const entry of path) { + if (entry.type === "ttsr_injection") { + for (const name of entry.injectedRules) { + ruleNames.add(name); + } + } + } + return Array.from(ruleNames); + } + // ========================================================================= // Tree Traversal // ========================================================================= diff --git a/packages/coding-agent/src/core/settings-manager.ts b/packages/coding-agent/src/core/settings-manager.ts index 5f65c207c..63005d7fe 100644 --- a/packages/coding-agent/src/core/settings-manager.ts +++ b/packages/coding-agent/src/core/settings-manager.ts @@ -68,6 +68,16 @@ export interface EditSettings { fuzzyMatch?: boolean; // default: true (accept high-confidence fuzzy matches for whitespace/indentation) } +export interface TtsrSettings { + enabled?: boolean; // default: true + /** What to do with partial output when TTSR triggers: "keep" shows interrupted attempt, "discard" removes it */ + contextMode?: "keep" | "discard"; // default: "discard" + /** How TTSR rules repeat: "once" = only trigger once per session, "after-gap" = can repeat after N messages */ + repeatMode?: "once" | "after-gap"; // default: "once" + /** Number of messages before a rule can trigger again (only used when repeatMode is "after-gap") */ + repeatGap?: number; // default: 10 +} + export interface Settings { lastChangelogVersion?: string; /** Model roles map: { default: "provider/modelId", small: "provider/modelId", ... } */ @@ -93,6 +103,7 @@ export interface Settings { mcp?: MCPSettings; lsp?: LspSettings; edit?: EditSettings; + ttsr?: TtsrSettings; disabledProviders?: string[]; // Discovery provider IDs that are disabled } @@ -582,4 +593,61 @@ export class SettingsManager { this.globalSettings.disabledProviders = providerIds; this.save(); } + + getTtsrSettings(): TtsrSettings { + return this.settings.ttsr ?? {}; + } + + setTtsrSettings(settings: TtsrSettings): void { + this.globalSettings.ttsr = { ...this.globalSettings.ttsr, ...settings }; + this.save(); + } + + getTtsrEnabled(): boolean { + return this.settings.ttsr?.enabled ?? true; + } + + setTtsrEnabled(enabled: boolean): void { + if (!this.globalSettings.ttsr) { + this.globalSettings.ttsr = {}; + } + this.globalSettings.ttsr.enabled = enabled; + this.save(); + } + + getTtsrContextMode(): "keep" | "discard" { + return this.settings.ttsr?.contextMode ?? "discard"; + } + + setTtsrContextMode(mode: "keep" | "discard"): void { + if (!this.globalSettings.ttsr) { + this.globalSettings.ttsr = {}; + } + this.globalSettings.ttsr.contextMode = mode; + this.save(); + } + + getTtsrRepeatMode(): "once" | "after-gap" { + return this.settings.ttsr?.repeatMode ?? "once"; + } + + setTtsrRepeatMode(mode: "once" | "after-gap"): void { + if (!this.globalSettings.ttsr) { + this.globalSettings.ttsr = {}; + } + this.globalSettings.ttsr.repeatMode = mode; + this.save(); + } + + getTtsrRepeatGap(): number { + return this.settings.ttsr?.repeatGap ?? 10; + } + + setTtsrRepeatGap(gap: number): void { + if (!this.globalSettings.ttsr) { + this.globalSettings.ttsr = {}; + } + this.globalSettings.ttsr.repeatGap = gap; + this.save(); + } } diff --git a/packages/coding-agent/src/core/ttsr.ts b/packages/coding-agent/src/core/ttsr.ts new file mode 100644 index 000000000..137f9aff9 --- /dev/null +++ b/packages/coding-agent/src/core/ttsr.ts @@ -0,0 +1,211 @@ +/** + * Time Traveling Stream Rules (TTSR) Manager + * + * Manages rules that get injected mid-stream when their trigger pattern matches + * the agent's output. When a match occurs, the stream is aborted, the rule is + * injected as a system reminder, and the request is retried. + */ + +import type { Rule } from "../capability/rule"; +import { logger } from "./logger"; +import type { TtsrSettings } from "./settings-manager"; + +interface TtsrEntry { + rule: Rule; + regex: RegExp; +} + +/** Tracks when a rule was last injected (for repeat-after-gap mode) */ +interface InjectionRecord { + /** Message count when the rule was last injected */ + lastInjectedAt: number; +} + +export interface TtsrManager { + /** Add a TTSR rule to be monitored */ + addRule(rule: Rule): void; + + /** Check if any uninjected TTSR matches the stream buffer. Returns matching rules. */ + check(streamBuffer: string): Rule[]; + + /** Mark rules as injected (won't trigger again until conditions allow) */ + markInjected(rules: Rule[]): void; + + /** Get names of all injected rules (for persistence) */ + getInjectedRuleNames(): string[]; + + /** Restore injected state from a list of rule names */ + restoreInjected(ruleNames: string[]): void; + + /** Reset stream buffer (called on new turn) */ + resetBuffer(): void; + + /** Get current stream buffer */ + getBuffer(): string; + + /** Append to stream buffer */ + appendToBuffer(text: string): void; + + /** Check if any TTSRs are registered */ + hasRules(): boolean; + + /** Increment message counter (call after each turn) */ + incrementMessageCount(): void; + + /** Get current message count */ + getMessageCount(): number; + + /** Get settings */ + getSettings(): Required; +} + +const DEFAULT_SETTINGS: Required = { + enabled: true, + contextMode: "discard", + repeatMode: "once", + repeatGap: 10, +}; + +export function createTtsrManager(settings?: TtsrSettings): TtsrManager { + /** Resolved settings with defaults */ + const resolvedSettings: Required = { + ...DEFAULT_SETTINGS, + ...settings, + }; + + /** Map of rule name -> { rule, compiled regex } */ + const rules = new Map(); + + /** Map of rule name -> injection record */ + const injectionRecords = new Map(); + + /** Current stream buffer for pattern matching */ + let buffer = ""; + + /** Message counter for tracking gap between injections */ + let messageCount = 0; + + /** Check if a rule can be triggered based on repeat settings */ + function canTrigger(ruleName: string): boolean { + const record = injectionRecords.get(ruleName); + if (!record) { + // Never injected, can trigger + return true; + } + + if (resolvedSettings.repeatMode === "once") { + // Once mode: never trigger again after first injection + return false; + } + + // After-gap mode: check if enough messages have passed + const gap = messageCount - record.lastInjectedAt; + return gap >= resolvedSettings.repeatGap; + } + + return { + addRule(rule: Rule): void { + // Only add rules that have a TTSR trigger pattern + if (!rule.ttsrTrigger) { + return; + } + + // Skip if already registered + if (rules.has(rule.name)) { + return; + } + + // Compile the regex pattern + try { + const regex = new RegExp(rule.ttsrTrigger); + rules.set(rule.name, { rule, regex }); + logger.debug("TTSR rule registered", { + ruleName: rule.name, + pattern: rule.ttsrTrigger, + }); + } catch (err) { + logger.warn("TTSR rule has invalid regex pattern, skipping", { + ruleName: rule.name, + pattern: rule.ttsrTrigger, + error: err instanceof Error ? err.message : String(err), + }); + } + }, + + check(streamBuffer: string): Rule[] { + const matches: Rule[] = []; + + for (const [name, entry] of rules) { + // Skip rules that can't trigger yet + if (!canTrigger(name)) { + continue; + } + + // Test the buffer against the rule's pattern + if (entry.regex.test(streamBuffer)) { + matches.push(entry.rule); + logger.debug("TTSR pattern matched", { + ruleName: name, + pattern: entry.rule.ttsrTrigger, + }); + } + } + + return matches; + }, + + markInjected(rulesToMark: Rule[]): void { + for (const rule of rulesToMark) { + injectionRecords.set(rule.name, { lastInjectedAt: messageCount }); + logger.debug("TTSR rule marked as injected", { + ruleName: rule.name, + messageCount, + repeatMode: resolvedSettings.repeatMode, + }); + } + }, + + getInjectedRuleNames(): string[] { + return Array.from(injectionRecords.keys()); + }, + + restoreInjected(ruleNames: string[]): void { + // When restoring, we don't know the original message count, so use 0 + // This means in "after-gap" mode, rules can trigger again after the gap + for (const name of ruleNames) { + injectionRecords.set(name, { lastInjectedAt: 0 }); + } + if (ruleNames.length > 0) { + logger.debug("TTSR injected state restored", { ruleNames }); + } + }, + + resetBuffer(): void { + buffer = ""; + }, + + getBuffer(): string { + return buffer; + }, + + appendToBuffer(text: string): void { + buffer += text; + }, + + hasRules(): boolean { + return rules.size > 0; + }, + + incrementMessageCount(): void { + messageCount++; + }, + + getMessageCount(): number { + return messageCount; + }, + + getSettings(): Required { + return resolvedSettings; + }, + }; +} diff --git a/packages/coding-agent/src/discovery/builtin.ts b/packages/coding-agent/src/discovery/builtin.ts index 262efdba9..006baea42 100644 --- a/packages/coding-agent/src/discovery/builtin.ts +++ b/packages/coding-agent/src/discovery/builtin.ts @@ -310,6 +310,7 @@ function loadRules(ctx: LoadContext): LoadResult { globs: frontmatter.globs as string[] | undefined, alwaysApply: frontmatter.alwaysApply as boolean | undefined, description: frontmatter.description as string | undefined, + ttsrTrigger: typeof frontmatter.ttsr_trigger === "string" ? frontmatter.ttsr_trigger : undefined, _source: source, }; }, diff --git a/packages/coding-agent/src/discovery/cline.ts b/packages/coding-agent/src/discovery/cline.ts index 0f6233a48..1e7cf9580 100644 --- a/packages/coding-agent/src/discovery/cline.ts +++ b/packages/coding-agent/src/discovery/cline.ts @@ -52,6 +52,7 @@ function loadRules(ctx: LoadContext): LoadResult { globs, alwaysApply: typeof frontmatter.alwaysApply === "boolean" ? frontmatter.alwaysApply : undefined, description: typeof frontmatter.description === "string" ? frontmatter.description : undefined, + ttsrTrigger: typeof frontmatter.ttsr_trigger === "string" ? frontmatter.ttsr_trigger : undefined, _source: source, }; }, @@ -85,6 +86,7 @@ function loadRules(ctx: LoadContext): LoadResult { globs, alwaysApply: typeof frontmatter.alwaysApply === "boolean" ? frontmatter.alwaysApply : undefined, description: typeof frontmatter.description === "string" ? frontmatter.description : undefined, + ttsrTrigger: typeof frontmatter.ttsr_trigger === "string" ? frontmatter.ttsr_trigger : undefined, _source: source, }); } diff --git a/packages/coding-agent/src/discovery/cursor.ts b/packages/coding-agent/src/discovery/cursor.ts index 29439d6f2..e9e46c562 100644 --- a/packages/coding-agent/src/discovery/cursor.ts +++ b/packages/coding-agent/src/discovery/cursor.ts @@ -163,6 +163,7 @@ function transformMDCRule( // Extract frontmatter fields const description = typeof frontmatter.description === "string" ? frontmatter.description : undefined; const alwaysApply = frontmatter.alwaysApply === true; + const ttsrTrigger = typeof frontmatter.ttsr_trigger === "string" ? frontmatter.ttsr_trigger : undefined; // Parse globs (can be array or single string) let globs: string[] | undefined; @@ -182,6 +183,7 @@ function transformMDCRule( description, alwaysApply, globs, + ttsrTrigger, _source: source, }; } diff --git a/packages/coding-agent/src/discovery/windsurf.ts b/packages/coding-agent/src/discovery/windsurf.ts index 90511c94b..0baca35d5 100644 --- a/packages/coding-agent/src/discovery/windsurf.ts +++ b/packages/coding-agent/src/discovery/windsurf.ts @@ -128,6 +128,7 @@ function loadRules(ctx: LoadContext): LoadResult { globs, alwaysApply: frontmatter.alwaysApply as boolean | undefined, description: frontmatter.description as string | undefined, + ttsrTrigger: typeof frontmatter.ttsr_trigger === "string" ? frontmatter.ttsr_trigger : undefined, _source: createSourceMeta(PROVIDER_ID, userPath, "user"), }); } @@ -157,6 +158,7 @@ function loadRules(ctx: LoadContext): LoadResult { globs, alwaysApply: frontmatter.alwaysApply as boolean | undefined, description: frontmatter.description as string | undefined, + ttsrTrigger: typeof frontmatter.ttsr_trigger === "string" ? frontmatter.ttsr_trigger : undefined, _source: source, }; }, @@ -187,6 +189,7 @@ function loadRules(ctx: LoadContext): LoadResult { globs, alwaysApply: frontmatter.alwaysApply as boolean | undefined, description: frontmatter.description as string | undefined, + ttsrTrigger: typeof frontmatter.ttsr_trigger === "string" ? frontmatter.ttsr_trigger : undefined, _source: createSourceMeta(PROVIDER_ID, legacyPath, "project"), }); } diff --git a/packages/coding-agent/src/modes/interactive/components/settings-defs.ts b/packages/coding-agent/src/modes/interactive/components/settings-defs.ts index b70b4d614..1ac2c23c6 100644 --- a/packages/coding-agent/src/modes/interactive/components/settings-defs.ts +++ b/packages/coding-agent/src/modes/interactive/components/settings-defs.ts @@ -155,6 +155,35 @@ export const SETTINGS_DEFS: SettingDef[] = [ get: (sm) => sm.getEditFuzzyMatch(), set: (sm, v) => sm.setEditFuzzyMatch(v), }, + { + id: "ttsrEnabled", + tab: "config", + type: "boolean", + label: "TTSR enabled", + description: "Time Traveling Stream Rules: interrupt agent when output matches rule patterns", + get: (sm) => sm.getTtsrEnabled(), + set: (sm, v) => sm.setTtsrEnabled(v), + }, + { + id: "ttsrContextMode", + tab: "config", + type: "enum", + label: "TTSR context mode", + description: "What to do with partial output when TTSR triggers", + values: ["discard", "keep"], + get: (sm) => sm.getTtsrContextMode(), + set: (sm, v) => sm.setTtsrContextMode(v as "keep" | "discard"), + }, + { + id: "ttsrRepeatMode", + tab: "config", + type: "enum", + label: "TTSR repeat mode", + description: "How rules can repeat: once per session or after a message gap", + values: ["once", "after-gap"], + get: (sm) => sm.getTtsrRepeatMode(), + set: (sm, v) => sm.setTtsrRepeatMode(v as "once" | "after-gap"), + }, { id: "thinkingLevel", tab: "config", diff --git a/packages/coding-agent/src/modes/interactive/components/ttsr-notification.ts b/packages/coding-agent/src/modes/interactive/components/ttsr-notification.ts new file mode 100644 index 000000000..75252cc4e --- /dev/null +++ b/packages/coding-agent/src/modes/interactive/components/ttsr-notification.ts @@ -0,0 +1,82 @@ +import { Box, Container, Spacer, Text } from "@oh-my-pi/pi-tui"; +import type { Rule } from "../../../capability/rule"; +import { theme } from "../theme/theme"; + +/** + * Component that renders a TTSR (Time Traveling Stream Rules) notification. + * Shows when a rule violation is detected and the stream is being rewound. + */ +export class TtsrNotificationComponent extends Container { + private rules: Rule[]; + private box: Box; + private _expanded = false; + + constructor(rules: Rule[]) { + super(); + this.rules = rules; + + this.addChild(new Spacer(1)); + + // Use inverse warning color for yellow background effect + this.box = new Box(1, 1, (t) => theme.inverse(theme.fg("warning", t))); + this.addChild(this.box); + + this.rebuild(); + } + + setExpanded(expanded: boolean): void { + if (this._expanded !== expanded) { + this._expanded = expanded; + this.rebuild(); + } + } + + isExpanded(): boolean { + return this._expanded; + } + + private rebuild(): void { + this.box.clear(); + + // Build header: ⚠ Injecting rule-name ↩ + const ruleNames = this.rules.map((r) => theme.bold(r.name)).join(", "); + const label = this.rules.length === 1 ? "rule" : "rules"; + const header = `\u26A0 Injecting ${label}: ${ruleNames}`; + + // Create header with rewind icon on the right + const rewindIcon = "\u21A9"; // ↩ + + this.box.addChild(new Text(`${header} ${rewindIcon}`, 0, 0)); + + // Show description(s) - italic and truncated + for (const rule of this.rules) { + const desc = rule.description || rule.content; + if (desc) { + this.box.addChild(new Spacer(1)); + + let displayText = desc.trim(); + if (!this._expanded) { + // Truncate to first 2 lines + const lines = displayText.split("\n"); + if (lines.length > 2) { + displayText = `${lines.slice(0, 2).join("\n")}...`; + } + } + + // Use italic for subtle distinction (fg colors conflict with inverse) + this.box.addChild(new Text(theme.italic(displayText), 0, 0)); + } + } + + // Show expand hint if collapsed and there's more content + if (!this._expanded) { + const hasMoreContent = this.rules.some((r) => { + const desc = r.description || r.content; + return desc && desc.split("\n").length > 2; + }); + if (hasMoreContent) { + this.box.addChild(new Text(theme.italic(" (ctrl+o to expand)"), 0, 0)); + } + } + } +} diff --git a/packages/coding-agent/src/modes/interactive/interactive-mode.ts b/packages/coding-agent/src/modes/interactive/interactive-mode.ts index 4ee852482..7cd1fa5f5 100644 --- a/packages/coding-agent/src/modes/interactive/interactive-mode.ts +++ b/packages/coding-agent/src/modes/interactive/interactive-mode.ts @@ -67,6 +67,7 @@ import { SessionSelectorComponent } from "./components/session-selector"; import { SettingsSelectorComponent } from "./components/settings-selector"; import { ToolExecutionComponent } from "./components/tool-execution"; import { TreeSelectorComponent } from "./components/tree-selector"; +import { TtsrNotificationComponent } from "./components/ttsr-notification"; import { UserMessageComponent } from "./components/user-message"; import { UserMessageSelectorComponent } from "./components/user-message-selector"; import { WelcomeComponent } from "./components/welcome"; @@ -1007,18 +1008,28 @@ export class InteractiveMode { if (event.message.role === "user") break; if (this.streamingComponent && event.message.role === "assistant") { this.streamingMessage = event.message; - this.streamingComponent.updateContent(this.streamingMessage); + // Don't show "Aborted" text for TTSR aborts - we'll show a nicer message + if (this.session.isTtsrAbortPending && this.streamingMessage.stopReason === "aborted") { + // TTSR abort - suppress the "Aborted" rendering in the component + const msgWithoutAbort = { ...this.streamingMessage, stopReason: "stop" as const }; + this.streamingComponent.updateContent(msgWithoutAbort); + } else { + this.streamingComponent.updateContent(this.streamingMessage); + } if (this.streamingMessage.stopReason === "aborted" || this.streamingMessage.stopReason === "error") { - const errorMessage = - this.streamingMessage.stopReason === "aborted" - ? "Operation aborted" - : this.streamingMessage.errorMessage || "Error"; - for (const [, component] of this.pendingTools.entries()) { - component.updateResult({ - content: [{ type: "text", text: errorMessage }], - isError: true, - }); + // Skip error handling for TTSR aborts + if (!this.session.isTtsrAbortPending) { + const errorMessage = + this.streamingMessage.stopReason === "aborted" + ? "Operation aborted" + : this.streamingMessage.errorMessage || "Error"; + for (const [, component] of this.pendingTools.entries()) { + component.updateResult({ + content: [{ type: "text", text: errorMessage }], + isError: true, + }); + } } this.pendingTools.clear(); } else { @@ -1188,6 +1199,15 @@ export class InteractiveMode { this.ui.requestRender(); break; } + + case "ttsr_triggered": { + // Show a fancy notification when TTSR rules are triggered + const component = new TtsrNotificationComponent(event.rules); + component.setExpanded(this.toolOutputExpanded); + this.chatContainer.addChild(component); + this.ui.requestRender(); + break; + } } } @@ -2490,9 +2510,7 @@ export class InteractiveMode { const nameRendered = sourceText ? `${theme.bold(nameText)} ${sourceRendered}` : theme.bold(nameText); const pad = Math.max(0, maxNameWidth - visibleWidth(nameWithSourcePlain)); const desc = line.desc; - const descPart = desc - ? ` ${theme.fg("dim", desc.slice(0, 50) + (desc.length > 50 ? "..." : ""))}` - : ""; + const descPart = desc ? ` ${theme.fg("dim", desc.slice(0, 50) + (desc.length > 50 ? "..." : ""))}` : ""; return ` ${nameRendered}${" ".repeat(pad)}${descPart}`; }); @@ -2536,15 +2554,13 @@ export class InteractiveMode { (s) => (s._source ? { provider: s._source.providerName, level: s._source.level } : "unknown"), ); if (skillWarnings.length > 0) { - sections.push( - { - kind: "text", - text: - theme.bold(theme.fg("warning", "Skill Warnings")) + - "\n" + - skillWarnings.map((w) => theme.fg("warning", ` ${w.skillPath}: ${w.message}`)).join("\n"), - }, - ); + sections.push({ + kind: "text", + text: + theme.bold(theme.fg("warning", "Skill Warnings")) + + "\n" + + skillWarnings.map((w) => theme.fg("warning", ` ${w.skillPath}: ${w.message}`)).join("\n"), + }); } } @@ -2633,14 +2649,12 @@ export class InteractiveMode { if (hookRunner) { const hookPaths = hookRunner.getHookPaths(); if (hookPaths.length > 0) { - sections.push( - { - kind: "text", - text: - `${theme.bold(theme.fg("accent", "Hooks"))}\n` + - hookPaths.map((p) => ` ${theme.bold(basename(p))} ${theme.fg("dim", "hook")}`).join("\n"), - }, - ); + sections.push({ + kind: "text", + text: + `${theme.bold(theme.fg("accent", "Hooks"))}\n` + + hookPaths.map((p) => ` ${theme.bold(basename(p))} ${theme.fg("dim", "hook")}`).join("\n"), + }); } } @@ -2652,9 +2666,7 @@ export class InteractiveMode { ? Math.min(60, Math.max(...allLines.map((line) => visibleWidth(line.nameWithSource)))) : 0; const renderedSections = sections - .map((section) => - section.kind === "lines" ? renderLineSection(section.section, maxNameWidth) : section.text, - ) + .map((section) => (section.kind === "lines" ? renderLineSection(section.section, maxNameWidth) : section.text)) .filter((section) => section.length > 0); if (renderedSections.length === 0) { diff --git a/packages/tui/src/terminal.ts b/packages/tui/src/terminal.ts index c2d1ee6a9..d733744b4 100644 --- a/packages/tui/src/terminal.ts +++ b/packages/tui/src/terminal.ts @@ -10,10 +10,10 @@ let activeTerminal: ProcessTerminal | null = null; * Resets terminal state without requiring access to the ProcessTerminal instance */ export function emergencyTerminalRestore(): void { - if (activeTerminal) { - activeTerminal.stop(); - activeTerminal.showCursor(); - activeTerminal = null; + const terminal = activeTerminal; + if (terminal) { + terminal.stop(); + terminal.showCursor(); } else { // Blind restore if no instance tracked - covers edge cases process.stdout.write(