fix(campaign): bound operator queue polling
This commit is contained in:
@@ -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<Record<string, Record<string, CampaignSummary | null>>>({});
|
||||
const versionsRef = useRef<Record<string, CampaignVersionListItem[]>>({});
|
||||
const loadingRef = useRef(false);
|
||||
const pendingFullDiscoveryRef = useRef(false);
|
||||
const backgroundLoadRef = useRef<() => Promise<void>>(async () => undefined);
|
||||
const fullDiscoveryLoadRef = useRef<() => Promise<void>>(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<OperatorRow[]> {
|
||||
async function loadCampaignVersionRows(campaigns: CampaignListItem[], fullDiscovery: boolean): Promise<OperatorRow[]> {
|
||||
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<string>();
|
||||
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<string, unknown>): 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 {
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
});
|
||||
|
||||
@@ -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, /<ConfirmDialog[\s\S]*cancelTarget/);
|
||||
assert.match(source, /includeVersions: true/);
|
||||
assert.match(source, /getCampaignSummary\(settings, campaign\.id, version\.id\)/);
|
||||
assert.match(source, /campaignVersionWorkItems/);
|
||||
assert.match(source, /operatorDeliveryStateCounts/);
|
||||
assert.doesNotMatch(source, /queuedOrActive: numberValue\(cards\.queued_or_active\)/);
|
||||
assert.match(source, /getCampaignJobs\(settings, selectedRow\.campaign\.id,[\s\S]*versionId: selectedVersionId/);
|
||||
assert.doesNotMatch(source, /getCampaignJobDetail|getCampaignJobsDelta|diagnostics/);
|
||||
assert.match(source, /mode: "server"/);
|
||||
|
||||
Reference in New Issue
Block a user