diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index a8912a45a..7265aaa6c 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -2,6 +2,11 @@ ## [Unreleased] +### Added + +- Added `RedisSessionStorage`, a `bun:redis`-backed implementation of the `SessionStorage` interface that lets API consumers route session JSONL through Redis instead of local disk. Pass a connected `Bun.RedisClient` (or any compatible adapter) to `RedisSessionStorage.create({ client, prefix? })` and hand the returned storage to `SessionManager.create(cwd, sessionDir, storage)` (or any other static factory that accepts a storage argument). An in-memory mirror is loaded on creation so the interface's synchronous methods (`existsSync`, `statSync`, `listFilesSync`, …) keep their contracts; `drain()` waits for queued background writes. Tool artifacts and image blobs still live on disk via `ArtifactManager`/`BlobStore` — Redis only owns the session JSONL keyspace under the configured prefix. +- Exported the `SessionStorage` / `SessionStorageWriter` / `FileSessionStorage` / `MemorySessionStorage` symbols (already reachable via the `./session/session-storage` subpath) from the package root so SDK consumers can construct alternative storage backends without deep-importing. + ## [15.5.10] - 2026-05-28 ### Added diff --git a/packages/coding-agent/examples/sdk/12-redis-sessions.ts b/packages/coding-agent/examples/sdk/12-redis-sessions.ts new file mode 100644 index 000000000..3db40045d --- /dev/null +++ b/packages/coding-agent/examples/sdk/12-redis-sessions.ts @@ -0,0 +1,54 @@ +/** + * Redis-Backed Sessions + * + * Store session JSONL in Redis (or Valkey) instead of the local filesystem. + * Useful when the agent runs in an ephemeral container, behind a load + * balancer, or anywhere a shared session store beats per-host disk state. + * + * The storage substrate is the only thing that changes — every other SDK + * surface (extensions, hooks, custom tools, slash commands, branching, + * `SessionManager.list`, …) continues to work unmodified. + * + * Tool artifacts and image blobs are out of scope: `ArtifactManager` / + * `BlobStore` keep writing to `~/.omp/agent/...`. Reach for an object store + * (S3, R2, GCS) if you need those off-host too. + */ + +import { createAgentSession, RedisSessionStorage, SessionManager } from "@oh-my-pi/pi-coding-agent"; +import { RedisClient } from "bun"; + +// `bun:redis` picks up `REDIS_URL` / `VALKEY_URL` from the environment, or +// you can pass an explicit `redis://`/`rediss://` URL. +const redis = new RedisClient(); +await redis.ping(); + +// `create()` warms an in-memory mirror with every existing key under the +// prefix so SessionManager's synchronous lookups (resume, recent sessions, +// list) work without per-call network round-trips. +const storage = await RedisSessionStorage.create({ + client: redis, + prefix: "omp:sessions:", // optional, this is the default +}); + +const sessionDir = "/sessions/my-project"; + +// 1) Fresh persistent session, JSONL backed by Redis. +const { session } = await createAgentSession({ + sessionManager: SessionManager.create(process.cwd(), sessionDir, storage), +}); +console.log("New Redis session:", session.sessionFile); + +// 2) Continue the most recent session for this `sessionDir`. +const { session: continued } = await createAgentSession({ + sessionManager: await SessionManager.continueRecent(process.cwd(), sessionDir, storage), +}); +console.log("Resumed:", continued.sessionFile); + +// 3) List every Redis-backed session under this directory key prefix. +const sessions = await SessionManager.list(process.cwd(), sessionDir, storage); +console.log(`Found ${sessions.length} sessions under ${sessionDir}`); + +// On graceful shutdown, drain any background writes the writer queued and +// close the Redis connection so containerized hosts can exit cleanly. +await storage.drain(); +redis.close(); diff --git a/packages/coding-agent/src/index.ts b/packages/coding-agent/src/index.ts index b17ec967c..7d13a5422 100644 --- a/packages/coding-agent/src/index.ts +++ b/packages/coding-agent/src/index.ts @@ -40,8 +40,10 @@ export * from "./session/agent-session"; // Auth and model registry export * from "./session/auth-storage"; export * from "./session/messages"; +export * from "./session/redis-session-storage"; export * from "./session/session-dump-format"; export * from "./session/session-manager"; +export * from "./session/session-storage"; export * from "./task/executor"; export type * from "./task/types"; // Tools (detail types and utilities) diff --git a/packages/coding-agent/src/session/redis-session-storage.ts b/packages/coding-agent/src/session/redis-session-storage.ts new file mode 100644 index 000000000..ef5762861 --- /dev/null +++ b/packages/coding-agent/src/session/redis-session-storage.ts @@ -0,0 +1,481 @@ +import { logger, toError } from "@oh-my-pi/pi-utils"; +import type { SessionStorage, SessionStorageStat, SessionStorageWriter } from "./session-storage"; + +/** + * Minimal subset of the `bun:redis` `RedisClient` surface used by + * {@link RedisSessionStorage}. Keeping the contract narrow (and accepting any + * client that conforms) lets callers swap in test doubles or shared clients + * without dragging the entire Bun typings into this module. + */ +export interface RedisSessionStorageClient { + get(key: string): Promise; + set(key: string, value: string): Promise; + append(key: string, value: string): Promise; + del(...keys: string[]): Promise; + rename(src: string, dst: string): Promise; + scan(cursor: string, ...args: string[]): Promise<[string, string[]]>; + hset(key: string, field: string, value: string): Promise; + hgetall(key: string): Promise>; + hdel(key: string, ...fields: string[]): Promise; +} + +export interface RedisSessionStorageOptions { + /** A connected `bun:redis` RedisClient (or any compatible adapter). */ + client: RedisSessionStorageClient; + /** + * Key prefix applied to every Redis key this storage owns. Default `omp:sessions:`. + * Trailing colon is preserved verbatim — set to a project-scoped prefix to share + * one Redis instance between multiple agents. + */ + prefix?: string; + /** + * Maximum number of keys returned per SCAN batch when warming the mirror. + * Default 500. + */ + scanCount?: number; +} + +interface MirrorEntry { + content: string; + mtimeMs: number; +} + +const DEFAULT_PREFIX = "omp:sessions:"; +const DEFAULT_SCAN_COUNT = 500; + +function enoent(p: string): NodeJS.ErrnoException { + const err = new Error(`ENOENT: no such file, '${p}'`) as NodeJS.ErrnoException; + err.code = "ENOENT"; + err.errno = -2; + err.path = p; + err.syscall = "open"; + return err; +} + +function matchesGlob(name: string, pattern: string): boolean { + if (pattern === "*") return true; + if (pattern.startsWith("*.")) return name.endsWith(pattern.slice(1)); + return name === pattern; +} + +/** + * Redis-backed implementation of {@link SessionStorage}. Each session JSONL + * file maps to a Redis STRING key, with per-key metadata (mtime) tracked in a + * single sibling HASH. An in-memory mirror is loaded on construction so the + * interface's synchronous methods (`existsSync`, `statSync`, `listFilesSync`, + * `readTextSync`, `writeTextSync`) keep their contracts — Bun's Redis client + * is async only, and the persist hot path (`writer.writeLineSync`) cannot + * wait on a network round-trip. + * + * Trade-offs vs `FileSessionStorage`: + * - Mirror state is process-local. Two processes writing the same session key + * will diverge until one of them reloads via {@link refresh}. This matches + * `FileSessionStorage`'s existing single-writer assumption. + * - `writeLineSync` updates the mirror synchronously and queues an async + * `APPEND`. The promise is awaited by `flush()` / `close()` / {@link drain}. + * A SIGKILL landing between the sync mirror update and the network round + * trip loses the last line; the file-backed implementation survives that + * window because bytes are handed to the kernel page cache before + * returning. + * - Blobs (image data) and tool artifact files still live on disk via + * `BlobStore` / `ArtifactManager`. Those are out of scope for this storage. + */ +export class RedisSessionStorage implements SessionStorage { + readonly #client: RedisSessionStorageClient; + readonly #prefix: string; + readonly #scanCount: number; + readonly #mirror = new Map(); + readonly #writers = new Set(); + #nextMtimeMs = 0; + #pendingTail: Promise = Promise.resolve(); + + private constructor(options: RedisSessionStorageOptions) { + this.#client = options.client; + this.#prefix = options.prefix ?? DEFAULT_PREFIX; + this.#scanCount = options.scanCount ?? DEFAULT_SCAN_COUNT; + } + + /** + * Warm the in-memory mirror with every existing session key under the + * configured prefix and return the ready-to-use storage. Must be awaited + * before passing the storage into `SessionManager.create()` so synchronous + * lookups (session resume, recent sessions, EPERM-backup recovery) see + * the existing keyspace. + */ + static async create(options: RedisSessionStorageOptions): Promise { + const storage = new RedisSessionStorage(options); + await storage.refresh(); + return storage; + } + + /** + * Re-scan Redis and replace the mirror's contents. Call this from a + * different process that took over a session keyspace, or after an + * out-of-band write made by another agent. + */ + async refresh(): Promise { + this.#mirror.clear(); + const filePrefix = this.#fileKey(""); + const metaRaw = await this.#client.hgetall(this.#metaKey()); + const meta: Record = metaRaw ?? {}; + + const seen = new Set(); + let cursor = "0"; + do { + const [next, batch] = await this.#client.scan( + cursor, + "MATCH", + `${filePrefix}*`, + "COUNT", + String(this.#scanCount), + ); + cursor = next; + for (const key of batch) seen.add(key); + } while (cursor !== "0"); + + await Promise.all( + Array.from(seen, async key => { + const path = key.slice(filePrefix.length); + const content = await this.#client.get(key); + if (content === null) return; + const mtimeRaw = meta[path]; + const mtimeMs = mtimeRaw ? Number(mtimeRaw) : Date.now(); + this.#mirror.set(path, { content, mtimeMs }); + if (mtimeMs > this.#nextMtimeMs) this.#nextMtimeMs = mtimeMs; + }), + ); + } + + /** + * Resolve once every pending background write (issued via `writeTextSync` + * or `writer.writeLineSync`) has been acknowledged by Redis. Throws if any + * background write failed since the last drain. + * + * Call this on graceful shutdown to avoid losing the last unflushed line. + * The session-manager's own `flush()` / `close()` already drain through + * the writer chain — this method exists for callers (test harnesses, + * subprocess-style consumers) that bypass the writer. + */ + async drain(): Promise { + // Take ownership of the current tail, then reset so subsequent + // operations start from a clean (resolved) chain. Without the reset, + // any failure observed here would also be re-thrown by every later + // write that piggybacks on the tail via `#trackPending`. + const tail = this.#pendingTail; + this.#pendingTail = Promise.resolve(); + await tail; + } + + #fileKey(path: string): string { + return `${this.#prefix}file:${path}`; + } + + #metaKey(): string { + return `${this.#prefix}meta`; + } + + /** + * Allocate a strictly monotonic mtime. Multiple writes within the same + * millisecond would otherwise yield identical `mtimeMs` values and break + * `getSortedSessions`' newest-first ordering. + */ + #allocMtimeMs(): number { + const now = Date.now(); + const next = now > this.#nextMtimeMs ? now : this.#nextMtimeMs + 1; + this.#nextMtimeMs = next; + return next; + } + + #trackPending(promise: Promise): void { + // `Promise.all` rejects if either input rejects, which is exactly + // what we want for `drain()`. The follow-up `.catch(() => {})` is + // attached only to silence the unhandled-rejection signal on the + // shared tail — `drain()` keeps its own handler chain and still + // observes the original error, because rejection delivery is + // per-handler-chain, not per-promise. + this.#pendingTail = Promise.all([this.#pendingTail, promise]).then(() => {}); + this.#pendingTail.catch(() => {}); + } + + // --- sync surface --------------------------------------------------------- + + ensureDirSync(_dir: string): void { + // Redis is flat: directories are derived from key prefixes. + } + + existsSync(path: string): boolean { + return this.#mirror.has(path); + } + + writeTextSync(path: string, content: string): void { + const mtimeMs = this.#allocMtimeMs(); + this.#mirror.set(path, { content, mtimeMs }); + this.#trackPending(this.#writeRemote(path, content, mtimeMs)); + } + + readTextSync(path: string): string { + const entry = this.#mirror.get(path); + if (!entry) throw enoent(path); + return entry.content; + } + + statSync(path: string): SessionStorageStat { + const entry = this.#mirror.get(path); + if (!entry) throw enoent(path); + return { + size: Buffer.byteLength(entry.content, "utf-8"), + mtimeMs: entry.mtimeMs, + mtime: new Date(entry.mtimeMs), + }; + } + + listFilesSync(dir: string, pattern: string): string[] { + const prefix = dir.endsWith("/") ? dir : `${dir}/`; + const out: string[] = []; + for (const path of this.#mirror.keys()) { + if (!path.startsWith(prefix)) continue; + const name = path.slice(prefix.length); + if (name.includes("/")) continue; + if (!matchesGlob(name, pattern)) continue; + out.push(path); + } + return out; + } + + // --- async surface -------------------------------------------------------- + + async exists(path: string): Promise { + // Mirror is the source of truth; checking Redis would only diverge + // when a peer process mutated the key, which is outside the + // storage's contract (see class JSDoc). + return this.#mirror.has(path); + } + + async readText(path: string): Promise { + const entry = this.#mirror.get(path); + if (!entry) throw enoent(path); + return entry.content; + } + + async readTextPrefix(path: string, maxBytes: number): Promise { + const entry = this.#mirror.get(path); + if (!entry) throw enoent(path); + if (maxBytes <= 0) return ""; + // `entry.content` is a JS string (UTF-16 code units), but the prefix + // contract is byte-oriented. Encode to UTF-8, slice, then decode — + // matching `peekFile`'s behaviour for the file-backed storage. + const bytes = Buffer.from(entry.content, "utf-8"); + const slice = bytes.subarray(0, Math.min(maxBytes, bytes.byteLength)); + return slice.toString("utf-8"); + } + + async writeText(path: string, content: string): Promise { + const mtimeMs = this.#allocMtimeMs(); + this.#mirror.set(path, { content, mtimeMs }); + await this.#writeRemote(path, content, mtimeMs); + } + + async rename(src: string, dst: string): Promise { + const entry = this.#mirror.get(src); + if (!entry) throw enoent(src); + // Update the mirror first so a synchronous existsSync() right after + // the await resolves consistently. If RENAME fails the mirror is + // rolled back below. + this.#mirror.delete(src); + this.#mirror.set(dst, entry); + + try { + await this.#client.rename(this.#fileKey(src), this.#fileKey(dst)); + } catch (err) { + this.#mirror.delete(dst); + this.#mirror.set(src, entry); + throw toError(err); + } + + // Move the mtime hash entry too. Failures here cause meta drift but + // the mirror cache keeps statSync accurate, so log and continue. + try { + await this.#client.hdel(this.#metaKey(), src); + await this.#client.hset(this.#metaKey(), dst, String(entry.mtimeMs)); + } catch (err) { + logger.warn("Redis session storage meta rename failed", { + src, + dst, + error: toError(err).message, + }); + } + } + + async unlink(path: string): Promise { + const existed = this.#mirror.delete(path); + await this.#client.del(this.#fileKey(path)); + await this.#client.hdel(this.#metaKey(), path); + if (!existed) { + throw enoent(path); + } + } + + async deleteSessionWithArtifacts(sessionPath: string): Promise { + await this.unlink(sessionPath); + + // Mirror artifacts live under `/...`. The + // Redis storage doesn't actually persist tool artifact bytes — those + // stay on disk via `ArtifactManager` — but a draft sidecar may have + // been written through `writeText`. Sweep any keys under that prefix. + const artifactsDir = sessionPath.slice(0, -6); + const prefix = artifactsDir.endsWith("/") ? artifactsDir : `${artifactsDir}/`; + const victims: string[] = []; + for (const key of this.#mirror.keys()) { + if (key.startsWith(prefix)) victims.push(key); + } + if (victims.length === 0) return; + + for (const key of victims) this.#mirror.delete(key); + await this.#client.del(...victims.map(v => this.#fileKey(v))); + await this.#client.hdel(this.#metaKey(), ...victims); + } + + openWriter(path: string, options?: { flags?: "a" | "w"; onError?: (err: Error) => void }): SessionStorageWriter { + const writer = new RedisSessionStorageWriter(this, path, options); + this.#writers.add(writer); + return writer; + } + + // --- writer support ------------------------------------------------------- + + _writerClosed(writer: RedisSessionStorageWriter): void { + this.#writers.delete(writer); + } + + /** Mirror-only mutation, no Redis call. Used by writers to update local state synchronously. */ + _mirrorAppend(path: string, line: string): void { + const existing = this.#mirror.get(path); + const content = existing ? existing.content + line : line; + this.#mirror.set(path, { content, mtimeMs: this.#allocMtimeMs() }); + } + + /** Mirror-only mutation, no Redis call. Used by writers opened with `flags: "w"` to truncate. */ + _mirrorTruncate(path: string): void { + this.#mirror.set(path, { content: "", mtimeMs: this.#allocMtimeMs() }); + } + + async _remoteTruncate(path: string): Promise { + const entry = this.#mirror.get(path); + const mtimeMs = entry?.mtimeMs ?? Date.now(); + await this.#client.set(this.#fileKey(path), ""); + await this.#client.hset(this.#metaKey(), path, String(mtimeMs)); + } + + async _remoteAppend(path: string, line: string): Promise { + await this.#client.append(this.#fileKey(path), line); + const entry = this.#mirror.get(path); + if (entry) { + await this.#client.hset(this.#metaKey(), path, String(entry.mtimeMs)); + } + } + + /** Record a writer's pending promise on the storage-level tail so `drain()` waits for it. */ + _attachPending(promise: Promise): void { + this.#trackPending(promise); + } + + async #writeRemote(path: string, content: string, mtimeMs: number): Promise { + await this.#client.set(this.#fileKey(path), content); + await this.#client.hset(this.#metaKey(), path, String(mtimeMs)); + } +} + +class RedisSessionStorageWriter implements SessionStorageWriter { + #storage: RedisSessionStorage; + #path: string; + #closed = false; + #error: Error | undefined; + #onError: ((err: Error) => void) | undefined; + #pendingChain: Promise = Promise.resolve(); + + constructor( + storage: RedisSessionStorage, + path: string, + options?: { flags?: "a" | "w"; onError?: (err: Error) => void }, + ) { + this.#storage = storage; + this.#path = path; + this.#onError = options?.onError; + const flags = options?.flags ?? "a"; + if (flags === "w") { + // "w" mirrors FileSessionStorageWriter passing `"w"` to + // `fs.openSync`: start from empty content. Materialize the + // truncate in the mirror synchronously so an immediate reader + // can't observe stale content, then queue the remote SET. + storage._mirrorTruncate(path); + this.#enqueueRaw(() => storage._remoteTruncate(path)); + } + } + + #recordError(err: unknown): Error { + const error = toError(err); + if (!this.#error) this.#error = error; + this.#onError?.(error); + return error; + } + + #enqueueRaw(task: () => Promise): Promise { + const next = this.#pendingChain.then(async () => { + if (this.#error) throw this.#error; + try { + await task(); + } catch (err) { + throw this.#recordError(err); + } + }); + this.#pendingChain = next.catch(() => { + // Errors are recorded on `this.#error`; subsequent enqueues + // throw from inside the wrapper above. The outer chain swallows + // to avoid surfacing as an unhandled promise rejection. + }); + // Storage-level drain() waits for every writer's pending work too. + this.#storage._attachPending(next); + return next; + } + + writeLineSync(line: string): void { + if (this.#closed) throw new Error("Writer closed"); + if (this.#error) throw this.#error; + this.#storage._mirrorAppend(this.#path, line); + this.#enqueueRaw(() => this.#storage._remoteAppend(this.#path, line)); + } + + async writeLine(line: string): Promise { + if (this.#closed) throw new Error("Writer closed"); + if (this.#error) throw this.#error; + this.#storage._mirrorAppend(this.#path, line); + await this.#enqueueRaw(() => this.#storage._remoteAppend(this.#path, line)); + } + + async flush(): Promise { + if (this.#error) throw this.#error; + await this.#enqueueRaw(async () => {}); + if (this.#error) throw this.#error; + } + + async fsync(): Promise { + // Bun's `RedisClient` has no fsync equivalent; APPEND/SET return only + // after the server has acknowledged the write. `flush()` already + // awaits that ack, so this collapses into a drain. + await this.flush(); + } + + async close(): Promise { + if (this.#closed) return; + this.#closed = true; + try { + await this.flush(); + } finally { + this.#storage._writerClosed(this); + } + } + + getError(): Error | undefined { + return this.#error; + } +} diff --git a/packages/coding-agent/test/session/redis-session-storage-manager.test.ts b/packages/coding-agent/test/session/redis-session-storage-manager.test.ts new file mode 100644 index 000000000..94dba35ce --- /dev/null +++ b/packages/coding-agent/test/session/redis-session-storage-manager.test.ts @@ -0,0 +1,206 @@ +/** + * Integration: `SessionManager` driven by `RedisSessionStorage` instead of a + * file-backed store. Verifies that the storage substrate is genuinely + * pluggable — message append, persistence, reload via `open()`, and + * `SessionManager.list()` all behave the same against Redis-backed keys. + * + * Driven by the same hand-rolled in-memory Redis double used in + * `redis-session-storage.test.ts`; we don't require a live server. + */ + +import { describe, expect, it } from "bun:test"; +import type { Usage } from "@oh-my-pi/pi-ai"; +import { + RedisSessionStorage, + type RedisSessionStorageClient, +} from "@oh-my-pi/pi-coding-agent/session/redis-session-storage"; +import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager"; + +interface FakeRedis extends RedisSessionStorageClient { + strings: Map; + hashes: Map>; +} + +function createFakeRedis(): FakeRedis { + const strings = new Map(); + const hashes = new Map>(); + + const getHash = (key: string): Map => { + let h = hashes.get(key); + if (!h) { + h = new Map(); + hashes.set(key, h); + } + return h; + }; + + return { + strings, + hashes, + async get(key) { + return strings.has(key) ? (strings.get(key) as string) : null; + }, + async set(key, value) { + strings.set(key, value); + return "OK"; + }, + async append(key, value) { + const current = strings.get(key) ?? ""; + const next = current + value; + strings.set(key, next); + return next.length; + }, + async del(...keys) { + let n = 0; + for (const k of keys) { + if (strings.delete(k)) n += 1; + } + return n; + }, + async rename(src, dst) { + if (!strings.has(src)) throw new Error("ERR no such key"); + strings.set(dst, strings.get(src) as string); + strings.delete(src); + return "OK"; + }, + async scan(_cursor, ...rest) { + let pattern = "*"; + for (let i = 0; i < rest.length; i++) { + if (String(rest[i]).toUpperCase() === "MATCH") { + pattern = String(rest[i + 1] ?? "*"); + } + } + const regex = new RegExp(`^${pattern.replace(/[.+?^${}()|[\]\\]/g, "\\$&").replace(/\*/g, ".*")}$`); + const matches = Array.from(strings.keys()).filter(k => regex.test(k)); + return ["0", matches]; + }, + async hset(key, field, value) { + getHash(key).set(field, value); + return 1; + }, + async hgetall(key) { + const h = hashes.get(key); + if (!h) return {}; + const out: Record = {}; + for (const [k, v] of h) out[k] = v; + return out; + }, + async hdel(key, ...fields) { + const h = hashes.get(key); + if (!h) return 0; + let n = 0; + for (const f of fields) { + if (h.delete(f)) n += 1; + } + return n; + }, + }; +} + +function fakeUsage(input: number, output: number): Usage { + return { + input, + output, + cacheRead: 0, + cacheWrite: 0, + totalTokens: input + output, + cost: { total: 0, input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + }; +} + +describe("SessionManager + RedisSessionStorage", () => { + it("persists appended assistant messages into Redis and reloads them via open()", async () => { + const redis = createFakeRedis(); + const storage = await RedisSessionStorage.create({ client: redis }); + const sessionDir = "/sessions/proj"; + + const manager = SessionManager.create("/cwd", sessionDir, storage); + manager.appendMessage({ + role: "assistant", + provider: "anthropic", + model: "claude-3-7-sonnet", + content: [{ type: "text", text: "hi" }], + usage: fakeUsage(10, 5), + api: "anthropic-messages", + stopReason: "stop", + timestamp: Date.now(), + }); + + const sessionFile = manager.getSessionFile(); + expect(sessionFile).toBeDefined(); + const sessionFilePath = sessionFile as string; + expect(sessionFilePath.startsWith(sessionDir)).toBe(true); + + // `appendMessage` queues the cold-path rewrite onto SessionManager's + // internal persist chain via a fire-and-forget call. `flush()` awaits + // that chain; `drain()` mops up the storage-level pending tail. + await manager.flush(); + await storage.drain(); + await manager.close(); + + // Redis now contains the JSONL — header + one message entry. + const stored = redis.strings.get(`omp:sessions:file:${sessionFilePath}`); + expect(stored).toBeDefined(); + const lines = (stored as string).trim().split("\n"); + expect(lines.length).toBeGreaterThanOrEqual(2); + const header = JSON.parse(lines[0]); + expect(header.type).toBe("session"); + const msg = JSON.parse(lines[lines.length - 1]); + expect(msg.type).toBe("message"); + expect(msg.message.role).toBe("assistant"); + expect(msg.message.content[0].text).toBe("hi"); + + // Reopening the session through SessionManager.open should recover the leaf. + const reopened = await SessionManager.open(sessionFilePath, sessionDir, storage); + const leaf = reopened.getLeafEntry(); + expect(leaf).toBeDefined(); + expect(leaf?.type).toBe("message"); + await reopened.close(); + }); + + it("SessionManager.list returns Redis-backed sessions for the cwd", async () => { + const redis = createFakeRedis(); + const storage = await RedisSessionStorage.create({ client: redis }); + const sessionDir = "/sessions/list-proj"; + + const a = SessionManager.create("/cwd", sessionDir, storage); + a.appendMessage({ + role: "assistant", + provider: "anthropic", + model: "claude-3-7-sonnet", + content: [{ type: "text", text: "alpha" }], + usage: fakeUsage(1, 1), + api: "anthropic-messages", + stopReason: "stop", + timestamp: Date.now(), + }); + await a.flush(); + await storage.drain(); + await a.close(); + + const b = SessionManager.create("/cwd", sessionDir, storage); + b.appendMessage({ + role: "assistant", + provider: "anthropic", + model: "claude-3-7-sonnet", + content: [{ type: "text", text: "beta" }], + usage: fakeUsage(1, 1), + api: "anthropic-messages", + stopReason: "stop", + timestamp: Date.now(), + }); + await b.flush(); + await storage.drain(); + await b.close(); + + const aFile = a.getSessionFile(); + const bFile = b.getSessionFile(); + expect(aFile).toBeDefined(); + expect(bFile).toBeDefined(); + + const sessions = await SessionManager.list("/cwd", sessionDir, storage); + const sessionFiles = sessions.map(s => s.path).sort(); + expect(sessionFiles).toContain(aFile as string); + expect(sessionFiles).toContain(bFile as string); + }); +}); diff --git a/packages/coding-agent/test/session/redis-session-storage.test.ts b/packages/coding-agent/test/session/redis-session-storage.test.ts new file mode 100644 index 000000000..03e27eda9 --- /dev/null +++ b/packages/coding-agent/test/session/redis-session-storage.test.ts @@ -0,0 +1,317 @@ +/** + * Functional tests for {@link RedisSessionStorage}. Driven by a hand-rolled + * fake Redis client so the suite runs without a live server. + * + * The harness mirrors only the surface the storage actually uses; it is *not* + * a general-purpose mock. Each test exercises one contract: + * + * - the mirror keeps `existsSync`/`statSync`/`readTextSync`/`listFilesSync` + * coherent with `writeText`/`writer.writeLineSync`; + * - `drain()` waits for fire-and-forget background writes; + * - `deleteSessionWithArtifacts` removes both the JSONL key and any sidecar + * keys under the artifacts prefix; + * - `refresh()` re-loads the keyspace, so a peer process's writes become + * visible after an explicit re-scan. + */ + +import { beforeEach, describe, expect, it } from "bun:test"; +import { + RedisSessionStorage, + type RedisSessionStorageClient, +} from "@oh-my-pi/pi-coding-agent/session/redis-session-storage"; + +interface FakeRedisCall { + method: string; + args: unknown[]; +} + +interface FakeRedis extends RedisSessionStorageClient { + calls: FakeRedisCall[]; + strings: Map; + hashes: Map>; + /** Override the next call to `method` to reject with `error`. */ + failNext(method: string, error: Error): void; +} + +function createFakeRedis(): FakeRedis { + const strings = new Map(); + const hashes = new Map>(); + const calls: FakeRedisCall[] = []; + const failures = new Map(); + + const checkFailure = (method: string): void => { + const queue = failures.get(method); + if (!queue || queue.length === 0) return; + throw queue.shift() as Error; + }; + + const record = (method: string, args: unknown[]): void => { + calls.push({ method, args }); + }; + + const getHash = (key: string): Map => { + let h = hashes.get(key); + if (!h) { + h = new Map(); + hashes.set(key, h); + } + return h; + }; + + const client: FakeRedis = { + calls, + strings, + hashes, + failNext(method: string, error: Error): void { + const queue = failures.get(method) ?? []; + queue.push(error); + failures.set(method, queue); + }, + async get(key) { + record("get", [key]); + checkFailure("get"); + return strings.has(key) ? (strings.get(key) as string) : null; + }, + async set(key, value) { + record("set", [key, value]); + checkFailure("set"); + strings.set(key, value); + return "OK"; + }, + async append(key, value) { + record("append", [key, value]); + checkFailure("append"); + const current = strings.get(key) ?? ""; + const next = current + value; + strings.set(key, next); + return next.length; + }, + async del(...keys) { + record("del", keys); + checkFailure("del"); + let deleted = 0; + for (const k of keys) { + if (strings.delete(k)) deleted += 1; + } + return deleted; + }, + async rename(src, dst) { + record("rename", [src, dst]); + checkFailure("rename"); + if (!strings.has(src)) { + throw new Error("ERR no such key"); + } + strings.set(dst, strings.get(src) as string); + strings.delete(src); + return "OK"; + }, + async scan(cursor, ...rest) { + record("scan", [cursor, ...rest]); + checkFailure("scan"); + let pattern = "*"; + for (let i = 0; i < rest.length; i++) { + if (String(rest[i]).toUpperCase() === "MATCH") { + pattern = String(rest[i + 1] ?? "*"); + } + } + const regex = new RegExp(`^${pattern.replace(/[.+?^${}()|[\]\\]/g, "\\$&").replace(/\*/g, ".*")}$`); + const matches = Array.from(strings.keys()).filter(k => regex.test(k)); + return ["0", matches]; + }, + async hset(key, field, value) { + record("hset", [key, field, value]); + checkFailure("hset"); + getHash(key).set(field, value); + return 1; + }, + async hgetall(key) { + record("hgetall", [key]); + checkFailure("hgetall"); + const h = hashes.get(key); + if (!h) return {}; + const out: Record = {}; + for (const [k, v] of h) out[k] = v; + return out; + }, + async hdel(key, ...fields) { + record("hdel", [key, ...fields]); + checkFailure("hdel"); + const h = hashes.get(key); + if (!h) return 0; + let n = 0; + for (const f of fields) { + if (h.delete(f)) n += 1; + } + return n; + }, + }; + + return client; +} + +describe("RedisSessionStorage", () => { + let redis: FakeRedis; + + beforeEach(() => { + redis = createFakeRedis(); + }); + + it("mirrors writeText into Redis and exposes content via sync reads", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await storage.writeText("/sessions/p/a.jsonl", "line1\nline2\n"); + + expect(storage.existsSync("/sessions/p/a.jsonl")).toBe(true); + expect(storage.readTextSync("/sessions/p/a.jsonl")).toBe("line1\nline2\n"); + expect(redis.strings.get("omp:sessions:file:/sessions/p/a.jsonl")).toBe("line1\nline2\n"); + + const stat = storage.statSync("/sessions/p/a.jsonl"); + expect(stat.size).toBe(12); + expect(typeof stat.mtimeMs).toBe("number"); + }); + + it("listFilesSync returns only direct children matching the glob", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await storage.writeText("/dir/a.jsonl", "x"); + await storage.writeText("/dir/b.jsonl", "y"); + await storage.writeText("/dir/sub/c.jsonl", "z"); // nested — not a direct child + await storage.writeText("/dir/note.bak", "skip"); + + const jsonl = storage.listFilesSync("/dir", "*.jsonl").sort(); + expect(jsonl).toEqual(["/dir/a.jsonl", "/dir/b.jsonl"]); + + const bak = storage.listFilesSync("/dir", "*.bak"); + expect(bak).toEqual(["/dir/note.bak"]); + }); + + it("statSync mtimes are strictly monotonic across rapid writes", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await storage.writeText("/s/a", "1"); + await storage.writeText("/s/b", "2"); + await storage.writeText("/s/c", "3"); + + const a = storage.statSync("/s/a").mtimeMs; + const b = storage.statSync("/s/b").mtimeMs; + const c = storage.statSync("/s/c").mtimeMs; + expect(b).toBeGreaterThan(a); + expect(c).toBeGreaterThan(b); + }); + + it("writer.writeLineSync 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'); + + // Mirror reflects the writes synchronously. + expect(storage.readTextSync("/sessions/p/session.jsonl")).toBe('{"type":"session"}\n{"type":"message"}\n'); + + // Redis has not necessarily caught up yet — drain to force. + await storage.drain(); + expect(redis.strings.get("omp:sessions:file:/sessions/p/session.jsonl")).toBe( + '{"type":"session"}\n{"type":"message"}\n', + ); + + await writer.close(); + }); + + it("flags='w' truncates both mirror and Redis", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + 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.close(); + + expect(storage.readTextSync("/sessions/p/keep.jsonl")).toBe("fresh\n"); + expect(redis.strings.get("omp:sessions:file:/sessions/p/keep.jsonl")).toBe("fresh\n"); + }); + + it("drain() surfaces writer errors so background failures are observable", async () => { + 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"); + + await expect(storage.drain()).rejects.toThrow("redis exploded"); + expect(writer.getError()?.message).toBe("redis exploded"); + }); + + it("deleteSessionWithArtifacts removes JSONL plus any sidecar keys", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await storage.writeText("/sessions/p/s1.jsonl", "session\n"); + await storage.writeText("/sessions/p/s1/draft.txt", "draft body"); + await storage.writeText("/sessions/p/s1/sub/notes", "more"); + await storage.writeText("/sessions/p/other.jsonl", "untouched\n"); + + await storage.deleteSessionWithArtifacts("/sessions/p/s1.jsonl"); + + expect(storage.existsSync("/sessions/p/s1.jsonl")).toBe(false); + expect(storage.existsSync("/sessions/p/s1/draft.txt")).toBe(false); + expect(storage.existsSync("/sessions/p/s1/sub/notes")).toBe(false); + expect(storage.existsSync("/sessions/p/other.jsonl")).toBe(true); + expect(redis.strings.has("omp:sessions:file:/sessions/p/s1.jsonl")).toBe(false); + expect(redis.strings.has("omp:sessions:file:/sessions/p/s1/draft.txt")).toBe(false); + expect(redis.strings.has("omp:sessions:file:/sessions/p/other.jsonl")).toBe(true); + }); + + it("rename moves content and meta atomically inside the mirror", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await storage.writeText("/sessions/p/orig.jsonl", "payload\n"); + const originalMtime = storage.statSync("/sessions/p/orig.jsonl").mtimeMs; + + await storage.rename("/sessions/p/orig.jsonl", "/sessions/p/renamed.jsonl"); + expect(storage.existsSync("/sessions/p/orig.jsonl")).toBe(false); + expect(storage.readTextSync("/sessions/p/renamed.jsonl")).toBe("payload\n"); + expect(storage.statSync("/sessions/p/renamed.jsonl").mtimeMs).toBe(originalMtime); + expect(redis.strings.get("omp:sessions:file:/sessions/p/renamed.jsonl")).toBe("payload\n"); + expect(redis.strings.has("omp:sessions:file:/sessions/p/orig.jsonl")).toBe(false); + }); + + it("rename rolls back the mirror when Redis RENAME fails", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await storage.writeText("/sessions/p/a.jsonl", "keep\n"); + redis.failNext("rename", new Error("ERR redis rejected rename")); + + await expect(storage.rename("/sessions/p/a.jsonl", "/sessions/p/b.jsonl")).rejects.toThrow( + "ERR redis rejected rename", + ); + + expect(storage.existsSync("/sessions/p/a.jsonl")).toBe(true); + expect(storage.existsSync("/sessions/p/b.jsonl")).toBe(false); + }); + + it("refresh() reloads the mirror from Redis after out-of-band writes", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + // Simulate a peer process writing directly to Redis. + redis.strings.set("omp:sessions:file:/peer/x.jsonl", "from peer\n"); + const peerHash = redis.hashes.get("omp:sessions:meta") ?? new Map(); + peerHash.set("/peer/x.jsonl", String(Date.now() + 5_000)); + redis.hashes.set("omp:sessions:meta", peerHash); + + expect(storage.existsSync("/peer/x.jsonl")).toBe(false); + await storage.refresh(); + expect(storage.existsSync("/peer/x.jsonl")).toBe(true); + expect(storage.readTextSync("/peer/x.jsonl")).toBe("from peer\n"); + }); + + it("readTextPrefix returns at most maxBytes from the head", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await storage.writeText("/sessions/p/big.jsonl", "abcdefghij"); + + expect(await storage.readTextPrefix("/sessions/p/big.jsonl", 4)).toBe("abcd"); + expect(await storage.readTextPrefix("/sessions/p/big.jsonl", 100)).toBe("abcdefghij"); + expect(await storage.readTextPrefix("/sessions/p/big.jsonl", 0)).toBe(""); + }); + + it("custom prefix isolates keyspaces", async () => { + const storage = await RedisSessionStorage.create({ client: redis, prefix: "proj-a:" }); + await storage.writeText("/sessions/x.jsonl", "hello\n"); + expect(redis.strings.has("proj-a:file:/sessions/x.jsonl")).toBe(true); + expect(redis.strings.has("omp:sessions:file:/sessions/x.jsonl")).toBe(false); + }); + + it("unlink on a missing key throws ENOENT", async () => { + const storage = await RedisSessionStorage.create({ client: redis }); + await expect(storage.unlink("/sessions/p/ghost.jsonl")).rejects.toMatchObject({ code: "ENOENT" }); + }); +});