Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import { resolveUserContextFromHeaders } from "@/middleware/ensure-user/resolve"
import { ProjectRepository } from "@/server/features/projects/repositories/ProjectRepository";
import { SamSessionRepository } from "@/server/features/sam/SamSessionRepository";
import { runScheduledRankChecks } from "@/server/features/rank-tracking/services/scheduledRankChecks";
import { reconcileStuckRankCheckRuns } from "@/server/features/rank-tracking/services/rankCheckReconciler";
import { reconcileStaleAudits } from "@/server/features/audit/services/auditReconciler";
import { getOrCreateOrganizationCustomer } from "@/server/billing/subscription";
import { isHostedServerAuthMode } from "@/server/lib/runtime-env";
Expand Down Expand Up @@ -193,6 +194,7 @@ export default {
let watchdogError: unknown;
try {
await withPgClient(() => reconcileStaleAudits());
await withPgClient(() => reconcileStuckRankCheckRuns());
} catch (err) {
watchdogError = err;
console.error("[cron] Stale-audit reconcile failed:", err);
Expand Down
107 changes: 107 additions & 0 deletions src/server/features/rank-tracking/services/rankCheckFinalize.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { completeRankCheckRunFromSnapshots } from "./rankCheckFinalize";

const mocks = vi.hoisted(() => ({
getSnapshotsForRun: vi.fn(),
updateRun: vi.fn(),
updateConfig: vi.fn(),
}));

vi.mock(
"@/server/features/rank-tracking/repositories/RankTrackingRepository",
() => ({ RankTrackingRepository: mocks }),
);

const baseRun = {
id: "run_1",
configId: "config_1",
projectId: "project_1",
status: "running" as const,
keywordsTotal: 2,
keywordsChecked: 2,
isSubsetRun: false,
errorMessage: null,
startedAt: "2026-09-12T06:00:00.000Z",
completedAt: null,
};

describe("completeRankCheckRunFromSnapshots", () => {
beforeEach(() => {
mocks.getSnapshotsForRun.mockReset();
mocks.updateRun.mockReset();
mocks.updateConfig.mockReset();
mocks.updateRun.mockResolvedValue(undefined);
mocks.updateConfig.mockResolvedValue(undefined);
});

it("completes a running run when every keyword has a snapshot", async () => {
mocks.getSnapshotsForRun.mockResolvedValue([
{ trackingKeywordId: "kw_1" },
{ trackingKeywordId: "kw_1" },
{ trackingKeywordId: "kw_2" },
]);

const result = await completeRankCheckRunFromSnapshots({ run: baseRun });

expect(result).toMatchObject({
keywordsChecked: 2,
keywordsTotal: 2,
});
expect(mocks.updateRun).toHaveBeenCalledWith(
"run_1",
expect.objectContaining({
status: "completed",
keywordsChecked: 2,
completedAt: expect.any(String),
}),
);
expect(mocks.updateConfig).toHaveBeenCalledWith(
"config_1",
"project_1",
expect.objectContaining({
lastCheckedAt: expect.any(String),
lastSkipReason: null,
}),
);
});

it("skips when snapshot coverage is still incomplete and requireFullCoverage is set", async () => {
mocks.getSnapshotsForRun.mockResolvedValue([
{ trackingKeywordId: "kw_1" },
]);

await expect(
completeRankCheckRunFromSnapshots({
run: baseRun,
requireFullCoverage: true,
}),
).resolves.toBeNull();
expect(mocks.updateRun).not.toHaveBeenCalled();
});

it("completes partial coverage when requireFullCoverage is off", async () => {
mocks.getSnapshotsForRun.mockResolvedValue([
{ trackingKeywordId: "kw_1" },
]);

const result = await completeRankCheckRunFromSnapshots({ run: baseRun });
expect(result).toMatchObject({ keywordsChecked: 1, keywordsTotal: 2 });
expect(mocks.updateRun).toHaveBeenCalledWith(
"run_1",
expect.objectContaining({
status: "completed",
keywordsChecked: 1,
errorMessage: "1 keyword(s) could not be checked",
}),
);
});

it("skips terminal runs", async () => {
await expect(
completeRankCheckRunFromSnapshots({
run: { ...baseRun, status: "completed" },
}),
).resolves.toBeNull();
expect(mocks.getSnapshotsForRun).not.toHaveBeenCalled();
});
});
73 changes: 73 additions & 0 deletions src/server/features/rank-tracking/services/rankCheckFinalize.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
import { RankTrackingRepository } from "@/server/features/rank-tracking/repositories/RankTrackingRepository";

type RunRow = NonNullable<
Awaited<ReturnType<typeof RankTrackingRepository.getRunById>>
>;

/**
* Flip an in-flight rank check to completed from already-written snapshots.
*
* Scheduled DFS checks write snapshots incrementally, then a separate finalize
* step flips status. If that step dies (telemetry hang, isolate kill, workflow
* retention) the API hides positions because results only read completed runs.
* Callers use this to finish the status flip without re-probing DataForSEO.
*
* When `requireFullCoverage` is true (watchdog / stale reclaim), returns null
* unless every expected keyword already has a snapshot — so a still-collecting
* run is left alone. The workflow finalize path passes false and always closes
* the run (matching prior finalize behavior, including partial/error cases).
*/
export async function completeRankCheckRunFromSnapshots(input: {
run: RunRow;
batchError?: string | null;
requireFullCoverage?: boolean;
}): Promise<{
keywordsChecked: number;
keywordsTotal: number;
completedAt: string;
} | null> {
const { run } = input;
if (run.status === "completed" || run.status === "failed") {
return null;
}

const snapshots = await RankTrackingRepository.getSnapshotsForRun(run.id);
const keywordsChecked = new Set(snapshots.map((s) => s.trackingKeywordId))
.size;
const keywordsTotal = run.keywordsTotal || keywordsChecked;

if (input.requireFullCoverage) {
if (keywordsChecked === 0 || keywordsChecked < keywordsTotal) {
return null;
}
}

const completedAt = new Date().toISOString();
const incompleteCount = Math.max(keywordsTotal - keywordsChecked, 0);

let errorMessage: string | undefined;
if (input.batchError) {
errorMessage = `Completed ${keywordsChecked} of ${keywordsTotal} keyword(s). Error: ${input.batchError}`;
} else if (incompleteCount > 0) {
errorMessage = `${incompleteCount} keyword(s) could not be checked`;
}

// Flipping status away from 'pending'/'running' releases the partial-index
// slot for the next run. Do this before any telemetry.
await RankTrackingRepository.updateRun(run.id, {
status: "completed",
keywordsChecked,
completedAt,
...(errorMessage ? { errorMessage } : {}),
});

// Clear any previous skip reason on success.
// Note: nextCheckAt is NOT set here — the cron handler advances it eagerly
// before starting the workflow to prevent retry storms.
await RankTrackingRepository.updateConfig(run.configId, run.projectId, {
lastCheckedAt: completedAt,
lastSkipReason: null,
});

return { keywordsChecked, keywordsTotal, completedAt };
}
73 changes: 73 additions & 0 deletions src/server/features/rank-tracking/services/rankCheckReconciler.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
import { and, asc, inArray, lt } from "drizzle-orm";
import { db } from "@/db";
import { rankCheckRuns } from "@/db/schema";
import { completeRankCheckRunFromSnapshots } from "@/server/features/rank-tracking/services/rankCheckFinalize";
import { failRunIfActive } from "@/server/features/rank-tracking/services/rankCheckRunGuards";

/**
* Snapshots for a finished DFS collect are usually present minutes before the
* finalize step runs. Give the live workflow that window, then complete from
* DB state so positions stop staying invisible behind status=running.
*/
const SNAPSHOT_FINALIZE_GRACE_MS = 3 * 60 * 1000;

/** In-flight runs older than this with no complete snapshots get failed. */
const STALE_INCOMPLETE_MS = 25 * 60 * 1000;

const WATCHDOG_BATCH_LIMIT = 100;

/**
* Cron watchdog: finish (or fail) rank_check_runs stuck in pending/running
* after their workflow should have finalized.
*/
export async function reconcileStuckRankCheckRuns() {
const snapshotGraceCutoff = new Date(
Date.now() - SNAPSHOT_FINALIZE_GRACE_MS,
).toISOString();
const incompleteCutoff = new Date(
Date.now() - STALE_INCOMPLETE_MS,
).toISOString();

const stuck = await db
.select()
.from(rankCheckRuns)
.where(
and(
inArray(rankCheckRuns.status, ["pending", "running"]),
lt(rankCheckRuns.startedAt, snapshotGraceCutoff),
),
)
.orderBy(asc(rankCheckRuns.startedAt))
.limit(WATCHDOG_BATCH_LIMIT);

for (const run of stuck) {
try {
const completed = await completeRankCheckRunFromSnapshots({
run,
requireFullCoverage: true,
});
if (completed) {
console.log(
`[rank-check] watchdog completed run ${run.id} from snapshots (${completed.keywordsChecked}/${completed.keywordsTotal})`,
);
continue;
}

if (run.startedAt < incompleteCutoff) {
await failRunIfActive(
run.id,
"Rank check timed out before finalizing",
run,
);
console.log(
`[rank-check] watchdog failed stale incomplete run ${run.id}`,
);
}
} catch (error) {
console.error(
`[rank-check] watchdog failed to reconcile ${run.id}:`,
error,
);
}
}
}
11 changes: 10 additions & 1 deletion src/server/features/rank-tracking/services/rankCheckRunGuards.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { env } from "cloudflare:workers";
import type { BillingCustomerContext } from "@/server/billing/subscription";
import { RankTrackingRepository } from "@/server/features/rank-tracking/repositories/RankTrackingRepository";
import { completeRankCheckRunFromSnapshots } from "@/server/features/rank-tracking/services/rankCheckFinalize";
import type {
RankCheckTriggerResult,
RankTrackingConfig,
Expand Down Expand Up @@ -212,7 +213,15 @@ export async function beginRankCheckRun(input: {
ageMs: Date.now() - new Date(blocker.startedAt).getTime(),
});
if (staleReason) {
await failRunIfActive(blocker.id, staleReason, blocker);
// Prefer completing from snapshots over failing — failing still hides
// paid DFS results (results SQL filters to status=completed).
const completed = await completeRankCheckRunFromSnapshots({
run: blocker,
requireFullCoverage: true,
});
if (!completed) {
await failRunIfActive(blocker.id, staleReason, blocker);
}
continue; // slot is free now — retry insert
}
}
Expand Down
12 changes: 10 additions & 2 deletions src/server/lib/posthog.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,12 @@ export async function captureServerError(
} catch (posthogError) {
console.error("posthog server capture failed", posthogError);
} finally {
await client.shutdown().catch(() => {});
// Bound shutdown: an indefinite PostHog flush has wedged Worker/workflow
// steps after the real work finished (rank-check finalize hung on this).
await Promise.race([
client.shutdown().catch(() => {}),
new Promise<void>((resolve) => setTimeout(resolve, 1500)),
]);
}
}

Expand Down Expand Up @@ -71,6 +76,9 @@ export async function captureServerEvent(args: {
} catch (posthogError) {
console.error("posthog server capture failed", posthogError);
} finally {
await client.shutdown().catch(() => {});
await Promise.race([
client.shutdown().catch(() => {}),
new Promise<void>((resolve) => setTimeout(resolve, 1500)),
]);
}
}
Loading