feat(session): added modular session APIs and rebuilt listing/persistence behavior

- Added session-domain modules and exports for session-entries, context, listing, loader, and migrations.
- Changed persistence to async append writes plus writeTextAtomic, removing sync line APIs.
- Added compaction-aware session context rebuild with dangling tool-call cleanup.
- Added resumable session resolution with status inference, id/stem/suffix matching, and backup recovery.
This commit is contained in:
can1357
2026-06-14 01:32:10 +02:00
parent 17ba881c5a
commit 24c8bb24c6
94 changed files with 3129 additions and 3688 deletions
+13
View File
@@ -1,9 +1,18 @@
# Changelog
## [Unreleased]
### Breaking Changes
- Removed the `writeLine` and `writeLineSync` methods from the public `SessionStorageWriter` contract, requiring custom `SessionStorage` backends to switch to the `append` API
### Added
- Added package-level exports for `SessionContext`, session entry types, session listing/loader helpers, and migration APIs via `session/session-context`, `session/session-entries`, `session/session-listing`, `session/session-loader`, and `session/session-migrations`
- Added asynchronous session-write `append(...)`-based persistence API in session storage implementations so callers can stream writes without sync line-appending methods
### Changed
- Changed session persistence internals to expose `writeTextAtomic(...)` on session storage writers for atomic whole-file replacements
- Changed online session-title generation to support tool-choice-less title models. Providers/models that cannot be forced to call a tool (chat-completions hosts without `tool_choice` support such as DeepSeek V4, and Claude Fable/Mythos) are now prompted to wrap the title in `<title>...</title>` markers instead of the `set_title` tool call; extraction is lenient, accepting a plain sentence or a truncated/unclosed tag. A `TITLE_SYSTEM.md` override is reused in this mode with the marker instruction appended.
### Fixed
@@ -12,6 +21,10 @@
- Fixed submitted user messages emitting OSC 133 command-start markers without a matching command-finished marker, which made some terminals group later transcript output under the first prompt instead of appending it normally.
- Fixed queued forced tool choices being rejected and requeued or dropped when their named tool is no longer active for the upcoming turn, preventing eager todo and pending-action reminders from forcing unavailable tools. ([#1701](https://github.com/can1357/oh-my-pi/issues/1701))
### Removed
- Removed the `re-roots past a cwd-less legacy session in a shared explicit sessionDir` relocation test case and the `stores symlink-equivalent home cwd sessions under home-relative directories` file-operations test case.
## [15.12.5] - 2026-06-13
### Changed
@@ -1,4 +1,4 @@
import type { SessionEntry } from "../session/session-manager";
import type { SessionEntry } from "../session/session-entries";
import { inferMetricUnitFromName, isBetter } from "./helpers";
import type { RunRow, SessionRow } from "./storage";
import type {
@@ -1,6 +1,6 @@
import type { AgentToolResult } from "@oh-my-pi/pi-agent-core";
import type { ExtensionAPI, ExtensionContext } from "../extensibility/extensions";
import type { SessionEntry } from "../session/session-manager";
import type { SessionEntry } from "../session/session-entries";
import type { TruncationResult } from "../session/streaming-output";
export type MetricDirection = "lower" | "higher";
@@ -2,7 +2,8 @@ import { ProcessTerminal, TUI } from "@oh-my-pi/pi-tui";
import { logger } from "@oh-my-pi/pi-utils";
import { SessionSelectorComponent } from "../modes/components/session-selector";
import { HistoryStorage } from "../session/history-storage";
import { type SessionInfo, SessionManager } from "../session/session-manager";
import type { SessionInfo } from "../session/session-listing";
import { SessionManager } from "../session/session-manager";
import { FileSessionStorage } from "../session/session-storage";
/**
+1 -1
View File
@@ -20,7 +20,7 @@ import { AgentLifecycleManager } from "../registry/agent-lifecycle";
import { AgentRegistry } from "../registry/agent-registry";
import type { AgentSessionEvent } from "../session/agent-session";
import { stripImagesFromMessage, USER_INTERRUPT_LABEL } from "../session/messages";
import type { SessionEntry as StoredSessionEntry } from "../session/session-manager";
import type { SessionEntry as StoredSessionEntry } from "../session/session-entries";
import { TASK_SUBAGENT_LIFECYCLE_CHANNEL, TASK_SUBAGENT_PROGRESS_CHANNEL } from "../task";
import { generateRoomKey, generateWriteToken, importRoomKey } from "./crypto";
import {
+1 -1
View File
@@ -25,7 +25,7 @@ import {
} from "@oh-my-pi/pi-wire";
import type { ContextUsage } from "../extensibility/extensions/types";
import type { AgentSessionEvent } from "../session/agent-session";
import type { SessionEntry, SessionHeader } from "../session/session-manager";
import type { SessionEntry, SessionHeader } from "../session/session-entries";
export type {
CollabPromptDetails,
@@ -1,6 +1,6 @@
import { describe, expect, it } from "bun:test";
import type { GoalModeState } from "../../goals/state";
import type { UsageStatistics } from "../../session/session-manager";
import type { UsageStatistics } from "../../session/session-entries";
import type { ToolSession } from "../../tools";
import { runEvalBudget } from "../budget-bridge";
@@ -3,12 +3,9 @@ import * as path from "node:path";
import type { AgentState } from "@oh-my-pi/pi-agent-core";
import { APP_NAME, isEnoent } from "@oh-my-pi/pi-utils";
import { getResolvedThemeColors, getThemeExportColors } from "../../modes/theme/theme";
import {
loadEntriesFromFile,
type SessionEntry,
type SessionHeader,
SessionManager,
} from "../../session/session-manager";
import type { SessionEntry, SessionHeader } from "../../session/session-entries";
import { loadEntriesFromFile } from "../../session/session-loader";
import { SessionManager } from "../../session/session-manager";
import templateCss from "./template.css" with { type: "text" };
import templateHtml from "./template.html" with { type: "text" };
import templateJs from "./template.js" with { type: "text" };
@@ -1,4 +1,5 @@
export type { ReadonlySessionManager, UsageStatistics } from "../../session/session-manager";
export type { UsageStatistics } from "../../session/session-entries";
export type { ReadonlySessionManager } from "../../session/session-manager";
export * from "./loader";
export * from "./runner";
export * from "./tool-wrapper";
@@ -17,7 +17,7 @@ import type { CompactionPreparation, CompactionResult } from "@oh-my-pi/pi-agent
import type { ImageContent, TextContent, ToolResultMessage } from "@oh-my-pi/pi-ai";
import type { Rule } from "../capability/rule";
import type { Goal, GoalModeState } from "../goals/state";
import type { BranchSummaryEntry, CompactionEntry, SessionEntry } from "../session/session-manager";
import type { BranchSummaryEntry, CompactionEntry, SessionEntry } from "../session/session-entries";
import type { TodoItem } from "../tools/todo";
// ============================================================================
+1 -1
View File
@@ -1,4 +1,4 @@
import type { UsageStatistics } from "../session/session-manager";
import type { UsageStatistics } from "../session/session-entries";
export type GoalStatus = "active" | "paused" | "budget-limited" | "complete" | "dropped";
@@ -8,7 +8,7 @@
*/
import type { AssistantMessage } from "@oh-my-pi/pi-ai";
import type { SessionEntry } from "../session/session-manager";
import type { SessionEntry } from "../session/session-entries";
import type { HindsightMessage } from "./content";
export interface ReadonlySessionManagerLike {
+5
View File
@@ -42,8 +42,13 @@ export * from "./session/auth-storage";
export * from "./session/indexed-session-storage";
export * from "./session/messages";
export * from "./session/redis-session-storage";
export * from "./session/session-context";
export * from "./session/session-dump-format";
export * from "./session/session-entries";
export * from "./session/session-listing";
export * from "./session/session-loader";
export * from "./session/session-manager";
export * from "./session/session-migrations";
export * from "./session/session-storage";
export * from "./session/sql-session-storage";
export * from "./task/executor";
@@ -12,7 +12,7 @@
import type { AgentRef } from "../registry/agent-registry";
import { AgentRegistry } from "../registry/agent-registry";
import { formatSessionHistoryMarkdown } from "../session/session-history-format";
import { loadSessionMessagesReadOnly } from "../session/session-manager";
import { loadSessionMessagesReadOnly } from "../session/session-loader";
import type { InternalResource, InternalUrl, ProtocolHandler, UrlCompletion } from "./types";
/** Humanize a last-activity timestamp as `Ns/Nm/Nh/Nd ago`. */
+2 -1
View File
@@ -64,7 +64,8 @@ import {
} from "./sdk";
import type { AgentSession } from "./session/agent-session";
import type { AuthStorage } from "./session/auth-storage";
import { resolveResumableSession, type SessionInfo, SessionManager } from "./session/session-manager";
import { resolveResumableSession, type SessionInfo } from "./session/session-listing";
import { SessionManager } from "./session/session-manager";
import { executeBuiltinSlashCommand } from "./slash-commands/builtin-registry";
import { discoverTitleSystemPromptFile, resolvePromptInput } from "./system-prompt";
import { initTelemetryExport, isTelemetryExportEnabled } from "./telemetry-export";
@@ -63,11 +63,9 @@ import { theme } from "../../modes/theme/theme";
import { type PlanApprovalDetails, resolveApprovedPlan } from "../../plan-mode/approved-plan";
import type { AgentSession, AgentSessionEvent } from "../../session/agent-session";
import { isSilentAbort, SKILL_PROMPT_MESSAGE_TYPE } from "../../session/messages";
import {
SessionManager,
type SessionInfo as StoredSessionInfo,
type UsageStatistics,
} from "../../session/session-manager";
import type { UsageStatistics } from "../../session/session-entries";
import type { SessionInfo as StoredSessionInfo } from "../../session/session-listing";
import { SessionManager } from "../../session/session-manager";
import { executeAcpBuiltinSlashCommand } from "../../slash-commands/acp-builtins";
import { buildAvailableSlashCommands, toAcpAvailableCommands } from "../../slash-commands/available-commands";
import { AUTO_THINKING, parseConfiguredThinkingLevel } from "../../thinking";
@@ -35,8 +35,8 @@ import {
type SkillPromptDetails,
USER_INTERRUPT_LABEL,
} from "../../session/messages";
import type { SessionMessageEntry } from "../../session/session-manager";
import { parseSessionEntries } from "../../session/session-manager";
import type { SessionMessageEntry } from "../../session/session-entries";
import { parseSessionEntries } from "../../session/session-loader";
import { createIrcMessageCard } from "../../tools/irc";
import { replaceTabs, TRUNCATE_LENGTHS, truncateToWidth } from "../../tools/render-utils";
import { hasVisibleThinking } from "../../utils/thinking-display";
@@ -15,7 +15,7 @@ import {
import { formatBytes } from "@oh-my-pi/pi-utils";
import { theme } from "../../modes/theme/theme";
import { matchesAppInterrupt, matchesSelectDown, matchesSelectUp } from "../../modes/utils/keybinding-matchers";
import type { SessionInfo, SessionStatus } from "../../session/session-manager";
import type { SessionInfo, SessionStatus } from "../../session/session-listing";
import { shortenPath } from "../../tools/render-utils";
import { DynamicBorder } from "./dynamic-border";
import { HookSelectorComponent } from "./hook-selector";
@@ -8,6 +8,7 @@ import {
Image,
ImageProtocol,
imageFallback,
type NativeScrollbackLiveRegion,
Spacer,
TERMINAL,
Text,
@@ -160,7 +161,7 @@ let toolExecutionInstanceSeq = 0;
/**
* Component that renders a tool call with its result (updateable)
*/
export class ToolExecutionComponent extends Container {
export class ToolExecutionComponent extends Container implements NativeScrollbackLiveRegion {
#contentBox: Box; // Used for custom tools and bash visual truncation
#contentText: Text; // For built-in tools (with its own padding/bg)
#multiFileBoxes: (Box | Spacer)[] = []; // Extra boxes for multi-file edit results
@@ -568,6 +569,17 @@ export class ToolExecutionComponent extends Container {
}
}
/**
* Standalone harnesses may mount a tool component directly under `TUI`
* instead of inside `TranscriptContainer`. In that shape the component must
* report its own live-region seam for provisional previews, or the core
* renderer treats it like shell output and commits tail-window edit/eval/bash
* previews to immutable native scrollback before the result replaces them.
*/
getNativeScrollbackLiveRegionStart(): number | undefined {
return !this.isTranscriptBlockFinalized() && !this.isTranscriptBlockCommitStable() ? 0 : undefined;
}
/**
* Whether this block has reached a terminal state for transcript freezing.
* Reports `false` while it can still visually change so the
@@ -591,28 +603,20 @@ export class ToolExecutionComponent extends Container {
/**
* Whether this still-live block's settled rows may enter native scrollback
* (see `FinalizableBlock.isTranscriptBlockCommitStable`). Classification is
* per renderer (`ToolRenderer.provisionalPendingPreview`): tail-window
* streaming views (edit's streamed-diff tail, bash/ssh command caps, eval
* cells) are re-anchored top-first by the result render, so promoting
* their visually static head — e.g. an edit preview idling on its last
* frame while the apply + LSP pass runs — would strand a stale copy of
* the call box above the final block the moment the result lands. Every
* other pending preview streams top-anchored append-shaped rows the
* result render preserves (a task call's context/assignment markdown, a
* write's content), so it stays commit-eligible — a call taller than the
* viewport scrolls into native history mid-stream instead of reading as
* cut off until the result. Expanded blocks always stream top-anchored
* (the over-tall write/eval scrollback contract). Displaceable waiting
* polls are removed wholesale by the next poll and must never commit.
* (see `FinalizableBlock.isTranscriptBlockCommitStable`). Renderers classify
* pending views by durability instead of by tool name: a provisional view is
* allowed to be useful on screen, but finalization may replace or re-anchor
* it wholesale, so committing any of its rows would strand stale preview
* bytes in immutable scrollback. Non-provisional views stream rows whose
* committed prefix survives the remaining transitions.
*/
isTranscriptBlockCommitStable(): boolean {
if (this.#displaceable) return false;
if (this.#expanded || this.isTranscriptBlockFinalized()) return true;
if ((this.#tool as { provisionalPendingPreview?: boolean } | undefined)?.provisionalPendingPreview) {
return false;
}
return !toolRenderers[this.#toolName]?.provisionalPendingPreview;
if (this.isTranscriptBlockFinalized()) return true;
const tool = this.#tool as { provisionalPendingPreview?: boolean | "collapsed" } | undefined;
const provisionalPendingPreview =
tool?.provisionalPendingPreview ?? toolRenderers[this.#toolName]?.provisionalPendingPreview;
return provisionalPendingPreview !== true && (provisionalPendingPreview !== "collapsed" || this.#expanded);
}
/**
@@ -15,7 +15,7 @@ import {
import type { TreeFilterMode } from "../../config/settings-schema";
import { theme } from "../../modes/theme/theme";
import { matchesAppInterrupt, matchesSelectDown, matchesSelectUp } from "../../modes/utils/keybinding-matchers";
import type { SessionTreeNode } from "../../session/session-manager";
import type { SessionTreeNode } from "../../session/session-entries";
import { shortenPath } from "../../tools/render-utils";
import { toPathList } from "../../tools/search";
import { DynamicBorder } from "./dynamic-border";
@@ -38,7 +38,7 @@ import { buildHotkeysMarkdown } from "../../modes/utils/hotkeys-markdown";
import { buildToolsMarkdown } from "../../modes/utils/tools-markdown";
import type { AsyncJobSnapshotItem } from "../../session/agent-session";
import type { AuthStorage, OAuthAccountIdentity } from "../../session/auth-storage";
import type { NewSessionOptions } from "../../session/session-manager";
import type { NewSessionOptions } from "../../session/session-entries";
import { formatShakeSummary, type ShakeMode, type ShakeResult } from "../../session/shake-types";
import { limitMatchesActiveAccount } from "../../slash-commands/helpers/active-oauth-account";
import { outputMeta } from "../../tools/output-meta";
@@ -28,7 +28,8 @@ import {
} from "../../modes/theme/theme";
import type { InteractiveModeContext } from "../../modes/types";
import type { ResetCreditRedeemOutcome } from "../../session/auth-storage";
import { type SessionInfo, SessionManager } from "../../session/session-manager";
import type { SessionInfo } from "../../session/session-listing";
import { SessionManager } from "../../session/session-manager";
import { FileSessionStorage } from "../../session/session-storage";
import { type LogoutAccount, toLogoutAccounts } from "../../slash-commands/helpers/logout";
import {
@@ -80,8 +80,9 @@ import planModeCompactInstructionsPrompt from "../prompts/system/plan-mode-compa
};
import type { AgentSession, AgentSessionEvent, ResolvedRoleModel } from "../session/agent-session";
import { HistoryStorage } from "../session/history-storage";
import type { SessionContext, SessionManager } from "../session/session-manager";
import { getRecentSessions } from "../session/session-manager";
import type { SessionContext } from "../session/session-context";
import { getRecentSessions } from "../session/session-listing";
import type { SessionManager } from "../session/session-manager";
import type { ShakeMode } from "../session/shake-types";
import { BUILTIN_SLASH_COMMAND_RESERVED_NAMES, BUILTIN_SLASH_COMMANDS } from "../slash-commands/builtin-registry";
import { formatDuration } from "../slash-commands/helpers/format";
@@ -1,7 +1,7 @@
import * as fs from "node:fs/promises";
import { isEnoent } from "@oh-my-pi/pi-utils";
import type { FileEntry, SessionMessageEntry } from "../../session/session-manager";
import { parseSessionEntries } from "../../session/session-manager";
import type { FileEntry, SessionMessageEntry } from "../../session/session-entries";
import { parseSessionEntries } from "../../session/session-loader";
import {
type AgentProgress,
type SubagentEventPayload,
@@ -10,7 +10,7 @@ import type { Effort, ImageContent, Model } from "@oh-my-pi/pi-ai";
import type { BashResult } from "../../exec/bash-executor";
import type { ContextUsage } from "../../extensibility/extensions/types";
import type { AgentSessionEvent, SessionStats } from "../../session/agent-session";
import type { FileEntry } from "../../session/session-manager";
import type { FileEntry } from "../../session/session-entries";
import type { AvailableSlashCommandSource } from "../../slash-commands/available-commands";
import type {
AgentProgress,
+2 -1
View File
@@ -18,7 +18,8 @@ import type { MCPManager } from "../mcp";
import type { PlanApprovalDetails } from "../plan-mode/approved-plan";
import type { AgentSession } from "../session/agent-session";
import type { HistoryStorage } from "../session/history-storage";
import type { SessionContext, SessionManager } from "../session/session-manager";
import type { SessionContext } from "../session/session-context";
import type { SessionManager } from "../session/session-manager";
import type { ShakeMode } from "../session/shake-types";
import type { LspStartupServerInfo } from "../tools";
import type { EventBus } from "../utils/event-bus";
@@ -36,7 +36,7 @@ import {
SKILL_PROMPT_MESSAGE_TYPE,
type SkillPromptDetails,
} from "../../session/messages";
import type { SessionContext } from "../../session/session-manager";
import type { SessionContext } from "../../session/session-context";
import { createIrcMessageCard } from "../../tools/irc";
import { formatBytes, formatDuration } from "../../tools/render-utils";
import { hasVisibleThinking } from "../../utils/thinking-display";
+2 -1
View File
@@ -128,7 +128,8 @@ import {
LSP_LATE_DIAGNOSTIC_MESSAGE_TYPE,
wrapSteeringForModel,
} from "./session/messages";
import { getRestorableSessionModels, SessionManager } from "./session/session-manager";
import { getRestorableSessionModels } from "./session/session-context";
import { SessionManager } from "./session/session-manager";
import { SnapcompactInlineTransformer } from "./session/snapcompact-inline";
import { closeAllConnections } from "./ssh/connection-manager";
import { unmountAll } from "./ssh/sshfs-mount";
@@ -1,6 +1,6 @@
import type { Context, Message, Tool } from "@oh-my-pi/pi-ai";
import { toolWireSchema } from "@oh-my-pi/pi-ai/utils/schema";
import type { SessionContext } from "../session/session-manager";
import type { SessionContext } from "../session/session-context";
import { compileSecretRegex } from "./regex";
// ═══════════════════════════════════════════════════════════════════════════
@@ -258,15 +258,12 @@ import {
SKILL_PROMPT_MESSAGE_TYPE,
stripImagesFromMessage,
} from "./messages";
import type { SessionContext } from "./session-context";
import { getLatestCompactionEntry, getRestorableSessionModels } from "./session-context";
import { formatSessionDumpText } from "./session-dump-format";
import type {
BranchSummaryEntry,
CompactionEntry,
NewSessionOptions,
SessionContext,
SessionManager,
} from "./session-manager";
import { EPHEMERAL_MODEL_CHANGE_ROLE, getLatestCompactionEntry, getRestorableSessionModels } from "./session-manager";
import type { BranchSummaryEntry, CompactionEntry, NewSessionOptions } from "./session-entries";
import { EPHEMERAL_MODEL_CHANGE_ROLE } from "./session-entries";
import type { SessionManager } from "./session-manager";
import type { ShakeMode, ShakeResult } from "./shake-types";
import { ToolChoiceQueue } from "./tool-choice-queue";
import { YieldQueue } from "./yield-queue";
@@ -931,6 +928,7 @@ export class AgentSession {
/** Messages queued to be included with the next user prompt as context ("asides"). */
#pendingNextTurnMessages: CustomMessage[] = [];
#scheduledHiddenNextTurnGeneration: number | undefined = undefined;
#queuedMessageDrainScheduled = false;
#planModeState: PlanModeState | undefined;
#goalModeState: GoalModeState | undefined;
#goalRuntime: GoalRuntime;
@@ -1085,6 +1083,7 @@ export class AgentSession {
#streamingEditFileCache = new Map<string, string>();
#promptInFlightCount = 0;
#abortInProgress = false;
// Wire-level agent_end emission deferred until #promptInFlightCount drops to 0.
// Internal extension hooks and post-emit work (auto-retry, auto-compaction, todo
// checks in #handleAgentEvent) still fire on the original schedule — only the
@@ -1160,10 +1159,8 @@ export class AgentSession {
* Runs whenever the session settles; the guard makes it a no-op when the
* queue was consumed normally or a new turn already started. */
#drainStrandedQueuedMessages(): void {
if (!this.agent.hasQueuedMessages()) return;
this.#scheduleAgentContinue({
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
});
if (this.#abortInProgress) return;
this.#scheduleQueuedMessageDrain();
}
#resetInFlight(): void {
@@ -1916,6 +1913,11 @@ export class AgentSession {
return;
}
// A deliberate abort should settle the current turn, not trigger queued continuations.
if (msg.stopReason === "aborted") {
this.#resolveRetry();
return;
}
// Check for retryable errors first (overloaded, rate limit, server errors)
if (this.#isRetryableError(msg)) {
const didRetry = await this.#handleRetryableError(msg);
@@ -1938,7 +1940,7 @@ export class AgentSession {
if (compactionDeferredHandoff) {
return;
}
if (msg.stopReason !== "error" && msg.stopReason !== "aborted") {
if (msg.stopReason !== "error") {
if (this.#enforceRewindBeforeYield()) {
return;
}
@@ -2034,13 +2036,13 @@ export class AgentSession {
onError?: () => void;
}): void {
this.#schedulePostPromptTask(
async () => {
async signal => {
// Defense in depth: if compaction/handoff slipped onto the post-prompt queue
// alongside us (e.g. via a scheduler we don't own), refuse to start a fresh
// streaming turn — agent.continue() here would race the handoff's session
// reset. The first-class fix is in #checkCompaction/the agent_end handler,
// but this guard catches anything that bypasses that path.
if (this.isCompacting || this.isGeneratingHandoff) {
if (signal.aborted || this.#isDisposed || this.isCompacting || this.isGeneratingHandoff) {
options?.onSkip?.();
return;
}
@@ -2051,6 +2053,10 @@ export class AgentSession {
this.#beginInFlight();
try {
await this.#maybeRestoreRetryFallbackPrimary();
if (signal.aborted || this.#isDisposed) {
options?.onSkip?.();
return;
}
await this.agent.continue();
} catch (error) {
logger.warn("agent.continue failed after scheduling", {
@@ -3162,8 +3168,9 @@ export class AgentSession {
// session's dispose.
this.abortRetry();
this.abortCompaction();
const postPromptDrain = this.#cancelPostPromptTasks();
this.agent.abort();
await this.#cancelPostPromptTasks();
await postPromptDrain;
// Cancel jobs this agent registered so a subagent's teardown doesn't
// leak its background bash/task work into the parent's manager. Only
// the session that owns the manager goes on to dispose it (which itself
@@ -5061,9 +5068,25 @@ export class AgentSession {
}
#scheduleIdleQueueDrain(): void {
if (!this.#canAutoContinueForFollowUp()) return;
this.#scheduleQueuedMessageDrain();
}
#scheduleQueuedMessageDrain(): void {
if (this.#queuedMessageDrainScheduled || !this.#canAutoContinueForFollowUp() || !this.agent.hasQueuedMessages()) {
return;
}
this.#queuedMessageDrainScheduled = true;
this.#scheduleAgentContinue({
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
shouldContinue: () => {
this.#queuedMessageDrainScheduled = false;
return this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages();
},
onSkip: () => {
this.#queuedMessageDrainScheduled = false;
},
onError: () => {
this.#queuedMessageDrainScheduled = false;
},
});
}
@@ -5391,28 +5414,37 @@ export class AgentSession {
* abort. Omit it for internal/lifecycle aborts.
*/
async abort(options?: { goalReason?: "interrupted" | "internal"; reason?: string }): Promise<void> {
this.abortRetry();
this.#promptGeneration++;
this.#scheduledHiddenNextTurnGeneration = undefined;
this.abortCompaction();
this.abortHandoff();
this.abortBash();
this.abortEval();
const postPromptDrain = this.#cancelPostPromptTasks();
this.agent.abort(options?.reason);
await postPromptDrain;
await this.agent.waitForIdle();
await this.#goalRuntime.onTaskAborted({ reason: options?.goalReason ?? "interrupted" });
// Clear prompt-in-flight state: waitForIdle resolves when the agent loop's finally
// block runs, but nested prompt setup/finalizers may still be unwinding. Without this,
// a subsequent prompt() can incorrectly observe the session as busy after an abort.
this.#resetInFlight();
// Safety net: if the agent loop aborted without producing an assistant
// message (e.g. failed before the first stream), the in-flight yield was
// never resolved or rejected by the normal message_end path. Reject it now
// so any requeue callback still fires and the queue stays consistent.
if (this.#toolChoiceQueue.hasInFlight) {
this.#toolChoiceQueue.reject("aborted");
// Session switch/compact paths disconnect first; explicit aborts should
// leave any queued steer/follow-up visible for the user rather than
// auto-starting a fresh turn during cleanup.
this.#abortInProgress = true;
try {
this.abortRetry();
this.#promptGeneration++;
this.#scheduledHiddenNextTurnGeneration = undefined;
this.abortCompaction();
this.abortHandoff();
this.abortBash();
this.abortEval();
const postPromptDrain = this.#cancelPostPromptTasks();
this.agent.abort(options?.reason);
await postPromptDrain;
await this.agent.waitForIdle();
await this.#goalRuntime.onTaskAborted({ reason: options?.goalReason ?? "interrupted" });
// Clear prompt-in-flight state: waitForIdle resolves when the agent loop's finally
// block runs, but nested prompt setup/finalizers may still be unwinding. Without this,
// a subsequent prompt() can incorrectly observe the session as busy after an abort.
this.#resetInFlight();
// Safety net: if the agent loop aborted without producing an assistant
// message (e.g. failed before the first stream), the in-flight yield was
// never resolved or rejected by the normal message_end path. Reject it now
// so any requeue callback still fires and the queue stays consistent.
if (this.#toolChoiceQueue.hasInFlight) {
this.#toolChoiceQueue.reject("aborted");
}
} finally {
this.#abortInProgress = false;
this.#drainStrandedQueuedMessages();
}
}
@@ -173,6 +173,10 @@ export class IndexedSessionStorage implements SessionStorage {
}
}
writeTextAtomic(path: string, content: string): Promise<void> {
return this.writeText(path, content);
}
async rename(src: string, dst: string): Promise<void> {
await this.#awaitPath(src);
await this.#awaitPath(dst);
@@ -390,14 +394,7 @@ class IndexedSessionStorageWriter implements SessionStorageWriter {
return next;
}
writeLineSync(line: string): void {
if (this.#closed) throw new Error("Writer closed");
if (this.#error) throw this.#error;
const mtimeMs = this.#storage._appendForWriter(this.#path, line);
this.#trackPromise(this.#storage._queueAppend(this.#path, line, mtimeMs, () => this.#error));
}
async writeLine(line: string): Promise<void> {
async append(line: string): Promise<void> {
if (this.#closed) throw new Error("Writer closed");
if (this.#error) throw this.#error;
const mtimeMs = this.#storage._appendForWriter(this.#path, line);
@@ -410,15 +407,8 @@ class IndexedSessionStorageWriter implements SessionStorageWriter {
if (this.#error) throw this.#error;
}
async fsync(): Promise<void> {
await this.flush();
}
fsyncSync(): void {
// Indexed storage has no real fd to fsync; drain the pending chain
// synchronously is not possible, so this is a no-op. The async flush()
// above already ensures durability for the indexed backend.
if (this.#error) throw this.#error;
isOpen(): boolean {
return !this.#closed;
}
async close(): Promise<void> {
@@ -0,0 +1,352 @@
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import type { ProviderPayload, ServiceTier } from "@oh-my-pi/pi-ai";
import * as snapcompact from "@oh-my-pi/snapcompact";
import { createBranchSummaryMessage, createCompactionSummaryMessage, createCustomMessage } from "./messages";
import { type CompactionEntry, EPHEMERAL_MODEL_CHANGE_ROLE, type SessionEntry } from "./session-entries";
export interface SessionContext {
messages: AgentMessage[];
thinkingLevel?: string;
serviceTier?: ServiceTier;
/** Model roles: { default: "provider/modelId", small: "provider/modelId", ... } */
models: Record<string, string>;
/** Names of TTSR rules that have been injected this session */
injectedTtsrRules: string[];
/** MCP tool names selected through discovery for this session branch. */
selectedMCPToolNames: string[];
/** Whether this branch contains an explicit persisted MCP selection entry. */
hasPersistedMCPToolSelection: boolean;
/** Active mode (e.g. "plan") or "none" if no special mode is active */
mode: string;
/** Mode-specific data from the last mode_change entry */
modeData?: Record<string, unknown>;
}
/** Lists session model strings to try when restoring, in fallback order. */
export function getRestorableSessionModels(
models: Readonly<Record<string, string>>,
lastModelChangeRole: string | undefined,
): string[] {
const defaultModel = models.default;
if (
!lastModelChangeRole ||
lastModelChangeRole === "default" ||
lastModelChangeRole === EPHEMERAL_MODEL_CHANGE_ROLE
) {
return defaultModel ? [defaultModel] : [];
}
const roleModel = models[lastModelChangeRole];
if (!roleModel) return defaultModel ? [defaultModel] : [];
if (!defaultModel || roleModel === defaultModel) return [roleModel];
return [roleModel, defaultModel];
}
export function getLatestCompactionEntry(entries: SessionEntry[]): CompactionEntry | null {
for (let i = entries.length - 1; i >= 0; i--) {
if (entries[i].type === "compaction") {
return entries[i] as CompactionEntry;
}
}
return null;
}
export interface BuildSessionContextOptions {
/**
* Build the full-history display transcript instead of the LLM context:
* every path entry in chronological order, with each compaction emitted
* inline as a `compactionSummary` message at the position it fired rather
* than replacing the history before it. Display-only — never send the
* result to a provider.
*/
transcript?: boolean;
}
/**
* Build the session context from entries using tree traversal.
* If leafId is provided, walks from that entry to root.
* Handles compaction and branch summaries along the path.
*/
export function buildSessionContext(
entries: SessionEntry[],
leafId?: string | null,
byId?: Map<string, SessionEntry>,
options?: BuildSessionContextOptions,
): SessionContext {
// Build uuid index if not available
if (!byId) {
byId = new Map<string, SessionEntry>();
for (const entry of entries) {
byId.set(entry.id, entry);
}
}
// Find leaf
let leaf: SessionEntry | undefined;
if (leafId === null) {
// Explicitly null - return no messages (navigated to before first entry)
return {
messages: [],
thinkingLevel: "off",
serviceTier: undefined,
models: {},
injectedTtsrRules: [],
selectedMCPToolNames: [],
hasPersistedMCPToolSelection: false,
mode: "none",
};
}
if (leafId) {
leaf = byId.get(leafId);
}
if (!leaf) {
// Fallback to last entry (when leafId is undefined)
leaf = entries[entries.length - 1];
}
if (!leaf) {
return {
messages: [],
thinkingLevel: "off",
serviceTier: undefined,
models: {},
injectedTtsrRules: [],
selectedMCPToolNames: [],
hasPersistedMCPToolSelection: false,
mode: "none",
};
}
// Walk from leaf to root, collecting path
const path: SessionEntry[] = [];
let current: SessionEntry | undefined = leaf;
while (current) {
path.unshift(current);
current = current.parentId ? byId.get(current.parentId) : undefined;
}
// Extract settings and find compaction
let thinkingLevel: string | undefined = "off";
let serviceTier: ServiceTier | undefined;
const models: Record<string, string> = {};
let compaction: CompactionEntry | null = null;
const injectedTtsrRulesSet = new Set<string>();
let selectedMCPToolNames: string[] = [];
let hasPersistedMCPToolSelection = false;
let mode = "none";
let modeData: Record<string, unknown> | undefined;
// Track whether an explicit `model_change` with role="default" has been
// seen on this path. Once a user (or the agent itself) records an
// explicit default, later assistant-message inference must NOT overwrite
// it: temporary fallbacks (retry fallback, context promotion) and
// server-side model downgrades both produce assistant messages tagged
// with the wrong model id, which previously clobbered the user's pick on
// resume (issue #849).
let hasExplicitDefaultModel = false;
for (const entry of path) {
if (entry.type === "thinking_level_change") {
thinkingLevel = entry.thinkingLevel ?? "off";
} else if (entry.type === "model_change") {
// New format: { model: "provider/id", role?: string }
if (entry.model) {
const role = entry.role ?? "default";
models[role] = entry.model;
if (role === "default") {
hasExplicitDefaultModel = true;
}
}
} else if (entry.type === "service_tier_change") {
serviceTier = entry.serviceTier ?? undefined;
} else if (entry.type === "message" && entry.message.role === "assistant") {
// Legacy fallback: infer default model from assistant messages only
// when no explicit `model_change` (role=default) entry has been
// recorded yet. Newer sessions always record an explicit default
// model_change at the start of the conversation, so this branch is
// only used to keep pre-model_change sessions working.
if (!hasExplicitDefaultModel) {
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);
}
} else if (entry.type === "mcp_tool_selection") {
selectedMCPToolNames = [...entry.selectedToolNames];
hasPersistedMCPToolSelection = true;
} else if (entry.type === "mode_change") {
mode = entry.mode;
modeData = entry.data;
}
}
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)
// 2. Emit kept messages (from firstKeptEntryId up to compaction)
// 3. Emit messages after compaction
const messages: AgentMessage[] = [];
const appendMessage = (entry: SessionEntry) => {
if (entry.type === "message") {
messages.push(entry.message);
} else if (entry.type === "custom_message") {
messages.push(
createCustomMessage(
entry.customType,
entry.content,
entry.display,
entry.details,
entry.timestamp,
entry.attribution,
),
);
} else if (entry.type === "branch_summary" && entry.summary) {
messages.push(createBranchSummaryMessage(entry.summary, entry.fromId, entry.timestamp));
}
};
if (options?.transcript) {
// Display transcript: every entry in chronological order. Compactions do
// not erase prior history here — each renders inline (as a divider in the
// TUI) at the point it fired, with any snapcompact frames re-attached so
// the component can report them.
for (const entry of path) {
if (entry.type === "compaction") {
const snapcompactArchive = snapcompact.getPreservedArchive(entry.preserveData);
messages.push(
createCompactionSummaryMessage(
entry.summary,
entry.tokensBefore,
entry.timestamp,
entry.shortSummary,
undefined,
snapcompactArchive ? snapcompact.images(snapcompactArchive) : undefined,
),
);
} else {
appendMessage(entry);
}
}
} else if (compaction) {
const providerPayload: ProviderPayload | undefined = (() => {
const candidate = compaction.preserveData?.openaiRemoteCompaction;
if (!candidate || typeof candidate !== "object") return undefined;
const remote = candidate as { provider?: unknown; replacementHistory?: unknown };
if (typeof remote.provider !== "string" || remote.provider.length === 0) return undefined;
if (!Array.isArray(remote.replacementHistory)) return undefined;
return {
type: "openaiResponsesHistory",
provider: remote.provider,
items: remote.replacementHistory as Array<Record<string, unknown>>,
};
})();
const remoteReplacementHistory = providerPayload?.items;
// Emit summary first; re-attach any archived snapcompact frames so the
// model can keep reading the archived history after every context rebuild.
const snapcompactArchive = snapcompact.getPreservedArchive(compaction.preserveData);
messages.push(
createCompactionSummaryMessage(
compaction.summary,
compaction.tokensBefore,
compaction.timestamp,
compaction.shortSummary,
providerPayload,
snapcompactArchive ? snapcompact.images(snapcompactArchive) : undefined,
),
);
// Find compaction index in path
const compactionIdx = path.findIndex(e => e.type === "compaction" && e.id === compaction.id);
if (!remoteReplacementHistory) {
// Emit kept messages (before compaction, starting from firstKeptEntryId)
let foundFirstKept = false;
for (let i = 0; i < compactionIdx; i++) {
const entry = path[i];
if (entry.id === compaction.firstKeptEntryId) {
foundFirstKept = true;
}
if (foundFirstKept) {
appendMessage(entry);
}
}
}
// Emit messages after compaction
for (let i = compactionIdx + 1; i < path.length; i++) {
const entry = path[i];
appendMessage(entry);
}
} else {
// No compaction - emit all messages, handle branch summaries and custom messages
for (const entry of path) {
appendMessage(entry);
}
}
// Strip dangling tool_use blocks — a tool_use with no matching tool_result on the
// resolved leaf→root path — from ANY assistant turn, not just the trailing one.
// This happens whenever the leaf (or a branch point) lands such that an assistant
// turn's tool results are off the selected path: its result children live on a
// sibling branch, or it is the leaf itself (results are children below it). Left
// in place, `transformMessages` fabricates one synthetic "aborted"/"No result
// provided" result per dangling call, which render as phantom failed calls and
// re-inject the failed batch into the model's
// context — the rewind/restore loop.
//
// Stripping is necessary but not sufficient: a *modified* assistant turn that still
// carries signed `thinking`/`redacted_thinking` is rejected by Anthropic — "thinking
// blocks in the latest assistant message cannot be modified", and signed thinking
// replayed out of its original turn shape can also fail signature validation (this
// bites the handoff/branch-summary request). So when we rewrite a turn we also
// neutralize its protected reasoning: drop `redactedThinking` (encrypted, no
// plaintext to keep) and clear `thinking` signatures so the provider encoder
// downgrades them to plain text (verified accepted by the live API), preserving the
// visible reasoning while removing the immutability/invalid-signature hazard. Drop a
// turn left with no content. (Live turns never qualify: their results are persisted
// on the same path before any context rebuild.)
const pairedToolResultIds = new Set<string>();
for (const message of messages) {
if (message.role === "toolResult") pairedToolResultIds.add(message.toolCallId);
}
for (let i = messages.length - 1; i >= 0; i--) {
const message = messages[i];
if (message.role !== "assistant") continue;
const hasDangling = message.content.some(
block => block.type === "toolCall" && !pairedToolResultIds.has(block.id),
);
if (!hasDangling) continue;
const normalized = message.content
.filter(
block =>
!(block.type === "toolCall" && !pairedToolResultIds.has(block.id)) && block.type !== "redactedThinking",
)
.map(block =>
block.type === "thinking" && block.thinkingSignature ? { ...block, thinkingSignature: undefined } : block,
);
if (normalized.length === 0) {
messages.splice(i, 1);
} else {
messages[i] = { ...message, content: normalized };
}
}
return {
messages,
thinkingLevel,
serviceTier,
models,
injectedTtsrRules,
selectedMCPToolNames,
hasPersistedMCPToolSelection,
mode,
modeData,
};
}
@@ -0,0 +1,194 @@
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import type { ImageContent, MessageAttribution, ServiceTier, TextContent } from "@oh-my-pi/pi-ai";
export const CURRENT_SESSION_VERSION = 3;
export const EPHEMERAL_MODEL_CHANGE_ROLE = "fallback";
export interface SessionHeader {
type: "session";
version?: number; // v1 sessions don't have this
id: string;
title?: string; // Auto-generated title from first message
titleSource?: "auto" | "user";
timestamp: string;
cwd: string;
parentSession?: string;
}
export interface NewSessionOptions {
parentSession?: string;
/** Skip flushing the current session and delete it instead of saving. */
drop?: boolean;
}
export interface SessionEntryBase {
type: string;
id: string;
parentId: string | null;
timestamp: string;
}
export interface SessionMessageEntry extends SessionEntryBase {
type: "message";
message: AgentMessage;
}
export interface ThinkingLevelChangeEntry extends SessionEntryBase {
type: "thinking_level_change";
thinkingLevel?: string | null;
}
export interface ModelChangeEntry extends SessionEntryBase {
type: "model_change";
/** Model in "provider/modelId" format */
model: string;
/** Role: "default", "smol", "slow", etc. Undefined treated as "default" */
role?: string;
}
export interface ServiceTierChangeEntry extends SessionEntryBase {
type: "service_tier_change";
serviceTier: ServiceTier | null;
}
export interface CompactionEntry<T = unknown> extends SessionEntryBase {
type: "compaction";
summary: string;
shortSummary?: string;
firstKeptEntryId: string;
tokensBefore: number;
/** Extension-specific data (e.g., ArtifactIndex, version markers for structured compaction) */
details?: T;
/** Hook-provided data to persist across compaction */
preserveData?: Record<string, unknown>;
/** True if generated by an extension, undefined/false if pi-generated (backward compatible) */
fromExtension?: boolean;
}
export interface BranchSummaryEntry<T = unknown> extends SessionEntryBase {
type: "branch_summary";
fromId: string;
summary: string;
/** Extension-specific data (not sent to LLM) */
details?: T;
/** True if generated by an extension, false if pi-generated */
fromExtension?: boolean;
}
/**
* Custom entry for extensions to store extension-specific data in the session.
* Use customType to identify your extension's entries.
*
* Purpose: Persist extension state across session reloads. On reload, extensions can
* scan entries for their customType and reconstruct internal state.
*
* Does NOT participate in LLM context (ignored by buildSessionContext).
* For injecting content into context, see CustomMessageEntry.
*/
export interface CustomEntry<T = unknown> extends SessionEntryBase {
type: "custom";
customType: string;
data?: T;
}
/** Label entry for user-defined bookmarks/markers on entries. */
export interface LabelEntry extends SessionEntryBase {
type: "label";
targetId: string;
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[];
}
/** Persisted MCP discovery selection state for a session branch. */
export interface MCPToolSelectionEntry extends SessionEntryBase {
type: "mcp_tool_selection";
/** MCP tool names selected for visibility in discovery mode. */
selectedToolNames: string[];
}
/** Session init entry - captures initial context for subagent sessions (debugging/replay). */
export interface SessionInitEntry extends SessionEntryBase {
type: "session_init";
/** Full system prompt sent to the model */
systemPrompt: string;
/** Initial task/user message */
task: string;
/** Tools available to the agent */
tools: string[];
/** Output schema if structured output was requested */
outputSchema?: unknown;
}
/** Mode change entry - tracks agent mode transitions (e.g. plan mode). */
export interface ModeChangeEntry extends SessionEntryBase {
type: "mode_change";
/** Current mode name, or "none" when exiting a mode */
mode: string;
/** Optional mode-specific data (e.g. plan file path) */
data?: Record<string, unknown>;
}
/**
* Custom message entry for extensions to inject messages into LLM context.
* Use customType to identify your extension's entries.
*
* Unlike CustomEntry, this DOES participate in LLM context.
* The content participates in LLM context through convertToLlm().
* Use details for extension-specific metadata (not sent to LLM).
*
* display controls TUI rendering:
* - false: hidden entirely
* - true: rendered with distinct styling (different from user messages)
*/
export interface CustomMessageEntry<T = unknown> extends SessionEntryBase {
type: "custom_message";
customType: string;
content: string | (TextContent | ImageContent)[];
details?: T;
display: boolean;
/** Who initiated this message for billing/attribution semantics. */
attribution?: MessageAttribution;
}
/** Session entry - has id/parentId for tree structure (returned by "read" methods in SessionManager) */
export type SessionEntry =
| SessionMessageEntry
| ThinkingLevelChangeEntry
| ModelChangeEntry
| ServiceTierChangeEntry
| CompactionEntry
| BranchSummaryEntry
| CustomEntry
| CustomMessageEntry
| LabelEntry
| TtsrInjectionEntry
| MCPToolSelectionEntry
| SessionInitEntry
| ModeChangeEntry;
/** Raw file entry (includes header) */
export type FileEntry = SessionHeader | SessionEntry;
/** Tree node for getTree() - defensive copy of session structure */
export interface SessionTreeNode {
entry: SessionEntry;
children: SessionTreeNode[];
/** Resolved label for this entry, if any */
label?: string;
}
export interface UsageStatistics {
input: number;
output: number;
cacheRead: number;
cacheWrite: number;
premiumRequests: number;
cost: number;
}
@@ -0,0 +1,588 @@
import * as os from "node:os";
import * as path from "node:path";
import type { Message, TextContent } from "@oh-my-pi/pi-ai";
import { getAgentDir as getDefaultAgentDir, logger, parseJsonlLenient, toError } from "@oh-my-pi/pi-utils";
import { computeDefaultSessionDir } from "./session-paths";
import { FileSessionStorage, type SessionStorage } from "./session-storage";
/**
* Coarse lifecycle status of a session, derived from its last persisted message.
*
* - `complete` — the last assistant turn ended with no unanswered tool calls, i.e.
* the agent yielded control back to the user.
* - `interrupted` — work was cut off mid-flight: a trailing assistant turn with
* pending tool calls, a trailing tool result the agent never continued from, or
* a length-truncated turn.
* - `aborted` — the last assistant turn was cancelled by the user.
* - `error` — the last assistant turn ended in an error.
* - `pending` — a trailing user message with no assistant reply persisted after it.
* - `unknown` — status could not be determined (empty/header-only session, or the
* final message was larger than the tail window that was read).
*/
export type SessionStatus = "complete" | "interrupted" | "aborted" | "error" | "pending" | "unknown";
export interface SessionInfo {
path: string;
id: string;
/** Working directory where the session was started. Empty string for old sessions. */
cwd: string;
title?: string;
/** Path to the parent session (if this session was forked). */
parentSessionPath?: string;
created: Date;
modified: Date;
messageCount: number;
/** File size in bytes on disk; used for compact list rendering. */
size: number;
firstMessage: string;
allMessagesText: string;
/**
* Coarse lifecycle status from the session's last persisted message. Optional:
* synthesized {@link SessionInfo}s (cross-project stubs, tests) leave it unset.
*/
status?: SessionStatus;
}
export interface ResolvedSessionMatch {
session: SessionInfo;
scope: "local" | "global";
}
/** Lightweight metadata for a recent session, used in welcome/picker UI. */
export interface RecentSessionInfo {
path: string;
name: string;
timeAgo: string;
}
const SESSION_LIST_PREFIX_BYTES = 4096;
/**
* Tail window read to derive {@link SessionStatus}. Large enough to capture a
* typical final assistant turn (thinking + text); when the final message exceeds
* it the status falls back to `unknown` rather than misreporting.
*/
const SESSION_LIST_SUFFIX_BYTES = 32_768;
const SESSION_LIST_PARALLEL_THRESHOLD = 64;
const SESSION_LIST_MAX_WORKERS = 16;
function sanitizeSessionName(value: string | undefined): string | undefined {
if (!value) return undefined;
const firstLine = value.split(/\r?\n/)[0] ?? "";
const stripped = firstLine.replace(/[\x00-\x1F\x7F]/g, "");
const trimmed = stripped.trim();
return trimmed.length > 0 ? trimmed : undefined;
}
/** Format a time difference as a human-readable string */
function formatTimeAgo(date: Date): string {
const now = Date.now();
const diffMs = now - date.getTime();
const diffMins = Math.floor(diffMs / 60000);
const diffHours = Math.floor(diffMs / 3600000);
const diffDays = Math.floor(diffMs / 86400000);
if (diffMins < 1) return "just now";
if (diffMins < 60) return `${diffMins}m ago`;
if (diffHours < 24) return `${diffHours}h ago`;
if (diffDays < 7) return `${diffDays}d ago`;
return date.toLocaleDateString();
}
/**
* Friendly display name for a session: explicit title, then first user prompt,
* then a timestamp-based label. The raw UUID `id` is intentionally never used —
* it is unfriendly and indistinguishable from neighboring sessions in the UI.
*/
function sessionDisplayName(info: SessionInfo): string {
const title = sanitizeSessionName(info.title);
if (title) return title;
const first =
info.firstMessage && info.firstMessage !== "(no messages)" ? sanitizeSessionName(info.firstMessage) : undefined;
if (first) return first;
const created = info.created.getTime();
const ts = Number.isFinite(created) ? created : info.modified.getTime();
const date = new Date(ts);
const time = date.toLocaleTimeString(undefined, { hour: "2-digit", minute: "2-digit" });
return `Untitled · ${time}`;
}
function extractTextFromContent(content: Message["content"]): string {
if (typeof content === "string") return content;
return content
.filter((block): block is TextContent => block.type === "text")
.map(block => block.text)
.join(" ");
}
/**
* Derive a {@link SessionStatus} from a tail window of a session file. Entries are
* newline-terminated on write, so within the window only the first line can be a
* partial fragment — it simply fails to parse and is skipped. We walk backwards to
* the last `message` entry and classify by its role / stop reason.
*/
function deriveSessionStatus(suffix: string): SessionStatus {
if (!suffix) return "unknown";
const lines = suffix.split("\n");
for (let i = lines.length - 1; i >= 0; i--) {
const line = lines[i];
// Every persisted entry is `JSON.stringify(obj)` → starts with `{`. This
// cheaply rejects blank lines and the leading partial fragment without
// attempting to parse a multi-KB tail of a truncated line.
if (line.charCodeAt(0) !== 123) continue;
let entry: { type?: string; message?: TailMessage };
try {
entry = JSON.parse(line);
} catch {
continue;
}
if (entry.type === "message" && entry.message) {
return statusFromTailMessage(entry.message);
}
}
return "unknown";
}
interface TailMessage {
role?: string;
stopReason?: string;
content?: unknown;
}
function isToolCallBlock(block: unknown): boolean {
return typeof block === "object" && block !== null && (block as { type?: unknown }).type === "toolCall";
}
function statusFromTailMessage(message: TailMessage): SessionStatus {
switch (message.role) {
case "assistant": {
switch (message.stopReason) {
case "error":
return "error";
case "aborted":
return "aborted";
case "length":
return "interrupted";
}
// A turn that ends without unanswered tool calls means the agent yielded
// control back to the user — complete. Trailing tool calls (no tool
// results after) mean the loop was cut off before running them.
const content = message.content;
if (Array.isArray(content) && content.some(isToolCallBlock)) return "interrupted";
return "complete";
}
case "toolResult":
// Tools ran but the agent never produced the following assistant turn.
return "interrupted";
case "user":
// User message with no assistant reply persisted after it.
return "pending";
default:
return "unknown";
}
}
function decodeJsonStringFragment(value: string): string {
const safeValue = value.endsWith("\\") ? value.slice(0, -1) : value;
try {
return JSON.parse(`"${safeValue}"`) as string;
} catch {
return safeValue
.replace(/\\n/g, "\n")
.replace(/\\r/g, "\r")
.replace(/\\t/g, "\t")
.replace(/\\"/g, '"')
.replace(/\\\\/g, "\\");
}
}
function extractStringProperty(source: string, name: string, startIndex = 0): string | undefined {
const propertyIndex = source.indexOf(`"${name}"`, startIndex);
if (propertyIndex === -1) return undefined;
const colonIndex = source.indexOf(":", propertyIndex + name.length + 2);
if (colonIndex === -1) return undefined;
let valueIndex = colonIndex + 1;
while (valueIndex < source.length) {
const char = source.charCodeAt(valueIndex);
if (char !== 32 && char !== 9 && char !== 10 && char !== 13) break;
valueIndex++;
}
if (source.charCodeAt(valueIndex) !== 34) return undefined;
const valueStart = valueIndex + 1;
let escaped = false;
for (let i = valueStart; i < source.length; i++) {
const char = source.charCodeAt(i);
if (escaped) {
escaped = false;
continue;
}
if (char === 92) {
escaped = true;
continue;
}
if (char === 34) {
return decodeJsonStringFragment(source.slice(valueStart, i));
}
}
return decodeJsonStringFragment(source.slice(valueStart));
}
function countMessageMarkers(content: string): number {
let count = 0;
let index = 0;
while (index < content.length) {
const typeIndex = content.indexOf('"type"', index);
if (typeIndex === -1) break;
const colonIndex = content.indexOf(":", typeIndex + 6);
if (colonIndex === -1) break;
const type = extractStringProperty(content, "type", typeIndex);
if (type === "message") count++;
index = colonIndex + 1;
}
return count;
}
function extractFirstUserMessageFromPrefix(content: string): string | undefined {
const roleIndex = content.indexOf('"role"');
if (roleIndex === -1) return undefined;
let index = roleIndex;
while (index !== -1) {
const role = extractStringProperty(content, "role", index);
if (role === "user") {
return extractStringProperty(content, "content", index) ?? extractStringProperty(content, "text", index);
}
index = content.indexOf('"role"', index + 6);
}
return undefined;
}
interface SessionListHeader {
type: "session";
id: string;
cwd?: string;
title?: string;
parentSession?: string;
timestamp?: string;
}
function parseSessionListHeader(
content: string,
entries: Array<Record<string, unknown>>,
): SessionListHeader | undefined {
const parsedHeader = entries[0];
if (parsedHeader?.type === "session" && typeof parsedHeader.id === "string") {
return {
type: "session",
id: parsedHeader.id,
cwd: typeof parsedHeader.cwd === "string" ? parsedHeader.cwd : undefined,
title: typeof parsedHeader.title === "string" ? parsedHeader.title : undefined,
parentSession: typeof parsedHeader.parentSession === "string" ? parsedHeader.parentSession : undefined,
timestamp: typeof parsedHeader.timestamp === "string" ? parsedHeader.timestamp : undefined,
};
}
const firstLineEnd = content.indexOf("\n");
const firstLine = firstLineEnd === -1 ? content : content.slice(0, firstLineEnd);
if (extractStringProperty(firstLine, "type") !== "session") return undefined;
const id = extractStringProperty(firstLine, "id");
if (!id) return undefined;
return {
type: "session",
id,
cwd: extractStringProperty(firstLine, "cwd"),
title: extractStringProperty(firstLine, "title"),
parentSession: extractStringProperty(firstLine, "parentSession"),
timestamp: extractStringProperty(firstLine, "timestamp"),
};
}
function getSessionListWorkerCount(fileCount: number): number {
if (fileCount <= SESSION_LIST_PARALLEL_THRESHOLD) return 1;
return Math.min(
SESSION_LIST_MAX_WORKERS,
os.availableParallelism(),
Math.ceil(fileCount / SESSION_LIST_PARALLEL_THRESHOLD),
);
}
/**
* Scan a single session file into a {@link SessionInfo}. Always reads the 4 KB
* header/first-message prefix; only reads the 32 KB tail window (and derives
* {@link SessionStatus}) when `withStatus` is set — the recent/most-recent
* lookups skip it.
*/
async function scanSessionFile(
file: string,
storage: SessionStorage,
withStatus: boolean,
): Promise<SessionInfo | undefined> {
try {
const stat = storage.statSync(file);
const [content, suffix] = await storage.readTextSlices(
file,
SESSION_LIST_PREFIX_BYTES,
withStatus ? SESSION_LIST_SUFFIX_BYTES : 0,
);
const { size, mtime } = stat;
const entries = parseJsonlLenient<Record<string, unknown>>(content);
const header = parseSessionListHeader(content, entries);
if (!header) return undefined;
let parsedMessageCount = 0;
let firstMessage = "";
const allMessages: string[] = [];
let shortSummary: string | undefined;
for (let i = 1; i < entries.length; i++) {
const entry = entries[i] as { type?: string; message?: Message; shortSummary?: string };
if (entry.type === "compaction" && typeof entry.shortSummary === "string") {
shortSummary = entry.shortSummary;
}
if (entry.type === "message" && entry.message) {
parsedMessageCount++;
if (entry.message.role === "user" || entry.message.role === "assistant") {
const textContent = extractTextFromContent(entry.message.content);
if (textContent) {
allMessages.push(textContent);
if (!firstMessage && entry.message.role === "user") {
firstMessage = textContent;
}
}
}
}
}
firstMessage ||= extractFirstUserMessageFromPrefix(content) ?? "";
const messageCount = Math.max(parsedMessageCount, countMessageMarkers(content));
return {
path: file,
id: header.id,
cwd: header.cwd ?? "",
title: header.title ?? shortSummary,
parentSessionPath: header.parentSession,
created: new Date(header.timestamp ?? ""),
modified: mtime,
messageCount,
size,
firstMessage: firstMessage || "(no messages)",
allMessagesText: allMessages.length > 0 ? allMessages.join(" ") : firstMessage,
status: withStatus ? deriveSessionStatus(suffix) : undefined,
};
} catch {
return undefined;
}
}
async function collectSessionsFromFileStride(
files: string[],
storage: SessionStorage,
startIndex: number,
stride: number,
withStatus: boolean,
): Promise<SessionInfo[]> {
const sessions: SessionInfo[] = [];
for (let i = startIndex; i < files.length; i += stride) {
const session = await scanSessionFile(files[i], storage, withStatus);
if (session) sessions.push(session);
}
return sessions;
}
async function collectSessionsFromFiles(
files: string[],
storage: SessionStorage,
withStatus: boolean,
): Promise<SessionInfo[]> {
const workerCount = getSessionListWorkerCount(files.length);
const sessions =
workerCount === 1
? await collectSessionsFromFileStride(files, storage, 0, 1, withStatus)
: (
await Promise.all(
Array.from({ length: workerCount }, (_, workerIndex) =>
collectSessionsFromFileStride(files, storage, workerIndex, workerCount, withStatus),
),
)
).flat();
sessions.sort((a, b) => b.modified.getTime() - a.modified.getTime());
return sessions;
}
/**
* Promote orphaned `<basename>.jsonl.<snowflake>.bak` backups created by the
* EPERM-rewrite path back to their primary path when the primary is missing.
* This runs once per session-dir scan, before the main `*.jsonl` glob, so a
* crash between the two renames in the EPERM-rewrite path does not leave the
* user's last good state stranded outside the loader's view.
*
* Exported for testing.
*/
export async function recoverOrphanedBackups(sessionDir: string, storage: SessionStorage): Promise<void> {
let backups: string[];
try {
backups = storage.listFilesSync(sessionDir, "*.bak");
} catch {
return;
}
if (backups.length === 0) return;
// For each primary path, pick the newest backup (highest mtime) as the recovery source.
const candidates = new Map<string, { backup: string; mtimeMs: number }>();
for (const backup of backups) {
const name = path.basename(backup);
// Expect "<primary>.<snowflake>.bak" where <primary> ends in ".jsonl".
if (!name.endsWith(".bak")) continue;
const trimmed = name.slice(0, -".bak".length);
const dotIdx = trimmed.lastIndexOf(".");
if (dotIdx <= 0) continue;
const primaryName = trimmed.slice(0, dotIdx);
if (!primaryName.endsWith(".jsonl")) continue;
const primaryPath = path.join(sessionDir, primaryName);
let mtimeMs = 0;
try {
mtimeMs = storage.statSync(backup).mtimeMs;
} catch {
continue;
}
const existing = candidates.get(primaryPath);
if (!existing || mtimeMs > existing.mtimeMs) {
candidates.set(primaryPath, { backup, mtimeMs });
}
}
for (const [primaryPath, { backup }] of candidates) {
if (storage.existsSync(primaryPath)) continue;
try {
await storage.rename(backup, primaryPath);
logger.warn("Recovered orphaned session backup", {
sessionFile: primaryPath,
backupPath: backup,
});
} catch (err) {
logger.warn("Failed to recover orphaned session backup", {
sessionFile: primaryPath,
backupPath: backup,
error: toError(err).message,
});
}
}
}
async function scanSessionDir(
sessionDir: string,
storage: SessionStorage,
withStatus: boolean,
): Promise<SessionInfo[]> {
try {
await recoverOrphanedBackups(sessionDir, storage);
const files = storage.listFilesSync(sessionDir, "*.jsonl");
return await collectSessionsFromFiles(files, storage, withStatus);
} catch {
return [];
}
}
/**
* List sessions in a resolved session directory (newest first), reading each
* file's lifecycle {@link SessionStatus}.
*/
export function listSessions(sessionDir: string, storage: SessionStorage): Promise<SessionInfo[]> {
return scanSessionDir(sessionDir, storage, true);
}
/** List all sessions across all project directories (newest first). */
export async function listAllSessions(storage: SessionStorage = new FileSessionStorage()): Promise<SessionInfo[]> {
const sessionsRoot = path.join(getDefaultAgentDir(), "sessions");
try {
const files = await Array.fromAsync(new Bun.Glob("*/*.jsonl").scan(sessionsRoot), name =>
path.join(sessionsRoot, name),
);
return await collectSessionsFromFiles(files, storage, true);
} catch {
return [];
}
}
/** Exported for testing */
export async function findMostRecentSession(
sessionDir: string,
storage: SessionStorage = new FileSessionStorage(),
): Promise<string | null> {
const sessions = await scanSessionDir(sessionDir, storage, false);
return sessions[0]?.path ?? null;
}
/** Get recent sessions for display in the welcome screen. */
export async function getRecentSessions(
sessionDir: string,
limit = 4,
storage: SessionStorage = new FileSessionStorage(),
): Promise<RecentSessionInfo[]> {
const sessions = await scanSessionDir(sessionDir, storage, false);
const recent: RecentSessionInfo[] = [];
for (let i = 0; i < sessions.length && i < limit; i++) {
const info = sessions[i];
recent.push({ path: info.path, name: sessionDisplayName(info), timeAgo: formatTimeAgo(info.modified) });
}
return recent;
}
function sessionMatchesResumeArg(session: SessionInfo, sessionArg: string): boolean {
const normalizedArg = sessionArg.toLowerCase();
const normalizedId = session.id.toLowerCase();
if (normalizedId.startsWith(normalizedArg)) {
return true;
}
const fileName = path.basename(session.path, ".jsonl").toLowerCase();
if (fileName.startsWith(normalizedArg)) {
return true;
}
const separator = fileName.lastIndexOf("_");
if (separator < 0) {
return false;
}
const fileSessionId = fileName.slice(separator + 1);
return fileSessionId.startsWith(normalizedArg);
}
export async function resolveResumableSession(
sessionArg: string,
cwd: string,
sessionDir?: string,
storage: SessionStorage = new FileSessionStorage(),
): Promise<ResolvedSessionMatch | undefined> {
const localSessionDir = sessionDir ?? computeDefaultSessionDir(cwd, storage);
const localSessions = await listSessions(localSessionDir, storage);
const localMatch = localSessions.find(session => sessionMatchesResumeArg(session, sessionArg));
if (localMatch) {
return { session: localMatch, scope: "local" };
}
if (sessionDir) {
return undefined;
}
const globalSessions = await listAllSessions(storage);
const globalMatch = globalSessions.find(session => sessionMatchesResumeArg(session, sessionArg));
if (!globalMatch) {
return undefined;
}
return { session: globalMatch, scope: "global" };
}
@@ -0,0 +1,106 @@
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { getBlobsDir, isEnoent, parseJsonlLenient } from "@oh-my-pi/pi-utils";
import { BlobStore, isBlobRef, resolveImageData, resolveImageDataUrl } from "./blob-store";
import { buildSessionContext } from "./session-context";
import type { FileEntry, SessionEntry, SessionHeader } from "./session-entries";
import { migrateToCurrentVersion } from "./session-migrations";
import { isImageBlock } from "./session-persistence";
import { FileSessionStorage, type SessionStorage } from "./session-storage";
/** Exported for compaction.test.ts */
export function parseSessionEntries(content: string): FileEntry[] {
return parseJsonlLenient<FileEntry>(content);
}
/** Exported for testing */
export async function loadEntriesFromFile(
filePath: string,
storage: SessionStorage = new FileSessionStorage(),
): Promise<FileEntry[]> {
let content: string;
try {
content = await storage.readText(filePath);
} catch (err) {
if (isEnoent(err)) return [];
throw err;
}
const entries = parseJsonlLenient<FileEntry>(content);
// Validate session header
if (entries.length === 0) return entries;
const header = entries[0] as SessionHeader;
if (header.type !== "session" || typeof header.id !== "string") {
return [];
}
return entries;
}
/**
* Resolve blob references in loaded entries, restoring both session image blocks and persisted
* provider image URLs back to the inline data expected by downstream transports. Mutates entries in place.
*/
function hasImageUrl(value: unknown): value is { image_url: string } {
return typeof value === "object" && value !== null && "image_url" in value && typeof value.image_url === "string";
}
async function resolvePersistedImageUrlRefs(value: unknown, blobStore: BlobStore): Promise<void> {
if (Array.isArray(value)) {
await Promise.all(value.map(item => resolvePersistedImageUrlRefs(item, blobStore)));
return;
}
if (typeof value !== "object" || value === null) return;
if (hasImageUrl(value) && isBlobRef(value.image_url)) {
value.image_url = await resolveImageDataUrl(blobStore, value.image_url);
}
await Promise.all(Object.values(value).map(item => resolvePersistedImageUrlRefs(item, blobStore)));
}
export async function resolveBlobRefsInEntries(entries: FileEntry[], blobStore: BlobStore): Promise<void> {
const promises: Promise<void>[] = [];
for (const entry of entries) {
if (entry.type === "session") continue;
let contentArray: unknown[] | undefined;
if (entry.type === "message" && "content" in entry.message && Array.isArray(entry.message.content)) {
contentArray = entry.message.content;
} else if (entry.type === "custom_message" && Array.isArray(entry.content)) {
contentArray = entry.content;
}
if (contentArray) {
for (const block of contentArray) {
if (isImageBlock(block) && isBlobRef(block.data)) {
promises.push(
resolveImageData(blobStore, block.data).then(resolved => {
block.data = resolved;
}),
);
}
}
}
promises.push(resolvePersistedImageUrlRefs(entry, blobStore));
}
await Promise.all(promises);
}
/**
* Read-only message view of a session file: load entries, migrate to the
* current version, resolve blob refs, and build the context along the
* persisted leaf path (last entry). Does NOT create a writer or take the
* session lock — safe to call against a file another session is writing.
*/
export async function loadSessionMessagesReadOnly(filePath: string): Promise<AgentMessage[]> {
const entries = await loadEntriesFromFile(filePath);
if (entries.length === 0) return [];
migrateToCurrentVersion(entries);
await resolveBlobRefsInEntries(entries, new BlobStore(getBlobsDir()));
const sessionEntries = entries.filter((e): e is SessionEntry => e.type !== "session");
return buildSessionContext(sessionEntries).messages;
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,78 @@
import { Snowflake } from "@oh-my-pi/pi-utils";
import { type CompactionEntry, CURRENT_SESSION_VERSION, type FileEntry, type SessionHeader } from "./session-entries";
/** Generate a unique short ID (8 hex chars, collision-checked) */
export function generateId(byId: { has(id: string): boolean }): string {
for (let i = 0; i < 100; i++) {
const id = crypto.randomUUID().slice(-8);
if (!byId.has(id)) return id;
}
return Snowflake.next(); // fallback to full snowflake id
}
/** Migrate v1 → v2: add id/parentId tree structure. Mutates in place. */
function migrateV1ToV2(entries: FileEntry[]): void {
const ids = new Set<string>();
let prevId: string | null = null;
for (const entry of entries) {
if (entry.type === "session") {
entry.version = 2;
continue;
}
entry.id = generateId(ids);
entry.parentId = prevId;
prevId = entry.id;
// Convert firstKeptEntryIndex to firstKeptEntryId for compaction
if (entry.type === "compaction") {
const comp = entry as CompactionEntry & { firstKeptEntryIndex?: number };
if (typeof comp.firstKeptEntryIndex === "number") {
const targetEntry = entries[comp.firstKeptEntryIndex];
if (targetEntry && targetEntry.type !== "session") {
comp.firstKeptEntryId = targetEntry.id;
}
delete comp.firstKeptEntryIndex;
}
}
}
}
/** Migrate v2 → v3: rename hookMessage role to custom. Mutates in place. */
function migrateV2ToV3(entries: FileEntry[]): void {
for (const entry of entries) {
if (entry.type === "session") {
entry.version = 3;
continue;
}
if (entry.type === "message") {
const msg = entry.message as { role?: string };
if (msg.role === "hookMessage") {
(entry.message as { role: string }).role = "custom";
}
}
}
}
/**
* Run all necessary migrations to bring entries to current version.
* Mutates entries in place. Returns true if any migration was applied.
*/
export function migrateToCurrentVersion(entries: FileEntry[]): boolean {
const header = entries.find(e => e.type === "session") as SessionHeader | undefined;
const version = header?.version ?? 1;
if (version >= CURRENT_SESSION_VERSION) return false;
if (version < 2) migrateV1ToV2(entries);
if (version < 3) migrateV2ToV3(entries);
return true;
}
/** Exported for testing */
export function migrateSessionEntries(entries: FileEntry[]): void {
migrateToCurrentVersion(entries);
}
@@ -0,0 +1,193 @@
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import { getTerminalId } from "@oh-my-pi/pi-tui";
import { getSessionsDir, getTerminalSessionsDir, isEnoent, logger, resolveEquivalentPath } from "@oh-my-pi/pi-utils";
import type { SessionStorage } from "./session-storage";
const migratedSessionRoots = new Set<string>();
/**
* Merge or rename a legacy session directory into its canonical target.
* Best effort: callers decide whether migration failures should surface.
*/
function migrateSessionDirPath(oldPath: string, newPath: string): void {
const existing = fs.statSync(newPath, { throwIfNoEntry: false });
if (existing?.isDirectory()) {
for (const file of fs.readdirSync(oldPath)) {
const src = path.join(oldPath, file);
const dst = path.join(newPath, file);
if (!fs.existsSync(dst)) {
fs.renameSync(src, dst);
}
}
fs.rmSync(oldPath, { recursive: true, force: true });
return;
}
if (existing) {
fs.rmSync(newPath, { recursive: true, force: true });
}
fs.renameSync(oldPath, newPath);
}
function encodeLegacyAbsoluteSessionDirName(cwd: string): string {
const resolvedCwd = path.resolve(cwd);
return `--${resolvedCwd.replace(/^[/\\]/, "").replace(/[/\\:]/g, "-")}--`;
}
function encodeRelativeSessionDirName(prefix: string, relative: string): string {
const encoded = relative.replace(/[/\\:]/g, "-");
return encoded ? (prefix.endsWith("-") ? `${prefix}${encoded}` : `${prefix}-${encoded}`) : prefix;
}
function getDefaultSessionDirName(cwd: string): { encodedDirName: string; resolvedCwd: string } {
const resolvedCwd = path.resolve(cwd);
const canonicalCwd = resolveEquivalentPath(resolvedCwd);
const home = os.homedir();
const canonicalHome = resolveEquivalentPath(home);
const tempRoot = os.tmpdir();
const canonicalTempRoot = resolveEquivalentPath(tempRoot);
const homeRelative = path.relative(canonicalHome, canonicalCwd);
const tempRelative = path.relative(canonicalTempRoot, canonicalCwd);
const encodedDirName =
homeRelative === "" || (!homeRelative.startsWith("..") && !path.isAbsolute(homeRelative))
? encodeRelativeSessionDirName("-", homeRelative)
: tempRelative === "" || (!tempRelative.startsWith("..") && !path.isAbsolute(tempRelative))
? encodeRelativeSessionDirName("-tmp", tempRelative)
: encodeLegacyAbsoluteSessionDirName(canonicalCwd);
return { encodedDirName, resolvedCwd };
}
/**
* Migrate old `--<home-encoded>-*--` session dirs to the new `-*` format.
* Runs once per sessions root on first access, best-effort.
*/
function migrateHomeSessionDirs(sessionsRoot: string): void {
if (migratedSessionRoots.has(sessionsRoot)) return;
migratedSessionRoots.add(sessionsRoot);
const home = os.homedir();
const homeEncoded = home.replace(/^[/\\]/, "").replace(/[/\\:]/g, "-");
const oldPrefix = `--${homeEncoded}-`;
const oldExact = `--${homeEncoded}--`;
let entries: string[];
try {
entries = fs.readdirSync(sessionsRoot);
} catch {
return;
}
for (const entry of entries) {
let remainder: string;
if (entry === oldExact) {
remainder = "";
} else if (entry.startsWith(oldPrefix) && entry.endsWith("--")) {
remainder = entry.slice(oldPrefix.length, -2);
} else {
continue;
}
const newName = remainder ? `-${remainder}` : "-";
const oldPath = path.join(sessionsRoot, entry);
const newPath = path.join(sessionsRoot, newName);
try {
migrateSessionDirPath(oldPath, newPath);
} catch {
// Best effort
}
}
}
function migrateLegacyAbsoluteSessionDir(cwd: string, sessionDir: string, sessionsRoot: string): void {
const legacyDir = path.join(sessionsRoot, encodeLegacyAbsoluteSessionDirName(cwd));
if (legacyDir === sessionDir || !fs.existsSync(legacyDir)) return;
try {
migrateSessionDirPath(legacyDir, sessionDir);
} catch {
// Best effort
}
}
export function resolveManagedSessionRoot(sessionDir: string, cwd: string): string | undefined {
const currentDirName = path.basename(sessionDir);
const { encodedDirName } = getDefaultSessionDirName(cwd);
if (currentDirName !== encodedDirName && currentDirName !== encodeLegacyAbsoluteSessionDirName(cwd)) {
return undefined;
}
return path.dirname(sessionDir);
}
/**
* Compute the default session directory for a cwd.
* Classifies cwd by canonical location so symlink/alias paths resolve to the
* same home-relative or temp-root directory names as their real targets.
*/
export function computeDefaultSessionDir(
cwd: string,
storage: SessionStorage,
sessionsRoot: string = getSessionsDir(),
): string {
const { encodedDirName, resolvedCwd } = getDefaultSessionDirName(cwd);
migrateHomeSessionDirs(sessionsRoot);
const sessionDir = path.join(sessionsRoot, encodedDirName);
migrateLegacyAbsoluteSessionDir(resolvedCwd, sessionDir, sessionsRoot);
storage.ensureDirSync(sessionDir);
return sessionDir;
}
// =============================================================================
// Terminal breadcrumbs: maps terminal (TTY) -> last session file for --continue
// =============================================================================
/**
* Write a breadcrumb linking the current terminal to a session file.
* The breadcrumb contains the cwd and session path so --continue can
* find "this terminal's last session" even when running concurrent instances.
*/
export function writeTerminalBreadcrumb(cwd: string, sessionFile: string): void {
const terminalId = getTerminalId();
if (!terminalId) return;
const breadcrumbDir = getTerminalSessionsDir();
const breadcrumbFile = path.join(breadcrumbDir, terminalId);
const content = `${cwd}\n${sessionFile}\n`;
// Best-effort — don't break session creation if breadcrumb fails
Bun.write(breadcrumbFile, content).catch(() => {});
}
export interface TerminalBreadcrumb {
cwd: string;
sessionFile: string;
}
/**
* Read the raw terminal breadcrumb for the current terminal.
* Returns the recorded cwd + session file (verified to exist) regardless of
* whether the recorded cwd still matches the current one. Callers decide how
* to interpret a cwd mismatch (e.g. a moved/renamed worktree).
*/
export async function readTerminalBreadcrumbEntry(): Promise<TerminalBreadcrumb | null> {
const terminalId = getTerminalId();
if (!terminalId) return null;
try {
const breadcrumbFile = path.join(getTerminalSessionsDir(), terminalId);
const content = await Bun.file(breadcrumbFile).text();
const lines = content.trim().split("\n");
if (lines.length < 2) return null;
const breadcrumbCwd = lines[0];
const sessionFile = lines[1];
// Verify the session file still exists
const stat = fs.statSync(sessionFile, { throwIfNoEntry: false });
if (stat?.isFile()) return { cwd: breadcrumbCwd, sessionFile };
} catch (err) {
if (!isEnoent(err)) logger.debug("Terminal breadcrumb read failed", { err });
// Breadcrumb doesn't exist or is corrupt — fall through
}
return null;
}
@@ -0,0 +1,131 @@
import {
type BlobStore,
externalizeImageDataSync,
externalizeImageDataUrlSync,
isBlobRef,
isImageDataUrl,
} from "./blob-store";
import type { FileEntry } from "./session-entries";
const MAX_PERSIST_CHARS = 500_000;
const TRUNCATION_NOTICE = "\n\n[Session persistence truncated large content]";
/** Minimum base64 length to externalize to blob store (skip tiny inline images) */
const BLOB_EXTERNALIZE_THRESHOLD = 1024;
const TEXT_CONTENT_KEY = "content";
function truncateString(value: string, maxLength: number): string {
if (value.length <= maxLength) return value;
let truncated = value.slice(0, maxLength);
if (truncated.length > 0) {
const last = truncated.charCodeAt(truncated.length - 1);
if (last >= 0xd800 && last <= 0xdbff) {
truncated = truncated.slice(0, -1);
}
}
return truncated;
}
export function isImageBlock(value: unknown): value is { type: "image"; data: string; mimeType?: string } {
return (
typeof value === "object" &&
value !== null &&
"type" in value &&
(value as { type?: string }).type === "image" &&
"data" in value &&
typeof (value as { data?: string }).data === "string"
);
}
/**
* Recursively truncate large strings in an object for session persistence.
* - Truncates any oversized string fields (key-agnostic)
* - Replaces oversized image blocks with text notices
* - Updates lineCount when content is truncated
* - Returns original object if no changes needed (structural sharing)
*
* Runs in one synchronous tick so an OOM/SIGKILL landing right after a persist
* call returns cannot lose the entry. Image externalization happens via the
* synchronous blob-store path (`fs.writeFileSync`), so blob bytes are in the
* kernel page cache before the JSONL line referencing them is written.
*/
function truncateForPersistence(obj: unknown, blobStore: BlobStore, key?: string): unknown {
if (obj === null || obj === undefined) return obj;
if (typeof obj === "string") {
if (key === "image_url" && isImageDataUrl(obj)) {
return externalizeImageDataUrlSync(blobStore, obj);
}
if (obj.length > MAX_PERSIST_CHARS) {
// Cryptographic signatures must be preserved exactly or cleared entirely — never truncated.
// Truncation would produce an invalid signature that the API rejects.
if (key === "thinkingSignature" || key === "thoughtSignature" || key === "textSignature") {
return "";
}
const limit = Math.max(0, MAX_PERSIST_CHARS - TRUNCATION_NOTICE.length);
return `${truncateString(obj, limit)}${TRUNCATION_NOTICE}`;
}
return obj;
}
if (Array.isArray(obj)) {
let changed = false;
const result: unknown[] = new Array(obj.length);
for (let i = 0; i < obj.length; i++) {
const item = obj[i];
if (
key === TEXT_CONTENT_KEY &&
isImageBlock(item) &&
!isBlobRef(item.data) &&
item.data.length >= BLOB_EXTERNALIZE_THRESHOLD
) {
changed = true;
result[i] = { ...item, data: externalizeImageDataSync(blobStore, item.data, item.mimeType) };
continue;
}
const newItem = truncateForPersistence(item, blobStore, key);
if (newItem !== item) changed = true;
result[i] = newItem;
}
return changed ? result : obj;
}
if (typeof obj === "object") {
let changed = false;
const entries: Array<readonly [string, unknown]> = [];
for (const [childKey, value] of Object.entries(obj)) {
// Strip transient/redundant properties that shouldn't be persisted.
// - partialJson: streaming accumulator for tool call JSON parsing
// - jsonlEvents: raw subprocess streaming events (already saved to artifact files)
if (childKey === "partialJson" || childKey === "jsonlEvents") {
changed = true;
continue;
}
const newValue = truncateForPersistence(value, blobStore, childKey);
if (newValue !== value) changed = true;
entries.push([childKey, newValue]);
}
if (!changed) return obj;
const contentEntry = entries.find(([childKey]) => childKey === "content");
const lineCountEntry = entries.find(([childKey]) => childKey === "lineCount");
if (
contentEntry &&
typeof contentEntry[1] === "string" &&
lineCountEntry &&
typeof lineCountEntry[1] === "number"
) {
const content = contentEntry[1];
const updatedEntries = entries.map(([childKey, value]) =>
childKey === "lineCount" ? ([childKey, content.split("\n").length] as const) : ([childKey, value] as const),
);
return Object.fromEntries(updatedEntries);
}
return Object.fromEntries(entries);
}
return obj;
}
export function prepareEntryForPersistence(entry: FileEntry, blobStore: BlobStore): FileEntry {
return truncateForPersistence(entry, blobStore) as FileEntry;
}
@@ -1,7 +1,7 @@
import * as fs from "node:fs";
import * as fsp from "node:fs/promises";
import * as path from "node:path";
import { isEnoent, peekFileEnds, toError } from "@oh-my-pi/pi-utils";
import { hasFsCode, isEnoent, logger, peekFileEnds, Snowflake, toError } from "@oh-my-pi/pi-utils";
const utf8Decoder = new TextDecoder("utf-8");
@@ -12,22 +12,17 @@ export interface SessionStorageStat {
}
export interface SessionStorageWriter {
writeLine(line: string): Promise<void>;
/**
* Synchronously append a single line. Returns once the bytes are handed to the kernel
* (page cache), so the data survives a non-graceful process death (OOM, SIGKILL, etc.)
* even though it has not yet been fsynced to the underlying disk.
* Append one newline-terminated line. File and memory storage perform the
* write synchronously in-body; indexed backends queue in call order.
*
* `line` MUST already include the trailing newline. Throws synchronously on I/O error.
* `line` MUST include the trailing newline.
*/
writeLineSync(line: string): void;
append(line: string): Promise<void>;
/** Resolve once all queued appends complete. No fsync. */
flush(): Promise<void>;
fsync(): Promise<void>;
/**
* Synchronously fsync the underlying file descriptor. Returns once the data
* is on the physical disk. Throws synchronously on I/O error.
*/
fsyncSync(): void;
/** False once close() has begun/finished. */
isOpen(): boolean;
close(): Promise<void>;
getError(): Error | undefined;
}
@@ -44,6 +39,7 @@ export interface SessionStorage {
/** Read the requested UTF-8 byte windows from the head and tail of the file. */
readTextSlices(path: string, prefixBytes: number, suffixBytes: number): Promise<[string, string]>;
writeText(path: string, content: string): Promise<void>;
writeTextAtomic(path: string, content: string): Promise<void>;
rename(path: string, nextPath: string): Promise<void>;
unlink(path: string): Promise<void>;
deleteSessionWithArtifacts(sessionPath: string): Promise<void>;
@@ -86,7 +82,7 @@ class FileSessionStorageWriter implements SessionStorageWriter {
return error;
}
writeLineSync(line: string): void {
async append(line: string): Promise<void> {
if (this.#closed) throw new Error("Writer closed");
if (this.#error) throw this.#error;
try {
@@ -104,33 +100,12 @@ class FileSessionStorageWriter implements SessionStorageWriter {
}
}
async writeLine(line: string): Promise<void> {
this.writeLineSync(line);
}
async flush(): Promise<void> {
if (this.#error) throw this.#error;
// OS buffers are flushed on fsync, nothing to do here
}
async fsync(): Promise<void> {
if (this.#closed) throw new Error("Writer closed");
if (this.#error) throw this.#error;
try {
fs.fsyncSync(this.#fd);
} catch (err) {
throw this.#recordError(err);
}
}
fsyncSync(): void {
if (this.#closed) throw new Error("Writer closed");
if (this.#error) throw this.#error;
try {
fs.fsyncSync(this.#fd);
} catch (err) {
throw this.#recordError(err);
}
isOpen(): boolean {
return !this.#closed;
}
async close(): Promise<void> {
@@ -204,6 +179,77 @@ export class FileSessionStorage implements SessionStorage {
await Bun.write(path, content, { createPath: true });
}
async writeTextAtomic(fpath: string, content: string): Promise<void> {
const dir = path.resolve(fpath, "..");
const tempPath = path.join(dir, `.${path.basename(fpath)}.${Snowflake.next()}.tmp`);
await fs.promises.mkdir(dir, { recursive: true });
try {
await fs.promises.writeFile(tempPath, content);
try {
await this.rename(tempPath, fpath);
return;
} catch (err) {
if (!hasFsCode(err, "EPERM")) throw toError(err);
await this.#replaceSessionFileAfterEperm(tempPath, fpath, err);
return;
}
} catch (err) {
try {
await this.unlink(tempPath);
} catch (cleanupErr) {
if (!isEnoent(cleanupErr)) {
logger.warn("Failed to remove session rewrite temp file", {
sessionFile: fpath,
tempPath,
error: toError(cleanupErr).message,
});
}
}
throw toError(err);
}
}
async #replaceSessionFileAfterEperm(tempPath: string, targetPath: string, renameError: unknown): Promise<void> {
const dir = path.resolve(targetPath, "..");
const backupPath = path.join(dir, `${path.basename(targetPath)}.${Snowflake.next()}.bak`);
try {
await this.rename(targetPath, backupPath);
} catch (moveAsideError) {
if (isEnoent(moveAsideError)) {
await this.rename(tempPath, targetPath);
return;
}
throw toError(renameError);
}
try {
await this.rename(tempPath, targetPath);
} catch (replaceError) {
try {
await this.rename(backupPath, targetPath);
} catch (rollbackErr) {
const rollbackError = toError(rollbackErr);
throw new Error(
`Failed to replace session file after EPERM (original: ${toError(renameError).message}; retry: ${
toError(replaceError).message
}; rollback: ${rollbackError.message})`,
{ cause: toError(renameError) },
);
}
throw toError(replaceError);
}
try {
await this.unlink(backupPath);
} catch (err) {
if (!isEnoent(err)) {
logger.warn("Failed to remove session rewrite backup", {
sessionFile: targetPath,
backupPath,
error: toError(err).message,
});
}
}
}
async rename(path: string, nextPath: string): Promise<void> {
try {
await fs.promises.rename(path, nextPath);
@@ -282,7 +328,7 @@ class MemorySessionStorageWriter implements SessionStorageWriter {
return error;
}
writeLineSync(line: string): void {
async append(line: string): Promise<void> {
if (this.#closed) throw new Error("Writer closed");
if (this.#error) throw this.#error;
try {
@@ -293,22 +339,12 @@ class MemorySessionStorageWriter implements SessionStorageWriter {
}
}
async writeLine(line: string): Promise<void> {
this.writeLineSync(line);
}
async flush(): Promise<void> {
if (this.#error) throw this.#error;
}
async fsync(): Promise<void> {
// No-op for in-memory storage
if (this.#error) throw this.#error;
}
fsyncSync(): void {
// No-op for in-memory storage
if (this.#error) throw this.#error;
isOpen(): boolean {
return !this.#closed;
}
async close(): Promise<void> {
@@ -527,6 +563,11 @@ export class MemorySessionStorage implements SessionStorage {
return Promise.resolve();
}
writeTextAtomic(path: string, content: string): Promise<void> {
this.writeTextSync(path, content);
return Promise.resolve();
}
rename(path: string, nextPath: string): Promise<void> {
const entry = this.#files.get(path);
if (!entry) return Promise.reject(new Error(`File not found: ${path}`));
+4 -3
View File
@@ -1385,9 +1385,10 @@ export function createShellRenderer<TArgs>(config: ShellRendererConfig<TArgs>) {
},
mergeCallAndResult: true,
inline: true,
// Pending preview caps the command to a viewport-sized tail window that
// shifts while args stream; keep it out of native scrollback mid-run.
provisionalPendingPreview: true,
// Collapsed pending preview caps the command to a viewport-sized tail
// window that shifts while args stream. Expanded output is top-anchored
// enough for the transcript to commit its settled prefix.
provisionalPendingPreview: "collapsed",
};
}
@@ -754,8 +754,9 @@ export const evalToolRenderer = {
mergeCallAndResult: true,
inline: true,
// Pending preview shows tail-window code cells; the result render
// Collapsed pending preview shows tail-window code cells; the result render
// interleaves each cell's output under its code, re-laying-out every row
// below the first cell. Keep the preview out of native scrollback mid-run.
provisionalPendingPreview: true,
// below the first cell. Expanded output is top-anchored enough for the
// transcript to commit its settled prefix.
provisionalPendingPreview: "collapsed",
};
+2 -1
View File
@@ -22,6 +22,7 @@ import type { AgentRegistry } from "../registry/agent-registry";
import type { ArtifactManager } from "../session/artifacts";
import type { ClientBridge } from "../session/client-bridge";
import type { CustomMessage } from "../session/messages";
import type { UsageStatistics } from "../session/session-entries";
import type { ToolChoiceQueue } from "../session/tool-choice-queue";
import { TaskTool } from "../task";
import type { AgentOutputManager } from "../task/output-manager";
@@ -255,7 +256,7 @@ export interface ToolSession {
/** Goal runtime for the active agent session. */
getGoalRuntime?: () => GoalRuntime | undefined;
/** Get cumulative session usage statistics (input/output tokens, cost). */
getUsageStatistics?: () => import("../session/session-manager").UsageStatistics;
getUsageStatistics?: () => UsageStatistics;
/** Current per-turn token budget {total, spent, hard} for the eval `budget` helper. */
getTurnBudget?: () => { total: number | null; spent: number; hard: boolean };
/** Record output tokens consumed by an eval-spawned subagent toward the current turn budget. */
+7 -11
View File
@@ -44,18 +44,14 @@ export type ToolRenderer = {
/** Render without background box, inline in the response flow */
inline?: boolean;
/**
* Collapsed pending preview is provisional — a tail-window or otherwise
* re-anchored view the result render replaces wholesale (an edit's
* streamed-diff tail, bash/ssh command caps, eval cells whose outputs
* interleave under each cell). Its rows must never commit to native
* scrollback mid-run; see
* `ToolExecutionComponent.isTranscriptBlockCommitStable`. Absent = the
* pending preview streams top-anchored append-shaped rows the result
* render preserves (task context/assignment, write content), which stay
* commit-eligible so a call taller than the viewport scrolls into history
* instead of reading as cut off.
* Whether pending-call rows are provisional: useful on screen while a tool is
* streaming, but not durable transcript history. `true` means every pending
* shape is provisional. `"collapsed"` means only the collapsed pending shape
* is provisional; expanded rendering is top-anchored/append-shaped enough to
* let the transcript commit its settled prefix. Absent = the pending preview
* streams rows the result render preserves.
*/
provisionalPendingPreview?: boolean;
provisionalPendingPreview?: boolean | "collapsed";
};
export const toolRenderers: Record<string, ToolRenderer> = {
+4 -3
View File
@@ -346,7 +346,8 @@ export const sshToolRenderer = {
});
},
mergeCallAndResult: true,
// Pending preview caps the command to a viewport-sized tail window that
// shifts while args stream; keep it out of native scrollback mid-run.
provisionalPendingPreview: true,
// Collapsed pending preview caps the command to a viewport-sized tail window
// that shifts while args stream. Expanded output is top-anchored enough for
// the transcript to commit its settled prefix.
provisionalPendingPreview: "collapsed",
};
+1 -1
View File
@@ -8,7 +8,7 @@ import type { RenderResultOptions } from "../extensibility/custom-tools/types";
import type { Theme } from "../modes/theme/theme";
import todoDescription from "../prompts/tools/todo.md" with { type: "text" };
import type { ToolSession } from "../sdk";
import type { SessionEntry } from "../session/session-manager";
import type { SessionEntry } from "../session/session-entries";
import { framedBlock, renderStatusLine, renderTreeList } from "../tui";
import { formatErrorDetail, PREVIEW_LIMITS } from "./render-utils";
@@ -46,6 +46,9 @@ describe("AgentSession concurrent prompt guard", () => {
const authStorages: AuthStorage[] = [];
beforeEach(() => {
// Collapse scheduler settle delays so the post-abort auto-continue and
// dispose teardown are deterministic instead of racing the wall clock.
collapseSchedulerSettleDelays();
tempDir = path.join(os.tmpdir(), `pi-concurrent-test-${Snowflake.next()}`);
fs.mkdirSync(tempDir, { recursive: true });
});
@@ -8,11 +8,9 @@ import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { type CreateAgentSessionResult, createAgentSession } from "@oh-my-pi/pi-coding-agent/sdk";
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
import {
EPHEMERAL_MODEL_CHANGE_ROLE,
getRestorableSessionModels,
SessionManager,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getRestorableSessionModels } from "@oh-my-pi/pi-coding-agent/session/session-context";
import { EPHEMERAL_MODEL_CHANGE_ROLE } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { TempDir } from "@oh-my-pi/pi-utils";
describe("AgentSession model persistence", () => {
@@ -17,11 +17,8 @@ import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { createAgentSession } from "@oh-my-pi/pi-coding-agent/sdk";
import type { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
import {
type SessionEntry,
SessionManager,
type SessionMessageEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry, SessionMessageEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { Snowflake } from "@oh-my-pi/pi-utils";
function createUsage(): Usage {
@@ -483,7 +480,7 @@ describe("AgentSession OpenAI Responses replay boundaries", () => {
sessionManager.appendCustomMessageEntry("proxy-details", "Proxy metadata", true, proxyDetails);
const snapshot = sessionManager.captureState();
const customEntry = snapshot.fileEntries.find(
const customEntry = snapshot.entries.find(
entry => entry.type === "custom_message" && entry.customType === "proxy-details",
);
if (customEntry?.type !== "custom_message") {
@@ -3,12 +3,12 @@ import * as os from "node:os";
import * as path from "node:path";
import type { AgentEvent, AgentMessage } from "@oh-my-pi/pi-agent-core";
import { RpcClient } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-client";
import {
type BranchSummaryEntry,
type CustomMessageEntry,
parseSessionEntries,
type SessionMessageEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type {
BranchSummaryEntry,
CustomMessageEntry,
SessionMessageEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { parseSessionEntries } from "@oh-my-pi/pi-coding-agent/session/session-loader";
function extractText(message: AgentMessage): string {
if (message.role !== "assistant") return "";
@@ -1,7 +1,8 @@
import { afterEach, describe, expect, it } from "bun:test";
import * as path from "node:path";
import { isBlobRef } from "@oh-my-pi/pi-coding-agent/session/blob-store";
import { type SessionEntry, SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { TempDir } from "@oh-my-pi/pi-utils";
const tempDirs: TempDir[] = [];
+10 -10
View File
@@ -15,16 +15,16 @@ import * as ai from "@oh-my-pi/pi-ai";
import { encodeTextSignatureV1 } from "@oh-my-pi/pi-ai/providers/openai-responses-shared";
import type { AssistantMessage, Model, ProviderPayload, Usage } from "@oh-my-pi/pi-ai/types";
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
import {
buildSessionContext,
type CompactionEntry,
type ModelChangeEntry,
migrateSessionEntries,
parseSessionEntries,
type SessionEntry,
type SessionMessageEntry,
type ThinkingLevelChangeEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { buildSessionContext } from "@oh-my-pi/pi-coding-agent/session/session-context";
import type {
CompactionEntry,
ModelChangeEntry,
SessionEntry,
SessionMessageEntry,
ThinkingLevelChangeEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { parseSessionEntries } from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { migrateSessionEntries } from "@oh-my-pi/pi-coding-agent/session/session-migrations";
import { mockFetch } from "./helpers/fetch-mock";
import { e2eApiKey } from "./utilities";
@@ -3,7 +3,7 @@ import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
import { UiHelpers } from "@oh-my-pi/pi-coding-agent/modes/utils/ui-helpers";
import { buildSessionContext, type SessionContext } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { buildSessionContext, type SessionContext } from "@oh-my-pi/pi-coding-agent/session/session-context";
import { type Component, Container } from "@oh-my-pi/pi-tui";
function renderLastLine(container: Container, width = 120): string {
@@ -15,7 +15,7 @@ import * as path from "node:path";
import { InternalUrlRouter } from "@oh-my-pi/pi-coding-agent/internal-urls";
import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry";
import type { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { CURRENT_SESSION_VERSION } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { CURRENT_SESSION_VERSION } from "@oh-my-pi/pi-coding-agent/session/session-entries";
async function withTempDir<T>(fn: (dir: string) => Promise<T>): Promise<T> {
const dir = await fs.mkdtemp(path.join(os.tmpdir(), "history-protocol-"));
@@ -1,11 +1,11 @@
import { describe, expect, it } from "bun:test";
import type { AssistantMessage } from "@oh-my-pi/pi-ai";
import {
buildSessionContext,
type ModelChangeEntry,
type SessionEntry,
type SessionMessageEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { buildSessionContext } from "@oh-my-pi/pi-coding-agent/session/session-context";
import type {
ModelChangeEntry,
SessionEntry,
SessionMessageEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-entries";
/**
* Issue #849: After a user explicitly switches to gpt-5.5, the session reverts
@@ -6,7 +6,7 @@ import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { ModelSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/model-selector";
import { SessionSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/session-selector";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
import { setKeybindings, type TUI } from "@oh-my-pi/pi-tui";
beforeAll(() => {
@@ -12,7 +12,8 @@ import { TreeSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/component
import { UserMessageSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/user-message-selector";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import { HistoryStorage } from "@oh-my-pi/pi-coding-agent/session/history-storage";
import type { SessionInfo, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
import { setKeybindings } from "@oh-my-pi/pi-tui";
const CTRL_N = "\x0e";
@@ -15,8 +15,11 @@ import * as path from "node:path";
import type { Args } from "@oh-my-pi/pi-coding-agent/cli/args";
import type { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { createSessionManager } from "@oh-my-pi/pi-coding-agent/main";
import type { SessionHeader, SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import * as sessionManagerModule from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionHeader } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
import * as sessionListingModule from "@oh-my-pi/pi-coding-agent/session/session-listing";
import { loadEntriesFromFile } from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
function buildArgs(resume: string, sessionDir?: string): Args {
return {
@@ -64,7 +67,7 @@ describe("createSessionManager — cross-project --resume cancellation (#1668)",
});
it("returns undefined when an interactive user declines the fork prompt instead of throwing", async () => {
vi.spyOn(sessionManagerModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(existingProject));
vi.spyOn(sessionListingModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(existingProject));
const result = await createSessionManager(
buildArgs("019e84ed"),
@@ -80,7 +83,7 @@ describe("createSessionManager — cross-project --resume cancellation (#1668)",
const originalIsTTY = process.stdin.isTTY;
Object.defineProperty(process.stdin, "isTTY", { value: false, configurable: true });
try {
vi.spyOn(sessionManagerModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(existingProject));
vi.spyOn(sessionListingModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(existingProject));
await expect(createSessionManager(buildArgs("019e84ed"), "/current/project", stubSettings)).rejects.toThrow(
`Session "019e84ed" is in another project (${existingProject}); run interactively to fork it into the current project.`,
@@ -106,7 +109,7 @@ describe("createSessionManager — cross-project --resume relocation (moved work
});
it("offers move (not fork) and returns undefined when the user declines", async () => {
vi.spyOn(sessionManagerModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(missingProject));
vi.spyOn(sessionListingModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(missingProject));
expect(fs.existsSync(missingProject)).toBe(false);
const forkPrompt = vi.fn(async () => "accepted" as const);
@@ -127,7 +130,7 @@ describe("createSessionManager — cross-project --resume relocation (moved work
const originalIsTTY = process.stdin.isTTY;
Object.defineProperty(process.stdin, "isTTY", { value: false, configurable: true });
try {
vi.spyOn(sessionManagerModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(missingProject));
vi.spyOn(sessionListingModule, "resolveResumableSession").mockResolvedValue(buildGlobalMatch(missingProject));
await expect(createSessionManager(buildArgs("019e84ed"), "/current/project", stubSettings)).rejects.toThrow(
`Session "019e84ed" belongs to a directory that no longer exists (${missingProject}); run interactively to move it into the current project.`,
@@ -142,7 +145,7 @@ describe("createSessionManager — cross-project --resume relocation (moved work
const explicitSessionDir = path.join(missingRoot, "sessions");
await fsp.mkdir(currentProject, { recursive: true });
const moved = sessionManagerModule.SessionManager.create(missingProject, explicitSessionDir);
const moved = SessionManager.create(missingProject, explicitSessionDir);
moved.appendMessage({ role: "user", content: "before local move", timestamp: 1 });
await moved.flush();
const oldFile = moved.getSessionFile();
@@ -162,7 +165,7 @@ describe("createSessionManager — cross-project --resume relocation (moved work
};
await moved.close();
expect(fs.existsSync(missingProject)).toBe(false);
vi.spyOn(sessionManagerModule, "resolveResumableSession").mockResolvedValue({
vi.spyOn(sessionListingModule, "resolveResumableSession").mockResolvedValue({
scope: "local",
session: sessionInfo,
});
@@ -181,7 +184,7 @@ describe("createSessionManager — cross-project --resume relocation (moved work
try {
expect(result.getSessionFile()).toBe(oldFile);
expect(result.getCwd()).toBe(path.resolve(currentProject));
const entries = await sessionManagerModule.loadEntriesFromFile(oldFile);
const entries = await loadEntriesFromFile(oldFile);
const header = entries.find(
(entry): entry is SessionHeader =>
typeof entry === "object" &&
@@ -9,7 +9,7 @@ import { describe, expect, it, vi } from "bun:test";
import type { Args } from "@oh-my-pi/pi-coding-agent/cli/args";
import type { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { createSessionManager, SessionResolutionError } from "@oh-my-pi/pi-coding-agent/main";
import * as sessionManagerModule from "@oh-my-pi/pi-coding-agent/session/session-manager";
import * as sessionListingModule from "@oh-my-pi/pi-coding-agent/session/session-listing";
function buildResumeArgs(resume: string): Args {
return {
@@ -36,7 +36,7 @@ const stubSettings = { get: () => undefined } as unknown as Settings;
describe("createSessionManager — missing session (#2084)", () => {
it("rejects --resume with SessionResolutionError carrying a usage hint", async () => {
vi.spyOn(sessionManagerModule, "resolveResumableSession").mockResolvedValue(undefined);
vi.spyOn(sessionListingModule, "resolveResumableSession").mockResolvedValue(undefined);
try {
await expect(
createSessionManager(
@@ -63,7 +63,7 @@ describe("createSessionManager — missing session (#2084)", () => {
});
it("rejects --fork with SessionResolutionError carrying a usage hint", async () => {
vi.spyOn(sessionManagerModule, "resolveResumableSession").mockResolvedValue(undefined);
vi.spyOn(sessionListingModule, "resolveResumableSession").mockResolvedValue(undefined);
try {
await expect(
createSessionManager(
@@ -2,14 +2,14 @@ import { describe, expect, test } from "bun:test";
import { MemorySessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
describe("MemorySessionStorage indexed mirror", () => {
test("writeLineSync builds the same content as a single writeTextSync of the join", async () => {
test("append builds the same content as a single writeTextSync of the join", async () => {
const storage = new MemorySessionStorage();
const path = "/virtual/session.jsonl";
const writer = storage.openWriter(path, { flags: "w" });
try {
const N = 1000;
for (let i = 0; i < N; i++) {
writer.writeLineSync(`{"i":${i}}\n`);
await writer.append(`{"i":${i}}\n`);
}
} finally {
await writer.close();
@@ -22,13 +22,13 @@ describe("MemorySessionStorage indexed mirror", () => {
expect(actual.length).toBe(expected.length);
});
test("statSync reports UTF-8 byte length, not character count", () => {
test("statSync reports UTF-8 byte length, not character count", async () => {
const storage = new MemorySessionStorage();
const path = "/virtual/unicode.jsonl";
const writer = storage.openWriter(path, { flags: "w" });
try {
writer.writeLineSync("héllo\n"); // é = 2 bytes in UTF-8
writer.writeLineSync("日本語\n"); // 3 chars × 3 bytes = 9
await writer.append("héllo\n"); // é = 2 bytes in UTF-8
await writer.append("日本語\n"); // 3 chars × 3 bytes = 9
} finally {
void writer.close();
}
@@ -42,9 +42,9 @@ describe("MemorySessionStorage indexed mirror", () => {
const path = "/virtual/prefix.jsonl";
const writer = storage.openWriter(path, { flags: "w" });
try {
writer.writeLineSync("alpha\n");
writer.writeLineSync("bravo\n");
writer.writeLineSync("charlie\n");
await writer.append("alpha\n");
await writer.append("bravo\n");
await writer.append("charlie\n");
} finally {
void writer.close();
}
@@ -59,9 +59,9 @@ describe("MemorySessionStorage indexed mirror", () => {
const path = "/virtual/suffix.jsonl";
const writer = storage.openWriter(path, { flags: "w" });
try {
writer.writeLineSync("alpha\n");
writer.writeLineSync("bravo\n");
writer.writeLineSync("charlie\n");
await writer.append("alpha\n");
await writer.append("bravo\n");
await writer.append("charlie\n");
} finally {
void writer.close();
}
@@ -78,9 +78,9 @@ describe("MemorySessionStorage indexed mirror", () => {
const path = "/virtual/both.jsonl";
const writer = storage.openWriter(path, { flags: "w" });
try {
writer.writeLineSync("alpha\n");
writer.writeLineSync("bravo\n");
writer.writeLineSync("charlie\n");
await writer.append("alpha\n");
await writer.append("bravo\n");
await writer.append("charlie\n");
} finally {
void writer.close();
}
@@ -93,8 +93,8 @@ describe("MemorySessionStorage indexed mirror", () => {
const path = "/virtual/unicode-slices.jsonl";
const writer = storage.openWriter(path, { flags: "w" });
try {
writer.writeLineSync("é\n");
writer.writeLineSync("日本\n");
await writer.append("é\n");
await writer.append("日本\n");
} finally {
void writer.close();
}
@@ -104,17 +104,17 @@ describe("MemorySessionStorage indexed mirror", () => {
expect(await storage.readTextSlices(path, 0, 4)).toEqual(["", "本\n"]);
});
test("subsequent writeLineSync after readText appends after materialized content", async () => {
test("subsequent append after readText appends after materialized content", async () => {
const storage = new MemorySessionStorage();
const path = "/virtual/cont.jsonl";
const writer = storage.openWriter(path, { flags: "w" });
try {
writer.writeLineSync("first\n");
writer.writeLineSync("second\n");
await writer.append("first\n");
await writer.append("second\n");
// Materialise once — implementation may collapse previous chunks into one
// string, but future appends must still retain content and byte accounting.
expect(await storage.readText(path)).toBe("first\nsecond\n");
writer.writeLineSync("third\n");
await writer.append("third\n");
expect(await storage.readText(path)).toBe("first\nsecond\nthird\n");
expect(storage.statSync(path).size).toBe(Buffer.byteLength("first\nsecond\nthird\n", "utf-8"));
} finally {
@@ -1,7 +1,7 @@
import { beforeAll, describe, expect, it } from "bun:test";
import { SessionSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/session-selector";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
beforeAll(() => {
initTheme();
@@ -1,7 +1,7 @@
import { beforeAll, describe, expect, it } from "bun:test";
import { SessionSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/session-selector";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
beforeAll(() => {
initTheme();
@@ -1,7 +1,7 @@
import { afterAll, beforeAll, describe, expect, it } from "bun:test";
import { SessionSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/session-selector";
import { initTheme, theme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionInfo, SessionStatus } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo, SessionStatus } from "@oh-my-pi/pi-coding-agent/session/session-listing";
beforeAll(async () => {
await initTheme();
@@ -1,7 +1,7 @@
import { beforeAll, describe, expect, it } from "bun:test";
import { SessionSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/session-selector";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
beforeAll(() => {
initTheme();
@@ -2,7 +2,7 @@ import { beforeAll, describe, expect, it } from "bun:test";
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { TreeSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/tree-selector";
import * as themeModule from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-entries";
let counter = 0;
function makeNode(role: "user" | "assistant", text: string, parentId: string | null = null): SessionTreeNode {
@@ -2,7 +2,7 @@ import { beforeAll, describe, expect, it } from "bun:test";
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { TreeSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/tree-selector";
import * as themeModule from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-entries";
let counter = 0;
function makeMessageNode(message: AgentMessage, parentId: string | null = null, label?: string): SessionTreeNode {
@@ -1,7 +1,7 @@
import { beforeAll, describe, expect, it } from "bun:test";
import { TreeSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/tree-selector";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-entries";
beforeAll(async () => {
await initTheme(false, undefined, undefined, "dark", "light");
@@ -2,7 +2,7 @@ import { beforeAll, describe, expect, it } from "bun:test";
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { TreeSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/tree-selector";
import * as themeModule from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-entries";
let counter = 0;
function makeNode(role: "user" | "assistant", text: string, parentId: string | null = null): SessionTreeNode {
@@ -2,7 +2,7 @@ import { beforeAll, describe, expect, it } from "bun:test";
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { TreeSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/tree-selector";
import * as themeModule from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-entries";
let counter = 0;
function makeUserNode(text: string, parentId: string | null = null): SessionTreeNode {
@@ -3,7 +3,7 @@ import { SessionSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/compon
import { SelectorController } from "@oh-my-pi/pi-coding-agent/modes/controllers/selector-controller";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { FileSessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
@@ -1,7 +1,7 @@
import { afterEach, beforeAll, describe, expect, it, vi } from "bun:test";
import { SessionSelectorComponent } from "@oh-my-pi/pi-coding-agent/modes/components/session-selector";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
beforeAll(() => {
initTheme();
@@ -16,7 +16,7 @@ import { beforeAll, describe, expect, it, type Mock, vi } from "bun:test";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
import { UiHelpers } from "@oh-my-pi/pi-coding-agent/modes/utils/ui-helpers";
import type { SessionContext } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionContext } from "@oh-my-pi/pi-coding-agent/session/session-context";
beforeAll(() => {
initTheme();
@@ -40,20 +40,14 @@ class CloseHoldingStorage implements SessionStorage {
const inner = this.#inner.openWriter(path, options);
const gates = this.#closeGates;
return {
writeLine(line) {
return inner.writeLine(line);
},
writeLineSync(line) {
inner.writeLineSync(line);
append(line) {
return inner.append(line);
},
flush() {
return inner.flush();
},
fsync() {
return inner.fsync();
},
fsyncSync() {
inner.fsyncSync();
isOpen() {
return inner.isOpen();
},
async close() {
const gate = Promise.withResolvers<void>();
@@ -106,6 +100,9 @@ class CloseHoldingStorage implements SessionStorage {
writeText(p: string, content: string): Promise<void> {
return this.#inner.writeText(p, content);
}
writeTextAtomic(p: string, content: string): Promise<void> {
return this.#inner.writeTextAtomic(p, content);
}
rename(p: string, nextPath: string): Promise<void> {
return this.#inner.rename(p, nextPath);
}
@@ -13,7 +13,8 @@
*/
import { describe, expect, it } from "bun:test";
import { type SkillPromptDetails, stripInternalDetailsFields } from "@oh-my-pi/pi-coding-agent/session/messages";
import { type CustomMessageEntry, SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { CustomMessageEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
const SKILL_TYPE = "skill-prompt";
@@ -1,13 +1,13 @@
import { describe, expect, it } from "bun:test";
import {
type BranchSummaryEntry,
buildSessionContext,
type CompactionEntry,
type ModelChangeEntry,
type SessionEntry,
type SessionMessageEntry,
type ThinkingLevelChangeEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { buildSessionContext } from "@oh-my-pi/pi-coding-agent/session/session-context";
import type {
BranchSummaryEntry,
CompactionEntry,
ModelChangeEntry,
SessionEntry,
SessionMessageEntry,
ThinkingLevelChangeEntry,
} from "@oh-my-pi/pi-coding-agent/session/session-entries";
function msg(id: string, parentId: string | null, role: "user" | "assistant", text: string): SessionMessageEntry {
const base = { type: "message" as const, id, parentId, timestamp: "2025-01-01T00:00:00Z" };
@@ -3,11 +3,9 @@ import * as fs from "node:fs";
import * as fsp from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import {
loadEntriesFromFile,
type SessionHeader,
SessionManager,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionHeader } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { loadEntriesFromFile } from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getTerminalId } from "@oh-my-pi/pi-tui";
import { getConfigRootDir, getTerminalSessionsDir, setAgentDir } from "@oh-my-pi/pi-utils";
@@ -243,17 +241,9 @@ describe("SessionManager.continueRecent relocation", () => {
});
it("re-roots past a cwd-less legacy session in a shared explicit sessionDir", async () => {
// Regression: SessionInfo.cwd is "" for sessions whose header has no cwd, and
// path.resolve("") === process.cwd(). A guard that only excluded `undefined`
// treated such a legacy session as "belongs to the current cwd" whenever
// --continue ran from process.cwd(), hijacking the moved session. Resume must
// be invoked with process.cwd() to reproduce the path.resolve("") collision.
const explicitSessionDir = path.join(testAgentDir, "shared-legacy-sessions");
const currentCwd = process.cwd();
// Older session with no recorded cwd (header cwd stripped → "" on load).
const legacy = SessionManager.create(cwdB, explicitSessionDir);
legacy.appendMessage({ role: "user", content: "legacy cwd-less", timestamp: 1 });
legacy.appendMessage({ role: "user", content: "legacy without cwd", timestamp: 1 });
legacy.appendMessage(makeAssistantMessage());
await legacy.flush();
const legacyFile = legacy.getSessionFile();
@@ -261,10 +251,10 @@ describe("SessionManager.continueRecent relocation", () => {
await legacy.close();
stripHeaderCwd(legacyFile);
// Newer moved session, recorded under the now-missing worktree cwd.
// Ensure the stale moved session is newer than the cwd-less legacy session.
await new Promise(resolve => setTimeout(resolve, 20));
const moved = SessionManager.create(cwdA, explicitSessionDir);
moved.appendMessage({ role: "user", content: "newer moved cwd", timestamp: 2 });
moved.appendMessage({ role: "user", content: "newer stale moved cwd", timestamp: 2 });
moved.appendMessage(makeAssistantMessage());
await moved.flush();
const movedFile = moved.getSessionFile();
@@ -274,12 +264,13 @@ describe("SessionManager.continueRecent relocation", () => {
writeBreadcrumb(cwdA, movedFile);
await fsp.rm(cwdA, { recursive: true, force: true });
const resumed = await SessionManager.continueRecent(currentCwd, explicitSessionDir);
const resumed = await SessionManager.continueRecent(cwdB, explicitSessionDir);
try {
// The moved session is re-rooted; the cwd-less legacy session is not hijacked.
expect(resumed.getSessionFile()).toBe(movedFile);
expect(resumed.getCwd()).toBe(path.resolve(currentCwd));
expect(resumed.getCwd()).toBe(path.resolve(cwdB));
expect(fs.existsSync(legacyFile)).toBe(true);
expect(getHeader(await loadEntriesFromFile(movedFile))?.cwd).toBe(path.resolve(cwdB));
} finally {
await resumed.close();
}
@@ -2,14 +2,10 @@ import { afterEach, beforeEach, describe, expect, it } from "bun:test";
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import {
type FileEntry,
findMostRecentSession,
loadEntriesFromFile,
resolveResumableSession,
type SessionHeader,
SessionManager,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { FileEntry, SessionHeader } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { findMostRecentSession, resolveResumableSession } from "@oh-my-pi/pi-coding-agent/session/session-listing";
import { loadEntriesFromFile } from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getConfigRootDir, getSessionsDir, Snowflake, setAgentDir } from "@oh-my-pi/pi-utils";
describe("loadEntriesFromFile", () => {
@@ -192,36 +188,6 @@ describe("SessionManager temp cwd session dirs", () => {
fs.rmSync(testAgentDir, { recursive: true, force: true });
});
it("stores symlink-equivalent home cwd sessions under home-relative directories", () => {
if (process.platform === "win32") return;
const projectsRoot = path.join(os.homedir(), "Projects");
fs.mkdirSync(projectsRoot, { recursive: true });
const realProjectDir = fs.mkdtempSync(path.join(projectsRoot, "omp-session-home-"));
const nestedDir = path.join(realProjectDir, "nested");
const aliasRoot = fs.mkdtempSync(path.join(os.tmpdir(), "omp-session-home-alias-"));
const homeAlias = path.join(aliasRoot, "home-link");
try {
fs.mkdirSync(nestedDir, { recursive: true });
fs.symlinkSync(os.homedir(), homeAlias, "dir");
const aliasedCwd = path.join(homeAlias, "Projects", path.basename(realProjectDir), "nested");
const session = SessionManager.create(aliasedCwd);
const sessionFile = session.getSessionFile();
if (!sessionFile) throw new Error("Expected session file path");
const expectedDir = path.join(
getSessionsDir(),
`-${path.relative(os.homedir(), fs.realpathSync(aliasedCwd)).replace(/[/\\:]/g, "-")}`,
);
expect(path.dirname(sessionFile)).toBe(expectedDir);
} finally {
fs.rmSync(aliasRoot, { recursive: true, force: true });
fs.rmSync(realProjectDir, { recursive: true, force: true });
}
});
it("stores temp-root cwd sessions under -tmp-prefixed directories", () => {
const tempCwd = path.join(testAgentDir, `temp-cwd-${Snowflake.next()}`);
fs.mkdirSync(tempCwd, { recursive: true });
@@ -1,5 +1,6 @@
import { describe, expect, it } from "bun:test";
import { type LabelEntry, SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { LabelEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
describe("SessionManager labels", () => {
it("sets and gets labels", () => {
@@ -1,5 +1,6 @@
import { describe, expect, it } from "bun:test";
import { type FileEntry, migrateSessionEntries } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { FileEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { migrateSessionEntries } from "@oh-my-pi/pi-coding-agent/session/session-migrations";
describe("migrateSessionEntries", () => {
it("should add id/parentId to v1 entries", () => {
@@ -3,11 +3,9 @@ import * as fs from "node:fs";
import * as fsp from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import {
loadEntriesFromFile,
type SessionHeader,
SessionManager,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionHeader } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { loadEntriesFromFile } from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { stripOuterDoubleQuotes } from "@oh-my-pi/pi-coding-agent/tools/path-utils";
import { getConfigRootDir, setAgentDir } from "@oh-my-pi/pi-utils";
@@ -1,6 +1,10 @@
import { describe, expect, it } from "bun:test";
import { recoverOrphanedBackups, SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { MemorySessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
import * as fsp from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import { recoverOrphanedBackups } from "@oh-my-pi/pi-coding-agent/session/session-listing";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { FileSessionStorage, MemorySessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
class FsCodeError extends Error {
code: string;
@@ -11,11 +15,15 @@ class FsCodeError extends Error {
}
}
class RenameEpermOnceStorage extends MemorySessionStorage {
// The atomic-write + EPERM `.bak` move-aside/rollback dance lives in
// FileSessionStorage.writeTextAtomic, so these tests must drive a real
// file-backed storage (with a temp dir) and override `rename` to simulate the
// Windows EPERM-on-replace failure.
class RenameEpermOnceStorage extends FileSessionStorage {
failNextSessionReplace = false;
backupCleanupPath: string | undefined;
rename(source: string, target: string): Promise<void> {
async rename(source: string, target: string): Promise<void> {
if (
this.failNextSessionReplace &&
source.includes(".tmp") &&
@@ -23,14 +31,12 @@ class RenameEpermOnceStorage extends MemorySessionStorage {
this.existsSync(target)
) {
this.failNextSessionReplace = false;
return Promise.reject(
new FsCodeError("EPERM", `EPERM: operation not permitted, rename '${source}' -> '${target}'`),
);
throw new FsCodeError("EPERM", `EPERM: operation not permitted, rename '${source}' -> '${target}'`);
}
return super.rename(source, target);
}
unlink(target: string): Promise<void> {
async unlink(target: string): Promise<void> {
if (target.endsWith(".bak")) {
this.backupCleanupPath = target;
}
@@ -39,9 +45,19 @@ class RenameEpermOnceStorage extends MemorySessionStorage {
}
describe("SessionManager rewrite EPERM replacement fallback", () => {
let sessionDir: string;
beforeEach(async () => {
sessionDir = await fsp.mkdtemp(path.join(os.tmpdir(), "omp-eperm-"));
});
afterEach(async () => {
await fsp.rm(sessionDir, { recursive: true, force: true });
});
it("keeps the active session healthy when replacing an existing file hits EPERM", async () => {
const storage = new RenameEpermOnceStorage();
const session = SessionManager.create("/cwd", "/sessions", storage);
const session = SessionManager.create(sessionDir, sessionDir, storage);
await session.ensureOnDisk();
const sessionFile = session.getSessionFile();
if (!sessionFile) throw new Error("Expected session file");
@@ -61,30 +77,40 @@ describe("SessionManager rewrite EPERM replacement fallback", () => {
});
describe("SessionManager rewrite EPERM rollback failure", () => {
let sessionDir: string;
beforeEach(async () => {
sessionDir = await fsp.mkdtemp(path.join(os.tmpdir(), "omp-eperm-"));
});
afterEach(async () => {
await fsp.rm(sessionDir, { recursive: true, force: true });
});
it("preserves the original EPERM as the thrown error's cause when rollback also fails", async () => {
class DoubleFailStorage extends MemorySessionStorage {
class DoubleFailStorage extends FileSessionStorage {
failureMode = false;
tempRenameAttempts = 0;
rename(source: string, target: string): Promise<void> {
async rename(source: string, target: string): Promise<void> {
if (!this.failureMode) return super.rename(source, target);
// Every temp -> target rename fails with EPERM (both the upstream attempt in
// #replaceSessionFile and the retry inside #replaceSessionFileAfterEperm).
// writeTextAtomic and the retry inside #replaceSessionFileAfterEperm).
if (source.includes(".tmp") && target.endsWith(".jsonl")) {
this.tempRenameAttempts++;
const tag = this.tempRenameAttempts === 1 ? "original" : "retry";
return Promise.reject(new FsCodeError("EPERM", `EPERM ${tag}: rename '${source}' -> '${target}'`));
throw new FsCodeError("EPERM", `EPERM ${tag}: rename '${source}' -> '${target}'`);
}
// The rollback rename (backup -> target) fails with a distinct code.
if (source.endsWith(".bak") && target.endsWith(".jsonl")) {
return Promise.reject(new FsCodeError("EIO", `EIO rollback: rename '${source}' -> '${target}'`));
throw new FsCodeError("EIO", `EIO rollback: rename '${source}' -> '${target}'`);
}
return super.rename(source, target);
}
}
const storage = new DoubleFailStorage();
const session = SessionManager.create("/cwd", "/sessions", storage);
const session = SessionManager.create(sessionDir, sessionDir, storage);
await session.ensureOnDisk();
storage.failureMode = true;
const sessionFile = session.getSessionFile();
@@ -1,5 +1,6 @@
import { describe, expect, it } from "bun:test";
import { type CustomEntry, SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { CustomEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
describe("SessionManager.saveCustomEntry", () => {
it("saves custom entries and includes them in tree traversal", () => {
@@ -2,7 +2,8 @@ import { describe, expect, it } from "bun:test";
import * as fs from "node:fs/promises";
import * as path from "node:path";
import type { AssistantMessage } from "@oh-my-pi/pi-ai";
import { SessionManager, type SessionMessageEntry } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionMessageEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getBlobsDir, TempDir } from "@oh-my-pi/pi-utils";
function isAssistantSessionEntry(entry: unknown): entry is SessionMessageEntry & { message: AssistantMessage } {
@@ -2,11 +2,9 @@ import { afterEach, beforeEach, describe, expect, it } from "bun:test";
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import {
loadEntriesFromFile,
type SessionHeader,
SessionManager,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionHeader } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { loadEntriesFromFile } from "@oh-my-pi/pi-coding-agent/session/session-loader";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getConfigRootDir, setAgentDir } from "@oh-my-pi/pi-utils";
import { makeAssistantMessage } from "./helpers";
@@ -1,5 +1,6 @@
import { describe, expect, it } from "bun:test";
import { type CustomEntry, SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { CustomEntry } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { assistantMsg, userMsg } from "../utilities";
describe("SessionManager append and tree traversal", () => {
@@ -3,7 +3,7 @@ import {
mergeSessionRanking,
rankSessionSearchMatches,
} from "@oh-my-pi/pi-coding-agent/modes/components/session-selector";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionInfo } from "@oh-my-pi/pi-coding-agent/session/session-listing";
function makeSession(id: string, overrides: Partial<SessionInfo> = {}): SessionInfo {
return {
@@ -6,7 +6,7 @@
* a general-purpose mock. Each test exercises one contract:
*
* - the metadata index keeps `existsSync`/`statSync`/`listFilesSync`
* coherent with `writeText`/`writer.writeLineSync`;
* coherent with `writeText`/`writer.append`;
* - `drain()` waits for fire-and-forget background writes;
* - `deleteSessionWithArtifacts` removes both the JSONL key and any sidecar
* keys under the artifacts prefix;
@@ -221,11 +221,11 @@ describe("RedisSessionStorage", () => {
expect(c).toBeGreaterThan(b);
});
it("writer.writeLineSync appends to Redis after drain", async () => {
it("writer.append appends to Redis after drain", async () => {
const storage = await RedisSessionStorage.create({ client: redis });
const writer = storage.openWriter("/sessions/p/session.jsonl");
writer.writeLineSync('{"type":"session"}\n');
writer.writeLineSync('{"type":"message"}\n');
await writer.append('{"type":"session"}\n');
await writer.append('{"type":"message"}\n');
// Reads await queued appends and fetch content from Redis.
expect(await storage.readText("/sessions/p/session.jsonl")).toBe('{"type":"session"}\n{"type":"message"}\n');
@@ -244,7 +244,7 @@ describe("RedisSessionStorage", () => {
await storage.writeText("/sessions/p/keep.jsonl", "old content\n");
const writer = storage.openWriter("/sessions/p/keep.jsonl", { flags: "w" });
writer.writeLineSync("fresh\n");
await writer.append("fresh\n");
await writer.close();
expect(await storage.readText("/sessions/p/keep.jsonl")).toBe("fresh\n");
@@ -255,7 +255,7 @@ describe("RedisSessionStorage", () => {
const storage = await RedisSessionStorage.create({ client: redis });
const writer = storage.openWriter("/sessions/p/fail.jsonl");
redis.failNext("append", new Error("redis exploded"));
writer.writeLineSync("doomed\n");
void writer.append("doomed\n").catch(() => {});
await expect(storage.drain()).rejects.toThrow("redis exploded");
expect(writer.getError()?.message).toBe("redis exploded");
@@ -1,11 +1,8 @@
import { describe, expect, it } from "bun:test";
import * as fs from "node:fs/promises";
import * as path from "node:path";
import {
CURRENT_SESSION_VERSION,
type SessionHeader,
SessionManager,
} from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { CURRENT_SESSION_VERSION, type SessionHeader } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getTerminalId } from "@oh-my-pi/pi-tui";
import { getAgentDir, getTerminalSessionsDir, setAgentDir, TempDir } from "@oh-my-pi/pi-utils";
@@ -1,5 +1,6 @@
import { describe, expect, it } from "bun:test";
import { SessionManager, type SessionStatus } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionStatus } from "@oh-my-pi/pi-coding-agent/session/session-listing";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { MemorySessionStorage } from "@oh-my-pi/pi-coding-agent/session/session-storage";
const SESSION_DIR = "/sessions/status-proj";
@@ -77,11 +77,11 @@ describe("SqlSessionStorage (SQLite backend)", () => {
await client.end();
});
it("writer.writeLineSync appends to SQL after drain", async () => {
it("writer.append appends to SQL after drain", async () => {
const { client, storage } = await createSqlite();
const writer = storage.openWriter("/sessions/p/session.jsonl");
writer.writeLineSync('{"type":"session"}\n');
writer.writeLineSync('{"type":"message"}\n');
await writer.append('{"type":"session"}\n');
await writer.append('{"type":"message"}\n');
// Reads await queued appends and fetch content from SQL.
expect(await storage.readText("/sessions/p/session.jsonl")).toBe('{"type":"session"}\n{"type":"message"}\n');
@@ -101,7 +101,7 @@ describe("SqlSessionStorage (SQLite backend)", () => {
await storage.writeText("/sessions/p/keep.jsonl", "old content\n");
const writer = storage.openWriter("/sessions/p/keep.jsonl", { flags: "w" });
writer.writeLineSync("fresh\n");
await writer.append("fresh\n");
await writer.close();
expect(await storage.readText("/sessions/p/keep.jsonl")).toBe("fresh\n");
@@ -132,7 +132,7 @@ describe("SqlSessionStorage (SQLite backend)", () => {
// Force a SQL error: drop the table so the next append throws.
await client.unsafe("DROP TABLE omp_session_files");
writer.writeLineSync("doomed\n");
void writer.append("doomed\n").catch(() => {});
await expect(storage.drain()).rejects.toThrow();
expect(writer.getError()).toBeDefined();
@@ -328,7 +328,7 @@ describe("SqlSessionStorage (dialect-specific SQL)", () => {
const { client, queries } = capturingClient("postgres");
const storage = await SqlSessionStorage.create({ client });
const writer = storage.openWriter("/s/p.jsonl");
writer.writeLineSync("chunk\n");
await writer.append("chunk\n");
await writer.close();
const ddl = queries.find(q => q.sql.startsWith("CREATE TABLE"));
@@ -352,7 +352,7 @@ describe("SqlSessionStorage (dialect-specific SQL)", () => {
const { client, queries } = capturingClient("mysql");
const storage = await SqlSessionStorage.create({ client });
const writer = storage.openWriter("/s/m.jsonl");
writer.writeLineSync("chunk\n");
await writer.append("chunk\n");
await writer.close();
const ddl = queries.find(q => q.sql.startsWith("CREATE TABLE"));
+2 -1
View File
@@ -2,7 +2,8 @@ import { describe, expect, test } from "bun:test";
import type { SessionData } from "../src/export/html";
import { buildShareSnapshot, normalizeShareServerUrl, SERVER_MAX_SEALED_BYTES, sealToFit } from "../src/export/share";
import { SecretObfuscator } from "../src/secrets/obfuscator";
import type { SessionEntry, SessionManager } from "../src/session/session-manager";
import type { SessionEntry } from "../src/session/session-entries";
import type { SessionManager } from "../src/session/session-manager";
const IV_LENGTH = 12;
@@ -16,7 +16,7 @@ import type {
import type { Theme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import type { SessionEntry, SessionTreeNode } from "@oh-my-pi/pi-coding-agent/session/session-entries";
import { ToolChoiceQueue } from "@oh-my-pi/pi-coding-agent/session/tool-choice-queue";
import { createTools, type ToolSession } from "@oh-my-pi/pi-coding-agent/tools";
import { searchToolRenderer } from "@oh-my-pi/pi-coding-agent/tools/search";
@@ -57,9 +57,9 @@ describe("write tool shebang chmod", () => {
const stat = await fs.stat(filePath);
// All three execute bits flipped on (chmod a+x semantics).
expect(stat.mode & 0o111).toBe(0o111);
// Notice surfaces on details, not in the model-facing text.
// Notice remains model-facing so callers see that chmod changed the file mode.
expect(details(result).madeExecutable).toBe(true);
expect(resultText(result)).not.toContain("executable");
expect(resultText(result)).toContain("[Notice: Made executable via chmod +x]");
});
it("does not chmod files without a shebang", async () => {
+28 -2
View File
@@ -762,6 +762,11 @@ export class TUI extends Container {
// fast path (`#renderResizeViewport`) instead of an authoritative full
// paint, and no commit/window/diff state is advanced.
#resizeViewportActive = false;
// Set only by the resize callback's cheap-paint request. A concurrent
// caller-forced render (tool finalization, reset, image reconciliation) must
// not be downgraded to the throwaway viewport path just because a resize
// settle window is active.
#resizeViewportPaintPending = false;
// Quiet-window timer that ends the drag: its callback clears the flag and
// drives the one authoritative full paint. Reset on every resize event so it
// only fires once the drag stops. Cancelled on stop().
@@ -1219,7 +1224,7 @@ export class TUI extends Container {
// request the cheap viewport-only paint. The authoritative full
// replay fires from the settle timer once the drag goes quiet.
this.#beginResizeViewport();
this.requestRender(true);
this.#requestResizeViewportPaint();
return;
}
this.#armMultiplexerResizeTimer(false);
@@ -1503,6 +1508,7 @@ export class TUI extends Container {
// Any non-component-scoped request makes the pending frame a full one.
this.#pendingRenderComponentsOnly = false;
if (force) {
this.#resizeViewportPaintPending = false;
// Forced repaints landing inside the multiplexer resize debounce
// (e.g. `#finishSixelProbe`, image-budget eviction, a programmatic
// `requestRender(true)`) would paint into a still-reflowing pane
@@ -2181,7 +2187,13 @@ export class TUI extends Container {
// ran. A visible overlay composites over the transcript and needs the
// whole window, so fall through to the normal forced paint when one is up
// (overlay resizes are not on the drag-cost hot path).
if (this.#resizeViewportActive && this.#hasEverRendered && this.#getTopmostVisibleOverlay() === undefined) {
if (
this.#resizeViewportPaintPending &&
this.#resizeViewportActive &&
this.#hasEverRendered &&
this.#getTopmostVisibleOverlay() === undefined
) {
this.#resizeViewportPaintPending = false;
this.#componentRenderTargets.clear();
this.#renderResizeViewport(width, height);
return;
@@ -2797,6 +2809,20 @@ export class TUI extends Container {
}, TUI.#RESIZE_VIEWPORT_SETTLE_MS);
}
#requestResizeViewportPaint(): void {
if (this.#stopped) return;
this.#resizeViewportPaintPending = true;
if (this.#renderRequested) return;
this.#renderRequested = true;
this.#renderScheduler.scheduleImmediate(() => {
if (this.#stopped || !this.#renderRequested) return;
this.#renderRequested = false;
this.#lastRenderAt = this.#renderScheduler.now();
this.#doRender();
if (this.#renderRequested) this.#scheduleRender();
});
}
/**
* Compose and paint only the viewport for one resize fast-path frame.
* State-isolated: advances no commit/window/diff field and calls neither