SPB Git forge
38commits 1branches 0releases
338.7 MBsize
maindefault branch
2 h agolast push
HTML 53.9% TypeScript 44.5% JavaScript 0.6% SQL 0.5%

Admin console (connector control center, dev tool, matches, documents, ops), partial worker metrics + pluggable budget store

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Simon-Pierre Boucher committed 21 days ago (Sep 11, 2026) parent 650b3c9

15 changed files +449 −29

added apps/web/qa/admin/connector-detail-1280.png +0 −0

Binary file not shown.

added apps/web/qa/admin/connector-detail-config-1280.png +0 −0

Binary file not shown.

added apps/web/qa/admin/connectors-1280.png +0 −0

Binary file not shown.

added apps/web/qa/admin/connectors-run-form-1280.png +0 −0

Binary file not shown.

added apps/web/qa/admin/documents-1280.png +0 −0

Binary file not shown.

added apps/web/qa/admin/login-error-1280.png +0 −0

Binary file not shown.

added apps/web/qa/admin/overview-1280.png +0 −0

Binary file not shown.

added apps/web/qa/admin/shoot-admin.mjs +252 −0
@@ -0,0 +1,252 @@
1 +// Admin console QA: node apps/web/qa/admin/shoot-admin.mjs (Playwright from ~/Desktop/uqo-eval)
2 +// env: BASE (http://localhost:8310), TOKEN (dev-admin-token), DOC_DIFF (document id with ≥ 2 versions), SKIP_DEVTOOL=1
3 +import { createRequire } from "node:module";
4 +import { mkdirSync } from "node:fs";
5 +import { dirname, resolve } from "node:path";
6 +import { fileURLToPath } from "node:url";
7 +
8 +const require = createRequire("/Users/simon-pierreboucher/Desktop/uqo-eval/package.json");
9 +const { chromium } = require("playwright");
10 +const out = dirname(fileURLToPath(import.meta.url));
11 +mkdirSync(out, { recursive: true });
12 +const BASE = process.env.BASE ?? "http://localhost:8310";
13 +const TOKEN = process.env.TOKEN ?? "dev-admin-token";
14 +const DOC_DIFF = process.env.DOC_DIFF ?? "";
15 +const SKIP_DEVTOOL = process.env.SKIP_DEVTOOL === "1";
16 +const report = [];
17 +const log = (...a) => console.log(new Date().toISOString().slice(11, 19), ...a);
18 +
19 +async function login(page) {
20 + await page.goto(`${BASE}/admin/login`, { waitUntil: "load" });
21 + await page.fill('input[name="token"]', TOKEN);
22 + await page.click('button[type="submit"]');
23 + await page.waitForURL((u) => u.pathname === "/admin", { timeout: 30000 });
24 +}
25 +
26 +async function shot(page, name, w, opts = {}) {
27 + await page.waitForTimeout(opts.wait ?? 800);
28 + const overflow = await page.evaluate(() => {
29 + const doc = document.documentElement;
30 + return { scrollWidth: doc.scrollWidth, clientWidth: doc.clientWidth };
31 + });
32 + report.push({ w, name, overflow: overflow.scrollWidth > overflow.clientWidth ? `OVERFLOW ${overflow.scrollWidth}>${overflow.clientWidth}` : "ok" });
33 + await page.screenshot({ path: resolve(out, `${name}-${w}.png`), fullPage: opts.fullPage ?? true });
34 +}
35 +
36 +async function api(page, path) {
37 + return page.evaluate(async (p) => {
38 + const r = await fetch(`/admin/api/${p}`, { credentials: "same-origin" });
39 + return { status: r.status, json: await r.json().catch(() => null) };
40 + }, path);
41 +}
42 +
43 +const browser = await chromium.launch();
44 +
45 +// ------------------------------------------------------------------ desktop
46 +{
47 + const w = 1280;
48 + const ctx = await browser.newContext({ viewport: { width: w, height: 900 }, deviceScaleFactor: 1 });
49 + const page = await ctx.newPage();
50 + page.on("dialog", (d) => d.accept());
51 + page.on("pageerror", (e) => log("PAGEERROR", String(e).slice(0, 200)));
52 +
53 + // wrong token first
54 + await page.goto(`${BASE}/admin/login`);
55 + await page.fill('input[name="token"]', "nope");
56 + await page.click('button[type="submit"]');
57 + await page.waitForSelector('[role="alert"]', { timeout: 20000 });
58 + log("login error shown:", (await page.textContent('[role="alert"]')).trim());
59 + await shot(page, "login-error", w, { fullPage: false });
60 +
61 + await login(page);
62 + log("logged in");
63 + await page.waitForSelector("text=Last 20 runs", { timeout: 30000 });
64 + await page.waitForTimeout(2500);
65 + await shot(page, "overview", w);
66 +
67 + // connectors: run azure-regions
68 + await page.goto(`${BASE}/admin/connectors?q=azure`);
69 + await page.waitForSelector("text=azure-regions", { timeout: 30000 });
70 + const row = page.locator("tr", { hasText: "azure-regions" }).first();
71 + await row.getByRole("button", { name: "Run", exact: true }).click();
72 + await page.waitForSelector("text=Task");
73 + await page.selectOption("select", { label: "discover — sitemap / RSS / seeds" }).catch(() => {});
74 + await shot(page, "connectors-run-form", w);
75 + await page.locator("tr.bg-surface-2\\/60 + tr button", { hasText: "Run" }).last().click();
76 + await page.waitForSelector("text=Run enqueued", { timeout: 20000 });
77 + log("run enqueued toast visible");
78 + await page.goto(`${BASE}/admin/connectors`);
79 + await page.waitForSelector("text=Connector", { timeout: 30000 });
80 + await page.waitForTimeout(2000);
81 + await shot(page, "connectors", w);
82 +
83 + await page.goto(`${BASE}/admin/connectors/azure-regions`);
84 + await page.waitForSelector("text=Health history", { timeout: 30000 });
85 + await page.waitForTimeout(1500);
86 + await shot(page, "connector-detail", w);
87 + await page.getByRole("tab", { name: "Config" }).click();
88 + await page.waitForTimeout(500);
89 + await shot(page, "connector-detail-config", w, { fullPage: false });
90 +
91 + // documents
92 + await page.goto(`${BASE}/admin/documents`);
93 + await page.waitForSelector("table tbody tr", { timeout: 30000 });
94 + await shot(page, "documents", w);
95 + const firstDoc = await page.locator("table tbody tr td a").first().getAttribute("href");
96 + await page.goto(`${BASE}${firstDoc}`);
97 + await page.waitForSelector("text=Metadata", { timeout: 30000 });
98 + await page.locator("iframe[title^='Rendered']").or(page.getByText("Raw body unavailable")).or(page.getByText("No stored versions")).first().waitFor({ timeout: 30000 }).catch(() => {});
99 + await page.waitForTimeout(2500);
100 + await shot(page, "document-inspect", w);
101 + if (DOC_DIFF) {
102 + await page.goto(`${BASE}/admin/documents/${DOC_DIFF}`);
103 + await page.waitForSelector("text=Diff", { timeout: 30000 });
104 + await page.getByText("Stored diff summary").or(page.getByText("No line differences")).first().waitFor({ timeout: 60000 });
105 + await page.waitForTimeout(1500);
106 + await shot(page, "document-diff", w);
107 + log("diff rendered for", DOC_DIFF);
108 + }
109 +
110 + // matches: reject the lowest-scoring pending
111 + await page.goto(`${BASE}/admin/matches`);
112 + await page.waitForSelector("text=Entity matches", { timeout: 30000 });
113 + await page.waitForTimeout(2500);
114 + await shot(page, "matches", w);
115 + const pending = await api(page, "matches?status=pending&per_page=200");
116 + const list = pending.json?.data ?? [];
117 + if (list.length) {
118 + const target = list[list.length - 1];
119 + log(`rejecting lowest-score pending match ${target.id} (${target.candidate?.name} → ${target.matched?.name}, score ${target.score})`);
120 + const card = page.locator("article", { hasText: target.candidateKey }).first();
121 + if (await card.count()) {
122 + await card.getByRole("button", { name: "Reject" }).click();
123 + await page.waitForSelector("text=Match rejected", { timeout: 20000 });
124 + await shot(page, "matches-rejected", w, { fullPage: false });
125 + const after = await api(page, `matches?status=rejected&per_page=200`);
126 + log("rejected now visible in rejected tab:", (after.json?.data ?? []).some((m) => m.id === target.id));
127 + } else log("target card not on first page; skipped UI reject");
128 + } else log("no pending matches to decide");
129 +
130 + // events
131 + await page.goto(`${BASE}/admin/events`);
132 + await page.waitForSelector("text=Events review", { timeout: 30000 });
133 + await page.waitForTimeout(2500);
134 + await shot(page, "events", w);
135 +
136 + // facility edit (no save)
137 + const fac = (pending.json?.data ?? []).find((m) => m.matched)?.matched?.id ?? "fac_2ezznz6uhn5b";
138 + await page.goto(`${BASE}/admin/facilities/${fac}/edit`);
139 + await page.waitForSelector("text=Curation note", { timeout: 30000 });
140 + await page.waitForTimeout(1000);
141 + await shot(page, "facility-edit", w);
142 +
143 + // sources
144 + await page.goto(`${BASE}/admin/sources`);
145 + await page.waitForSelector("table tbody tr", { timeout: 30000 });
146 + await shot(page, "sources", w);
147 +
148 + // dev tool
149 + if (!SKIP_DEVTOOL) {
150 + await page.goto(`${BASE}/admin/devtool`);
151 + await page.waitForSelector("text=Connector development tool", { timeout: 30000 });
152 + await page.evaluate(() => window.localStorage.removeItem("dci:devtool:draft:v1"));
153 + await page.reload();
154 + await page.waitForSelector("text=Source URL and fetch options", { timeout: 30000 });
155 + await page.fill('input[placeholder^="https://www.example.com"]', "https://www.digitalrealty.com/data-centers");
156 + await page.getByRole("button", { name: "Fetch sample" }).click();
157 + await page.waitForSelector("text=HTML source", { timeout: 120000 });
158 + await page.waitForTimeout(1500);
159 + log("devtool fetch done");
160 + await shot(page, "devtool-1-fetch", w);
161 + await page.getByRole("tab", { name: "HTML source" }).click();
162 + await page.fill('input[placeholder^="Find"]', "h1");
163 + await page.waitForTimeout(600);
164 + await shot(page, "devtool-1-source-search", w, { fullPage: false });
165 +
166 + await page.getByRole("button", { name: /Extraction rules →/ }).click();
167 + await page.waitForSelector("text=YAML preview", { timeout: 20000 });
168 + await shot(page, "devtool-2-rules", w);
169 +
170 + await page.getByRole("button", { name: /Preview fields →/ }).click();
171 + await page.waitForSelector("text=Preview extraction", { timeout: 20000 });
172 + await page.getByRole("button", { name: "Preview fields", exact: true }).click();
173 + await page.waitForSelector("text=Records (raw extraction)", { timeout: 120000 });
174 + await page.waitForTimeout(1000);
175 + log("devtool preview done");
176 + await shot(page, "devtool-3-preview", w);
177 +
178 + await page.getByRole("button", { name: /Save connector →/ }).click();
179 + await page.waitForSelector("text=Connector YAML", { timeout: 20000 });
180 + await page.fill('input[placeholder="digitalrealty"]', "devtool-test");
181 + await page.getByRole("button", { name: "Regenerate from steps" }).click();
182 + await page.waitForTimeout(400);
183 + await page.getByRole("button", { name: "Validate" }).click();
184 + await page.waitForSelector("text=valid", { timeout: 30000 });
185 + await shot(page, "devtool-4-validated", w);
186 + await page.getByRole("button", { name: "Save connector", exact: true }).click();
187 + await page.waitForSelector("text=Saved config", { timeout: 30000 }).catch(() => {});
188 + await page.waitForSelector("text=/Saved .*devtool-test\\.yaml/", { timeout: 30000 });
189 + log("devtool save-config done");
190 + await shot(page, "devtool-4-saved", w);
191 +
192 + await page.getByRole("button", { name: /Run on pages →/ }).click();
193 + await page.waitForSelector("text=Run on multiple pages", { timeout: 20000 });
194 + await page.fill("textarea", "https://www.digitalrealty.com/data-centers\nhttps://www.digitalrealty.com/data-centers/americas");
195 + await page.getByRole("button", { name: "Run sample" }).click();
196 + await page.waitForSelector("text=Comparison — first entity per URL", { timeout: 280000 });
197 + await page.waitForTimeout(1000);
198 + log("devtool run-sample done");
199 + await shot(page, "devtool-5-sample", w);
200 +
201 + await page.getByRole("button", { name: /Compare & deploy →/ }).click();
202 + await page.waitForSelector("text=Compare with the saved connector", { timeout: 20000 });
203 + await page.waitForSelector("text=Saved parser version", { timeout: 30000 });
204 + await page.waitForTimeout(800);
205 + await shot(page, "devtool-6-compare", w);
206 + await page.getByLabel(/enqueue reprocess/).uncheck();
207 + await page.getByRole("button", { name: "Deploy", exact: true }).click();
208 + await page.waitForSelector("text=save-config", { timeout: 60000 });
209 + await page.waitForTimeout(1500);
210 + log("devtool deploy done");
211 + await shot(page, "devtool-6-deployed", w);
212 + }
213 + await ctx.close();
214 +}
215 +
216 +// ------------------------------------------------------------------ phone
217 +{
218 + const w = 390;
219 + const ctx = await browser.newContext({ viewport: { width: w, height: 844 }, deviceScaleFactor: 1, isMobile: true, hasTouch: true });
220 + const page = await ctx.newPage();
221 + page.on("pageerror", (e) => log("PAGEERROR(m)", String(e).slice(0, 200)));
222 + await page.goto(`${BASE}/admin/login`);
223 + await shot(page, "login", w, { fullPage: false });
224 + await login(page);
225 + await page.waitForSelector("text=Last 20 runs", { timeout: 30000 });
226 + await page.waitForTimeout(2500);
227 + await shot(page, "overview", w);
228 + for (const [name, path, waitFor] of [
229 + ["connectors", "/admin/connectors", "text=Connector"],
230 + ["connector-detail", "/admin/connectors/azure-regions", "text=Health history"],
231 + ["documents", "/admin/documents", "table tbody tr"],
232 + ["matches", "/admin/matches", "text=Entity matches"],
233 + ["events", "/admin/events", "text=Events review"],
234 + ["sources", "/admin/sources", "table tbody tr"],
235 + ["devtool", "/admin/devtool", "text=Source URL"],
236 + ]) {
237 + await page.goto(`${BASE}${path}`);
238 + await page.waitForSelector(waitFor, { timeout: 30000 }).catch(() => {});
239 + await page.waitForTimeout(2000);
240 + await shot(page, name, w);
241 + }
242 + if (DOC_DIFF) {
243 + await page.goto(`${BASE}/admin/documents/${DOC_DIFF}`);
244 + await page.waitForSelector("text=Metadata", { timeout: 30000 });
245 + await page.waitForTimeout(3000);
246 + await shot(page, "document-inspect", w);
247 + }
248 + await ctx.close();
249 +}
250 +
251 +await browser.close();
252 +for (const r of report) console.log(r.w, r.name, r.overflow);
modified apps/web/src/components/admin/devtool/step-deploy.tsx +2 −1
@@ -39,7 +39,8 @@ export function StepDeploy({ s, configs, onReloadConfigs }: { s: DevtoolShared;
39 39 const [id, setId] = useState(draft.meta.id);
40 40 const [reloadToken, setReloadToken] = useState(0);
41 41 const savedQ = useSavedYaml(id, reloadToken);
42 − const saved = savedQ.text != null ? { id, text: savedQ.text } : null;
42 + const savedText = savedQ.text;
43 + const saved = useMemo(() => (savedText != null ? { id, text: savedText } : null), [id, savedText]);
43 44 const [versionInput, setVersionInput] = useState<string | null>(null);
44 45 const [reprocess, setReprocess] = useState(true);
45 46 const [busy, setBusy] = useState(false);
modified apps/web/src/components/admin/devtool/step-save.tsx +4 −2
@@ -99,8 +99,10 @@ export function StepSave({ s }: { s: DevtoolShared }) {
99 99 toast({ tone: "ok", title: `Saved ${id}.yaml`, description: `parser ${r.data.parserVersion}${r.data.backup ? " · backup kept" : ""}${r.data.enabled ? "" : " · enabled: false"}` });
100 100 };
101 101
102 − const errors = validation && !validation.ok ? validation.errors : [];
103 − const errorLines = useMemo(() => new Set(errors.map((e) => lineForPath(draft.yaml, e.path)).filter((n): n is number => n != null)), [errors, draft.yaml]);
102 + const errorLines = useMemo(() => {
103 + const errs = validation && !validation.ok ? validation.errors : [];
104 + return new Set(errs.map((e) => lineForPath(draft.yaml, e.path)).filter((n): n is number => n != null));
105 + }, [validation, draft.yaml]);
104 106 const jumpTo = (line: number | null) => {
105 107 if (!line || !ta.current) return;
106 108 setView("edit");
modified apps/web/src/components/admin/use-admin-query.ts +19 −17
@@ -53,32 +53,36 @@ export function useAdminQuery<T>(path: string | null, opts: Options = {}): Admin
53 53 keyRef.current = key;
54 54 }, [key]);
55 55
56 − const load = useCallback(
57 − async (background: boolean) => {
58 − if (!path || !enabled) return;
59 − const myKey = `${path}|${qs}`;
60 − abortRef.current?.abort();
61 − const ctrl = new AbortController();
62 − abortRef.current = ctrl;
63 − if (background) setSlot((s) => ({ ...s, refreshing: true }));
64 − const res: AdminResult<T> = await adminFetch<T>(path, { query: (JSON.parse(qs) as AdminQuery | null) ?? undefined, signal: ctrl.signal, timeoutMs });
56 + // `load` never writes state before its first await (it is invoked from an effect); the `refreshing` flag is set
57 + // by `refresh()` / the poller only, which run from timers and event handlers.
58 + const load = useCallback((): Promise<void> => {
59 + if (!path || !enabled) return Promise.resolve();
60 + const myKey = `${path}|${qs}`;
61 + abortRef.current?.abort();
62 + const ctrl = new AbortController();
63 + abortRef.current = ctrl;
64 + return adminFetch<T>(path, { query: (JSON.parse(qs) as AdminQuery | null) ?? undefined, signal: ctrl.signal, timeoutMs }).then((res: AdminResult<T>) => {
65 65 if (ctrl.signal.aborted || keyRef.current !== myKey) return;
66 66 if (res.ok) setSlot({ key: myKey, data: res.data, meta: res.meta, error: null, status: res.status, updatedAt: Date.now(), refreshing: false });
67 67 else setSlot((s) => ({ key: myKey, data: s.key === myKey ? s.data : null, meta: s.key === myKey ? s.meta : undefined, error: res.error, status: res.status, updatedAt: s.key === myKey ? s.updatedAt : null, refreshing: false }));
68 − },
69 − [path, enabled, qs, timeoutMs],
70 − );
68 + });
69 + }, [path, enabled, qs, timeoutMs]);
71 70
72 71 useEffect(() => {
73 72 if (!active) return;
74 − void load(false);
73 + void load();
75 74 return () => abortRef.current?.abort();
76 75 }, [load, active]);
77 76
77 + const refresh = useCallback(() => {
78 + setSlot((s) => ({ ...s, refreshing: true }));
79 + return load();
80 + }, [load]);
81 +
78 82 useEffect(() => {
79 83 if (!refreshMs || !active) return;
80 84 const tick = () => {
81 − if (document.visibilityState === "visible") void load(true);
85 + if (document.visibilityState === "visible") void refresh();
82 86 };
83 87 const t = setInterval(tick, refreshMs);
84 88 document.addEventListener("visibilitychange", tick);
@@ -86,9 +90,7 @@ export function useAdminQuery<T>(path: string | null, opts: Options = {}): Admin
86 90 clearInterval(t);
87 91 document.removeEventListener("visibilitychange", tick);
88 92 };
89 − }, [refreshMs, load, active]);
90 −
91 − const refresh = useCallback(() => load(true), [load]);
93 + }, [refreshMs, refresh, active]);
92 94 const mutate = useCallback((fn: (prev: T | null) => T | null) => setSlot((s) => ({ ...s, data: fn(s.data) })), []);
93 95
94 96 const current = slot.key === key;
modified apps/web/src/lib/admin-api.ts +2 −1
@@ -90,7 +90,8 @@ export async function adminFetch<T>(path: string, init: AdminFetchInit = {}): Pr
90 90 if (res.status === 401 && init.redirectOn401 !== false && typeof window !== "undefined" && !redirecting) {
91 91 redirecting = true;
92 92 const next = encodeURIComponent(window.location.pathname + window.location.search);
93 − window.location.assign(`/admin/login?next=${next}&reason=expired`);
93 + // full navigation on purpose: the login page is server-rendered and client state must be dropped
94 + window.location.assign(new URL(`/admin/login?next=${next}&reason=expired`, window.location.origin).href);
94 95 }
95 96 const text = await res.text();
96 97 let json: unknown = null;
added apps/worker/src/prom.ts +141 −0
@@ -0,0 +1,141 @@
1 +/**
2 + * Tiny in-process Prometheus registry (text exposition format 0.0.4) — no prom-client dependency.
3 + * Counters, gauges and a fixed-bucket histogram with string labels. Shared by the worker (`main.ts`) and the
4 + * standalone scheduler (`scheduler.ts`); `renderMetrics()` produces the `/metrics` body.
5 + *
6 + * Metric contract (docs/DEPLOY.md, Grafana "DataCenterIndex — overview", deploy/monitoring/alerts.yml):
7 + * dci_crawl_fetches_total{connector,level,outcome=ok|not_modified|error|blocked}
8 + * dci_crawl_credits_total{provider=scrapfly|firecrawl}
9 + * dci_crawl_daily_budget{provider} dci_crawl_daily_credits_used{provider}
10 + * dci_crawl_fetch_duration_seconds (histogram)
11 + * dci_ingest_entities_total{connector,result=created|updated|unchanged|rejected}
12 + * dci_events_total
13 + * dci_queue_jobs{queue,state} dci_worker_running_jobs dci_worker_up
14 + */
15 +
16 +type Labels = Record<string, string | number>;
17 +
18 +function esc(v: string | number): string { return String(v).replace(/\\/g, "\\\\").replace(/"/g, '\\"').replace(/\n/g, "\\n"); }
19 +function fmt(v: number): string { return Number.isFinite(v) ? (Number.isInteger(v) ? String(v) : v.toString()) : "0"; }
20 +
21 +function key(names: readonly string[], labels: Labels | undefined): string {
22 + if (!names.length) return "";
23 + return names.map((n) => `${n}="${esc(labels?.[n] ?? "")}"`).join(",");
24 +}
25 +
26 +interface Metric { render(): string[]; reset(): void }
27 +const registry = new Map<string, Metric>();
28 +
29 +function register<T extends Metric>(name: string, m: T): T {
30 + if (registry.has(name)) throw new Error(`metric ${name} already registered`);
31 + registry.set(name, m);
32 + return m;
33 +}
34 +
35 +export class Counter implements Metric {
36 + private readonly values = new Map<string, number>();
37 + constructor(readonly name: string, readonly help: string, readonly labelNames: readonly string[] = []) {}
38 + inc(labels?: Labels, n = 1): void {
39 + if (!Number.isFinite(n) || n === 0) return;
40 + const k = key(this.labelNames, labels);
41 + this.values.set(k, (this.values.get(k) ?? 0) + n);
42 + }
43 + get(labels?: Labels): number { return this.values.get(key(this.labelNames, labels)) ?? 0; }
44 + reset(): void { this.values.clear(); }
45 + render(): string[] {
46 + const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} counter`];
47 + // an unlabelled, never-incremented counter renders as 0; a labelled one renders no sample (Prometheus convention)
48 + if (!this.values.size && !this.labelNames.length) out.push(`${this.name} 0`);
49 + for (const [k, v] of this.values) out.push(`${this.name}${k ? `{${k}}` : ""} ${fmt(v)}`);
50 + return out;
51 + }
52 +}
53 +
54 +export class Gauge implements Metric {
55 + private readonly values = new Map<string, number>();
56 + constructor(readonly name: string, readonly help: string, readonly labelNames: readonly string[] = []) {}
57 + set(value: number, labels?: Labels): void { this.values.set(key(this.labelNames, labels), Number.isFinite(value) ? value : 0); }
58 + inc(labels?: Labels, n = 1): void { const k = key(this.labelNames, labels); this.values.set(k, (this.values.get(k) ?? 0) + n); }
59 + dec(labels?: Labels, n = 1): void { this.inc(labels, -n); }
60 + get(labels?: Labels): number { return this.values.get(key(this.labelNames, labels)) ?? 0; }
61 + /** drop every series (used before re-publishing a full snapshot such as queue depths) */
62 + reset(): void { this.values.clear(); }
63 + render(): string[] {
64 + const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} gauge`];
65 + if (!this.values.size && !this.labelNames.length) out.push(`${this.name} 0`);
66 + for (const [k, v] of this.values) out.push(`${this.name}${k ? `{${k}}` : ""} ${fmt(v)}`);
67 + return out;
68 + }
69 +}
70 +
71 +export class Histogram implements Metric {
72 + private readonly series = new Map<string, { counts: number[]; sum: number; count: number }>();
73 + constructor(readonly name: string, readonly help: string, readonly buckets: readonly number[], readonly labelNames: readonly string[] = []) {
74 + if (!buckets.length || buckets.some((b, i) => i > 0 && b <= buckets[i - 1]!)) throw new Error(`histogram ${name}: buckets must be strictly increasing`);
75 + }
76 + observe(value: number, labels?: Labels): void {
77 + if (!Number.isFinite(value)) return;
78 + const k = key(this.labelNames, labels);
79 + let s = this.series.get(k);
80 + if (!s) { s = { counts: new Array<number>(this.buckets.length).fill(0), sum: 0, count: 0 }; this.series.set(k, s); }
81 + for (let i = 0; i < this.buckets.length; i++) if (value <= this.buckets[i]!) s.counts[i]!++;
82 + s.sum += value; s.count++;
83 + }
84 + reset(): void { this.series.clear(); }
85 + render(): string[] {
86 + const out = [`# HELP ${this.name} ${this.help}`, `# TYPE ${this.name} histogram`];
87 + const emit = (k: string, s: { counts: number[]; sum: number; count: number }) => {
88 + const sep = k ? "," : "";
89 + for (let i = 0; i < this.buckets.length; i++) out.push(`${this.name}_bucket{${k}${sep}le="${fmt(this.buckets[i]!)}"} ${s.counts[i]}`);
90 + out.push(`${this.name}_bucket{${k}${sep}le="+Inf"} ${s.count}`, `${this.name}_sum${k ? `{${k}}` : ""} ${fmt(s.sum)}`, `${this.name}_count${k ? `{${k}}` : ""} ${s.count}`);
91 + };
92 + if (!this.series.size && !this.labelNames.length) emit("", { counts: new Array<number>(this.buckets.length).fill(0), sum: 0, count: 0 });
93 + for (const [k, s] of this.series) emit(k, s);
94 + return out;
95 + }
96 +}
97 +
98 +export function counter(name: string, help: string, labelNames: readonly string[] = []): Counter { return register(name, new Counter(name, help, labelNames)); }
99 +export function gauge(name: string, help: string, labelNames: readonly string[] = []): Gauge { return register(name, new Gauge(name, help, labelNames)); }
100 +export function histogram(name: string, help: string, buckets: readonly number[], labelNames: readonly string[] = []): Histogram { return register(name, new Histogram(name, help, buckets, labelNames)); }
101 +
102 +/** Full exposition body (trailing newline included). */
103 +export function renderMetrics(): string {
104 + const lines: string[] = [];
105 + for (const m of registry.values()) lines.push(...m.render());
106 + return lines.join("\n") + "\n";
107 +}
108 +
109 +/** Tests only. */
110 +export function resetMetrics(): void { for (const m of registry.values()) m.reset(); }
111 +
112 +export type FetchOutcome = "ok" | "not_modified" | "error" | "blocked";
113 +
114 +/* ---------- the worker's metrics ---------- */
115 +export const metrics = {
116 + fetches: counter("dci_crawl_fetches_total", "fetch attempts that returned a document, by connector, final level and outcome", ["connector", "level", "outcome"]),
117 + fetchDuration: histogram("dci_crawl_fetch_duration_seconds", "wall time of one fetch (all escalation levels included)", [0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 120]),
118 + credits: counter("dci_crawl_credits_total", "premium credits spent since process start", ["provider"]),
119 + dailyBudget: gauge("dci_crawl_daily_budget", "daily premium credit limit (DCI_<PROVIDER>_DAILY_BUDGET)", ["provider"]),
120 + dailyUsed: gauge("dci_crawl_daily_credits_used", "premium credits used today, shared across workers (Redis dci:budget:<provider>:<day>)", ["provider"]),
121 + ingest: counter("dci_ingest_entities_total", "normalized entities handed to ingest, by result", ["connector", "result"]),
122 + events: counter("dci_events_total", "change events emitted by ingest"),
123 + runs: counter("dci_crawl_runs_total", "connector runs finished, by status", ["status"]),
124 + jobs: counter("dci_worker_jobs_total", "BullMQ jobs processed by this worker", ["queue", "result"]),
125 + queueJobs: gauge("dci_queue_jobs", "jobs per queue and state (snapshot at scrape time)", ["queue", "state"]),
126 + runningJobs: gauge("dci_worker_running_jobs", "jobs currently processed by this worker"),
127 + up: gauge("dci_worker_up", "1 while the process accepts work (0 while shutting down)"),
128 + uptime: gauge("dci_worker_uptime_seconds", "seconds since process start"),
129 + schedulerTicks: counter("dci_scheduler_ticks_total", "scheduler passes, by result", ["result"]),
130 + schedulerEnqueued: counter("dci_scheduler_enqueued_total", "crawl jobs enqueued by the scheduler"),
131 + schedulerLastTick: gauge("dci_scheduler_last_tick_timestamp_seconds", "unix time of the last successful scheduler pass"),
132 +};
133 +
134 +/** Outcome label for a fetched RawDocument-like object. */
135 +export function fetchOutcome(doc: { notModified: boolean; status: number; error?: { code: string } | null }): FetchOutcome {
136 + if (doc.notModified) return "not_modified";
137 + if (doc.error) return doc.error.code === "robots_disallow" || doc.error.code === "ssrf_blocked" ? "blocked" : "error";
138 + if (doc.status === 401 || doc.status === 403 || doc.status === 429 || doc.status === 503) return "blocked";
139 + if (doc.status >= 400 || doc.status === 0) return "error";
140 + return "ok";
141 +}
modified packages/connectors/src/config.ts +2 −0
@@ -75,6 +75,8 @@ export const connectorConfigSchema = z.object({
75 75 respectRobots: z.boolean().default(true),
76 76 /** premium credit budget per run */
77 77 maxCreditsPerRun: z.number().int().default(200),
78 + /** premium credit budget per UTC day for this connector (all runs, all workers — Redis counter); unset = only the global daily budgets apply */
79 + maxCreditsPerDay: z.number().int().min(0).optional(),
78 80 })
79 81 .prefault({}),
80 82 discovery: z
modified packages/connectors/src/fetchers.ts +27 −8
@@ -120,19 +120,38 @@ export class DirectFetcher implements Fetcher {
120 120 }
121 121 }
122 122
123 +export type PremiumProvider = "scrapfly" | "firecrawl";
124 +
125 +/**
126 + * Pluggable shared store for the daily premium budgets. The in-process counter is always kept; a store lets the
127 + * runtime persist credits (e.g. Redis `dci:budget:<provider>:<day>`) so budgets survive restarts and are shared by
128 + * several worker processes. `used()` returns the latest known shared value (may lag) or null when unknown;
129 + * `spend()` must never throw (fire-and-forget is fine).
130 + */
131 +export interface BudgetStore {
132 + used(provider: PremiumProvider, day: string): number | null;
133 + spend(provider: PremiumProvider, day: string, credits: number): void;
134 +}
135 +let budgetStore: BudgetStore | null = null;
136 +export function setBudgetStore(store: BudgetStore | null): void { budgetStore = store; }
137 +/** UTC day key used by the budgets and the shared store. */
138 +export function budgetDay(now = new Date()): string { return now.toISOString().slice(0, 10); }
139 +
123 140 /** Daily budget bookkeeping shared by premium fetchers. */
124 141 class Budget {
125 142 private day = "";
126 − private used = 0;
127 − constructor(private readonly envKey: string, private readonly fallback: number) {}
128 − private roll() { const d = new Date().toISOString().slice(0, 10); if (d !== this.day) { this.day = d; this.used = 0; } }
143 + private local = 0;
144 + constructor(readonly provider: PremiumProvider, private readonly envKey: string, private readonly fallback: number) {}
145 + private roll() { const d = budgetDay(); if (d !== this.day) { this.day = d; this.local = 0; } }
129 146 limit(): number { return Number(process.env[this.envKey] ?? this.fallback); }
130 − remaining(): number { this.roll(); return Math.max(0, this.limit() - this.used); }
131 − spend(n = 1) { this.roll(); this.used += n; }
132 − snapshot() { this.roll(); return { day: this.day, used: this.used, limit: this.limit() }; }
147 + /** credits used today: the shared store's view when available, never below what this process spent itself */
148 + used(): number { this.roll(); return Math.max(this.local, budgetStore?.used(this.provider, this.day) ?? 0); }
149 + remaining(): number { return Math.max(0, this.limit() - this.used()); }
150 + spend(n = 1) { this.roll(); this.local += n; try { budgetStore?.spend(this.provider, this.day, n); } catch { /* store must not break fetching */ } }
151 + snapshot() { this.roll(); return { day: this.day, used: this.used(), local: this.local, limit: this.limit() }; }
133 152 }
134 −export const scrapflyBudget = new Budget("DCI_SCRAPFLY_DAILY_BUDGET", 400);
135 −export const firecrawlBudget = new Budget("DCI_FIRECRAWL_DAILY_BUDGET", 200);
153 +export const scrapflyBudget = new Budget("scrapfly", "DCI_SCRAPFLY_DAILY_BUDGET", 400);
154 +export const firecrawlBudget = new Budget("firecrawl", "DCI_FIRECRAWL_DAILY_BUDGET", 200);
136 155
137 156 export class FirecrawlFetcher implements Fetcher {
138 157 readonly name = "firecrawl" as const;
139 158