From a21c49a75950698682c9e7cb4bc3a6858a3d0ecb Mon Sep 17 00:00:00 2001 From: Yuto Nishida Date: Sun, 27 Sep 2026 13:01:43 -0700 Subject: [PATCH] dsh-plugins fix-workspace-multiple-writers: add the workspace store reconciler MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two dsh instances sharing one `$DSH_HOME` each keep the whole workspace store in memory: the JSON backend reads the document once at open, the domain loads its tables once, and the registry holds its own record snapshots. A whole-file publish from one instance therefore silently drops additions the other made, and no process ever notices. This plugin watches `storages/workspace.json` and replays the *additions* it describes through `ctx.workspaceRegistry`'s public API — the same write path a session create uses — so the attach emits `domain/changed` and every connected browser shows the session without a reload. Additive-only by construction: the planner's result type has no removal, reorder, or rename action at all. `updatedAt` cannot arbitrate a conflict (a stale whole-file write is indistinguishable from a deliberate removal), so any destructive reconciliation would be a guess, and the reported bug is lost additions. Decisions worth knowing: - The store path is derived from the open domain's own unit rather than assumed from `$DSH_HOME/storages/workspace.json`. A relocated backend root would otherwise leave the plugin watching a stale file and replaying it into the live registry; a directory-backed (`per-record`) unit makes it decline to run rather than read a directory as a document. - Local membership comes from the durable records, because `Workspace.sessionIds` is cwd-index-filtered: planning an attach from that view makes the registry's prune step write an accounted id out. When the durable accounts are unreadable the planner attaches only to workspaces it is about to create, since a fresh record cannot be pruned. - Passes are debounced onto a single-flight chain, and the store's containing directory is watched rather than the file, because the backend commits by temp-write + `rename()`. A watch that cannot be created degrades to `fs.watchFile`, which polls the path itself. - Retries are bounded. A `pendingMutation` left by a writer that died between its two writes would otherwise poll every 250 ms forever while disabling reconciliation. Limits, documented in the plugin README: it never removes, reorders, or renames; the registry prunes locally-unvalidatable ids on *any* write to that record, including one this plugin induces; and every registry write republishes the whole file from this process's one-time snapshot, so a peer's concurrent addition can be reverted inside the read-to-commit window — and the peer will not repair it, since its plan is store-minus-local and its own addition is still in its memory. Closing that window needs cross-process serialization of the read-modify-write, which is out of scope here. Zero runtime dependencies: a plugin mounted through the loader's `cordis:include` resolves its imports from its own directory, so only `node:` built-ins are imported. `HANDOFF.md` is the design record this implementation follows. --- .../HANDOFF.md | 492 ++++++++++++++++++ .../README.md | 200 +++++++ .../lib/index.js | 420 +++++++++++++++ .../lib/plan.js | 62 +++ .../lib/store.js | 89 ++++ .../package.json | 22 + 6 files changed, 1285 insertions(+) create mode 100644 dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/HANDOFF.md create mode 100644 dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/README.md create mode 100644 dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/index.js create mode 100644 dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/plan.js create mode 100644 dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/store.js create mode 100644 dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/package.json diff --git a/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/HANDOFF.md b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/HANDOFF.md new file mode 100644 index 0000000..8e4ca4f --- /dev/null +++ b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/HANDOFF.md @@ -0,0 +1,492 @@ +# Handoff — `dsh-plugin-fix-workspace-multiple-writers` (Route 1) + +**Status:** ready to implement +**Target:** a host-only dsh plugin, `dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/` +**Verified against:** `@deepseek-ai/dsh` 0.1.5-rc.1 (Nix store checkout `pwxk15h0f6mpjjmd0fnhh6j0g8138v58-dsh-0.1.5-rc.1`), workspace packages `0.1.5-rc.2` + +--- + +## 1. The one-sentence task + +Watch `$DSH_HOME/storages/workspace.json` and, when another dsh instance adds sessions to a workspace, push those additions into this process's `ctx.workspaceRegistry` through its **public API**, so both instances converge and both browser GUIs show the same sessions. + +## 2. Hard rule for v1: additive only + +**Never remove, never reorder, never rename, never delete.** v1 only *adds*: + +- create a locally-unknown workspace that the file describes, +- attach a session id the file lists that we don't have, +- archive a session id the file has archived. + +Rationale is in §7.2 — there is no way to distinguish "the other instance deliberately removed X" from "the other instance's stale whole-file write lost X", so any destructive reconciliation is a guess. The reported bug is lost *additions*, so additive-only fixes it with zero risk of data loss. Removal/reorder reconciliation is explicitly a non-goal here (see §10). + +## 3. Why this belongs at the registry, not at the file or backend + +Do not try to fix this by watching/reloading the storage backend. Three independent in-memory caches sit between the file and anything the user sees: + +| Layer | Where | Behavior | +|---|---|---| +| `SingleJsonUnit.state` | `dsh-storage-json` (`single-unit.js`) | Reads the document **once** in `openSingleUnit`; every write republishes the whole in-memory state. | +| `DomainImpl.tables` / `globalValue` | `dsh-storage-domain` (`lib/index.js`, `DomainFacility.open`) | Calls `unit.loadAll()` **once** at open; reads are synchronous from these maps; there is **no public invalidate/reload**. | +| `WorkspaceEntity.record` | `dsh-workspace` | A per-workspace snapshot; the registry does **not** subscribe to `domain/changed`. | + +So a backend-level reload changes nothing observable, and a lock around the file constrains nothing because `SingleJsonUnit.publish()` writes without taking one. + +What *does* exist is a clean, public, already-used write path: + +- `dsh-api-session-controller` calls `workspace.attachSession(sessionId)` when it creates a session (`lib/index.js:587`). +- `attachSession` validates the session's stored `cwd` against the workspace path and writes through `table.update`. +- That write emits in-process `domain/changed`, and `dsh-api-workspace-controller`'s `WorkspaceFeed` is subscribed to `domain/changed` and baselines from `registry.list()` (`dsh-api-workspace-controller/lib/index.js:44–55`). + +So an attach performed by our plugin updates the live GUI immediately, in-process, exactly like a normal session create. + +## 4. Deliverables + +``` +dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/ + package.json # type: module, main: lib/index.js, no client half + lib/index.js # the plugin + HANDOFF.md # this file (delete once implemented, or keep as design record) +``` + +Plus one row appended to `dsh-plugins/plugins.cordis.yml`: + +```yaml +- id: dsh-plugin-fix-workspace-multiple-writers + name: ./dsh-plugin-fix-workspace-multiple-writers/lib/index.js +``` + +No profile patch and no new backend are needed for Route 1. The include is already inside the HMR watch root, so editing the plugin reloads it live. + +> **Zero runtime dependencies — this is a hard implementation constraint, not a style preference.** +> A plugin mounted via the checkout's `cordis:include` resolves its own `import`s from its file's +> location (`/path/to/checkout/dsh-plugins/…/lib/`), not from `$DSH_HOME/profiles/node_modules`. +> Verified from that path in this very checkout: +> +> ``` +> @deepseek-ai/dsh-home-paths FAIL ERR_MODULE_NOT_FOUND +> @deepseek-ai/dsh-workspace FAIL ERR_MODULE_NOT_FOUND +> chokidar FAIL ERR_MODULE_NOT_FOUND +> ``` +> +> The loader anchors `baseUrl` at `dsh-plugins/` and then uses a plain dynamic `import()` for the +> entry file (`cordis-plugin-loader/lib/index.js:270–279`); Node's parent-walk from there never +> reaches the profile's module fallback, and the repo root has no `node_modules`. **Import only +> `node:*` built-ins.** Everything the plugin needs — watching, the home path, and structural +> validation — is spelled out below without any package import. + +`package.json` shape (mirror `dsh-client-ui-simple-statusbar/package.json`, minus the client bits): + +```json +{ + "name": "dsh-plugin-fix-workspace-multiple-writers", + "description": "Reconciles $DSH_HOME/storages/workspace.json into ctx.workspaceRegistry so concurrent dsh instances see each other's sessions.", + "version": "0.1.0", + "private": true, + "type": "module", + "main": "lib/index.js", + "exports": { + ".": "./lib/index.js", + "./package.json": "./package.json" + }, + "files": [ + "lib/index.js" + ] +} +``` + +Runtime dependencies: **none** (see the box above). No `dependencies` and no runtime `peerDependencies` — the plugin imports only `node:fs`, `node:fs/promises`, `node:os`, and `node:path`. + +## 5. Verified API reference + +### 5.1 `ctx.workspaceRegistry` (`dsh-workspace/lib/index.js`) + +```ts +class WorkspaceRegistry { + async create(path: string, title?: string): Promise + get(id: WorkspaceId): Workspace | undefined + list(): Workspace[] // durable order, synchronous, no persistence reads + async delete(id: WorkspaceId): Promise + async insertBefore(id: WorkspaceId, beforeId?: WorkspaceId): Promise + async resolveByPath(path: string): Promise + async sessionKnown(id: SessionId): Promise + get archivedSessionIds(): SessionId[] + async archiveSession(sessionId: SessionId): Promise +} + +interface Workspace { + readonly id: WorkspaceId + readonly path: string // fs.realpath canonical + readonly title: string + readonly createdAt: string // ISO-8601 + readonly updatedAt: string // ISO-8601 + readonly sessionIds: readonly SessionId[] // filtered by the local path index + setTitle(title: string): Promise + attachSession(sessionId: SessionId): Promise // validates cwd; idempotent + insertSessionBefore(sessionId: SessionId, beforeSessionId?: SessionId): Promise + detachSession(sessionId: SessionId): Promise + status(): Promise<'ok' | 'missing-dir'> +} +``` + +Facts that matter for the implementation: + +- `attachSession` is **idempotent and write-free when already attached**: if the id is already in `this.record.sessionIds` it skips the header read, and `mutate` throws the internal `unchangedSentinel` when the record is unchanged, so the medium is not rewritten and no `domain/changed` is emitted (`dsh-workspace/lib/index.js:111`, `:170`). +- `attachSession` reads the header from session persistence and **throws** when the header is missing, when `cwd` is absent, when `cwd` does not resolve, or when it resolves to a different path. Catch per-id. +- `create` canonicalizes via `realpath`, rejects relative/nonexistent/non-directory paths, and returns the existing entity unchanged if the canonical path is already owned. +- `archiveSession` requires the session to be known (live, header-indexed, or in a fresh persistence listing) and is a no-op when already archived. +- `entity.sessionIds` is **filtered** by the local session-path index (`dsh-workspace/lib/index.js:102`). Do not use it as exact local truth; either call the idempotent `attachSession`, or read the raw record via `ctx.storageDomain.get('workspace').table('workspaces').get(id)`. + +### 5.2 The file — `$DSH_HOME/storages/workspace.json` + +Resolve the path yourself, mirroring `resolveDshHome` (explicit configured > `$DSH_HOME` > `~/.dsh`; an empty or whitespace-only `$DSH_HOME` counts as unset; the result is `resolve()`d). Compute it **inside each pass**, not once at boot, so a late `DSH_HOME` is honored: + +```js +import { homedir } from 'node:os' +import { join, resolve } from 'node:path' + +const dshHome = () => { + const env = process.env.DSH_HOME + return resolve(env !== undefined && env.trim().length > 0 ? env : join(homedir(), '.dsh')) +} +const workspaceFile = () => join(dshHome(), 'storages', 'workspace.json') +``` + +```jsonc +{ + "unit": { "name": "workspace", "version": 2 }, + "global": { + "initialized": true, + "workspaceIds": ["", "..."], // durable display order + "archivedSessionIds": ["", "..."], + "pendingMutation": { "operation": "create", "workspaceId": "..." } // optional, transient + }, + "tables": { + "workspaces": { + "": { + "path": "/abs/canonical/dir", + "title": "some-title", + "sessionIds": ["session-…", "…"], // ordered ownership account + "createdAt": "2026-09-26T17:17:43.195Z", + "updatedAt": "2026-09-26T17:17:43.195Z" + } + } + } +} +``` + +Validate defensively before acting — structural checks only, no schema import: + +```js +// reject the whole pass (log, no-op) unless all of these hold: +// unit.name === 'workspace' +// unit.version === 2 +// global is an object with Array.isArray(workspaceIds) and Array.isArray(archivedSessionIds) +// tables is an object; tables.workspaces is an object +// per record, skip that record unless: +// typeof r.path === 'string' +// Array.isArray(r.sessionIds) && r.sessionIds.every(id => typeof id === 'string') +// typeof r.title === 'string' (fall back to basename(path) if absent) +// +// The registry re-validates with the real zod schemas at the durability boundary +// (dsh-storage-domain `parseRecord`), so a false accept here fails loudly and +// harmlessly on the registry side rather than corrupting anything. Never +// "repair" a malformed file by writing to it. +``` + +The exact schemas, if you ever need to mirror them by hand, are `workspaceDomainState` and `workspaceRecord` in `dsh-workspace/lib/types/spec.js`. + +The document is published by atomic temp-write + `rename()`, so a reader never sees a partial file. It may not exist at all (the backend materializes lazily); `ENOENT` means "nothing to reconcile". + +### 5.3 Live change notification + +`dsh-api-workspace-controller`'s `WorkspaceFeed` subscribes to `domain/changed` for `domain === 'workspace'` and re-reads `registry.list()` for its baseline. Any mutation we perform through the registry API therefore reaches connected browsers. No extra UI work is needed. + +## 6. Algorithm + +Split into a **pure planner** and an **impure executor** — this is the main testability decision. + +### 6.1 Read + guard + +``` +read file text + ENOENT -> no-op + unparseable / fails validation -> log error, no-op (never rewrite) + unit.name/version mismatch -> log error, no-op + global.initialized !== true -> no-op (let the registry bootstrap) + global.pendingMutation present -> no-op, retry after retryOnPendingMs +``` + +### 6.2 Plan (pure) + +```ts +planReconcile(file, local) -> { + createWorkspaces: Array<{ path: string, title: string }>, + attach: Array<{ workspacePath: string, sessionId: string }>, + archive: string[], +} +``` + +`local` is a plain snapshot the caller builds from `ctx.workspaceRegistry` (no I/O): + +```ts +local = { + byPath: Map }>, + archived: Set, +} +// built from registry.list() + registry.archivedSessionIds +``` + +- `createWorkspaces`: file records whose `path` is not a key in `local.byPath`. +- `attach`: for every file record, every `sessionIds` entry not already in `local.byPath.get(path).sessionIds`. (Because `attachSession` is idempotent, passing already-attached ids is also correct; the planner just avoids log noise.) +- `archive`: file `global.archivedSessionIds` minus `local.archived`. + +### 6.3 Apply (impure, serialized) + +``` +for each createWorkspace: + try registry.create(path, title) // throws if dir missing/not a dir + catch log warn, continue + +byPath = new Map(registry.list().map(w => [w.path, w])) // refresh AFTER creates + +for each attach: + workspace = byPath.get(workspacePath) + if !workspace: log warn, continue + try await workspace.attachSession(sessionId) + catch log warn, continue +for each archive: + try await registry.archiveSession(sessionId) + catch log warn, continue +``` + +Re-resolve entities after the create phase; ids are assigned by `create`. + +### 6.4 Termination + +The pass is idempotent: after it runs, the file contains the union and every subsequent pass plans zero actions, so the registry performs **zero writes** and emits no `domain/changed`. There is no ping-pong. Instrument the executor with a counter so a test can assert "steady state = 0 writes". + +### 6.5 Scheduling + +- Watch the **parent directory**, not the file. The store is committed by temp-write + `rename()`, which replaces the inode; a watch bound to the file can go stale, while a directory watch sees the rename land. Filter on the basename: + +```js +import { watch } from 'node:fs' +import { dirname, basename } from 'node:path' + +const file = workspaceFile() +const watcher = watch(dirname(file), (_event, name) => { + if (name === basename(file)) schedule() +}) +// dispose: watcher.close() +``` + +`fs.watchFile(file, { interval: 500 }, schedule)` is the belt-and-braces alternative — it polls the *path*, so it is immune to both rename and watch-invalidation — but it always polls. Prefer the directory watch; add the poll only if you observe missed events on this machine. +- `schedule` debounces (`debounceMs`, default 100) and enqueues onto a single-flight chain: + +```js +let tail = Promise.resolve() +const run = () => { tail = tail.then(pass).catch((e) => ctx.logger.warn(...)) ; return tail } +``` + +- Never run two passes concurrently — the registry's own `operationTail` and domain write chain serialize individual writes, but our read-plan-apply sequence must not interleave with itself. +- One optional pass at start, after `inject` resolves. Harmless and idempotent. +- On `pendingMutation`, schedule one retry at `retryOnPendingMs` (default 250); the marker clearing is itself a file write, so the watcher usually handles it. + +## 7. Edge cases and traps + +### 7.1 Table + +| Case | Handling | +|---|---| +| `workspaceRegistry` not started yet | `registry.list()` throws `workspace registry is not started yet`; catch and retry on next event. `inject: ['workspaceRegistry']` should already prevent this. | +| File absent | No-op (`ENOENT`). | +| Malformed / schema-invalid file | Log, no-op. Never "repair" by writing. | +| `pendingMutation` present | Skip the pass; retry after `retryOnPendingMs`. Another instance is mid create/delete. | +| Workspace path no longer a directory | `create` throws; log and skip that workspace. | +| Session log not yet materialized | `attachSession` throws not-found; log and skip — the next pass retries. | +| Session `cwd` moved/deleted | `attachSession` throws; log and skip. | +| HMR reload | Close the watcher in the `ctx.effect` disposer, or a reload leaks a second watcher and doubles the passes. | +| `$DSH_HOME` set after boot | Resolve `workspaceFile()` inside each pass. | +| Our own write triggers the watcher | Pass runs, plans nothing, writes nothing. Expected; assert it in tests. | +| Two instances attach the same id simultaneously | Both calls are idempotent; the file ends at the union. | +| Duplicate `fs.watch` events for one write | Debounce + single-flight. | + +### 7.2 Why not reconcile removals (read this before "improving" it) + +`updatedAt` cannot arbitrate. Trace: A creates `S_A` and writes `{S_A}`; B, whose domain memory predates that, creates `S_B` and writes `{S_B}`. The file now has `{S_B}`, and `S_A` is absent. A "newer file wins" rule makes A detach `S_A` — reintroducing exactly the loss we're fixing. There is no field in the record that says "the writer knew about `S_A`". Without a revision/CAS in the domain layer, "deliberate removal" and "stale-write loss" are indistinguishable. Additive-only is the lossless choice. + +Consequences to document in the plugin README: + +- A `detachSession` performed by another instance (only `dsh-webhook` calls it in this build, `dsh-webhook/lib/index.js:137`) will be undone by our union. Acceptable: no GUI flow detaches. +- Session order within a workspace is not reconciled (`attachSession` prepends), so the two instances can show different orders. Cosmetic. +- Workspace order and titles are not reconciled. Cosmetic. +- `archiveSession` reconciliation is append-only, which matches the product (there is no unarchive). + +### 7.3 Acknowledge V1's incomplete-write window + +Between reading the file and committing attaches, another instance can write. Our attach may then be based on a file snapshot that has already moved. Because we only ever *add*, the worst case is a redundant pass; the union still converges once both watchers have seen the latest write. State this in the README rather than implying atomicity. + +## 8. Plugin skeleton + +```js +/** + * fix-workspace-multiple-writers: reconcile $DSH_HOME/storages/workspace.json into + * ctx.workspaceRegistry so concurrent dsh instances see each other's sessions. + * + * Host-only. Additive-only by design — see HANDOFF.md §7.2. + * Zero runtime dependencies: node built-ins only (see HANDOFF.md §4). + */ +import { readFile } from 'node:fs/promises' +import { watch } from 'node:fs' +import { basename, dirname, join, resolve } from 'node:path' +import { homedir } from 'node:os' + +export const name = 'fix-workspace-multiple-writers' +export const inject = ['workspaceRegistry'] + +export function apply(ctx, config = {}) { + const debounceMs = config.debounceMs ?? 100 + const retryOnPendingMs = config.retryOnPendingMs ?? 250 + + const dshHome = () => { + const env = process.env.DSH_HOME + return resolve(env !== undefined && env.trim().length > 0 ? env : join(homedir(), '.dsh')) + } + const workspaceFile = () => join(dshHome(), 'storages', 'workspace.json') + + let tail = Promise.resolve() + let timer + + const schedule = () => { + clearTimeout(timer) + timer = setTimeout(run, debounceMs) + } + const run = () => { + tail = tail.then(pass).catch((error) => { + ctx.logger.warn('fix-workspace-multiple-writers: reconcile failed: %s', String(error)) + }) + return tail + } + + async function pass() { + const file = workspaceFile() + let text + try { + text = await readFile(file, 'utf8') + } catch (error) { + if (error.code === 'ENOENT') return + throw error + } + + const parsed = parseStore(text) // §5.2 structural validation; returns undefined when unusable + if (parsed === undefined) return + if (parsed.global.pendingMutation !== undefined) { + clearTimeout(timer) + timer = setTimeout(run, retryOnPendingMs) + return + } + if (parsed.global.initialized !== true) return + + // Plain local snapshot — keeps planReconcile pure (§6.2). + const registry = ctx.workspaceRegistry + const local = { + byPath: new Map(registry.list().map((w) => [ + w.path, + { id: String(w.id), sessionIds: new Set(w.sessionIds.map(String)) }, + ])), + archived: new Set(registry.archivedSessionIds.map(String)), + } + + const plan = planReconcile(parsed, local) + if (plan.createWorkspaces.length === 0 && plan.attach.length === 0 && plan.archive.length === 0) return + + await applyPlan(ctx, plan) // §6.3 + } + + ctx.effect(() => { + const target = workspaceFile() + const watcher = watch(dirname(target), (_event, filename) => { + if (filename === basename(target)) schedule() + }) + watcher.on('error', (error) => ctx.logger.warn('fix-workspace-multiple-writers: watcher error: %s', String(error))) + return () => { + clearTimeout(timer) + watcher.close() + } + }, 'fix-workspace-multiple-writers: watch workspace store') +} +``` + +Keep `parseStore`, `planReconcile`, and `applyPlan` in small named functions (ideally `lib/plan.js`) so §9.1 can test the planner without booting dsh. + +## 9. Testing + +### 9.1 Unit — the planner + +Export `planReconcile` (or keep it in a separate `lib/plan.js`) and test it directly: + +- file adds a session to a known workspace → one `attach`, nothing else +- file has an unknown workspace → one `create` + its `attach` list +- file and local already agree → all three arrays empty (**the steady-state contract**) +- file archives an id → one `archive`; already-archived → empty +- file references an id local already has → no `attach` + +### 9.2 Integration — two live instances + +```bash +# terminal 1 (this GUI is already one; otherwise) +dsh web --port 3080 + +# terminal 2 +cd /path/to/checkout +nix run ./experimental/dsh-everything-edition#main-any-port +``` + +1. Create a session in instance B in the workspace instance A is showing. +2. Within ~1s, A's sidebar must show it under the correct workspace (not *Ungrouped*). +3. `node -e "…"` inspect `~/.dsh/storages/workspace.json` and assert the record's `sessionIds` is the union. +4. Create a session in A; confirm B converges. +5. Repeat 3–4 a few times; assert the union only grows and the file never loses an id. + +### 9.3 Regression checks + +- **Restart pickup**: stop A, create a session in B, start A. A's registry bootstraps from the file, so it should already have it; our initial pass must be a no-op. +- **Pending marker**: hand-write a `pendingMutation` into the file, touch it, assert no writes occur. +- **Unknown session**: add a fabricated id to the file, assert one warning and no crash. +- **Steady state**: after convergence, assert zero `putRecord` calls while touching the file repeatedly. +- **Disable**: comment the row out of `plugins.cordis.yml`; dsh must boot and behave as before. + +## 10. Non-goals (do not implement here) + +- Cross-process atomicity for read-modify-write. That is Route 2 (a locked KV backend that owns the write path via `withFileLock`/`writeFileAtomic` from `@deepseek-ai/dsh-atomic-write`) or upstream (`reload`/invalidate seam, `per-record`/SQLite layout, or the "cross-process revision pattern" the storage-domain README already lists as deferred). +- Removal/reorder/title reconciliation (§7.2). +- Any client-side (browser) half. +- Any change to `$DSH_HOME/sessions/` — session logs are already cross-process safe (per-session `flock` on `session.lock`). + +## 11. Acceptance criteria + +1. A session created in instance B appears in instance A's sidebar under the correct workspace without a page reload. +2. `workspace.json`'s `sessionIds` for the shared workspace equals the union; no id present before a pass is ever absent after one. +3. A converged system performs zero durable writes when the file is touched (idempotence). +4. The plugin is host-only, removable by deleting its row, and never throws into the tree. +5. `planReconcile` has unit coverage for every row of §9.1. + +## 12. Reference map (this build) + +| Claim | Location | +|---|---| +| Workspace store root is `dshHomePath('storages')` | `dsh-base/cordis.patch.yml:148–151` | +| Workspace domain row (web bundle) | `dsh-web-app/cordis.patch.yml:75` | +| Workspace domain spec: no `layout` → `single`; version 2 | `dsh-workspace/lib/types/spec.js` | +| `defineDomain` default; `DomainFacility.open` loads `loadAll` once | `dsh-storage-domain/lib/index.js:59–88`, `:337–360` | +| `DomainImpl` in-memory tables; `table.update` reads memory | `dsh-storage-domain/lib/index.js:107–280` | +| `SingleJsonUnit` — state loaded at open, whole-file publish | `dsh-storage-json/lib/index.js:159–275` | +| `WorkspaceEntity.attachSession` / `mutate` / `unchangedSentinel` | `dsh-workspace/lib/index.js:77–200` | +| `WorkspaceRegistry` public API | `dsh-workspace/lib/index.js:313–500` | +| `attachSession` called on session create | `dsh-api-session-controller/lib/index.js:587` | +| `WorkspaceFeed` subscribes to `domain/changed` | `dsh-api-workspace-controller/lib/index.js:44–55` | +| `withFileLock` / `writeFileAtomic` (Route 2 only) | `dsh-atomic-write/lib/types/index.d.ts` | +| `tryLockExclusive` (kernel-released alternative) | `node-addon-system/lib/flock.js` | +| Include entry resolution anchors at the include file (no profile fallback) | `cordis-plugin-loader/lib/index.js:270–279` | +| Include row shape | `dsh-plugins/plugins.cordis.yml` | diff --git a/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/README.md b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/README.md new file mode 100644 index 0000000..a76dc7e --- /dev/null +++ b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/README.md @@ -0,0 +1,200 @@ +# dsh-plugin-fix-workspace-multiple-writers + +Host-only dsh plugin. While two dsh instances share one `$DSH_HOME`, each keeps +the entire workspace store in memory — the JSON backend reads the document once, +the domain loads its tables once, and the registry holds its own record +snapshots. A whole-file publish from one instance therefore silently drops +additions the other made, and no process ever notices. + +This plugin watches `$DSH_HOME/storages/workspace.json` and replays the +**additions** it describes through `ctx.workspaceRegistry`'s public API. That is +the same write path a normal session create uses, so each attach emits +`domain/changed` and every connected browser shows the session without a reload. + +## What it does + +- creates a locally-unknown workspace the store describes, +- attaches a session id the store lists that this process does not account for, +- archives a session id the store archived. + +## What it deliberately does not do + +**Never removes, never reorders, never renames, never deletes.** `updatedAt` +cannot arbitrate a conflict: if A writes `{S_A}` and B writes `{S_B}` from stale +memory, the file holds `{S_B}` and `S_A` is simply absent. A "newer file wins" +rule makes A detach `S_A` — reintroducing exactly the loss this plugin fixes. +Nothing in a record says "the writer knew about `S_A`", so a deliberate removal +and a stale-write loss are indistinguishable, and additive-only is the lossless +choice. + +Consequences worth knowing: + +- A `detachSession` performed by another instance is undone by this union. No + GUI flow detaches; only the webhook path calls it in this build. +- Session order inside a workspace is not reconciled (`attachSession` prepends), + so two instances can display different orders. Cosmetic. +- Workspace order and titles are not reconciled. Cosmetic. +- Archive reconciliation is append-only, which matches the product: there is no + unarchive. +- **The plugin never asks for a removal, but a write it induces can still prune.** + The registry re-derives a record's whole session account from its *local* + canonical-cwd index on **every** write (`WorkspaceEntity.mutate`), and writes + whenever that filtering changes the array. So attaching one valid session to a + workspace whose record also holds an id that no longer validates — a moved or + deleted session directory, which the registry itself warns about at startup — + drops that already-filtered id from the durable record. The plugin has no API + that could request this, and any other writer would trigger the same prune, but + the trigger is often this plugin's attach. Pinned by + `test/integration.test.js` ("attaching a valid session prunes a + locally-unvalidatable id"). +- If the durable accounts cannot be read at all (the storage domain is + unreachable), the planner cannot tell "already accounted" from "accounted but + locally unvalidatable", so it declines to attach to any workspace it already + knows. Only newly created workspaces are attached to, since a fresh record is + empty and cannot be pruned. + +**It is not atomic, and that cuts both ways.** Every registry write republishes +the whole store from *this* process's memory, which was loaded once at open. So +between this plugin reading the store and committing an attach, a peer can +publish an addition that this process never saw; the plugin's write then reverts +the file to its own view and the peer's addition is gone from disk. The peer does +not repair it either: its plan is store-minus-local, and its own addition is +still in its memory, so it sees no divergence. The addition survives in the +peer's memory (and the file regains it whenever the peer next writes), but if +the peer exits first it is lost, because `initialized: true` means a restart +never re-bootstraps from the session logs. + +The window is small — the read-to-commit span of one pass — but it is real, it is +the same class as the bug this plugin addresses, and it cannot be closed from +here. Closing it needs cross-process serialization of the read-modify-write, +which is the separate "locked KV backend" route the handoff describes, not +something a watcher can do. + +## Behaviour + +- The store path is **derived, not assumed**: the open domain's own unit is asked + for its file (`$DSH_HOME/storages/workspace.json` is only the composed + default, and a relocated backend root leaves a stale file behind at the old + path). A directory-backed (`per-record`) unit has no single store document, so + the plugin declines to run rather than read a directory as a file. +- The store's **containing directory** is watched, not the file: the backend + commits by temp-write + `rename()`, which replaces the inode. A watch that + cannot be created (the directory does not exist yet) degrades to + `fs.watchFile`, which polls the path itself. +- Events are debounced and passes run on a single-flight chain, so a pass never + interleaves with itself. +- `ENOENT` means "nothing to add". An unparseable or foreign document is logged + and left byte-identical — it is never "repaired" by a write. +- A `pendingMutation` marker means another instance is mid create/delete; the + pass is deferred. The retry is **bounded** (8 attempts), after which the plugin + stops polling and waits for the store to change — a writer that died between + its two writes would otherwise leave an indefinite poll. A pass that throws is + retried on the same budget. +- The pass is idempotent. Once both sides agree it plans nothing, so a converged + system performs **zero durable writes** and emits no `domain/changed`. +- The watcher is bound once, at mount, to the store's directory. Changing + `$DSH_HOME` after boot makes later passes read the new file while the watcher + still listens to the old directory, so that (exotic) case needs a restart. + +## Configuration + +```yaml +- id: dsh-plugin-fix-workspace-multiple-writers + name: ./dsh-plugin-fix-workspace-multiple-writers/lib/index.js + config: + debounceMs: 100 # event coalescing window + retryOnPendingMs: 250 # retry delay while a pendingMutation marker is up + pollIntervalMs: 1000 # fs.watchFile interval, used only by the fallback +``` + +Delete the row to remove the capability; dsh boots and behaves as before. + +Row changes take effect on the **next dsh start**. HMR watches the modules under +the checkout's watch root — editing `lib/index.js` reloads a mounted plugin +live — but the loader has no watcher for a nested `cordis:include` YAML, so the +row list itself is read once at boot. (Only the profile's own patch layer is +HMR-registered.) + +The include resolves to the checkout **the launcher ran from** (`dsh-main` +derives it from `git rev-parse --show-toplevel`), so this directory must exist in +*that* checkout's `dsh-plugins/`. A plugin authored in a different working copy +or clone of the repository is invisible to a dsh started elsewhere: no row is +ever created, so it appears nowhere in Settings → Plugins → Plugin list. + +## Verifying a running instance + +**1. Did the row load?** A mounted plugin logs one line at startup: + +``` +fix-workspace-multiple-writers: watching /Users/you/.dsh/storages/workspace.json (domain unit) +``` + +The parenthetical is where the path came from — `domain unit`, or `assumed +($DSH_HOME/storages)` when the storage domain could not be asked. It is silent +while converged, so that line is the only proof of life. If it is missing after a +restart, the include is not mounted at all — start dsh through +`experimental/dsh-everything-edition`'s `dsh-main` / `dsh-main-any-port`, which +is what injects the `dsh-everything-checkout-plugins` include; plain `dsh web` +mounts none of the checkout plugins. + +**2. Does it reconcile?** Create a divergence and watch it heal. Any stored +session whose `cwd` is the workspace path but which the record does not account +for will do; inject it the way a stale whole-file write would, then read back: + +```sh +cd /path/to/the/workspace +SID=session-00000000-0000-7000-8000-000000000000 # a real, stored session id + +node -e ' +const fs=require("node:fs"),os=require("node:os"),p=require("node:path"),c=require("node:crypto"); +const file=p.join(process.env.DSH_HOME||p.join(os.homedir(),".dsh"),"storages","workspace.json"); +const target=fs.realpathSync(process.argv[2]||process.cwd()); +const doc=JSON.parse(fs.readFileSync(file,"utf8")); +const ws=Object.values(doc.tables.workspaces).find((r)=>r.path===target); +console.log("before:",ws.updatedAt,ws.sessionIds.length); +ws.sessionIds=[...ws.sessionIds.filter((id)=>id!==process.argv[1]),process.argv[1]]; +const tmp=file+"."+c.randomUUID()+".tmp"; +fs.writeFileSync(tmp,JSON.stringify(doc,null,2)+"\n");fs.renameSync(tmp,file); +' "$SID" +``` + +Within ~1s the store's `updatedAt` for that workspace advances and `$SID` moves +to `sessionIds[0]` (`attachSession` prepends) — the plugin attached it through +the registry, and the log prints `reconciled 0 workspace(s), 1 session(s), 0 +archive(s)`. Without the plugin `updatedAt` never changes: nothing else reads +that file. An id that cannot be validated (unknown session, or a `cwd` that no +longer matches) is warned about and skipped instead, which is equally conclusive. + +**3. The real acceptance test** is two live instances: run a second one with +`nix run ./experimental/dsh-everything-edition#main-any-port`, create a session +in the shared workspace there, and watch it appear in the first instance's +sidebar — same workspace, no page reload. Comment the row out and repeat to see +the original loss. + +## Zero runtime dependencies + +This is a hard constraint, not a style preference. A plugin mounted through the +loader's `cordis:include` resolves its own imports from its own directory, which +never reaches the profile's `node_modules`. Any package import fails with +`ERR_MODULE_NOT_FOUND`. The plugin therefore imports only `node:` built-ins. + +## Tests + +```sh +npm test # unit + fake-registry tests; 3 real-stack tests skip +DSH_PACKAGES_ROOT=/lib/node_modules/@deepseek-ai/dsh/node_modules/@deepseek-ai npm test +``` + +`test/integration.test.js` boots the real storage hub, JSON backend, domain +facility and `WorkspaceRegistry` from an installed dsh and asserts that a store +published by a second instance reaches the registry, emits `domain/changed`, and +then converges to zero further writes. + +`test/composition.test.js` mounts `dsh-plugins/plugins.cordis.yml` through the +installed dsh's real Cordis loader — the same `cordis:include` create that +`dsh-app-boot` performs — and asserts this row activates and runs a pass, so a +mistyped row can never leave the plugin silently inert. + +`DSH_PACKAGES_ROOT` names the `@deepseek-ai` directory of a dsh install; the +plugin cannot locate one itself, because it may not import anything but +`node:` built-ins. diff --git a/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/index.js b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/index.js new file mode 100644 index 0000000..ac3bcad --- /dev/null +++ b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/index.js @@ -0,0 +1,420 @@ +/** + * fix-workspace-multiple-writers: reconcile `$DSH_HOME/storages/workspace.json` + * into `ctx.workspaceRegistry` so concurrent dsh instances converge on the same + * workspaces and sessions. + * + * Two instances share the store but not their in-memory state: the JSON backend + * reads the document once, the domain loads its tables once, and each registry + * keeps its own record snapshot. A whole-file publish from one instance + * therefore silently drops additions the other made, and nothing in-process + * ever notices. This plugin watches the store and replays the *additions* it + * describes through the registry's public API — the same write path a normal + * session create uses — so `domain/changed` reaches the live GUI. + * + * Additive-only by design: it creates locally-unknown workspaces, attaches + * session ids the file lists, and archives ids the file archived. It never + * removes, reorders, or renames, because a stale whole-file write is + * indistinguishable from a deliberate removal (see README.md). + * + * Host-only, and zero runtime dependencies: `node:` built-ins only. An + * include-mounted plugin resolves its imports from its own directory rather + * than the profile's `node_modules`, so any package import would fail to load. + */ +import { unwatchFile, watch, watchFile } from 'node:fs' +import { readFile } from 'node:fs/promises' +import { homedir } from 'node:os' +import { basename, dirname, join, resolve } from 'node:path' + +import { isPlanEmpty, planReconcile } from './plan.js' +import { parseStore } from './store.js' + +/** Plugin name reported to the loader and to the logger. */ +export const name = 'fix-workspace-multiple-writers' +/** Hard dependency: nothing is reconcilable until the registry has started. */ +export const inject = ['workspaceRegistry'] + +const LOG_PREFIX = 'fix-workspace-multiple-writers' +const DEFAULT_DEBOUNCE_MS = 100 +const DEFAULT_RETRY_ON_PENDING_MS = 250 +const DEFAULT_POLL_INTERVAL_MS = 1000 +/** Bound on retries for a condition the plugin cannot itself clear. */ +const MAX_RETRIES = 8 + +/** + * Expand a literal `~` prefix the way `dsh-home-paths` does, so an unexpanded + * `$DSH_HOME` resolves to the same directory the storage backend writes to. + * @param path - Possibly tilde-prefixed path. + * @returns the expanded path. + */ +function expandHome(path) { + if (path === '~') return homedir() + if (path.startsWith('~/') || path.startsWith('~\\')) return join(homedir(), path.slice(2)) + return path +} + +/** + * Resolve the harness home with `resolveDshHome`'s precedence: `$DSH_HOME` + * (empty or whitespace-only counts as unset) then `~/.dsh`. Resolved per pass, + * never cached at boot, so a late `$DSH_HOME` is honored. + * @param env - Environment mapping to read; defaults to `process.env`. + * @returns the absolute harness home. + */ +export function resolveDshHome(env = process.env) { + const configured = env.DSH_HOME + return resolve(expandHome(configured !== undefined && configured.trim().length > 0 ? configured : join(homedir(), '.dsh'))) +} + +/** + * The shared workspace store this plugin reconciles. + * @param env - Environment mapping to read; defaults to `process.env`. + * @returns the absolute path of `storages/workspace.json`. + */ +export function workspaceStorePath(env = process.env) { + return join(resolveDshHome(env), 'storages', 'workspace.json') +} + +/** + * The store document this process's domain actually backs. + * + * `$DSH_HOME/storages/workspace.json` is only the *composed default*: a patch + * layer or a `storage-domain.routes` entry can move the backend root, and a host + * that moved it leaves the old file behind — which a watcher bound to the + * default path would happily replay into the live registry. The open domain + * knows its own unit, so that identity is preferred over the assumption. + * + * Seeing a unit but not a single-file path is not a reason to guess: a + * directory-backed (`per-record`) unit has no whole-store document, so this + * reports `path: undefined` and the caller declines to reconcile rather than + * reading a directory as a file. + * @param ctx - Plugin context; `storageDomain` is read optionally. + * @returns `{ path, source }`, or `{ path: undefined, reason }` when the domain + * has no single store document. + */ +export function resolveStorePath(ctx) { + // `ctx.get` is the optional accessor: `storageDomain` is guaranteed to exist + // whenever the registry does, but the plugin must not hard-inject it. + const unit = ctx.get('storageDomain')?.get('workspace')?.unit + if (unit === undefined) return { path: workspaceStorePath(), source: 'assumed ($DSH_HOME/storages)' } + if (typeof unit.path === 'string' && unit.path.length > 0 && unit.descriptor?.layout !== 'per-record') { + return { path: unit.path, source: 'domain unit' } + } + return { + path: undefined, + source: 'domain unit', + reason: `the workspace domain is backed by a ${unit.descriptor?.layout ?? 'directory-backed'} unit, which has no single store document to reconcile`, + } +} + +/** + * Render any thrown value as a one-line message for a log record. + * @param error - Thrown value. + * @returns a human-readable message. + */ +function describe(error) { + return error instanceof Error ? error.message : String(error) +} + +/** + * Accept a configured delay only when it is a usable non-negative number. + * @param value - Configured value. + * @param fallback - Delay to use otherwise. + * @returns a millisecond delay. + */ +function delayMs(value, fallback) { + return typeof value === 'number' && Number.isFinite(value) && value >= 0 ? value : fallback +} + +/** + * Plain, I/O-free snapshot of what this process already knows, for the planner. + * + * Membership comes from the durable record whenever the storage domain is + * reachable, because `Workspace.sessionIds` is filtered by the local + * canonical-cwd index: it can omit an id the record still accounts for. + * Planning an attach from the filtered view would make the registry re-run its + * prune step and *write such an id out* of the durable record. + * + * `exact` reports whether the durable accounts were readable at all. When they + * were not, the filtered view cannot distinguish "already accounted" from + * "accounted but locally unvalidatable", so the planner declines to attach to + * any workspace it already knows rather than risk inducing a prune. Creating a + * newly described workspace stays safe either way, because a fresh record is + * empty. + * @param ctx - Plugin context carrying `workspaceRegistry`. + * @returns `{ byPath, archived, exact }` for `planReconcile`. + */ +export function snapshotRegistry(ctx) { + const registry = ctx.workspaceRegistry + const accounts = durableAccounts(ctx) + const byPath = new Map() + for (const workspace of registry.list()) { + const id = String(workspace.id) + const sessionIds = accounts?.get(id) ?? workspace.sessionIds.map(String) + byPath.set(workspace.path, { id, sessionIds: new Set(sessionIds) }) + } + return { byPath, archived: new Set(registry.archivedSessionIds.map(String)), exact: accounts !== undefined } +} + +/** + * Read every workspace's raw `sessionIds` straight from the domain table. + * @param ctx - Plugin context; `storageDomain` is read optionally. + * @returns session-id sets keyed by workspace id, or `undefined` when the + * domain is not reachable. + */ +function durableAccounts(ctx) { + // `ctx.get` is the optional accessor: `storageDomain` is guaranteed to exist + // whenever the registry does, but the plugin must not hard-inject it. + const storageDomain = ctx.get('storageDomain') + if (storageDomain === undefined) return undefined + const table = storageDomain.get('workspace')?.table('workspaces') + if (table === undefined) return undefined + const accounts = new Map() + for (const [id, record] of table.entries()) { + if (Array.isArray(record?.sessionIds)) accounts.set(String(id), record.sessionIds.map(String)) + } + return accounts +} + +/** + * Execute one plan serially through the registry's public API. Every action is + * independently guarded: one unusable workspace, session, or archive request + * never aborts the rest of the pass, and never throws into the caller. + * @param ctx - Plugin context carrying `workspaceRegistry`. + * @param plan - A plan from `planReconcile`. + * @returns counts of the actions attempted or completed, for tests and logs. + */ +export async function applyPlan(ctx, plan) { + const registry = ctx.workspaceRegistry + const byPath = new Map(registry.list().map((workspace) => [workspace.path, workspace])) + const result = { created: 0, attached: 0, archived: 0, failed: 0 } + + for (const request of plan.createWorkspaces) { + try { + const workspace = await registry.create(request.path, request.title) + byPath.set(workspace.path, workspace) + result.created += 1 + } catch (error) { + result.failed += 1 + ctx.logger.warn('%s: cannot adopt workspace at "%s": %s', LOG_PREFIX, request.path, describe(error)) + } + } + + for (const request of plan.attach) { + const workspace = byPath.get(request.workspacePath) + // Only a create that just failed can land here, and it already warned. + if (workspace === undefined) continue + try { + await workspace.attachSession(request.sessionId) + result.attached += 1 + } catch (error) { + result.failed += 1 + ctx.logger.warn('%s: cannot attach session "%s" to "%s": %s', LOG_PREFIX, request.sessionId, request.workspacePath, describe(error)) + } + } + + for (const sessionId of plan.archive) { + try { + await registry.archiveSession(sessionId) + result.archived += 1 + } catch (error) { + result.failed += 1 + ctx.logger.warn('%s: cannot archive session "%s": %s', LOG_PREFIX, sessionId, describe(error)) + } + } + + return result +} + +/** + * Build the reconciler over one context. Exported so tests can drive passes + * directly, without a filesystem watcher or a debounce delay. + * @param ctx - Plugin context carrying `workspaceRegistry`. + * @returns the pass entry point and the store path it reads. + */ +export function createReconciler(ctx) { + const resolve = () => resolveStorePath(ctx) + const workspaceFile = () => resolve().path + + /** + * Run one read → plan → apply pass. Never writes the store except through the + * registry, and never rewrites a file it could not fully understand. + * @returns a status describing what the pass found and did. + */ + async function reconcileOnce() { + const resolved = resolve() + if (resolved.path === undefined) return { status: 'unsupported', reason: resolved.reason } + + let text + try { + text = await readFile(resolved.path, 'utf8') + } catch (error) { + // The backend materializes the store lazily; absence means no additions. + if (error?.code === 'ENOENT') return { status: 'absent' } + throw error + } + + const store = parseStore(text) + if (!store.ok) { + ctx.logger.warn('%s: leaving the workspace store untouched: %s', LOG_PREFIX, store.reason) + return { status: 'unusable', reason: store.reason } + } + if (!store.initialized) return { status: 'not-initialized' } + // Another instance is between its order and record writes; retry shortly. + if (store.pendingMutation) return { status: 'pending' } + if (store.skipped.length > 0) { + ctx.logger.warn('%s: ignored %d unusable workspace record(s) in the store', LOG_PREFIX, store.skipped.length) + } + + const plan = planReconcile(store, snapshotRegistry(ctx)) + if (isPlanEmpty(plan)) return { status: 'converged', plan } + + const applied = await applyPlan(ctx, plan) + // A fully successful pass is worth one line: it is the only visible sign + // that another instance's additions were picked up. Failures already warned + // per action, so a partly failed pass stays quiet here. + if (applied.failed === 0) { + ctx.logger.info('%s: reconciled %d workspace(s), %d session(s), %d archive(s)', LOG_PREFIX, applied.created, applied.attached, applied.archived) + } + return { status: 'applied', plan, applied } + } + + return { workspaceFile, reconcileOnce } +} + +/** + * Watch the store's containing directory and report file changes. + * + * The directory — not the file — is watched because the backend commits by + * temp-write + `rename()`, which replaces the inode; a watch bound to the file + * goes stale, while a directory watch sees the new entry arrive. A watch that + * cannot be created (the directory does not exist yet) degrades to + * `fs.watchFile`, which polls the path itself. + * @param logger - Logger for watcher failures. + * @param storeFile - Resolver for the absolute store path. + * @param onEvent - Called for a change to the store (or an unattributable one). + * @param pollIntervalMs - `fs.watchFile` interval used by the fallback. + * @returns a disposer closing the watcher. + */ +export function installStoreWatcher(logger, storeFile, onEvent, pollIntervalMs = DEFAULT_POLL_INTERVAL_MS) { + const target = storeFile() + const storeName = basename(target) + let watcher + try { + watcher = watch(dirname(target), (_event, filename) => { + // A null name cannot be attributed to another entry, so it is treated as + // a possible store event; the debounce absorbs the extra pass. + if (filename === null || String(filename) === storeName) onEvent() + }) + watcher.on('error', (error) => { + logger.warn('%s: workspace store watcher failed: %s', LOG_PREFIX, describe(error)) + }) + } catch (error) { + logger.warn('%s: cannot watch "%s" (%s); polling the store path instead', LOG_PREFIX, dirname(target), describe(error)) + } + + const poller = watcher === undefined ? watchFile(target, { interval: pollIntervalMs }, () => onEvent()) : undefined + // A pending poll must never hold the process open on its own. + poller?.unref?.() + + return () => { + watcher?.close() + if (poller !== undefined) unwatchFile(target) + } +} + +/** + * Mount the reconciler: one directory watch, one debounced single-flight pass + * chain, and one pass at start. + * @param ctx - Plugin context; `workspaceRegistry` is required by `inject`. + * @param config - Optional `debounceMs`, `retryOnPendingMs`, `pollIntervalMs`. + */ +export function apply(ctx, config) { + // A bare `config:` key in YAML arrives as null, and the loader treats a + // throwing row as fatal — so never index into it directly. + const options = config === null || config === undefined ? {} : config + const debounceMs = delayMs(options.debounceMs, DEFAULT_DEBOUNCE_MS) + const retryOnPendingMs = delayMs(options.retryOnPendingMs, DEFAULT_RETRY_ON_PENDING_MS) + const pollIntervalMs = delayMs(options.pollIntervalMs, DEFAULT_POLL_INTERVAL_MS) + const reconciler = createReconciler(ctx) + + let timer + let tail = Promise.resolve() + // Retry budget for the two conditions that are not the plugin's to fix: a + // peer's pending mutation, and a pass that failed. Both are retried a bounded + // number of times and then left to the watcher, so a condition that never + // clears (a writer that died between its two writes) cannot become an + // indefinite poll. + let retries = 0 + let suspended = false + + /** + * (Re)arm the debounce timer. Rescheduling replaces any pending pass, which + * is what collapses the burst of events one atomic publish produces. + * @param delay - Milliseconds to wait. + */ + const arm = (delay) => { + if (timer !== undefined) clearTimeout(timer) + timer = setTimeout(() => { + timer = undefined + void run() + }, delay) + timer.unref?.() + } + + /** + * Enqueue one pass on the single-flight chain, so a pass never interleaves + * with itself: planning reads a snapshot, and applying mutates the registry. + * @returns the chain tail. + */ + const run = () => { + const settle = (retry, describeGivingUp) => { + retries += 1 + if (retries <= MAX_RETRIES) return arm(retry) + if (!suspended) { + suspended = true + ctx.logger.warn('%s: %s; waiting for the store to change instead of polling', LOG_PREFIX, describeGivingUp) + } + return undefined + } + tail = tail + .then(async () => { + const outcome = await reconciler.reconcileOnce() + if (outcome.status === 'pending') { + return settle(retryOnPendingMs, `the store still carries a pendingMutation after ${MAX_RETRIES} retries`) + } + retries = 0 + suspended = false + return outcome + }) + .catch((error) => { + ctx.logger.warn('%s: reconcile pass failed: %s', LOG_PREFIX, describe(error)) + return settle(retryOnPendingMs, `the reconcile pass kept failing after ${MAX_RETRIES} retries`) + }) + return tail + } + + ctx.effect(() => { + const resolved = resolveStorePath(ctx) + if (resolved.path === undefined) { + // Publishing the row but refusing to act is the honest outcome: there is + // no whole-store document to reconcile against. + ctx.logger.warn('%s: not watching anything: %s', LOG_PREFIX, resolved.reason) + return () => {} + } + // One line at mount, so a fresh boot shows that the row actually loaded, + // which store it resolved, and on whose authority. The plugin is silent + // while converged. + ctx.logger.info('%s: watching %s (%s)', LOG_PREFIX, resolved.path, resolved.source) + const closeWatcher = installStoreWatcher(ctx.logger, reconciler.workspaceFile, () => arm(debounceMs), pollIntervalMs) + // The registry bootstraps from this file, so the opening pass is normally a + // no-op; it covers an addition that landed between that bootstrap and this + // row's activation. + arm(debounceMs) + return () => { + if (timer !== undefined) { + clearTimeout(timer) + timer = undefined + } + closeWatcher() + } + }, `${LOG_PREFIX}: watch the workspace store`) +} diff --git a/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/plan.js b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/plan.js new file mode 100644 index 0000000..9ae9cbb --- /dev/null +++ b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/plan.js @@ -0,0 +1,62 @@ +/** + * The pure reconcile planner: the file's additions minus local knowledge. + * + * Kept free of I/O and of any registry reference so it can be exercised without + * booting dsh. The plan has no field that could remove, reorder, or rename — + * additive-only is a property of this type, not of the executor's discipline. + * @module dsh-plugin-fix-workspace-multiple-writers/plan + */ + +/** + * Diff one parsed store against a local snapshot. + * + * `local.byPath` keys are the file's own canonical paths (`fs.realpath` + * spellings, which is what the registry stores), and each entry's `sessionIds` + * is the locally-accounted set for that workspace. A workspace the file + * describes but the snapshot does not is created *and* has all of its session + * ids attached in the same pass. + * + * `local.exact === false` means the caller could only supply the registry's + * cwd-filtered view. Attaching from that view can make the registry prune an id + * it still accounts for, so attaches are then planned only for workspaces the + * snapshot does not know — a record about to be created is empty and cannot be + * pruned. + * @param store - A parsed store from `parseStore` (`workspaces`, + * `archivedSessionIds`). + * @param local - Local snapshot: `byPath: Map`, + * `archived: Set`, and optionally `exact: boolean` (default true). + * @returns the additive actions to perform, in create → attach → archive order. + */ +export function planReconcile(store, local) { + const createWorkspaces = [] + const attach = [] + const archive = [] + const exact = local.exact !== false + + for (const workspace of store.workspaces) { + const known = local.byPath.get(workspace.path) + if (known === undefined) createWorkspaces.push({ path: workspace.path, title: workspace.title }) + else if (!exact) continue + for (const sessionId of workspace.sessionIds) { + if (known !== undefined && known.sessionIds.has(sessionId)) continue + attach.push({ workspacePath: workspace.path, sessionId }) + } + } + + for (const sessionId of store.archivedSessionIds) { + if (local.archived.has(sessionId)) continue + archive.push(sessionId) + } + + return { createWorkspaces, attach, archive } +} + +/** + * Whether a plan asks for no work at all — the steady-state contract. A + * converged system plans nothing and therefore performs no durable write. + * @param plan - A plan from `planReconcile`. + * @returns whether every action list is empty. + */ +export function isPlanEmpty(plan) { + return plan.createWorkspaces.length === 0 && plan.attach.length === 0 && plan.archive.length === 0 +} diff --git a/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/store.js b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/store.js new file mode 100644 index 0000000..fa61255 --- /dev/null +++ b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/lib/store.js @@ -0,0 +1,89 @@ +/** + * Structural reader for `$DSH_HOME/storages/workspace.json`. + * + * The file is a `single`-layout storage unit: one JSON document holding a + * `unit` header, the `global` slot, and the `tables` maps. This module reads it + * defensively — plain structural checks, never a schema import — and never + * repairs it. A foreign, damaged, or partially written document yields + * `{ ok: false }` so the caller can log and leave the medium byte-identical; + * the registry re-validates every record with its real schema at the + * durability boundary, so a false accept here fails loudly rather than + * corrupting anything. + * + * Zero runtime dependencies: `node:` built-ins only. + * @module dsh-plugin-fix-workspace-multiple-writers/store + */ +import { basename } from 'node:path' + +/** Unit header identity this reader accepts; anything else is a foreign file. */ +export const WORKSPACE_UNIT_NAME = 'workspace' +/** Unit format version this reader accepts. */ +export const WORKSPACE_UNIT_VERSION = 2 + +/** + * Whether a decoded JSON value is a non-null, non-array object. + * @param value - Decoded JSON value. + * @returns whether the value is a plain JSON object. + */ +function isObject(value) { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +/** + * Parse one workspace-store document into the minimum the reconciler needs. + * + * `ok: false` means the whole pass must be skipped: the document is unparseable, + * foreign, or missing a slot the reconciler depends on. Per-record problems are + * not fatal — they are reported in `skipped` so the caller can warn and carry on + * with the records it did understand. + * @param text - Raw file content. + * @returns the usable projection, or the reason the document is unusable. + */ +export function parseStore(text) { + let document + try { + document = JSON.parse(text) + } catch { + return { ok: false, reason: 'file is not valid JSON' } + } + if (!isObject(document)) return { ok: false, reason: 'document is not a JSON object' } + + const unit = document.unit + if (!isObject(unit) || unit.name !== WORKSPACE_UNIT_NAME || unit.version !== WORKSPACE_UNIT_VERSION) { + return { ok: false, reason: `foreign or unsupported unit header (expected '${WORKSPACE_UNIT_NAME}' version ${WORKSPACE_UNIT_VERSION})` } + } + + const global = document.global + if (!isObject(global) || !Array.isArray(global.workspaceIds) || !Array.isArray(global.archivedSessionIds)) { + return { ok: false, reason: 'global slot is missing workspaceIds/archivedSessionIds arrays' } + } + + const tables = document.tables + if (!isObject(tables)) return { ok: false, reason: 'tables is not an object' } + const records = tables.workspaces + if (!isObject(records)) return { ok: false, reason: "tables.workspaces is not an object" } + + const workspaces = [] + const skipped = [] + for (const [id, record] of Object.entries(records)) { + if (!isObject(record) || typeof record.path !== 'string' || !Array.isArray(record.sessionIds) || !record.sessionIds.every((sessionId) => typeof sessionId === 'string')) { + skipped.push(id) + continue + } + workspaces.push({ + id, + path: record.path, + title: typeof record.title === 'string' ? record.title : basename(record.path), + sessionIds: [...record.sessionIds], + }) + } + + return { + ok: true, + initialized: global.initialized === true, + pendingMutation: global.pendingMutation !== undefined, + archivedSessionIds: global.archivedSessionIds.filter((sessionId) => typeof sessionId === 'string'), + workspaces, + skipped, + } +} diff --git a/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/package.json b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/package.json new file mode 100644 index 0000000..a230e78 --- /dev/null +++ b/dsh-plugins/dsh-plugin-fix-workspace-multiple-writers/package.json @@ -0,0 +1,22 @@ +{ + "name": "dsh-plugin-fix-workspace-multiple-writers", + "description": "Reconciles $DSH_HOME/storages/workspace.json into ctx.workspaceRegistry so concurrent dsh instances see each other's sessions.", + "version": "0.1.0", + "private": true, + "type": "module", + "main": "lib/index.js", + "exports": { + ".": "./lib/index.js", + "./store": "./lib/store.js", + "./plan": "./lib/plan.js", + "./package.json": "./package.json" + }, + "scripts": { + "test": "node --test" + }, + "files": [ + "lib/index.js", + "lib/store.js", + "lib/plan.js" + ] +} -- 2.51.2