| 1 | /** Owner-local plugin state, never session history or a credential service. */ |
| 2 | import { constants } from 'node:fs' |
| 3 | import { lstat, mkdir, open, readdir, realpath, rename, unlink } from 'node:fs/promises' |
| 4 | import { isAbsolute, join } from 'node:path' |
| 5 | import { createHash, randomUUID } from 'node:crypto' |
| 6 | import { setTimeout as delay } from 'node:timers/promises' |
| 7 | import { isJson } from '../json.ts' |
| 8 | import type { Json } from '../protocol.ts' |
| 9 | |
| 10 | export const STORAGE_LIMITS = Object.freeze({ keyBytes: 128, valueBytes: 128 * 1024, totalBytes: 4 * 1024 * 1024, keys: 1024 }) |
| 11 | const RECORD_NAME = /^[a-f0-9]{64}\.json$/u |
| 12 | const MAX_RECORD_BYTES = STORAGE_LIMITS.valueBytes + STORAGE_LIMITS.keyBytes * 6 + 128 |
| 13 | const queues = new Map<string, Promise<void>>() |
| 14 | |
| 15 | export interface PluginStorage { |
| 16 | get(key: string): Promise<Json | undefined> |
| 17 | set(key: string, value: Json): Promise<void> |
| 18 | delete(key: string): Promise<boolean> |
| 19 | } |
| 20 | |
| 21 | export interface StorageOptions { |
| 22 | /** The exact directory Rust already assigned to this owner. */ |
| 23 | dataDir: string |
| 24 | /** Must become false as soon as this owner starts disposal or is revoked. */ |
| 25 | isActive: () => boolean |
| 26 | /** Optional host diagnostic when directory fsync is unavailable after publication. */ |
| 27 | onWarning?: (message: string) => void |
| 28 | } |
| 29 | |
| 30 | export class StorageError extends Error { |
| 31 | readonly code: 'not_available' | 'invalid' | 'limit' | 'corrupt' | 'io' |
| 32 | constructor(code: StorageError['code'], message: string) { |
| 33 | super(message) |
| 34 | this.code = code |
| 35 | this.name = 'StorageError' |
| 36 | } |
| 37 | } |
| 38 | |
| 39 | function checkKey(key: string): void { |
| 40 | if (typeof key !== 'string' || !key.length || key.includes('\0') || Buffer.byteLength(key) > STORAGE_LIMITS.keyBytes) { |
| 41 | throw new StorageError('invalid', `storage key must contain 1 to ${STORAGE_LIMITS.keyBytes} UTF-8 bytes and no NUL`) |
| 42 | } |
| 43 | } |
| 44 | |
| 45 | /** Reject accessors, symbols, sparse arrays and extra properties before isJson reads values. */ |
| 46 | function snapshot(value: unknown): Json { |
| 47 | if (!isJson(value)) throw new StorageError('invalid', 'storage value must be plain JSON') |
| 48 | const encoded = JSON.stringify(value) |
| 49 | if (Buffer.byteLength(encoded) > STORAGE_LIMITS.valueBytes) throw new StorageError('limit', 'storage value exceeds its byte limit') |
| 50 | return JSON.parse(encoded) |
| 51 | } |
| 52 | |
| 53 | function recordName(key: string): string { |
| 54 | return `${createHash('sha256').update(key).digest('hex')}.json` |
| 55 | } |
| 56 | |
| 57 | function fsCode(error: unknown): string | undefined { |
| 58 | return (error as NodeJS.ErrnoException)?.code |
| 59 | } |
| 60 | |
| 61 | // Windows may briefly deny access while another process replaces a record. |
| 62 | // Retry only sharing-related I/O errors; permanent denial still fails closed. |
| 63 | async function retryWindowsSharing<T>(operation: () => Promise<T>, beforeRetry?: () => unknown | Promise<unknown>): Promise<T> { |
| 64 | for (let retry = 0; ; retry++) { |
| 65 | if (retry) await beforeRetry?.() |
| 66 | try { return await operation() } catch (error) { |
| 67 | if (process.platform !== 'win32' || retry >= 10 || !['EACCES', 'EBUSY', 'EPERM'].includes(fsCode(error) ?? '')) throw error |
| 68 | await delay(50) |
| 69 | } |
| 70 | } |
| 71 | } |
| 72 | |
| 73 | interface RecordSnapshot { key: string; value: Json; bytes: number } |
| 74 | |
| 75 | /** Read one bounded, complete record without following a link. */ |
| 76 | async function readRecordOnce(directory: string, name: string): Promise<RecordSnapshot | undefined> { |
| 77 | const path = join(directory, name) |
| 78 | let file |
| 79 | try { |
| 80 | const before = await lstat(path) |
| 81 | if (!before.isFile() || before.isSymbolicLink() || before.nlink > 1) throw new StorageError('corrupt', 'plugin storage record is not a single-link regular file') |
| 82 | file = await open(path, constants.O_RDONLY | (constants.O_NOFOLLOW ?? 0)) |
| 83 | const stat = await file.stat() |
| 84 | // Another host may atomically replace this key between lstat and open, |
| 85 | // or unlink the old inode after we opened it. Both complete versions are |
| 86 | // valid snapshots; retain the no-follow/single-link checks on the version |
| 87 | // actually opened rather than calling an ordinary replacement corruption. |
| 88 | const current = await lstat(path).catch((error) => { if (fsCode(error) !== 'ENOENT') throw error; return undefined }) |
| 89 | if (!stat.isFile() || stat.nlink > 1 || current?.isSymbolicLink()) throw new StorageError('corrupt', 'plugin storage record changed while opening') |
| 90 | if (stat.size > MAX_RECORD_BYTES) throw new StorageError('corrupt', 'plugin storage record exceeds its byte limit') |
| 91 | const bytes = Buffer.alloc(stat.size + 1) |
| 92 | let length = 0 |
| 93 | while (length < bytes.length) { |
| 94 | const read = await file.read(bytes, length, bytes.length - length, length) |
| 95 | if (!read.bytesRead) break |
| 96 | length += read.bytesRead |
| 97 | } |
| 98 | if (length > stat.size) throw new StorageError('corrupt', 'plugin storage record changed while reading') |
| 99 | const stored = JSON.parse(new TextDecoder('utf-8', { fatal: true }).decode(bytes.subarray(0, length))) |
| 100 | if (!stored || stored.version !== 1 || Object.keys(stored).length !== 3 || !Object.hasOwn(stored, 'value')) throw new StorageError('corrupt', 'plugin storage record format is invalid') |
| 101 | checkKey(stored.key) |
| 102 | if (recordName(stored.key) !== name) throw new StorageError('corrupt', 'plugin storage record does not match its key') |
| 103 | return { key: stored.key, value: snapshot(stored.value), bytes: length } |
| 104 | } catch (error) { |
| 105 | if (fsCode(error) === 'ENOENT') return undefined |
| 106 | if (error instanceof StorageError) throw new StorageError('corrupt', error.message) |
| 107 | if (error instanceof SyntaxError || error instanceof TypeError) throw new StorageError('corrupt', 'plugin storage record is invalid; existing state was preserved') |
| 108 | throw error |
| 109 | } finally { |
| 110 | await file?.close() |
| 111 | } |
| 112 | } |
| 113 | |
| 114 | /** |
| 115 | * One frozen API per owner generation; isActive is checked on every admission |
| 116 | * and again before publication. Unload never deletes state. A write admitted |
| 117 | * at that final check may finish its atomic rename during disposal. |
| 118 | * |
| 119 | * Completed records are atomic per key. Same-key writes from different host |
| 120 | * processes are last-writer-wins; different keys cannot overwrite one another. |
| 121 | * A shared in-process queue checks the observed 4 MiB / 1024-record quota before |
| 122 | * each set. This is not a strict shared disk cap: independent-process races, |
| 123 | * pending files and crash-orphaned temporary files are outside that accounting. |
| 124 | * Key/value limits remain strict. No lock survives a crash, no orphan is guessed |
| 125 | * away, and reads/deletes remain usable when the observed quota is exceeded. |
| 126 | * This API does not contain arbitrary co-resident native plugin code. |
| 127 | */ |
| 128 | export function createStorage({ dataDir, isActive, onWarning }: StorageOptions): PluginStorage { |
| 129 | if (typeof dataDir !== 'string' || !isAbsolute(dataDir)) throw new StorageError('invalid', 'plugin storage needs its assigned absolute dataDir') |
| 130 | |
| 131 | function active(): void { |
| 132 | if (!isActive()) throw new StorageError('not_available', 'plugin storage owner is no longer active') |
| 133 | } |
| 134 | |
| 135 | async function readRecord(directory: string, name: string): Promise<RecordSnapshot | undefined> { |
| 136 | return retryWindowsSharing(() => readRecordOnce(directory, name), active) |
| 137 | } |
| 138 | |
| 139 | async function directory(): Promise<string> { |
| 140 | const stat = await lstat(dataDir) |
| 141 | if (!stat.isDirectory() || stat.isSymbolicLink()) throw new StorageError('invalid', 'plugin storage dataDir must be a real directory') |
| 142 | return join(await realpath(dataDir), 'storage-v1') |
| 143 | } |
| 144 | |
| 145 | async function prepare(directory: string): Promise<void> { |
| 146 | await mkdir(directory, { mode: 0o700 }).catch((error) => { if (fsCode(error) !== 'EEXIST') throw error }) |
| 147 | const stat = await lstat(directory) |
| 148 | if (!stat.isDirectory() || stat.isSymbolicLink()) throw new StorageError('corrupt', 'plugin storage record directory must be a real directory') |
| 149 | active() |
| 150 | } |
| 151 | |
| 152 | async function syncDirectory(directory: string): Promise<void> { |
| 153 | // File fsync is required before publication. Directory fsync is OS-best- |
| 154 | // effort after publication; reporting failure as a rejected set would lie |
| 155 | // about the complete record that already replaced the previous version. |
| 156 | let dir |
| 157 | try { dir = await open(directory, constants.O_RDONLY); await dir.sync() } catch (error) { |
| 158 | try { onWarning?.(`plugin storage record committed; directory fsync unavailable (${fsCode(error) ?? 'unknown filesystem error'})`) } catch { /* Diagnostics cannot undo publication. */ } |
| 159 | } finally { |
| 160 | await dir?.close().catch(() => undefined) |
| 161 | } |
| 162 | } |
| 163 | |
| 164 | async function checkQuota(directory: string, name: string, candidateBytes: number): Promise<void> { |
| 165 | const names = (await readdir(directory)).filter((entry) => RECORD_NAME.test(entry)) |
| 166 | let count = 0, bytes = 0 |
| 167 | for (const entry of names) { |
| 168 | if (entry === name) continue |
| 169 | const record = await readRecord(directory, entry) |
| 170 | if (!record) continue // A separate process may have deleted this key. |
| 171 | count++ |
| 172 | bytes += record.bytes |
| 173 | if (count + 1 > STORAGE_LIMITS.keys || bytes + candidateBytes > STORAGE_LIMITS.totalBytes) throw new StorageError('limit', 'plugin storage exceeds its observed owner quota') |
| 174 | } |
| 175 | if (count + 1 > STORAGE_LIMITS.keys || bytes + candidateBytes > STORAGE_LIMITS.totalBytes) throw new StorageError('limit', 'plugin storage exceeds its observed owner quota') |
| 176 | } |
| 177 | |
| 178 | async function write(directory: string, key: string, value: Json): Promise<void> { |
| 179 | const name = recordName(key) |
| 180 | const encoded = JSON.stringify({ version: 1, key, value }) + '\n' |
| 181 | // Refuse silent replacement of a corrupt or linked record. |
| 182 | await readRecord(directory, name) |
| 183 | await checkQuota(directory, name, Buffer.byteLength(encoded)) |
| 184 | const temporary = join(directory, `.pending-${randomUUID()}.tmp`) |
| 185 | let file |
| 186 | let published = false |
| 187 | try { |
| 188 | file = await open(temporary, 'wx', 0o600) |
| 189 | await file.writeFile(encoded) |
| 190 | await file.sync() |
| 191 | await file.close() |
| 192 | file = undefined |
| 193 | await retryWindowsSharing(async () => { |
| 194 | active() |
| 195 | await rename(temporary, join(directory, name)) |
| 196 | }, async () => { |
| 197 | // A retry never bypasses a newly corrupt/linked destination or revocation. |
| 198 | active() |
| 199 | await readRecord(directory, name) |
| 200 | }) |
| 201 | published = true |
| 202 | await syncDirectory(directory) |
| 203 | } finally { |
| 204 | await file?.close() |
| 205 | // This invocation's unpredictable file only; never remove crash orphans. |
| 206 | if (!published) await unlink(temporary).catch((error) => { if (fsCode(error) !== 'ENOENT') throw error }) |
| 207 | } |
| 208 | } |
| 209 | |
| 210 | async function queue<T>(operation: (directory: string) => Promise<T>): Promise<T> { |
| 211 | active() |
| 212 | try { |
| 213 | const dir = await directory() |
| 214 | const next = (queues.get(dir) ?? Promise.resolve()).then(async () => { active(); await prepare(dir); return operation(dir) }) |
| 215 | const tail = next.then(() => undefined, () => undefined) |
| 216 | queues.set(dir, tail) |
| 217 | void tail.then(() => { if (queues.get(dir) === tail) queues.delete(dir) }) |
| 218 | return await next |
| 219 | } catch (error) { |
| 220 | if (error instanceof StorageError) throw error |
| 221 | throw new StorageError('io', `plugin storage operation failed (${fsCode(error) ?? 'unknown filesystem error'})`) |
| 222 | } |
| 223 | } |
| 224 | |
| 225 | return Object.freeze({ |
| 226 | get(key: string) { |
| 227 | return queue(async (dir) => { checkKey(key); const record = await readRecord(dir, recordName(key)); active(); return record?.value }) |
| 228 | }, |
| 229 | set(key: string, value: Json) { |
| 230 | let copy: Json |
| 231 | try { active(); checkKey(key); copy = snapshot(value) } catch (error) { return Promise.reject(error) } |
| 232 | return queue((dir) => write(dir, key, copy)) |
| 233 | }, |
| 234 | delete(key: string) { |
| 235 | return queue(async (dir) => { |
| 236 | checkKey(key) |
| 237 | const path = join(dir, recordName(key)) |
| 238 | const stat = await lstat(path).catch((error) => { if (fsCode(error) !== 'ENOENT') throw error; return undefined }) |
| 239 | if (!stat) return false |
| 240 | if (!stat.isFile() || stat.isSymbolicLink() || stat.nlink !== 1) throw new StorageError('corrupt', 'plugin storage record is not a single-link regular file') |
| 241 | active() |
| 242 | await unlink(path) |
| 243 | await syncDirectory(dir) |
| 244 | return true |
| 245 | }) |
| 246 | }, |
| 247 | }) |
| 248 | } |
| 249 |