From 7ddf13053f284ab8286304b76c4f99f8ebf33645 Mon Sep 17 00:00:00 2001 From: xCyanGrizzly Date: Thu, 23 Jul 2026 14:23:53 +0200 Subject: [PATCH] feat(worker): opportunistic provenance backfill during scan Wire tryProvenanceBackfill into processOneArchiveSet: before downloading a scanned ZIP/RAR/7Z, check whether it's the true origin of a placeholder-provenance package in the destination channel and backfill in place, skipping the download. Add the zipsBackfilled counter through PipelineContext, updateRunActivity, and completeIngestionRun. Co-Authored-By: Claude Opus 4.8 (1M context) --- worker/src/db/queries.ts | 3 +++ worker/src/worker.ts | 53 ++++++++++++++++++++++++++++++++++++++++ 2 files changed, 56 insertions(+) diff --git a/worker/src/db/queries.ts b/worker/src/db/queries.ts index ab39cbc..ef3101f 100644 --- a/worker/src/db/queries.ts +++ b/worker/src/db/queries.ts @@ -376,6 +376,7 @@ export interface ActivityUpdate { zipsFound?: number; zipsDuplicate?: number; zipsIngested?: number; + zipsBackfilled?: number; } export async function updateRunActivity( @@ -399,6 +400,7 @@ export async function updateRunActivity( ...(activity.zipsFound !== undefined && { zipsFound: activity.zipsFound }), ...(activity.zipsDuplicate !== undefined && { zipsDuplicate: activity.zipsDuplicate }), ...(activity.zipsIngested !== undefined && { zipsIngested: activity.zipsIngested }), + ...(activity.zipsBackfilled !== undefined && { zipsBackfilled: activity.zipsBackfilled }), ...(activity.currentTopicId !== undefined && { currentTopicId: activity.currentTopicId }), ...(activity.currentAccountChannelMapId !== undefined && { currentAccountChannelMapId: activity.currentAccountChannelMapId, @@ -429,6 +431,7 @@ export async function completeIngestionRun( zipsFound: number; zipsDuplicate: number; zipsIngested: number; + zipsBackfilled: number; } ) { return db.ingestionRun.update({ diff --git a/worker/src/worker.ts b/worker/src/worker.ts index 3c3971f..430283e 100644 --- a/worker/src/worker.ts +++ b/worker/src/worker.ts @@ -63,6 +63,7 @@ import { hashParts } from "./archive/hash.js"; import { readZipCentralDirectory } from "./archive/zip-reader.js"; import { readRarContents } from "./archive/rar-reader.js"; import { read7zContents } from "./archive/sevenz-reader.js"; +import { tryProvenanceBackfill } from "./provenance-backfill.js"; import { byteLevelSplit, concatenateFiles } from "./archive/split.js"; import { uploadToChannel, UploadStallError } from "./upload/channel.js"; import { processAlbumGroups, detectGroupingConflicts, type IndexedPackageRef } from "./grouping.js"; @@ -314,6 +315,7 @@ interface PipelineContext { zipsFound: number; zipsDuplicate: number; zipsIngested: number; + zipsBackfilled: number; }; /** Creator from forum topic name (null for non-forum). */ topicCreator: string | null; @@ -422,6 +424,7 @@ export async function runWorkerForAccount( zipsFound: 0, zipsDuplicate: 0, zipsIngested: 0, + zipsBackfilled: 0, }; try { @@ -1645,6 +1648,56 @@ async function processOneArchiveSet( return null; } + // ── Cross-channel provenance backfill ── + // The same-channel checks above missed. Before downloading, see if this + // archive is the true origin of a placeholder-source package (manual upload + // / rebuild record whose sourceChannelId == destChannelId). If so, backfill + // its real provenance and skip the download entirely. + const archType = archiveSet.type === "7Z" ? "SEVEN_Z" : archiveSet.type; + if (destChannelId && (archType === "ZIP" || archType === "RAR" || archType === "SEVEN_Z")) { + try { + const derivedCreator = + topicCreator && topicCreator !== "General" + ? topicCreator + : (extractCreatorFromFileName(archiveName) ?? topicCreator ?? null); + const preview = previewMatches.get(archiveSet.baseName); + const result = await tryProvenanceBackfill({ + client, + destChannelId, + scannedSourceChannelId: channel.id, + fileName: archiveName, + fileSize: totalArchiveSize, + archiveType: archType, + sourceMessageId: archiveSet.parts[0].id, + sourceTopicId, + sourceCaption: archiveSet.parts[0].caption ?? null, + remoteUniqueId: archiveSet.parts[0].remoteUniqueId ?? null, + creator: derivedCreator, + scannedFileId: archiveSet.parts[archiveSet.parts.length - 1].fileId, + previewData: null, + previewMsgId: preview?.id ?? null, + }); + if (result.backfilled) { + counters.zipsBackfilled++; + accountLog.info( + { fileName: archiveName, sourceMessageId: Number(archiveSet.parts[0].id), confidence: result.confidence }, + "Backfilled provenance for placeholder package — skipping download", + ); + await updateRunActivity(runId, { + currentActivity: `Backfilled provenance for ${archiveName}`, + currentStep: "backfilling", + currentFile: archiveName, + currentFileNum: setIdx + 1, + totalFiles: totalSets, + zipsBackfilled: counters.zipsBackfilled, + }); + return null; + } + } catch (err) { + accountLog.warn({ err, fileName: archiveName }, "Provenance backfill attempt failed (non-fatal), continuing to normal ingestion"); + } + } + // ── Size guard: skip archives that exceed WORKER_MAX_ZIP_SIZE_MB ── const maxSizeBytes = BigInt(config.maxZipSizeMB) * 1024n * 1024n; if (totalArchiveSize > maxSizeBytes) {