Compare commits

..
10 Commits
Author SHA1 Message Date
gitea-actions 7981c715ce chore(release): v1.0.15
Build and Push Images / Build jorgecuadros-web (push) Successful in 1m49s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m6s
Deploy on tag / Deploy to galactus (push) Successful in 23s
Cut by rmancinas via the "Cut release" workflow. Pushing the tag triggers build.yml; deploy separately with tag=1.0.15.
2026-08-05 07:37:01 +00:00
rmancinasandClaude Opus 5 d173c9e9a0 fix(billing): stop double-counting history a BALANCE FORWARD already carries
Build and Push Images / Build jorgecuadros-web (push) Successful in 1m52s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m4s
BALANCE FORWARD rows are not movements. Access materialized one per
customer per year, dated Jan 1, holding the closing balance of everything
before it — that is what let the portal keep each year in its own table and
still show a correct running balance from one year's rows.

The platform imported those rows AND the real pre-cutover history they
summarize, and every balance aggregate summed the lot. NUMid 501 read
-10,469.29 on the receivables worklist against -14,065.29 on the customer's
own statement and on the legacy portal; the gap was two cash receipts from
2009 and 2012 that the 2026 opening balance had already absorbed.

The scale settles what it is: summed the old way the whole book came to
+20,605,447.86 MXN — the office owing its customers 20.6 million pesos.
Floored, it is -56,855.90. A receivables ledger cannot be 20M in credit.

Adds BALANCE_FLOOR_JOIN + NOT_SUPERSEDED and applies them to balances()
(page and count queries, which must agree), to stats()'s per-currency and
per-domain figures, and to the owing/in-credit split. The four stats()
aggregates moved from Prisma groupBy to raw SQL because groupBy cannot
express a per-customer floor.

statement() takes the same floor as a scalar, which is also what stops
FEE ANUAL and fee15 leaking in. Those are not in
STATEMENT_EXCLUDED_SOURCE_TABLES — that list reproduces legacy's
DATOS2-only `datosfreak` — and they were putting 2,092 pre-cutover fee rows
across 1,062 customers into the statement, skewing it by -5,129,764 against
the number those customers have been quoted for years. Dating rather than
source is the right test: a fee row *after* the opening balance is a real
charge and still counts.

movements() is deliberately left alone. It is a browser over captured rows
— "how much water did we capture in April" — and staff need the historical
rows visible, so it keeps totalling everything, the same asymmetry
NOT_OUTSTANDING already has.

stats() now separates the two questions it was mixing: movements,
ledgerCustomers, crossLineCustomers and the date range stay unfloored
inventory; everything under byCurrency/byDomain is a balance and is floored.

BillingService had no tests. Adds 13 covering the floor's failure modes —
it fails silently, so MIN-vs-MAX, `>` vs `>=`, the NULL branch for customers
with no opening balance, and the join/predicate alias pairing are each
pinned, plus the 501 arithmetic as a regression.

Verified through the real service against the live ledger: balances() and
statement() both return -14,065.29 for 501, matching the portal.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-05 00:34:53 -07:00
rmancinasandClaude Opus 5 e9a5ee9e90 fix(migration): carry NOPAGO into transactions.outstanding
datosfreak's NOPAGO is the legacy "still owed" flag, and the website reads
it directly — account.statement.php splits the statement on NOPAGO = 0 vs
NOPAGO = 1 and renders the latter as "Outstanding Bills Requiring
Attention". transform_transactions.py hardcoded 0, so all 40,421 rows came
across settled and that section renders empty for anyone served off the
platform. Not a missing column: a missing section, with no error.

Only the three DATOS2-shaped tables carry the flag (76 rows set in datos2,
0 in FEE ANUAL and fee15); the EFECTIVO/FM3 cash streams have no such
column and keep the 0 default. Sync mode gets outstanding=VALUES(...) too,
so an additive sync corrects rows already loaded rather than leaving them
settled forever.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-04 23:53:45 -07:00
gitea-actions 458e67340c chore(release): v1.0.14
Build and Push Images / Build jorgecuadros-web (push) Successful in 1m46s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m45s
Deploy on tag / Deploy to galactus (push) Successful in 1m5s
Cut by rmancinas via the "Cut release" workflow. Pushing the tag triggers build.yml; deploy separately with tag=1.0.14.
2026-08-05 05:33:49 +00:00
rmancinasandClaude Opus 5 b12382b436 ci: move the galactus deploy chain to a tag-only workflow
Build and Push Images / Build jorgecuadros-web (push) Successful in 1m39s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m17s
The deploy was a job inside build.yml gated by
`if: startsWith(github.ref, 'refs/tags/v')`. Gitea draws every job into the
run graph before it evaluates that `if`, so an ordinary push to master showed
a pending "Deploy to galactus" — indistinguishable from prod being about to be
redeployed off an unreleased commit, and the only safe reaction is to cancel
the run, which takes the images down with it.

The gate itself was never wrong (no deploy-galactus run has ever been created
from a branch ref), but a guarantee you cannot see is not much of a guarantee.
`on: push: tags: ["v*"]` in a workflow of its own makes it structural: the
deploy cannot appear on a master build because the workflow does not exist
there.

It replaces `needs: build` by polling the Actions API for the build.yml run at
this tag and requiring it green, so both images are still known to be in the
registry before anything is pulled. AUTO_DEPLOY_GALACTUS still cuts the chain.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-04 22:28:02 -07:00
gitea-actions 2620559975 chore(release): v1.0.13
Build and Push Images / Build jorgecuadros-web (push) Successful in 2m1s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m15s
Build and Push Images / Deploy to galactus (push) Successful in 8s
Cut by rmancinas via the "Cut release" workflow. Pushing the tag triggers build.yml; deploy separately with tag=1.0.13.
2026-08-05 05:23:57 +00:00
rmancinasandClaude Opus 5 ed19f51a52 fix(ops): re-stage before an additive sync
Build and Push Images / Deploy to galactus (push) Canceled after 0s
Build and Push Images / Build jorgecuadros-web (push) Canceled after 1m26s
Build and Push Images / Build jorgecuadros-api (push) Canceled after 1m27s
SYNC ran `run_all.py --sync` without `--stage`, so it depended on staged
Parquet under migration/output. That directory is part of the image, not a
volume, so any redeploy wiped it and the job died on the first transform:

    FileNotFoundError: '/repo/migration/output/stg_utilities/datgral.parquet'

Re-staging is also what makes the job's own label true — without it a sync
would replay whatever upload staged last, not the files currently in the
ingest folder.

Staging now counts as a numbered step when it runs, so the Operaciones
progress bar moves during the slowest phase instead of sitting empty.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-04 22:20:30 -07:00
rmancinasandClaude Opus 5 e85db73dbc ci: deploy to galactus automatically when a tag build goes green
Build and Push Images / Build jorgecuadros-web (push) Successful in 1m55s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m23s
Build and Push Images / Deploy to galactus (push) Skipped
Cutting a release then had one manual step left: watch build.yml and
dispatch "Deploy to galactus" by hand with the version. Chain it.

build.yml gains a `deploy` job, `needs: build` and gated on
refs/tags/v*, that dispatches deploy-galactus.yml against the tag with
tag=<version> scope=app bootstrap=false skip_migrate=false. `needs`
waits for both matrix legs, so api and web are both in the registry
before prod pulls either — deploy-galactus.yml only pulls, and a
half-pushed pair leaves prod running one new image and one old one.

A dispatch rather than a `workflow_run:` trigger (which Gitea has
supported since 1.24) because deploy-galactus.yml reads
github.event.inputs.* in ten places; under workflow_run all of them are
empty strings, so the deploy would run with no tag. The dispatch keeps
that workflow's contract intact and keeps it hand-runnable, which is
how rollbacks work.

The dispatch is confirmed the same way release.yml confirms the build
started: snapshot the existing deploy-galactus run ids first, then
require a new one to appear. An accepted dispatch that creates no run
is the failure mode that cost v1.0.3 its images, and a plain "is there
a deploy run" check would be satisfied by the previous release.

Kill switch: repo variable AUTO_DEPLOY_GALACTUS=false prints the manual
command instead of deploying. Needs the existing RELEASE_TOKEN secret;
preflight fails loudly and names the manual command if it is unset.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-04 21:48:23 -07:00
gitea-actions d38bbc52ec chore(release): v1.0.12
Build and Push Images / Build jorgecuadros-web (push) Successful in 2m30s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m48s
Cut by rmancinas via the "Cut release" workflow. Pushing the tag triggers build.yml; deploy separately with tag=1.0.12.
2026-08-05 04:40:07 +00:00
rmancinasandClaude Opus 5 fe761e119e feat(ops): show relay apply progress on the replication card
Build and Push Images / Build jorgecuadros-web (push) Successful in 2m0s
Build and Push Images / Build jorgecuadros-api (push) Successful in 2m15s
Seconds_Behind_Source cannot answer "is it moving?". While the SQL thread
works through one large transaction the lag counter holds still — often at
0 — even though the replica is not caught up. The relay backlog does move,
and it comes out of the SHOW REPLICA STATUS the panel already runs, so this
costs no extra query and no connection to the source.

Adds applyProgress(), which reads Source_Log_File / Read_Source_Log_Pos vs
Relay_Source_Log_File / Exec_Source_Log_Pos and reports the fetched-but-not-
applied byte delta plus a percentage. Both positions are source binlog
coordinates, so they are only comparable while the two threads are on the
same file; across files the delta is meaningless (positions restart at ~4 in
each new file) and is reported as null rather than as a huge negative number.

The percentage deliberately stops at 99.99 while any backlog remains —
binlog positions are large enough that a real backlog of a few KB rounds to
100% and would render a lagging replica as caught up.

Not folded into `healthy`: a non-zero backlog is the normal state of a
working replica between fetch and apply, so alarming on it would cry wolf.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-04 21:35:39 -07:00
16 changed files with 865 additions and 75 deletions
+10 -1
View File
@@ -14,7 +14,16 @@
# ARG/ENV (APP_VERSION / GIT_SHA / BUILD_DATE) and as OCI labels, so a running # ARG/ENV (APP_VERSION / GIT_SHA / BUILD_DATE) and as OCI labels, so a running
# container can report exactly what is deployed. # container can report exactly what is deployed.
# #
# Release flow: git tag v1.2.0 && git push origin v1.2.0 -> versioned images. # Release flow: git tag v1.2.0 && git push origin v1.2.0 -> versioned images
# -> deploy-on-tag.yml waits for this run to go green and then
# dispatches deploy-galactus.yml.
#
# This workflow BUILDS ONLY — it never deploys. The deploy chain used to be a
# job here, gated to tag refs, but Gitea draws every job of a workflow into the
# run graph before it evaluates the job's `if`: a routine master build showed a
# pending "Deploy to galactus" and looked like prod was about to be redeployed
# off an unreleased commit. Keeping the deploy in a `on: push: tags` workflow of
# its own makes that structurally impossible.
name: Build and Push Images name: Build and Push Images
+199
View File
@@ -0,0 +1,199 @@
# Chain the PROD deploy onto a green tag build.
#
# This is a SEPARATE workflow, not a job in build.yml, and the trigger is the
# whole point: `on: push: tags` cannot fire on a push to master. When this was a
# `deploy` job inside build.yml gated by `if: startsWith(github.ref,
# 'refs/tags/v')`, Gitea still drew "Deploy to galactus" into the job graph of
# every ordinary master build — the `if` is not evaluated until `needs` resolve,
# so the job sits there looking like an imminent production deploy on a commit
# nobody released. That is indistinguishable from a real misfire, and the only
# safe reaction is to cancel the run, which kills the images with it.
#
# What it does NOT do is build. build.yml already builds and pushes both images
# from one run; this waits for that run to go green and then dispatches
# deploy-galactus.yml, which only pulls.
#
# Why wait for the build run rather than just dispatching: deploy-galactus.yml
# pulls api and web at the same tag, and a half-pushed pair is exactly the state
# that leaves prod running one new image and one old one. The build run turning
# green is the signal that both are in the registry.
#
# Why a dispatch and not a `workflow_run:` trigger, which Gitea does support as
# of 1.24: deploy-galactus.yml reads `github.event.inputs.*` in ten places (tag,
# scope, bootstrap, skip_migrate). Under workflow_run every one of them is the
# empty string, so the deploy would silently run with no tag and scope != 'full'.
# A dispatch keeps that workflow's contract intact and keeps it hand-runnable for
# rollbacks, which is the whole point of it.
#
# Kill switch: set the repo variable AUTO_DEPLOY_GALACTUS to `false` to cut the
# chain and go back to dispatching the deploy by hand. Anything else (including
# unset) deploys.
name: Deploy on tag
on:
push:
tags: ["v*"]
jobs:
deploy:
name: Deploy to galactus
runs-on: docker
container:
image: node:20-alpine
steps:
- name: Preflight — RELEASE_TOKEN
env:
RELEASE_TOKEN: ${{ secrets.RELEASE_TOKEN }}
run: |
set -eu
if [ -z "${RELEASE_TOKEN:-}" ]; then
echo "::error::Secret RELEASE_TOKEN is not set, so this cannot wait"
echo "::error::for the build or dispatch the deploy. Once build.yml"
V=${GITHUB_REF#refs/tags/}
echo "::error::is green, run 'Deploy to galactus' by hand with tag=${V#v}."
exit 1
fi
- name: Wait for the tag build, then dispatch deploy-galactus.yml
env:
RELEASE_TOKEN: ${{ secrets.RELEASE_TOKEN }}
AUTO_DEPLOY: ${{ vars.AUTO_DEPLOY_GALACTUS }}
TAG_REF: ${{ github.ref }}
BUILD_SHA: ${{ github.sha }}
run: |
node -e '
const base = `${process.env.GITHUB_SERVER_URL}/api/v1/repos/${process.env.GITHUB_REPOSITORY}`;
const headers = { Authorization: `token ${process.env.RELEASE_TOKEN}` };
// refs/tags/v1.2.3 — derived from github.ref rather than ref_name so
// it does not depend on how Gitea populates GITHUB_REF_NAME.
const tagRef = process.env.TAG_REF;
const tag = tagRef.replace(/^refs\/tags\//, "");
// The git tag carries the leading v; the image tag does not.
const version = tag.replace(/^v/, "");
const sha = process.env.BUILD_SHA;
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
const runs = async () => {
const r = await fetch(`${base}/actions/runs?limit=50`, { headers });
if (!r.ok) throw new Error(`runs query failed: HTTP ${r.status}`);
return (await r.json()).workflow_runs || [];
};
// The release commit and its tag are the SAME sha, and build.yml
// skips the master run by design — so a sha match alone can latch
// onto that skipped run and call the build green when no image was
// ever pushed. Require the tag ref when the API reports one.
const isTagRun = (r) => {
const ref = r.head_branch || r.ref || "";
return !ref || ref === tag || ref === tagRef;
};
const buildRun = async () =>
(await runs()).find(
(r) =>
r.head_sha === sha &&
String(r.path || "").includes("build.yml") &&
isTagRun(r),
);
// Gitea reports a run as `status` and mirrors it into `conclusion`;
// read whichever is populated rather than betting on one field.
const outcome = (r) => String(r.conclusion || r.status || "").toLowerCase();
const DONE = ["success", "failure", "cancelled", "canceled", "skipped"];
// Every deploy-galactus run id visible right now. A dispatch is only
// confirmed by an id that is NOT in here — a plain "is there a deploy
// run" check is satisfied by the PREVIOUS release run, and would
// report success for a dispatch that never took.
const deployRunIds = async () =>
new Set(
(await runs())
.filter((r) => String(r.path || "").includes("deploy-galactus.yml"))
.map((r) => r.id),
);
(async () => {
if (process.env.AUTO_DEPLOY === "false") {
console.log("AUTO_DEPLOY_GALACTUS=false — not deploying.");
console.log(`Deploy by hand with tag=${version} when ready.`);
return;
}
// ~20 min. A build is about 90s; the rest is queue time behind
// other runs on a single runner.
let run = null;
for (let i = 0; i < 80; i++) {
run = await buildRun();
if (run && DONE.includes(outcome(run))) break;
if (!run && i === 11) {
// Two minutes with no run at all. The post-receive hook drops
// runs silently when it errors (this cost v1.0.3 its images),
// so say so rather than timing out with no explanation.
console.log(`::warning::No build.yml run for ${tag} yet after 2 min.`);
console.log(`::warning::If the Gitea post-receive hook is broken, dispatch`);
console.log(`::warning::"Build and Push Images" by hand with ref=${tag}.`);
}
await sleep(15_000);
}
if (!run) {
console.log(`::error::No build.yml run for ${tag} (${sha}) after 20 min.`);
console.log(`::error::Dispatch "Build and Push Images" with ref=${tag} (the`);
console.log(`::error::tag, not master), then deploy by hand with tag=${version}.`);
process.exit(1);
}
const result = outcome(run);
if (result !== "success") {
console.log(`::error::build.yml for ${tag} ended as "${result}" — not deploying.`);
console.log(`::error::Fix the build, re-run it, then deploy by hand with tag=${version}.`);
process.exit(1);
}
console.log(`build.yml for ${tag} is green (run ${run.id}). Deploying ${version}.`);
const before = await deployRunIds();
// Dispatch against the TAG, not master: the deploy applies the
// compose files under deploy/galactus/ from whatever ref it runs
// on, and those must be the ones this release was cut with.
const res = await fetch(
`${base}/actions/workflows/deploy-galactus.yml/dispatches`,
{
method: "POST",
headers: { ...headers, "Content-Type": "application/json" },
body: JSON.stringify({
ref: tagRef,
inputs: {
tag: version,
scope: "app",
bootstrap: "false",
skip_migrate: "false",
},
}),
},
);
if (!res.ok) {
console.log(`::error::Dispatch returned HTTP ${res.status}: ${await res.text()}`);
console.log(`::error::Images for ${version} are published. Run`);
console.log(`::error::"Deploy to galactus" by hand with tag=${version}.`);
process.exit(1);
}
// A 204 only means Gitea accepted the request. Confirm a NEW run
// exists — an accepted call that creates no run is the failure mode
// that cost v1.0.3 its images.
for (let i = 0; i < 3; i++) {
await sleep(5_000);
const fresh = [...(await deployRunIds())].filter((id) => !before.has(id));
if (fresh.length) {
console.log(`Deploy of ${version} to galactus is running (run ${fresh[0]}).`);
return;
}
}
console.log(`::error::Dispatch was accepted but no deploy run appeared.`);
console.log(`::error::Run "Deploy to galactus" by hand with tag=${version}.`);
process.exit(1);
})();
'
+17 -7
View File
@@ -1,10 +1,16 @@
# Cut a release: stamp the version across every package.json, commit, tag, push. # Cut a release: stamp the version across every package.json, commit, tag, push.
# #
# This does NOT build and does NOT deploy. Pushing the `vX.Y.Z` tag is what # This does NOT build and does NOT deploy itself. Pushing the `vX.Y.Z` tag is
# triggers build.yml, which publishes `X.Y.Z`, `X.Y`, `sha-<short>` and `latest` # what triggers both build.yml, which publishes the `X.Y.Z`, `X.Y`,
# image tags. Deploying stays a separate, deliberate act: once the build is # `sha-<short>` and `latest` image tags, and deploy-on-tag.yml, which waits for
# green, dispatch deploy-galactus.yml with `tag=X.Y.Z` (no leading v — the tag # that build to go green and then dispatches deploy-galactus.yml with
# carries the `v`, the image tag does not). # `tag=X.Y.Z scope=app` (no leading v — the git tag carries the `v`, the image
# tag does not). A tag is the only ref that starts either chain; pushing to
# master builds images and stops there.
#
# So cutting a release DOES reach prod. To cut a version without deploying it,
# set the repo variable AUTO_DEPLOY_GALACTUS=false first; deploy-on-tag.yml then
# prints the manual command instead of running it.
# #
# Why a workflow instead of three local commands: the release commit is the one # Why a workflow instead of three local commands: the release commit is the one
# thing that must be identical every time, and cutting it from a laptop is how # thing that must be identical every time, and cutting it from a laptop is how
@@ -266,5 +272,9 @@ jobs:
echo "Released v${VERSION}." echo "Released v${VERSION}."
echo "" echo ""
echo "build.yml is now building git.mancinas.io/rmancinas/jorgecuadros-{api,web}:${VERSION}." echo "build.yml is now building git.mancinas.io/rmancinas/jorgecuadros-{api,web}:${VERSION}."
echo "When it is green, dispatch 'Deploy to galactus' with:" echo "deploy-on-tag.yml is watching that build; when it goes green it dispatches"
echo " tag=${VERSION} scope=app bootstrap=false skip_migrate=false" echo "'Deploy to galactus' with tag=${VERSION} scope=app bootstrap=false skip_migrate=false."
echo ""
echo "Watch that run. If it did not start (or AUTO_DEPLOY_GALACTUS=false),"
echo "dispatch 'Deploy to galactus' by hand with the same inputs."
echo "Rollback = re-dispatch it with an older tag."
+1 -1
View File
@@ -1,6 +1,6 @@
{ {
"name": "@jorgecuadros/api", "name": "@jorgecuadros/api",
"version": "1.0.11", "version": "1.0.15",
"private": true, "private": true,
"scripts": { "scripts": {
"build": "nest build", "build": "nest build",
+180
View File
@@ -0,0 +1,180 @@
import { Prisma } from "@jorgecuadros/database";
import {
BALANCE_FLOOR_JOIN,
BALANCE_FORWARD_TYPE,
BillingService,
NOT_SUPERSEDED,
} from "./billing.service";
/**
* The balance floor drops rows a later BALANCE FORWARD already accounts for.
*
* It is worth testing because it fails silently: nothing throws, the numbers are
* just wrong, and they were wrong for years — the whole book read +20.6M MXN in
* credit because every customer's pre-cutover history was counted twice, once
* inside their opening balance and once as itself.
*/
describe("balance floor", () => {
describe("SQL fragments", () => {
it("binds the type name rather than interpolating it", () => {
// A literal would be a second place to edit if the label ever changes,
// and this string reaches SQL from a module constant.
expect(BALANCE_FLOOR_JOIN.values).toEqual([BALANCE_FORWARD_TYPE]);
});
it("keys the floor to the row's own customer", () => {
// Without this the derived table cross-joins and every customer inherits
// the earliest BALANCE FORWARD in the book.
expect(BALANCE_FLOOR_JOIN.sql).toContain(
"bfloor ON bfloor.customerId = t.customerId",
);
});
it("takes the most recent opening balance, not the first", () => {
// A customer accumulates one BALANCE FORWARD per year. MIN would floor at
// the oldest and leave every intervening year double-counted.
expect(BALANCE_FLOOR_JOIN.sql).toContain("MAX(bf.transactionDate)");
expect(BALANCE_FLOOR_JOIN.sql).not.toContain("MIN(bf.transactionDate)");
});
it("ignores voided opening balances when locating the floor", () => {
expect(BALANCE_FLOOR_JOIN.sql).toContain("bf.voidedAt IS NULL");
});
it("is inclusive of the opening balance row itself", () => {
// `>` instead of `>=` would drop the carried balance and understate every
// customer by exactly that amount.
expect(NOT_SUPERSEDED.sql).toContain("t.transactionDate >= bfloor.floorDate");
expect(NOT_SUPERSEDED.sql).not.toMatch(/transactionDate\s*>\s*bfloor/);
});
it("leaves customers with no opening balance untouched", () => {
// NULL comparisons are never true, so without the explicit IS NULL branch
// a customer who has no BALANCE FORWARD row loses their entire ledger.
expect(NOT_SUPERSEDED.sql).toContain("bfloor.floorDate IS NULL");
});
it("only ever references the alias the join defines", () => {
// The predicate is useless without the join; pairing them wrongly is a
// runtime "unknown column", so keep the alias identical in both.
const aliases = NOT_SUPERSEDED.sql.match(/bfloor\.\w+/g) ?? [];
expect(aliases.length).toBeGreaterThan(0);
for (const ref of aliases) {
expect(BALANCE_FLOOR_JOIN.sql).toContain(ref.split(".")[1]);
}
});
});
describe("statement()", () => {
/**
* One customer means one floor date, so the statement uses a scalar lookup
* instead of the join. Asserting on the `where` Prisma is handed is the only
* way to see it without a database.
*/
function serviceWith(floor: Date | null) {
const findMany = jest.fn().mockResolvedValue([]);
const prisma = {
customer: {
findUnique: jest.fn().mockResolvedValue({
id: "c1",
name: "CUADROS, JORGE H.",
preferredCurrency: "USD",
_count: { properties: 0, policies: 0 },
}),
},
transaction: {
findFirst: jest
.fn()
.mockResolvedValue(floor ? { transactionDate: floor } : null),
findMany,
},
};
return {
service: new BillingService(prisma as never),
prisma,
findMany,
};
}
it("looks the floor up from the customer's newest opening balance", async () => {
const { service, prisma } = serviceWith(new Date("2026-01-01T00:00:00Z"));
await service.statement("c1");
expect(prisma.transaction.findFirst).toHaveBeenCalledWith(
expect.objectContaining({
where: {
customerId: "c1",
voidedAt: null,
type: { nameEn: BALANCE_FORWARD_TYPE },
},
orderBy: { transactionDate: "desc" },
select: { transactionDate: true },
}),
);
});
it("bounds the statement at the floor, inclusive", async () => {
const floor = new Date("2026-01-01T00:00:00Z");
const { service, findMany } = serviceWith(floor);
await service.statement("c1");
expect(findMany.mock.calls[0][0].where).toMatchObject({
customerId: "c1",
transactionDate: { gte: floor },
});
});
it("applies no date bound when the customer has no opening balance", async () => {
const { service, findMany } = serviceWith(null);
await service.statement("c1");
expect(findMany.mock.calls[0][0].where).not.toHaveProperty(
"transactionDate",
);
});
it("keeps the source-table exclusion alongside the floor", async () => {
// The two guards answer different questions — one reproduces legacy's
// DATOS2-only materialization, the other drops superseded history — and
// dropping either one changes the customer's balance.
const { service, findMany } = serviceWith(new Date("2026-01-01T00:00:00Z"));
await service.statement("c1");
const where = findMany.mock.calls[0][0].where;
expect(where.OR).toEqual([
{ legacySourceTable: null },
{ legacySourceTable: { notIn: expect.arrayContaining(["EFECTIVO"]) } },
]);
});
});
describe("regression: NUMid 501", () => {
/**
* The arithmetic that exposed the bug, pinned so it cannot silently return.
* Figures measured against the live ledger on 2026-08-05.
*/
const openingBalance = new Prisma.Decimal("-6732.29");
const activitySinceOpening = new Prisma.Decimal("-7333.00");
const preCutoverCashAlreadyInOpening = new Prisma.Decimal("3596.00");
it("matches the legacy portal once superseded rows are dropped", () => {
expect(openingBalance.plus(activitySinceOpening).toFixed(2)).toBe(
"-14065.29",
);
});
it("reproduces the wrong figure when they are not", () => {
expect(
openingBalance
.plus(activitySinceOpening)
.plus(preCutoverCashAlreadyInOpening)
.toFixed(2),
).toBe("-10469.29");
});
});
});
+159 -52
View File
@@ -104,24 +104,31 @@ interface BalanceRow {
nameMissing: number; nameMissing: number;
city: string | null; city: string | null;
state: string | null; state: string | null;
movements: bigint | number | string; movements: RawCount;
balanceMxn: Prisma.Decimal | null; balanceMxn: Prisma.Decimal | null;
balanceUsd: Prisma.Decimal | null; balanceUsd: Prisma.Decimal | null;
chargesMxn: Prisma.Decimal | null; chargesMxn: Prisma.Decimal | null;
creditsMxn: Prisma.Decimal | null; creditsMxn: Prisma.Decimal | null;
chargesUsd: Prisma.Decimal | null; chargesUsd: Prisma.Decimal | null;
creditsUsd: Prisma.Decimal | null; creditsUsd: Prisma.Decimal | null;
utilityMovements: bigint | number | string; utilityMovements: RawCount;
insuranceMovements: bigint | number | string; insuranceMovements: RawCount;
lastMovement: Date | null; lastMovement: Date | null;
} }
/** /**
* Raw-query counts come back in three shapes depending on the aggregate: * Every shape a raw-query count can arrive in. `COUNT(*)` is a bigint,
* `COUNT(*)` as bigint, `SUM(bool)` as a decimal *string*, and plain numbers. * `SUM(bool)` is a Prisma.Decimal, and plain numbers occur too — none of which
* Normalize all of them before they reach the client as JSON. * survive JSON serialization the way the client expects.
*/ */
function num(v: bigint | number | string | null | undefined): number { type RawCount = bigint | number | string | Prisma.Decimal;
/**
* Normalizes a raw-query count before it reaches the client as JSON. A bigint
* throws on JSON.stringify and a Decimal serializes to a *string*, so counts
* must not be passed through untouched.
*/
function num(v: RawCount | null | undefined): number {
if (v === null || v === undefined) return 0; if (v === null || v === undefined) return 0;
return typeof v === "number" ? v : Number(v); return typeof v === "number" ? v : Number(v);
} }
@@ -152,6 +159,56 @@ const NOT_VOIDED: Prisma.TransactionWhereInput = { voidedAt: null };
*/ */
const NOT_OUTSTANDING: Prisma.TransactionWhereInput = { outstanding: false }; const NOT_OUTSTANDING: Prisma.TransactionWhereInput = { outstanding: false };
/**
* The legacy type name for a carried-forward opening balance.
*
* These rows are not movements. Access materialized one per customer per year,
* dated Jan 1, holding the closing balance of everything before it — that is
* what let the portal keep each year in its own table (`datosfreak` = current,
* `2025`, `2024`, ...) and still show a correct running balance from a single
* year's rows.
*/
export const BALANCE_FORWARD_TYPE = "BALANCE FORWARD";
/**
* Per-customer date of the most recent BALANCE FORWARD row.
*
* Joined rather than correlated: one small derived table (1,170 rows) beats a
* subquery evaluated per ledger row.
*/
export const BALANCE_FLOOR_JOIN = Prisma.sql`
LEFT JOIN (
SELECT bf.customerId, MAX(bf.transactionDate) AS floorDate
FROM transactions bf
JOIN type_transactions bft ON bft.id = bf.typeId
WHERE bft.nameEn = ${BALANCE_FORWARD_TYPE} AND bf.voidedAt IS NULL
GROUP BY bf.customerId
) bfloor ON bfloor.customerId = t.customerId`;
/**
* Excludes rows a later BALANCE FORWARD already accounts for.
*
* WHY THIS EXISTS. The platform holds both the synthetic BALANCE FORWARD rows
* and the real pre-cutover history they summarize, so summing a customer's
* whole ledger counts that history twice — once inside the opening balance,
* once as itself. NUMid 501 read -10,469.29 on the worklist against -14,065.29
* on the customer's own statement and on the legacy portal, the gap being two
* cash receipts from 2009 and 2012 that the 2026 opening balance had already
* absorbed.
*
* The scale is what settles it: summed the old way the entire book came to
* +20,605,447.86 MXN — the office owing its customers 20.6 million pesos.
* Floored, it is -56,855.90, a modest net receivable. A receivables ledger
* cannot be 20M in credit.
*
* Applies to BALANCES ONLY, in the same spirit as NOT_OUTSTANDING: the movement
* browser still totals every captured row, because "how much water did we
* capture in April" is a question about what was recorded, not about what is
* owed. Customers with no BALANCE FORWARD row (the floor is NULL) are
* unaffected.
*/
export const NOT_SUPERSEDED = Prisma.sql`(bfloor.floorDate IS NULL OR t.transactionDate >= bfloor.floorDate)`;
/** /**
* Source tables excluded from the customer-facing statement. * Source tables excluded from the customer-facing statement.
* *
@@ -402,22 +459,24 @@ export class BillingService {
MAX(t.transactionDate) AS lastMovement MAX(t.transactionDate) AS lastMovement
FROM customers c FROM customers c
JOIN transactions t ON t.customerId = c.id JOIN transactions t ON t.customerId = c.id
WHERE t.voidedAt IS NULL AND t.outstanding = 0 ${nameFilter} ${txFilter} ${BALANCE_FLOOR_JOIN}
WHERE t.voidedAt IS NULL AND t.outstanding = 0 AND ${NOT_SUPERSEDED} ${nameFilter} ${txFilter}
GROUP BY c.id, c.name, c.nameSource, c.nameMissing, c.city, c.state GROUP BY c.id, c.name, c.nameSource, c.nameMissing, c.city, c.state
${having} ${having}
${orderBy} ${orderBy}
LIMIT ${pageSize} OFFSET ${(page - 1) * pageSize} LIMIT ${pageSize} OFFSET ${(page - 1) * pageSize}
`; `;
const counted = await this.prisma.$queryRaw<{ total: bigint | number | string }[]>` const counted = await this.prisma.$queryRaw<{ total: RawCount }[]>`
SELECT COUNT(*) AS total FROM ( SELECT COUNT(*) AS total FROM (
SELECT c.id SELECT c.id
FROM customers c FROM customers c
JOIN transactions t ON t.customerId = c.id JOIN transactions t ON t.customerId = c.id
${BALANCE_FLOOR_JOIN}
-- Must match the page query's filters exactly, or the total disagrees -- Must match the page query's filters exactly, or the total disagrees
-- with the rows. (The void exclusion was missing here before the -- with the rows. (The void exclusion was missing here before the
-- outstanding work; a voided-only customer inflated the count.) -- outstanding work; a voided-only customer inflated the count.)
WHERE t.voidedAt IS NULL AND t.outstanding = 0 ${nameFilter} ${txFilter} WHERE t.voidedAt IS NULL AND t.outstanding = 0 AND ${NOT_SUPERSEDED} ${nameFilter} ${txFilter}
GROUP BY c.id GROUP BY c.id
${having} ${having}
) x ) x
@@ -458,9 +517,18 @@ export class BillingService {
}; };
} }
/** Top-line figures for the billing page header. */ /**
* Top-line figures for the billing page header.
*
* Two different questions live here and they use different row sets.
* `movements`, `ledgerCustomers`, `crossLineCustomers` and the date range are
* INVENTORY — what is stored — and count everything not voided. Everything
* under `byCurrency` / `byDomain` is a BALANCE, so it applies NOT_SUPERSEDED
* and drops rows an opening balance already accounts for. The four aggregates
* moved from Prisma groupBy to raw SQL to express that join; groupBy cannot.
*/
async stats() { async stats() {
const [movements, ledgerCustomers, byCurrency, byDomain] = await Promise.all([ const [movements, ledgerCustomers] = await Promise.all([
this.prisma.transaction.count({ where: NOT_VOIDED }), this.prisma.transaction.count({ where: NOT_VOIDED }),
this.prisma.transaction this.prisma.transaction
.findMany({ .findMany({
@@ -469,34 +537,47 @@ export class BillingService {
select: { customerId: true }, select: { customerId: true },
}) })
.then((r) => r.length), .then((r) => r.length),
this.prisma.transaction.groupBy({
by: ["currency"],
where: NOT_VOIDED,
_sum: { amount: true },
_count: { _all: true },
}),
this.prisma.transaction.groupBy({
by: ["domain", "currency"],
where: NOT_VOIDED,
_sum: { amount: true },
_count: { _all: true },
}),
]); ]);
const charges = await this.prisma.transaction.groupBy({ const byCurrency = await this.prisma.$queryRaw<
by: ["currency"], {
where: { AND: [{ amount: { lt: 0 } }, NOT_VOIDED] }, currency: string;
_sum: { amount: true }, net: Prisma.Decimal | null;
_count: { _all: true }, count: RawCount;
}); charges: Prisma.Decimal | null;
const credits = await this.prisma.transaction.groupBy({ chargeCount: RawCount;
by: ["currency"], credits: Prisma.Decimal | null;
where: { AND: [{ amount: { gt: 0 } }, NOT_VOIDED] }, creditCount: RawCount;
_sum: { amount: true }, }[]
_count: { _all: true }, >`
}); SELECT t.currency AS currency,
const chargeMap = new Map(charges.map((c) => [c.currency, c])); SUM(t.amount) AS net,
const creditMap = new Map(credits.map((c) => [c.currency, c])); COUNT(*) AS count,
SUM(CASE WHEN t.amount < 0 THEN t.amount ELSE 0 END) AS charges,
SUM(t.amount < 0) AS chargeCount,
SUM(CASE WHEN t.amount > 0 THEN t.amount ELSE 0 END) AS credits,
SUM(t.amount > 0) AS creditCount
FROM transactions t
${BALANCE_FLOOR_JOIN}
WHERE t.voidedAt IS NULL AND ${NOT_SUPERSEDED}
GROUP BY t.currency
`;
const byDomain = await this.prisma.$queryRaw<
{
domain: string;
currency: string;
net: Prisma.Decimal | null;
count: RawCount;
}[]
>`
SELECT t.domain AS domain, t.currency AS currency,
SUM(t.amount) AS net, COUNT(*) AS count
FROM transactions t
${BALANCE_FLOOR_JOIN}
WHERE t.voidedAt IS NULL AND ${NOT_SUPERSEDED}
GROUP BY t.domain, t.currency
`;
// How many customers sit on each side of the line, per currency — the // How many customers sit on each side of the line, per currency — the
// headline for a receivables view. Counted in SQL; a customer can be // headline for a receivables view. Counted in SQL; a customer can be
@@ -504,16 +585,19 @@ export class BillingService {
const sides = await this.prisma.$queryRaw< const sides = await this.prisma.$queryRaw<
{ {
currency: string; currency: string;
owing: bigint | number | string; owing: RawCount;
inCredit: bigint | number | string; inCredit: RawCount;
}[] }[]
>` >`
SELECT currency, SELECT currency,
SUM(bal < -0.005) AS owing, SUM(bal < -0.005) AS owing,
SUM(bal > 0.005) AS inCredit SUM(bal > 0.005) AS inCredit
FROM ( FROM (
SELECT customerId, currency, SUM(amount) AS bal SELECT t.customerId, t.currency, SUM(t.amount) AS bal
FROM transactions WHERE voidedAt IS NULL GROUP BY customerId, currency FROM transactions t
${BALANCE_FLOOR_JOIN}
WHERE t.voidedAt IS NULL AND ${NOT_SUPERSEDED}
GROUP BY t.customerId, t.currency
) x ) x
GROUP BY currency GROUP BY currency
`; `;
@@ -534,7 +618,7 @@ export class BillingService {
// Customers whose ledger spans both business lines — the whole reason this // Customers whose ledger spans both business lines — the whole reason this
// module is one view instead of two. // module is one view instead of two.
const crossLine = await this.prisma.$queryRaw<{ n: bigint | number | string }[]>` const crossLine = await this.prisma.$queryRaw<{ n: RawCount }[]>`
SELECT COUNT(*) AS n FROM ( SELECT COUNT(*) AS n FROM (
SELECT customerId FROM transactions WHERE voidedAt IS NULL SELECT customerId FROM transactions WHERE voidedAt IS NULL
GROUP BY customerId HAVING COUNT(DISTINCT domain) > 1 GROUP BY customerId HAVING COUNT(DISTINCT domain) > 1
@@ -549,20 +633,20 @@ export class BillingService {
lastMovement: lastRow?.transactionDate ?? null, lastMovement: lastRow?.transactionDate ?? null,
byCurrency: byCurrency.map((c) => ({ byCurrency: byCurrency.map((c) => ({
currency: c.currency, currency: c.currency,
net: c._sum.amount, net: c.net,
count: c._count._all, count: num(c.count),
charges: chargeMap.get(c.currency)?._sum.amount ?? null, charges: c.charges,
chargeCount: chargeMap.get(c.currency)?._count._all ?? 0, chargeCount: num(c.chargeCount),
credits: creditMap.get(c.currency)?._sum.amount ?? null, credits: c.credits,
creditCount: creditMap.get(c.currency)?._count._all ?? 0, creditCount: num(c.creditCount),
owing: num(sideMap.get(c.currency)?.owing), owing: num(sideMap.get(c.currency)?.owing),
inCredit: num(sideMap.get(c.currency)?.inCredit), inCredit: num(sideMap.get(c.currency)?.inCredit),
})), })),
byDomain: byDomain.map((d) => ({ byDomain: byDomain.map((d) => ({
domain: d.domain, domain: d.domain,
currency: d.currency, currency: d.currency,
net: d._sum.amount, net: d.net,
count: d._count._all, count: num(d.count),
})), })),
}; };
} }
@@ -589,7 +673,7 @@ export class BillingService {
}); });
const years = await this.prisma.$queryRaw< const years = await this.prisma.$queryRaw<
{ year: number; count: bigint | number | string }[] { year: number; count: RawCount }[]
>` >`
SELECT YEAR(transactionDate) AS year, COUNT(*) AS count SELECT YEAR(transactionDate) AS year, COUNT(*) AS count
FROM transactions WHERE voidedAt IS NULL GROUP BY year ORDER BY year DESC FROM transactions WHERE voidedAt IS NULL GROUP BY year ORDER BY year DESC
@@ -647,9 +731,32 @@ export class BillingService {
throw new NotFoundException(`Customer ${customerId} not found`); throw new NotFoundException(`Customer ${customerId} not found`);
} }
// One customer, so the balance floor is a single date rather than the
// derived table the aggregate queries join. See NOT_SUPERSEDED: rows before
// the opening balance are already inside it, and showing them would both
// double the total and make every balanceAfter below wrong.
//
// This is also what stops FEE ANUAL and fee15 leaking in. They are not in
// STATEMENT_EXCLUDED_SOURCE_TABLES — that list exists to reproduce legacy's
// DATOS2-only `datosfreak`, and it was letting 2,092 pre-cutover fee rows
// across 1,062 customers through, skewing the statement by -5,129,764
// against the number those customers have been quoted for years. Dating
// rather than source is the right test: a FEE ANUAL row *after* the opening
// balance is a real charge and still counts.
const floor = await this.prisma.transaction.findFirst({
where: {
customerId,
voidedAt: null,
type: { nameEn: BALANCE_FORWARD_TYPE },
},
orderBy: { transactionDate: "desc" },
select: { transactionDate: true },
});
const rows = await this.prisma.transaction.findMany({ const rows = await this.prisma.transaction.findMany({
where: { where: {
customerId, customerId,
...(floor ? { transactionDate: { gte: floor.transactionDate } } : {}),
// NULL-safe exclusion. `notIn` alone compiles to SQL `NOT IN`, and // NULL-safe exclusion. `notIn` alone compiles to SQL `NOT IN`, and
// `NULL NOT IN (...)` is NULL, not true — so every app-captured row // `NULL NOT IN (...)` is NULL, not true — so every app-captured row
// (which has no legacySourceTable) silently vanished from the // (which has no legacySourceTable) silently vanished from the
+6 -1
View File
@@ -400,7 +400,12 @@ export class OpsService implements OnModuleInit {
`${PIPEFAIL}echo '== Respaldo de seguridad previo ==' && ` + `${PIPEFAIL}echo '== Respaldo de seguridad previo ==' && ` +
`${this.dumpCommand(flags, db, out)} && ` + `${this.dumpCommand(flags, db, out)} && ` +
`echo '== Sincronización aditiva desde carpeta de ingesta ==' && ` + `echo '== Sincronización aditiva desde carpeta de ingesta ==' && ` +
`${shq(py)} ${runAll} --env ${shq(this.migrationEnv)} --sync`; // --stage is not optional here. The staged Parquet lives in the image
// at migration/output, NOT on a volume, so every redeploy wipes it and
// a sync without --stage dies on a missing stg_*/*.parquet. Re-staging
// is also the only thing that makes "desde carpeta de ingesta" true:
// stale Parquet would sync the previous upload, not the current one.
`${shq(py)} ${runAll} --env ${shq(this.migrationEnv)} --stage --sync`;
return { cmd, resolvedParams: { safetyBackup: file } }; return { cmd, resolvedParams: { safetyBackup: file } };
} }
+90
View File
@@ -4,6 +4,45 @@ import { promisify } from "node:util";
const exec = promisify(execFile); const exec = promisify(execFile);
/**
* How far the SQL thread is behind the I/O thread, in source binlog bytes.
*
* This is a different question from `secondsBehind`, and it answers the case
* that lag hides: while the SQL thread grinds through one huge transaction,
* `Seconds_Behind_Source` can sit still or even read 0, but the relay backlog
* is plainly shrinking (or not). It costs nothing extra — every field here
* comes out of the same `SHOW REPLICA STATUS` the panel already runs.
*
* Both positions are coordinates in the SOURCE's binlog, so they are only
* comparable while both threads are working on the SAME source file. When they
* are not, the replica is whole files behind and the byte delta is meaningless
* (positions restart at ~4 in each new file), so `backlogBytes` and `percent`
* are null and `sameFile` says why.
*/
export interface ApplyProgress {
/** Source binlog file the I/O thread is currently reading. */
sourceLogFile: string | null;
/** Position in `sourceLogFile` that the I/O thread has fetched up to. */
readPos: number;
/** Source binlog file the SQL thread is currently applying. */
relayLogFile: string | null;
/** Position in `relayLogFile` that the SQL thread has applied up to. */
execPos: number;
/** True while both threads are on the same source file. */
sameFile: boolean;
/** Fetched-but-not-yet-applied bytes. Null when the files differ. */
backlogBytes: number | null;
/**
* `execPos / readPos` as a percentage, null when the files differ.
*
* Deliberately never rounded up to 100 while any backlog remains: binlog
* positions are large, so a real backlog of a few KB is 99.99% of the file
* and would render as "caught up" when it is not. Read `backlogBytes === 0`
* for actually caught up.
*/
percent: number | null;
}
export interface ReplicationStatus { export interface ReplicationStatus {
/** false when the replica is not configured for this environment at all. */ /** false when the replica is not configured for this environment at all. */
configured: boolean; configured: boolean;
@@ -17,6 +56,8 @@ export interface ReplicationStatus {
lastIoError: string | null; lastIoError: string | null;
lastSqlError: string | null; lastSqlError: string | null;
sourceHost: string | null; sourceHost: string | null;
/** Relay-log apply progress. Null when the status output has no positions. */
apply: ApplyProgress | null;
/** Human-readable reason when healthy is false. */ /** Human-readable reason when healthy is false. */
problem: string | null; problem: string | null;
checkedAt: string; checkedAt: string;
@@ -57,6 +98,7 @@ export class ReplicationService {
lastIoError: null, lastIoError: null,
lastSqlError: null, lastSqlError: null,
sourceHost: null, sourceHost: null,
apply: null,
problem: null, problem: null,
checkedAt: now, checkedAt: now,
}; };
@@ -142,12 +184,60 @@ export class ReplicationService {
lastIoError, lastIoError,
lastSqlError, lastSqlError,
sourceHost: field("Source_Host"), sourceHost: field("Source_Host"),
// Reported, never folded into `healthy`: a non-zero backlog is the normal
// state of a working replica for the instant between fetch and apply, so
// alarming on it would cry wolf. It is here to answer "is it moving?"
// when the lag counter is stuck.
apply: applyProgress(raw),
problem, problem,
checkedAt: now, checkedAt: now,
}; };
} }
} }
/**
* Derive relay-apply progress from `SHOW REPLICA STATUS\G` output.
*
* Exported for testing. Free in query terms — it re-reads four more fields from
* the output the caller already has, with no second round trip to the replica
* and no connection to the source.
*
* @returns null when either position is missing or unparseable, which is what
* happens on a server that is not a replica at all.
*/
export function applyProgress(raw: string): ApplyProgress | null {
const num = (name: string): number | null => {
const v = replicaField(raw, name);
if (v === null || v === "NULL") return null;
const n = Number(v);
return Number.isFinite(n) ? n : null;
};
const readPos = num("Read_Source_Log_Pos");
const execPos = num("Exec_Source_Log_Pos");
if (readPos === null || execPos === null) return null;
const sourceLogFile = replicaField(raw, "Source_Log_File");
const relayLogFile = replicaField(raw, "Relay_Source_Log_File");
const sameFile =
sourceLogFile !== null && relayLogFile !== null && sourceLogFile === relayLogFile;
// Clamped at 0: the SQL thread cannot be ahead of the I/O thread, but the two
// fields are sampled independently, so a rotation racing this read can print
// a momentarily negative delta. Zero is the honest floor, not a bug.
const backlogBytes = sameFile ? Math.max(0, readPos - execPos) : null;
let percent: number | null = null;
if (backlogBytes !== null && readPos > 0) {
// Truncate rather than round, and hold short of 100 while bytes remain —
// see the doc on ApplyProgress.percent.
const p = Math.floor((execPos / readPos) * 10_000) / 100;
percent = backlogBytes === 0 ? 100 : Math.min(p, 99.99);
}
return { sourceLogFile, readPos, relayLogFile, execPos, sameFile, backlogBytes, percent };
}
/** /**
* Read one field out of `SHOW REPLICA STATUS\G` output. * Read one field out of `SHOW REPLICA STATUS\G` output.
* *
+94 -1
View File
@@ -1,4 +1,4 @@
import { replicaField } from "./replication.service"; import { applyProgress, replicaField } from "./replication.service";
/** /**
* Verbatim shape of `SHOW REPLICA STATUS\G` from the live replica, trimmed to * Verbatim shape of `SHOW REPLICA STATUS\G` from the live replica, trimmed to
@@ -13,6 +13,10 @@ const HEALTHY = [
" Replica_IO_State: Waiting for source to send event", " Replica_IO_State: Waiting for source to send event",
" Source_Host: 100.103.77.46", " Source_Host: 100.103.77.46",
" Source_User: repl", " Source_User: repl",
" Source_Log_File: binlog.000042",
" Read_Source_Log_Pos: 194884231",
" Relay_Source_Log_File: binlog.000042",
" Exec_Source_Log_Pos: 194884231",
" Replica_IO_Running: Yes", " Replica_IO_Running: Yes",
" Replica_SQL_Running: Yes", " Replica_SQL_Running: Yes",
" Replicate_Do_DB: ", " Replicate_Do_DB: ",
@@ -85,3 +89,92 @@ describe("replicaField", () => {
expect(replicaField(raw, "Last_SQL_Error")).toBe("boom"); expect(replicaField(raw, "Last_SQL_Error")).toBe("boom");
}); });
}); });
/** Builds the four position fields the apply-progress reader cares about. */
function positions(
sourceFile: string,
readPos: number | string,
relayFile: string,
execPos: number | string,
): string {
return [
` Source_Log_File: ${sourceFile}`,
` Read_Source_Log_Pos: ${readPos}`,
` Relay_Source_Log_File: ${relayFile}`,
` Exec_Source_Log_Pos: ${execPos}`,
].join("\n");
}
describe("applyProgress", () => {
it("reports zero backlog and 100% when both positions match", () => {
const p = applyProgress(HEALTHY)!;
expect(p.sameFile).toBe(true);
expect(p.sourceLogFile).toBe("binlog.000042");
expect(p.readPos).toBe(194884231);
expect(p.execPos).toBe(194884231);
expect(p.backlogBytes).toBe(0);
expect(p.percent).toBe(100);
});
it("reports the byte delta when the SQL thread trails inside one file", () => {
const p = applyProgress(positions("binlog.000042", 2_000_000, "binlog.000042", 1_500_000))!;
expect(p.backlogBytes).toBe(500_000);
expect(p.percent).toBe(75);
});
/**
* The reason the byte delta exists at all. `Seconds_Behind_Source` holds at 0
* while the SQL thread is mid-transaction, so the backlog is the only field
* that moves — and the only one that says the replica is not caught up.
*/
it("shows a backlog even when the lag counter reads zero", () => {
const raw = [
" Seconds_Behind_Source: 0",
positions("binlog.000042", 900, "binlog.000042", 400),
].join("\n");
expect(replicaField(raw, "Seconds_Behind_Source")).toBe("0");
expect(applyProgress(raw)!.backlogBytes).toBe(500);
});
/**
* Positions restart near 4 in every new binlog file, so subtracting across
* files produces a number that is not a backlog — here it would be a large
* NEGATIVE one, which would render as "ahead of the source".
*/
it("refuses to compare positions across different binlog files", () => {
const p = applyProgress(positions("binlog.000043", 500, "binlog.000042", 194_000_000))!;
expect(p.sameFile).toBe(false);
expect(p.backlogBytes).toBeNull();
expect(p.percent).toBeNull();
expect(p.sourceLogFile).toBe("binlog.000043");
expect(p.relayLogFile).toBe("binlog.000042");
});
/**
* Percent must not round up to 100 while bytes remain: binlog positions are
* large, so a genuine backlog is a rounding error away from the whole file
* and would otherwise render as "caught up" on a replica that is not.
*/
it("stops short of 100% while any backlog remains", () => {
const p = applyProgress(positions("binlog.000042", 194_884_231, "binlog.000042", 194_884_230))!;
expect(p.backlogBytes).toBe(1);
expect(p.percent).toBe(99.99);
});
/** Sampled independently, so a rotation racing the read can invert them. */
it("clamps a momentarily negative delta to zero", () => {
const p = applyProgress(positions("binlog.000042", 400, "binlog.000042", 500))!;
expect(p.backlogBytes).toBe(0);
expect(p.percent).toBe(100);
});
it("returns null when the server is not a replica and prints no positions", () => {
expect(applyProgress("")).toBeNull();
expect(applyProgress(BROKEN)).toBeNull();
});
/** A stopped thread makes MySQL print NULL, which is not a position. */
it("returns null when a position is NULL", () => {
expect(applyProgress(positions("binlog.000042", "NULL", "binlog.000042", 400))).toBeNull();
});
});
+1 -1
View File
@@ -1,6 +1,6 @@
{ {
"name": "@jorgecuadros/web", "name": "@jorgecuadros/web",
"version": "1.0.11", "version": "1.0.15",
"private": true, "private": true,
"scripts": { "scripts": {
"dev": "next dev -p 4500", "dev": "next dev -p 4500",
+51
View File
@@ -23,6 +23,7 @@ import {
} from "@/lib/api"; } from "@/lib/api";
import type { UploadProgress } from "@/lib/api"; import type { UploadProgress } from "@/lib/api";
import type { import type {
ApplyProgress,
BackupFile, BackupFile,
IngestFile, IngestFile,
OpsJob, OpsJob,
@@ -685,8 +686,58 @@ function ReplicationCard() {
label="Retraso" label="Retraso"
value={status.secondsBehind === null ? "sin dato" : `${status.secondsBehind} s`} value={status.secondsBehind === null ? "sin dato" : `${status.secondsBehind} s`}
/> />
<KV label="Pendiente de aplicar" value={backlogLabel(status.apply)} />
<KV label="Consultado" value={formatDateTime(status.checkedAt)} /> <KV label="Consultado" value={formatDateTime(status.checkedAt)} />
</div> </div>
<ApplyProgressBar apply={status.apply} />
</div>
);
}
/**
* Bytes the replica has fetched but not yet applied.
*
* Kept separate from the lag figure because it answers a question the lag
* cannot: while the SQL thread chews through one big transaction, the seconds
* counter can hold still, but this number visibly falls.
*/
function backlogLabel(apply: ApplyProgress | null): string {
if (!apply) return "sin dato";
// Different source binlog files means the replica is whole files behind and
// the byte delta is not a delta at all — positions restart in each new file.
if (!apply.sameFile) return "más de un archivo de binlog";
if (apply.backlogBytes === 0) return "al día";
return formatBytes(apply.backlogBytes);
}
/**
* Applied-vs-fetched bar. Rendered only when both threads are on the same
* source binlog file, because that is the only case where the percentage is
* arithmetic rather than a guess.
*/
function ApplyProgressBar({ apply }: { apply: ApplyProgress | null }) {
if (!apply || !apply.sameFile || apply.percent === null) return null;
return (
<div className="upload-progress" style={{ marginTop: 12 }}>
<div
className="progress-track"
role="progressbar"
aria-valuenow={apply.percent}
aria-valuemin={0}
aria-valuemax={100}
aria-label="Eventos aplicados de los recibidos"
>
<div className="progress-fill" style={{ width: `${apply.percent}%` }} />
</div>
<div className="upload-progress-stats mono">
<span>{apply.percent}% aplicado</span>
<span>
{apply.sourceLogFile} · {apply.execPos.toLocaleString("es-MX")} /{" "}
{apply.readPos.toLocaleString("es-MX")}
</span>
</div>
</div> </div>
); );
} }
+22
View File
@@ -99,10 +99,32 @@ export interface ReplicationStatus {
lastIoError: string | null; lastIoError: string | null;
lastSqlError: string | null; lastSqlError: string | null;
sourceHost: string | null; sourceHost: string | null;
apply: ApplyProgress | null;
problem: string | null; problem: string | null;
checkedAt: string; checkedAt: string;
} }
/**
* Relay-log apply progress, in source binlog bytes.
*
* Answers "is it moving?" when `secondsBehind` cannot: the lag counter sits
* still while the SQL thread works through one large transaction, but the
* backlog visibly shrinks. `backlogBytes === 0` is the only reading that means
* caught up — `percent` deliberately stops at 99.99 while bytes remain.
*
* Null fields when the two threads are on different source binlog files
* (`sameFile === false`), because the positions are then not comparable.
*/
export interface ApplyProgress {
sourceLogFile: string | null;
readPos: number;
relayLogFile: string | null;
execPos: number;
sameFile: boolean;
backlogBytes: number | null;
percent: number | null;
}
/** One of the four legacy Access files expected in the ingest folder. */ /** One of the four legacy Access files expected in the ingest folder. */
export interface IngestFile { export interface IngestFile {
name: string; name: string;
+17 -5
View File
@@ -15,6 +15,11 @@ Then:
./.venv/bin/python run_all.py --env dev # data only (staging already present) ./.venv/bin/python run_all.py --env dev # data only (staging already present)
./.venv/bin/python run_all.py --env prod --stage # re-extract from Access first, then load ./.venv/bin/python run_all.py --env prod --stage # re-extract from Access first, then load
--sync swaps the truncate+rebuild steps for the additive upsert ones. It reads
the same staged Parquet, so it needs --stage too unless a previous run left
migration/output populated on this machine — which is never true in a
container, where that directory is part of the image and dies with it.
Reproducing dev -> prod is exactly `--env prod` (plus --stage if the staged Reproducing dev -> prod is exactly `--env prod` (plus --stage if the staged
Parquet isn't present on the machine running it). Parquet isn't present on the machine running it).
@@ -100,15 +105,22 @@ def main() -> None:
help="upsert legacy rows and archive removed legacy rows; preserve manual rows") help="upsert legacy rows and archive removed legacy rows; preserve manual rows")
args = ap.parse_args() args = ap.parse_args()
if args.stage:
run([PY, str(HERE / "load_staging.py"), "--output-dir", str(HERE / "output")])
steps = SYNC_STEPS if args.sync else STEPS steps = SYNC_STEPS if args.sync else STEPS
for i, step in enumerate(steps, start=1): # Staging counts as a step when it runs: it is the slowest part of the pass
# (mdbtools re-reads every Access file), so leaving it outside the numbering
# would park the Operaciones progress bar at "nothing yet" for minutes.
total = len(steps) + (1 if args.stage else 0)
offset = 1 if args.stage else 0
if args.stage:
run([PY, str(HERE / "load_staging.py"), "--output-dir", str(HERE / "output")],
step=1, total=total)
for i, step in enumerate(steps, start=1 + offset):
cmd = [PY, str(HERE / step), "--env", args.env] cmd = [PY, str(HERE / step), "--env", args.env]
if args.sync: if args.sync:
cmd.append("--sync") cmd.append("--sync")
run(cmd, step=i, total=len(steps)) run(cmd, step=i, total=total)
print(f"\n✓ migration complete for env={args.env}") print(f"\n✓ migration complete for env={args.env}")
+16 -4
View File
@@ -149,10 +149,11 @@ def main():
skip_cust = skip_date = skip_dupe = 0 skip_cust = skip_date = skip_dupe = 0
def add(cid, domain, tdate, amount, currency, *, period=None, reference=None, def add(cid, domain, tdate, amount, currency, *, period=None, reference=None,
typeid=None, check=None, message=None, src_db=None, src_tbl=None, legacy=None): typeid=None, check=None, message=None, src_db=None, src_tbl=None, legacy=None,
outstanding=0):
tx.append((str(uuid.uuid4()), cid, domain, typeid, tdate, period, reference, tx.append((str(uuid.uuid4()), cid, domain, typeid, tdate, period, reference,
amount if amount is not None else Decimal(0), currency, None, check, amount if amount is not None else Decimal(0), currency, None, check,
message, 0, src_db, src_tbl, legacy)) message, outstanding, src_db, src_tbl, legacy))
# Business key of a real cash payment. `folio` is deliberately excluded: it # Business key of a real cash payment. `folio` is deliberately excluded: it
# is a per-table sequential number that collides between EFECTIVO and # is a per-table sequential number that collides between EFECTIVO and
@@ -230,6 +231,16 @@ def main():
src_db="UTILITIES", src_tbl=legacy_tbl, legacy=str(int(r["_row_num"]))) src_db="UTILITIES", src_tbl=legacy_tbl, legacy=str(int(r["_row_num"])))
def billing(name, legacy_tbl): def billing(name, legacy_tbl):
"""Load a DATOS2-shaped billing ledger.
NOPAGO is the legacy "still owed" flag. The website reads it directly —
`account.statement.php` splits the statement on `NOPAGO = 0` vs
`NOPAGO = 1` and renders the latter as the "Outstanding Bills Requiring
Attention" table — so dropping it does not merely lose a column, it
silently empties that whole section for anyone served off the platform.
Only these three tables carry it (76 rows set in DATOS2 today); the
EFECTIVO/FM3 cash streams have no such column and stay 0.
"""
nonlocal skip_cust, skip_date nonlocal skip_cust, skip_date
df = load("stg_utilities", name) df = load("stg_utilities", name)
for _, r in df.iterrows(): for _, r in df.iterrows():
@@ -243,7 +254,8 @@ def main():
add(cid, "UTILITY", td, dec(r["chargecredit"], Decimal(0)), "MXN", add(cid, "UTILITY", td, dec(r["chargecredit"], Decimal(0)), "MXN",
period=s(r["period"]), reference=s(r["refer"]), typeid=tid, period=s(r["period"]), reference=s(r["refer"]), typeid=tid,
check=s(r["cheque"]), src_db="UTILITIES", src_tbl=legacy_tbl, check=s(r["cheque"]), src_db="UTILITIES", src_tbl=legacy_tbl,
legacy=str(int(r["_row_num"]))) legacy=str(int(r["_row_num"])),
outstanding=1 if s(r["nopago"]) == "1" else 0)
def iva(): def iva():
nonlocal skip_cust nonlocal skip_cust
@@ -297,7 +309,7 @@ def main():
if new_types: if new_types:
c.executemany("INSERT INTO type_transactions (id,nameEn,nameEs,isService) VALUES (%s,%s,%s,%s)", new_types) c.executemany("INSERT INTO type_transactions (id,nameEn,nameEs,isService) VALUES (%s,%s,%s,%s)", new_types)
tx = [(t[0], t[1], t[2], (db_types.get(fresh_name.get(t[3])) if t[3] else None), *t[4:]) for t in tx] tx = [(t[0], t[1], t[2], (db_types.get(fresh_name.get(t[3])) if t[3] else None), *t[4:]) for t in tx]
c.executemany("INSERT INTO transactions (id,customerId,domain,typeId,transactionDate,period,reference,amount,currency,exchangeRate,checkNumber,message,outstanding,legacySourceDb,legacySourceTable,legacyId) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) ON DUPLICATE KEY UPDATE customerId=VALUES(customerId),domain=VALUES(domain),typeId=VALUES(typeId),transactionDate=VALUES(transactionDate),period=VALUES(period),reference=VALUES(reference),amount=VALUES(amount),currency=VALUES(currency),checkNumber=VALUES(checkNumber),message=VALUES(message),voidedAt=NULL", tx) c.executemany("INSERT INTO transactions (id,customerId,domain,typeId,transactionDate,period,reference,amount,currency,exchangeRate,checkNumber,message,outstanding,legacySourceDb,legacySourceTable,legacyId) VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) ON DUPLICATE KEY UPDATE customerId=VALUES(customerId),domain=VALUES(domain),typeId=VALUES(typeId),transactionDate=VALUES(transactionDate),period=VALUES(period),reference=VALUES(reference),amount=VALUES(amount),currency=VALUES(currency),checkNumber=VALUES(checkNumber),message=VALUES(message),outstanding=VALUES(outstanding),voidedAt=NULL", tx)
else: else:
c.execute("SET FOREIGN_KEY_CHECKS=0") c.execute("SET FOREIGN_KEY_CHECKS=0")
for t in ("transactions", "type_transactions", "exchange_rates"): for t in ("transactions", "type_transactions", "exchange_rates"):
+1 -1
View File
@@ -1,6 +1,6 @@
{ {
"name": "jorgecuadros-platform", "name": "jorgecuadros-platform",
"version": "1.0.11", "version": "1.0.15",
"private": true, "private": true,
"workspaces": [ "workspaces": [
"apps/*", "apps/*",
+1 -1
View File
@@ -1,6 +1,6 @@
{ {
"name": "@jorgecuadros/database", "name": "@jorgecuadros/database",
"version": "1.0.11", "version": "1.0.15",
"private": true, "private": true,
"main": "generated/client/index.js", "main": "generated/client/index.js",
"types": "generated/client/index.d.ts", "types": "generated/client/index.d.ts",