diff --git a/public/index.html b/public/index.html index 69d77ef..0ffc024 100644 --- a/public/index.html +++ b/public/index.html @@ -93,7 +93,29 @@ } .table > :not(caption) > * > * { padding: 1rem; } - + + .status-dot { + width: 8px; + height: 8px; + border-radius: 999px; + display: inline-block; + margin-right: 8px; + background: #10b981; + box-shadow: 0 0 0 3px rgba(16, 185, 129, 0.15); + } + .status-queued { + background: #f59e0b; + box-shadow: 0 0 0 3px rgba(245, 158, 11, 0.15); + } + .status-active { + background: #10b981; + box-shadow: 0 0 0 3px rgba(16, 185, 129, 0.15); + } + .status-backfilling { + background: #f97316; + box-shadow: 0 0 0 3px rgba(249, 115, 22, 0.18); + } + .tweet-preview { max-width: 300px; white-space: nowrap; @@ -213,8 +235,8 @@ useEffect(() => { if (view !== 'dashboard' || !token) return; - const statusTimer = setInterval(fetchStatus, 5000); - const activityTimer = setInterval(fetchActivity, 10000); + const statusTimer = setInterval(fetchStatus, 2000); + const activityTimer = setInterval(fetchActivity, 7000); return () => { clearInterval(statusTimer); clearInterval(activityTimer); @@ -327,7 +349,13 @@ } }; - const runBackfill = async (id) => { + const runBackfill = async (id, label = 'Backfill') => { + const hasQueue = status.pendingBackfills && status.pendingBackfills.length > 0; + const isActiveBackfill = status.currentStatus?.state === 'backfilling'; + if (hasQueue || isActiveBackfill) { + const confirmMsg = `${label} is already running or queued. This will add a new request and replace any existing one for this account. Continue?`; + if (!confirm(confirmMsg)) return; + } const limit = prompt(`How many tweets to backfill per account?`, "15"); if (limit === null) return; try { @@ -349,6 +377,12 @@ }; const resetAndBackfill = async (id) => { + const hasQueue = status.pendingBackfills && status.pendingBackfills.length > 0; + const isActiveBackfill = status.currentStatus?.state === 'backfilling'; + if (hasQueue || isActiveBackfill) { + const confirmMsg = 'Backfill is already running or queued. This will add a new request and replace any existing one for this account. Continue?'; + if (!confirm(confirmMsg)) return; + } const limit = prompt(`Reset cache and backfill how many tweets?`, "15"); if (limit === null) return; try { @@ -484,6 +518,9 @@ } const isBackfillQueued = (id) => status.pendingBackfills?.some(b => (b.id || b) === id); + const backfillEntry = (id) => status.pendingBackfills?.find(b => (b.id || b) === id); + const activeBackfillId = status.currentStatus?.backfillMappingId; + const isBackfillActive = (id) => status.currentStatus?.state === 'backfilling' && activeBackfillId === id; // Check if configs are set to collapse by default const hasTwitterConfig = twitterConfig.authToken && twitterConfig.ct0; @@ -601,8 +638,14 @@ {m.bskyIdentifier}
- - {isBackfillQueued(m.id) ? 'Backfilling' : 'Active'} + + {isBackfillActive(m.id) ? ( + Backfilling + ) : isBackfillQueued(m.id) ? ( + Queued #{backfillEntry(m.id)?.position || ''} + ) : ( + Active + )}
@@ -614,7 +657,7 @@ {isAdmin && ( <>
  • -
  • +

  • diff --git a/src/index.ts b/src/index.ts index f1bb13c..2efe658 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1422,7 +1422,14 @@ async function processTweets( import { getAgent } from './bsky.js'; -async function importHistory(twitterUsername: string, bskyIdentifier: string, limit = 15, dryRun = false, ignoreCancellation = false): Promise { +async function importHistory( + twitterUsername: string, + bskyIdentifier: string, + limit = 15, + dryRun = false, + ignoreCancellation = false, + requestId?: string, +): Promise { const config = getConfig(); const mapping = config.mappings.find((m) => m.twitterUsernames.map(u => u.toLowerCase()).includes(twitterUsername.toLowerCase())); if (!mapping) { @@ -1468,11 +1475,12 @@ async function importHistory(twitterUsername: string, bskyIdentifier: string, li for await (const scraperTweet of generator) { if (!ignoreCancellation) { - const stillPending = getPendingBackfills().some(b => b.id === mapping.id); - if (!stillPending) { - console.log(`[${twitterUsername}] 🛑 Backfill cancelled.`); - break; - } + const stillPending = getPendingBackfills().some(b => b.id === mapping.id && (!requestId || b.requestId === requestId)); + if (!stillPending) { + console.log(`[${twitterUsername}] 🛑 Backfill cancelled.`); + break; + } + } const t = mapScraperTweetToLocalTweet(scraperTweet); @@ -1501,7 +1509,7 @@ async function importHistory(twitterUsername: string, bskyIdentifier: string, li // Task management const activeTasks = new Map>(); -async function runAccountTask(mapping: AccountMapping, forceBackfill = false, dryRun = false) { +async function runAccountTask(mapping: AccountMapping, backfillRequest?: PendingBackfill, dryRun = false) { if (activeTasks.has(mapping.id)) return; // Already running const task = (async () => { @@ -1509,23 +1517,52 @@ async function runAccountTask(mapping: AccountMapping, forceBackfill = false, dr const agent = await getAgent(mapping); if (!agent) return; - const backfillReq = getPendingBackfills().find(b => b.id === mapping.id); + const backfillReq = backfillRequest ?? getPendingBackfills().find(b => b.id === mapping.id); - if (forceBackfill || backfillReq) { - const limit = backfillReq?.limit || 15; + if (backfillReq) { + const limit = backfillReq.limit || 15; console.log(`[${mapping.bskyIdentifier}] Running backfill for ${mapping.twitterUsernames.length} accounts (limit ${limit})...`); + updateAppStatus({ + state: 'backfilling', + currentAccount: mapping.twitterUsernames[0], + message: `Starting backfill (limit ${limit})...`, + backfillMappingId: mapping.id, + backfillRequestId: backfillReq.requestId, + }); for (const twitterUsername of mapping.twitterUsernames) { + const stillPending = getPendingBackfills().some( + (b) => b.id === mapping.id && b.requestId === backfillReq.requestId, + ); + if (!stillPending) { + console.log(`[${mapping.bskyIdentifier}] 🛑 Backfill request replaced; stopping.`); + break; + } + try { - updateAppStatus({ state: 'backfilling', currentAccount: twitterUsername, message: `Starting backfill (limit ${limit})...` }); - await importHistory(twitterUsername, mapping.bskyIdentifier, limit, dryRun); + updateAppStatus({ + state: 'backfilling', + currentAccount: twitterUsername, + message: `Starting backfill (limit ${limit})...`, + backfillMappingId: mapping.id, + backfillRequestId: backfillReq.requestId, + }); + await importHistory(twitterUsername, mapping.bskyIdentifier, limit, dryRun, false, backfillReq.requestId); } catch (err) { console.error(`❌ Error backfilling ${twitterUsername}:`, err); } } - clearBackfill(mapping.id); + clearBackfill(mapping.id, backfillReq.requestId); + updateAppStatus({ + state: 'idle', + message: `Backfill complete for ${mapping.bskyIdentifier}`, + backfillMappingId: undefined, + backfillRequestId: undefined, + }); console.log(`[${mapping.bskyIdentifier}] Backfill complete.`); } else { + updateAppStatus({ backfillMappingId: undefined, backfillRequestId: undefined }); + // Pre-load processed IDs for optimization const processedMap = loadProcessedTweets(mapping.bskyIdentifier); const processedIds = new Set(Object.keys(processedMap)); @@ -1533,7 +1570,13 @@ async function runAccountTask(mapping: AccountMapping, forceBackfill = false, dr for (const twitterUsername of mapping.twitterUsernames) { try { console.log(`[${twitterUsername}] 🏁 Starting check for new tweets...`); - updateAppStatus({ state: 'checking', currentAccount: twitterUsername, message: 'Fetching latest tweets...' }); + updateAppStatus({ + state: 'checking', + currentAccount: twitterUsername, + message: 'Fetching latest tweets...', + backfillMappingId: undefined, + backfillRequestId: undefined, + }); // Use fetchUserTweets with early stopping optimization // Increase limit slightly since we have early stopping now @@ -1570,6 +1613,7 @@ import { getNextCheckTime, updateAppStatus, } from './server.js'; +import type { PendingBackfill } from './server.js'; import { AccountMapping } from './config-manager.js'; async function main(): Promise { @@ -1631,44 +1675,65 @@ async function main(): Promise { // Concurrency limit for processing accounts const runLimit = pLimit(3); + const findMappingById = (mappings: AccountMapping[], id: string) => + mappings.find((mapping) => mapping.id === id); + // Main loop while (true) { const now = Date.now(); const config = getConfig(); // Reload config to get new mappings/settings const nextTime = getNextCheckTime(); - + // Check if it's time for a scheduled run OR if we have pending backfills const isScheduledRun = now >= nextTime; const pendingBackfills = getPendingBackfills(); - + if (isScheduledRun) { - console.log(`[${new Date().toISOString()}] ⏰ Scheduled check triggered.`); - updateLastCheckTime(); + console.log(`[${new Date().toISOString()}] ⏰ Scheduled check triggered.`); + updateLastCheckTime(); } const tasks: Promise[] = []; - for (const mapping of config.mappings) { - if (!mapping.enabled) continue; - - const hasPendingBackfill = pendingBackfills.some(b => b.id === mapping.id); - - // Run if scheduled OR backfill requested - if (isScheduledRun || hasPendingBackfill) { - // Queue task with concurrency limit - tasks.push(runLimit(async () => { - await runAccountTask(mapping, hasPendingBackfill, options.dryRun); - })); + if (pendingBackfills.length > 0) { + const [nextBackfill, ...rest] = pendingBackfills; + if (nextBackfill) { + const mapping = findMappingById(config.mappings, nextBackfill.id); + if (mapping && mapping.enabled) { + console.log(`[Scheduler] 🚧 Backfill priority: ${mapping.bskyIdentifier}`); + await runAccountTask(mapping, nextBackfill, options.dryRun); + } else { + clearBackfill(nextBackfill.id, nextBackfill.requestId); } - } + } + if (pendingBackfills.length === 0 && getPendingBackfills().length === 0) { + updateAppStatus({ + state: 'idle', + message: 'Backfill queue empty', + backfillMappingId: undefined, + backfillRequestId: undefined, + }); + } + nextCheckTime = Date.now() + (config.checkIntervalMinutes || 5) * 60 * 1000; + } else if (isScheduledRun) { + for (const mapping of config.mappings) { + if (!mapping.enabled) continue; + + tasks.push(runLimit(async () => { + await runAccountTask(mapping, undefined, options.dryRun); + })); + } - if (tasks.length > 0) { + if (tasks.length > 0) { await Promise.all(tasks); console.log(`[Scheduler] ✅ All tasks for this cycle complete.`); + } + + updateAppStatus({ state: 'idle', message: 'Scheduled checks complete' }); } - + // Sleep for 5 seconds - await new Promise(resolve => setTimeout(resolve, 5000)); + await new Promise((resolve) => setTimeout(resolve, 5000)); } } diff --git a/src/server.ts b/src/server.ts index 34504ad..478937b 100644 --- a/src/server.ts +++ b/src/server.ts @@ -18,11 +18,15 @@ const JWT_SECRET = process.env.JWT_SECRET || 'fallback-secret'; // In-memory state for triggers and scheduling let lastCheckTime = Date.now(); let nextCheckTime = Date.now() + (getConfig().checkIntervalMinutes || 5) * 60 * 1000; -interface PendingBackfill { +export interface PendingBackfill { id: string; limit?: number; + queuedAt: number; + sequence: number; + requestId: string; } let pendingBackfills: PendingBackfill[] = []; +let backfillSequence = 0; interface AppStatus { state: 'idle' | 'checking' | 'backfilling' | 'pacing' | 'processing'; @@ -30,6 +34,8 @@ interface AppStatus { processedCount?: number; totalCount?: number; message?: string; + backfillMappingId?: string; + backfillRequestId?: string; lastUpdate: number; } @@ -285,7 +291,13 @@ app.get('/api/status', authenticateToken, (_req, res) => { nextCheckTime, nextCheckMinutes: Math.ceil(nextRunMs / 60000), checkIntervalMinutes: config.checkIntervalMinutes, - pendingBackfills, + pendingBackfills: pendingBackfills + .slice() + .sort((a, b) => a.sequence - b.sequence) + .map((backfill, index) => ({ + ...backfill, + position: index + 1, + })), currentStatus: currentAppStatus, }); }); @@ -307,12 +319,25 @@ app.post('/api/backfill/:id', authenticateToken, requireAdmin, (req, res) => { return; } - if (!pendingBackfills.find((b) => b.id === id)) { - pendingBackfills.push({ id, limit: limit ? Number(limit) : undefined }); - } + const queuedAt = Date.now(); + const sequence = backfillSequence++; + const requestId = Math.random().toString(36).slice(2); + pendingBackfills = pendingBackfills.filter((b) => b.id !== id); + pendingBackfills.push({ + id, + limit: limit ? Number(limit) : undefined, + queuedAt, + sequence, + requestId, + }); + pendingBackfills.sort((a, b) => a.sequence - b.sequence); // Do not force a global run; the scheduler loop will pick up the pending backfill in ~5s - res.json({ success: true, message: `Backfill queued for @${mapping.twitterUsernames.join(', ')}` }); + res.json({ + success: true, + message: `Backfill queued for @${mapping.twitterUsernames.join(', ')}`, + requestId, + }); }); app.delete('/api/backfill/:id', authenticateToken, (req, res) => { @@ -323,6 +348,12 @@ app.delete('/api/backfill/:id', authenticateToken, (req, res) => { app.post('/api/backfill/clear-all', authenticateToken, requireAdmin, (_req, res) => { pendingBackfills = []; + updateAppStatus({ + state: 'idle', + message: 'All backfills cleared', + backfillMappingId: undefined, + backfillRequestId: undefined, + }); res.json({ success: true, message: 'All backfills cleared' }); }); @@ -392,14 +423,18 @@ export function updateAppStatus(status: Partial) { } export function getPendingBackfills(): PendingBackfill[] { - return [...pendingBackfills]; + return [...pendingBackfills].sort((a, b) => a.sequence - b.sequence); } export function getNextCheckTime(): number { return nextCheckTime; } -export function clearBackfill(id: string) { +export function clearBackfill(id: string, requestId?: string) { + if (requestId) { + pendingBackfills = pendingBackfills.filter((bid) => !(bid.id === id && bid.requestId === requestId)); + return; + } pendingBackfills = pendingBackfills.filter((bid) => bid.id !== id); }