feat: introduced TTSR for real-time pattern matching during streaming

- Added Time Traveling Streamed Rules (TTSR) feature for real-time pattern matching and rule injection during streaming.
- Added TtsrManager with configurable settings for context mode, repeat mode, and repeat gap.
- Added TTSR trigger support across all rule discovery sources (builtin, cline, cursor, windsurf).
- Added abort signal handling in agent loop to support stream interruption for TTSR.
- Added TtsrNotificationComponent for displaying rule violation notifications in interactive mode.
This commit is contained in:
can1357
2026-01-03 19:43:03 +01:00
parent f8e642adce
commit 13f3eae8c0
20 changed files with 690 additions and 55 deletions
+29 -13
View File
@@ -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)
<p align="center">
<img src="https://github.com/can1357/oh-my-pi/blob/main/assets/ttsr.webp?raw=true" alt="ttsr">
</p>
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
<p align="center">
<img src="https://github.com/can1357/oh-my-pi/blob/main/assets/review.webp?raw=true" alt="review">
</p>
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)
<p align="center">
@@ -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
<p align="center">
<img src="https://github.com/can1357/oh-my-pi/blob/main/assets/review.webp?raw=true" alt="review">
</p>
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
<p align="center">
BIN
View File
Binary file not shown.

After

Width:  |  Height:  |  Size: 3.4 MiB

+8
View File
@@ -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
+21
View File
@@ -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();
+8
View File
@@ -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();
}
+7
View File
@@ -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
@@ -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;
}
@@ -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<SkillsSettings>;
/** 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) =>
`<system_interrupt reason="rule_violation" rule="${r.name}" path="${r.path}">\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</system_interrupt>`,
)
.join("\n\n");
this._pendingTtsrInjections = [];
return content;
}
/** Extract text content from a message */
private _getUserMessageText(message: Message): string {
if (message.role !== "user") return "";
+14
View File
@@ -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<Rule>(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");
@@ -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<string, string>;
/** 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<string, string> = {};
let compaction: CompactionEntry | null = null;
const injectedTtsrRulesSet = new Set<string>();
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<string>();
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
// =========================================================================
@@ -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();
}
}
+211
View File
@@ -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<TtsrSettings>;
}
const DEFAULT_SETTINGS: Required<TtsrSettings> = {
enabled: true,
contextMode: "discard",
repeatMode: "once",
repeatGap: 10,
};
export function createTtsrManager(settings?: TtsrSettings): TtsrManager {
/** Resolved settings with defaults */
const resolvedSettings: Required<TtsrSettings> = {
...DEFAULT_SETTINGS,
...settings,
};
/** Map of rule name -> { rule, compiled regex } */
const rules = new Map<string, TtsrEntry>();
/** Map of rule name -> injection record */
const injectionRecords = new Map<string, InjectionRecord>();
/** 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<TtsrSettings> {
return resolvedSettings;
},
};
}
@@ -310,6 +310,7 @@ function loadRules(ctx: LoadContext): LoadResult<Rule> {
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,
};
},
@@ -52,6 +52,7 @@ function loadRules(ctx: LoadContext): LoadResult<Rule> {
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<Rule> {
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,
});
}
@@ -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,
};
}
@@ -128,6 +128,7 @@ function loadRules(ctx: LoadContext): LoadResult<Rule> {
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<Rule> {
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<Rule> {
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"),
});
}
@@ -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",
@@ -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 <bold>rule-name</bold> ↩
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));
}
}
}
}
@@ -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) {
+4 -4
View File
@@ -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(