diff --git a/.github/workflows/reuse-quality.yml b/.github/workflows/reuse-quality.yml index 66450e5..8084551 100644 --- a/.github/workflows/reuse-quality.yml +++ b/.github/workflows/reuse-quality.yml @@ -32,7 +32,28 @@ jobs: run: docker compose up -d --wait - name: Start dev API - run: NODE_ENV=development NODE_CONFIG_DIR=$PWD/api/config/ node api/index.ts & + run: | + set -a && . ./.env && set +a + NODE_ENV=development NODE_CONFIG_DIR=$PWD/api/config/ node api/index.ts & + + # Sync mirrors artefacts *from* an upstream registry, so exercising it needs a + # second one. Same code, different port / database / data dir. + - name: Start upstream registry API + run: | + set -a && . ./.env && set +a + NODE_ENV=development NODE_CONFIG_DIR=$PWD/api/config/ \ + PORT=$DEV_UPSTREAM_API_PORT \ + MONGO_URL=mongodb://localhost:$MONGO_PORT/data-fair-registry-upstream \ + DATA_DIR=./data-upstream \ + node api/index.ts & + + - name: Wait for both APIs + run: | + set -a && . ./.env && set +a + for p in $DEV_API_PORT $DEV_UPSTREAM_API_PORT; do + timeout 60 bash -c "until curl -sf http://localhost:$p/api/ping >/dev/null; do sleep 1; done" + echo "api on :$p is up" + done - name: Run tests run: npm run test-unit && npm run test-api diff --git a/.gitignore b/.gitignore index 97d8157..e76fcab 100644 --- a/.gitignore +++ b/.gitignore @@ -19,6 +19,7 @@ dev/logs/ dev/fixtures-output.json tmp/ data/ +data-upstream/ lib-node/**/*.js lib-node/**/*.d.ts lib-node/**/*.js.map diff --git a/.zellij.kdl b/.zellij.kdl index a63dc4a..9c13615 100644 --- a/.zellij.kdl +++ b/.zellij.kdl @@ -20,6 +20,10 @@ layout { command "bash" args "-ic" "nvm use > /dev/null 2>&1 && npm -w api run dev" } + pane name="api-upstream" { + command "bash" + args "-ic" "nvm use > /dev/null 2>&1 && npm run dev-api-upstream" + } } pane size=2 borderless=true { command "bash" diff --git a/AGENTS.md b/AGENTS.md index 2f62ab8..c0537cd 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -21,7 +21,9 @@ The development processes are managed by the user using zellij and docker compos Check if services are running: `bash dev/status.sh` -Log files are in `dev/logs/` (dev-api.log, dev-ui.log, docker-compose.log). +Four dev processes run under zellij: `api` (the registry, port `DEV_API_PORT`), `api-upstream` (a second registry used as a federation mirror source, port `DEV_UPSTREAM_API_PORT`), `ui`, and `deps` (docker compose). + +Log files are in `dev/logs/` (dev-api.log, dev-api-upstream.log, dev-ui.log, docker-compose.log). ### Vulnerability scanner (osv-scanner) @@ -43,12 +45,22 @@ Run all tests: npm run test -Run a specific test file: +Run a specific test file, or a single test. Note `npm run test` is a compound `a && b && c` script, so appending a file to it only filters the *last* project (e2e) and runs the others in full — use the per-project scripts: - npm run test tests/artefacts.api.spec.ts + npm run test-api -- tests/artefacts.api.spec.ts + npm run test-api -- tests/artefacts.api.spec.ts -g "upload happy path" + npm run test-unit -- tests/artefacts-operations.unit.spec.ts Test users are defined in @dev/resources/users.json and organizations in @dev/resources/organizations.json. +### Exercising federation sync + +Sync mirrors artefacts *from* an upstream registry, so it needs two registries. `api-upstream` is a second registry process (same code, `PORT`/`MONGO_URL`/`DATA_DIR` overridden). Pointing a registry at itself cannot work: selecting an artefact that already exists locally without an `origin` returns 409. + +`npm run dev:fixtures` seeds the upstream, mints a read key owned by org `test1`, registers the mirror and selects two artefacts. It stops short of syncing — click **Sync now** in the admin UI. + +`tests/remote-registries-sync.api.spec.ts` covers the mirror path end to end against the same upstream. + ## Vulnerability scanning Advisory, admin-only vulnerability scanning of `npm` artefacts (bundled `node_modules` tarballs) via the bundled **osv-scanner v2** binary. It never blocks uploads or downloads, and scan data is stripped from responses for non-admin callers. diff --git a/Dockerfile b/Dockerfile index 49bdaad..b7938ff 100644 --- a/Dockerfile +++ b/Dockerfile @@ -72,6 +72,12 @@ RUN npm -w ui run build ########################## FROM installer AS api-installer +# The installer stage's `npm install -w ui` rewrites the lock file against fresh +# registry metadata (it can bump optional transitive edges without re-locking +# their entries), which `npm ci` then rejects as out of sync. Restore the +# pristine lock so `npm ci` validates the committed one. +COPY --from=package-strip /app/package-lock.json package-lock.json + # `npm ci --omit=optional` + clean-modules drop sharp's platform-specific # binaries (@img/sharp-*). Stash the musl/x64 ones from the full install and # restore them afterwards so the runtime image can load sharp on alpine. diff --git a/api/package.json b/api/package.json index c996279..d674816 100644 --- a/api/package.json +++ b/api/package.json @@ -4,7 +4,8 @@ "type": "module", "license": "MIT", "scripts": { - "dev": "mkdir -p ../dev/logs && NODE_ENV=development EVENTS_LOG_LEVEL=warn nodemon -e js,ts,json index.ts 2>&1 | tee ../dev/logs/dev-api.log" + "dev": "mkdir -p ../dev/logs && NODE_ENV=development EVENTS_LOG_LEVEL=warn nodemon -e js,ts,json index.ts 2>&1 | tee ../dev/logs/dev-api.log", + "dev-upstream": "mkdir -p ../dev/logs && NODE_ENV=development EVENTS_LOG_LEVEL=warn nodemon -e js,ts,json index.ts 2>&1 | tee ../dev/logs/dev-api-upstream.log" }, "imports": { "#config": "./src/config.ts", @@ -38,6 +39,7 @@ "semver": "^7.6.0", "sharp": "^0.33.5", "tar": "^7.5.16", - "tar-stream": "^3.1.0" + "tar-stream": "^3.1.0", + "ws": "^8.21.0" } } diff --git a/api/src/app.ts b/api/src/app.ts index 8282d87..364ee44 100644 --- a/api/src/app.ts +++ b/api/src/app.ts @@ -51,6 +51,7 @@ if (process.env.NODE_ENV === 'development') { await mongo.accessGrants.deleteMany({ 'account.id': { $regex: /^test/ } }) await mongo.thumbnails.deleteMany({}) await mongo.remoteRegistries.deleteMany({}) + await mongo.db.collection('locks').deleteMany({}) await filesStorage.clean() res.send() }) @@ -66,6 +67,34 @@ if (process.env.NODE_ENV === 'development') { res.send() }) + // Hold a lock under a FOREIGN pid so the API process can neither release nor prolong it. + // Lets a test observe `running` deterministically instead of racing a real sync. + // The row still expires on its own via the locks TTL index if a test forgets to clean up. + app.put('/api/test-env/locks/:id', async (req, res) => { + assertReqInternal(req) + const now = new Date() + await mongo.db.collection('locks').insertOne({ + _id: req.params.id as any, + pid: 'test-env', + hostname: 'test-env', + createdAt: now, + updatedAt: now + }) + res.send() + }) + + app.delete('/api/test-env/locks/:id', async (req, res) => { + assertReqInternal(req) + await mongo.db.collection('locks').deleteOne({ _id: req.params.id as any }) + res.send() + }) + + app.get('/api/test-env/locks/:id', async (req, res) => { + assertReqInternal(req) + const lock = await mongo.db.collection('locks').findOne({ _id: req.params.id as any }) + res.json({ exists: !!lock }) + }) + app.put('/api/test-env/artefacts/:id/scan', async (req, res) => { assertReqInternal(req) const id = decodeURIComponent(req.params.id) diff --git a/api/src/remote-registries/operations.ts b/api/src/remote-registries/operations.ts index bd04310..c962418 100644 --- a/api/src/remote-registries/operations.ts +++ b/api/src/remote-registries/operations.ts @@ -15,3 +15,26 @@ export const filterSuggestedArtefacts = ( const results = listing.results.filter(a => !a.deprecated || selected.has(a._id)) return { results, count: results.length } } + +export type SyncState = 'running' | 'interrupted' | 'idle' + +// The lock row id in the shared `locks` collection. Holding it means a sync is in flight. +export const syncLockId = (registryId: string) => `sync-remote-${registryId}` + +// The ws channel a registry's sync progress is published on. The registry id IS a url, +// so it must be encoded — a raw `/` would shred this `/`-delimited channel name. +export const syncChannel = (registryId: string) => + `remote-registries/${encodeURIComponent(registryId)}/sync` + +// Running state is derived, never stored: a stored `running` flag becomes a lie the moment +// a process is killed. `interrupted` means the final write that sets lastSyncAt never +// happened, stranding the attempt's startedAt ahead of it. +export const syncState = ( + locked: boolean, + registry: { syncProgress?: { startedAt: string }, lastSyncAt?: string } +): SyncState => { + if (locked) return 'running' + const startedAt = registry.syncProgress?.startedAt + if (startedAt && (!registry.lastSyncAt || startedAt > registry.lastSyncAt)) return 'interrupted' + return 'idle' +} diff --git a/api/src/remote-registries/router.ts b/api/src/remote-registries/router.ts index 8439ea5..d7db332 100644 --- a/api/src/remote-registries/router.ts +++ b/api/src/remote-registries/router.ts @@ -4,8 +4,8 @@ import { httpError } from '@data-fair/lib-utils/http-errors.js' import { axiosBuilder } from '@data-fair/lib-node/axios.js' import mongo from '#mongo' import { cipher, decipher } from '../cipher.ts' -import { syncRemoteRegistry } from './sync.ts' -import { filterSuggestedArtefacts } from './operations.ts' +import { startSync } from './sync.ts' +import { filterSuggestedArtefacts, syncLockId, syncState } from './operations.ts' import * as postReqBody from '#doc/remote-registries/post-req/index.ts' import * as patchReqBody from '#doc/remote-registries/patch-req/index.ts' @@ -17,6 +17,15 @@ const extractShortId = (apiKey: string): string => { return match ? match[1] : apiKey.slice(0, 12) } +// One query for the whole page, not one per row. +const lockedLockIds = async (registryIds: string[]): Promise> => { + if (registryIds.length === 0) return new Set() + const rows = await mongo.db.collection('locks') + .find({ _id: { $in: registryIds.map(syncLockId) as any } }, { projection: { _id: 1 } }) + .toArray() + return new Set(rows.map(row => String(row._id))) +} + // Create remote registry router.post('/', async (req, res, next) => { try { @@ -48,7 +57,11 @@ router.get('/', async (req, res, next) => { try { await session.reqAdminMode(req) const results = await mongo.remoteRegistries.find({}, { projection: { apiKey: 0 } }).toArray() - res.json({ results, count: results.length }) + const locked = await lockedLockIds(results.map(r => r._id)) + res.json({ + results: results.map(r => ({ ...r, syncState: syncState(locked.has(syncLockId(r._id)), r) })), + count: results.length + }) } catch (err) { next(err) } }) @@ -58,7 +71,8 @@ router.get('/:id', async (req, res, next) => { await session.reqAdminMode(req) const doc = await mongo.remoteRegistries.findOne({ _id: req.params.id }, { projection: { apiKey: 0 } }) if (!doc) throw httpError(404, 'remote registry not found') - res.json(doc) + const locked = await lockedLockIds([doc._id]) + res.json({ ...doc, syncState: syncState(locked.has(syncLockId(doc._id)), doc) }) } catch (err) { next(err) } }) @@ -191,9 +205,7 @@ router.post('/:id/sync', async (req, res, next) => { const doc = await mongo.remoteRegistries.findOne({ _id: req.params.id }) if (!doc) throw httpError(404, 'remote registry not found') - syncRemoteRegistry(req.params.id).catch(err => { - console.error(`[sync] Manual sync error for ${req.params.id}:`, err.message || err) - }) + if (!await startSync(req.params.id)) throw httpError(409, 'sync already running') res.status(202).json({ message: 'sync started' }) } catch (err) { next(err) } }) diff --git a/api/src/remote-registries/sync.ts b/api/src/remote-registries/sync.ts index 129c0cf..d12f68d 100644 --- a/api/src/remote-registries/sync.ts +++ b/api/src/remote-registries/sync.ts @@ -1,11 +1,35 @@ import { randomUUID } from 'node:crypto' import locks from '@data-fair/lib-node/locks.js' import { axiosBuilder } from '@data-fair/lib-node/axios.js' +import { internalError } from '@data-fair/lib-node/observer.js' +import * as wsEmitter from '@data-fair/lib-node/ws-emitter.js' import type { AxiosInstance } from 'axios' import mongo from '#mongo' import { decipher } from '../cipher.ts' import { filesStorage } from '../files-storage/index.ts' import type { Artefact } from '#types/artefact/index.ts' +import { syncLockId, syncChannel } from './operations.ts' + +export type SyncEvent = { + running: boolean + startedAt: string + done: number + total: number + currentArtefact?: string + lastSyncAt?: string + lastSyncStatus?: 'success' | 'error' + lastSyncError?: string +} + +// A dropped progress frame is cosmetic — the next frame supersedes it — so an emit +// failure must never abort a sync. +const emitSync = async (remoteRegistryId: string, event: SyncEvent) => { + try { + await wsEmitter.emit(syncChannel(remoteRegistryId), event) + } catch (err) { + internalError('sync-ws-emit', err) + } +} const syncNpmArtefact = async (ax: AxiosInstance, remoteUrl: string, artefactId: string) => { const encodedId = encodeURIComponent(artefactId) @@ -123,67 +147,119 @@ const syncFileArtefact = async (ax: AxiosInstance, remoteUrl: string, artefactId } } -export const syncRemoteRegistry = async (remoteRegistryId: string) => { - const lockId = `sync-remote-${remoteRegistryId}` - const acquired = await locks.acquire(lockId) - if (!acquired) { - console.log(`[sync] Lock already held for ${remoteRegistryId}, skipping`) - return - } +// The actual work. Callers own the lock. +const runSync = async (remoteRegistryId: string) => { + const remote = await mongo.remoteRegistries.findOne({ _id: remoteRegistryId }) + if (!remote) return - try { - const remote = await mongo.remoteRegistries.findOne({ _id: remoteRegistryId }) - if (!remote) return + const startedAt = new Date().toISOString() + const total = remote.selectedArtefacts.length + let done = 0 - const apiKey = decipher(remote.apiKey) - const ax = axiosBuilder({ - baseURL: remote._id, - headers: { 'x-api-key': apiKey } - }) + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { $set: { syncProgress: { startedAt, done, total } } } + ) + await emitSync(remoteRegistryId, { running: true, startedAt, done, total }) - let hasErrors = false - let lastError = '' - - for (const artefactId of remote.selectedArtefacts) { - try { - // Fetch remote artefact to determine format - const encodedId = encodeURIComponent(artefactId) - const detailRes = await ax.get(`/api/v1/artefacts/${encodedId}`) - const format: Artefact['format'] = detailRes.data.format - - if (format === 'npm') { - // We already fetched detail, but syncNpmArtefact re-fetches for simplicity - await syncNpmArtefact(ax, remote._id, artefactId) - } else { - await syncFileArtefact(ax, remote._id, artefactId) - } - } catch (err: any) { - hasErrors = true - lastError = `${artefactId}: ${err.message || err}` - console.error(`[sync] Error syncing ${artefactId} from ${remote._id}:`, err.message || err) + const apiKey = decipher(remote.apiKey) + const ax = axiosBuilder({ + baseURL: remote._id, + headers: { 'x-api-key': apiKey } + }) + + let hasErrors = false + let lastError = '' + + for (const artefactId of remote.selectedArtefacts) { + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { $set: { 'syncProgress.currentArtefact': artefactId } } + ) + await emitSync(remoteRegistryId, { running: true, startedAt, done, total, currentArtefact: artefactId }) + + try { + const encodedId = encodeURIComponent(artefactId) + const detailRes = await ax.get(`/api/v1/artefacts/${encodedId}`) + const format: Artefact['format'] = detailRes.data.format + + if (format === 'npm') { + await syncNpmArtefact(ax, remote._id, artefactId) + } else { + await syncFileArtefact(ax, remote._id, artefactId) } + } catch (err: any) { + hasErrors = true + lastError = `${artefactId}: ${err.message || err}` + console.error(`[sync] Error syncing ${artefactId} from ${remote._id}:`, err.message || err) } + done++ await mongo.remoteRegistries.updateOne( { _id: remoteRegistryId }, - { - $set: { - lastSyncAt: new Date().toISOString(), - lastSyncStatus: hasErrors ? 'error' : 'success', - ...(hasErrors ? { lastSyncError: lastError } : {}), - ...(!hasErrors ? {} : {}) - }, - ...(!hasErrors ? { $unset: { lastSyncError: '' } } : {}) - } + { $set: { 'syncProgress.done': done } } ) + await emitSync(remoteRegistryId, { running: true, startedAt, done, total, currentArtefact: artefactId }) + } + + const lastSyncAt = new Date().toISOString() + const lastSyncStatus = hasErrors ? 'error' as const : 'success' as const + + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { + $set: { + lastSyncAt, + lastSyncStatus, + ...(hasErrors ? { lastSyncError: lastError } : {}) + }, + $unset: { + 'syncProgress.currentArtefact': '', + ...(hasErrors ? {} : { lastSyncError: '' }) + } + } + ) + + // The end event carries the terminal state, so the UI never refetches to learn the outcome. + await emitSync(remoteRegistryId, { + running: false, + startedAt, + done, + total, + lastSyncAt, + lastSyncStatus, + ...(hasErrors ? { lastSyncError: lastError } : {}) + }) +} + +// Returns as soon as the lock is taken; the work continues in the background. +// A held lock is a conflict the caller (a human clicking a button) should see. +export const startSync = async (remoteRegistryId: string): Promise => { + const lockId = syncLockId(remoteRegistryId) + if (!await locks.acquire(lockId)) return false + runSync(remoteRegistryId) + .catch(err => internalError('sync-remote-registry', err)) + .finally(() => locks.release(lockId).catch(err => internalError('sync-remote-registry-release', err))) + return true +} + +// Awaits completion. Used by the daily job, which syncs registries one at a time. +export const syncRemoteRegistry = async (remoteRegistryId: string): Promise => { + const lockId = syncLockId(remoteRegistryId) + if (!await locks.acquire(lockId)) return false + try { + await runSync(remoteRegistryId) } finally { await locks.release(lockId) } + return true } export const syncAllRemoteRegistries = async () => { const remotes = await mongo.remoteRegistries.find({}).toArray() for (const remote of remotes) { + // A held lock means a peer replica is already syncing this registry. That is the + // normal outcome of N replicas firing the same daily timer — not an error. await syncRemoteRegistry(remote._id).catch(err => { console.error(`[sync] Failed to sync ${remote._id}:`, err.message || err) }) diff --git a/api/src/server.ts b/api/src/server.ts index ba504da..760afd6 100644 --- a/api/src/server.ts +++ b/api/src/server.ts @@ -4,6 +4,8 @@ import { startObserver, stopObserver, internalError } from '@data-fair/lib-node/ import eventPromise from '@data-fair/lib-utils/event-promise.js' import eventsQueue from '@data-fair/lib-node/events-queue.js' import locks from '@data-fair/lib-node/locks.js' +import * as wsEmitter from '@data-fair/lib-node/ws-emitter.js' +import * as wsServer from '@data-fair/lib-express/ws-server.js' import { createHttpTerminator } from 'http-terminator' import { app } from './app.ts' import config from '#config' @@ -31,6 +33,11 @@ export const start = async () => { await mongo.init() await renameFilePathToPath(mongo.db) await locks.start(mongo.db) + await wsEmitter.init(mongo.db) + // `canSubscribe` is never reached for admins: ws-server short-circuits on + // sessionState.user?.adminMode before calling it. Remote-registry sync is an + // admin-only surface, so returning false refuses precisely everyone else. + await wsServer.start(server, mongo.db, async () => false) if (config.privateEventsUrl) { if (!config.secretKeys?.events) { @@ -65,6 +72,7 @@ export const stop = async () => { if (syncTimer) clearInterval(syncTimer) if (rescanTimer) clearInterval(rescanTimer) await httpTerminator.terminate() + await wsServer.stop() if (config.observer?.active) await stopObserver() await locks.stop() await mongo.client.close() diff --git a/api/types/remote-registry/schema.js b/api/types/remote-registry/schema.js index e637593..a1e1214 100644 --- a/api/types/remote-registry/schema.js +++ b/api/types/remote-registry/schema.js @@ -26,6 +26,17 @@ export default { lastSyncAt: { type: 'string', format: 'date-time' }, lastSyncStatus: { type: 'string', enum: ['success', 'error'] }, lastSyncError: { type: 'string' }, + syncProgress: { + type: 'object', + additionalProperties: false, + required: ['startedAt', 'done', 'total'], + properties: { + startedAt: { type: 'string', format: 'date-time' }, + done: { type: 'integer' }, + total: { type: 'integer' }, + currentArtefact: { type: 'string' } + } + }, createdAt: { type: 'string', format: 'date-time', readOnly: true }, updatedAt: { type: 'string', format: 'date-time', readOnly: true } } diff --git a/dev/fixtures.ts b/dev/fixtures.ts index 1a6f630..6939adf 100644 --- a/dev/fixtures.ts +++ b/dev/fixtures.ts @@ -8,7 +8,10 @@ import { readFile, writeFile } from 'node:fs/promises' import { existsSync } from 'node:fs' import { join } from 'node:path' import FormData from 'form-data' -import { superAdmin, axiosWithApiKey, baseURL } from '../tests/support/axios.ts' +import { + superAdmin, axiosWithApiKey, baseURL, axios as axiosFactory, + upstreamBaseURL, upstreamSuperAdmin, upstreamAxiosAuth, upstreamAxiosWithApiKey +} from '../tests/support/axios.ts' import { createTestTarball } from '../tests/support/test-tarball.ts' const OUTPUT_PATH = join(import.meta.dirname, 'fixtures-output.json') @@ -37,6 +40,10 @@ const saveOutput = async (out: OutputFile) => { const isHttp404 = (err: any) => err?.response?.status === 404 || err?.status === 404 const isHttp409 = (err: any) => err?.response?.status === 409 || err?.status === 409 +const anonymousPing = async (url: string) => { + await axiosFactory({ baseURL: url }).get('/api/ping') +} + async function main () { console.log(`→ Connecting to ${baseURL}`) const admin = await superAdmin @@ -210,6 +217,141 @@ async function main () { } } + // --- Federation upstream ------------------------------------------------ + // Sync mirrors artefacts *from* an upstream registry, so exercising it needs a + // second registry. Requires the `api-upstream` process (npm run dev-api-upstream); + // hard-fails if it is down, since that pane is part of the standard dev layout. + const upstreamUrl = upstreamBaseURL() + console.log(`\n→ Upstream registry ${upstreamUrl}`) + try { + await anonymousPing(upstreamUrl) + } catch { + throw new Error( + `upstream registry not reachable at ${upstreamUrl} — is the \`api-upstream\` zellij pane running? (npm run dev-api-upstream)` + ) + } + + const upstreamAdmin = await upstreamSuperAdmin() + + const upstreamKeys = await upstreamAdmin.get('/api/v1/api-keys?type=upload') + const upstreamKeyNames = new Set(upstreamKeys.data.results.map((k: any) => k.name)) + if (upstreamKeyNames.has('dev-upstream-upload') && output.keys['dev-upstream-upload']) { + console.log(' ✓ upstream api-key dev-upstream-upload (skipped)') + } else if (upstreamKeyNames.has('dev-upstream-upload')) { + throw new Error('upstream upload key exists but its raw value is lost — wipe the upstream db and re-run') + } else { + const res = await upstreamAdmin.post('/api/v1/api-keys', { type: 'upload', name: 'dev-upstream-upload' }) + output.keys['dev-upstream-upload'] = res.data.key + console.log(' + upstream api-key dev-upstream-upload') + await saveOutput(output) + } + const upstreamUpload = upstreamAxiosWithApiKey(output.keys['dev-upstream-upload']) + + const upstreamArtefactExists = async (id: string) => { + try { + await upstreamAdmin.get(`/api/v1/artefacts/${encodeURIComponent(id)}`) + return true + } catch (err) { + if (isHttp404(err)) return false + throw err + } + } + + // Ids deliberately disjoint from the local fixtures above: selecting an id that + // already exists locally without an `origin` is rejected with 409. + const UPSTREAM_NPM_ID = '@upstream/processing-remote@1' + const UPSTREAM_FILE_ID = 'upstream-terrain' + + if (await upstreamArtefactExists(UPSTREAM_NPM_ID)) { + console.log(` ✓ upstream npm ${UPSTREAM_NPM_ID} (skipped)`) + } else { + const tarball = await createTestTarball({ name: '@upstream/processing-remote', version: '1.0.0', licence: 'MIT' }) + const form = new FormData() + form.append('file', tarball, { filename: 'package.tgz', contentType: 'application/gzip' }) + form.append('category', 'processing') + await upstreamUpload.post(`/api/v1/artefacts/npm/${encodeURIComponent(UPSTREAM_NPM_ID)}`, form, { headers: form.getHeaders() }) + console.log(` + upstream npm ${UPSTREAM_NPM_ID}`) + } + + if (await upstreamArtefactExists(UPSTREAM_FILE_ID)) { + console.log(` ✓ upstream file ${UPSTREAM_FILE_ID} (skipped)`) + } else { + const form = new FormData() + form.append('file', Buffer.from('upstream-tileset-bytes'), { filename: 'upstream-terrain.mbtiles', contentType: 'application/octet-stream' }) + form.append('category', 'tileset') + await upstreamUpload.post(`/api/v1/artefacts/file/${encodeURIComponent(UPSTREAM_FILE_ID)}`, form, { headers: form.getHeaders() }) + console.log(` + upstream file ${UPSTREAM_FILE_ID}`) + } + + // Public so any read key may download them. + await upstreamAdmin.patch(`/api/v1/artefacts/${encodeURIComponent(UPSTREAM_NPM_ID)}`, { public: true }) + await upstreamAdmin.patch(`/api/v1/artefacts/${encodeURIComponent(UPSTREAM_FILE_ID)}`, { public: true }) + + // A read key requires an access grant for its owner account, minted by an admin of it. + try { + await upstreamAdmin.post('/api/v1/access-grants', { account: { type: 'organization', id: 'test1' } }) + console.log(' + upstream access-grant organization:test1') + } catch (err) { + if (!isHttp409(err)) throw err + console.log(' ✓ upstream access-grant organization:test1 (skipped)') + } + + // Consult the upstream itself, not just our output file — the api-test suite + // wipes the upstream between runs (cleanUpstream), so a raw key cached here can + // reference a key the upstream no longer has. Mirrors the upload-key logic above. + // + // The read-key LIST endpoint filters by the caller's own account, so a read key + // owned by org test1 is only visible through a test1-scoped session — never to a + // plain superadmin session. List and mint through the same org admin. + const orgAdmin = await upstreamAxiosAuth('test1-admin1', { org: 'test1' }) + const upstreamReadKeys = await orgAdmin.get('/api/v1/api-keys?type=read') + const hasFederationKey = upstreamReadKeys.data.results.some((k: any) => k.name === 'dev-federation') + if (hasFederationKey && output.keys['dev-upstream-read']) { + console.log(' ✓ upstream read-key dev-federation (skipped)') + } else { + const res = await orgAdmin.post('/api/v1/api-keys', { + type: 'read', + name: 'dev-federation', + owner: { type: 'organization', id: 'test1' } + }) + output.keys['dev-upstream-read'] = res.data.key + console.log(' + upstream read-key dev-federation (owner organization:test1)') + await saveOutput(output) + } + + // Downstream: register the mirror and select both artefacts. Deliberately does NOT + // sync — leave the "Sync now" button for a human to click. + try { + await admin.post('/api/v1/remote-registries', { + url: upstreamUrl, + name: 'Dev upstream', + apiKey: output.keys['dev-upstream-read'] + }) + console.log(` + remote-registry ${upstreamUrl}`) + } catch (err) { + if (!isHttp409(err)) throw err + console.log(` ✓ remote-registry ${upstreamUrl} (skipped)`) + } + + // Always refresh the stored key: on a re-run the read key above may have been + // re-minted, and the existing registry doc would otherwise keep ciphering the + // old, now-invalid key — the cause of a 401 when browsing the mirror. + await admin.patch(`/api/v1/remote-registries/${encodeURIComponent(upstreamUrl)}`, { + apiKey: output.keys['dev-upstream-read'] + }) + + for (const artefactId of [UPSTREAM_NPM_ID, UPSTREAM_FILE_ID]) { + try { + await admin.post(`/api/v1/remote-registries/${encodeURIComponent(upstreamUrl)}/selected-artefacts`, { artefactId }) + console.log(` + selected ${artefactId}`) + } catch (err) { + if (!isHttp409(err)) throw err + console.log(` ✓ selected ${artefactId} (skipped)`) + } + } + + console.log('\n → Sync is wired but not run. Click "Sync now" in the admin UI to exercise it.') + console.log(`\n✔ Fixtures applied. API keys written to ${OUTPUT_PATH}`) } diff --git a/dev/init-env.sh b/dev/init-env.sh index b3e95ac..c480bc2 100755 --- a/dev/init-env.sh +++ b/dev/init-env.sh @@ -11,6 +11,7 @@ DEV_UI_PORT=$((RANDOM_NB + 2)) DEV_UI_HMR_PORT=$((RANDOM_NB + 3)) MAILDEV_UI_PORT=$((RANDOM_NB + 4)) MAILDEV_SMTP_PORT=$((RANDOM_NB + 5)) +DEV_UPSTREAM_API_PORT=$((RANDOM_NB + 6)) MONGO_PORT=$((RANDOM_NB + 10)) diff --git a/dev/status.sh b/dev/status.sh index 91c22fb..a79e6f7 100755 --- a/dev/status.sh +++ b/dev/status.sh @@ -65,6 +65,13 @@ echo "" echo -e "${BOLD}Dev processes:${RESET}" check_http "dev-api" "$NGINX/registry/api/ping" check_http "dev-ui" "$NGINX/registry" +# The upstream registry has no nginx route by design — probe it directly. +# Guarded: this script runs under `set -u`, and older .env files lack the var. +if [ -n "${DEV_UPSTREAM_API_PORT:-}" ]; then + check_http "dev-api-upstream" "http://localhost:${DEV_UPSTREAM_API_PORT}/api/ping" +else + printf "%-20s MISSING DEV_UPSTREAM_API_PORT not set in .env\n" "dev-api-upstream" +fi echo "" # --- Docker compose services (probed through nginx where possible) --- diff --git a/docs/superpowers/plans/2026-07-09-federation-sync-ui.md b/docs/superpowers/plans/2026-07-09-federation-sync-ui.md new file mode 100644 index 0000000..804815d --- /dev/null +++ b/docs/superpowers/plans/2026-07-09-federation-sync-ui.md @@ -0,0 +1,1284 @@ +# Federation Sync Observable State Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Make remote-registry sync observable — the UI knows when a sync is running, how far it has got, and refuses a second click, without polling. + +**Architecture:** Three mechanisms, one job each. The Mongo `locks` collection answers *is it running* (cross-replica, self-healing via its 60s TTL). A `syncProgress` field on the registry doc answers *how far did it get* (survives the process, renders on page load). A WebSocket channel answers *what just changed* (push, no polling). `syncState` (`running` / `interrupted` / `idle`) is derived from the first two, never stored. + +**Tech Stack:** Express + MongoDB (api), Vue 3 + Vuetify (ui), `@data-fair/lib-node/locks.js`, `@data-fair/lib-node/ws-emitter.js`, `@data-fair/lib-express/ws-server.js`, `useWS` from `@data-fair/lib-vue`, Playwright test runner. + +**Spec:** `docs/superpowers/specs/2026-07-09-federation-sync-ui-design.md` + +## Global Constraints + +- Types are generated from JSON schemas. After editing `api/types/remote-registry/schema.js` you MUST run `npm run build-types`. Never hand-edit `api/types/remote-registry/.type/*`. +- After `build-types`, the running dev API may still execute stale generated code. Touch `api/index.ts` to force a nodemon reload. +- Never start, stop, or restart dev processes. Check them with `bash dev/status.sh`. Logs are in `dev/logs/`. +- Quality gates: `npm run lint-fix`, `npm run check-types`, `npm run test`. All three via `npm run quality`. +- Test projects are matched by filename: `*.unit.spec.ts`, `*.api.spec.ts`, `*.e2e.spec.ts`. +- **Running one file or one test.** `npm run test` is a compound `a && b && c` script, so `npm run test ` appends the file to the *last* command only (the e2e project) and runs unit and api in full. `AGENTS.md` documents this broken form. Use the per-project scripts instead: + - `npm run test-api -- tests/remote-registries.api.spec.ts` + - `npm run test-api -- tests/remote-registries.api.spec.ts -g "held lock rejects"` + - `npm run test-unit -- tests/remote-registries-operations.unit.spec.ts` +- `syncState` is **never persisted**. It is computed on read from the lock plus `syncProgress`/`lastSyncAt`. +- The registry `_id` **is a URL**. Every channel name, lock id, and route param built from it must be `encodeURIComponent`-ed at the boundary that needs it. +- No upgrade script. `syncProgress` is absent on existing docs and every derivation treats absent as "no attempt recorded". +- **Dependencies (landed before Task 3, commit `878eb8f`).** `ws@8` is a peerDependency of `@data-fair/lib-express` (needed by `ws-server.js`) and `reconnecting-websocket@4` is a peerDependency of `@data-fair/lib-vue` (needed by `ws.js`). Both were unmet. They are now in `api/package.json` and `ui/package.json` respectively. Neither is imported directly by our code, so no `@types/*` package is needed. + +## Deviation from the spec + +The spec's testing section says the emit call site is "unit-asserted (once per artefact, right channel)". Counting emits requires injecting the emitter into `runSync` — dependency injection added purely for a test, on a function whose real behaviour is already covered by the `syncProgress` assertions in Task 6. **This plan unit-tests `syncChannel()`'s encoding instead** (the part that can actually go wrong) and does not assert emit counts. Everything else follows the spec. + +## File Structure + +| File | Responsibility | +|---|---| +| `api/src/remote-registries/operations.ts` | **Modify.** Add pure helpers: `syncLockId`, `syncChannel`, `syncState`. No I/O. | +| `api/types/remote-registry/schema.js` | **Modify.** Add the `syncProgress` object. | +| `api/src/server.ts` | **Modify.** Boot `wsEmitter` + `wsServer`; stop `wsServer`. | +| `api/src/remote-registries/sync.ts` | **Modify.** Split `syncRemoteRegistry` into `runSync` (progress + emit), `syncRemoteRegistry` (awaits, for the daily job), `startSync` (returns on lock acquisition, for the route). | +| `api/src/remote-registries/router.ts` | **Modify.** Enrich reads with `syncState`; `409` on a held lock. | +| `api/src/app.ts` | **Modify.** Dev-only `test-env` lock endpoints; `clean()` drops `locks`. | +| `tests/support/axios.ts` | **Modify.** `holdSyncLock` / `releaseSyncLock` helpers. | +| `tests/remote-registries-operations.unit.spec.ts` | **Modify.** Cover `syncState` and `syncChannel`. | +| `tests/remote-registries.api.spec.ts` | **Modify.** Cover `syncState` on reads, `409`, and progress persistence. | +| `ui/src/composables/registry-sync.ts` | **Create.** Owns the socket subscription; folds events into the registry ref. | +| `ui/src/pages/admin/remote-registries/[id].vue` | **Modify.** Live progress, three states, disabled button, `409` handling. | +| `ui/src/components/admin/remote-registries-section.vue` | **Modify.** `syncState` chip per row. | + +--- + +### Task 1: Pure sync-state helpers + +The whole derivation lives here so it has one home, is trivially testable, and both the router and the sync module import it rather than re-deriving. + +**Files:** +- Modify: `api/src/remote-registries/operations.ts` +- Test: `tests/remote-registries-operations.unit.spec.ts` + +**Interfaces:** +- Consumes: nothing. +- Produces: + - `type SyncState = 'running' | 'interrupted' | 'idle'` + - `syncLockId(registryId: string): string` + - `syncChannel(registryId: string): string` + - `syncState(locked: boolean, registry: { syncProgress?: { startedAt: string }, lastSyncAt?: string }): SyncState` + +- [ ] **Step 1: Write the failing tests** + +Append to `tests/remote-registries-operations.unit.spec.ts`: + +```ts +import { filterSuggestedArtefacts, syncLockId, syncChannel, syncState } from '../api/src/remote-registries/operations.ts' + +test.describe('syncLockId', () => { + test('namespaces the registry url', () => { + expect(syncLockId('https://up.example.com/registry')).toBe('sync-remote-https://up.example.com/registry') + }) +}) + +test.describe('syncChannel', () => { + test('encodes the registry url so slashes do not split the channel', () => { + expect(syncChannel('https://up.example.com/registry')) + .toBe('remote-registries/https%3A%2F%2Fup.example.com%2Fregistry/sync') + }) + + test('the encoded channel has exactly three segments', () => { + expect(syncChannel('https://up.example.com/registry').split('/')).toHaveLength(3) + }) +}) + +test.describe('syncState', () => { + test('a held lock means running, whatever the doc says', () => { + expect(syncState(true, {})).toBe('running') + expect(syncState(true, { syncProgress: { startedAt: '2026-07-09T10:00:00.000Z' }, lastSyncAt: '2026-07-09T11:00:00.000Z' })).toBe('running') + }) + + test('no progress recorded means idle', () => { + expect(syncState(false, {})).toBe('idle') + expect(syncState(false, { lastSyncAt: '2026-07-09T11:00:00.000Z' })).toBe('idle') + }) + + test('an attempt that finished is idle', () => { + expect(syncState(false, { + syncProgress: { startedAt: '2026-07-09T10:00:00.000Z' }, + lastSyncAt: '2026-07-09T10:00:05.000Z' + })).toBe('idle') + }) + + test('an attempt stranded ahead of the last completed sync is interrupted', () => { + expect(syncState(false, { + syncProgress: { startedAt: '2026-07-09T12:00:00.000Z' }, + lastSyncAt: '2026-07-09T10:00:05.000Z' + })).toBe('interrupted') + }) + + test('a first-ever attempt that never finished is interrupted', () => { + expect(syncState(false, { syncProgress: { startedAt: '2026-07-09T12:00:00.000Z' } })).toBe('interrupted') + }) +}) +``` + +Note the import line **replaces** the existing `import { filterSuggestedArtefacts } ...` line at the top of the file. Leave the existing `filterSuggestedArtefacts` describe block untouched. + +- [ ] **Step 2: Run the tests to verify they fail** + +Run: `npm run test-unit -- tests/remote-registries-operations.unit.spec.ts` +Expected: FAIL — `syncLockId is not a function` (or a TypeScript "has no exported member" error). + +- [ ] **Step 3: Write the implementation** + +Append to `api/src/remote-registries/operations.ts`: + +```ts +export type SyncState = 'running' | 'interrupted' | 'idle' + +// The lock row id in the shared `locks` collection. Holding it means a sync is in flight. +export const syncLockId = (registryId: string) => `sync-remote-${registryId}` + +// The ws channel a registry's sync progress is published on. The registry id IS a url, +// so it must be encoded — a raw `/` would shred this `/`-delimited channel name. +export const syncChannel = (registryId: string) => + `remote-registries/${encodeURIComponent(registryId)}/sync` + +// Running state is derived, never stored: a stored `running` flag becomes a lie the moment +// a process is killed. `interrupted` means the final write that sets lastSyncAt never +// happened, stranding the attempt's startedAt ahead of it. +export const syncState = ( + locked: boolean, + registry: { syncProgress?: { startedAt: string }, lastSyncAt?: string } +): SyncState => { + if (locked) return 'running' + const startedAt = registry.syncProgress?.startedAt + if (startedAt && (!registry.lastSyncAt || startedAt > registry.lastSyncAt)) return 'interrupted' + return 'idle' +} +``` + +- [ ] **Step 4: Run the tests to verify they pass** + +Run: `npm run test-unit -- tests/remote-registries-operations.unit.spec.ts` +Expected: PASS — 11 tests (3 pre-existing + 8 new). + +- [ ] **Step 5: Commit** + +```bash +git add api/src/remote-registries/operations.ts tests/remote-registries-operations.unit.spec.ts +git commit -m "feat(api): pure sync-state derivation helpers" +``` + +--- + +### Task 2: `syncProgress` on the registry schema + +**Files:** +- Modify: `api/types/remote-registry/schema.js` +- Regenerates: `api/types/remote-registry/.type/index.d.ts` + +**Interfaces:** +- Consumes: nothing. +- Produces: `RemoteRegistry['syncProgress']?: { startedAt: string, done: number, total: number, currentArtefact?: string }` + +- [ ] **Step 1: Add the field to the schema** + +In `api/types/remote-registry/schema.js`, insert after the `lastSyncError` property and before `createdAt`: + +```js + lastSyncError: { type: 'string' }, + syncProgress: { + type: 'object', + additionalProperties: false, + required: ['startedAt', 'done', 'total'], + properties: { + startedAt: { type: 'string', format: 'date-time' }, + done: { type: 'integer' }, + total: { type: 'integer' }, + currentArtefact: { type: 'string' } + } + }, + createdAt: { type: 'string', format: 'date-time', readOnly: true }, +``` + +Do **not** add `syncProgress` to the root `required` array — it is absent on every existing doc. + +- [ ] **Step 2: Regenerate types** + +Run: `npm run build-types` +Expected: exits 0. + +- [ ] **Step 3: Verify the generated type carries the field** + +Run: `grep -A6 'syncProgress' api/types/remote-registry/.type/index.d.ts` +Expected: shows `syncProgress?:` with `startedAt`, `done`, `total`, `currentArtefact`. + +- [ ] **Step 4: Force the dev API to reload the generated code** + +Run: `touch api/index.ts` + +This is not optional — nodemon can otherwise keep executing the pre-generation module. + +- [ ] **Step 5: Type-check** + +Run: `npm run check-types` +Expected: exits 0. + +- [ ] **Step 6: Commit** + +```bash +git add api/types/remote-registry/ +git commit -m "feat(api): add syncProgress to the remote-registry schema" +``` + +--- + +### Task 3: Boot the websocket server + +No behaviour is observable yet — this task's deliverable is that the socket accepts an authenticated admin connection and nothing regresses. It is separated from Task 6 because a reviewer could reasonably reject the boot wiring while approving the emit call sites. + +**Files:** +- Modify: `api/src/server.ts` + +**Interfaces:** +- Consumes: nothing. +- Produces: a live `WebSocketServer` on the API's http server; `wsEmitter.emit(channel, data)` becomes usable process-wide. + +- [ ] **Step 1: Add the imports** + +In `api/src/server.ts`, after the existing `import locks from '@data-fair/lib-node/locks.js'`: + +```ts +import * as wsEmitter from '@data-fair/lib-node/ws-emitter.js' +import * as wsServer from '@data-fair/lib-express/ws-server.js' +``` + +- [ ] **Step 2: Start them** + +In `start()`, immediately after `await locks.start(mongo.db)`: + +```ts + await locks.start(mongo.db) + await wsEmitter.init(mongo.db) + // `canSubscribe` is never reached for admins: ws-server short-circuits on + // sessionState.user?.adminMode before calling it. Remote-registry sync is an + // admin-only surface, so returning false refuses precisely everyone else. + await wsServer.start(server, mongo.db, async () => false) +``` + +- [ ] **Step 3: Stop it** + +In `stop()`, between `httpTerminator.terminate()` and `locks.stop()`: + +```ts + await httpTerminator.terminate() + await wsServer.stop() + if (config.observer?.active) await stopObserver() + await locks.stop() +``` + +- [ ] **Step 4: Verify no regression and that the socket is up** + +Run: `npm run check-types && npm run test-api -- tests/ping.api.spec.ts` +Expected: both exit 0. + +Then confirm the server upgraded a websocket rather than 404ing it: + +```bash +set -a && . ./.env && set +a +curl -s -o /dev/null -w '%{http_code}' \ + -H 'Connection: Upgrade' -H 'Upgrade: websocket' \ + -H 'Sec-WebSocket-Version: 13' -H 'Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==' \ + "http://localhost:${DEV_API_PORT}/api/" +``` +Expected: `101`. A `404` or `426` means `wsServer.start` did not attach to the http server. (`wsServer` attaches to the whole http server, not a path, so any path upgrades — `/api/` is just convenient.) + +- [ ] **Step 5: Commit** + +```bash +git add api/src/server.ts +git commit -m "feat(api): boot the ws pub/sub server and emitter" +``` + +--- + +### Task 4: Reads expose `syncState` + +The dev-only lock endpoints are folded in here because this task's tests are the first thing that needs them: a real sync finishes far too fast to observe `running` by racing it. We hold the lock directly instead. + +**Files:** +- Modify: `api/src/app.ts` (dev-only `test-env` routes; `clean()` drops `locks`) +- Modify: `api/src/remote-registries/router.ts` +- Modify: `tests/support/axios.ts` +- Test: `tests/remote-registries.api.spec.ts` + +**Interfaces:** +- Consumes: `syncLockId`, `syncState` (Task 1). +- Produces: + - `GET /api/v1/remote-registries` → each result carries `syncState: SyncState` + - `GET /api/v1/remote-registries/:id` → carries `syncState: SyncState` + - `holdSyncLock(registryId: string): Promise` and `releaseSyncLock(registryId: string): Promise` in `tests/support/axios.ts` + +- [ ] **Step 1: Add the dev-only lock endpoints** + +In `api/src/app.ts`, inside the existing `if (process.env.NODE_ENV === 'development') {` block, add the `locks` wipe to the existing `DELETE /api/test-env` handler (after `await mongo.remoteRegistries.deleteMany({})`): + +```ts + await mongo.remoteRegistries.deleteMany({}) + await mongo.db.collection('locks').deleteMany({}) +``` + +Then add two new routes inside the same block, after the existing `PUT /api/test-env/artefacts/:id/origin` handler: + +```ts + // Hold a lock under a FOREIGN pid so the API process can neither release nor prolong it. + // Lets a test observe `running` deterministically instead of racing a real sync. + // The row still expires on its own via the locks TTL index if a test forgets to clean up. + app.put('/api/test-env/locks/:id', async (req, res) => { + assertReqInternal(req) + const now = new Date() + await mongo.db.collection('locks').insertOne({ + _id: req.params.id as any, + pid: 'test-env', + hostname: 'test-env', + createdAt: now, + updatedAt: now + }) + res.send() + }) + + app.delete('/api/test-env/locks/:id', async (req, res) => { + assertReqInternal(req) + await mongo.db.collection('locks').deleteOne({ _id: req.params.id as any }) + res.send() + }) +``` + +- [ ] **Step 2: Add the test helpers** + +Append to `tests/support/axios.ts`: + +```ts +const testEnvUrl = `http://localhost:${process.env.DEV_API_PORT}/api/test-env` + +// The lock id embeds the registry url, so it must be encoded as a single path segment. +const syncLockPath = (registryId: string) => `${testEnvUrl}/locks/${encodeURIComponent('sync-remote-' + registryId)}` + +export const holdSyncLock = async (registryId: string) => { + await anonymousAx.put(syncLockPath(registryId), {}) +} + +export const releaseSyncLock = async (registryId: string) => { + await anonymousAx.delete(syncLockPath(registryId)) +} +``` + +- [ ] **Step 3: Write the failing tests** + +Append a new describe block to `tests/remote-registries.api.spec.ts`, and add `holdSyncLock, releaseSyncLock` to the existing import from `./support/axios.ts`: + +```ts + test.describe('Sync state on reads', () => { + const url = 'https://upstream.example.com' + + test.beforeEach(async () => { + const ax = await superAdmin + await ax.post('/api/v1/remote-registries', { url, name: 'Upstream', apiKey: 'reg_abc_secretkey123' }) + }) + + test('a registry with no attempt recorded is idle', async () => { + const ax = await superAdmin + const res = await ax.get('/api/v1/remote-registries/' + encodeURIComponent(url)) + expect(res.data.syncState).toBe('idle') + }) + + test('a held lock reports running on the detail endpoint', async () => { + const ax = await superAdmin + await holdSyncLock(url) + try { + const res = await ax.get('/api/v1/remote-registries/' + encodeURIComponent(url)) + expect(res.data.syncState).toBe('running') + } finally { + await releaseSyncLock(url) + } + const after = await ax.get('/api/v1/remote-registries/' + encodeURIComponent(url)) + expect(after.data.syncState).toBe('idle') + }) + + test('a held lock reports running on the list endpoint', async () => { + const ax = await superAdmin + await holdSyncLock(url) + try { + const res = await ax.get('/api/v1/remote-registries') + expect(res.data.results.find((r: any) => r._id === url).syncState).toBe('running') + } finally { + await releaseSyncLock(url) + } + }) + + test('one registry running does not mark its siblings as running', async () => { + const ax = await superAdmin + const other = 'https://other.example.com' + await ax.post('/api/v1/remote-registries', { url: other, name: 'Other', apiKey: 'reg_def_otherkey' }) + await holdSyncLock(url) + try { + const res = await ax.get('/api/v1/remote-registries') + const byId = Object.fromEntries(res.data.results.map((r: any) => [r._id, r.syncState])) + expect(byId[url]).toBe('running') + expect(byId[other]).toBe('idle') + } finally { + await releaseSyncLock(url) + } + }) + }) +``` + +- [ ] **Step 4: Run the tests to verify they fail** + +Run: `npm run test-api -- tests/remote-registries.api.spec.ts` +Expected: FAIL — `expected 'idle', received undefined` (the field does not exist yet). + +- [ ] **Step 5: Enrich the read endpoints** + +In `api/src/remote-registries/router.ts`, extend the existing import from `./operations.ts`: + +```ts +import { filterSuggestedArtefacts, syncLockId, syncState } from './operations.ts' +``` + +Add this helper below `extractShortId`: + +```ts +// One query for the whole page, not one per row. +const lockedLockIds = async (registryIds: string[]): Promise> => { + if (registryIds.length === 0) return new Set() + const rows = await mongo.db.collection('locks') + .find({ _id: { $in: registryIds.map(syncLockId) as any } }, { projection: { _id: 1 } }) + .toArray() + return new Set(rows.map(row => String(row._id))) +} +``` + +Replace the body of `GET /`: + +```ts +router.get('/', async (req, res, next) => { + try { + await session.reqAdminMode(req) + const results = await mongo.remoteRegistries.find({}, { projection: { apiKey: 0 } }).toArray() + const locked = await lockedLockIds(results.map(r => r._id)) + res.json({ + results: results.map(r => ({ ...r, syncState: syncState(locked.has(syncLockId(r._id)), r) })), + count: results.length + }) + } catch (err) { next(err) } +}) +``` + +Replace the body of `GET /:id`: + +```ts +router.get('/:id', async (req, res, next) => { + try { + await session.reqAdminMode(req) + const doc = await mongo.remoteRegistries.findOne({ _id: req.params.id }, { projection: { apiKey: 0 } }) + if (!doc) throw httpError(404, 'remote registry not found') + const locked = await lockedLockIds([doc._id]) + res.json({ ...doc, syncState: syncState(locked.has(syncLockId(doc._id)), doc) }) + } catch (err) { next(err) } +}) +``` + +- [ ] **Step 6: Run the tests to verify they pass** + +Run: `npm run test-api -- tests/remote-registries.api.spec.ts` +Expected: PASS — all pre-existing tests plus the 4 new ones. + +- [ ] **Step 7: Commit** + +```bash +git add api/src/app.ts api/src/remote-registries/router.ts tests/support/axios.ts tests/remote-registries.api.spec.ts +git commit -m "feat(api): expose derived syncState on remote-registry reads" +``` + +--- + +### Task 5: `409` on a concurrent manual sync + +**Files:** +- Modify: `api/src/remote-registries/sync.ts` +- Modify: `api/src/remote-registries/router.ts:188-199` +- Test: `tests/remote-registries.api.spec.ts` + +**Interfaces:** +- Consumes: `syncLockId` (Task 1); `holdSyncLock`/`releaseSyncLock` (Task 4). +- Produces: + - `startSync(registryId: string): Promise` — returns as soon as the lock is taken; `false` if already held. + - `syncRemoteRegistry(registryId: string): Promise` — awaits completion; `false` if already held. Used by the daily job. + - `runSync(registryId: string): Promise` — module-private, the actual work. + +- [ ] **Step 1: Write the failing tests** + +Append to the `Sync state on reads` describe block's parent (a new describe block) in `tests/remote-registries.api.spec.ts`: + +```ts + test.describe('Manual sync trigger', () => { + const url = 'https://upstream.example.com' + + test.beforeEach(async () => { + const ax = await superAdmin + await ax.post('/api/v1/remote-registries', { url, name: 'Upstream', apiKey: 'reg_abc_secretkey123' }) + }) + + test('a free lock accepts the sync with 202', async () => { + const ax = await superAdmin + const res = await ax.post('/api/v1/remote-registries/' + encodeURIComponent(url) + '/sync') + expect(res.status).toBe(202) + }) + + test('a held lock rejects the sync with 409', async () => { + const ax = await superAdmin + await holdSyncLock(url) + try { + await ax.post('/api/v1/remote-registries/' + encodeURIComponent(url) + '/sync') + expect(true).toBe(false) + } catch (err: any) { + expect(err.status).toBe(409) + } finally { + await releaseSyncLock(url) + } + }) + + // Asserting only `404` here would pass under a wrong-order (lock-then-404) + // implementation too. Pre-holding the lock makes the orders diverge: correct + // order 404s before ever touching the lock; inverted order fails to acquire + // and returns 409. Checking the lock row *after* a plain 404 does NOT work — + // the background release wins the race against the follow-up request. + test('an unknown registry is still 404, not 409, even while a lock is held', async () => { + const ax = await superAdmin + const unknown = 'https://nope.example.com' + await holdSyncLock(unknown) + try { + await ax.post('/api/v1/remote-registries/' + encodeURIComponent(unknown) + '/sync') + expect(true).toBe(false) + } catch (err: any) { + expect(err.status).toBe(404) + } finally { + await releaseSyncLock(unknown) + } + }) + }) +``` + +- [ ] **Step 2: Run the tests to verify they fail** + +Run: `npm run test-api -- tests/remote-registries.api.spec.ts -g "held lock rejects"` +Expected: FAIL — the request returns `202`, so `expect(true).toBe(false)` trips. + +- [ ] **Step 3: Split the sync entry points** + +In `api/src/remote-registries/sync.ts`, add the imports: + +```ts +import { internalError } from '@data-fair/lib-node/observer.js' +import { syncLockId } from './operations.ts' +``` + +Rename the existing exported `syncRemoteRegistry` to a private `runSync` and strip its locking. It becomes exactly the old body **minus** the `locks.acquire` guard and the `try`/`finally` release — the lock is now the caller's responsibility: + +```ts +// The actual work. Callers own the lock. +const runSync = async (remoteRegistryId: string) => { + const remote = await mongo.remoteRegistries.findOne({ _id: remoteRegistryId }) + if (!remote) return + + const apiKey = decipher(remote.apiKey) + const ax = axiosBuilder({ + baseURL: remote._id, + headers: { 'x-api-key': apiKey } + }) + + let hasErrors = false + let lastError = '' + + for (const artefactId of remote.selectedArtefacts) { + try { + const encodedId = encodeURIComponent(artefactId) + const detailRes = await ax.get(`/api/v1/artefacts/${encodedId}`) + const format: Artefact['format'] = detailRes.data.format + + if (format === 'npm') { + await syncNpmArtefact(ax, remote._id, artefactId) + } else { + await syncFileArtefact(ax, remote._id, artefactId) + } + } catch (err: any) { + hasErrors = true + lastError = `${artefactId}: ${err.message || err}` + console.error(`[sync] Error syncing ${artefactId} from ${remote._id}:`, err.message || err) + } + } + + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { + $set: { + lastSyncAt: new Date().toISOString(), + lastSyncStatus: hasErrors ? 'error' : 'success', + ...(hasErrors ? { lastSyncError: lastError } : {}) + }, + ...(!hasErrors ? { $unset: { lastSyncError: '' } } : {}) + } + ) +} + +// Returns as soon as the lock is taken; the work continues in the background. +// A held lock is a conflict the caller (a human clicking a button) should see. +// +// Both the work AND the release are guarded. `.finally(cb)` propagates a rejection +// thrown by `cb` into a promise nobody consumes here — an unguarded `locks.release` +// failure would surface as an unhandled rejection and kill the process. +export const startSync = async (remoteRegistryId: string): Promise => { + const lockId = syncLockId(remoteRegistryId) + if (!await locks.acquire(lockId)) return false + runSync(remoteRegistryId) + .catch(err => internalError('sync-remote-registry', err)) + .finally(() => locks.release(lockId).catch(err => internalError('sync-remote-registry-release', err))) + return true +} + +// Awaits completion. Used by the daily job, which syncs registries one at a time. +export const syncRemoteRegistry = async (remoteRegistryId: string): Promise => { + const lockId = syncLockId(remoteRegistryId) + if (!await locks.acquire(lockId)) return false + try { + await runSync(remoteRegistryId) + } finally { + await locks.release(lockId) + } + return true +} + +export const syncAllRemoteRegistries = async () => { + const remotes = await mongo.remoteRegistries.find({}).toArray() + for (const remote of remotes) { + // A held lock means a peer replica is already syncing this registry. That is the + // normal outcome of N replicas firing the same daily timer — not an error. + await syncRemoteRegistry(remote._id).catch(err => { + console.error(`[sync] Failed to sync ${remote._id}:`, err.message || err) + }) + } +} +``` + +Note the old `console.log('[sync] Lock already held...')` line disappears: the boolean return replaces it. + +- [ ] **Step 4: Make the route honour the lock** + +In `api/src/remote-registries/router.ts`, change the import `import { syncRemoteRegistry } from './sync.ts'` to `import { startSync } from './sync.ts'`, and replace the `POST /:id/sync` handler body: + +```ts +router.post('/:id/sync', async (req, res, next) => { + try { + await session.reqAdminMode(req) + const doc = await mongo.remoteRegistries.findOne({ _id: req.params.id }) + if (!doc) throw httpError(404, 'remote registry not found') + + if (!await startSync(req.params.id)) throw httpError(409, 'sync already running') + res.status(202).json({ message: 'sync started' }) + } catch (err) { next(err) } +}) +``` + +The 404 lookup stays **before** the lock attempt, so an unknown registry never acquires a lock it would then have to release. + +- [ ] **Step 5: Run the tests to verify they pass** + +Run: `npm run test-api -- tests/remote-registries.api.spec.ts` +Expected: PASS. + +- [ ] **Step 6: Commit** + +```bash +git add api/src/remote-registries/sync.ts api/src/remote-registries/router.ts tests/remote-registries.api.spec.ts +git commit -m "feat(api): reject a concurrent manual sync with 409" +``` + +--- + +### Task 6: Persist and publish sync progress + +**Files:** +- Modify: `api/src/remote-registries/sync.ts` +- Test: `tests/remote-registries.api.spec.ts` + +**Interfaces:** +- Consumes: `syncChannel` (Task 1); `syncProgress` type (Task 2); `wsEmitter.init` (Task 3); `runSync` (Task 5). +- Produces: `syncProgress` written on the doc at start, per artefact, and at the end; a `SyncEvent` published on `syncChannel(id)` at each of those points. + +- [ ] **Step 1: Write the failing test** + +Append to the `Manual sync trigger` describe block in `tests/remote-registries.api.spec.ts`: + +```ts + test('a sync over zero selected artefacts records progress and succeeds', async () => { + const ax = await superAdmin + await ax.post('/api/v1/remote-registries/' + encodeURIComponent(url) + '/sync') + + // the sync runs in the background; poll the doc until it settles (test-only wait) + let doc: any + for (let i = 0; i < 50; i++) { + const res = await ax.get('/api/v1/remote-registries/' + encodeURIComponent(url)) + doc = res.data + if (doc.lastSyncStatus) break + await new Promise(resolve => setTimeout(resolve, 100)) + } + + expect(doc.lastSyncStatus).toBe('success') + expect(doc.syncProgress.total).toBe(0) + expect(doc.syncProgress.done).toBe(0) + expect(doc.syncProgress.currentArtefact).toBeUndefined() + expect(doc.syncProgress.startedAt <= doc.lastSyncAt).toBe(true) + expect(doc.syncState).toBe('idle') + }) +``` + +The poll loop here is a *test* waiting on a background job, not the production polling we removed. + +- [ ] **Step 2: Run the test to verify it fails** + +Run: `npm run test-api -- tests/remote-registries.api.spec.ts -g "zero selected artefacts"` +Expected: FAIL — `Cannot read properties of undefined (reading 'total')`; `syncProgress` is never written. + +- [ ] **Step 3: Add the emitter and progress writes** + +In `api/src/remote-registries/sync.ts`, add the import and the event type: + +```ts +import * as wsEmitter from '@data-fair/lib-node/ws-emitter.js' +import { syncLockId, syncChannel } from './operations.ts' + +export type SyncEvent = { + running: boolean + startedAt: string + done: number + total: number + currentArtefact?: string + lastSyncAt?: string + lastSyncStatus?: 'success' | 'error' + lastSyncError?: string +} + +// A dropped progress frame is cosmetic — the next frame supersedes it — so an emit +// failure must never abort a sync. +const emitSync = async (remoteRegistryId: string, event: SyncEvent) => { + try { + await wsEmitter.emit(syncChannel(remoteRegistryId), event) + } catch (err) { + internalError('sync-ws-emit', err) + } +} +``` + +Then rewrite `runSync` to bracket the loop with progress writes: + +```ts +const runSync = async (remoteRegistryId: string) => { + const remote = await mongo.remoteRegistries.findOne({ _id: remoteRegistryId }) + if (!remote) return + + const startedAt = new Date().toISOString() + const total = remote.selectedArtefacts.length + let done = 0 + + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { $set: { syncProgress: { startedAt, done, total } } } + ) + await emitSync(remoteRegistryId, { running: true, startedAt, done, total }) + + const apiKey = decipher(remote.apiKey) + const ax = axiosBuilder({ + baseURL: remote._id, + headers: { 'x-api-key': apiKey } + }) + + let hasErrors = false + let lastError = '' + + for (const artefactId of remote.selectedArtefacts) { + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { $set: { 'syncProgress.currentArtefact': artefactId } } + ) + await emitSync(remoteRegistryId, { running: true, startedAt, done, total, currentArtefact: artefactId }) + + try { + const encodedId = encodeURIComponent(artefactId) + const detailRes = await ax.get(`/api/v1/artefacts/${encodedId}`) + const format: Artefact['format'] = detailRes.data.format + + if (format === 'npm') { + await syncNpmArtefact(ax, remote._id, artefactId) + } else { + await syncFileArtefact(ax, remote._id, artefactId) + } + } catch (err: any) { + hasErrors = true + lastError = `${artefactId}: ${err.message || err}` + console.error(`[sync] Error syncing ${artefactId} from ${remote._id}:`, err.message || err) + } + + done++ + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { $set: { 'syncProgress.done': done } } + ) + await emitSync(remoteRegistryId, { running: true, startedAt, done, total, currentArtefact: artefactId }) + } + + const lastSyncAt = new Date().toISOString() + const lastSyncStatus = hasErrors ? 'error' as const : 'success' as const + + await mongo.remoteRegistries.updateOne( + { _id: remoteRegistryId }, + { + // NB: `syncProgress.done` is deliberately NOT re-written here. The loop's + // per-artefact write is its only writer when total > 0 — a redundant write + // here would mask a regression in the loop, making any test of `done` vacuous. + $set: { + lastSyncAt, + lastSyncStatus, + ...(hasErrors ? { lastSyncError: lastError } : {}) + }, + $unset: { + 'syncProgress.currentArtefact': '', + ...(hasErrors ? {} : { lastSyncError: '' }) + } + } + ) + + // The end event carries the terminal state, so the UI never refetches to learn the outcome. + await emitSync(remoteRegistryId, { + running: false, + startedAt, + done, + total, + lastSyncAt, + lastSyncStatus, + ...(hasErrors ? { lastSyncError: lastError } : {}) + }) +} +``` + +`syncProgress` is deliberately **not** unset at the end: `syncState` needs `startedAt` to sit behind `lastSyncAt` to conclude `idle`, and a stranded `startedAt` is exactly what makes a crashed run render as `interrupted`. + +- [ ] **Step 4: Run the tests to verify they pass** + +Run: `npm run test-api -- tests/remote-registries.api.spec.ts` +Expected: PASS. + +- [ ] **Step 5: Verify a crashed run renders as interrupted** + +There is no automated test for this (it requires killing the API mid-sync). Assert the derivation by hand against the shipped helper: + +Run: `npm run test-unit -- tests/remote-registries-operations.unit.spec.ts -g "stranded ahead"` +Expected: PASS. This is the same function the router calls. + +- [ ] **Step 6: Commit** + +```bash +git add api/src/remote-registries/sync.ts tests/remote-registries.api.spec.ts +git commit -m "feat(api): persist and publish sync progress" +``` + +--- + +### Task 7: UI sync composable + +**Files:** +- Create: `ui/src/composables/registry-sync.ts` + +**Interfaces:** +- Consumes: `syncChannel`'s wire format (Task 1); `SyncEvent` (Task 6). +- Produces: `useRegistrySync(registryId: string, registry: Ref): void` — subscribes and folds events into `registry.value`. + +- [ ] **Step 1: Create the composable** + +```ts +import { type Ref } from 'vue' +import useWS from '@data-fair/lib-vue/ws.js' +import { $apiPath } from '~/context' + +// Mirrors SyncEvent in api/src/remote-registries/sync.ts +export type SyncEvent = { + running: boolean + startedAt: string + done: number + total: number + currentArtefact?: string + lastSyncAt?: string + lastSyncStatus?: 'success' | 'error' + lastSyncError?: string +} + +// Subscribes to a registry's sync channel and folds each event into the registry ref. +// `subscribe` registers its own onScopeDispose teardown, so callers need no onUnmounted. +// useWS returns undefined when the browser has no WebSocket: the page then renders correct +// state at load and simply does not animate. +export const useRegistrySync = (registryId: string, registry: Ref) => { + const ws = useWS($apiPath + '/') + ws?.subscribe(`remote-registries/${encodeURIComponent(registryId)}/sync`, (event) => { + const reg = registry.value + // an event can land before the initial fetch resolves; the next one supersedes it + if (!reg) return + + reg.syncProgress = { + startedAt: event.startedAt, + done: event.done, + total: event.total, + currentArtefact: event.currentArtefact + } + + if (event.running) { + reg.syncState = 'running' + return + } + + reg.syncState = 'idle' + reg.lastSyncAt = event.lastSyncAt + reg.lastSyncStatus = event.lastSyncStatus + reg.lastSyncError = event.lastSyncError + }) +} +``` + +Imports are explicit rather than relying on auto-import: a newly created file's own auto-imports are not picked up by an already-running vite dev server. + +- [ ] **Step 2: Type-check** + +Run: `npm run check-types` +Expected: exits 0. + +- [ ] **Step 3: Lint** + +Run: `npm run lint-fix` +Expected: exits 0, no remaining errors. + +- [ ] **Step 4: Commit** + +```bash +git add ui/src/composables/registry-sync.ts +git commit -m "feat(ui): composable subscribing to a registry's sync channel" +``` + +--- + +### Task 8: Live sync panel on the detail page + +**Files:** +- Modify: `ui/src/pages/admin/remote-registries/[id].vue` + +**Interfaces:** +- Consumes: `useRegistrySync` (Task 7); `syncState` / `syncProgress` on `GET /:id` (Tasks 4, 6); `409` from `POST /:id/sync` (Task 5). +- Produces: nothing downstream. + +- [ ] **Step 1: Wire the composable and the derived flags** + +In the `