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" + ] +}