diff --git a/README.md b/README.md index eb2fcf7..5462987 100644 --- a/README.md +++ b/README.md @@ -64,10 +64,11 @@ both stacks, in the same partition), checks match `(_id, cd)` pairs — the retry copy's cd can never equal the migrated copy's. Preflight verifies the boundary is trustworthy (source frozen, clocks sane) before anything runs. -This README covers what you need BEFORE the dashboard exists (installing, -env vars, starting the service, automation reference). Everything after — -running, monitoring, troubleshooting, verifying — lives in the dashboard, -with `docs/RUNBOOK.md` as the cross-system procedure (cutover choreography, -Kafka retention, incident tables) for operators. - -## Architecture \ No newline at end of file +This README covers what you need BEFORE the dashboard exists: installing +and starting the service. `.env.example` is the commented configuration +reference (the two required variables and every optional one). Everything +after — running, monitoring, troubleshooting, verifying — lives in the +dashboard's **Migration Guide** and **Help & Recovery** tabs, with +`docs/RUNBOOK.md` as the standalone operations manual (terms, cutover +scenarios, incident table, sign-off procedure, curl reference) — start +there if you are planning a migration from scratch. \ No newline at end of file diff --git a/docs/RUNBOOK.md b/docs/RUNBOOK.md index 5a06e4d..f23deb3 100644 --- a/docs/RUNBOOK.md +++ b/docs/RUNBOOK.md @@ -1,17 +1,32 @@ # Migration Runbook -Operational procedure for migrating a customer's `drill_events` from MongoDB -to ClickHouse with this service. The guiding property: **after cutover, no +Operational procedure for migrating a deployment's `drill_events` data from +MongoDB to ClickHouse with this service. It assumes no prior knowledge of the +tool — terms are defined below, and every action is available both in the +dashboard and as a `curl` command. The guiding property: **after cutover, no failure anywhere in this flow can touch live data** — every incident response is *restart or resume*, never clean up or restore. Ingestion pauses exactly once, for minutes, at cutover — never for the migration. +## Terms used throughout + +| Term | Meaning | +|---|---| +| **Old cluster / source** | The MongoDB holding the `drill_events*` collections being migrated (sometimes a frozen clone of it — see the clone-source variant). | +| **New stack / target** | The new Countly architecture whose ClickHouse holds the `drill_events` table this service fills. | +| **cd** | Each document's server-side creation timestamp. The migration chunks, verifies and audits by cd; migrated rows keep their historical cd, live-ingested rows get post-cutover cds. | +| **Chunk** | One cd range of one collection — the unit of work, retry and verification. Chunk state lives in `mig_ranges` (the *ledger*) in `MANIFEST_DB`. | +| **DLQ** | Dead-letter queue (`mig_dlq_docs`): documents that could not or should not be migrated, stored with their full raw source so nothing is silently dropped. | +| **Tee / mirror** | A reverse-proxy (e.g. nginx) duplicating incoming SDK requests to both stacks; each side re-ingests independently, so the same event gets DIFFERENT `_id`/`cd` on each side. | +| **Bound** | `LEDGER_CD_UPPER_BOUND`: a cd ceiling — documents at/after it are never migrated. Required exactly when a tee is active (see the scenario table). | +| **Pod** | One instance of this service. Pods coordinate through chunk leases in MongoDB; any pod's dashboard shows the whole run. | + ## The flow -1. **Prepare** (old cluster still live, no customer impact) +1. **Prepare** (old cluster still live, no user-facing impact) - Deploy the new stack alongside the old. - Set Kafka `drill-events` retention to cover the migration window - (14 days default). Replication factor is the customer's redundancy + (14 days default). Replication factor is a redundancy choice — RF≥2 recommended for large instances; if RF=1, record the accepted risk (one broker disk loss forfeits the replay guarantee). - Bulk pre-copy the stateful set: apps & app keys, `app_users`, event @@ -25,7 +40,7 @@ once, for minutes, at cutover — never for the migration. 3. **Rehearse** — dry run with `DRY_RUN=1` (≤5% stratified sample against a Null-engine clone; full ClickHouse validation, nothing stored). Review - `GET /report` (skips, coercions per key, DLQ) with the customer, sign off. + `GET /report` (skips, coercions per key, DLQ) with whoever owns sign-off. 4. **Cutover** — stop old ingestion → sync the stateful-set delta since the pre-copy (changed users via last-seen; aggregated data must land BEFORE @@ -41,9 +56,27 @@ once, for minutes, at cutover — never for the migration. live ingestion. Watch `/viz`; the invariant monitor spot-checks continuously. -6. **Finish** — all chunks done → final `GET /report` → customer sign-off → +6. **Finish** — all chunks done → Final check green → sign-off → revert Kafka retention → decommission old cluster. +## Choose your scenario first + +The one decision that changes the configuration is whether a TEE mirrors +the same requests into both stacks. Everything else is shared machinery. + +| # | Topology | LEDGER_CD_UPPER_BOUND | Ingestion switch | New data arriving in old Mongo | Sign-off | +|---|---|---|---|---|---| +| 1 | Two clusters, **no mirroring** (plain switch) | **UNSET** | Before the migration (cutover-first) or after the bulk (bulk-before-cutover + final drain) | **Migrated** — top-up passes chase it until the drain finds nothing | Verify + audits, DLQ = 0 | +| 2 | Two clusters, **mirror old → new** (old primary) | **SET** = tee flip | At sign-off | **Never migrated past the bound** — it is the tee's copy (different _id/cd; duplicates would be undetectable) | Verify + audits for pre-bound; dashboard comparison + sync parity for post-bound | +| 3 | Two clusters, **mirror new → old** (new primary, old = rollback net) | **SET** = the moment new became primary | Already happened at the flip | Same as 2 — post-flip old-side docs are mirror copies | Same as 2 | +| 4 | **Single cluster, in-place upgrade** (drill mongo → ClickHouse in background) | **UNSET** | The upgrade itself is the switch; old drill collections freeze | Transition tail drained by top-up; no tee → nothing to duplicate | Verify + audits, DLQ = 0 (live-parallel path; backpressure protects prod CH) | + +Scenario is also selectable on the dashboard's **Migration Guide** tab — +it renders the per-scenario checklist and states the bound requirement. +For 2 and 3: use **Detect boundary** + **Apply this bound to the run** +(one click covers all pods), verify the `bounded · cd < …` badge on every +pod, and keep re-running sync parity during the validation window. + ## Incident responses | Incident | What happens | Operator action | @@ -55,10 +88,111 @@ once, for minutes, at cutover — never for the migration. | A doc CRASHES the process every time (poison pill) | After 3 crash-retries the chunk is auto-split instead of retried; repeated splitting converges on a ≤1-min window quarantined as a tiny failed chunk — everything else migrates (verified: 20k-doc drill localized 1 poison doc to a 2-doc window in 25 restarts) | Inspect the few source docs in the failed chunk's cd window; fix/remove them, then `POST /control/retry-failed` | | Live ClickHouse itself must be rebuilt | Live events still sit in the Kafka log; history still sits in frozen Mongo | Recreate table → reset ONLY the ClickHouse-sink connector's offsets to earliest (aggregator groups untouched) → re-run the migrator | +## Final check — the one-click sign-off + +Don't interpret audit buckets by hand: the **Final check** runs everything +(chunk states, DLQ, full source recount, cd-checksum fingerprints, sampled +content comparison), applies the tee/cutover rules itself, and answers the +only question that matters — *is it safe to decommission the old cluster?* — +as **PASS / PASS WITH NOTES / FAIL** in plain sentences with the action named +on every red line. + +Two tiers: + +- **Quick** (default — minutes): chunk states + DLQ + target-vs-ledger + verification (every migrated window's live count against the recorded + count, plus duplicate attribution) + random content samples against the + source. Catches everything that can happen AFTER reading. Capped at + PASS WITH NOTES, and its headline never authorizes teardown — the note + names what it did not re-prove. +- **Deep** (opt-in — hours on large runs): additionally recounts EVERY + window against the source with cd-checksum fingerprints and sampled + identity coverage. It is not distrust of the ledger — chunk reads are + already recounted against the source at migration time — it is the only + check that derives everything from the two databases alone, with zero + reliance on the tool's own records. Run it once, as the gate before the + source is deleted; while the source exists, quick is enough. + +What each layer can and cannot see: + +| Failure class | Caught by | +|---|---| +| Under-read at read time (source count ≠ read tally) | the migration itself, per chunk (source-count guard) | +| Rows lost or duplicated in ClickHouse after attach | quick — ledger-vs-target verify | +| Skipped documents | DLQ accounting (both tiers; unresolved = FAIL) | +| Wrong content in migrated rows | quick — random content samples vs source | +| Docs written into already-done windows later (imports, restores, backdated cds) | deep only — the ledger is blind to them by design | +| Count-preserving identity swaps (same count, same cd-sum, different docs) | deep only — checksum + sampled id coverage | +| Routine retention deleting source docs (drift) | deep classifies it exactly (spot-checked, never assumed benign) | + +- Dashboard: the **Final check** card → *Run final check* (tick *deep source + recount* for the pre-teardown gate). On tee/mirror runs without a stored + bound, type the cutover time into the field first. +- SSH-only: + +```bash +# quick (add {"cutoverMs": } for mirror runs without a stored bound) +curl -s -X POST localhost:PORT/control/final-check -H 'content-type: application/json' -d '{}' +# deep — before deleting the source +curl -s -X POST localhost:PORT/control/final-check -H 'content-type: application/json' -d '{"deep": true}' +# read the verdict (re-run until it says PASS/FAIL; shows progress while running) +curl -s localhost:PORT/final-check.txt +``` + +A stored/env cd bound is picked up automatically as the cutover. Post-cutover +source windows are excluded and explained in a note — divergence there is the +mirror still feeding the old side, not data loss. Run it while the old +cluster is still up: the source is the reference. + +## Tee-overlap dedupe — fixing a missing bound after the fact + +A mirrored cutover migrated WITHOUT `LEDGER_CD_UPPER_BOUND` copies the +mirror's re-ingested docs on top of natively ingested rows: every event in +the overlap window (tee flip → migration completion) exists twice in +ClickHouse. The copies are separable — the migrated copy's `_id` exists in +the old cluster's Mongo; the native one's doesn't. **Must run before the old +cluster is decommissioned** (old Mongo is the separator). + +An id match alone is not proof of duplication: if the tee (or the new +side's ingestion) dropped a request, the migrated row is the ONLY copy of +that event. Every hour bucket therefore needs count-evidence of native +counterparts — `native = live − matched` must roughly cover `matched` — +before anything in it is deleted. Buckets that fall short are skipped and +reported (`unsafe` in the result); review those hours (tee outage? wrong +start time?) instead of forcing them. Collections without their own +(a,e,n) scope (e.g. a base `drill_events` collection) have no usable +native-counterpart evidence — their matches are always reported as unsafe +and never deleted. The check is strict (zero slack) by +default; `slackPct` (≤5) may be passed consciously to absorb ingest-timing +straddle at bucket edges. Known limit: a loss exactly offset by +mirror-dropped natives in the same hour is invisible to count evidence — +an EMPTY dry run means no duplicates (skip the step; never widen the window +to make it match something). Both dedupe (dry run included) and the Final +check refuse while any pod still holds an active chunk claim. + +There is a dashboard card for this (Overview → **Tee-overlap dedupe**: +enter the window, *Dry run* first — *Delete duplicates* unlocks only after +it) as well as the endpoints below. + +```bash +# 1. DRY RUN (counts only): fromMs = tee flip / IP swap, toMs = migration completion +curl -s -X POST localhost:PORT/control/dedupe-overlap -H 'content-type: application/json' \ + -d '{"fromMs": 1789700000000, "toMs": 1789794970435}' +curl -s localhost:PORT/api/dedupe-overlap # totals.chMatched = the duplicates +# 2. EXECUTE (refused unless the dry run over the SAME window completed first) +curl -s -X POST localhost:PORT/control/dedupe-overlap -H 'content-type: application/json' \ + -d '{"fromMs": 1789700000000, "toMs": 1789794970435, "execute": true}' +# 3. re-run the Final check with the same cutover to confirm +``` + ## Verification cheat sheet ```sql --- exactness (instant, exact): +-- quick orientation only — on a LIVE target this total moves with ingestion +-- and proves nothing about the migration. Sign-off relies on the windowed, +-- scoped checks (Final check / verify): migrated rows keep historical cd, +-- live rows are insert-stamped, so every audited window excludes live data +-- by construction. SELECT count() AS total, uniqExact(_id) AS distinct_ids FROM countly_drill.drill_events; -- full re-verification of the whole migration in minutes: -- grouped count per chunk window vs the ledger's rows_expected (mig_ranges) @@ -78,33 +212,39 @@ old per-event collections only after their chunks are done and signed off), and hard memory limits on the new components — an OOM there is a production incident. -## Validation before a customer run +## Validation before a production run `bench/README.md`: seed → straight run (counts must be exact) → SIGKILL crash drill → optionally `bench/seed-failures.ts` for a full failure-scenario drill (breaker, DLQ, monitor, retry-failed). -## Choose your scenario first - -The one decision that changes the configuration is whether a TEE mirrors -the same requests into both stacks. Everything else is shared machinery. - -| # | Topology | LEDGER_CD_UPPER_BOUND | Ingestion switch | New data arriving in old Mongo | Sign-off | -|---|---|---|---|---|---| -| 1 | Two clusters, **no mirroring** (plain switch) | **UNSET** | Before the migration (cutover-first) or after the bulk (bulk-before-cutover + final drain) | **Migrated** — top-up passes chase it until the drain finds nothing | Verify + audits, DLQ = 0 | -| 2 | Two clusters, **mirror old → new** (old primary) | **SET** = tee flip | At customer sign-off | **Never migrated past the bound** — it is the tee's copy (different _id/cd; duplicates would be undetectable) | Verify + audits for pre-bound; dashboard comparison + sync parity for post-bound | -| 3 | Two clusters, **mirror new → old** (new primary, old = rollback net) | **SET** = the moment new became primary | Already happened at the flip | Same as 2 — post-flip old-side docs are mirror copies | Same as 2 | -| 4 | **Single cluster, in-place upgrade** (drill mongo → ClickHouse in background) | **UNSET** | The upgrade itself is the switch; old drill collections freeze | Transition tail drained by top-up; no tee → nothing to duplicate | Verify + audits, DLQ = 0 (live-parallel path; backpressure protects prod CH) | - -Scenario is also selectable on the dashboard's **Migration Guide** tab — -it renders the per-scenario checklist and states the bound requirement. -For 2 and 3: use **Detect boundary** + **Apply this bound to the run** -(one click covers all pods), verify the `bounded · cd < …` badge on every -pod, and keep re-running sync parity during the validation window. - -## Tee-mirror cutover (customer keeps the old architecture until sign-off) - -For customers who require approval before switching: the old arch stays +## Clone-source variant (migrate from a frozen copy) + +An OPTIONAL variant of scenarios 1–3 — most migrations run against the +LIVE source cluster, which is fully supported (in unbounded modes, top-up +passes keep chasing data that arrives during the run). The variant: pause +old ingestion, clone the source MongoDB, resume ingestion on the NEW stack +(optionally mirroring back to the old one as the rollback net), and migrate +from the clone. SDK offline queues absorb the pause. Properties: + +- The source is frozen at the clone moment, so no bound is needed and top-up + finds nothing — the startup guard asks its one question at start: + **Proceed unbounded is correct** here. +- Parity/audit tables compare against the CLONE: zeros after the clone + moment mean "clone taken here", not a dead mirror. The live old-side + MongoDB is invisible to the tool. +- Cloned INSIDE the ingestion pause (the sequence above) → the clone can + never hold a natively-ingested event's mirror copy: **no duplicates, no + bound, no dedupe** — the cleanest possible run. Only a clone taken AFTER + ingestion resumed has a duplicated tail: dedupe with exactly + [ingestion-resume, clone-moment], never earlier. +- Any doc-count comparison against the live old-arch Mongo will drift by + everything ingested after T-clone — compare against the clone, or scope + counts to cd < T-clone. + +## Tee-mirror cutover (keep the old architecture until sign-off) + +When approval is required before switching: the old stack stays authoritative, nginx TEES the same SDK requests to the new architecture (which re-ingests them with its own logic — drill, sessions, aggregations, profiles all populate natively), and the bulk migration backfills history @@ -132,7 +272,7 @@ can deduplicate across that seam — the ONLY protection is the time bound. region; post-bound windows show as pending/uncovered, never as defects. The post-bound region is the tee's responsibility and is validated by comparing dashboards between the two systems, not by this tool. -5. Customer validates side-by-side as long as needed; both systems ingest +5. Validate side-by-side as long as needed; both systems ingest the same requests the whole time. 6. On approval: point SDK traffic solely at the new arch, drop the tee, decommission old ingestion on its own schedule. @@ -147,6 +287,52 @@ Caveats: - Retention TTL keeps deleting on the old side throughout — the source audit reports that as deletion drift, not as a defect. +### One-call boundary setting (SSH / API) + +The whole detect-and-apply flow is a single endpoint: + +```bash +# detect, and apply automatically when the seam is an exact ingestion-pause gap +curl -s -X POST localhost:PORT/control/set-boundary -H 'content-type: application/json' -d '{}' +# read the outcome — the apply receipt lands in .applied +curl -s localhost:PORT/api/boundary +# no exact gap? review the report, then accept the anchor explicitly… +curl -s -X POST localhost:PORT/control/set-boundary -H 'content-type: application/json' -d '{"acceptAnchor": true}' +# …or set the bound to a known timestamp directly (applies immediately) +curl -s -X POST localhost:PORT/control/set-boundary -H 'content-type: application/json' -d '{"boundMs": 1789966140000}' +``` + +An exact gap applies unattended; an anchor (quantified ambiguity) is never +auto-applied without `acceptAnchor`. The dashboard flow and the separate +`/control/detect-boundary` + `/control/apply-bound` endpoints keep working. + +### The startup guard — the bound mistake, made impossible to miss + +A FRESH run that finds its target ClickHouse already receiving live data, +with no cd bound set, **holds before mapping** (`pauseReason: +boundary-unset`). That is exactly the setup where an unset bound either +duplicates the overlap window (mirror active) or is a deliberate choice +(cutover-first / in-place, where new data must still be migrated). The tool +cannot tell those apart from data alone, so it asks — once: + +- mirror active → apply the bound (`POST /control/set-boundary`, the + dashboard card, or `LEDGER_CD_UPPER_BOUND`); the run releases itself, or +- nothing mirrors traffic → click **Proceed unbounded** in the banner, or + `curl -X POST localhost:PORT/control/allow-unbounded` (cluster-wide, + releases every held pod), or deploy with `LEDGER_UNBOUNDED_OK=1`. + +A plain Resume is deliberately ignored while the question is open — only a +bound or the no-mirror answer releases the hold. There is NO automatic +verdict: even a provably empty target proves only current emptiness (a tee +that has not carried its first request yet looks identical to no-mirror), +so the question is answered exactly once per run, explicitly — the +dashboard's Proceed unbounded, `POST /control/allow-unbounded`, a bound, or +`LEDGER_UNBOUNDED_OK=1` in the deployment for topologies known to have no +mirror (required for fire-and-forget Jobs). The answer is stored +cluster-wide, releasing every held pod. The corollary: a mirror enabled +AFTER the answer is invisible to the guard by construction — enabling any +mirror is exactly when to run the sync-parity card. + ### Bound is opt-in — pick the mode deliberately | Situation | LEDGER_CD_UPPER_BOUND | Behavior | @@ -169,7 +355,7 @@ in either the old or the new system. silently dropped), the waive is recorded and counted, and the source audit attributes each window's shortfall to its waived docs — sign-off stays exact. Only consider a sentinel-uid replay instead if the affected volume -is large enough to distort historical event totals for a customer AND the +is large enough to distort historical event totals for an app AND the docs carry usable ts/did (check a few samples in the DLQ panel first). Chunks that were 100% such docs complete as done (structured skips do not @@ -194,8 +380,8 @@ failed, DLQ, status + pause reason). No network access needed: kubectl logs -f deploy/drill-migrator | grep 'progress heartbeat' docker logs -f drill-migrator-p1 2>&1 | grep 'progress heartbeat' ``` -These lines also flow into the stack's log pipeline (alloy → Loki), so -Grafana log panels/alerts work with zero extra plumbing. +Because they go to stdout, they flow into whatever log pipeline collects +container output (Loki, ELK, CloudWatch, …) with zero extra plumbing. **3. Actions via curl** (same endpoints the buttons call; POSTs need the JSON content type): diff --git a/src/config/loader.ts b/src/config/loader.ts index 6b436d9..0dc0fe0 100644 --- a/src/config/loader.ts +++ b/src/config/loader.ts @@ -26,6 +26,7 @@ function envToRawConfig(env: NodeJS.ProcessEnv) { cdUpperBoundMs: env.LEDGER_CD_UPPER_BOUND, captureTransformErrors: env.LEDGER_CAPTURE_TRANSFORM_ERRORS, startPaused: env.LEDGER_START_PAUSED, + unboundedOk: env.LEDGER_UNBOUNDED_OK, dryRun: env.DRY_RUN, dryRunSamplePct: env.DRY_RUN_SAMPLE_PCT, }, diff --git a/src/config/schema.ts b/src/config/schema.ts index 9605a2f..672ec70 100644 --- a/src/config/schema.ts +++ b/src/config/schema.ts @@ -82,6 +82,8 @@ export const configSchema = z.object({ // click starts the whole fleet, pods that join later start // immediately, and a pod that restarts after Start stays started. startPaused: booleanFromEnv.default(false), + /** Explicit no-mirror declaration: skips the unbounded-with-live-target startup guard. */ + unboundedOk: booleanFromEnv.default(false), // Dry run: sampled rehearsal against a Null-engine clone. dryRun: booleanFromEnv.default(false), dryRunSamplePct: numberFromEnv.default(2).pipe(z.number().min(0.1).max(5)), diff --git a/src/http/ledger-viz-route.ts b/src/http/ledger-viz-route.ts index b336a87..ee9d693 100644 --- a/src/http/ledger-viz-route.ts +++ b/src/http/ledger-viz-route.ts @@ -300,6 +300,16 @@ const PAGE = ` Destructive actions ask for a second click. Every action shows a receipt. +
+

Final check — one click answers: is it safe to decommission the old source? (chunks + DLQ + full source recount + checksums + content samples — interpreted for you)

+
+ + + +
+
Not run. Run it after the migration completes — it recounts every window against the source, so give it time on big runs; progress shows here. SSH-only: curl -X POST :PORT/control/final-check then curl :PORT/final-check.txt
+
+

Collections

Waiting for first chunk…
@@ -333,6 +343,17 @@ const PAGE = `
+
+

Tee-overlap dedupe (fix a mirrored run that migrated WITHOUT the cd bound: remove the migrated copies of events the new cluster already ingested natively)

+
+ + + + +
+
Only for runs that migrated a mirrored setup unbounded. Every hour bucket is checked for count-evidence of native counterparts before anything is deleted — buckets where migrated rows are the ONLY copy are skipped and reported. Old-cluster Mongo must still be up. Start must be AT or AFTER the actual flip: too early deletes real data, too late only leaves a few duplicates.
+
+

Dead-letter queue (unmigratable docs, stored with their full raw source — replay after a fix, or waive)

@@ -427,7 +448,7 @@ const PAGE = `

-

Then: final report (/report), customer sign-off, revert Kafka retention, decommission the old cluster.

+

Then: final report (/report), sign-off, revert Kafka retention, decommission the old cluster.

@@ -513,7 +534,7 @@ async function control(action, okMsg, btn, needsConfirm) { btn.dataset.label = btn.textContent; btn.textContent = 'Click again to confirm'; btn.classList.add('armed'); - setTimeout(() => { armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); }, 4000); + setTimeout(() => { armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); }, 8000); return; } if (btn) { armed.delete(btn); if (btn.dataset.label) { btn.textContent = btn.dataset.label; btn.classList.remove('armed'); } btn.disabled = true; } @@ -534,7 +555,7 @@ async function startRebuild(btn, force) { btn.dataset.label = btn.textContent; btn.textContent = 'Click again to confirm'; btn.classList.add('armed'); - setTimeout(() => { armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); }, 4000); + setTimeout(() => { armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); }, 8000); return; } armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); btn.disabled = true; @@ -720,6 +741,138 @@ function updatePhaseBadges() { } updatePhaseBadges(); +async function allowUnbounded(btn) { + if (!armed.get(btn)) { + armed.set(btn, true); + btn.dataset.label = btn.textContent; + btn.textContent = 'Click again to confirm: NOTHING mirrors traffic'; + btn.classList.add('armed'); + setTimeout(function () { armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); }, 8000); + return; + } + armed.delete(btn); btn.disabled = true; + try { + var res = await fetch('/control/allow-unbounded', { method: 'POST' }); + var out = await res.json(); + toast(out.allowed ? '\u2705 no-mirror declared \u2014 held pods release within seconds' : '\u274c ' + (out.reason || res.status)); + } catch (e) { toast('\u274c ' + e.message); } +} + +function ddWindow() { + var f = Date.parse((document.getElementById('dd-from').value || '').trim()); + var t = Date.parse((document.getElementById('dd-to').value || '').trim()); + if (isNaN(f) || isNaN(t) || !(f < t)) { toast('Enter both times as ISO (e.g. 2026-09-18T18:00Z), start before end'); return null; } + return { fromMs: f, toMs: t }; +} +async function startDedupe(btn, execute) { + var w = ddWindow(); + if (!w) return; + if (execute && !armed.get(btn)) { + armed.set(btn, true); + btn.dataset.label = btn.textContent; + btn.textContent = 'Click again to DELETE the counted duplicates'; + btn.classList.add('armed'); + setTimeout(function () { armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); }, 8000); + return; + } + if (execute) { armed.delete(btn); btn.textContent = btn.dataset.label; btn.classList.remove('armed'); } + btn.disabled = true; + try { + var res = await fetch('/control/dedupe-overlap', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify({ fromMs: w.fromMs, toMs: w.toMs, execute: !!execute }) }); + var out = await res.json(); + if (!out.started) toast('Not started: ' + (out.reason || 'unknown')); + else toast(execute ? 'Deleting duplicates\u2026' : 'Dry run started \u2014 counting duplicates'); + } catch (e) { toast('failed: ' + e.message); } + btn.disabled = false; + pollDedupe(); +} +var ddTimer = null; +async function pollDedupe() { + try { + var dd = await fetch('/api/dedupe-overlap').then(function (r) { return r.json(); }); + renderDedupe(dd); + if (dd.status === 'running') { clearTimeout(ddTimer); ddTimer = setTimeout(pollDedupe, 2000); } + } catch (e) { /* engine restarting */ } +} +function renderDedupe(dd) { + var el = document.getElementById('dedupe-out'); + if (!el || !dd || dd.status === 'not_run') return; + var execBtn = document.getElementById('btn-dd-exec'); + if (dd.status === 'running') { el.innerHTML = '
Running \u2014 ' + fcEsc(dd.phase) + '
'; return; } + if (dd.status === 'failed') { el.innerHTML = '
Failed: ' + fcEsc(dd.error) + '
'; return; } + var t = dd.totals || {}; + var unsafeN = 0; + (dd.collections || []).forEach(function (c) { unsafeN += (c.unsafe || []).length; }); + var html = '

' + (dd.execute + ? '\u2705 Deleted ' + fmt(t.deleted) + ' duplicate row(s).' + : 'Dry run: ' + fmt(t.chMatched) + ' migrated row(s) match old-cluster ids in the window (' + fmt(t.mongoDocsInWindow) + ' old-side docs scanned). Nothing deleted.') + '

'; + if (t.unsafeMatched > 0) { + html += '

\u26a0 ' + fmt(t.unsafeMatched) + ' matched row(s) in ' + unsafeN + ' bucket(s) were NOT ' + (dd.execute ? 'deleted' : 'counted as deletable') + ': they lack count-evidence of a native counterpart, or belong to a collection without its own (a,e,n) scope \u2014 there the migrated row may be the ONLY copy. Review those buckets (tee outage / wrong start time / base collection?) before touching them.

'; + } + if (!dd.execute && dd.lastDryRun) { + html += '

Window measured \u2014 the Delete button is now enabled for this exact window.

'; + if (execBtn) { execBtn.disabled = false; execBtn.title = ''; } + } + html += '

window: ' + new Date(dd.fromMs).toISOString() + ' \u2192 ' + new Date(dd.toMs).toISOString() + ' \u00b7 ' + (dd.collections || []).length + ' collection(s) with matches

'; + el.innerHTML = html; +} + +async function startFinalCheck(btn) { + var body = {}; + var deepEl = document.getElementById('fc-deep'); + if (deepEl && deepEl.checked) body.deep = true; + var cutRaw = (document.getElementById('fc-cutover').value || '').trim(); + if (cutRaw) { + var ms = Date.parse(cutRaw); + if (isNaN(ms)) { toast('Could not parse the cutover time — use ISO like 2026-09-18T18:00Z'); return; } + body.cutoverMs = ms; + } + btn.disabled = true; + try { + var res = await fetch('/control/final-check', { method: 'POST', headers: { 'content-type': 'application/json' }, body: JSON.stringify(body) }); + var out = await res.json(); + if (!out.started) toast('Not started: ' + (out.reason || 'unknown')); + else toast('Final check started — the verdict will appear below'); + } catch (e) { toast('failed: ' + e.message); } + btn.disabled = false; + pollFinalCheck(); +} +var fcTimer = null; +async function pollFinalCheck() { + try { + var fc = await fetch('/api/final-check').then(function (r) { return r.json(); }); + renderFinalCheck(fc); + if (fc.status === 'running') { clearTimeout(fcTimer); fcTimer = setTimeout(pollFinalCheck, 2000); } + } catch (e) { /* engine restarting — next poll or reload recovers */ } +} +function fcEsc(s) { var d = document.createElement('div'); d.textContent = s == null ? '' : String(s); return d.innerHTML; } +function renderFinalCheck(fc) { + var el = document.getElementById('finalcheck-out'); + if (!el || !fc || fc.status === 'not_run') return; + if (fc.status === 'running') { + var a = fc.audit || {}; + el.innerHTML = '
Running — ' + fcEsc(fc.phase) + (a.collectionsTotal ? ' (' + a.collectionsDone + '/' + a.collectionsTotal + ' collections)' : '') + '
'; + return; + } + if (fc.status === 'failed') { + el.innerHTML = '
The check itself failed to complete: ' + fcEsc(fc.error) + ' — a tooling error, not a data verdict. Re-run it.
'; + return; + } + var pal = fc.verdict === 'PASS' ? ['#E4F6EC', '#157A45'] : fc.verdict === 'PASS_WITH_NOTES' ? ['#FDEEDD', '#A05A16'] : ['#FDECEC', '#B3261E']; + var badge = (fc.verdict === 'PASS' ? 'PASS' : fc.verdict === 'PASS_WITH_NOTES' ? 'PASS WITH NOTES' : 'FAIL') + + (fc.mode === 'deep' ? ' \u00b7 deep (full source recount)' : ' \u00b7 quick (ledger verify + samples)'); + var html = '
' + + '
' + badge + '
' + + '
' + fcEsc(fc.headline) + '
' + + ''; + if (fc.cutoverMs) html += '

cutover applied: source compared only for cd < ' + new Date(fc.cutoverMs).toISOString() + '

'; + el.innerHTML = html; +} + async function tick() { try { const [stats, chunkResp] = await Promise.all([ @@ -760,23 +913,39 @@ async function tick() { if (changes >= 3 && (last.t - windowStart.t) >= 10000) break; } var rspan = (last.t - windowStart.t) / 1000; - var liveRate = rspan >= 10 ? Math.max(0, (last.d - windowStart.d) / rspan) : null; - var multiPod = stats.cluster && stats.cluster.pods > 1; + // the client window is only trustworthy once it has WITNESSED chunk + // completions — a freshly opened tab on a huge-chunk run showed "0" + // (field report); until then the server's 10-min ledger window is truth + var clientReady = changes >= 3 && rspan >= 10; + var slowRate = stats.clusterSlow && stats.clusterSlow.docsPerSecond > 0 ? stats.clusterSlow.docsPerSecond : null; + var liveRate = clientReady ? Math.max(0, (last.d - windowStart.d) / rspan) + : slowRate !== null ? slowRate + : null; + var podsSeen = Math.max(stats.cluster ? stats.cluster.pods : 0, stats.clusterSlow ? stats.clusterSlow.pods : 0); + var multiPod = podsSeen > 1; var effRate = liveRate !== null ? liveRate - : multiPod && stats.status === 'running' ? stats.cluster.docsPerSecond + : multiPod && stats.status === 'running' ? (stats.cluster ? stats.cluster.docsPerSecond : 0) : stats.docsPerSecond; if (stats.status === 'completed') { - // a pod restarted after completion migrated nothing itself — its - // local average is 0 and would read as an anomaly - dpsEl.textContent = stats.docsPerSecond >= 1 ? fmt(stats.docsPerSecond) + ' avg' : '\u2013'; + // whole-RUN average from ledger docs + run timeline; the pod's own + // lifetime counter is only its share of a multi-pod run (field: a + // 4-pod run showed 5,060 instead of the run's ~20,300), and a pod + // restarted after completion migrated nothing at all + var rt = stats.runTimes || {}; + var runSec = rt.startedAtMs && rt.completedAtMs ? (rt.completedAtMs - rt.startedAtMs) / 1000 : 0; + var runAvg = runSec > 0 && sum.docsDone > 0 ? sum.docsDone / runSec : stats.docsPerSecond; + dpsEl.textContent = runAvg >= 1 ? fmt(Math.round(runAvg)) + ' avg' : '\u2013'; } else if (liveRate !== null) { - dpsEl.textContent = fmt(Math.round(liveRate)) + (multiPod ? ' \u00b7 ' + stats.cluster.pods + ' pods' : ''); + dpsEl.textContent = fmt(Math.round(liveRate)) + (multiPod ? ' \u00b7 ' + podsSeen + ' pods' : ''); } else if (stats.status === 'running') { dpsEl.textContent = 'measuring\u2026'; } else { dpsEl.textContent = '\u2013'; } - document.getElementById('s-skipped').textContent = fmt(stats.totalDocsSkipped); + // ledger truth — each pod's in-memory counter only knows its own share + // (field: a 3-pod run showed 100,623 while the DLQ held 314,125) + document.getElementById('s-skipped').textContent = + fmt(Math.max(sum.docsSkipped || 0, stats.totalDocsSkipped || 0)); // ledger truth, not this pod's counter — in multi-pod each pod only // counts its own failures, so the card under-reported cluster-wide document.getElementById('s-failed').textContent = fmt((sum.byStatus || {}).failed || 0); @@ -858,17 +1027,31 @@ async function tick() { var hint = document.getElementById('pause-hint'); if (isPaused) { hint.style.display = ''; - hint.textContent = stats.pauseReason === 'not-started' - ? '\u23f8 NOT STARTED \u2014 deployed and waiting. Nothing has been read, mapped or indexed yet; ' - + 'run preflight, build indexes and rehearse first, then click Start to begin the run (all pods).' - : '\u23f8 ENGINE PAUSED' + - (stats.pauseReason === 'breaker-transient' ? ' (backend outage \u2014 auto-resume armed)' : - stats.pauseReason === 'breaker-data' ? ' (systematic data problem \u2014 needs you)' : ' (by operator)') + - ' \u2014 Retry / Replay / Waive only QUEUE work; click Resume to process it.'; + if (stats.pauseReason === 'boundary-unset') { + hint.innerHTML = '\u26a0 HELD BY THE BOUNDARY GUARD \u2014 the target ClickHouse is already receiving live data and no cd bound is set. ' + + 'If a mirror re-ingests the same requests on both sides, running unbounded WILL duplicate the overlap window. ' + + 'Either apply a bound (Tee boundary card below), or \u2014 if NOTHING mirrors traffic between the stacks \u2014 ' + + ''; + } else { + hint.textContent = stats.pauseReason === 'not-started' + ? '\u23f8 NOT STARTED \u2014 deployed and waiting. Nothing has been read, mapped or indexed yet; ' + + 'run preflight, build indexes and rehearse first, then click Start to begin the run (all pods).' + : '\u23f8 ENGINE PAUSED' + + (stats.pauseReason === 'breaker-transient' ? ' (backend outage \u2014 auto-resume armed)' : + stats.pauseReason === 'breaker-data' ? ' (systematic data problem \u2014 needs you)' : ' (by operator)') + + ' \u2014 Retry / Replay / Waive only QUEUE work; click Resume to process it.'; + } } else { hint.style.display = 'none'; } var prBtn = document.getElementById('btn-pauseresume'); if (prBtn) { - if (isPaused) { + if (isPaused && stats.pauseReason === 'boundary-unset') { + // a plain Resume cannot answer the mirror question — the banner + // above carries the two real actions (bound / proceed unbounded) + prBtn.dataset.action = ''; + prBtn.innerHTML = '\u25b6 Resume'; + prBtn.classList.remove('primary'); + prBtn.disabled = true; + } else if (isPaused) { prBtn.dataset.action = 'resume'; prBtn.innerHTML = stats.pauseReason === 'not-started' ? '\u25b6 Start' : '\u25b6 Resume'; prBtn.classList.add('primary'); @@ -998,6 +1181,7 @@ var SCENARIOS = [ '
  • LEDGER_CD_UPPER_BOUND: LEAVE UNSET. The migration must take everything, including data still arriving in the old cluster \u2014 top-up passes chase it until the final drain finds nothing new.
  • ' + '
  • Cutover-first: switch SDK ingestion to the new cluster, then run the migration (old drill data is frozen). Bulk-before-cutover: run the bulk first, switch ingestion, then let the final top-up pass drain the tail.
  • ' + '
  • Ignore the Tee boundary card \u2014 it is for mirrored setups only. Applying a bound here would ORPHAN newly arrived data.
  • ' + + '
  • If ingestion already switched to the new cluster before the run starts, the startup guard will hold and ask \u2014 Proceed unbounded is the correct answer for this scenario.
  • ' + '
  • Sign-off: Verify + Audit vs source + content audit, DLQ pending = 0.
  • ' + '' }, { id: 'tee-old', name: '2 \u00b7 Mirror old \u2192 new', @@ -1012,7 +1196,7 @@ var SCENARIOS = [ '' }, { id: 'tee-new', name: '3 \u00b7 Mirror new \u2192 old', bound: true, - html: '

    New cluster is already primary; nginx mirrors back to the old stack as the customer\u2019s rollback safety net during validation.

    ' + + html: '

    New cluster is already primary; nginx mirrors back to the old stack as the rollback safety net during validation.

    ' + '