fix(session-persistence): handle Windows directory fsync
This commit is contained in:
@@ -155,18 +155,20 @@ describe.skipIf(!existsSync(cliBin))('dsh-cli-demo BUILT bin', () => {
|
|||||||
}
|
}
|
||||||
}, 30_000)
|
}, 30_000)
|
||||||
|
|
||||||
it.each([
|
describe.skipIf(process.platform === 'win32')('POSIX signal delivery', () => {
|
||||||
['SIGINT', 130],
|
it.each([
|
||||||
['SIGTERM', 143],
|
['SIGINT', 130],
|
||||||
] as const)('cancels and disposes on %s with exit %i', async (signal, code) => {
|
['SIGTERM', 143],
|
||||||
consumer = await makeConsumer()
|
] as const)('cancels and disposes on %s with exit %i', async (signal, code) => {
|
||||||
const result = await runBuiltBin(
|
consumer = await makeConsumer()
|
||||||
consumer,
|
const result = await runBuiltBin(
|
||||||
['--config', './cordis.yml', '--output-format', 'stream-json', 'hang'],
|
consumer,
|
||||||
signal,
|
['--config', './cordis.yml', '--output-format', 'stream-json', 'hang'],
|
||||||
)
|
signal,
|
||||||
expect(result, JSON.stringify(result)).toMatchObject({ code, signal: null })
|
)
|
||||||
expect(result.stdout).toContain('"kind":"aborted"')
|
expect(result, JSON.stringify(result)).toMatchObject({ code, signal: null })
|
||||||
expect(result.stderr).toContain(`received ${signal}`)
|
expect(result.stdout).toContain('"kind":"aborted"')
|
||||||
}, 30_000)
|
expect(result.stderr).toContain(`received ${signal}`)
|
||||||
|
}, 30_000)
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -21,7 +21,7 @@ The JSONL durable session-persistence backend — a concrete `SessionPersistence
|
|||||||
|
|
||||||
## Durability and crash semantics
|
## Durability and crash semantics
|
||||||
|
|
||||||
- **Lazy materialization.** `create(meta)` writes nothing; on the first `append`, the backend writes and `fsync`s a temporary file, publishes it without overwrite via a hard link, then `fsync`s the directory. A created-but-never-appended session leaves nothing on disk and is absent from `list`.
|
- **Lazy materialization.** `create(meta)` writes nothing; on the first `append`, the backend writes and `fsync`s a temporary file, publishes it without overwrite via a hard link, then `fsync`s the directory when the host supports it. A created-but-never-appended session leaves nothing on disk and is absent from `list`.
|
||||||
- **Append-only.** Committed events (at or below a flushed `turn/end`) are never rewritten. Subsequent appends are line appends at EOF + `fsync`.
|
- **Append-only.** Committed events (at or below a flushed `turn/end`) are never rewritten. Subsequent appends are line appends at EOF + `fsync`.
|
||||||
- **Crash recovery — preserve valid tail work.** `load` keeps the contiguous valid prefix of an interrupted final turn. It truncates from the first unparsable or sequence-gapped uncommitted record, then appends the synthetic tool, step, and turn closers required by the shared [persistence contract](../../../docs/rfc/implemented/architecture/2026-06-14-session-persistence.md); the same defect at or before the last committed `turn/end` rejects.
|
- **Crash recovery — preserve valid tail work.** `load` keeps the contiguous valid prefix of an interrupted final turn. It truncates from the first unparsable or sequence-gapped uncommitted record, then appends the synthetic tool, step, and turn closers required by the shared [persistence contract](../../../docs/rfc/implemented/architecture/2026-06-14-session-persistence.md); the same defect at or before the last committed `turn/end` rejects.
|
||||||
- **Contiguous-seq.** `append` rejects a batch whose first `seq` does not continue the stored log, and rejects non-JSON-serializable `event.data` naming the offending event type.
|
- **Contiguous-seq.** `append` rejects a batch whose first `seq` does not continue the stored log, and rejects non-JSON-serializable `event.data` naming the offending event type.
|
||||||
@@ -44,3 +44,4 @@ The plugin buffers frozen session events and drains them on flush or disposal. A
|
|||||||
- **Nothing deletes session files** — logs accumulate under `root` until removed externally (the seam has no deletion surface).
|
- **Nothing deletes session files** — logs accumulate under `root` until removed externally (the seam has no deletion surface).
|
||||||
- **Single-process assumption** — per-session serialization and the write cursor live in this process; two processes appending to the same `root` are not coordinated.
|
- **Single-process assumption** — per-session serialization and the write cursor live in this process; two processes appending to the same `root` are not coordinated.
|
||||||
- **Initial materialization requires hard-link support** — first append uses `link()` so same-id races fail instead of overwriting a committed log; a filesystem that cannot create hard links cannot host this backend.
|
- **Initial materialization requires hard-link support** — first append uses `link()` so same-id races fail instead of overwriting a committed log; a filesystem that cannot create hard links cannot host this backend.
|
||||||
|
- **Windows cannot `fsync` directory handles through Node** — the backend tolerates only Windows `EPERM` from directory `fsync`; file-content `fsync` remains mandatory, but a crash can lose a newly published directory entry on a host without an equivalent directory-sync primitive.
|
||||||
|
|||||||
@@ -56,6 +56,9 @@ export class SessionPersistenceJsonl extends SessionPersistence implements Persi
|
|||||||
private root: string
|
private root: string
|
||||||
private coordinator: PersistenceCoordinator<number>
|
private coordinator: PersistenceCoordinator<number>
|
||||||
|
|
||||||
|
/** Runtime-only host-platform seam for directory-sync compatibility tests. */
|
||||||
|
readonly internals: { platform: NodeJS.Platform } = { platform: process.platform }
|
||||||
|
|
||||||
constructor(ctx: Context, public config: Config) {
|
constructor(ctx: Context, public config: Config) {
|
||||||
super(ctx)
|
super(ctx)
|
||||||
// Resolve once so later process.cwd() changes cannot split one backend across roots.
|
// Resolve once so later process.cwd() changes cannot split one backend across roots.
|
||||||
@@ -202,11 +205,18 @@ export class SessionPersistenceJsonl extends SessionPersistence implements Persi
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/** fsync a directory so a just-created or published entry inside it is crash-durable. */
|
/** fsync a directory when the host exposes that durability primitive. */
|
||||||
private async syncDir(dir: string): Promise<void> {
|
private async syncDir(dir: string): Promise<void> {
|
||||||
const handle = await open(dir, 'r')
|
const handle = await open(dir, 'r')
|
||||||
try {
|
try {
|
||||||
await handle.sync()
|
try {
|
||||||
|
await handle.sync()
|
||||||
|
} catch (error: unknown) {
|
||||||
|
const code = (error as NodeJS.ErrnoException | null)?.code
|
||||||
|
// Node opens directories on Windows but its fsync binding rejects them.
|
||||||
|
// File-content fsync remains mandatory; only this unsupported primitive is skipped.
|
||||||
|
if (this.internals.platform !== 'win32' || code !== 'EPERM') throw error
|
||||||
|
}
|
||||||
} finally {
|
} finally {
|
||||||
await handle.close()
|
await handle.close()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||||
import { Context } from 'cordis'
|
import { Context } from 'cordis'
|
||||||
import { appendFile, mkdtemp, mkdir, rm, readFile, writeFile, readdir, stat } from 'node:fs/promises'
|
import { appendFile, mkdtemp, mkdir, open, rm, readFile, writeFile, readdir, stat } from 'node:fs/promises'
|
||||||
|
import type { FileHandle } from 'node:fs/promises'
|
||||||
import { tmpdir } from 'node:os'
|
import { tmpdir } from 'node:os'
|
||||||
import { join } from 'node:path'
|
import { join } from 'node:path'
|
||||||
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
|
||||||
@@ -40,9 +41,25 @@ async function freshRoot(): Promise<string> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
afterEach(async () => {
|
afterEach(async () => {
|
||||||
|
vi.restoreAllMocks()
|
||||||
for (const d of dirs.splice(0)) await rm(d, { recursive: true, force: true })
|
for (const d of dirs.splice(0)) await rm(d, { recursive: true, force: true })
|
||||||
})
|
})
|
||||||
|
|
||||||
|
async function rejectDirectorySync(code: string): Promise<void> {
|
||||||
|
const handle = await open(root, 'r')
|
||||||
|
const proto = Object.getPrototypeOf(handle) as { sync: () => Promise<void> }
|
||||||
|
await handle.close()
|
||||||
|
const realSync = proto.sync
|
||||||
|
vi.spyOn(proto, 'sync').mockImplementation(async function (this: FileHandle) {
|
||||||
|
if ((await this.stat()).isDirectory()) {
|
||||||
|
const error = new Error(`simulated directory fsync ${code}`) as NodeJS.ErrnoException
|
||||||
|
error.code = code
|
||||||
|
throw error
|
||||||
|
}
|
||||||
|
return realSync.call(this)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
function appendClosedTurn(session: Session): void {
|
function appendClosedTurn(session: Session): void {
|
||||||
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
|
||||||
session.append('user/message', {
|
session.append('user/message', {
|
||||||
@@ -264,6 +281,28 @@ describe('SessionPersistenceJsonl: durability and crash semantics', () => {
|
|||||||
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
|
expect(loaded.events.map(e => e.seq)).toEqual([0, 1, 2, 3, 4, 5, 6, 7])
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('keeps file fsync mandatory while tolerating unsupported Windows directory fsync', async () => {
|
||||||
|
await rejectDirectorySync('EPERM')
|
||||||
|
const backend = ctx.sessionPersistence as SessionPersistenceJsonl
|
||||||
|
backend.internals.platform = 'win32'
|
||||||
|
const m = meta('windows-directory-sync')
|
||||||
|
await ctx.sessionPersistence.create(m)
|
||||||
|
await expect(ctx.sessionPersistence.append(m.id, oneTurnLog())).resolves.toBeUndefined()
|
||||||
|
expect((await ctx.sessionPersistence.load(m.id)).events).toEqual(oneTurnLog())
|
||||||
|
})
|
||||||
|
|
||||||
|
it.each([
|
||||||
|
['linux', 'EPERM'],
|
||||||
|
['win32', 'EIO'],
|
||||||
|
] as const)('surfaces directory fsync errors on %s with %s', async (platform, code) => {
|
||||||
|
await rejectDirectorySync(code)
|
||||||
|
const backend = ctx.sessionPersistence as SessionPersistenceJsonl
|
||||||
|
backend.internals.platform = platform
|
||||||
|
const m = meta(`directory-sync-${platform}-${code}`)
|
||||||
|
await ctx.sessionPersistence.create(m)
|
||||||
|
await expect(ctx.sessionPersistence.append(m.id, oneTurnLog())).rejects.toMatchObject({ code })
|
||||||
|
})
|
||||||
|
|
||||||
it('load returns a meta copy: mutating it does not corrupt backend pathing', async () => {
|
it('load returns a meta copy: mutating it does not corrupt backend pathing', async () => {
|
||||||
const m = meta('meta-copy', '/proj')
|
const m = meta('meta-copy', '/proj')
|
||||||
await ctx.sessionPersistence.create(m)
|
await ctx.sessionPersistence.create(m)
|
||||||
|
|||||||
Reference in New Issue
Block a user