The recovery fold scanned the generation log from generation 1 at every open, on the open's foreground — O(whole history) on long-lived brains (measured at two minutes of a large brain's open). Now an advisory mark records the log's head whenever the pending set drains to empty (and at clean close when empty); recovery scans from the mark + 1. The mark is advisory and monotone-safe: stale-low costs a longer scan, never a marker. The fold itself moves behind the doors as a latched background task — the embed worker starts when it settles, and awaitPendingEmbeds() and close() wait on the latch first, so no caller can observe a half-recovered set. A pending embed's outcome was always eventual; moving its recovery off the foreground changes when the worker starts, never whether a marker is honored. Pinned in tests/integration/pending-embed-low-water.test.ts: the drain writes the mark and the next open scans from mark + 1; a pending embed enqueued after the mark survives an unclean stop; open arms the fold as a background latch the barrier waits on; a clean close writes the mark even without a drain.
145 lines
5.7 KiB
TypeScript
145 lines
5.7 KiB
TypeScript
/**
|
|
* @module tests/integration/pending-embed-low-water
|
|
* @description The pending-embed recovery fold is bounded and background (10.4.9).
|
|
*
|
|
* The fold used to scan the generation log from generation 1 at EVERY open,
|
|
* on the open's foreground — O(whole history) per open on long-lived brains.
|
|
* Now: an advisory low-water mark (`_system/pending_embeds_lowwater.json`)
|
|
* records the committed generation whenever the pending set drains to empty,
|
|
* recovery scans from `mark + 1`, and the fold runs behind the doors as a
|
|
* latched background task the worker, `awaitPendingEmbeds()` and `close()`
|
|
* wait on. The mark is advisory: stale-low costs a longer scan, never a
|
|
* marker — a pending embed enqueued before a crash is still recovered.
|
|
*/
|
|
import { describe, it, expect, afterEach, vi } from 'vitest'
|
|
import { mkdtempSync, rmSync } from 'node:fs'
|
|
import { tmpdir } from 'node:os'
|
|
import { join } from 'node:path'
|
|
import { Brainy } from '../../src/brainy'
|
|
import { NounType } from '../../src/types/graphTypes'
|
|
|
|
const LOWWATER_PATH = '_system/pending_embeds_lowwater.json'
|
|
|
|
describe('pending-embed recovery: bounded by the low-water mark, behind the doors', () => {
|
|
const roots: string[] = []
|
|
const dir = (): string => {
|
|
const d = mkdtempSync(join(tmpdir(), 'brainy-lowwater-'))
|
|
roots.push(d)
|
|
return d
|
|
}
|
|
const open = async (root: string): Promise<Brainy<any>> => {
|
|
const brain = new Brainy<any>({
|
|
requireSubtype: false,
|
|
storage: { type: 'filesystem', path: root }
|
|
})
|
|
await brain.init()
|
|
return brain
|
|
}
|
|
|
|
afterEach(() => {
|
|
for (const d of roots.splice(0)) rmSync(d, { recursive: true, force: true })
|
|
})
|
|
|
|
it('drain-to-empty writes the mark, and the next open scans from mark + 1', async () => {
|
|
const root = dir()
|
|
const brain = await open(root)
|
|
// Hold the worker so the pending state is observable, then release it.
|
|
const realKick = (brain as any).kickEmbedWorker.bind(brain)
|
|
;(brain as any).kickEmbedWorker = () => {}
|
|
await brain.add({
|
|
id: 'row-1',
|
|
data: 'the first deferred row',
|
|
type: NounType.Thing,
|
|
deferEmbedding: true
|
|
})
|
|
expect(brain.pendingEmbedCount()).toBeGreaterThan(0)
|
|
;(brain as any).kickEmbedWorker = realKick
|
|
await brain.awaitPendingEmbeds()
|
|
// The drain wrote the advisory mark (fire-and-forget: settle the microtask).
|
|
await new Promise((r) => setTimeout(r, 50))
|
|
const mark = (await (brain as any).storage.readRawObject(LOWWATER_PATH)) as {
|
|
generation: number
|
|
} | null
|
|
expect(mark).not.toBeNull()
|
|
expect(mark!.generation).toBeGreaterThan(0)
|
|
await brain.close()
|
|
|
|
const brain2 = await open(root)
|
|
const log = (brain2 as any).generationStore.getFactLog()
|
|
const scanSpy = vi.spyOn(log, 'scanFacts')
|
|
try {
|
|
await (brain2 as any).recoverPendingEmbedsFromLog()
|
|
expect(scanSpy).toHaveBeenCalledTimes(1)
|
|
const opts = scanSpy.mock.calls[0][0] as { fromGeneration?: number }
|
|
expect(opts.fromGeneration).toBeGreaterThanOrEqual(mark!.generation + 1)
|
|
} finally {
|
|
scanSpy.mockRestore()
|
|
await brain2.close()
|
|
}
|
|
})
|
|
|
|
it('a pending embed enqueued after the mark survives an unclean stop', async () => {
|
|
const root = dir()
|
|
const brain = await open(root)
|
|
await brain.add({ id: 'settled', data: 'lands before the mark', type: NounType.Thing })
|
|
await brain.awaitPendingEmbeds()
|
|
await new Promise((r) => setTimeout(r, 50))
|
|
|
|
// A deferred write whose embed never lands: block the worker, then drop
|
|
// the instance without close() — the unclean-stop shape.
|
|
;(brain as any).kickEmbedWorker = () => {}
|
|
await brain.add({
|
|
id: 'orphan',
|
|
data: 'enqueued then abandoned',
|
|
type: NounType.Thing,
|
|
deferEmbedding: true
|
|
})
|
|
expect(brain.pendingEmbedCount()).toBeGreaterThan(0)
|
|
// No close(): simulate the crash by releasing only the writer lock so the
|
|
// next open can proceed.
|
|
await (brain as any).storage.releaseWriterLock()
|
|
|
|
const brain2 = await open(root)
|
|
await (brain2 as any)._pendingEmbedRecovery
|
|
expect(brain2.pendingEmbedCount()).toBeGreaterThan(0)
|
|
await brain2.awaitPendingEmbeds()
|
|
expect(brain2.pendingEmbedCount()).toBe(0)
|
|
await brain2.close()
|
|
// Reap the crashed instance: its fence is gone, so close() fails loudly —
|
|
// swallow that here; the point is clearing its watchers and registry entry.
|
|
await brain.close().catch(() => undefined)
|
|
})
|
|
|
|
it('open arms the fold as a background latch; awaitPendingEmbeds waits on it', async () => {
|
|
const root = dir()
|
|
const brain = await open(root)
|
|
await brain.add({ id: 'a-row', data: 'some data', type: NounType.Thing })
|
|
await brain.awaitPendingEmbeds()
|
|
await brain.close()
|
|
|
|
const brain2 = await open(root)
|
|
// The latch exists the moment init() returns (writable filesystem brain)…
|
|
expect((brain2 as any)._pendingEmbedRecovery).not.toBeNull()
|
|
// …and the barrier settles it before answering.
|
|
await brain2.awaitPendingEmbeds()
|
|
expect(brain2.pendingEmbedCount()).toBe(0)
|
|
await brain2.close()
|
|
})
|
|
|
|
it('a clean close with an empty set writes the mark even if no drain happened', async () => {
|
|
const root = dir()
|
|
const brain = await open(root)
|
|
await brain.add({ id: 'r1', data: 'row one', type: NounType.Thing })
|
|
await brain.awaitPendingEmbeds()
|
|
await brain.close()
|
|
// Read the mark back through the storage door (the adapter owns the
|
|
// on-disk encoding), on a fresh instance.
|
|
const brain2 = await open(root)
|
|
const mark = (await (brain2 as any).storage.readRawObject(LOWWATER_PATH)) as {
|
|
generation: number
|
|
} | null
|
|
expect(mark).not.toBeNull()
|
|
expect(mark!.generation).toBeGreaterThan(0)
|
|
await brain2.close()
|
|
})
|
|
})
|