Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
24 kB · 568 lines
JavaScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569import { test } from 'node:test'import assert from 'node:assert/strict'import { mkdir, mkdtemp, open, readFile, realpath, rename, rm, writeFile } from 'node:fs/promises'import { randomUUID } from 'node:crypto'import { tmpdir } from 'node:os'import { dirname, join } from 'node:path'
import { apply, resolveDshHome, resolveStorePath, workspaceStorePath } from '../lib/index.js'import { parseStore } from '../lib/store.js'
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms))
/** * Poll a predicate until it holds or the deadline passes. * @param predicate - Condition to evaluate. * @param timeoutMs - Budget. * @returns whether the predicate held. */async function until(predicate, timeoutMs = 3000) { const deadline = Date.now() + timeoutMs while (Date.now() < deadline) { if (await predicate()) return true await sleep(5) } return predicate()}
/** * Publish a store document the way the JSON backend does: temp file plus * `rename()`, which replaces the inode and so is only visible to a watcher * bound to the containing directory. * @param path - Store path to replace. * @param document - Document to serialize. */async function publishAtomic(path, document) { const tmp = join(dirname(path), `.${randomUUID()}.tmp`) const handle = await open(tmp, 'wx', 0o600) try { await handle.writeFile(`${JSON.stringify(document, null, 2)}\n`, 'utf8') } finally { await handle.close() } await rename(tmp, path)}
/** * One workspace entity with the real registry's observable semantics: an * attach for an already-accounted id is write-free and idempotent, and an * attach whose stored header cannot validate throws before touching anything. * @param registry - Owning fake registry. * @param id - Workspace id. * @param path - Canonical workspace path. * @param title - Display title. * @returns the entity. */function createWorkspace(registry, id, path, title) { let sessionIds = [] return { id, path, title, get sessionIds() { return [...sessionIds] }, adopt(sessionId) { sessionIds = [sessionId, ...sessionIds] }, async attachSession(sessionId) { if (sessionIds.includes(sessionId)) return const header = registry.headers.get(sessionId) if (header === undefined) throw new Error(`cannot validate session '${sessionId}': session persistence holds no such session`) if (header.cwd !== path) throw new Error(`cannot attach session '${sessionId}' to workspace '${path}': its cwd resolves to '${header.cwd}'`) sessionIds = [sessionId, ...sessionIds] await registry.record('attach', `${path}#${sessionId}`) }, }}
/** * A stand-in for `ctx.workspaceRegistry` that records every durable mutation * and republishes the whole store file after each one, exactly as the domain * layer does — that republish is what makes the convergence property testable. * @param options - Fixture state. * @returns the fake registry. */function createFakeRegistry(options) { const registry = { headers: new Map(Object.entries(options.headers ?? {})), directories: new Set(options.directories ?? []), archives: new Set(options.archived ?? []), workspaces: new Map(), writes: [], listThrows: false, /** Completed durable mutations, in order — the republish has landed. */ async record(operation, detail) { await registry.publish() registry.writes.push({ operation, detail }) }, publish() { if (options.storeFile === undefined) return Promise.resolve() const workspaces = {} for (const workspace of registry.workspaces.values()) { workspaces[workspace.id] = { path: workspace.path, title: workspace.title, sessionIds: workspace.sessionIds, createdAt: '2026-01-01T00:00:00.000Z', updatedAt: '2026-01-01T00:00:00.000Z', } } const document = { unit: { name: 'workspace', version: 2 }, global: { initialized: true, workspaceIds: [...registry.workspaces.keys()], archivedSessionIds: [...registry.archives], }, tables: { workspaces }, } return publishAtomic(options.storeFile, document) }, list() { if (registry.listThrows) throw new Error('workspace registry is not started yet') // The real `list()` returns live entities, not snapshots: callers drive // `attachSession` on exactly these objects. return [...registry.workspaces.values()] }, get archivedSessionIds() { return [...registry.archives] }, async create(path, title) { const existing = [...registry.workspaces.values()].find((workspace) => workspace.path === path) if (existing !== undefined) return existing if (!registry.directories.has(path)) throw new Error(`cannot create a workspace at '${path}': path is not a directory`) const workspace = createWorkspace(registry, `ws-${registry.workspaces.size + 1}`, path, title ?? path) registry.workspaces.set(workspace.id, workspace) await registry.record('create', path) return workspace }, async archiveSession(sessionId) { if (registry.archives.has(sessionId)) return if (!registry.headers.has(sessionId)) throw new Error(`cannot archive session '${sessionId}': live sessions and session persistence hold no such session`) registry.archives.add(sessionId) await registry.record('archive', sessionId) }, }
for (const seed of options.workspaces ?? []) { const workspace = createWorkspace(registry, seed.id, seed.path, seed.title ?? seed.path) for (const sessionId of [...(seed.sessionIds ?? [])].reverse()) workspace.adopt(sessionId) registry.workspaces.set(workspace.id, workspace) } return registry}
/** * A minimal plugin context: the fake registry, a capturing logger, and an * `effect` that runs the callback and remembers its disposer. * @param registry - Fake registry to expose as `ctx.workspaceRegistry`. * @returns the context plus test affordances. */function createFakeContext(registry) { const warnings = [] const disposers = [] const ctx = { workspaceRegistry: registry, logger: { debug: () => {}, info: () => {}, warn: (...args) => warnings.push(args.map(String).join(' ')), error: () => {}, }, // Production always reaches the durable accounts through the storage // domain (the registry injects it), so the fake exposes the same shape: // raw records, not the cwd-filtered `Workspace.sessionIds` view. get: (name) => (name !== 'storageDomain' ? undefined : { get: (domain) => (domain !== 'workspace' ? undefined : { table: (table) => { if (table !== 'workspaces') throw new Error(`domain 'workspace' declares no table '${table}'`) return { entries: () => [...registry.workspaces.values()] .map((workspace) => [workspace.id, { path: workspace.path, sessionIds: workspace.sessionIds }]) [Symbol.iterator](), } }, }), }), effect: (callback) => { disposers.push(callback()) return () => {} }, } return { ctx, warnings, async dispose() { for (const dispose of disposers.splice(0).reverse()) await dispose() }, }}
/** * Create a temp harness home, point `$DSH_HOME` at it, and register cleanup. * @param t - Test context. * @param options - `workspaces`, `headers`, `archived`, `directories`, `file`. * @returns the registry, context, store path, and mutation log. */async function fixture(t, options = {}) { const home = await realpath(await mkdtemp(join(tmpdir(), 'fix-workspace-'))) const storages = join(home, 'storages') await mkdir(storages, { recursive: true }) const storeFile = join(storages, 'workspace.json')
const previousHome = process.env.DSH_HOME process.env.DSH_HOME = home t.after(async () => { if (previousHome === undefined) delete process.env.DSH_HOME else process.env.DSH_HOME = previousHome await rm(home, { recursive: true, force: true }) })
const directories = new Set(options.directories ?? options.workspaces?.map((workspace) => workspace.path) ?? []) const registry = createFakeRegistry({ ...options, directories, storeFile }) if (options.file !== undefined) await writeFile(storeFile, options.file, 'utf8')
const harness = createFakeContext(registry) t.after(() => harness.dispose())
return { home, storages, storeFile, registry, ...harness }}
/** * Serialize a store document for seeding. * @param workspaces - Workspace records keyed by id. * @param global - Global slot overrides. * @returns the file content. */function storeText(workspaces, global = {}) { return `${JSON.stringify({ unit: { name: 'workspace', version: 2 }, global: { initialized: true, workspaceIds: Object.keys(workspaces), archivedSessionIds: [], ...global }, tables: { workspaces }, }, null, 2)}\n`}
const record = (path, sessionIds, title = path) => ({ path, title, sessionIds })
test('the store path comes from the domain unit, and a directory-backed unit is refused', () => { const domainCtx = (unit) => ({ get: (name) => (name === 'storageDomain' ? { get: () => ({ unit }) } : undefined) })
// The open domain's own unit is authoritative; a relocated backend root or a // stale file at the default path can then never be reconciled by mistake. assert.deepEqual(resolveStorePath(domainCtx({ path: '/elsewhere/workspace.json', descriptor: {} })), { path: '/elsewhere/workspace.json', source: 'domain unit', })
// A per-record unit stores records in a directory: refuse rather than read it. const perRecord = resolveStorePath(domainCtx({ descriptor: { layout: 'per-record' }, dir: '/elsewhere/workspace' })) assert.equal(perRecord.path, undefined) assert.match(perRecord.reason, /per-record/)
// A unit we cannot identify is also refused, not guessed at. assert.equal(resolveStorePath(domainCtx({ descriptor: {} })).path, undefined) assert.equal(resolveStorePath(domainCtx({ descriptor: {} })).reason !== undefined, true)
// Only a missing domain falls back to the composed default. const assumed = resolveStorePath({ get: () => undefined }) assert.equal(assumed.path, workspaceStorePath()) assert.match(assumed.source, /assumed/)})
test('$DSH_HOME resolution mirrors resolveDshHome', () => { assert.equal(workspaceStorePath({ DSH_HOME: '/tmp/home' }), '/tmp/home/storages/workspace.json') assert.equal(workspaceStorePath({ DSH_HOME: ' ' }), join(resolveDshHome({}), 'storages', 'workspace.json')) const previous = process.env.DSH_HOME process.env.DSH_HOME = '/tmp/late-home' try { assert.equal(workspaceStorePath(), '/tmp/late-home/storages/workspace.json') } finally { if (previous === undefined) delete process.env.DSH_HOME else process.env.DSH_HOME = previous }})
test('attaches a session another instance published into the store', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, warnings, storeFile } = await fixture(t, { headers: { 'session-a': { cwd: ws }, 'session-b': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: ['session-a'] }], // What instance B published: its union includes session-b, ours does not. file: storeText({ w1: record(ws, ['session-b', 'session-a']) }), })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => registry.writes.length > 0), true, 'the pass never ran') assert.deepEqual(registry.writes, [{ operation: 'attach', detail: `${ws}#session-b` }]) assert.deepEqual([...registry.workspaces.values()][0].sessionIds, ['session-b', 'session-a']) assert.deepEqual(warnings, [])
const published = parseStore(await readFile(storeFile, 'utf8')) assert.deepEqual(published.workspaces[0].sessionIds, ['session-b', 'session-a'])})
test('converges, then performs no durable write for repeated store events', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, storeFile } = await fixture(t, { headers: { 'session-a': { cwd: ws }, 'session-b': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: ['session-a'] }], file: storeText({ w1: record(ws, ['session-b', 'session-a']) }), })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => registry.writes.length === 1), true) // The attach republished the store; give the watcher and the follow-up pass // time to run and prove they add nothing. await sleep(250) assert.equal(registry.writes.length, 1, `unexpected extra writes: ${JSON.stringify(registry.writes)}`)
// Touch the file with byte-identical content several times. const content = await readFile(storeFile, 'utf8') for (let index = 0; index < 3; index += 1) { await publishAtomic(storeFile, JSON.parse(content)) await sleep(30) } await sleep(250) assert.equal(registry.writes.length, 1, `steady state wrote: ${JSON.stringify(registry.writes)}`)})
test('adopts a workspace the file describes but this process does not own', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, storeFile } = await fixture(t, { headers: { 'session-a': { cwd: ws } }, directories: [ws], file: storeText({ w9: record(ws, ['session-a'], 'remote-title') }), })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => registry.writes.length >= 2), true) assert.deepEqual(registry.writes.map((write) => write.operation), ['create', 'attach']) const workspace = [...registry.workspaces.values()][0] assert.equal(workspace.path, ws) assert.equal(workspace.title, 'remote-title') assert.deepEqual(workspace.sessionIds, ['session-a']) assert.deepEqual(parseStore(await readFile(storeFile, 'utf8')).workspaces[0].sessionIds, ['session-a'])})
test('archives a session the file archived', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx } = await fixture(t, { headers: { 'session-a': { cwd: ws }, 'session-b': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: ['session-a'] }], archived: ['session-a'], file: storeText({ w1: record(ws, ['session-a']) }, { archivedSessionIds: ['session-a', 'session-b'] }), })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => registry.writes.length === 1), true) assert.deepEqual(registry.writes, [{ operation: 'archive', detail: 'session-b' }]) assert.deepEqual(registry.archivedSessionIds, ['session-a', 'session-b'])})
test('never writes when the store is malformed, and says why', async (t) => { const { registry, ctx, warnings } = await fixture(t, { file: '{ this is not json' })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => warnings.length > 0), true) await sleep(200) assert.deepEqual(registry.writes, []) assert.match(warnings[0], /leaving the workspace store untouched/)})
test('never writes when the store has a foreign or older unit header', async (t) => { const { registry, ctx, warnings } = await fixture(t, { file: `${JSON.stringify({ unit: { name: 'workspace', version: 1 }, global: { initialized: true, workspaceIds: [], archivedSessionIds: [] }, tables: { workspaces: {} } })}\n` })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => warnings.length > 0), true) assert.match(warnings[0], /foreign or unsupported unit header/) assert.deepEqual(registry.writes, [])})
test('defers while another instance holds a pending mutation, then applies once it clears', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, storeFile } = await fixture(t, { headers: { 'session-b': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: [] }], file: storeText({ w1: record(ws, ['session-b']) }, { pendingMutation: { operation: 'create', workspaceId: 'w2' } }), })
apply(ctx, { debounceMs: 10, retryOnPendingMs: 20 }) await sleep(150) assert.deepEqual(registry.writes, [], 'a pending mutation must suspend reconciliation')
await writeFile(storeFile, storeText({ w1: record(ws, ['session-b']) }), 'utf8') assert.equal(await until(() => registry.writes.length === 1), true) assert.deepEqual(registry.writes, [{ operation: 'attach', detail: `${ws}#session-b` }])})
test('a pendingMutation marker that never clears is retried a bounded number of times, then left to the watcher', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, warnings, storeFile } = await fixture(t, { headers: { 'session-b': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: [] }], file: storeText({ w1: record(ws, ['session-b']) }, { pendingMutation: { operation: 'create', workspaceId: 'w2' } }), })
apply(ctx, { debounceMs: 10, retryOnPendingMs: 20 }) // The retry budget (8 attempts at 20ms) is spent well inside this window. assert.equal(await until(() => warnings.some((warning) => warning.includes('still carries a pendingMutation'))), true) await sleep(250) assert.deepEqual(registry.writes, [], 'a pending mutation must suspend reconciliation entirely') assert.equal( warnings.filter((warning) => warning.includes('still carries a pendingMutation')).length, 1, 'the suspension must be reported exactly once, not per retry', )
// The marker clearing is itself a write, so the watcher — not a poll — is what // resumes reconciliation after the budget is spent. await writeFile(storeFile, storeText({ w1: record(ws, ['session-b']) }), 'utf8') assert.equal(await until(() => registry.writes.length === 1), true) assert.deepEqual(registry.writes, [{ operation: 'attach', detail: `${ws}#session-b` }])})
test('leaves a store alone while it is still bootstrapping', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx } = await fixture(t, { headers: { 'session-a': { cwd: ws } }, directories: [ws], file: storeText({ w1: record(ws, ['session-a']) }, { initialized: false }), })
apply(ctx, { debounceMs: 10 }) await sleep(200) assert.deepEqual(registry.writes, [])})
test('survives an unknown session and still attaches the sessions it can validate', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, warnings } = await fixture(t, { headers: { 'session-real': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: [] }], file: storeText({ w1: record(ws, ['session-ghost', 'session-real']) }), })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => registry.writes.length > 0), true) await sleep(200) assert.deepEqual(registry.writes, [{ operation: 'attach', detail: `${ws}#session-real` }]) assert.equal(warnings.some((warning) => warning.includes('session-ghost')), true)})
test('warns and skips a workspace whose directory no longer exists', async (t) => { const gone = join(tmpdir(), `fix-workspace-gone-${randomUUID()}`) const { registry, ctx, warnings } = await fixture(t, { headers: { 'session-a': { cwd: gone } }, file: storeText({ w1: record(gone, ['session-a']) }), })
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => warnings.length > 0), true) await sleep(200) assert.deepEqual(registry.writes, []) assert.equal(warnings.some((warning) => warning.includes('cannot adopt workspace')), true)})
test('treats a missing store as nothing to do, then reconciles once it appears', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, storeFile } = await fixture(t, { headers: { 'session-a': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: [] }], })
apply(ctx, { debounceMs: 10 }) await sleep(120) assert.deepEqual(registry.writes, [])
await writeFile(storeFile, storeText({ w1: record(ws, ['session-a']) }), 'utf8') assert.equal(await until(() => registry.writes.length === 1), true) assert.deepEqual(registry.writes, [{ operation: 'attach', detail: `${ws}#session-a` }])})
test('stops reconciling once its effect is disposed', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, dispose, storeFile } = await fixture(t, { headers: { 'session-a': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: [] }], file: storeText({ w1: record(ws, []) }), })
apply(ctx, { debounceMs: 10 }) await sleep(120) assert.deepEqual(registry.writes, []) await dispose()
await writeFile(storeFile, storeText({ w1: record(ws, ['session-a']) }), 'utf8') await sleep(250) assert.deepEqual(registry.writes, [], 'a disposed plugin must not keep watching')})
test('never throws into the tree when the registry cannot be read', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const { registry, ctx, warnings } = await fixture(t, { headers: { 'session-a': { cwd: ws } }, directories: [ws], workspaces: [{ id: 'w1', path: ws, sessionIds: [] }], file: storeText({ w1: record(ws, ['session-a']) }), }) registry.listThrows = true
apply(ctx, { debounceMs: 10 }) assert.equal(await until(() => warnings.some((warning) => warning.includes('reconcile pass failed'))), true) assert.deepEqual(registry.writes, [])})
test('falls back to polling the store path when its directory cannot be watched', async (t) => { const ws = await realpath(await mkdtemp(join(tmpdir(), 'ws-'))) t.after(() => rm(ws, { recursive: true, force: true })) const home = await realpath(await mkdtemp(join(tmpdir(), 'fix-workspace-'))) t.after(() => rm(home, { recursive: true, force: true }))
const previousHome = process.env.DSH_HOME process.env.DSH_HOME = home t.after(() => { if (previousHome === undefined) delete process.env.DSH_HOME else process.env.DSH_HOME = previousHome })
// No `storages/` directory yet: the directory watch cannot be installed. const storeFile = join(home, 'storages', 'workspace.json') const registry = createFakeRegistry({ headers: { 'session-a': { cwd: ws } }, directories: [ws], storeFile: undefined, workspaces: [{ id: 'w1', path: ws, sessionIds: [] }], }) const { ctx, warnings, dispose } = createFakeContext(registry) t.after(dispose)
apply(ctx, { debounceMs: 10, pollIntervalMs: 20 }) assert.equal(await until(() => warnings.some((warning) => warning.includes('polling the store path'))), true) await sleep(50) assert.deepEqual(registry.writes, [])
await mkdir(join(home, 'storages'), { recursive: true }) await writeFile(storeFile, storeText({ w1: record(ws, ['session-a']) }), 'utf8') assert.equal(await until(() => registry.writes.length === 1), true, warnings.join('\n')) await dispose()})