feat(session): added RedisSessionStorage with redis-backed session persistence

- Added RedisSessionStorage with Redis-backed session JSONL persistence and storage options.
- Added create/refresh/list/read behavior with in-memory mirror, SCAN hydration, and monotonic mtime tracking.
- Added package exports for SessionStorage backends and a Redis session SDK example with usage guidance.
- Added in-memory fake Redis tests covering persistence restore, session listing, rename, and error-injection paths.
This commit is contained in:
can1357
2026-05-28 16:58:57 +02:00
parent e2ddf1a0d1
commit 47bd4fd766
6 changed files with 1065 additions and 0 deletions
+5
View File
@@ -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
@@ -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();
+2
View File
@@ -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)
@@ -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<string | null>;
set(key: string, value: string): Promise<unknown>;
append(key: string, value: string): Promise<number>;
del(...keys: string[]): Promise<number>;
rename(src: string, dst: string): Promise<unknown>;
scan(cursor: string, ...args: string[]): Promise<[string, string[]]>;
hset(key: string, field: string, value: string): Promise<unknown>;
hgetall(key: string): Promise<Record<string, string>>;
hdel(key: string, ...fields: string[]): Promise<unknown>;
}
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<string, MirrorEntry>();
readonly #writers = new Set<RedisSessionStorageWriter>();
#nextMtimeMs = 0;
#pendingTail: Promise<void> = 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<RedisSessionStorage> {
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<void> {
this.#mirror.clear();
const filePrefix = this.#fileKey("");
const metaRaw = await this.#client.hgetall(this.#metaKey());
const meta: Record<string, string> = metaRaw ?? {};
const seen = new Set<string>();
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<void> {
// 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>): 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<boolean> {
// 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<string> {
const entry = this.#mirror.get(path);
if (!entry) throw enoent(path);
return entry.content;
}
async readTextPrefix(path: string, maxBytes: number): Promise<string> {
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<void> {
const mtimeMs = this.#allocMtimeMs();
this.#mirror.set(path, { content, mtimeMs });
await this.#writeRemote(path, content, mtimeMs);
}
async rename(src: string, dst: string): Promise<void> {
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<void> {
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<void> {
await this.unlink(sessionPath);
// Mirror artifacts live under `<sessionPath without .jsonl>/...`. 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<void> {
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<void> {
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>): void {
this.#trackPending(promise);
}
async #writeRemote(path: string, content: string, mtimeMs: number): Promise<void> {
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<void> = 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<void>): Promise<void> {
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<void> {
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<void> {
if (this.#error) throw this.#error;
await this.#enqueueRaw(async () => {});
if (this.#error) throw this.#error;
}
async fsync(): Promise<void> {
// 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<void> {
if (this.#closed) return;
this.#closed = true;
try {
await this.flush();
} finally {
this.#storage._writerClosed(this);
}
}
getError(): Error | undefined {
return this.#error;
}
}
@@ -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<string, string>;
hashes: Map<string, Map<string, string>>;
}
function createFakeRedis(): FakeRedis {
const strings = new Map<string, string>();
const hashes = new Map<string, Map<string, string>>();
const getHash = (key: string): Map<string, string> => {
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<string, string> = {};
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);
});
});
@@ -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<string, string>;
hashes: Map<string, Map<string, string>>;
/** Override the next call to `method` to reject with `error`. */
failNext(method: string, error: Error): void;
}
function createFakeRedis(): FakeRedis {
const strings = new Map<string, string>();
const hashes = new Map<string, Map<string, string>>();
const calls: FakeRedisCall[] = [];
const failures = new Map<string, Error[]>();
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<string, string> => {
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<string, string> = {};
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<string, string>();
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" });
});
});