Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,9 @@ Test users are defined in @dev/resources/users.json and organizations in @dev/re

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.
`npm run dev:fixtures` seeds the upstream, mints a read key owned by org `test1`, registers the mirror and selects two artefacts. Selecting an artefact queues a background sync of that artefact alone (`pendingSync` on the registry doc, drained under the sync lock), so the mirrors land without clicking **Sync now** — that button runs a full sync. Sync progress is published on the registry's ws channel; the admin table gets each row's local state from `GET …/remote-artefacts` and loads upstream thumbnails through `GET …/remote-thumbnails/:id/data` (the browser cannot reach the upstream directly).

Unselecting keeps the local copy but unlocks it (drops `origin`); it then blocks re-mirroring of that id (409) until it is deleted. Deleting a registry unlocks all its mirrors the same way.

`tests/remote-registries-sync.api.spec.ts` covers the mirror path end to end against the same upstream.

Expand Down
39 changes: 31 additions & 8 deletions api/src/artefacts/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -96,21 +96,44 @@ const tryInternalSecret = (req: import('express').Request): boolean => {
return timingSafeEqual(received, expected)
}

// `?sort=<key>` with an optional leading `-` for descending. Keys map to the
// Mongo sort they stand for; a name tie-breaker keeps pages stable when many
// docs share a value. `vulnerabilities` and `public` are admin-only: scan data
// is stripped for non-admins (ordering by it would leak) and visibility is not
// shown to them — they silently get the default sort.
const sortKeys: Record<string, { sort: Record<string, 1 | -1>, admin?: boolean }> = {
dataUpdatedAt: { sort: { dataUpdatedAt: -1 } },
name: { sort: { name: 1 } },
category: { sort: { category: 1 } },
'group.en': { sort: { 'group.en': 1 } },
'group.fr': { sort: { 'group.fr': 1 } },
version: { sort: { version: 1 } },
size: { sort: { size: 1 } },
public: { sort: { public: 1 }, admin: true },
vulnerabilities: { sort: { 'scan.summary.critical': -1, 'scan.summary.high': -1, 'scan.summary.medium': -1 }, admin: true }
}
const parseSort = (raw: unknown, caller: Caller): Record<string, 1 | -1> => {
const fallback = { dataUpdatedAt: -1 as const, name: 1 as const }
if (typeof raw !== 'string' || !raw) return fallback
const desc = raw.startsWith('-')
const key = desc ? raw.slice(1) : raw
const def = sortKeys[key]
if (!def) throw httpError(400, `invalid sort, must be one of: ${Object.keys(sortKeys).join(', ')} (prefix with - for descending)`)
if (def.admin && !caller.admin) return fallback
const sort: Record<string, 1 | -1> = {}
for (const [field, dir] of Object.entries(def.sort)) sort[field] = desc ? (dir === 1 ? -1 : 1) : dir
if (!('name' in sort)) sort.name = 1
return sort
}

// List artefacts (filtered by access)
router.get('/', async (req, res, next) => {
try {
const caller = await resolveCaller(req)
const filter = artefactAccessFilter(caller)
const skip = Math.max(0, Math.min(parseInt(req.query.skip as string) || 0, 100000))
const size = Math.min(parseInt(req.query.size as string) || 10, 100)
// `vulnerabilities` is admin-only: scan data is stripped for non-admins, so
// ordering by it would leak. Non-admins silently get the default sort.
let sort: Record<string, 1 | -1> = { dataUpdatedAt: -1 }
if (req.query.sort === 'name') {
sort = { name: 1 }
} else if (req.query.sort === 'vulnerabilities' && caller.admin) {
sort = { 'scan.summary.critical': -1, 'scan.summary.high': -1, 'scan.summary.medium': -1, dataUpdatedAt: -1 }
}
const sort = parseSort(req.query.sort, caller)

// Text search on name
if (req.query.q) {
Expand Down
30 changes: 30 additions & 0 deletions api/src/remote-registries/operations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,3 +38,33 @@ export const syncState = (
if (startedAt && (!registry.lastSyncAt || startedAt > registry.lastSyncAt)) return 'interrupted'
return 'idle'
}

// Local mirror state of each remote artefact, for the admin's selection table.
// `upToDate` compares the upstream dataUpdatedAt with the mirrored copy's — the
// same fast-path key the npm sync uses, so a stale row is exactly one the next
// sync would re-download.
// `conflict` flags a local artefact with the same id that this registry does
// not own (uploaded locally, or unlocked when a registry was deleted): the
// select endpoint would answer 409, so the UI can say so up front.
export type LocalState = { synced: boolean, upToDate: boolean, dataUpdatedAt?: string, conflict?: boolean }

export const annotateLocalState = <T extends { _id: string, dataUpdatedAt?: string }> (
remoteArtefacts: T[],
localArtefacts: { _id: string, dataUpdatedAt?: string, origin?: string }[],
registryId: string
): (T & { local: LocalState })[] => {
const locals = new Map(localArtefacts.map(a => [a._id, a]))
return remoteArtefacts.map(remote => {
const local = locals.get(remote._id)
if (!local) return { ...remote, local: { synced: false, upToDate: false } }
if (local.origin !== registryId) return { ...remote, local: { synced: false, upToDate: false, conflict: true } }
return {
...remote,
local: {
synced: true,
upToDate: !!local.dataUpdatedAt && local.dataUpdatedAt === remote.dataUpdatedAt,
dataUpdatedAt: local.dataUpdatedAt
}
}
})
}
56 changes: 51 additions & 5 deletions api/src/remote-registries/router.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
import { Router } from 'express'
import { pipeline } from 'node:stream/promises'
import { session } from '@data-fair/lib-express/index.js'
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 { startSync } from './sync.ts'
import { filterSuggestedArtefacts, syncLockId, syncState } from './operations.ts'
import { startSync, enqueueArtefactSync } from './sync.ts'
import { filterSuggestedArtefacts, annotateLocalState, 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'

Expand Down Expand Up @@ -135,10 +136,52 @@ router.get('/:id/remote-artefacts', async (req, res, next) => {
// not already selected — a deprecated artefact is not suggested for new
// mirroring but stays visible if it is already mirrored.
const params: Record<string, string> = { size: String(size), skip: String(skip), includeDeprecated: 'true' }
if (req.query.q) params.q = req.query.q as string
for (const key of ['q', 'category', 'format'] as const) {
if (typeof req.query[key] === 'string' && req.query[key]) params[key] = req.query[key] as string
}

const remote = await ax.get('/api/v1/artefacts', { params })
res.json(filterSuggestedArtefacts(remote.data, doc.selectedArtefacts))
const suggested = filterSuggestedArtefacts(remote.data, doc.selectedArtefacts)
// One local read for the page: the admin table shows, per row, whether the
// mirror exists and whether it is behind the upstream.
const locals = await mongo.artefacts
.find({ _id: { $in: suggested.results.map((a: { _id: string }) => a._id) } }, { projection: { dataUpdatedAt: 1, origin: 1 } })
.toArray()
res.json({ ...suggested, results: annotateLocalState(suggested.results, locals, doc._id) })
} catch (err) { next(err) }
})

const remoteThumbnailTypes = new Set(['image/webp', 'image/svg+xml'])

// Proxy an upstream thumbnail for the admin's selection table. The upstream
// url is the one the *server* reaches (possibly an internal one), and the
// upstream sets Cross-Origin-Resource-Policy: same-origin on its assets, so the
// browser cannot load them directly.
router.get('/:id/remote-thumbnails/:thumbnailId/data', 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 (!/^[\w-]+$/.test(req.params.thumbnailId)) throw httpError(400, 'invalid thumbnail id')

const ax = axiosBuilder({ baseURL: doc._id, headers: { 'x-api-key': decipher(doc.apiKey) } })
const remote = await ax.get(`/api/v1/thumbnails/${req.params.thumbnailId}/data`, { responseType: 'stream', validateStatus: () => true })
if (remote.status !== 200) throw httpError(404, 'thumbnail not found')
// The upstream is another deployment: never relay its Content-Type blindly
// onto our origin. Only the types a registry stores are accepted, and SVG
// is served sandboxed so it cannot run scripts if opened as a document.
const contentType = String(remote.headers['content-type'] || '').split(';')[0].trim().toLowerCase()
if (!remoteThumbnailTypes.has(contentType)) {
remote.data.destroy()
throw httpError(415, 'unsupported thumbnail type')
}
res.set('Content-Type', contentType)
res.set('X-Content-Type-Options', 'nosniff')
res.set('Content-Security-Policy', "default-src 'none'; sandbox")
if (remote.headers['content-length']) res.set('Content-Length', remote.headers['content-length'])
// Upstream ids change on every replace, so the bytes behind one id never do.
res.set('Cache-Control', 'private, max-age=31536000, immutable')
await pipeline(remote.data, res)
} catch (err) { next(err) }
})

Expand Down Expand Up @@ -171,6 +214,9 @@ router.post('/:id/selected-artefacts', async (req, res, next) => {
$set: { updatedAt: new Date().toISOString() }
}
)
// Mirror it right away rather than waiting for the daily job or a manual
// full sync; progress is published on the registry's ws channel.
await enqueueArtefactSync(req.params.id, artefactId)
res.status(201).json({ artefactId })
} catch (err) { next(err) }
})
Expand All @@ -185,7 +231,7 @@ router.delete('/:id/selected-artefacts/:artefactId', async (req, res, next) => {
await mongo.remoteRegistries.updateOne(
{ _id: req.params.id },
{
$pull: { selectedArtefacts: req.params.artefactId },
$pull: { selectedArtefacts: req.params.artefactId, pendingSync: req.params.artefactId },
$set: { updatedAt: new Date().toISOString() }
}
)
Expand Down
90 changes: 83 additions & 7 deletions api/src/remote-registries/sync.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { randomUUID } from 'node:crypto'
import { Binary } from 'mongodb'
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'
Expand Down Expand Up @@ -31,6 +32,40 @@ const emitSync = async (remoteRegistryId: string, event: SyncEvent) => {
}
}

type RemoteThumbnail = NonNullable<Artefact['thumbnail']>

// Mirror the upstream thumbnail, keeping its id: a thumbnail's id changes on
// every replace upstream, so comparing ids is enough to know whether the local
// copy is current. Runs after the artefact doc exists locally, and outside the
// tarball fast path — a thumbnail change never bumps dataUpdatedAt.
const syncThumbnail = async (
ax: AxiosInstance,
artefactId: string,
remoteThumbnail: RemoteThumbnail | undefined,
localThumbnail: RemoteThumbnail | undefined
) => {
if (remoteThumbnail?.id === localThumbnail?.id) return
if (!remoteThumbnail) {
await mongo.thumbnails.deleteMany({ artefactId })
await mongo.artefacts.updateOne({ _id: artefactId }, { $unset: { thumbnail: '' } })
return
}
const res = await ax.get(`/api/v1/thumbnails/${remoteThumbnail.id}/data`, { responseType: 'arraybuffer' })
const data = Buffer.from(res.data)
await mongo.thumbnails.deleteMany({ artefactId })
await mongo.thumbnails.insertOne({
_id: remoteThumbnail.id,
artefactId,
data: new Binary(data),
mimeType: res.headers['content-type'] === 'image/svg+xml' ? 'image/svg+xml' : 'image/webp',
width: remoteThumbnail.width,
height: remoteThumbnail.height,
byteSize: data.byteLength,
createdAt: new Date().toISOString()
})
await mongo.artefacts.updateOne({ _id: artefactId }, { $set: { thumbnail: remoteThumbnail } })
}

const syncNpmArtefact = async (ax: AxiosInstance, remoteUrl: string, artefactId: string) => {
const encodedId = encodeURIComponent(artefactId)
const remoteRes = await ax.get(`/api/v1/artefacts/${encodedId}`)
Expand All @@ -40,6 +75,7 @@ const syncNpmArtefact = async (ax: AxiosInstance, remoteUrl: string, artefactId:

// Fast path: same upstream dataUpdatedAt means no new upload to mirror.
if (local?.path && local.dataUpdatedAt === remoteArtefact.dataUpdatedAt) {
await syncThumbnail(ax, artefactId, remoteArtefact.thumbnail, local.thumbnail)
return
}

Expand Down Expand Up @@ -87,6 +123,7 @@ const syncNpmArtefact = async (ax: AxiosInstance, remoteUrl: string, artefactId:
if (oldPath && oldPath !== localPath) {
await filesStorage.delete(oldPath).catch(() => {})
}
await syncThumbnail(ax, artefactId, remoteArtefact.thumbnail, local?.thumbnail)
}

const syncFileArtefact = async (ax: AxiosInstance, remoteUrl: string, artefactId: string) => {
Expand Down Expand Up @@ -145,15 +182,32 @@ const syncFileArtefact = async (ax: AxiosInstance, remoteUrl: string, artefactId
{ $set: { origin: remoteUrl } }
)
}
await syncThumbnail(ax, artefactId, remoteArtefact.thumbnail, local?.thumbnail)
}

// Atomically take the queued selections, so two drains can't sync the same id twice.
const drainPendingSync = async (remoteRegistryId: string): Promise<string[]> => {
const doc = await mongo.remoteRegistries.findOneAndUpdate(
{ _id: remoteRegistryId },
{ $unset: { pendingSync: '' } },
{ returnDocument: 'before', projection: { pendingSync: 1 } }
)
return doc?.pendingSync ?? []
}

// The actual work. Callers own the lock.
const runSync = async (remoteRegistryId: string) => {
export type SyncScope = 'all' | 'pending'

// The actual work. Callers own the lock. `pending` syncs only the artefacts
// queued by selections (see pendingSync); `all` walks every selected artefact.
// Either way, selections queued while this run was in flight are drained before
// returning, so the lock is only released once nothing is left to sync.
const runSync = async (remoteRegistryId: string, scope: SyncScope = 'all') => {
const remote = await mongo.remoteRegistries.findOne({ _id: remoteRegistryId })
if (!remote) return

const artefactIds = scope === 'pending' ? await drainPendingSync(remoteRegistryId) : remote.selectedArtefacts
const startedAt = new Date().toISOString()
const total = remote.selectedArtefacts.length
const total = artefactIds.length
let done = 0

await mongo.remoteRegistries.updateOne(
Expand All @@ -171,7 +225,7 @@ const runSync = async (remoteRegistryId: string) => {
let hasErrors = false
let lastError = ''

for (const artefactId of remote.selectedArtefacts) {
for (const artefactId of artefactIds) {
await mongo.remoteRegistries.updateOne(
{ _id: remoteRegistryId },
{ $set: { 'syncProgress.currentArtefact': artefactId } }
Expand Down Expand Up @@ -230,19 +284,41 @@ const runSync = async (remoteRegistryId: string) => {
lastSyncStatus,
...(hasErrors ? { lastSyncError: lastError } : {})
})

// A full run already covered anything queued meanwhile only if it was
// selected before the run read selectedArtefacts; draining is cheap (the
// dataUpdatedAt fast path) and keeps the queue semantics simple.
const pending = await mongo.remoteRegistries.findOne({ _id: remoteRegistryId }, { projection: { pendingSync: 1 } })
if (pending?.pendingSync?.length) await runSync(remoteRegistryId, 'pending')
}

// 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<boolean> => {
export const startSync = async (remoteRegistryId: string, scope: SyncScope = 'all'): Promise<boolean> => {
const lockId = syncLockId(remoteRegistryId)
if (!await locks.acquire(lockId)) return false
runSync(remoteRegistryId)
runSync(remoteRegistryId, scope)
.catch(err => internalError('sync-remote-registry', err))
.finally(() => locks.release(lockId).catch(err => internalError('sync-remote-registry-release', err)))
.finally(async () => {
await locks.release(lockId).catch(err => internalError('sync-remote-registry-release', err))
// A selection can land between the final drain and the release above;
// its own startSync lost the lock race, so pick it up here.
const doc = await mongo.remoteRegistries.findOne({ _id: remoteRegistryId }, { projection: { pendingSync: 1 } })
if (doc?.pendingSync?.length) await startSync(remoteRegistryId, 'pending')
})
return true
}

// Queue one freshly selected artefact and sync it in the background. If a sync
// already holds the lock, that run drains the queue before releasing it.
export const enqueueArtefactSync = async (remoteRegistryId: string, artefactId: string) => {
await mongo.remoteRegistries.updateOne(
{ _id: remoteRegistryId },
{ $addToSet: { pendingSync: artefactId } }
)
await startSync(remoteRegistryId, 'pending')
}

// Awaits completion. Used by the daily job, which syncs registries one at a time.
export const syncRemoteRegistry = async (remoteRegistryId: string): Promise<boolean> => {
const lockId = syncLockId(remoteRegistryId)
Expand Down
Loading
Loading