diff --git a/worker/src/db/queries.ts b/worker/src/db/queries.ts index 09590ad..ab39cbc 100644 --- a/worker/src/db/queries.ts +++ b/worker/src/db/queries.ts @@ -1,5 +1,6 @@ import { db } from "./client.js"; import type { ArchiveType, FetchStatus } from "@prisma/client"; +import type { FileEntry } from "../archive/zip-reader.js"; export async function getActiveAccounts() { return db.telegramAccount.findMany({ @@ -1005,3 +1006,91 @@ export async function createAutoGroup(input: { return group.id; } + +// ── Provenance backfill ── + +export async function findPlaceholderCandidate( + destChannelId: string, + fileName: string, + fileSize: bigint, +): Promise<{ id: string; archiveType: string; fileCount: number } | null> { + return db.package.findFirst({ + where: { + fileName, + fileSize, + destMessageId: { not: null }, + // Placeholder provenance (spec §1): manual-upload (source == destination) + // OR rebuild record (sourceMessageId == 0 "unknown" sentinel). + OR: [ + { sourceChannelId: destChannelId }, + { sourceMessageId: 0n }, + ], + }, + select: { id: true, archiveType: true, fileCount: true }, + orderBy: { indexedAt: "asc" }, + }); +} + +export async function getPackageFileCrcs(packageId: string): Promise<(string | null)[]> { + const rows = await db.packageFile.findMany({ + where: { packageId }, + select: { crc32: true }, + }); + return rows.map((r) => r.crc32); +} + +export interface BackfillProvenanceInput { + packageId: string; + destChannelId: string; // to re-check placeholder status in-txn + sourceChannelId: string; + sourceMessageId: bigint; + sourceTopicId: bigint | null; + sourceCaption: string | null; + remoteUniqueId: string | null; + creator: string | null; // always set (re-derived by caller) + entries?: FileEntry[]; // set only if candidate had fileCount === 0 + previewData?: Buffer | null; // set only if provided and candidate lacks one + previewMsgId?: bigint | null; +} + +export async function backfillProvenance(input: BackfillProvenanceInput): Promise { + return db.$transaction(async (tx) => { + const current = await tx.package.findUnique({ + where: { id: input.packageId }, + select: { sourceChannelId: true, sourceMessageId: true, previewData: true, fileCount: true }, + }); + // Re-check placeholder status inside the txn (another worker may have won). + // Placeholder = manual-upload (source==dest) OR rebuild (sourceMessageId==0). + const stillPlaceholder = + !!current && + (current.sourceChannelId === input.destChannelId || current.sourceMessageId === 0n); + if (!stillPlaceholder) return false; + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const data: any = { + sourceChannelId: input.sourceChannelId, + sourceMessageId: input.sourceMessageId, + sourceTopicId: input.sourceTopicId, + sourceCaption: input.sourceCaption, + remoteUniqueId: input.remoteUniqueId, + creator: input.creator, + }; + if (input.entries && current.fileCount === 0) { + await tx.packageFile.deleteMany({ where: { packageId: input.packageId } }); + await tx.packageFile.createMany({ + data: input.entries.map((e) => ({ + packageId: input.packageId, + path: e.path, fileName: e.fileName, extension: e.extension, + compressedSize: e.compressedSize, uncompressedSize: e.uncompressedSize, crc32: e.crc32, + })), + }); + data.fileCount = input.entries.length; + } + if (input.previewData && !current.previewData) { + data.previewData = input.previewData; + data.previewMsgId = input.previewMsgId ?? null; + } + await tx.package.update({ where: { id: input.packageId }, data }); + return true; + }); +}