| 1 | import test from 'node:test' |
| 2 | import assert from 'node:assert/strict' |
| 3 | import { mkdtemp, mkdir, readFile, writeFile, readdir, stat, symlink, rm } from 'node:fs/promises' |
| 4 | import { createHash } from 'node:crypto' |
| 5 | import { spawn } from 'node:child_process' |
| 6 | import { once } from 'node:events' |
| 7 | import { tmpdir } from 'node:os' |
| 8 | import { join } from 'node:path' |
| 9 | import { createStorage, STORAGE_LIMITS } from '../src/shims/storage.ts' |
| 10 | |
| 11 | const moduleUrl = new URL('../src/shims/storage.ts', import.meta.url).href |
| 12 | const keyPath = (dir, key) => join(dir, 'storage-v1', `${createHash('sha256').update(key).digest('hex')}.json`) |
| 13 | const code = (expected) => (error) => error.code === expected |
| 14 | |
| 15 | async function fixture(t) { |
| 16 | const root = await mkdtemp(join(tmpdir(), 'codewhale-plugin-storage-')) |
| 17 | t.after(() => rm(root, { recursive: true, force: true })) |
| 18 | const a = join(root, 'a'), b = join(root, 'b') |
| 19 | await Promise.all([mkdir(a, { mode: 0o700 }), mkdir(b, { mode: 0o700 })]) |
| 20 | let live = true |
| 21 | return { a, b, store: createStorage({ dataDir: a, isActive: () => live }), revoke: () => { live = false } } |
| 22 | } |
| 23 | |
| 24 | function worker(t, directory, script) { |
| 25 | const child = spawn(process.execPath, ['-e', script, directory], { stdio: ['ignore', 'ignore', 'pipe', 'ipc'] }) |
| 26 | let stderr = '' |
| 27 | child.stderr.on('data', (chunk) => { stderr = (stderr + chunk).slice(-4000) }) |
| 28 | const exited = once(child, 'exit') |
| 29 | t.after(async () => { if (child.exitCode === null && child.signalCode === null) child.kill('SIGKILL'); await exited }) |
| 30 | return { child, exited, stderr: () => stderr } |
| 31 | } |
| 32 | |
| 33 | test('storage persists JSON across owner generations and deletes exact keys', async (t) => { |
| 34 | const { a, store } = await fixture(t) |
| 35 | assert.ok(Object.isFrozen(store)) |
| 36 | assert.equal(await store.get('absent'), undefined) |
| 37 | await store.set('profile', { enabled: true, list: [null, 3, '你好'] }) |
| 38 | const restarted = createStorage({ dataDir: a, isActive: () => true }) |
| 39 | assert.deepEqual(await restarted.get('profile'), { enabled: true, list: [null, 3, '你好'] }) |
| 40 | const copy = await restarted.get('profile'); copy.enabled = false |
| 41 | assert.equal((await store.get('profile')).enabled, true) |
| 42 | assert.equal(await restarted.delete('profile'), true) |
| 43 | assert.equal(await restarted.delete('profile'), false) |
| 44 | assert.equal(await store.get('profile'), undefined) |
| 45 | await store.set('permissions', true) |
| 46 | if (process.platform !== 'win32') assert.equal((await stat(keyPath(a, 'permissions'))).mode & 0o777, 0o600) |
| 47 | }) |
| 48 | |
| 49 | test('owner directories isolate equal keys and keys never become paths', async (t) => { |
| 50 | const { a, b, store } = await fixture(t) |
| 51 | const other = createStorage({ dataDir: b, isActive: () => true }) |
| 52 | await store.set('../outside', 'a') |
| 53 | await other.set('../outside', 'b') |
| 54 | await store.set('__proto__', { owner: 'a' }) |
| 55 | assert.equal(await store.get('../outside'), 'a') |
| 56 | assert.equal(await other.get('../outside'), 'b') |
| 57 | assert.deepEqual(await store.get('__proto__'), { owner: 'a' }) |
| 58 | assert.equal(await other.get('__proto__'), undefined) |
| 59 | assert.ok((await readdir(join(a, 'storage-v1'))).every((name) => /^[a-f0-9]{64}\.json$/u.test(name))) |
| 60 | }) |
| 61 | |
| 62 | test('storage refuses invalid JSON, oversized values and keys without damaging state', async (t) => { |
| 63 | const { a, store } = await fixture(t) |
| 64 | await store.set('kept', 'original') |
| 65 | const before = await readFile(keyPath(a, 'kept')) |
| 66 | const sparse = new Array(2), cyclic = {}; cyclic.self = cyclic |
| 67 | const symbol = { [Symbol('hidden')]: 1 }, extra = [1]; extra.custom = 2 |
| 68 | const getter = Object.defineProperty({}, 'x', { enumerable: true, get() { throw new Error('must not run') } }) |
| 69 | for (const value of [undefined, NaN, Infinity, 1n, new Date(), () => 1, sparse, cyclic, symbol, getter, extra]) await assert.rejects(store.set('bad', value), code('invalid')) |
| 70 | for (const key of ['', 'nul\0key', '鲸'.repeat(STORAGE_LIMITS.keyBytes)]) await assert.rejects(store.set(key, 1), code('invalid')) |
| 71 | await assert.rejects(store.set('big', 'x'.repeat(STORAGE_LIMITS.valueBytes)), code('limit')) |
| 72 | assert.deepEqual(await readFile(keyPath(a, 'kept')), before) |
| 73 | }) |
| 74 | |
| 75 | test('observed owner quota refuses the next write and allows replacement or deletion', async (t) => { |
| 76 | const { store } = await fixture(t) |
| 77 | const value = 'x'.repeat(STORAGE_LIMITS.valueBytes - 2) |
| 78 | let count = 0 |
| 79 | for (;;) { |
| 80 | try { await store.set(`key${count}`, value); count++ } catch (error) { assert.equal(error.code, 'limit'); break } |
| 81 | } |
| 82 | assert.ok(count > 0) |
| 83 | assert.equal(await store.get('key0'), value) |
| 84 | assert.equal(await store.get(`key${count}`), undefined) |
| 85 | await store.set('key0', 'smaller') |
| 86 | await store.set(`key${count}`, value) |
| 87 | await store.delete('key1') |
| 88 | }) |
| 89 | |
| 90 | test('observed key-count quota is shared between API instances in one host', async (t) => { |
| 91 | const { a, store } = await fixture(t) |
| 92 | await store.set('kept', true) |
| 93 | await Promise.all(Array.from({ length: STORAGE_LIMITS.keys - 1 }, (_, i) => writeFile(keyPath(a, `k${i}`), JSON.stringify({ version: 1, key: `k${i}`, value: null }), { mode: 0o600 }))) |
| 94 | const other = createStorage({ dataDir: a, isActive: () => true }) |
| 95 | await assert.rejects(other.set('extra', 1), code('limit')) |
| 96 | await other.set('kept', false) |
| 97 | await store.delete('k1') |
| 98 | await other.set('extra', 2) |
| 99 | assert.equal(await store.get('extra'), 2) |
| 100 | }) |
| 101 | |
| 102 | test('queued operations refuse a revoked owner and state survives revocation', async (t) => { |
| 103 | const { a, store, revoke } = await fixture(t) |
| 104 | await store.set('kept', 7) |
| 105 | const pending = store.set('late', 8) |
| 106 | revoke() |
| 107 | await assert.rejects(pending, code('not_available')) |
| 108 | await assert.rejects(store.get('kept'), code('not_available')) |
| 109 | await assert.rejects(store.delete('kept'), code('not_available')) |
| 110 | const restarted = createStorage({ dataDir: a, isActive: () => true }) |
| 111 | assert.equal(await restarted.get('kept'), 7) |
| 112 | assert.equal(await restarted.get('late'), undefined) |
| 113 | await restarted.set('after-reload', 9) |
| 114 | assert.equal(await restarted.get('after-reload'), 9) |
| 115 | }) |
| 116 | |
| 117 | test('corrupt records fail closed, preserve old bytes, and support explicit key deletion', async (t) => { |
| 118 | const { a, store } = await fixture(t) |
| 119 | await store.set('key', true) |
| 120 | for (const bytes of ['{broken', '{"version":2,"key":"key","value":null}', '{"version":1,"key":"other","value":null}']) { |
| 121 | await writeFile(keyPath(a, 'key'), bytes) |
| 122 | await assert.rejects(store.get('key'), code('corrupt')) |
| 123 | await assert.rejects(store.set('key', 1), code('corrupt')) |
| 124 | assert.equal(await readFile(keyPath(a, 'key'), 'utf8'), bytes) |
| 125 | } |
| 126 | assert.equal(await store.delete('key'), true) |
| 127 | await store.set('key', 'recovered') |
| 128 | assert.equal(await store.get('key'), 'recovered') |
| 129 | }) |
| 130 | |
| 131 | test('storage refuses a linked record instead of reading or changing its target', { skip: process.platform === 'win32' }, async (t) => { |
| 132 | const { a, b, store } = await fixture(t) |
| 133 | await store.get('absent') |
| 134 | const outside = join(b, 'untouched.json'), bytes = '{"version":1,"key":"secret","value":"untouched"}' |
| 135 | await writeFile(outside, bytes) |
| 136 | await symlink(outside, keyPath(a, 'secret')) |
| 137 | await assert.rejects(store.get('secret'), code('corrupt')) |
| 138 | await assert.rejects(store.set('secret', 'changed'), code('corrupt')) |
| 139 | assert.equal(await readFile(outside, 'utf8'), bytes) |
| 140 | }) |
| 141 | |
| 142 | test('concurrent calls across API instances share the directory queue', async (t) => { |
| 143 | const { a, store } = await fixture(t) |
| 144 | const other = createStorage({ dataDir: a, isActive: () => true }) |
| 145 | await Promise.all(Array.from({ length: 24 }, (_, i) => (i % 2 ? store : other).set(`k${i}`, i))) |
| 146 | assert.deepEqual(await Promise.all(Array.from({ length: 24 }, (_, i) => store.get(`k${i}`))), Array.from({ length: 24 }, (_, i) => i)) |
| 147 | }) |
| 148 | |
| 149 | test('set snapshots input at invocation before callers can mutate it', async (t) => { |
| 150 | const { store } = await fixture(t) |
| 151 | const value = { counter: 1 }, pending = store.set('value', value) |
| 152 | value.counter = 99 |
| 153 | await pending |
| 154 | assert.deepEqual(await store.get('value'), { counter: 1 }) |
| 155 | }) |
| 156 | |
| 157 | test('storage requires an assigned absolute directory and refuses directory links', { skip: process.platform === 'win32' }, async (t) => { |
| 158 | const { a, b } = await fixture(t) |
| 159 | assert.throws(() => createStorage({ dataDir: 'relative', isActive: () => true }), code('invalid')) |
| 160 | const linked = join(a, 'linked') |
| 161 | await symlink(b, linked) |
| 162 | await assert.rejects(createStorage({ dataDir: linked, isActive: () => true }).set('key', 1), code('invalid')) |
| 163 | }) |
| 164 | |
| 165 | test('killing a writer before publication preserves old state and reload stays writable', async (t) => { |
| 166 | const { a, store } = await fixture(t) |
| 167 | await store.set('kept', 'original') |
| 168 | const run = worker(t, a, ` |
| 169 | const fs = require('node:fs'); |
| 170 | fs.promises.rename = async () => { setInterval(() => {}, 1000); process.send('publishing'); await new Promise(() => {}); }; |
| 171 | require('node:module').syncBuiltinESMExports(); |
| 172 | import(${JSON.stringify(moduleUrl)}).then(({ createStorage }) => createStorage({ dataDir: process.argv[1], isActive: () => true }).set('kept', 'interrupted')).catch(e => { console.error(e); process.exit(1); }); |
| 173 | `) |
| 174 | const ready = await Promise.race([once(run.child, 'message', { signal: AbortSignal.timeout(5000) }), run.exited.then(() => { throw new Error(run.stderr()) })]) |
| 175 | assert.equal(ready[0], 'publishing') |
| 176 | run.child.kill('SIGKILL') |
| 177 | const [, signal] = await run.exited |
| 178 | assert.equal(signal, 'SIGKILL') |
| 179 | const names = await readdir(join(a, 'storage-v1')) |
| 180 | assert.ok(names.some((name) => name.startsWith('.pending-'))) |
| 181 | const reloaded = createStorage({ dataDir: a, isActive: () => true }) |
| 182 | assert.equal(await reloaded.get('kept'), 'original') |
| 183 | await reloaded.set('kept', 'after-crash') |
| 184 | await reloaded.set('new', true) |
| 185 | assert.equal(await reloaded.get('kept'), 'after-crash') |
| 186 | assert.equal(await reloaded.get('new'), true) |
| 187 | }) |
| 188 | |
| 189 | test('different-key writers in separate processes cannot erase each other', async (t) => { |
| 190 | const { a, store } = await fixture(t) |
| 191 | const makeScript = (prefix) => `import(${JSON.stringify(moduleUrl)}).then(async ({ createStorage }) => { |
| 192 | const store = createStorage({ dataDir: process.argv[1], isActive: () => true }); |
| 193 | for (let i = 0; i < 12; i++) await store.set(${JSON.stringify(prefix)} + i, i); |
| 194 | process.disconnect(); |
| 195 | }).catch(e => { console.error(e); process.exit(1); });` |
| 196 | const first = worker(t, a, makeScript('first')), second = worker(t, a, makeScript('second')) |
| 197 | for (const run of [first, second]) assert.equal((await run.exited)[0], 0, run.stderr()) |
| 198 | for (let i = 0; i < 12; i++) { |
| 199 | assert.equal(await store.get(`first${i}`), i) |
| 200 | assert.equal(await store.get(`second${i}`), i) |
| 201 | } |
| 202 | }) |
| 203 | |
| 204 | test('same-key writers publish complete last-writer-wins records during concurrent reads', async (t) => { |
| 205 | const { a, store } = await fixture(t) |
| 206 | await store.set('shared', { writer: 'seed', payload: 'seed' }) |
| 207 | const makeScript = (writer) => `import(${JSON.stringify(moduleUrl)}).then(async ({ createStorage }) => { |
| 208 | const store = createStorage({ dataDir: process.argv[1], isActive: () => true }); |
| 209 | for (let i = 0; i < 12; i++) await store.set('shared', { writer: ${JSON.stringify(writer)}, payload: ${JSON.stringify(writer)}.repeat(10000) }); |
| 210 | process.disconnect(); |
| 211 | }).catch(e => { console.error(e); process.exit(1); });` |
| 212 | const first = worker(t, a, makeScript('first')), second = worker(t, a, makeScript('second')) |
| 213 | while (first.child.exitCode === null || second.child.exitCode === null) { |
| 214 | const value = await store.get('shared') |
| 215 | assert.ok(['seed', 'first', 'second'].includes(value.writer)) |
| 216 | assert.equal(value.payload, value.writer === 'seed' ? 'seed' : value.writer.repeat(10000)) |
| 217 | } |
| 218 | for (const run of [first, second]) assert.equal((await run.exited)[0], 0, run.stderr()) |
| 219 | assert.ok(['first', 'second'].includes((await store.get('shared')).writer)) |
| 220 | }) |
| 221 | |
| 222 | // Faults are injected before the first storage import in an owned subprocess. |
| 223 | // Bun's syncBuiltinESMExports is a no-op, so installing hooks after import would |
| 224 | // leave named filesystem bindings untouched. Seed records through the parent |
| 225 | // store first; every child operation then observes the actual injected failure. |
| 226 | // These prove retry/guard behavior, not Windows kernel semantics; the unchanged |
| 227 | // concurrent-writer case runs there. |
| 228 | async function sharingWorker(t, directory, body) { |
| 229 | const run = worker(t, directory, ` |
| 230 | (async () => { |
| 231 | const fs = require('node:fs'), { syncBuiltinESMExports } = require('node:module'); |
| 232 | ${body} |
| 233 | process.disconnect(); |
| 234 | })().catch(e => { console.error(e); process.exit(1); }); |
| 235 | `) |
| 236 | const message = once(run.child, 'message', { signal: AbortSignal.timeout(5000) }) |
| 237 | const [report] = await Promise.race([message, run.exited.then(() => { throw new Error(run.stderr()) })]) |
| 238 | assert.equal((await run.exited)[0], 0, run.stderr()) |
| 239 | return report |
| 240 | } |
| 241 | |
| 242 | test('Windows sharing retries keep complete records and retain atomic replacement', async (t) => { |
| 243 | const { a, store } = await fixture(t) |
| 244 | await store.set('shared', 'original') |
| 245 | const report = await sharingWorker(t, a, ` |
| 246 | const realOpen = fs.promises.open, realRename = fs.promises.rename; |
| 247 | let reads = 0, publications = 0; |
| 248 | fs.promises.open = async (...args) => { |
| 249 | if (String(args[0]).endsWith('.json') && reads++ < 2) throw Object.assign(new Error('sharing'), { code: 'EPERM' }); |
| 250 | return realOpen(...args); |
| 251 | }; |
| 252 | fs.promises.rename = async (...args) => { |
| 253 | if (publications++ < 2) throw Object.assign(new Error('sharing'), { code: 'EACCES' }); |
| 254 | return realRename(...args); |
| 255 | }; |
| 256 | syncBuiltinESMExports(); |
| 257 | const { createStorage } = await import(${JSON.stringify(moduleUrl)}); |
| 258 | const store = createStorage({ dataDir: process.argv[1], isActive: () => true }); |
| 259 | Object.defineProperty(process, 'platform', { value: 'win32' }); |
| 260 | await store.set('shared', 'complete replacement'); |
| 261 | const value = await store.get('shared'); |
| 262 | process.send({ value, publications, reads, names: await fs.promises.readdir(require('node:path').join(process.argv[1], 'storage-v1')) }); |
| 263 | `) |
| 264 | assert.equal(report.value, 'complete replacement') |
| 265 | assert.equal(report.publications, 3) |
| 266 | assert.ok(report.reads >= 3) |
| 267 | assert.equal(report.names.length, 1) |
| 268 | assert.match(report.names[0], /^[a-f0-9]{64}\.json$/u) |
| 269 | }) |
| 270 | |
| 271 | test('Windows sharing retries refuse revocation, corrupt replacements and permanent denial', async (t) => { |
| 272 | const { a } = await fixture(t) |
| 273 | for (const mode of ['revoked', 'corrupt', 'denied']) { |
| 274 | const directory = join(a, mode) |
| 275 | await mkdir(directory, { mode: 0o700 }) |
| 276 | await createStorage({ dataDir: directory, isActive: () => true }).set('shared', 'original') |
| 277 | const report = await sharingWorker(t, directory, ` |
| 278 | let live = true; |
| 279 | const path = require('node:path').join(process.argv[1], 'storage-v1', require('node:crypto').createHash('sha256').update('shared').digest('hex') + '.json'); |
| 280 | const original = await fs.promises.readFile(path, 'utf8'); |
| 281 | let publications = 0; |
| 282 | fs.promises.rename = async () => { |
| 283 | publications++; |
| 284 | if (${JSON.stringify(mode)} === 'revoked') live = false; |
| 285 | if (${JSON.stringify(mode)} === 'corrupt') await fs.promises.writeFile(path, '{broken'); |
| 286 | throw Object.assign(new Error('sharing'), { code: 'EPERM' }); |
| 287 | }; |
| 288 | syncBuiltinESMExports(); |
| 289 | const { createStorage } = await import(${JSON.stringify(moduleUrl)}); |
| 290 | const store = createStorage({ dataDir: process.argv[1], isActive: () => live }); |
| 291 | Object.defineProperty(process, 'platform', { value: 'win32' }); |
| 292 | let code; |
| 293 | try { await store.set('shared', 'must not publish'); throw new Error('unexpected publication'); } catch (error) { code = error.code; } |
| 294 | process.send({ code, publications, bytes: await fs.promises.readFile(path, 'utf8'), original, names: await fs.promises.readdir(require('node:path').dirname(path)) }); |
| 295 | `) |
| 296 | assert.equal(report.code, mode === 'revoked' ? 'not_available' : mode === 'corrupt' ? 'corrupt' : 'io') |
| 297 | assert.equal(report.publications, mode === 'denied' ? 11 : 1) |
| 298 | assert.equal(report.bytes, mode === 'corrupt' ? '{broken' : report.original) |
| 299 | assert.equal(report.names.length, 1, 'only this invocation temporary file is cleaned up') |
| 300 | } |
| 301 | }) |
| 302 | |
| 303 | test('sharing failures on other platforms remain visible without replaying publication', async (t) => { |
| 304 | const { a, store } = await fixture(t) |
| 305 | await store.set('shared', 'original') |
| 306 | const report = await sharingWorker(t, a, ` |
| 307 | let publications = 0; |
| 308 | fs.promises.rename = async () => { publications++; throw Object.assign(new Error('denied'), { code: 'EPERM' }); }; |
| 309 | syncBuiltinESMExports(); |
| 310 | const { createStorage } = await import(${JSON.stringify(moduleUrl)}); |
| 311 | const store = createStorage({ dataDir: process.argv[1], isActive: () => true }); |
| 312 | Object.defineProperty(process, 'platform', { value: 'linux' }); |
| 313 | let code; |
| 314 | try { await store.set('shared', 'must not publish'); throw new Error('unexpected publication'); } catch (error) { code = error.code; } |
| 315 | process.send({ code, publications, value: await store.get('shared') }); |
| 316 | `) |
| 317 | assert.equal(report.code, 'io') |
| 318 | assert.equal(report.publications, 1) |
| 319 | assert.equal(report.value, 'original') |
| 320 | }) |
| 321 | |
| 322 | test('Windows sharing read retries refuse revocation before reopening private state', async (t) => { |
| 323 | const { a, store } = await fixture(t) |
| 324 | await store.set('shared', 'private value') |
| 325 | const report = await sharingWorker(t, a, ` |
| 326 | let live = true; |
| 327 | const realOpen = fs.promises.open; |
| 328 | let opens = 0; |
| 329 | fs.promises.open = async (...args) => { |
| 330 | if (String(args[0]).endsWith('.json')) { |
| 331 | opens++; |
| 332 | live = false; |
| 333 | throw Object.assign(new Error('sharing'), { code: 'EPERM' }); |
| 334 | } |
| 335 | return realOpen(...args); |
| 336 | }; |
| 337 | syncBuiltinESMExports(); |
| 338 | const { createStorage } = await import(${JSON.stringify(moduleUrl)}); |
| 339 | const store = createStorage({ dataDir: process.argv[1], isActive: () => live }); |
| 340 | Object.defineProperty(process, 'platform', { value: 'win32' }); |
| 341 | let code, value; |
| 342 | try { value = await store.get('shared'); } catch (error) { code = error.code; } |
| 343 | process.send({ code, opens, value }); |
| 344 | `) |
| 345 | assert.equal(report.code, 'not_available') |
| 346 | assert.equal(report.opens, 1) |
| 347 | assert.equal(Object.hasOwn(report, 'value'), false, 'revoked state is never returned over IPC') |
| 348 | }) |
| 349 |