From 735e874bd03c55c626347f5356301fe221145b98 Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Wed, 22 Jul 2026 09:39:03 +0200 Subject: [PATCH] fix(campaign): bound operator queue polling --- .../features/operator/OperatorQueuePage.tsx | 88 ++++++++++++++----- .../features/operator/operatorQueueModel.ts | 28 ++++++ webui/tests/operator-queue-model.test.ts | 37 ++++++++ .../operator-queue-ui-structure.test.mjs | 4 + 4 files changed, 136 insertions(+), 21 deletions(-) diff --git a/webui/src/features/operator/OperatorQueuePage.tsx b/webui/src/features/operator/OperatorQueuePage.tsx index 0db4b4f..865a19f 100644 --- a/webui/src/features/operator/OperatorQueuePage.tsx +++ b/webui/src/features/operator/OperatorQueuePage.tsx @@ -46,7 +46,9 @@ import type { CampaignJobSortColumn } from "../campaigns/utils/jobListQuery"; import { campaignLifecycleTotals, campaignVersionWorkItems, + operatorDeliveryStateCounts, operatorQueueActionBlocks, + shouldRefreshOperatorVersionSummary, type OperatorQueueAction, type OperatorQueuePermissions } from "./operatorQueueModel"; @@ -123,7 +125,9 @@ export default function OperatorQueuePage({ settings, auth }: {settings: ApiSett const summariesRef = useRef>>({}); const versionsRef = useRef>({}); const loadingRef = useRef(false); + const pendingFullDiscoveryRef = useRef(false); const backgroundLoadRef = useRef<() => Promise>(async () => undefined); + const fullDiscoveryLoadRef = useRef<() => Promise>(async () => undefined); const [loading, setLoading] = useState(true); const [error, setError] = useState(""); const [message, setMessage] = useState(""); @@ -199,11 +203,12 @@ export default function OperatorQueuePage({ settings, auth }: {settings: ApiSett summariesRef.current = {}; versionsRef.current = {}; resetDeltaWatermark(); - void load(); + void load(false, true); }, [settingsKey, resetDeltaWatermark]); useEffect(() => { - backgroundLoadRef.current = () => load(true); + backgroundLoadRef.current = () => load(true, false); + fullDiscoveryLoadRef.current = () => load(true, true); backgroundJobsLoadRef.current = () => loadSelectedJobs(true); }); @@ -238,34 +243,41 @@ export default function OperatorQueuePage({ settings, auth }: {settings: ApiSett selectedRow && (selectedRow.queuedOrActive > 0 || selectedRow.needsAttention > 0) ) || selectedJobsActive; useEffect(() => { - if (!hasActiveDelivery && !selectedQueueActive) return; + if (!hasActiveDelivery && !selectedRow) return; const handle = window.setInterval(() => { void backgroundLoadRef.current(); if (selectedQueueActive) void backgroundJobsLoadRef.current(); }, 10_000); return () => window.clearInterval(handle); - }, [hasActiveDelivery, selectedQueueActive]); + }, [hasActiveDelivery, selectedQueueActive, selectedRow]); useEffect(() => { const handle = window.setInterval(() => { - void backgroundLoadRef.current(); + void fullDiscoveryLoadRef.current(); }, 60_000); return () => window.clearInterval(handle); }, []); - async function load(background = false) { - if (loadingRef.current) return; + async function load(background = false, fullDiscovery = !background) { + if (loadingRef.current) { + if (fullDiscovery) pendingFullDiscoveryRef.current = true; + return; + } loadingRef.current = true; if (!background) setLoading(true); setError(""); try { const campaigns = await loadCampaignsDelta(); - setRows(await loadCampaignVersionRows(campaigns)); + setRows(await loadCampaignVersionRows(campaigns, fullDiscovery)); } catch (err) { setError(err instanceof Error ? err.message : String(err)); } finally { if (!background) setLoading(false); loadingRef.current = false; + if (pendingFullDiscoveryRef.current) { + pendingFullDiscoveryRef.current = false; + void fullDiscoveryLoadRef.current(); + } } } @@ -296,7 +308,7 @@ export default function OperatorQueuePage({ settings, auth }: {settings: ApiSett } async function refreshAll() { - await load(); + await load(false, true); await backgroundJobsLoadRef.current(); } @@ -363,13 +375,14 @@ export default function OperatorQueuePage({ settings, auth }: {settings: ApiSett return campaigns; } - async function loadCampaignVersionRows(campaigns: CampaignListItem[]): Promise { + async function loadCampaignVersionRows(campaigns: CampaignListItem[], fullDiscovery: boolean): Promise { const summaries = { ...summariesRef.current }; const versions = { ...versionsRef.current }; await Promise.all(campaigns.map(async (campaign) => { const key = operatorCampaignWorkspaceKey(campaign.id); let nextWatermark = getDeltaWatermark(key); let campaignVersions = versions[campaign.id] ?? []; + const changedVersionIds = new Set(); let hasMore = false; try { do { @@ -379,6 +392,7 @@ export default function OperatorQueuePage({ settings, auth }: {settings: ApiSett includeSummary: false, since: nextWatermark }); + for (const version of response.versions) changedVersionIds.add(version.id); campaignVersions = response.full ? response.versions : mergeDeltaRows( @@ -398,7 +412,33 @@ export default function OperatorQueuePage({ settings, auth }: {settings: ApiSett versions[campaign.id] = campaignVersions; const campaignSummaries = { ...(summaries[campaign.id] ?? {}) }; - await Promise.all(campaignVersions.map(async (version) => { + const summariesToRefresh = campaignVersions.filter((version) => { + const cached = campaignSummaries[version.id] ?? null; + const cards = asRecord(cached?.cards); + const statusCounts = asRecord(cached?.status_counts); + const queueCounts = asRecord(statusCounts.queue); + const sendCounts = asRecord(statusCounts.send); + const deliveryState = operatorDeliveryStateCounts( + { queued: numberValue(queueCounts.queued) }, + { + claimed: numberValue(sendCounts.claimed), + sending: numberValue(sendCounts.sending) + } + ); + return shouldRefreshOperatorVersionSummary({ + fullDiscovery, + versionChanged: changedVersionIds.has(version.id), + selected: selectedCampaignId === campaign.id && ( + selectedVersionParam + ? selectedVersionParam === version.id + : campaign.current_version_id === version.id + ), + cached: cached !== null, + queuedOrActive: deliveryState.queuedOrActive, + needsAttention: numberValue(cards.needs_attention) + }); + }); + await Promise.all(summariesToRefresh.map(async (version) => { try { campaignSummaries[version.id] = await getCampaignSummary(settings, campaign.id, version.id); } catch { @@ -847,10 +887,15 @@ function toRow( const statusCounts = asRecord(summary?.status_counts); const queueCounts = asRecord(statusCounts.queue); const sendCounts = asRecord(statusCounts.send); - const queued = numberValue(sendCounts.queued); - const claimed = numberValue(sendCounts.claimed); - const sending = numberValue(sendCounts.sending); - const completed = numberValue(sendCounts.smtp_accepted) + numberValue(sendCounts.sent); + const deliveryState = operatorDeliveryStateCounts( + { queued: numberValue(queueCounts.queued) }, + { + claimed: numberValue(sendCounts.claimed), + sending: numberValue(sendCounts.sending), + smtpAccepted: numberValue(sendCounts.smtp_accepted), + sent: numberValue(sendCounts.sent) + } + ); const failed = numberValue(cards.failed); const notAttempted = numberValue(cards.not_attempted); return { @@ -864,11 +909,11 @@ function toRow( retryable: numberValue(cards.retryable), outcomeUnknown: numberValue(cards.outcome_unknown), notAttempted, - queuedOrActive: queued + claimed + sending, - queued, - claimed, - sending, - completed, + queuedOrActive: deliveryState.queuedOrActive, + queued: deliveryState.queued, + claimed: deliveryState.claimed, + sending: deliveryState.sending, + completed: deliveryState.completed, pausable: numberValue(queueCounts.queued), paused: numberValue(queueCounts.paused), campaignPausable: lifecycle.pausable, @@ -982,7 +1027,8 @@ function dataGridQueriesEqual(left: DataGridQueryState, right: DataGridQueryStat } function isActiveJob(job: Record): boolean { - return ["queued", "claimed", "sending"].includes(String(job.send_status ?? "")); + return String(job.queue_status ?? "") === "queued" + || ["claimed", "sending"].includes(String(job.send_status ?? "")); } function deliveryStatusOptionLabel(value: string): string { diff --git a/webui/src/features/operator/operatorQueueModel.ts b/webui/src/features/operator/operatorQueueModel.ts index b98ff7d..bcc7995 100644 --- a/webui/src/features/operator/operatorQueueModel.ts +++ b/webui/src/features/operator/operatorQueueModel.ts @@ -70,6 +70,34 @@ export function campaignLifecycleTotals(items: CampaignLifecycleTotals[]): Campa }), { pausable: 0, paused: 0, cancellable: 0 }); } +export function operatorDeliveryStateCounts( + queue: {queued?: number;}, + // send.queued intentionally does not contribute: paused queue rows retain that send state. + send: {queued?: number;claimed?: number;sending?: number;smtpAccepted?: number;sent?: number;} +) { + const queued = queue.queued ?? 0; + const claimed = send.claimed ?? 0; + const sending = send.sending ?? 0; + const completed = (send.smtpAccepted ?? 0) + (send.sent ?? 0); + return { queued, claimed, sending, completed, queuedOrActive: queued + claimed + sending }; +} + +export function shouldRefreshOperatorVersionSummary(input: { + fullDiscovery: boolean; + versionChanged: boolean; + selected: boolean; + cached: boolean; + queuedOrActive: number; + needsAttention: number; +}): boolean { + return input.fullDiscovery + || input.versionChanged + || !input.cached + || input.selected + || input.queuedOrActive > 0 + || input.needsAttention > 0; +} + /** * Return the reason each fixed queue action is unavailable. A null result means * the action is permitted and meaningful for the current queue projection. diff --git a/webui/tests/operator-queue-model.test.ts b/webui/tests/operator-queue-model.test.ts index 3cf85af..d3d46b2 100644 --- a/webui/tests/operator-queue-model.test.ts +++ b/webui/tests/operator-queue-model.test.ts @@ -5,7 +5,9 @@ import { OPERATOR_QUEUE_REASON, campaignLifecycleTotals, campaignVersionWorkItems, + operatorDeliveryStateCounts, operatorQueueActionBlocks, + shouldRefreshOperatorVersionSummary, type OperatorQueueFacts, type OperatorQueuePermissions } from "../src/features/operator/operatorQueueModel.ts"; @@ -112,3 +114,38 @@ test("campaign-wide lifecycle totals are not multiplied across version rows", () { pausable: 0, paused: 2, cancellable: 2 } ]), { pausable: 3, paused: 3, cancellable: 6 }); }); + +test("paused jobs do not overlap with queued or active delivery counts", () => { + assert.deepEqual(operatorDeliveryStateCounts( + { queued: 0 }, + { queued: 4, claimed: 0, sending: 0, smtpAccepted: 2, sent: 1 } + ), { + queued: 0, + claimed: 0, + sending: 0, + completed: 3, + queuedOrActive: 0 + }); +}); + +test("fast polling refreshes only selected, changed, uncached, active, or attention versions", () => { + const baseline = { + fullDiscovery: false, + versionChanged: false, + selected: false, + cached: true, + queuedOrActive: 0, + needsAttention: 0 + }; + assert.equal(shouldRefreshOperatorVersionSummary(baseline), false); + for (const override of [ + { fullDiscovery: true }, + { versionChanged: true }, + { selected: true }, + { cached: false }, + { queuedOrActive: 1 }, + { needsAttention: 1 } + ]) { + assert.equal(shouldRefreshOperatorVersionSummary({ ...baseline, ...override }), true); + } +}); diff --git a/webui/tests/operator-queue-ui-structure.test.mjs b/webui/tests/operator-queue-ui-structure.test.mjs index 38310e8..72f9228 100644 --- a/webui/tests/operator-queue-ui-structure.test.mjs +++ b/webui/tests/operator-queue-ui-structure.test.mjs @@ -22,10 +22,14 @@ assert.match(source, /next\.set\("campaign", campaignId\)[\s\S]*next\.set\("vers assert.match(source, /setSearchParams\(next, \{ replace: true \}\)/); assert.match(source, /window\.setInterval[\s\S]*10_000/); assert.match(source, /window\.setInterval[\s\S]*60_000/); +assert.match(source, /pendingFullDiscoveryRef\.current = true/); +assert.match(source, /if \(pendingFullDiscoveryRef\.current\)[\s\S]*fullDiscoveryLoadRef\.current\(\)/); assert.match(source, /