poh and poi were full truncate-reloads: 187k rows / 30 s and 529k / 75 s. Both now pull only changed POs -- 1,201 rows / 5.5 s and 3,660 / 10 s -- and reconcile.py --quick reads IN SYNC on all 115 / 120 metrics. POH turns out to carry its own change stamp in its "future" fields: KAFUT12 is Date Created and KAFUT20 is Date Updated, both CHAR(10) ISO text. Predicate is those two against the watermark, plus the open hot set KACRCM = '1' -- 1,035 of 187k headers, and KACRCM leads most of POH's logicals, so it is a keyed read. POI has no change stamp of its own, so it joins the header's changed set at PO grain; its merge key is kbpo#, so staging only the changed lines would drop the rest of the PO. CMS logs this family (POHL, POIL), but both are redundant for change detection: over 7 days POHL and the KAFUT20/12 window agree exactly, and POIL reports no PO that POHL does not, so a line edit always stamps the header. The logs' one unique contribution is deletes, which an incremental merge structurally cannot see -- delete-by-key only touches staged keys, so a purged PO would sit in the dest forever (~88 POs/yr). poh_deleted closes that: DQ0ACTN = '3' is the delete action, and the NOT EXISTS against live POH makes the list exact rather than a guess about action ordering, since a re-used PO number is back in POH and drops out on its own. Three dest hooks, in order -- clearing the tombstone of any PO that is live again must happen before the deletes, or the second hook could remove a row poh had just refreshed. Whenever these are scheduled, run order must be poh -> poi -> poh_deleted: poi's watermark resolves off cms.poh, and poh_deleted's first hook reads cms.poh expecting it to be fresh. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
16 KiB
CMS (DB2) → pipekit migration
Scope and plan for moving the remaining CMS/DB2 table syncs from the legacy
/opt/sync (jrunner + shell + cron) tool onto pipekit modules.
Decisions
- Destination schema:
cms.*. All migrated CMS tables land incms.<table>(pipekit convention), not the legacylgdat.*. Downstreamrlarp.*views that currently readlgdat.*must be repointed tocms.*as each table is cut over. - Sync tables as they reside on DB2. Column set and types come from fresh
introspection of the DB2 source (
QSYS2.SYSCOLUMNS), not from the legacypull.sql/build.sql. Legacy queries are not ported; we regenerate clean. - One table at a time on cutover. Generate/verify the module, repoint the
downstream view, then retire the
/opt/syncmodule.
Current state
- pipekit already covers 39 CMS tables (
source_connection: S78030956,FROM LGDAT.*, landing incms.*). /opt/sync/s7830956/(92 module dirs) holds the legacy CMS syncs. 24 overlap the 39; the rest is the migration backlog below.
The module recipe (pipekit)
A module is two sidecar files under config/modules/, named for the module:
<name>.json— definition. Required fields:name,source_connection,dest_connection,dest_table(schema-qualified),staging_table(pipekit_staging.<name>),merge_strategy(full|incremental|append),merge_key(required for incremental, elsenull),enabled,columns[],watermarks[],hooks[].dest_descriptionoptional.columns[]entries:source_name,source_type,dest_name,dest_type,description.
<name>.sql— the extractSELECT, run on the source connection in its dialect. DB2:RTRIM(col)on char cols,CASE WHEN col IN (DATE('0001-01-01'), DATE('9999-12-31')) THEN NULL ELSE col ENDon dates,"DOUBLE#QUOTES"for#-bearing identifiers. Watermarks are{name}placeholders substituted at run.
Group assignment + cron live in config/groups.json (add {"module","run_order"}
to a group's members). Source connections in config/connections.json
(S78030956 = DB2/AS400). pipekit apply hydrates pipekit.db from these files.
Merge semantics: staging recreated LIKE dest each run, then full = TRUNCATE +
INSERT; incremental = DELETE by merge_key + INSERT; append = INSERT only.
Introspection method (how drafts are generated)
pipekit exposes introspection two ways:
- HTTP API
GET /api/introspect/columns(:8200, Basic Auth). - Driver layer
pipekit.drivers.db2.DB2Driver.get_columns()— same code, in-process, needs onlyDB2PWfrom/etc/pipekit/secrets.env.
Drafts use the driver layer, reusing pipekit's own _TYPE_MAP, default_expression
(RTRIM / date-sentinel), and quote_identifier (# handling) — so output matches
pipekit style exactly, including real column/table descriptions from the DB2 catalog.
Type mapping: CHAR/VARCHAR/GRAPHIC/CLOB → text, DECIMAL/NUMERIC(p,s) preserved,
DATE → date, TIME → time, TIMESTAMP → text (pipekit convention), FLOAT → double precision, BIGINT/INTEGER/SMALLINT preserved.
Reproduce: drafts/bucket1/_generate.py (reads DB2 read-only, writes .json/.sql).
Bucket 1 — full truncate-reload (~41 tables) — DRAFTED
No watermark; legacy insert.sql was delete-all + reinsert. Simplest module:
merge_strategy: "full", merge_key: null.
Status: drafted, not activated. Files in /opt/pipekit/drafts/bucket1/
(41 .json + 41 .sql). Nothing wired to a group; pipekit apply not run.
Tables (source library in parens where not LGDAT):
adrs, cret, depts, fresre, resre, ftcstm, ftcstp, ftcstr, glie, icstm, icstp, icstr,
iprca, iprcb, iprcc, iprcctn, iprccto, irea, macgrp, majg, methdm, methdo, methdr,
mmgp, mmsl, opcode, plnt, punit, sach, sscc, stka, stkmm, stkmp, usrd, vend,
color / colorb / colortier / iprcbhc ("CMS.CUSLG"), usrcust (LGPGM),
ffpdglr1 (FANALDEV).
Activation (per table, when ready):
- Review the drafted
.json/.sql. mv drafts/bucket1/<name>.{json,sql} config/modules/.- Add
{"module":"<name>","run_order":N}to the right group ingroups.json. pipekit apply, thenpipekit run <name>to smoke-test.- Repoint any
rlarp.*view fromlgdat.<name>tocms.<name>.
Post-load hooks needed (legacy refreshed a mat-view; add a hooks[] entry once
the downstream view is repointed to cms.*): icstr → rlarp.icstx,
ftcstr → rlarp.ftcstx, ffpdglr1 → rlarp.pdglr1, stkmp → CALL rlarp.itemm_ps_build().
Bucket 2 — watermarked incremental (6 tables) — TODO
Need merge_key + a watermark resolver query (SELECT MAX(...) on the dest,
substituted into {wm} in the .sql). Templates: config/modules/ocrs.json,
gtran.json.
| Table | merge_key | watermark |
|---|---|---|
| qcrh | dcord# |
order # high-water |
| qcri | ddord# |
order # |
| qtnote | ggkey |
key |
| methh | anpart+anplnt (EXISTS merge) |
andate |
| iprcct | date | tadate |
| icstt | date | jhdate |
Each has a *_full unbounded counterpart in /opt/sync that can be dropped.
Price-list change log (iprcct / iprcctn / iprccto) — DONE 2026-08-05.
Append-only log, ~6.4–6.9M rows each, ~100 rows/day. All three: watermark
MAX(<date>) - 7 on the dest (guarded <= current_date), source
WHERE <date> BETWEEN '{wm}' AND DATE('9998-12-31'), incremental merge keyed on
the date alone — the staged set is every row on/after the watermark, so
delete-by-date replaces whole day partitions and re-running is a no-op. The
9998-12-31 upper bound matters: the source transform NULLs the 9999-12-31
sentinel, and a NULL merge key never matches the DELETE, so such a row would
duplicate on every run. Grain is (plcd, part, unit, date, time) for the header
and the same plus voll for the detail tables. All three reconcile exactly to
source counts. Members of group 14 (Price Lists), scheduled 0 3 * * * — local
time, after Sales Matrix at 02:20. The group is ~3.5 min, nearly all of it iprcc
(908k rows, full reload); the three change-log modules are seconds each.
Quote family (qcrh / qcri / qtnote) — DONE 2026-08-05. The order-number
watermark sketched in the table above is wrong and was replaced: quotes are edited
long after creation (of 644 headers updated in a 30-day sample, 25 were below
MAX(dcord#) - 1000), so an order-number window silently drops edits. QCRH shares
OCRH's DC* layout, so all three now use the ocrh/ocri pattern — watermark
MAX(dcudat) - 7 on cms.qcrh (guarded <= current_date) for all three
modules, header WHERE unioning dcudat / dcodat / dccdat, details joined to a
changed CTE over lgdat.qcrh. The union is load-bearing: 36,339 rows have
dccdat after dcudat, 270 have dcodat bumped without dcudat. Because the
detail watermarks resolve off the header table, run order within group 16 must
stay header-first and the lookback must exceed the sync interval.
qcri(492k rows, grainddord#+dditm#+dddes#) keys onddord#. Verified the changed-set catches line activity: all 5,807 lines created in a 30-day window (DDCTMS) belong to a header whose dates moved, 0 orphans.DDCTMSis a creation timestamp, not last-changed — unusable as a watermark.qtnote(253k rows) has no change signal of its own — no dates, no timestamp.GGKEYis9-digit quote number(header note) or+ 3-digit item(line note); all 50,876 note keys resolve to QCRH, only 6 also exist in OCRH, so it is purely quote notes. Keying onggkeywould leave a cleared note stranded in the dest forever, so a derived tail columncms.qtnote.ggord=SUBSTR(ggkey,1,9)was added (and backfilled) and is the merge key — delete now replaces all notes for a changed quote. Join is string-on-string viaDIGITS(dcord#), deliberately avoiding a numeric CAST: one junk key (19 spaces +.) would raise SQL0420 on cast, and staying textual leaves it harmlessly excluded (it survives in the dest withggord IS NULL, so dest count ties to source exactly at 252,721).
One-time rebaseline was required: a quote-level line-count diff found 273 quotes
disagreeing with source (in both directions, back to 2023) plus 22 missing
entirely — accumulated debt from the old order-number window. Rebaselined by
temporarily pointing each resolver at DATE '1900-01-01' and running, then
restoring; a wide date window provably covers all 82,224 headers, and delete-by-key
avoids the ACCESS EXCLUSIVE lock a full TRUNCATE would take on tables the 15-min
Quote Review group reads. All three now tie exactly to source. Group 16 (Quotes) is
scheduled */15 * * * * (~20 s per pass) with explicit run_order 1/2/3 — members
sort run_order, name, so header-first was previously only an alphabetical
accident, and both detail modules read their changed-set from cms.qcrh.
Residual risk to watch on qtnote: it inherits change detection entirely from the
header, so a note edited without the quote header's dates moving is invisible. Worth
a periodic wide-window rebaseline (the same widen/run/restore steps) or a
reconcile.py check.
icstt (Cost History, module 129) — reviewed 2026-08-05; KNOWN DRIFT ACCEPTED.
3.77M rows, watermark now MAX(jhdate) - 7 guarded <= current_date (the guard
matters — source holds 2 future-dated rows out to 2031). Scheduled 30 3 * * *
local as the sole member of group 13; before this it ran only when clicked.
Two defects found, one fixed, one deliberately left:
- Fixed: the resolver had no lookback, so anything back-dated below the high-water was invisible forever.
- Accepted: the source is purged at row level, and a date-keyed incremental
structurally cannot see it — delete-by-
jhdateonly touches dates in the staged set, and old dates never enter it, so every upstream deletion strands a dest row permanently. Currently +2,241 rows across 149 dates from 2010 onward (0.06% of the table). Verified by key-level diff on one date: 20 dest-only keys, 0 source-only, i.e. dest ⊇ source. This grows over time and silently inflates historical cost aggregates. Only a wide-window pass removes them.
Also 5 source rows dated 1953–1966 (1 row each, data-entry junk) have never synced and are intentionally not synced — a 7-day lookback starts ~2026-07-29 and never reaches them.
Consequence for monitoring: a bare reconcile.py icstt will always report DIVERGED.
Use the floor filter so the comparison is apples-to-apples — dates then tie exactly
(5,156 = 5,156) and the only mismatch is the purge overage:
PIPEKIT_SECRETS=/etc/pipekit/secrets.env .venv/bin/python reconcile.py icstt \
--super-quick --source-from "(SELECT * FROM LGDAT.ICSTT WHERE JHDATE >= DATE('2000-01-01')) t"
To rebaseline when there's a ~45 min window (est. from a 417k-row run at 5m06s):
point watermark 32 at SELECT DATE '2000-01-01' (that floor stages all 5,156 real
dates and clears the orphans while still skipping the 1953–66 junk), pipekit run icstt, restore the MAX(jhdate) - 7 resolver, then re-run the check above — it
should read IN SYNC. Delete-by-date, so no TRUNCATE and no ACCESS EXCLUSIVE lock.
Purchase order family (poh 138 / poi 139 / poh_deleted 144) — DONE
2026-08-18. Converted from full (187k rows / 30 s and 529k / 75 s) to
incremental: 1,201 rows / 5.5 s and 3,660 rows / 10 s. reconcile.py --quick reads IN SYNC on all 115 / 120 metrics, and a repeat run is a no-op.
- POH carries its own change stamp in its "future" fields —
KAFUT12= Date Created,KAFUT20= Date Updated, bothCHAR(10)ISOYYYY-MM-DD(so watermark compares are lexical and blanks sort below any floor).KAFUT20is populated on 113k of 187k rows — the gap is history that predates the field, irrelevant to a forward incremental. Watermarkpoh_wm/poi_wm=MAX(kafut20)::date - 7offcms.poh, capped atCURRENT_DATE. - Predicate:
KAFUT20 >= wm OR KAFUT12 >= wm OR KACRCM = '1'.KACRCM(1-Current / 2-Complete) leads most of POH's logicals, so the open hot set — only 1,035 of 187k headers, 2,374 of 529k lines — is a keyed read. - POI has no change stamp at all (every
KBFUT*is blank or a business attribute), so it follows the header on a PO-grain join, plus 4 open lines that hang off completed headers.merge_keyiskbpo#, so the pull must stay PO-grain — staging only the changed lines would drop the rest of the PO. - CMS logs this family too, and unlike
OCRITthe header log is usable:POHL(4.6M rows, back to 2015) has logicalPOHLX1keyed onDQ0TMSPalone;POILhas onlyPOILX1keyedDQ1PO#, DQ1TMSP. Both were measured as redundant for change detection: over 7 days POHL and theKAFUT20/12window agree exactly (365 POs each, 0 either direction), and POIL reports 0 POs POHL does not — a line edit always stamps the header. So no log join is in the predicate. DQ0ACTN = '3'is the delete action (verified on PO 618973: last event is an action 3, row gone from POH; 1/2 and 4/5 are before/after pairs). This is the one thing the date fields structurally cannot see — an incremental merge only deletes keys present in staging, so a purged PO would sit in the dest forever, the same shape as the acceptedicsttdrift. Rate is ~88 POs/yr (142,009 all-time since 2015, 140,924 currently absent from POH).poh_deleted(144) closes it:DQ0ACTN = '3' AND DQ0TMSP >= '{wm}' AND NOT EXISTS (… LGDAT.POH …),incrementalondq0po#, watermarkMAX(SUBSTRING(dq0tmsp,1,10))::date - 7off its own dest table. TheNOT EXISTSis what makes it exact instead of a guess about action ordering — a re-used PO number is back in POH and drops out of the list on its own. Three dest hooks, order matters: (1)DELETE FROM cms.poh_deleted WHERE dq0po# IN (SELECT kapo# FROM cms.poh)clears the tombstone of any PO that is live again before (2) deletes fromcms.pohand (3) fromcms.poi— without hook 1, hook 2 could remove a rowpohhad just refreshed. First run backfilled 140,924 rows in 13 s; steady state is a keyed range read of a handful. Hooks need an explicitconnection_idor the run fails withhook connection None not found.- Line-level deletes need no POIL module. Over 60 days, 66 POs had line
deletes: 36 still have a header and all 36 are already caught by
poi's predicate (the PO-grain delete+reinsert removes the vanished line), and the other 30 are whole-PO purges handled bypoh_deletedkeyed on the PO number. - All three are in no group — manual runs for now. Whenever they are
scheduled, order must be
poh→poi→poh_deleted:poi's watermark resolves offcms.poh, andpoh_deleted's hook 1 readscms.pohexpecting it to be fresh. - Note
WHERE TRUEcannot be used in these extracts — DB2 for i has no boolean literal and raises SQL0199.
Bucket 3 — special cases — decide, don't port as-is
- Reverse pushes PG→DB2 (
osmfs,reprice,reprice_test): go Postgres→AS400, opposite of pipekit's source→PG model. Leave in/opt/syncor redesign separately. - Stored-proc modules (
sb_ud_r2,sb_gj_r1,osm_sync):CALLprocs, not table extracts. Convert as hooks or leave out. - Derived RLARP reads (
family,ffterr,qrh,ctqpor): read derived AS400rlarp.*tables, not raw CMS masters. Confirm the upstream proc still runs first.
Sequencing recommendation
Prioritize by what the marts consume: the cost/pricing cluster (icst*, iprc*,
ftcst*, punit) and vendor/plant/method dims are highest value. Explicitly defer
or retire the reverse pushes and derived reads rather than porting them.