mirror of
https://github.com/Crosstalk-Solutions/project-nomad.git
synced 2026-10-02 03:24:51 +08:00
fix(drug-reference): backfill the missing dataset install row on boot
On 1.34.0 every FDA ingest finished with its labels loaded but lost the
installed_resources 'dataset' row, because resource_type was still
enum('zim','map'). The enum migration in #1243 fixes new writes but can't recreate
the lost row. Re-selecting the tier doesn't fix it either, because ZimService
skips a dataset that already has rows. So Medicine never resolved its
Standard tier.
- Add DrugInstallRowProvider. On each boot it backfills the row when an
ingest has finished (rows present, export date set, no download
marker, no drug job running) and no row exists.
- Put the decision in decideDrugRowReconcile() so it has unit tests.
- Add DrugReferenceService.recordInstalledRow() as the one writer used
by the ingest job, the provider, and ZimService.
- ZimService now writes the row when a Medicine tier re-select finds a
finished ingest with no row.
Manual ingests and reset-and-reingest runs, which dispatch without
resourceMeta, also get the row on the next boot.
This commit is contained in:
@@ -58,6 +58,7 @@ export default defineConfig({
|
||||
() => import('#providers/qdrant_restart_policy_provider'),
|
||||
() => import('#providers/version_check_provider'),
|
||||
() => import('#providers/gpu_passthrough_remediation_provider'),
|
||||
() => import('#providers/drug_install_row_provider'),
|
||||
],
|
||||
|
||||
/*
|
||||
|
||||
@@ -749,37 +749,25 @@ export class IngestDrugDataJob {
|
||||
// Install-state write-back: when this ingest was kicked off by a curated-tier
|
||||
// install, record an `installed_resources` row so the tier-status math (which
|
||||
// is row-driven for ZIM/map) recognizes the dataset uniformly — the home-tile
|
||||
// gate and the tier "installed" badge both read these rows. resource_type
|
||||
// 'dataset' was widened onto the model in the foundational slice. The row's
|
||||
// gate and the tier "installed" badge both read these rows. The row's
|
||||
// `version` is the openFDA export_date (the real freshness key), not the
|
||||
// manifest placeholder. Manual downloads pass no resourceMeta → no row, which
|
||||
// is correct: install-state belongs to the curated-tier path only.
|
||||
// manifest placeholder. Manual ingests pass no resourceMeta and write no row
|
||||
// here; DrugInstallRowProvider backfills that case on the next admin boot.
|
||||
if (!resourceMeta) return
|
||||
try {
|
||||
const { default: InstalledResource } = await import('#models/installed_resource')
|
||||
const { DateTime } = await import('luxon')
|
||||
const totalBytes = await this.totalDownloadedBytes()
|
||||
await InstalledResource.updateOrCreate(
|
||||
{ resource_id: resourceMeta.resourceId, resource_type: 'dataset' },
|
||||
{
|
||||
version: exportDate,
|
||||
collection_ref: resourceMeta.collectionRef,
|
||||
url: 'https://api.fda.gov/download.json',
|
||||
// No single on-disk file — the parts are deleted after ingest; the
|
||||
// installed artifact is the DB table. Record the staging dir for
|
||||
// provenance; uninstall keys off resource_id, never this path.
|
||||
file_path: STORAGE_BASE,
|
||||
file_size_bytes: totalBytes,
|
||||
installed_at: DateTime.now(),
|
||||
}
|
||||
)
|
||||
const { DrugReferenceService } = await import('#services/drug_reference_service')
|
||||
await new DrugReferenceService().recordInstalledRow({
|
||||
version: exportDate,
|
||||
collectionRef: resourceMeta.collectionRef,
|
||||
fileSizeBytes: await this.totalDownloadedBytes(),
|
||||
})
|
||||
logger.info(
|
||||
`[IngestDrugDataJob] Wrote installed_resources row for ${resourceMeta.resourceId} (export_date=${exportDate})`
|
||||
)
|
||||
} catch (err) {
|
||||
// A failed row write must NOT abort a completed ingest — the data is
|
||||
// already searchable. Log loud; the tier badge will simply read
|
||||
// not-installed until the next install reconcile.
|
||||
// not-installed until DrugInstallRowProvider backfills it on the next boot.
|
||||
logger.error(
|
||||
`[IngestDrugDataJob] Failed to write installed_resources row for ${resourceMeta.resourceId}: ${
|
||||
err instanceof Error ? err.message : String(err)
|
||||
|
||||
@@ -566,6 +566,60 @@ export class DrugReferenceService {
|
||||
return { ...result, nothingDownloaded: false }
|
||||
}
|
||||
|
||||
/**
|
||||
* Write (or refresh) the `installed_resources` 'dataset' row that the tier-
|
||||
* status math and the curated-install filter read. Shared by the ingest job's
|
||||
* final pass, the ZimService tier-select path, and DrugInstallRowProvider's
|
||||
* boot backfill so all three write the same shape. Throws; callers log.
|
||||
*/
|
||||
async recordInstalledRow(opts: {
|
||||
version: string
|
||||
collectionRef: string | null
|
||||
fileSizeBytes: number | null
|
||||
}): Promise<void> {
|
||||
const { default: InstalledResource } = await import('#models/installed_resource')
|
||||
const { DateTime } = await import('luxon')
|
||||
await InstalledResource.updateOrCreate(
|
||||
{ resource_id: DRUG_DATASET_RESOURCE_ID, resource_type: 'dataset' },
|
||||
{
|
||||
version: opts.version,
|
||||
collection_ref: opts.collectionRef,
|
||||
url: 'https://api.fda.gov/download.json',
|
||||
// No single on-disk file — the parts are deleted after ingest; the
|
||||
// installed artifact is the DB table. Record the staging dir for
|
||||
// provenance; uninstall keys off resource_id, never this path.
|
||||
file_path: STORAGE_BASE,
|
||||
file_size_bytes: opts.fileSizeBytes,
|
||||
installed_at: DateTime.now(),
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
/** Whether the `installed_resources` 'dataset' row exists. */
|
||||
async hasInstalledRow(): Promise<boolean> {
|
||||
const { default: InstalledResource } = await import('#models/installed_resource')
|
||||
const row = await InstalledResource.query()
|
||||
.where('resource_type', 'dataset')
|
||||
.where('resource_id', DRUG_DATASET_RESOURCE_ID)
|
||||
.first()
|
||||
return row !== null
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether either drug queue holds an active/waiting/delayed job. Counts the
|
||||
* whole queue rather than the deterministic jobIds, since ingest continuations
|
||||
* run without one. Throws when Redis is unreachable.
|
||||
*/
|
||||
async isJobInFlight(): Promise<boolean> {
|
||||
for (const queueName of [DownloadDrugDataJob.queue, IngestDrugDataJob.queue]) {
|
||||
const counts = await QueueService.getInstance()
|
||||
.getQueue(queueName)
|
||||
.getJobCounts('active', 'waiting', 'delayed')
|
||||
if ((counts.active ?? 0) + (counts.waiting ?? 0) + (counts.delayed ?? 0) > 0) return true
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
/**
|
||||
* Uninstall the offline FDA drug dataset — the curated-tier "remove" path.
|
||||
*
|
||||
|
||||
@@ -359,6 +359,26 @@ export class ZimService {
|
||||
const status = await drugReferenceService.getIngestStatus()
|
||||
if (status.phase === 'ready' || status.rowCount > 0) {
|
||||
logger.info('[ZimService] Drug dataset already ingested, skipping dispatch.')
|
||||
// Reaching here means no 'dataset' row exists (the installed filter
|
||||
// above would have dropped the resource). A finished ingest that lost
|
||||
// its row write (1.34.0's narrow enum, or a manual ingest) gets one now,
|
||||
// so re-selecting the tier resolves it without waiting for a reboot.
|
||||
if (status.phase === 'ready' && status.lastUpdated) {
|
||||
try {
|
||||
await drugReferenceService.recordInstalledRow({
|
||||
version: status.lastUpdated,
|
||||
collectionRef: categorySlug,
|
||||
fileSizeBytes: null,
|
||||
})
|
||||
logger.info('[ZimService] Backfilled installed_resources row for the drug dataset.')
|
||||
} catch (err) {
|
||||
logger.error(
|
||||
`[ZimService] Failed to backfill drug dataset row: ${
|
||||
err instanceof Error ? err.message : String(err)
|
||||
}`
|
||||
)
|
||||
}
|
||||
}
|
||||
continue
|
||||
}
|
||||
// DownloadDrugDataJob.dispatch() is idempotent on its deterministic
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
import type { ZimCategoriesSpec } from '../../types/collections.js'
|
||||
|
||||
/**
|
||||
* Observed state of the FDA drug dataset, gathered at admin boot by
|
||||
* DrugInstallRowProvider.
|
||||
*/
|
||||
export interface DrugRowReconcileInput {
|
||||
/** Rows in `drug_labels`. */
|
||||
rowCount: number
|
||||
/** Whether the `installed_resources` 'dataset' row already exists. */
|
||||
hasInstallRow: boolean
|
||||
/** KV `drugReference.lastUpdatedExportDate`; only the final ingest pass sets it. */
|
||||
lastUpdatedExportDate: string | null
|
||||
/** Whether KV `drugReference.downloadState` still holds a marker. */
|
||||
hasDownloadMarker: boolean
|
||||
/** Whether either drug queue has an active/waiting/delayed job. */
|
||||
jobInFlight: boolean
|
||||
}
|
||||
|
||||
export type DrugRowReconcileDecision =
|
||||
| { action: 'backfill'; version: string }
|
||||
| { action: 'skip'; reason: string }
|
||||
|
||||
/**
|
||||
* Decide whether a completed drug ingest is missing its `installed_resources`
|
||||
* row and should get one.
|
||||
*
|
||||
* On 1.34.0 the column was still enum('zim','map'), so every ingest finished
|
||||
* `ready` with 261k labels and a failed row write. The tier-status math reads
|
||||
* that row, so Medicine never resolved its Standard tier. The same state also
|
||||
* follows a manual ingest (no resourceMeta), including a reset-and-reingest.
|
||||
*
|
||||
* Backfill only when the ingest provably finished: the export-date marker is
|
||||
* written by the final pass, and the download marker is cleared right after it
|
||||
* (the same "completed" inference getIngestStatus() makes). A running job is
|
||||
* left alone, since it writes its own row on `ready`.
|
||||
*/
|
||||
export function decideDrugRowReconcile(input: DrugRowReconcileInput): DrugRowReconcileDecision {
|
||||
if (input.hasInstallRow) return { action: 'skip', reason: 'install row present' }
|
||||
if (input.rowCount <= 0) return { action: 'skip', reason: 'no drug labels ingested' }
|
||||
if (input.jobInFlight) return { action: 'skip', reason: 'drug download/ingest in flight' }
|
||||
if (input.hasDownloadMarker) {
|
||||
return { action: 'skip', reason: 'download marker present (ingest not finished)' }
|
||||
}
|
||||
const version = input.lastUpdatedExportDate?.trim()
|
||||
if (!version) return { action: 'skip', reason: 'no export date (ingest never completed)' }
|
||||
return { action: 'backfill', version }
|
||||
}
|
||||
|
||||
/**
|
||||
* Slug of the curated category whose tiers carry the given dataset resource,
|
||||
* for the row's `collection_ref`. Null when the spec is absent or no tier lists it.
|
||||
*/
|
||||
export function findDatasetCategorySlug(
|
||||
spec: ZimCategoriesSpec | null,
|
||||
resourceId: string
|
||||
): string | null {
|
||||
for (const category of spec?.categories ?? []) {
|
||||
for (const tier of category.tiers ?? []) {
|
||||
if ((tier.resources ?? []).some((r) => r.type === 'dataset' && r.id === resourceId)) {
|
||||
return category.slug
|
||||
}
|
||||
}
|
||||
}
|
||||
return null
|
||||
}
|
||||
@@ -0,0 +1,91 @@
|
||||
import logger from '@adonisjs/core/services/logger'
|
||||
import type { ApplicationService } from '@adonisjs/core/types'
|
||||
import type { ZimCategoriesSpec } from '../types/collections.js'
|
||||
|
||||
/**
|
||||
* Backfills the FDA drug dataset's `installed_resources` row when a completed
|
||||
* ingest has none.
|
||||
*
|
||||
* On 1.34.0 `installed_resources.resource_type` was still enum('zim','map'), so
|
||||
* every drug ingest landed its 261k labels, reported `ready`, and lost the
|
||||
* 'dataset' row write. The migration that widens the enum can't recreate the
|
||||
* row, and re-selecting the tier won't either (ZimService skips a dataset that
|
||||
* already has rows), so Medicine stayed stuck below its Standard tier. Manual
|
||||
* ingests and reset-and-reingest runs dispatch without resourceMeta and end in
|
||||
* the same state.
|
||||
*
|
||||
* Runs once per admin boot. The decision lives in decideDrugRowReconcile(); once
|
||||
* the row exists every later boot is a single skip.
|
||||
*/
|
||||
export default class DrugInstallRowProvider {
|
||||
constructor(protected app: ApplicationService) {}
|
||||
|
||||
async boot() {
|
||||
if (this.app.getEnvironment() !== 'web') return
|
||||
|
||||
setImmediate(async () => {
|
||||
try {
|
||||
const KVStore = (await import('#models/kv_store')).default
|
||||
const { DrugReferenceService, DRUG_DATASET_RESOURCE_ID } = await import(
|
||||
'#services/drug_reference_service'
|
||||
)
|
||||
const { CollectionManifestService } = await import(
|
||||
'#services/collection_manifest_service'
|
||||
)
|
||||
const { parseDownloadState } = await import('../util/drug_labels.js')
|
||||
const { decideDrugRowReconcile, findDatasetCategorySlug } = await import(
|
||||
'../app/utils/drug_install_row_reconcile.js'
|
||||
)
|
||||
|
||||
const drugService = new DrugReferenceService()
|
||||
|
||||
// Redis down reads as "in flight": leaving the tier stuck one more boot
|
||||
// beats racing a live ingest.
|
||||
let jobInFlight = true
|
||||
try {
|
||||
jobInFlight = await drugService.isJobInFlight()
|
||||
} catch (err: any) {
|
||||
logger.warn(
|
||||
`[DrugInstallRowProvider] Could not read drug queues (${err?.message ?? err}) — treating as in flight.`
|
||||
)
|
||||
}
|
||||
|
||||
const [rowCount, hasInstallRow, lastUpdated, rawMarker] = await Promise.all([
|
||||
drugService.rowCount(),
|
||||
drugService.hasInstalledRow(),
|
||||
KVStore.getValue('drugReference.lastUpdatedExportDate'),
|
||||
KVStore.getValue('drugReference.downloadState'),
|
||||
])
|
||||
|
||||
const decision = decideDrugRowReconcile({
|
||||
rowCount,
|
||||
hasInstallRow,
|
||||
lastUpdatedExportDate: lastUpdated ? String(lastUpdated) : null,
|
||||
hasDownloadMarker: parseDownloadState(rawMarker) !== null,
|
||||
jobInFlight,
|
||||
})
|
||||
|
||||
if (decision.action === 'skip') {
|
||||
logger.info(`[DrugInstallRowProvider] No backfill needed: ${decision.reason}.`)
|
||||
return
|
||||
}
|
||||
|
||||
// Cached spec only: no network at boot, and the box may be offline.
|
||||
const spec = await new CollectionManifestService().getCachedSpec<ZimCategoriesSpec>('zim_categories')
|
||||
const collectionRef = findDatasetCategorySlug(spec, DRUG_DATASET_RESOURCE_ID)
|
||||
|
||||
await drugService.recordInstalledRow({
|
||||
version: decision.version,
|
||||
collectionRef,
|
||||
fileSizeBytes: null,
|
||||
})
|
||||
logger.warn(
|
||||
`[DrugInstallRowProvider] ${rowCount} drug labels were ingested without an installed_resources row. ` +
|
||||
`Backfilled ${DRUG_DATASET_RESOURCE_ID} (version=${decision.version}, collection_ref=${collectionRef ?? 'null'}).`
|
||||
)
|
||||
} catch (err: any) {
|
||||
logger.error(`[DrugInstallRowProvider] Backfill check failed: ${err?.message ?? err}`)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,93 @@
|
||||
import * as assert from 'node:assert/strict'
|
||||
import { test } from 'node:test'
|
||||
|
||||
import {
|
||||
decideDrugRowReconcile,
|
||||
findDatasetCategorySlug,
|
||||
type DrugRowReconcileInput,
|
||||
} from '../../app/utils/drug_install_row_reconcile.js'
|
||||
import type { ZimCategoriesSpec } from '../../types/collections.js'
|
||||
|
||||
// The state a 1.34.0 box is left in: full ingest, failed row write.
|
||||
const stuck: DrugRowReconcileInput = {
|
||||
rowCount: 261671,
|
||||
hasInstallRow: false,
|
||||
lastUpdatedExportDate: '2026-08-01',
|
||||
hasDownloadMarker: false,
|
||||
jobInFlight: false,
|
||||
}
|
||||
|
||||
test('a completed ingest with no install row is backfilled at its export date', () => {
|
||||
assert.deepEqual(decideDrugRowReconcile(stuck), { action: 'backfill', version: '2026-08-01' })
|
||||
})
|
||||
|
||||
test('an existing install row is left alone', () => {
|
||||
assert.equal(decideDrugRowReconcile({ ...stuck, hasInstallRow: true }).action, 'skip')
|
||||
})
|
||||
|
||||
test('an empty drug_labels table is not backfilled', () => {
|
||||
assert.equal(decideDrugRowReconcile({ ...stuck, rowCount: 0 }).action, 'skip')
|
||||
})
|
||||
|
||||
test('a partial ingest (no export date) is not backfilled', () => {
|
||||
assert.equal(decideDrugRowReconcile({ ...stuck, lastUpdatedExportDate: null }).action, 'skip')
|
||||
assert.equal(decideDrugRowReconcile({ ...stuck, lastUpdatedExportDate: ' ' }).action, 'skip')
|
||||
})
|
||||
|
||||
test('a surviving download marker means the ingest has not finished', () => {
|
||||
assert.equal(decideDrugRowReconcile({ ...stuck, hasDownloadMarker: true }).action, 'skip')
|
||||
})
|
||||
|
||||
test('an in-flight job is left to write its own row', () => {
|
||||
assert.equal(decideDrugRowReconcile({ ...stuck, jobInFlight: true }).action, 'skip')
|
||||
})
|
||||
|
||||
const spec = {
|
||||
spec_version: '1',
|
||||
categories: [
|
||||
{
|
||||
name: 'Survival',
|
||||
slug: 'survival',
|
||||
icon: '',
|
||||
description: '',
|
||||
language: 'en',
|
||||
tiers: [{ name: 'Essential', slug: 'survival-essential', description: '', resources: [] }],
|
||||
},
|
||||
{
|
||||
name: 'Medicine',
|
||||
slug: 'medicine',
|
||||
icon: '',
|
||||
description: '',
|
||||
language: 'en',
|
||||
tiers: [
|
||||
{ name: 'Essential', slug: 'medicine-essential', description: '', resources: [] },
|
||||
{
|
||||
name: 'Standard',
|
||||
slug: 'medicine-standard',
|
||||
description: '',
|
||||
includesTier: 'medicine-essential',
|
||||
resources: [
|
||||
{
|
||||
id: 'openfda-drug-labels',
|
||||
type: 'dataset',
|
||||
version: '2025-01',
|
||||
title: 'FDA Drug Reference',
|
||||
description: '',
|
||||
url: 'https://api.fda.gov/download.json',
|
||||
size_mb: 1700,
|
||||
},
|
||||
],
|
||||
},
|
||||
],
|
||||
},
|
||||
],
|
||||
} as unknown as ZimCategoriesSpec
|
||||
|
||||
test('the dataset resolves to the category whose tier lists it', () => {
|
||||
assert.equal(findDatasetCategorySlug(spec, 'openfda-drug-labels'), 'medicine')
|
||||
})
|
||||
|
||||
test('an unknown dataset or missing spec yields null', () => {
|
||||
assert.equal(findDatasetCategorySlug(spec, 'nope'), null)
|
||||
assert.equal(findDatasetCategorySlug(null, 'openfda-drug-labels'), null)
|
||||
})
|
||||
Reference in New Issue
Block a user