mirror of
https://github.com/xCyanGrizzly/DragonsStash.git
synced 2026-09-21 05:21:43 +00:00
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) <noreply@anthropic.com>
This commit is contained in:
@@ -376,6 +376,7 @@ export interface ActivityUpdate {
|
|||||||
zipsFound?: number;
|
zipsFound?: number;
|
||||||
zipsDuplicate?: number;
|
zipsDuplicate?: number;
|
||||||
zipsIngested?: number;
|
zipsIngested?: number;
|
||||||
|
zipsBackfilled?: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function updateRunActivity(
|
export async function updateRunActivity(
|
||||||
@@ -399,6 +400,7 @@ export async function updateRunActivity(
|
|||||||
...(activity.zipsFound !== undefined && { zipsFound: activity.zipsFound }),
|
...(activity.zipsFound !== undefined && { zipsFound: activity.zipsFound }),
|
||||||
...(activity.zipsDuplicate !== undefined && { zipsDuplicate: activity.zipsDuplicate }),
|
...(activity.zipsDuplicate !== undefined && { zipsDuplicate: activity.zipsDuplicate }),
|
||||||
...(activity.zipsIngested !== undefined && { zipsIngested: activity.zipsIngested }),
|
...(activity.zipsIngested !== undefined && { zipsIngested: activity.zipsIngested }),
|
||||||
|
...(activity.zipsBackfilled !== undefined && { zipsBackfilled: activity.zipsBackfilled }),
|
||||||
...(activity.currentTopicId !== undefined && { currentTopicId: activity.currentTopicId }),
|
...(activity.currentTopicId !== undefined && { currentTopicId: activity.currentTopicId }),
|
||||||
...(activity.currentAccountChannelMapId !== undefined && {
|
...(activity.currentAccountChannelMapId !== undefined && {
|
||||||
currentAccountChannelMapId: activity.currentAccountChannelMapId,
|
currentAccountChannelMapId: activity.currentAccountChannelMapId,
|
||||||
@@ -429,6 +431,7 @@ export async function completeIngestionRun(
|
|||||||
zipsFound: number;
|
zipsFound: number;
|
||||||
zipsDuplicate: number;
|
zipsDuplicate: number;
|
||||||
zipsIngested: number;
|
zipsIngested: number;
|
||||||
|
zipsBackfilled: number;
|
||||||
}
|
}
|
||||||
) {
|
) {
|
||||||
return db.ingestionRun.update({
|
return db.ingestionRun.update({
|
||||||
|
|||||||
@@ -63,6 +63,7 @@ import { hashParts } from "./archive/hash.js";
|
|||||||
import { readZipCentralDirectory } from "./archive/zip-reader.js";
|
import { readZipCentralDirectory } from "./archive/zip-reader.js";
|
||||||
import { readRarContents } from "./archive/rar-reader.js";
|
import { readRarContents } from "./archive/rar-reader.js";
|
||||||
import { read7zContents } from "./archive/sevenz-reader.js";
|
import { read7zContents } from "./archive/sevenz-reader.js";
|
||||||
|
import { tryProvenanceBackfill } from "./provenance-backfill.js";
|
||||||
import { byteLevelSplit, concatenateFiles } from "./archive/split.js";
|
import { byteLevelSplit, concatenateFiles } from "./archive/split.js";
|
||||||
import { uploadToChannel, UploadStallError } from "./upload/channel.js";
|
import { uploadToChannel, UploadStallError } from "./upload/channel.js";
|
||||||
import { processAlbumGroups, detectGroupingConflicts, type IndexedPackageRef } from "./grouping.js";
|
import { processAlbumGroups, detectGroupingConflicts, type IndexedPackageRef } from "./grouping.js";
|
||||||
@@ -314,6 +315,7 @@ interface PipelineContext {
|
|||||||
zipsFound: number;
|
zipsFound: number;
|
||||||
zipsDuplicate: number;
|
zipsDuplicate: number;
|
||||||
zipsIngested: number;
|
zipsIngested: number;
|
||||||
|
zipsBackfilled: number;
|
||||||
};
|
};
|
||||||
/** Creator from forum topic name (null for non-forum). */
|
/** Creator from forum topic name (null for non-forum). */
|
||||||
topicCreator: string | null;
|
topicCreator: string | null;
|
||||||
@@ -422,6 +424,7 @@ export async function runWorkerForAccount(
|
|||||||
zipsFound: 0,
|
zipsFound: 0,
|
||||||
zipsDuplicate: 0,
|
zipsDuplicate: 0,
|
||||||
zipsIngested: 0,
|
zipsIngested: 0,
|
||||||
|
zipsBackfilled: 0,
|
||||||
};
|
};
|
||||||
|
|
||||||
try {
|
try {
|
||||||
@@ -1645,6 +1648,56 @@ async function processOneArchiveSet(
|
|||||||
return null;
|
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 ──
|
// ── Size guard: skip archives that exceed WORKER_MAX_ZIP_SIZE_MB ──
|
||||||
const maxSizeBytes = BigInt(config.maxZipSizeMB) * 1024n * 1024n;
|
const maxSizeBytes = BigInt(config.maxZipSizeMB) * 1024n * 1024n;
|
||||||
if (totalArchiveSize > maxSizeBytes) {
|
if (totalArchiveSize > maxSizeBytes) {
|
||||||
|
|||||||
Reference in New Issue
Block a user