Building a Sample Manager: Reworking the Analysis Pipeline into a Job Queue

Building a Sample Manager: Reworking the Analysis Pipeline into a Job Queue

29/09/2026

|

Projects

Update

This update is less about new features on the surface and more about the plumbing underneath finally catching up with itself: two real ownership bugs, a full rebuild of how analysis actually runs, merging two repos into one, and the first proper spec for the next big feature.

Two ownership gaps, closed

Went back through collections and uploads looking for security holes and found two real ones. The upload action was filtering requested collection ids by existence only, with no check that they actually belonged to the uploading user, so a crafted request could attach a sample to someone else’s collection. The collection detail route had a similar gap: it loaded a collection by id without checking ownership, so although other users’ samples were correctly filtered out of the list, the collection’s name and colour were visible to anyone who had the id. Both got a proper userId check, and both now return a 404 rather than a 403 for “doesn’t exist” and “not yours” alike, the same convention GitHub uses for private repos, so a 403 can’t be used to confirm something exists that you shouldn’t be able to see.

Rethinking how analysis actually runs

Trying to design a re-analyse feature exposed how shaky the existing analysis flow really was. Failures disappeared silently, since the upload action never checked the response and just left rows stuck on pending. Uploads waited on the full analysis to finish before returning, which is a real delay for a long pad. The analysis endpoint had no authentication at all, trusting a filename and sample id straight from the request body. And nothing recorded which version of the analysis code had produced a given row, so after the key detection changes there was no way to tell which rows were stale.

The fix was building a proper job queue rather than continuing to defer it. Went with pg-boss, running on the same Postgres rather than adding Redis, mainly because the sample insert and the job enqueue can now happen inside the same database transaction, atomically. With Redis, those two writes can’t share a transaction, so a crash between them can leave a row that nothing will ever pick up.

The queue sits behind a small module that exposes just a handful of functions (enqueueAnalysis, startAnalysisWorker, stopJobs), so the choice of pg-boss over anything else stays swappable later. The handler itself follows a simple contract: returning means the job succeeded, throwing means it failed and gets retried, and the row is written to processing, then complete or failed with the actual error recorded, rather than silently sitting on pending forever. FastAPI stays completely stateless, file in, numbers out, with all database access staying in SvelteKit.

export async function enqueueAnalysis(sampleId: string, tx?: Transaction) {
    const boss = await getBoss();
    const jobId = await boss.send(ANALYSIS_QUEUE, { sampleId } satisfies AnalysisJobData, tx ? { db: fromDrizzle(tx, sql) } : {})
    return jobId
}

async function registerWorker() {
    const boss = await getBoss();

    return boss.work(
        ANALYSIS_QUEUE,
        { includeMetadata: true, localConcurrency: 3 },
        analyseSample
    );
}

export async function startAnalysisWorker() {
    if (!globalThis.__sampleManagerAnalysisWorker) {
        globalThis.__sampleManagerAnalysisWorker = registerWorker().catch((err) => {
            globalThis.__sampleManagerAnalysisWorker = undefined;
            throw err;
        });
    }
    return globalThis.__sampleManagerAnalysisWorker   
}

export async function stopJobs() {
    const cached = globalThis.__sampleManagerBoss;

    // Nothing was started, so there's nothing to stop
    if (!cached) return;

    globalThis.__sampleManagerBoss = undefined;
    globalThis.__sampleManagerAnalysisWorker = undefined;

    const boss = await cached.catch(() => undefined);
    if (!boss) return;
    await boss.stop({ graceful: true });
}

Where previously the API call to initiate a samples analysis had to be fired off during upload, now I have a reusable function that can be called from anywhere in the application on demand to add jobs into the queue

export async function analyseSample([job]: JobWithMetadata<AnalysisJobData>[]) {
    // Checks if the job count has reached its retry limit to use in the clean up
    const isFinalAttempt = job.retryCount >= job.retryLimit;

    const [sample] = await db.select().from(samples).where(eq(samples.id, job.data.sampleId));

    if (!sample) return console.warn(`Sample ${job.data.sampleId} not found`);

    try {
        // Status updated ahead of analysis
        await db.update(samples).set({
            status: 'processing'
        }).where(eq(samples.id, job.data.sampleId))

        const filename= posix.basename(sample.sampleUrl)

        const librosa = await fetch(`${env.FASTAPI_URL}/files/analyse/${sample.userId}/${filename}`)

        if (!librosa.ok) {
            const body = await librosa.text();
            throw new Error(
                `FastAPI analysis failed for sample ${sample.id}: ${librosa.status} ${body.slice(0, 200)}`
            )
        }

        const librosaRes = await librosa.json()

        // Completed analysis
        await db.update(samples).set({
            sampleBpm: librosaRes.bpm,
            sampleRate: librosaRes.sampleRate,
            duration: librosaRes.duration,
            estimatedKey: librosaRes.key,
            harmonicRatio: librosaRes.harmonicRatio,
            tonality: librosaRes.tonality,
            status: 'complete',
            analysisVersion: librosaRes.analysisVersion,
            analysedAt: new Date(),
            analysisError: null
        }).where(eq(samples.id, job.data.sampleId))
    } catch (err) {
        const message = err instanceof Error ? err.message : String(err);
        try {
            await db.update(samples).set({
                status: isFinalAttempt ? 'failed' : 'pending',
                analysisError: message
            }).where(eq(samples.id, sample.id));
        } catch (updateErr) {
            console.error(`Failed to record analysis error for sample ${sample.id}`, updateErr);
        }

        throw err
    }
}

A few real bugs turned up building this: an unawaited enqueueAnalysis call that could let a transaction commit before the job was actually queued, a worker accidentally calling itself recursively on setup, and a handler pattern copied from documentation that silently marked failed jobs as completed because it returned an object instead of throwing. Also hit a classic hot-reload trap: changing the worker’s code doesn’t take effect until the dev server restarts, since the old handler function stays cached across reloads, which briefly made jobs look like they were completing in milliseconds when nothing was actually running.

Live status in the UI

With jobs now tracked properly, the sample list shows live status, pending, processing, failed, complete, and updates itself while analysis is running without a manual refresh. Getting this right took a couple of wrong turns: reaching for the browser’s invalidate function from inside the server-side worker (which doesn’t exist there and would have caused every job to fail and retry six times before the mistake was caught), then a polling interval that kept stacking new intervals on every reload instead of cleaning up after itself. The final version drives a single interval off a derived boolean that only flips when something’s actually still processing, so there’s exactly one timer running for as long as analysis is in flight and none once everything’s settled.

One repo instead of two

The FastAPI service had been living in its own git repository since the start, which meant any change touching both sides, like the analysis version bump, needed two separate commits in two separate places. Merged it into the main repo with its full commit history intact via git subtree.

What’s next

The near-term backlog is mostly tidying: moving the last few duplicated bulk actions into shared server helpers, and actually testing the failure path properly since no job has failed yet in practice. But the bigger news is the next real feature has a proper spec now: similarity search, split into “sounds like” (timbre), “grooves like” (rhythm), and “vibe” (overall character), built in that order. The first stage uses the existing analysis pipeline to compute a feature fingerprint per sample (MFCCs plus spectral features, standardised against a fixed reference rather than recalculated per user), stored in pgvector, with a “similar” panel per sample showing the closest matches by cosine distance rather than a meaningless percentage. Evaluation gets the same treatment as the key detection work: labelled groups of samples I consider genuinely alike, scored automatically, rather than eyeballing a handful of results and calling it done.