persistCounts() was write-through on every count change with no serialization, and the atomic writer named its temp file with millisecond granularity. Two persists inside one millisecond shared the temp path: both wrote it, the first rename consumed it, the second rename found nothing — ENOENT, roughly 1,500 times a day on a busy production brain, with a full ledger write per change behind it. No data was lost (the surviving rename carried a complete ledger and the next change re-persisted), but the race was real and the write rate absurd. flushCounts() now runs exactly one persist at a time; requests arriving during it collapse into one trailing pass that carries the burst's final state — N changes cost at most two writes. writeFileAtomic() adds a per-process sequence to the temp name so no two writes can share a path. Pinned: a 25-change burst → ≤2 ledger writes, zero errors, ledger equal to memory; parallel real writes land complete; three same-instant atomic writes own three distinct temp paths.
111 lines
4.6 KiB
TypeScript
111 lines
4.6 KiB
TypeScript
/**
|
|
* @module tests/integration/counts-persist-single-flight
|
|
* @description Regression for a production race in FileSystemStorage's
|
|
* counts ledger: `persistCounts()` was write-through on every count change
|
|
* with no serialization, and the atomic writer named its temp file with
|
|
* millisecond granularity (`.tmp-<pid>-<ms>`). Two persists inside one
|
|
* millisecond shared the temp path — both wrote it, the first rename
|
|
* consumed it, the second rename found nothing: ENOENT, ~1,500 times a day
|
|
* on a busy production brain, with a full ledger write per change behind it.
|
|
*
|
|
* Under pin: persists are single-flight and coalesced — one in flight, at
|
|
* most one trailing pass carrying the burst's final state — and every atomic
|
|
* write owns a unique temp path. A burst of N count changes costs at most
|
|
* two ledger writes, never errors, and leaves a ledger equal to memory.
|
|
*/
|
|
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest'
|
|
import * as fs from 'node:fs'
|
|
import * as os from 'node:os'
|
|
import * as path from 'node:path'
|
|
import { Brainy } from '../../src/brainy.js'
|
|
import { NounType } from '../../src/types/graphTypes.js'
|
|
|
|
describe('counts persistence is single-flight, coalesced, and never races its own temp file', () => {
|
|
let dir: string
|
|
let brain: any
|
|
|
|
beforeEach(async () => {
|
|
process.env.BRAINY_DETERMINISTIC_EMBEDDINGS = 'true'
|
|
dir = fs.mkdtempSync(path.join(os.tmpdir(), 'brainy-counts-race-'))
|
|
brain = new Brainy({
|
|
requireSubtype: false,
|
|
storage: { type: 'filesystem', path: dir },
|
|
dimensions: 384,
|
|
silent: true
|
|
})
|
|
await brain.init()
|
|
})
|
|
|
|
afterEach(async () => {
|
|
vi.restoreAllMocks()
|
|
await brain.close()
|
|
fs.rmSync(dir, { recursive: true, force: true })
|
|
})
|
|
|
|
it('a burst of concurrent count changes → at most two ledger writes, zero errors, ledger == memory', async () => {
|
|
const storage = brain.storage
|
|
const countsPath: string = storage.countsFilePath
|
|
expect(countsPath, 'the filesystem adapter persists a counts ledger').toBeTruthy()
|
|
|
|
// Let init's own persists settle so the burst is measured alone.
|
|
await storage.flushCounts?.()
|
|
|
|
const renameSpy = vi.spyOn(fs.promises, 'rename')
|
|
const errorSpy = vi.spyOn(console, 'error')
|
|
|
|
// Twenty-five concurrent count changes — the shape of a write burst; each
|
|
// used to launch its own persist.
|
|
const BURST = 25
|
|
await Promise.all(
|
|
Array.from({ length: BURST }, () => storage.scheduleCountPersist())
|
|
)
|
|
|
|
const ledgerRenames = renameSpy.mock.calls.filter(([, to]) => String(to) === countsPath)
|
|
expect(ledgerRenames.length, 'single-flight + one trailing pass').toBeLessThanOrEqual(2)
|
|
expect(ledgerRenames.length, 'the burst was persisted at all').toBeGreaterThanOrEqual(1)
|
|
|
|
const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts'))
|
|
expect(persistErrors).toEqual([])
|
|
|
|
const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8'))
|
|
expect(ledger.totalNounCount).toBe(storage.totalNounCount)
|
|
expect(ledger.totalVerbCount).toBe(storage.totalVerbCount)
|
|
})
|
|
|
|
it('real writes in parallel: the ledger lands complete and no persist error is logged', async () => {
|
|
const storage = brain.storage
|
|
const countsPath: string = storage.countsFilePath
|
|
const errorSpy = vi.spyOn(console, 'error')
|
|
|
|
await Promise.all(
|
|
Array.from({ length: 12 }, (_, i) =>
|
|
brain.add({ data: `burst row ${i}`, type: NounType.Thing })
|
|
)
|
|
)
|
|
await storage.flushCounts?.()
|
|
|
|
const persistErrors = errorSpy.mock.calls.filter((args) => String(args[0]).includes('persisting counts'))
|
|
expect(persistErrors).toEqual([])
|
|
const ledger = JSON.parse(fs.readFileSync(countsPath, 'utf-8'))
|
|
expect(ledger.totalNounCount).toBe(storage.totalNounCount)
|
|
expect(await brain.getNounCount()).toBe(ledger.totalNounCount)
|
|
})
|
|
|
|
it('every atomic write owns its own temp path — two writes in one millisecond never collide', async () => {
|
|
const storage = brain.storage
|
|
const tmpNames: string[] = []
|
|
vi.spyOn(fs.promises, 'writeFile').mockImplementation(async (p: any) => {
|
|
tmpNames.push(String(p))
|
|
})
|
|
vi.spyOn(fs.promises, 'rename').mockImplementation(async () => undefined)
|
|
const target = path.join(dir, 'probe.json')
|
|
await Promise.all([
|
|
storage.writeFileAtomic(target, '{"a":1}'),
|
|
storage.writeFileAtomic(target, '{"a":2}'),
|
|
storage.writeFileAtomic(target, '{"a":3}')
|
|
])
|
|
const probeTmps = tmpNames.filter((n) => n.startsWith(`${target}.tmp-`))
|
|
expect(probeTmps.length).toBe(3)
|
|
expect(new Set(probeTmps).size, 'no two writes shared a temp path').toBe(3)
|
|
})
|
|
})
|