| 1 | import { EntryGroup, EntryTree, isJsExpr, type EntryOptions } from '@deepseek-ai/cordis-plugin-loader' |
| 2 | import { Context, Service } from '@deepseek-ai/cordis' |
| 3 | import { extname } from 'node:path' |
| 4 | import { access, constants, readFile, rename, writeFile } from 'node:fs/promises' |
| 5 | import { setTimeout as delay } from 'node:timers/promises' |
| 6 | import { fileURLToPath, pathToFileURL } from 'node:url' |
| 7 | import * as yaml from 'js-yaml' |
| 8 | |
| 9 | const JsExpr = new yaml.Type('tag:yaml.org,2002:js', { |
| 10 | kind: 'scalar', |
| 11 | resolve: (data) => typeof data === 'string', |
| 12 | construct: (data) => ({ __jsExpr: data }), |
| 13 | predicate: isJsExpr, |
| 14 | represent: (data) => (data as { __jsExpr: string }).__jsExpr, |
| 15 | }) |
| 16 | |
| 17 | /** |
| 18 | * The entry-list YAML dialect: `!!js` scalars round-trip as expression nodes |
| 19 | * the Loader evaluates at entry activation. Exported so config tooling |
| 20 | * (`dsh --dump-config`) parses and prints exactly the dialect this include |
| 21 | * mounts. |
| 22 | */ |
| 23 | export const entryListSchema = yaml.JSON_SCHEMA.extend(JsExpr) |
| 24 | |
| 25 | const schema = entryListSchema |
| 26 | |
| 27 | const writable: Record<string, string> = { |
| 28 | '.json': 'application/json', |
| 29 | '.yaml': 'application/yaml', |
| 30 | '.yml': 'application/yaml', |
| 31 | } |
| 32 | |
| 33 | const supported = new Set(Object.keys(writable)) |
| 34 | |
| 35 | const WRITE_RETRY_LIMIT = 10 |
| 36 | const WRITE_RETRY_DELAY_MS = 50 |
| 37 | |
| 38 | function retryableWriteError(error: unknown): boolean { |
| 39 | const code = (error as NodeJS.ErrnoException | null)?.code |
| 40 | return code === 'EACCES' || code === 'EBUSY' || code === 'EPERM' |
| 41 | } |
| 42 | |
| 43 | /** |
| 44 | * Apply patch lists to an entry list — THE patch semantics of this include, |
| 45 | * shared by mounting (`applyPatches`) and offline config tooling |
| 46 | * (`dsh --dump-config`) so a dump can never drift from what boots. The input |
| 47 | * is never mutated: patching shared entry objects would bake earlier patch |
| 48 | * values into the cached parse, so repeated application (config hot-reloads) |
| 49 | * could never revert a removed or changed patch. Inserted entries are indexed |
| 50 | * as they are added, so a later patch in the same list can target a row an |
| 51 | * earlier patch inserted. A patch that matches nothing warns and is skipped. |
| 52 | * @param data - the parsed entry list (JSON-safe plain data). |
| 53 | * @param patches - the patch list to apply, in order. |
| 54 | * @param warn - sink for skipped-patch diagnostics (printf-style, `%C` = code). |
| 55 | * @returns a detached entry list with every applicable patch applied. |
| 56 | */ |
| 57 | export function applyEntryPatches( |
| 58 | data: EntryOptions[], |
| 59 | patches: PatchOptions[] | undefined, |
| 60 | warn: (message: string, ...args: any[]) => void, |
| 61 | ): EntryOptions[] { |
| 62 | if (!patches?.length) return [...data] |
| 63 | data = structuredClone(data) |
| 64 | |
| 65 | const entryMap = new Map<string, EntryOptions>() |
| 66 | const buildMap = (entries: EntryOptions[]) => { |
| 67 | for (const entry of entries) { |
| 68 | if (entry.id) entryMap.set(entry.id, entry) |
| 69 | if (entry.group && Array.isArray(entry.config)) { |
| 70 | buildMap(entry.config) |
| 71 | } |
| 72 | } |
| 73 | } |
| 74 | buildMap(data) |
| 75 | |
| 76 | for (const patch of patches) { |
| 77 | const { id, insert, name, ...overrides } = patch |
| 78 | |
| 79 | if (insert) { |
| 80 | if (id) { |
| 81 | const target = entryMap.get(id) |
| 82 | if (!target) { |
| 83 | warn('patch insert: entry %C not found', id) |
| 84 | continue |
| 85 | } |
| 86 | if (!target.group) { |
| 87 | warn('patch insert: entry %C is not a group', id) |
| 88 | continue |
| 89 | } |
| 90 | if (!Array.isArray(target.config)) target.config = [] |
| 91 | target.config.push(...insert) |
| 92 | } else { |
| 93 | data.push(...insert) |
| 94 | } |
| 95 | // Index what this patch added so a LATER patch in the same list can |
| 96 | // target it. Patch lists compose one layer per source (each bundle |
| 97 | // layer, then the user's, then `--patch` overlays), and a layer must be |
| 98 | // able to configure or disable a row an earlier layer inserted; without |
| 99 | // this, inserted rows were silently unpatchable. |
| 100 | buildMap(insert) |
| 101 | continue |
| 102 | } |
| 103 | |
| 104 | if (!id) { |
| 105 | warn('patch: id is required for non-insert patches') |
| 106 | continue |
| 107 | } |
| 108 | |
| 109 | const target = entryMap.get(id) |
| 110 | if (!target) { |
| 111 | warn('patch: entry %C not found', id) |
| 112 | continue |
| 113 | } |
| 114 | |
| 115 | if (name && name !== target.name) { |
| 116 | warn('patch: name mismatch for %C (expected %C, got %C), skipping', id, target.name, name) |
| 117 | continue |
| 118 | } |
| 119 | |
| 120 | for (const [key, value] of Object.entries(overrides)) { |
| 121 | if (key === 'id') continue |
| 122 | (target as unknown as Record<string, unknown>)[key] = value |
| 123 | } |
| 124 | } |
| 125 | |
| 126 | return data |
| 127 | } |
| 128 | |
| 129 | /** Runtime patch applied to entries loaded from an included config file. */ |
| 130 | export interface PatchOptions { |
| 131 | id?: string |
| 132 | insert?: EntryOptions[] |
| 133 | name?: string |
| 134 | config?: any |
| 135 | group?: boolean | null |
| 136 | disabled?: boolean | null |
| 137 | inject?: any |
| 138 | intercept?: any |
| 139 | isolate?: any |
| 140 | [key: string]: any |
| 141 | } |
| 142 | |
| 143 | /** Config namespace for the file-backed include loader. */ |
| 144 | export namespace Include { |
| 145 | /** Config for a file-backed loader subtree. */ |
| 146 | export interface Config { |
| 147 | /** YAML or JSON path resolved from `ctx.baseUrl`. */ |
| 148 | path: string |
| 149 | /** Entry list written when the file does not already exist. */ |
| 150 | initial?: any[] |
| 151 | /** Runtime patches applied after reading the file. */ |
| 152 | patches?: PatchOptions[] |
| 153 | /** Enables loader apply/reload/unload logs for this subtree. */ |
| 154 | enableLogs?: boolean |
| 155 | } |
| 156 | } |
| 157 | |
| 158 | /** Loader entry tree backed by a YAML or JSON file. */ |
| 159 | export class Include extends EntryTree { |
| 160 | static inject = ['loader'] |
| 161 | |
| 162 | // Tree-carrier marker (the Group plugin declares the same): this config is |
| 163 | // entry and patch lists, so the Loader's `internal/config` interpolation |
| 164 | // keeps it literal — a `!!js` expression inside a nested row's config |
| 165 | // belongs to that row's fiber, resolving lazily in the row's own context. |
| 166 | // Include's own fields (`path`, `enableLogs`) therefore stay literal too. |
| 167 | static readonly [EntryGroup.key] = true |
| 168 | |
| 169 | public filename: string |
| 170 | private type?: string |
| 171 | private readonly: boolean |
| 172 | private content?: string |
| 173 | private data?: EntryOptions[] |
| 174 | private writeTask?: NodeJS.Timeout | undefined |
| 175 | private pendingWrite?: EntryOptions[] |
| 176 | private writeQueue: Promise<void> = Promise.resolve() |
| 177 | |
| 178 | constructor(ctx: Context, public config: Include.Config) { |
| 179 | super(ctx) |
| 180 | this.enableLogs = config.enableLogs ?? ctx.fiber.entry?.parent.tree.enableLogs ?? false |
| 181 | this.filename = fileURLToPath(new URL(this.config.path, this.ctx.baseUrl)) |
| 182 | const ext = extname(this.filename) |
| 183 | if (!supported.has(ext)) { |
| 184 | throw new Error(`extension "${ext}" not supported`) |
| 185 | } |
| 186 | this.type = writable[ext] |
| 187 | this.readonly = !this.type |
| 188 | this.ctx.baseUrl = new URL('.', pathToFileURL(this.filename)).href |
| 189 | |
| 190 | ctx.on('internal/update', (config, _, next) => { |
| 191 | if (config.path !== this.config.path) return next() |
| 192 | // Veto the fiber restart (children update in place), but persist the new |
| 193 | // config ourselves — `Fiber.update` only assigns `this.config` behind |
| 194 | // `next()`, and a stale `this.config.patches` would make the next |
| 195 | // `refresh()` re-apply the old overlay. |
| 196 | this.config = config |
| 197 | this.root.update(this.applyPatches(this.data!, config.patches)).catch((error) => { |
| 198 | this.ctx.logger.warn('config update at %C failed', this.filename) |
| 199 | this.ctx.logger.warn(error) |
| 200 | }) |
| 201 | }) |
| 202 | } |
| 203 | |
| 204 | private async checkAccess() { |
| 205 | if (!this.type) return |
| 206 | try { |
| 207 | await access(this.filename, constants.W_OK) |
| 208 | } catch { |
| 209 | this.readonly = true |
| 210 | } |
| 211 | } |
| 212 | |
| 213 | private async read(forced = false) { |
| 214 | const content = await readFile(this.filename, 'utf8') |
| 215 | if (!forced && this.content === content) return false |
| 216 | let data: any |
| 217 | if (this.type === 'application/yaml') { |
| 218 | data = yaml.load(content, { schema: entryListSchema }) |
| 219 | } else if (this.type === 'application/json') { |
| 220 | data = JSON.parse(content) |
| 221 | } else { |
| 222 | const module = await import(/* @vite-ignore */ this.filename) |
| 223 | data = module.default || module |
| 224 | } |
| 225 | // An empty or truncated file (common mid-edit: editors and `sed -i` write |
| 226 | // through temp states) parses to `undefined`, not an error; reject every |
| 227 | // non-array shape here so callers see one "invalid file" signal. Content |
| 228 | // and data commit only on success, so an edit that is later reverted to |
| 229 | // the exact last good content correctly reads as "unchanged". |
| 230 | if (!Array.isArray(data)) { |
| 231 | throw new TypeError(`config file must be a top-level array of entries: ${this.filename}`) |
| 232 | } |
| 233 | this.content = content |
| 234 | this.data = data |
| 235 | await this.checkAccess() |
| 236 | return true |
| 237 | } |
| 238 | |
| 239 | private applyPatches(data: EntryOptions[], patches = this.config.patches): EntryOptions[] { |
| 240 | return applyEntryPatches(data, patches, (message, ...args) => { |
| 241 | this.ctx.root.logger?.('loader').warn(message, ...args) |
| 242 | }) |
| 243 | } |
| 244 | |
| 245 | async* [Service.init]() { |
| 246 | try { |
| 247 | await this.read() |
| 248 | } catch (error) { |
| 249 | // Only a missing file falls back to `initial` (or the not-found error): |
| 250 | // an existing-but-invalid file must fail loud with its real parse error, |
| 251 | // never be mislabelled as absent or silently overwritten. |
| 252 | if ((error as NodeJS.ErrnoException | null)?.code !== 'ENOENT') throw error |
| 253 | if (this.config.initial) { |
| 254 | await this._writeFile(this.config.initial as any) |
| 255 | await this.read(true) |
| 256 | } else { |
| 257 | throw new Error(`config file not found: ${this.filename}`) |
| 258 | } |
| 259 | } |
| 260 | |
| 261 | yield () => this.stop() |
| 262 | await this.root.update(this.applyPatches(this.data!)) |
| 263 | } |
| 264 | |
| 265 | async stop() { |
| 266 | try { |
| 267 | await this.flushWrite() |
| 268 | } finally { |
| 269 | this.root.stop() |
| 270 | await this.flushWrite() |
| 271 | } |
| 272 | } |
| 273 | |
| 274 | /** |
| 275 | * Re-read the file and refresh child entries when content changed. An |
| 276 | * unreadable or unparsable file logs a warning and keeps the last good |
| 277 | * tree: a hot-reload of a live app must never take the process down. |
| 278 | */ |
| 279 | async refresh() { |
| 280 | try { |
| 281 | if (!await this.read()) return |
| 282 | await this.root.update(this.applyPatches(this.data!)) |
| 283 | } catch (error) { |
| 284 | this.ctx.logger.warn('config reload at %C failed; keeping the running tree', this.filename) |
| 285 | this.ctx.logger.warn(error) |
| 286 | } |
| 287 | } |
| 288 | |
| 289 | private async _writeFile(config: EntryOptions[]) { |
| 290 | if (this.readonly) { |
| 291 | throw new Error(`cannot overwrite readonly config`) |
| 292 | } |
| 293 | if (this.type === 'application/yaml') { |
| 294 | this.content = yaml.dump(config, { schema }) |
| 295 | } else if (this.type === 'application/json') { |
| 296 | this.content = JSON.stringify(config, null, 2) |
| 297 | } |
| 298 | await writeFile(this.filename + '.tmp', this.content!) |
| 299 | for (let retry = 0; ; retry++) { |
| 300 | try { |
| 301 | await rename(this.filename + '.tmp', this.filename) |
| 302 | return |
| 303 | } catch (error) { |
| 304 | if (!retryableWriteError(error) || retry >= WRITE_RETRY_LIMIT) throw error |
| 305 | await delay((retry + 1) * WRITE_RETRY_DELAY_MS) |
| 306 | } |
| 307 | } |
| 308 | } |
| 309 | |
| 310 | private writeFile(config: EntryOptions[]) { |
| 311 | clearTimeout(this.writeTask) |
| 312 | this.pendingWrite = config |
| 313 | this.writeTask = setTimeout(() => { |
| 314 | void this.flushWrite() |
| 315 | }, 0) |
| 316 | } |
| 317 | |
| 318 | private flushWrite(): Promise<void> { |
| 319 | clearTimeout(this.writeTask) |
| 320 | this.writeTask = undefined |
| 321 | const config = this.pendingWrite |
| 322 | this.pendingWrite = undefined |
| 323 | if (config === undefined) return this.writeQueue |
| 324 | const run = this.writeQueue.then( |
| 325 | () => this._writeFile(config), |
| 326 | () => this._writeFile(config), |
| 327 | ) |
| 328 | this.writeQueue = run |
| 329 | void run.catch((error) => { |
| 330 | this.ctx.root.logger?.('loader').warn('failed to write config file %C', this.filename) |
| 331 | this.ctx.root.logger?.('loader').warn(error) |
| 332 | }) |
| 333 | return run |
| 334 | } |
| 335 | |
| 336 | /** Schedule a write of the current root entry data. */ |
| 337 | write() { |
| 338 | this.context.emit('loader/config-update') |
| 339 | return this.writeFile(this.root.data) |
| 340 | } |
| 341 | } |
| 342 | |
| 343 | export default Include |
| 344 |