From 23f8e91c5042bf18e87ca4387f8908a2fedba21a Mon Sep 17 00:00:00 2001 From: xCyanGrizzly Date: Mon, 22 Jun 2026 12:10:56 +0200 Subject: [PATCH] feat(telegram): add "skip & disable topic" on the live worker panel Track the forum topic currently being processed on IngestionRun (currentTopicId + currentAccountChannelMapId; additive migration) and expose it on the live status. processArchiveSets gains an optional shouldStop callback polled before each archive set; the forum branch passes a live isTopicFetchEnabled check, so disabling a topic mid-run lets the in-flight file finish, then skips the rest of that topic. A new disableActiveTopic server action sets the topic's fetchEnabled false (upsert), and the worker status panel shows a "Skip & disable this topic" button while a topic is being processed. Future runs skip the topic via the existing live per-topic read. Co-Authored-By: Claude Opus 4.8 --- .../migration.sql | 6 +++ prisma/schema.prisma | 2 + .../_components/worker-status-panel.tsx | 42 +++++++++++++++++- src/app/(app)/telegram/actions.ts | 43 +++++++++++++++++++ src/lib/telegram/queries.ts | 2 + src/lib/telegram/types.ts | 2 + worker/src/db/queries.ts | 8 ++++ worker/src/worker.ts | 29 ++++++++++++- 8 files changed, 131 insertions(+), 3 deletions(-) create mode 100644 prisma/migrations/20260622120000_add_run_current_topic/migration.sql diff --git a/prisma/migrations/20260622120000_add_run_current_topic/migration.sql b/prisma/migrations/20260622120000_add_run_current_topic/migration.sql new file mode 100644 index 0000000..aa6a42d --- /dev/null +++ b/prisma/migrations/20260622120000_add_run_current_topic/migration.sql @@ -0,0 +1,6 @@ +-- AlterTable: track the forum topic currently being processed on the live run +-- so the worker status panel can offer a "skip & disable topic" action. +-- Additive, nullable — no data change for existing rows. +ALTER TABLE "ingestion_runs" + ADD COLUMN "currentTopicId" BIGINT, + ADD COLUMN "currentAccountChannelMapId" TEXT; diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 2e26e8d..679a74e 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -582,6 +582,8 @@ model IngestionRun { totalBytes BigInt? // Total size of current download downloadPercent Int? // 0-100 lastActivityAt DateTime? // When activity was last updated + currentTopicId BigInt? // Forum topic currently being processed (live "skip topic") + currentAccountChannelMapId String? // AccountChannelMap owning the topic being processed account TelegramAccount @relation(fields: [accountId], references: [id]) packages Package[] diff --git a/src/app/(app)/telegram/_components/worker-status-panel.tsx b/src/app/(app)/telegram/_components/worker-status-panel.tsx index 215ecc3..580a3a0 100644 --- a/src/app/(app)/telegram/_components/worker-status-panel.tsx +++ b/src/app/(app)/telegram/_components/worker-status-panel.tsx @@ -9,13 +9,14 @@ import { Radio, AlertTriangle, RefreshCw, + SkipForward, } from "lucide-react"; import { Card, CardContent } from "@/components/ui/card"; import { Badge } from "@/components/ui/badge"; import { Button } from "@/components/ui/button"; import { cn } from "@/lib/utils"; import { toast } from "sonner"; -import { triggerIngestion } from "../actions"; +import { triggerIngestion, disableActiveTopic } from "../actions"; import type { IngestionAccountStatus } from "@/lib/telegram/types"; interface WorkerStatusPanelProps { @@ -218,6 +219,24 @@ function RunningStatus({ }: { run: NonNullable; }) { + const [isDisabling, startDisable] = useTransition(); + + const handleSkipTopic = () => { + const acmId = run.currentAccountChannelMapId; + const topicId = run.currentTopicId; + if (!acmId || !topicId) return; + startDisable(async () => { + const result = await disableActiveTopic(acmId, topicId); + if (result.success) { + toast.success( + "Topic disabled — the current file finishes, then the rest is skipped" + ); + } else { + toast.error(result.error ?? "Failed to disable topic"); + } + }); + }; + return (
@@ -273,6 +292,27 @@ function RunningStatus({ )}
+ + {/* Skip & disable the topic currently being processed */} + {run.currentTopicId && run.currentAccountChannelMapId && ( +
+ +
+ )}
); } diff --git a/src/app/(app)/telegram/actions.ts b/src/app/(app)/telegram/actions.ts index 2ba86ab..dac1794 100644 --- a/src/app/(app)/telegram/actions.ts +++ b/src/app/(app)/telegram/actions.ts @@ -462,6 +462,49 @@ export async function setTopicFetchEnabled( } } +/** + * Disable the topic currently being processed by the worker, identified by its + * account-channel mapping + Telegram topic id (as exposed on the live run + * status). Upserts so a disabled row exists even if the worker hasn't persisted + * the topic yet. The worker honours this live: it finishes the in-flight file + * then skips the rest of the topic, and future runs skip it entirely. + */ +export async function disableActiveTopic( + accountChannelMapId: string, + topicId: string +): Promise { + const admin = await requireAdmin(); + if (!admin.success) return admin; + + let topicIdBig: bigint; + try { + topicIdBig = BigInt(topicId); + } catch { + return { success: false, error: "Invalid topic id" }; + } + + try { + await prisma.topicProgress.upsert({ + where: { + accountChannelMapId_topicId: { + accountChannelMapId, + topicId: topicIdBig, + }, + }, + create: { + accountChannelMapId, + topicId: topicIdBig, + fetchEnabled: false, + }, + update: { fetchEnabled: false }, + }); + revalidatePath(REVALIDATE_PATH); + return { success: true, data: undefined }; + } catch { + return { success: false, error: "Failed to disable topic" }; + } +} + // ── Account-Channel link actions ── export async function linkChannel( diff --git a/src/lib/telegram/queries.ts b/src/lib/telegram/queries.ts index 898a016..fbf6c3a 100644 --- a/src/lib/telegram/queries.ts +++ b/src/lib/telegram/queries.ts @@ -550,6 +550,8 @@ export async function getIngestionStatus(): Promise { totalBytes: currentRun.totalBytes?.toString() ?? null, downloadPercent: currentRun.downloadPercent, lastActivityAt: currentRun.lastActivityAt?.toISOString() ?? null, + currentTopicId: currentRun.currentTopicId?.toString() ?? null, + currentAccountChannelMapId: currentRun.currentAccountChannelMapId, } : null, }); diff --git a/src/lib/telegram/types.ts b/src/lib/telegram/types.ts index c36185c..a0f0134 100644 --- a/src/lib/telegram/types.ts +++ b/src/lib/telegram/types.ts @@ -123,5 +123,7 @@ export interface IngestionAccountStatus { totalBytes: string | null; // BigInt serialized as string downloadPercent: number | null; lastActivityAt: string | null; + currentTopicId: string | null; // BigInt serialized as string + currentAccountChannelMapId: string | null; } | null; } diff --git a/worker/src/db/queries.ts b/worker/src/db/queries.ts index 3ae23f8..09590ad 100644 --- a/worker/src/db/queries.ts +++ b/worker/src/db/queries.ts @@ -369,6 +369,8 @@ export interface ActivityUpdate { downloadedBytes?: bigint | null; totalBytes?: bigint | null; downloadPercent?: number | null; + currentTopicId?: bigint | null; + currentAccountChannelMapId?: string | null; messagesScanned?: number; zipsFound?: number; zipsDuplicate?: number; @@ -396,6 +398,10 @@ export async function updateRunActivity( ...(activity.zipsFound !== undefined && { zipsFound: activity.zipsFound }), ...(activity.zipsDuplicate !== undefined && { zipsDuplicate: activity.zipsDuplicate }), ...(activity.zipsIngested !== undefined && { zipsIngested: activity.zipsIngested }), + ...(activity.currentTopicId !== undefined && { currentTopicId: activity.currentTopicId }), + ...(activity.currentAccountChannelMapId !== undefined && { + currentAccountChannelMapId: activity.currentAccountChannelMapId, + }), }, }); } @@ -410,6 +416,8 @@ const CLEAR_ACTIVITY = { downloadedBytes: null, totalBytes: null, downloadPercent: null, + currentTopicId: null, + currentAccountChannelMapId: null, lastActivityAt: new Date(), }; diff --git a/worker/src/worker.ts b/worker/src/worker.ts index fcd1be1..793ba0b 100644 --- a/worker/src/worker.ts +++ b/worker/src/worker.ts @@ -538,6 +538,8 @@ export async function runWorkerForAccount( currentActivity: `Enumerating topics in "${channelLabel}"`, currentStep: "scanning", currentChannel: channelLabel, + currentTopicId: null, + currentAccountChannelMapId: null, currentFile: null, currentFileNum: null, totalFiles: null, @@ -750,6 +752,8 @@ export async function runWorkerForAccount( currentActivity: `Scanning "${topicLabel}"${topicProgress}`, currentStep: "scanning", currentChannel: channelLabel, + currentTopicId: topic.topicId, + currentAccountChannelMapId: mapping.id, currentFile: null, currentFileNum: null, totalFiles: null, @@ -830,7 +834,11 @@ export async function runWorkerForAccount( topic.name, messageId ); - } + }, + // shouldStop: re-read the live fetch flag before each archive + // set. A mid-run "disable topic" lets the current file finish, + // then skips the rest of this topic's archives. + async () => !(await isTopicFetchEnabled(mapping.id, topic.topicId)) ); // Sync client back in case it was recreated during upload stall recovery client = pipelineCtx.client; @@ -922,6 +930,8 @@ export async function runWorkerForAccount( currentActivity: `Scanning "${channelLabel}" for new archives`, currentStep: "scanning", currentChannel: channelLabel, + currentTopicId: null, + currentAccountChannelMapId: null, currentFile: null, currentFileNum: null, totalFiles: null, @@ -1203,7 +1213,12 @@ async function processArchiveSets( * below any failed message ID in this scan). Used by the caller to * advance the channel/topic watermark incrementally — otherwise a long * scan that gets killed by worker restart loses all progress. */ - onWatermarkAdvance?: (messageId: bigint) => Promise + onWatermarkAdvance?: (messageId: bigint) => Promise, + /** Optional cancellation check, polled before each archive set. When it + * resolves true, processing stops after the set currently in flight (that + * one completes; remaining sets in this scan are skipped). Used by the + * forum branch to honour a mid-run "disable topic". */ + shouldStop?: () => Promise ): Promise<{ maxProcessedId: bigint | null; minFailedId: bigint | null }> { const { client, runId, channelTitle, channel, throttled, counters, accountLog } = ctx; @@ -1281,6 +1296,16 @@ async function processArchiveSets( const indexedPackageRefs: IndexedPackageRef[] = []; for (let setIdx = 0; setIdx < archiveSets.length; setIdx++) { + // Cooperative cancellation: if the caller signals stop (e.g. the topic was + // disabled mid-run), skip the remaining archive sets in this scan. The set + // processed in the previous iteration has already completed. + if (shouldStop && (await shouldStop())) { + accountLog.info( + { channel: channelTitle, processed: setIdx, total: archiveSets.length }, + "Stop signal received (topic disabled) — skipping remaining archive sets in this scan" + ); + break; + } try { const packageId = await processOneArchiveSet( ctx,