feat(scripts): phase 07 — reverse-backfill scripts + delete guard

Pre-execution prerequisites for the Phase 07 cutover. Stage 2 of the
cutover keeps DUAL_WRITE=0 for ~6 days; if anything regresses during
that window the operator MUST be able to roll back to KV/D1 with the
last N days of Mongo-only writes recovered. Pre-building these scripts
(per code-reviewer #4) eliminates "draft a backfill under outage
pressure" — the anti-pattern of writing untested code at 4am.

Reverse-backfill
- scripts/backfill-mongo-to-kv.js: full-scan Mongo collection per module,
  PUT each doc back to CF KV via REST. expiresAt → expirationTtl (clamped
  to 60s minimum per CF KV); already-expired docs are skipped (won't
  resurrect dead state). 50 ops/sec throttle. --dry-run + --module flags.
- scripts/backfill-mongo-to-d1.js: full-scan trading_trades, build INSERT
  SQL preserving legacy_id where present (round-trips D1 autoincrement IDs
  preserved by phase-05 forward backfill). Sequential int generation for
  any docs without legacy_id. Pipes through wrangler d1 execute.
- scripts/lib/migration-helpers.js: cfKvPut helper added.

Delete guard (debugger #12)
- scripts/wrangler-delete-guard.sh: interactive CONFIRM wrapper around
  wrangler kv namespace delete + wrangler d1 delete. Exits 3 when stdin
  is not a tty so it cannot run in CI. Documented: never run in CI.

package.json: backfill:mongo:kv[:dry] + backfill:mongo:d1[:dry] scripts
wired.

Tests: 697 → 733 (+36).
- 7 cfKvPut tests (REST URL, querystring, body, expiration_ttl param).
- 10 reverse-KV TTL math tests (expired sentinel, future seconds, no-TTL,
  CF 60s minimum clamp).
- 9 reverse-D1 SQL construction tests (escaping, legacy_id preservation,
  sequential generation).

Lint clean. No Worker code touched. Stage 1 cutover, 7-day soak,
snapshots, and Stage 3 cleanup (delete CFKVStore + simplify factories +
edit package.json deploy chain) remain operator-driven and will be
committed separately after binding deletion.
This commit is contained in:
tiennm99 committed 2026-04-26 09:29:14 +07:00
1 parent 439e874a11
commit 54f4c28807
8 files changed
+854

No files matched your search

+4
View File
@@ -22,6 +22,10 @@
"backfill:kv:dry": "node --env-file-if-exists=.env.deploy scripts/backfill-kv-to-mongo.js --dry-run",
"backfill:d1": "node --env-file-if-exists=.env.deploy scripts/backfill-d1-to-mongo.js",
"backfill:d1:dry": "node --env-file-if-exists=.env.deploy scripts/backfill-d1-to-mongo.js --dry-run",
"backfill:mongo:kv": "node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-kv.js",
"backfill:mongo:kv:dry": "node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-kv.js --dry-run",
"backfill:mongo:d1": "node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-d1.js",
"backfill:mongo:d1:dry": "node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-d1.js --dry-run",
"verify:mongo": "node --env-file-if-exists=.env.deploy scripts/verify-mongo-parity.js",
"wipe:mongo": "node --env-file-if-exists=.env.deploy scripts/wipe-mongo.js",
"analyze:soak": "node scripts/analyze-soak.js",
+235
View File
@@ -0,0 +1,235 @@
#!/usr/bin/env node
/**
* @file backfill-mongo-to-d1 — emergency reverse-backfill: MongoDB → Cloudflare D1.
*
* Reads trading_trades from Mongo (sorted by legacy_id) and writes them back
* into D1 via `wrangler d1 execute --remote --file=<tmp.sql>`.
*
* Preserves legacy_id when present (written by phase-05 forward backfill).
* When legacy_id is absent, generates sequential IDs from max(existing_d1_id)+1.
*
* Use this ONLY during a Stage-2 rollback (phase-07 debugger #14).
* Operator MUST inform users that N days of Mongo-only writes will revert.
*
* Flags:
* --dry-run Print SQL to stdout; do not execute wrangler.
* --force Proceed even if D1 trading_trades already has rows.
*
* Required env (loaded via --env-file-if-exists=.env.deploy):
* MONGODB_URI
*
* Usage:
* node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-d1.js
* node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-d1.js --dry-run
* node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-d1.js --force
*/
import { execSync } from "node:child_process";
import { rmSync, writeFileSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { closeMongoClient, getMongoClient } from "./lib/migration-helpers.js";
const DB_NAME = "miti99bot-db";
const COLLECTION = "trading_trades";
const dryRun = process.argv.includes("--dry-run");
const force = process.argv.includes("--force");
// ─── Preflight ────────────────────────────────────────────────────────────────
function validateEnv() {
const missingVars = [["MONGODB_URI", "Atlas connection string"]].filter(([k]) => !process.env[k]);
if (missingVars.length) {
for (const [k, desc] of missingVars)
console.error(`[backfill-mongo-d1] Missing: ${k} (${desc})`);
console.error(" Copy .env.deploy.example → .env.deploy and fill in values.");
process.exit(1);
}
}
// ─── D1 helpers ───────────────────────────────────────────────────────────────
/**
* Execute a single SQL command against remote D1 and return parsed rows.
*
* @param {string} sql
* @returns {any[]}
*/
function queryD1(sql) {
const cmd = `npx wrangler d1 execute ${DB_NAME} --remote --command "${sql.replace(/"/g, '\\"')}" --json`;
let stdout;
try {
stdout = execSync(cmd, { stdio: ["ignore", "pipe", "pipe"] }).toString();
} catch (err) {
const stderr = /** @type {any} */ (err).stderr?.toString() ?? "";
throw new Error(
`wrangler d1 execute failed:\n${stderr || /** @type {any} */ (err).stdout?.toString()}`,
);
}
const parsed = JSON.parse(stdout);
return Array.isArray(parsed) ? (parsed[0]?.results ?? []) : [];
}
/**
* Execute a SQL file against remote D1 via wrangler.
*
* @param {string} filePath — absolute path to .sql file
*/
function executeD1File(filePath) {
const cmd = `npx wrangler d1 execute ${DB_NAME} --remote --file="${filePath}" --json`;
try {
execSync(cmd, { stdio: ["ignore", "pipe", "pipe"] });
} catch (err) {
const stderr = /** @type {any} */ (err).stderr?.toString() ?? "";
throw new Error(
`wrangler d1 execute failed:\n${stderr || /** @type {any} */ (err).stdout?.toString()}`,
);
}
}
// ─── SQL generation ──────────────────────────────────────────────────────────
/**
* Escape a SQL string value (single-quote doubling).
*
* @param {string|null|undefined} v
* @returns {string} SQL literal (including surrounding quotes)
*/
export function sqlStr(v) {
if (v == null) return "NULL";
return `'${String(v).replace(/'/g, "''")}'`;
}
/**
* Build an INSERT statement for one trading trade row.
*
* @param {{
* id: number,
* user_id: string,
* symbol: string,
* side: string,
* qty: number,
* price_vnd: number,
* ts: number
* }} row
* @returns {string}
*/
export function buildInsertSql(row) {
return (
`INSERT INTO ${COLLECTION} (id, user_id, symbol, side, qty, price_vnd, ts) VALUES (` +
[
row.id,
sqlStr(row.user_id),
sqlStr(row.symbol),
sqlStr(row.side),
row.qty,
row.price_vnd,
row.ts,
].join(", ") +
");"
);
}
// ─── Main ─────────────────────────────────────────────────────────────────────
async function main() {
validateEnv();
if (dryRun) console.log("[backfill-mongo-d1] DRY RUN — SQL will be printed to stdout");
const client = await getMongoClient(/** @type {string} */ (process.env.MONGODB_URI));
const db = client.db();
const coll = db.collection(COLLECTION);
// Sort by legacy_id ascending so we maintain original D1 row order.
const docs = await coll.find({}).sort({ legacy_id: 1 }).toArray();
console.log(`[backfill-mongo-d1] Found ${docs.length} document(s) in Mongo ${COLLECTION}`);
if (docs.length === 0) {
console.log("[backfill-mongo-d1] Nothing to restore.");
await closeMongoClient();
return;
}
if (!dryRun) {
// Pre-flight: abort if D1 already has rows unless --force.
const existingRows = queryD1(`SELECT COUNT(*) AS cnt FROM ${COLLECTION}`);
const existingCount = existingRows[0]?.cnt ?? 0;
if (existingCount > 0 && !force) {
console.error(
`[backfill-mongo-d1] ABORT: D1 ${COLLECTION} already has ${existingCount} row(s).`,
);
console.error(" Run with --force to bypass (e.g. after wiping D1 manually).");
await closeMongoClient();
process.exit(1);
}
if (existingCount > 0 && force) {
console.log(
`[backfill-mongo-d1] --force: proceeding despite ${existingCount} existing D1 row(s).`,
);
}
}
// Determine starting ID for docs without a legacy_id.
let nextId = 1;
const hasLegacyIds = docs.some((d) => d.legacy_id != null);
if (!dryRun && !hasLegacyIds) {
// Fetch max id from D1 to avoid collisions.
const maxRows = queryD1(`SELECT MAX(id) AS m FROM ${COLLECTION}`);
const maxId = maxRows[0]?.m ?? 0;
nextId = maxId + 1;
}
// Build SQL statements.
/** @type {string[]} */
const statements = [];
for (const doc of docs) {
const id = doc.legacy_id != null ? Number(doc.legacy_id) : nextId++;
statements.push(
buildInsertSql({
id,
user_id: String(doc.user_id ?? ""),
symbol: String(doc.symbol ?? ""),
side: String(doc.side ?? ""),
qty: Number(doc.qty ?? 0),
price_vnd: Number(doc.price_vnd ?? 0),
ts: Number(doc.ts ?? 0),
}),
);
}
if (dryRun) {
console.log("[backfill-mongo-d1] Generated SQL:");
for (const stmt of statements) console.log(stmt);
console.log(`[backfill-mongo-d1] DRY RUN complete — ${statements.length} statement(s).`);
await closeMongoClient();
return;
}
// Write SQL to a temp file and pipe through wrangler.
const tmpFile = join(tmpdir(), `backfill-mongo-d1-${Date.now()}.sql`);
try {
writeFileSync(tmpFile, statements.join("\n"), "utf8");
console.log(`[backfill-mongo-d1] Executing ${statements.length} INSERT(s) via wrangler...`);
executeD1File(tmpFile);
console.log(`[backfill-mongo-d1] Done — ${statements.length} row(s) restored to D1.`);
} finally {
try {
rmSync(tmpFile);
} catch {
// Non-fatal cleanup.
}
}
await closeMongoClient();
}
// Run only when invoked directly (not when imported by tests).
const isMain = process.argv[1]?.endsWith("backfill-mongo-to-d1.js");
if (isMain) {
main().catch((err) => {
console.error("[backfill-mongo-d1] Fatal:", /** @type {Error} */ (err).message ?? err);
process.exit(1);
});
}
+201
View File
@@ -0,0 +1,201 @@
#!/usr/bin/env node
/**
* @file backfill-mongo-to-kv — emergency reverse-backfill: MongoDB → Cloudflare KV.
*
* Reads each per-module Mongo collection and writes every non-expired doc back
* into CF KV via the REST API with correct TTL derived from expiresAt.
*
* Use this ONLY during a Stage-2 rollback. Operator must inform users that
* N days of Mongo-only writes will revert (phase-07 debugger #14).
*
* Flags:
* --dry-run Log summary without writing to KV.
* --module <name> Restore only a single module (default: all).
*
* Required env (loaded via --env-file-if-exists=.env.deploy):
* MONGODB_URI, CLOUDFLARE_ACCOUNT_ID, CLOUDFLARE_API_TOKEN,
* KV_NAMESPACE_ID, MODULES (comma-separated)
*
* Usage:
* node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-kv.js
* node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-kv.js --dry-run
* node --env-file-if-exists=.env.deploy scripts/backfill-mongo-to-kv.js --module wordle
*/
import { cfKvPut, closeMongoClient, getMongoClient, sleep } from "./lib/migration-helpers.js";
// ─── Config ───────────────────────────────────────────────────────────────────
const {
MONGODB_URI,
CLOUDFLARE_ACCOUNT_ID,
CLOUDFLARE_API_TOKEN,
KV_NAMESPACE_ID,
MODULES: MODULES_ENV,
} = process.env;
const dryRun = process.argv.includes("--dry-run");
const moduleFlag = (() => {
const idx = process.argv.indexOf("--module");
return idx !== -1 ? process.argv[idx + 1] : null;
})();
/** Throttle: 50 ops/sec → 20 ms between writes. */
const THROTTLE_MS = 20;
/** Minimum TTL accepted by CF KV REST API. */
const MIN_TTL_SECS = 60;
/** Normalize module name to Mongo collection name (mirrors mongo-kv-store.js). */
const toCollName = (mod) => mod.replace(/-/g, "_");
// ─── Preflight ────────────────────────────────────────────────────────────────
function validateEnv() {
const missing = [
"MONGODB_URI",
"CLOUDFLARE_ACCOUNT_ID",
"CLOUDFLARE_API_TOKEN",
"KV_NAMESPACE_ID",
"MODULES",
].filter((k) => !process.env[k]);
if (missing.length) {
console.error(`[backfill-mongo-kv] Missing required env vars: ${missing.join(", ")}`);
console.error(" Copy .env.deploy.example → .env.deploy and fill in values.");
process.exit(1);
}
}
// ─── TTL computation ─────────────────────────────────────────────────────────
/**
* Compute CF KV expirationTtl (seconds from now) from an absolute expiresAt Date.
* Returns null if the doc has no TTL (key should be persistent).
* Returns undefined (sentinel) if the doc is already expired (key must be SKIPPED).
*
* @param {Date|null|undefined} expiresAt
* @param {number} nowMs — Date.now() at call time (injectable for tests)
* @returns {{ ttl: number }|null|"expired"}
*/
export function computeTtl(expiresAt, nowMs = Date.now()) {
if (!expiresAt) return null; // no TTL — write as persistent
const remainingMs = expiresAt.getTime() - nowMs;
if (remainingMs <= 0) return "expired"; // already past expiry → skip
const secs = Math.floor(remainingMs / 1000);
return { ttl: Math.max(MIN_TTL_SECS, secs) };
}
// ─── Per-module restore ───────────────────────────────────────────────────────
/**
* Restore one module's Mongo collection back into CF KV.
*
* @param {import("mongodb").Db} db
* @param {string} mod
* @returns {Promise<{restored: number, skipped: number, failed: number}>}
*/
async function restoreModule(db, mod) {
const coll = db.collection(toCollName(mod));
const docs = await coll.find({}).toArray();
let restored = 0;
let skipped = 0;
let failed = 0;
for (const doc of docs) {
const key = /** @type {string} */ (doc._id);
const value = /** @type {string} */ (doc.value ?? "");
const ttlResult = computeTtl(doc.expiresAt ?? null);
if (ttlResult === "expired") {
skipped++;
continue; // expired in Mongo → do not surface a stale key in KV
}
if (dryRun) {
restored++;
continue;
}
try {
/** @type {{ expirationTtl?: number }} */
const opts = ttlResult ? { expirationTtl: ttlResult.ttl } : {};
await cfKvPut(
/** @type {string} */ (CLOUDFLARE_ACCOUNT_ID),
/** @type {string} */ (KV_NAMESPACE_ID),
/** @type {string} */ (CLOUDFLARE_API_TOKEN),
key,
value,
opts,
);
restored++;
await sleep(THROTTLE_MS);
} catch (err) {
failed++;
// Log key hash only — never log plaintext keys (may encode user IDs).
const { sha256 } = await import("./lib/migration-helpers.js");
console.error(
`[${mod}] ERROR key_sha256=${sha256(String(key)).slice(0, 16)}: ${/** @type {Error} */ (err).message}`,
);
}
}
return { restored, skipped, failed };
}
// ─── Main ─────────────────────────────────────────────────────────────────────
async function main() {
validateEnv();
const allModules = /** @type {string} */ (MODULES_ENV)
.split(",")
.map((m) => m.trim())
.filter(Boolean);
const modules = moduleFlag ? [moduleFlag] : allModules;
if (moduleFlag && !allModules.includes(moduleFlag)) {
console.error(
`[backfill-mongo-kv] Unknown module "${moduleFlag}". Available: ${allModules.join(", ")}`,
);
process.exit(1);
}
if (dryRun) console.log("[backfill-mongo-kv] DRY RUN — no writes to KV");
console.log(`[backfill-mongo-kv] Modules: ${modules.join(", ")}`);
const client = await getMongoClient(/** @type {string} */ (MONGODB_URI));
const db = client.db();
let totalFailed = 0;
for (const mod of modules) {
const { restored, skipped, failed } = await restoreModule(db, mod);
const verb = dryRun ? "(dry-run)" : `${restored} restored, ${skipped} skipped (expired)`;
console.log(`[${mod}] ${docs_label(restored + skipped + failed)} docs: ${verb}`);
totalFailed += failed;
}
await closeMongoClient();
if (totalFailed > 0) {
console.error(
`[backfill-mongo-kv] Completed with ${totalFailed} failed write(s). Check logs above.`,
);
process.exit(1);
}
console.log("[backfill-mongo-kv] All modules restored.");
}
/** @param {number} n @returns {string} */
function docs_label(n) {
return `${n} doc${n !== 1 ? "s" : ""}`;
}
// Run only when invoked directly (not when imported by tests).
const isMain = process.argv[1]?.endsWith("backfill-mongo-to-kv.js");
if (isMain) {
main().catch((err) => {
console.error("[backfill-mongo-kv] Fatal:", /** @type {Error} */ (err).message ?? err);
process.exit(1);
});
}
+35
View File
@@ -128,6 +128,41 @@ export async function cfKvGet(accountId, nsId, token, key) {
return res.text();
}
/**
* Write a string value into CF KV via the REST API.
*
* CF KV REST PUT: PUT /accounts/{id}/storage/kv/namespaces/{nsid}/values/{key}
* Optional query param `expiration_ttl` (seconds from now, minimum 60).
*
* @param {string} accountId
* @param {string} nsId
* @param {string} token
* @param {string} key
* @param {string} value
* @param {{ expirationTtl?: number }} [opts]
* @returns {Promise<void>}
*/
export async function cfKvPut(accountId, nsId, token, key, value, opts = {}) {
const url = new URL(
`${CF_API_BASE}/accounts/${accountId}/storage/kv/namespaces/${nsId}/values/${encodeURIComponent(key)}`,
);
if (opts.expirationTtl != null) {
url.searchParams.set("expiration_ttl", String(opts.expirationTtl));
}
const res = await fetch(url.toString(), {
method: "PUT",
headers: {
Authorization: `Bearer ${token}`,
"Content-Type": "text/plain",
},
body: value,
});
if (!res.ok) {
const body = await res.text().catch(() => "");
throw new Error(`CF KV put "${key}" ${res.status}: ${body}`);
}
}
// ─── MongoDB singleton ────────────────────────────────────────────────────────
/** @type {MongoClient|null} */
+59
View File
@@ -0,0 +1,59 @@
#!/usr/bin/env bash
# wrangler-delete-guard.sh — interactive CONFIRM wrapper around irreversible
# wrangler delete commands. NEVER use in CI.
#
# Usage:
# bash scripts/wrangler-delete-guard.sh kv <namespace-id>
# bash scripts/wrangler-delete-guard.sh d1 <database-name>
#
# Exits:
# 0 — command executed successfully
# 1 — user aborted (did not type CONFIRM)
# 2 — bad arguments
# 3 — stdin not a tty (CI safety)
set -euo pipefail
if [[ $# -lt 2 ]]; then
echo "usage: $0 {kv|d1} <id-or-name>" >&2
exit 2
fi
KIND="$1"
TARGET="$2"
case "$KIND" in
kv)
DESC="KV namespace $TARGET"
CMD=(npx wrangler kv namespace delete --namespace-id "$TARGET")
;;
d1)
DESC="D1 database $TARGET"
CMD=(npx wrangler d1 delete "$TARGET")
;;
*)
echo "unknown kind: $KIND (expected 'kv' or 'd1')" >&2
exit 2
;;
esac
# Refuse to run non-interactively — prevents accidental CI execution.
if [[ ! -t 0 ]]; then
echo "stdin not a tty — refusing to run non-interactively (CI safety)" >&2
exit 3
fi
echo ""
echo "ABOUT TO DELETE: $DESC"
echo "This is IRREVERSIBLE. Backup files must already be on local disk."
echo ""
read -r -p "Type CONFIRM to proceed: " CONFIRM
if [[ "$CONFIRM" != "CONFIRM" ]]; then
echo "aborted" >&2
exit 1
fi
echo ""
echo "executing: ${CMD[*]}"
"${CMD[@]}"
+122
View File
@@ -0,0 +1,122 @@
/**
* @file backfill-mongo-to-d1.test.js — unit tests for SQL-statement construction.
*
* Tests `buildInsertSql` and `sqlStr` exported from backfill-mongo-to-d1.js.
* No execSync, no wrangler, no real Mongo connection.
*/
import { describe, expect, it } from "vitest";
import { buildInsertSql, sqlStr } from "../../scripts/backfill-mongo-to-d1.js";
// ─── sqlStr ───────────────────────────────────────────────────────────────────
describe("sqlStr", () => {
it("wraps a plain string in single quotes", () => {
expect(sqlStr("hello")).toBe("'hello'");
});
it("escapes internal single quotes by doubling them", () => {
expect(sqlStr("it's")).toBe("'it''s'");
});
it("escapes multiple single quotes", () => {
expect(sqlStr("a'b'c")).toBe("'a''b''c'");
});
it("returns NULL for null", () => {
expect(sqlStr(null)).toBe("NULL");
});
it("returns NULL for undefined", () => {
expect(sqlStr(undefined)).toBe("NULL");
});
it("handles empty string", () => {
expect(sqlStr("")).toBe("''");
});
it("coerces non-string values via String()", () => {
// @ts-ignore — testing runtime coercion
expect(sqlStr(42)).toBe("'42'");
});
});
// ─── buildInsertSql ───────────────────────────────────────────────────────────
describe("buildInsertSql", () => {
/** @type {Parameters<typeof buildInsertSql>[0]} */
const BASE_ROW = {
id: 1,
user_id: "u123",
symbol: "BTC",
side: "buy",
qty: 0.5,
price_vnd: 1500000,
ts: 1700000000,
};
it("produces a valid INSERT statement", () => {
const sql = buildInsertSql(BASE_ROW);
expect(sql).toMatch(
/^INSERT INTO trading_trades \(id, user_id, symbol, side, qty, price_vnd, ts\) VALUES \(/,
);
expect(sql).toMatch(/\);$/);
});
it("includes all column values in correct order", () => {
const sql = buildInsertSql(BASE_ROW);
expect(sql).toContain("1,"); // id
expect(sql).toContain("'u123'"); // user_id
expect(sql).toContain("'BTC'"); // symbol
expect(sql).toContain("'buy'"); // side
expect(sql).toContain("0.5,"); // qty
expect(sql).toContain("1500000,"); // price_vnd
expect(sql).toContain("1700000000"); // ts
});
it("preserves legacy_id as the INSERT id", () => {
const sql = buildInsertSql({ ...BASE_ROW, id: 42 });
// id=42 should be the first value
expect(sql).toMatch(/VALUES \(42,/);
});
it("escapes single quotes in user_id", () => {
const sql = buildInsertSql({ ...BASE_ROW, user_id: "o'brien" });
expect(sql).toContain("'o''brien'");
});
it("escapes single quotes in symbol", () => {
const sql = buildInsertSql({ ...BASE_ROW, symbol: "it's" });
expect(sql).toContain("'it''s'");
});
it("handles decimal qty correctly", () => {
const sql = buildInsertSql({ ...BASE_ROW, qty: 1.23456789 });
expect(sql).toContain("1.23456789");
});
it("handles zero values without coercion errors", () => {
const sql = buildInsertSql({
id: 0,
user_id: "",
symbol: "",
side: "",
qty: 0,
price_vnd: 0,
ts: 0,
});
expect(sql).toMatch(/VALUES \(0, '', '', '', 0, 0, 0\);/);
});
it("sequential id generation: each row gets unique ascending id", () => {
// Simulate the script logic of incrementing nextId for docs without legacy_id
let nextId = 1;
const rows = [
{ user_id: "a", symbol: "BTC", side: "buy", qty: 1, price_vnd: 100, ts: 1 },
{ user_id: "b", symbol: "ETH", side: "sell", qty: 2, price_vnd: 200, ts: 2 },
];
const stmts = rows.map((r) => buildInsertSql({ ...r, id: nextId++ }));
expect(stmts[0]).toMatch(/VALUES \(1,/);
expect(stmts[1]).toMatch(/VALUES \(2,/);
});
});
+118
View File
@@ -0,0 +1,118 @@
/**
* @file backfill-mongo-to-kv.test.js — unit tests for reverse-backfill TTL math.
*
* Tests the `computeTtl` helper exported from backfill-mongo-to-kv.js.
* No real Atlas connection, no real CF REST calls.
*/
import { afterEach, describe, expect, it, vi } from "vitest";
import { computeTtl } from "../../scripts/backfill-mongo-to-kv.js";
// ─── computeTtl ───────────────────────────────────────────────────────────────
describe("computeTtl", () => {
const NOW = 1_700_000_000_000; // fixed epoch ms for deterministic tests
it("returns null when expiresAt is null (persistent key)", () => {
expect(computeTtl(null, NOW)).toBeNull();
});
it("returns null when expiresAt is undefined (persistent key)", () => {
expect(computeTtl(undefined, NOW)).toBeNull();
});
it("returns 'expired' when expiresAt is in the past", () => {
const past = new Date(NOW - 1000); // 1 second ago
expect(computeTtl(past, NOW)).toBe("expired");
});
it("returns 'expired' when expiresAt equals now exactly", () => {
const now = new Date(NOW);
expect(computeTtl(now, NOW)).toBe("expired");
});
it("returns correct ttl (seconds) for a future expiresAt", () => {
const future = new Date(NOW + 120_000); // 120 seconds from now
const result = computeTtl(future, NOW);
expect(result).not.toBeNull();
expect(result).not.toBe("expired");
expect(/** @type {{ttl: number}} */ (result).ttl).toBe(120);
});
it("floors partial seconds (no rounding up)", () => {
// 90.9 seconds remaining → floor → 90
const future = new Date(NOW + 90_900);
const result = computeTtl(future, NOW);
expect(/** @type {{ttl: number}} */ (result).ttl).toBe(90);
});
it("clamps ttl to MIN_TTL_SECS (60) when remaining < 60s", () => {
// 30 seconds remaining — below CF KV minimum of 60
const future = new Date(NOW + 30_000);
const result = computeTtl(future, NOW);
expect(/** @type {{ttl: number}} */ (result).ttl).toBe(60);
});
it("clamps ttl to MIN_TTL_SECS (60) when remaining is 1s", () => {
const future = new Date(NOW + 1_000);
const result = computeTtl(future, NOW);
expect(/** @type {{ttl: number}} */ (result).ttl).toBe(60);
});
it("does NOT clamp when remaining is exactly 60s", () => {
const future = new Date(NOW + 60_000);
const result = computeTtl(future, NOW);
expect(/** @type {{ttl: number}} */ (result).ttl).toBe(60);
});
it("does NOT clamp for large TTL values", () => {
// 7 days = 604800s
const future = new Date(NOW + 7 * 24 * 3600 * 1000);
const result = computeTtl(future, NOW);
expect(/** @type {{ttl: number}} */ (result).ttl).toBe(604800);
});
});
// ─── cfKvPut integration (fetch mock) ────────────────────────────────────────
// Verify that computeTtl=null produces a PUT with no expiration_ttl query param,
// and computeTtl={ttl:N} appends the param.
describe("cfKvPut TTL query param (fetch mock)", () => {
afterEach(() => vi.restoreAllMocks());
async function mockPut(expirationTtl) {
const fetchMock = vi.fn().mockResolvedValue({
ok: true,
status: 200,
text: () => Promise.resolve(""),
});
vi.stubGlobal("fetch", fetchMock);
const { cfKvPut } = await import("../../scripts/lib/migration-helpers.js");
const opts = expirationTtl != null ? { expirationTtl } : {};
await cfKvPut("acct", "ns", "tok", "mod:key", "value", opts);
return fetchMock.mock.calls[0][0]; // the URL string
}
it("omits expiration_ttl param when no TTL provided", async () => {
const url = await mockPut(undefined);
expect(url).not.toContain("expiration_ttl");
});
it("appends expiration_ttl param when TTL provided", async () => {
const url = await mockPut(300);
expect(url).toContain("expiration_ttl=300");
});
it("URL-encodes the key in the PUT request URL", async () => {
const fetchMock = vi.fn().mockResolvedValue({
ok: true,
text: () => Promise.resolve(""),
});
vi.stubGlobal("fetch", fetchMock);
const { cfKvPut } = await import("../../scripts/lib/migration-helpers.js");
await cfKvPut("acct", "ns", "tok", "mod:key with spaces", "v", {});
const url = fetchMock.mock.calls[0][0];
expect(url).toContain("mod%3Akey%20with%20spaces");
});
});
+80
View File
@@ -13,6 +13,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import {
cfKvGet,
cfKvList,
cfKvPut,
clearCheckpoint,
closeMongoClient,
countDiffRatio,
@@ -257,6 +258,85 @@ describe("cfKvGet (fetch mock)", () => {
});
});
// ─── cfKvPut ──────────────────────────────────────────────────────────────────
describe("cfKvPut (fetch mock)", () => {
afterEach(() => vi.restoreAllMocks());
function mockFetch(status = 200) {
vi.stubGlobal(
"fetch",
vi.fn().mockResolvedValue({
ok: status >= 200 && status < 300,
status,
text: () => Promise.resolve(""),
}),
);
}
it("issues a PUT request to the correct URL", async () => {
const fetchMock = vi.fn().mockResolvedValue({ ok: true, text: () => Promise.resolve("") });
vi.stubGlobal("fetch", fetchMock);
await cfKvPut("acct", "ns", "tok", "my:key", "myvalue");
const [url, init] = fetchMock.mock.calls[0];
expect(url).toContain("/accounts/acct/storage/kv/namespaces/ns/values/my%3Akey");
expect(init.method).toBe("PUT");
});
it("sends Authorization Bearer header", async () => {
const fetchMock = vi.fn().mockResolvedValue({ ok: true, text: () => Promise.resolve("") });
vi.stubGlobal("fetch", fetchMock);
await cfKvPut("acct", "ns", "TOKEN", "k", "v");
const init = fetchMock.mock.calls[0][1];
expect(init.headers.Authorization).toBe("Bearer TOKEN");
});
it("sends the value as the request body", async () => {
const fetchMock = vi.fn().mockResolvedValue({ ok: true, text: () => Promise.resolve("") });
vi.stubGlobal("fetch", fetchMock);
await cfKvPut("acct", "ns", "tok", "k", "hello-value");
const init = fetchMock.mock.calls[0][1];
expect(init.body).toBe("hello-value");
});
it("appends expiration_ttl query param when provided", async () => {
const fetchMock = vi.fn().mockResolvedValue({ ok: true, text: () => Promise.resolve("") });
vi.stubGlobal("fetch", fetchMock);
await cfKvPut("acct", "ns", "tok", "k", "v", { expirationTtl: 300 });
const url = fetchMock.mock.calls[0][0];
expect(url).toContain("expiration_ttl=300");
});
it("omits expiration_ttl query param when opts is empty", async () => {
const fetchMock = vi.fn().mockResolvedValue({ ok: true, text: () => Promise.resolve("") });
vi.stubGlobal("fetch", fetchMock);
await cfKvPut("acct", "ns", "tok", "k", "v", {});
const url = fetchMock.mock.calls[0][0];
expect(url).not.toContain("expiration_ttl");
});
it("omits expiration_ttl when opts is omitted entirely", async () => {
const fetchMock = vi.fn().mockResolvedValue({ ok: true, text: () => Promise.resolve("") });
vi.stubGlobal("fetch", fetchMock);
await cfKvPut("acct", "ns", "tok", "k", "v");
const url = fetchMock.mock.calls[0][0];
expect(url).not.toContain("expiration_ttl");
});
it("throws on HTTP error response", async () => {
mockFetch(403);
await expect(cfKvPut("acct", "ns", "tok", "k", "v")).rejects.toThrow("403");
});
it("URL-encodes special characters in key", async () => {
const fetchMock = vi.fn().mockResolvedValue({ ok: true, text: () => Promise.resolve("") });
vi.stubGlobal("fetch", fetchMock);
await cfKvPut("acct", "ns", "tok", "key with spaces/slash", "v");
const url = fetchMock.mock.calls[0][0];
expect(url).toContain("key%20with%20spaces%2Fslash");
});
});
// ─── getMongoClient teardown ──────────────────────────────────────────────────
// Ensure any singleton client opened during tests is closed (prevents open handles).