From b3b669c84abefb4439b7b6cf471cc7c69c8eecba Mon Sep 17 00:00:00 2001 From: bodymindarts Date: Wed, 16 Sep 2026 09:09:06 +0200 Subject: [PATCH 1/2] chore(deps)!: bump obix to 0.11.0, job to 0.14.0 obix 0.11.0 adds commit-ordered delivery as an opt-in second lane (#151); job 0.14.0 returns the existing job for duplicate requests (#214). Both are breaking 0.x minor bumps. Migration updated in place from the obix 0.11.0 template (nothing is live yet, so no drop/add migration is needed). The regenerated file is a faithful cala_-prefixed copy of the upstream template, verified by round-trip: stripping the prefix reproduces it byte for byte. The new schema adds commit_xid to the events table plus the commit-log lane (partitioned log, its partitions, and the singleton state row). SingletonSubscriber::handle_persistent now receives the event as the shared Arc the outbox decoded once, so the EC rollup handler's signature is updated to match; Arc derefs, so the body is unchanged. .sqlx cache regenerated against the new schema. Co-Authored-By: Claude Opus 5 --- Cargo.lock | 12 +- Cargo.toml | 4 +- ...4ea4849f4601100f4c8d1ace48c4e37a1db5c.json | 18 +++ ...1d9debb5554d37a615774de2b55d9af6261bf.json | 32 +++++ ...cfc68f043e2992f67db202753f0124faefd8.json} | 10 +- ...ff8ac22071b70a47e63338653932656f3b2b.json} | 6 +- ...0bb473e308127b8bf6803e545e4c06b1d9a3.json} | 4 +- ...f06ce14994f8a53a2a5c951c67d00b41e0f41.json | 24 ---- ...9539172c6a59e9b4b545fbe7f3196f2c91a16.json | 24 ---- ...ab4b8fcd29e6b3c5aa49823784906e2b23c14.json | 29 +++++ ...508343dbb97dcd3df85becfb2d9324ebea1c7.json | 16 +++ ...434f57c5bd08f43e030be61e4bd840d0323fa.json | 32 ----- ...8cf6ee47b75cfb271e79b72607b939ea034a2.json | 24 ++++ ...6247a55fc96787207d5a50b4a44314f8d97d.json} | 10 +- ...06b2b35d28d97bb7ac3c3cf022a697f6111a2.json | 16 --- ...8a73432f4681c3e5d8da21702343f8b68b64c.json | 38 ------ ...bfb9b11aab3cb86b23465eaa41cee0f311079.json | 65 ++++++++++ ...f3d08e6413d2e1dfd92af5d22f9ae9dfc8e82.json | 17 --- ...98082e2b9c48b06347736dd0c4e8b7649f91.json} | 10 +- ...a5f6bf7a8da2bee286f5851c3c5470ef80584.json | 24 ++++ ...22a07e1c0eebdc0262b4dea9e63fe90468e0f.json | 17 +++ ...fbedd8a29b74b0c6e88da16fc0fbdaa5a3ff9.json | 26 ++++ ...a7dfed872f1aa8be8c117d0c1fbb02f4fa4aa.json | 16 --- ...ad6af995efc7f5c314d842740e903ee8fc9d.json} | 10 +- ...70be5173a6b992936215300ad0737c46bfeec.json | 17 --- ...024c3a9fa3ea73e81f7561de13404ebe9a5d8.json | 30 +++++ ...dfc3b514858c2d6e83d69e328b9cf57e56df.json} | 6 +- ...ad69d3c3abb79c1fe9127c13fc7bc3e7c303.json} | 10 +- ...1072abface895a3d1631aefa66773b7b7759c.json | 29 ----- ...4d4f0e0c0613ed3eefa50a6295d02c3dbc907.json | 29 ----- ...4152ca649b225518216413ad8aae57d1098c1.json | 17 --- ...1e564551d49b17045ad7a3c5c4006c6e4b08.json} | 10 +- ...2da5bac04d17ecb2ca6590d273514a28e4156.json | 29 +++++ ...f1685ee85f8a00888655cbb92ff35c5e3ba24.json | 25 ++++ ...d01169264da82273629d14afe509d1b6005e0.json | 65 ++++++++++ ...fc59ba9158c3e1c361dd5841b4ad435fb9428.json | 30 +++++ ...86cd4e433c8d84930496e6f5487acc9f2fd93.json | 30 +++++ ...859510d17971a4923bc4dd4cd95a1512748c.json} | 10 +- ...9c7e90c6648c648d09a3f101edee12de354f2.json | 58 --------- ...efdaae473009d3dc398aacff73a58eb0e0ad7.json | 16 --- ...e15fad82ebfca75ef553f9b3f0c165ccfb55c.json | 24 ++++ ...5bda267c3c163e2735aac61d4de4dcc8c6b0c.json | 39 ++++++ ...25b024b41060b447355a4ff9a8a244527a601.json | 71 ++++++++++ ...cff19f4760d88c2934249dfd2a0285be3f3e5.json | 24 ++++ ...23c9a01ad5840700c7ba784f0cd7d8864cfc3.json | 29 ----- .../20251204130226_cala_obix_setup.sql | 122 +++++++++++++++--- cala-ledger/src/ec_rollup.rs | 3 +- 47 files changed, 799 insertions(+), 408 deletions(-) create mode 100644 cala-ledger/.sqlx/query-0288ab7b063f1bf967ede15fdc04ea4849f4601100f4c8d1ace48c4e37a1db5c.json create mode 100644 cala-ledger/.sqlx/query-0ad6aea4a22619fd708d1b2a8ea1d9debb5554d37a615774de2b55d9af6261bf.json rename cala-ledger/.sqlx/{query-fd1773cb0943111b2424b64c1c8199b2ea77590bdbd3916beadfcb78471bedb7.json => query-18b66d955afb396bee82661e6f0dcfc68f043e2992f67db202753f0124faefd8.json} (68%) rename cala-ledger/.sqlx/{query-35ffe70c0d0b7c1d63e9cb42577e9f3359c7eb9c9a0a6230d0e3e75daaabf1d2.json => query-1a9c66dc201a67b5b3714d98e9f6ff8ac22071b70a47e63338653932656f3b2b.json} (89%) rename cala-ledger/.sqlx/{query-bd38c42c1cca0e9b77744eb33bccabad662ce9f76958829910114fdf993dd170.json => query-25b23d03c8f4c17eaef798fa5fe60bb473e308127b8bf6803e545e4c06b1d9a3.json} (89%) delete mode 100644 cala-ledger/.sqlx/query-31384a875f0f689ff3e6140992df06ce14994f8a53a2a5c951c67d00b41e0f41.json delete mode 100644 cala-ledger/.sqlx/query-340045a3e3b2b59a47401a1a84f9539172c6a59e9b4b545fbe7f3196f2c91a16.json create mode 100644 cala-ledger/.sqlx/query-44feb771a63f2a8c914513a865aab4b8fcd29e6b3c5aa49823784906e2b23c14.json create mode 100644 cala-ledger/.sqlx/query-51d00fc3a4fe9a2efc41a68cf70508343dbb97dcd3df85becfb2d9324ebea1c7.json delete mode 100644 cala-ledger/.sqlx/query-51eda6b786a3139160d8b6c06e4434f57c5bd08f43e030be61e4bd840d0323fa.json create mode 100644 cala-ledger/.sqlx/query-525ff205aae7044cc1acac956d98cf6ee47b75cfb271e79b72607b939ea034a2.json rename cala-ledger/.sqlx/{query-eac483be72a79861524aebcd30b62948a4e87df694cd9cbe4a941879481756ae.json => query-588a243c4c09448d1abc2b8f21936247a55fc96787207d5a50b4a44314f8d97d.json} (65%) delete mode 100644 cala-ledger/.sqlx/query-5a71ac3f9107bd2bbd5cd5cd54b06b2b35d28d97bb7ac3c3cf022a697f6111a2.json delete mode 100644 cala-ledger/.sqlx/query-5cec2bcadeb25847daa2bb53b728a73432f4681c3e5d8da21702343f8b68b64c.json create mode 100644 cala-ledger/.sqlx/query-606f573c3c6cd3898ce64ad6fe3bfb9b11aab3cb86b23465eaa41cee0f311079.json delete mode 100644 cala-ledger/.sqlx/query-61f22dc336d52d7f3898e28088df3d08e6413d2e1dfd92af5d22f9ae9dfc8e82.json rename cala-ledger/.sqlx/{query-828c03ee80336e13e64238f149cf80d842631584afb5f5ae32560355b246fce4.json => query-6f028dbcb7d30f9dc4a744936f6198082e2b9c48b06347736dd0c4e8b7649f91.json} (80%) create mode 100644 cala-ledger/.sqlx/query-80bfbce31b7d3fb529766b245efa5f6bf7a8da2bee286f5851c3c5470ef80584.json create mode 100644 cala-ledger/.sqlx/query-80ea63f5ddd05bbf7519cf6c46222a07e1c0eebdc0262b4dea9e63fe90468e0f.json create mode 100644 cala-ledger/.sqlx/query-91c71c304b86ba87df527714d44fbedd8a29b74b0c6e88da16fc0fbdaa5a3ff9.json delete mode 100644 cala-ledger/.sqlx/query-932da1ca0fcd83172902956efaaa7dfed872f1aa8be8c117d0c1fbb02f4fa4aa.json rename cala-ledger/.sqlx/{query-00b6ca945b532fe362dbd8eac3a08e168ebc55089279f47f29d8e9f78ce8fda5.json => query-9463b386cdfbeaa41e6a0e5b8b95ad6af995efc7f5c314d842740e903ee8fc9d.json} (66%) delete mode 100644 cala-ledger/.sqlx/query-9ab25bea85f68a606013136bf1d70be5173a6b992936215300ad0737c46bfeec.json create mode 100644 cala-ledger/.sqlx/query-9e4fa6286210c1c103985279291024c3a9fa3ea73e81f7561de13404ebe9a5d8.json rename cala-ledger/.sqlx/{query-992aef4c4cb3b91da4d7cb0572f1fbf7447c3d4fd92d0cdf83541a58afd8e570.json => query-aca7a9b6f0dd28404dd43921b767dfc3b514858c2d6e83d69e328b9cf57e56df.json} (70%) rename cala-ledger/.sqlx/{query-bbbfd153dfa0453dc65549c83b34340b67b1347ec85793928f2658d1e7c63897.json => query-ad3164a1a7689fced3d6108d0eebad69d3c3abb79c1fe9127c13fc7bc3e7c303.json} (69%) delete mode 100644 cala-ledger/.sqlx/query-b01670f025421b5761f4907d81e1072abface895a3d1631aefa66773b7b7759c.json delete mode 100644 cala-ledger/.sqlx/query-b03d25d915a2be8e8270b8ea4d54d4f0e0c0613ed3eefa50a6295d02c3dbc907.json delete mode 100644 cala-ledger/.sqlx/query-bd553029c1b59d0df10b2e681344152ca649b225518216413ad8aae57d1098c1.json rename cala-ledger/.sqlx/{query-7089e2251160ce1d7266f49467204e83806992fa6f11588c25f26f5b90f23af9.json => query-bdfa672fbc2595d6eed0c82d542e1e564551d49b17045ad7a3c5c4006c6e4b08.json} (50%) create mode 100644 cala-ledger/.sqlx/query-bfc3d07b931304ab6d1a4a290ff2da5bac04d17ecb2ca6590d273514a28e4156.json create mode 100644 cala-ledger/.sqlx/query-cd605cbf07c5268c51f02f05d88f1685ee85f8a00888655cbb92ff35c5e3ba24.json create mode 100644 cala-ledger/.sqlx/query-d5a15e7d6001914866c9fd9149bd01169264da82273629d14afe509d1b6005e0.json create mode 100644 cala-ledger/.sqlx/query-d6964f14852c2b5d45efd02bde1fc59ba9158c3e1c361dd5841b4ad435fb9428.json create mode 100644 cala-ledger/.sqlx/query-da2a3f0d5bb15136393ada4a31886cd4e433c8d84930496e6f5487acc9f2fd93.json rename cala-ledger/.sqlx/{query-f1125bc627825540d799741f15703f39401b3fb50755e188e6626e08c89e871c.json => query-e16b163bdb81c2afe37ccd396207859510d17971a4923bc4dd4cd95a1512748c.json} (80%) delete mode 100644 cala-ledger/.sqlx/query-e66100b21243cdecb6046f3748e9c7e90c6648c648d09a3f101edee12de354f2.json delete mode 100644 cala-ledger/.sqlx/query-e67fb92be8782f30d872b37b7a8efdaae473009d3dc398aacff73a58eb0e0ad7.json create mode 100644 cala-ledger/.sqlx/query-ef771e1a68cc5b3722dcef68621e15fad82ebfca75ef553f9b3f0c165ccfb55c.json create mode 100644 cala-ledger/.sqlx/query-f004ea481cacdb54848f6a506295bda267c3c163e2735aac61d4de4dcc8c6b0c.json create mode 100644 cala-ledger/.sqlx/query-f22cd97b0bf5369bd0c072c69a625b024b41060b447355a4ff9a8a244527a601.json create mode 100644 cala-ledger/.sqlx/query-ff58bbfb397e2bbe3108ef363d3cff19f4760d88c2934249dfd2a0285be3f3e5.json delete mode 100644 cala-ledger/.sqlx/query-ffeaefb95f238936233fd691b2923c9a01ad5840700c7ba784f0cd7d8864cfc3.json diff --git a/Cargo.lock b/Cargo.lock index 7bdeff37f..2d724ecf9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1417,9 +1417,9 @@ checksum = "92ecc6618181def0457392ccd0ee51198e065e016d1d527a7ac1b6dc7c1f09d2" [[package]] name = "job" -version = "0.13.17" +version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "77590af8cc81bf200ba18c0fda79a0bef3e3fd60094effe6178a8bc76bfeb6ad" +checksum = "db855e958c34617af794ce40309f05d1e03294362d833f6fcde922b49dd44dde" dependencies = [ "async-trait", "chrono", @@ -1623,9 +1623,9 @@ dependencies = [ [[package]] name = "obix" -version = "0.10.0" +version = "0.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1cda18e4f942ca698312e4c394a3665b69866ce2c634a65ec06411219bf227d1" +checksum = "0618bac68a5e7182f1d2db974e5dd6f082c2c9ef875e6fa4006344bc46d5173b" dependencies = [ "async-trait", "chrono", @@ -1646,9 +1646,9 @@ dependencies = [ [[package]] name = "obix-macros" -version = "0.10.0" +version = "0.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "09c4ac967bde98de98a4a496a50edd17b4b568a1ca8489791ea5a77b07fcb1a8" +checksum = "ddf5f258599e1fd62ec84b7046f0dde129f08fce5a3d6ba3c31cdbce39844194" dependencies = [ "darling 0.24.0", "proc-macro2", diff --git a/Cargo.toml b/Cargo.toml index 61f8ef2f0..4ff46e6d6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -16,8 +16,8 @@ cala-ledger = { path = "cala-ledger", version = "0.28.1-dev" } cel = "0.14.5" es-entity = "0.12.21" -job = { version = "0.13.17", features = ["es-entity"] } -obix = { version = "0.10.0", default-features = false } +job = { version = "0.14.0", features = ["es-entity"] } +obix = { version = "0.11.0", default-features = false } anyhow = "1.0.99" cached = { version = "3.1", features = ["async"] } diff --git a/cala-ledger/.sqlx/query-0288ab7b063f1bf967ede15fdc04ea4849f4601100f4c8d1ace48c4e37a1db5c.json b/cala-ledger/.sqlx/query-0288ab7b063f1bf967ede15fdc04ea4849f4601100f4c8d1ace48c4e37a1db5c.json new file mode 100644 index 000000000..0f6245be2 --- /dev/null +++ b/cala-ledger/.sqlx/query-0288ab7b063f1bf967ede15fdc04ea4849f4601100f4c8d1ace48c4e37a1db5c.json @@ -0,0 +1,18 @@ +{ + "db_name": "PostgreSQL", + "query": "\n INSERT INTO job_executions\n (id, job_type, queue_id, unique_key, execute_at, alive_at, created_at)\n SELECT t.id, $2, NULL, t.unique_key, t.execute_at,\n COALESCE($5, NOW()), COALESCE($5, NOW())\n FROM UNNEST($1::uuid[], $3::text[], $4::timestamptz[])\n AS t(id, unique_key, execute_at)\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "UuidArray", + "Varchar", + "TextArray", + "TimestamptzArray", + "Timestamptz" + ] + }, + "nullable": [] + }, + "hash": "0288ab7b063f1bf967ede15fdc04ea4849f4601100f4c8d1ace48c4e37a1db5c" +} diff --git a/cala-ledger/.sqlx/query-0ad6aea4a22619fd708d1b2a8ea1d9debb5554d37a615774de2b55d9af6261bf.json b/cala-ledger/.sqlx/query-0ad6aea4a22619fd708d1b2a8ea1d9debb5554d37a615774de2b55d9af6261bf.json new file mode 100644 index 000000000..158de5b79 --- /dev/null +++ b/cala-ledger/.sqlx/query-0ad6aea4a22619fd708d1b2a8ea1d9debb5554d37a615774de2b55d9af6261bf.json @@ -0,0 +1,32 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT s.last_commit_seq AS \"last_commit_seq!: i64\",\n s.logged_through_sequence AS \"logged_through_sequence!: i64\",\n l.sequence AS \"logged_ahead?: i64\"\n FROM cala_persistent_outbox_commit_log_state s\n LEFT JOIN cala_persistent_outbox_commit_log l\n ON l.sequence > s.logged_through_sequence\n WHERE s.singleton\n ORDER BY l.sequence", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "last_commit_seq!: i64", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "logged_through_sequence!: i64", + "type_info": "Int8" + }, + { + "ordinal": 2, + "name": "logged_ahead?: i64", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + false, + false, + false + ] + }, + "hash": "0ad6aea4a22619fd708d1b2a8ea1d9debb5554d37a615774de2b55d9af6261bf" +} diff --git a/cala-ledger/.sqlx/query-fd1773cb0943111b2424b64c1c8199b2ea77590bdbd3916beadfcb78471bedb7.json b/cala-ledger/.sqlx/query-18b66d955afb396bee82661e6f0dcfc68f043e2992f67db202753f0124faefd8.json similarity index 68% rename from cala-ledger/.sqlx/query-fd1773cb0943111b2424b64c1c8199b2ea77590bdbd3916beadfcb78471bedb7.json rename to cala-ledger/.sqlx/query-18b66d955afb396bee82661e6f0dcfc68f043e2992f67db202753f0124faefd8.json index 8c26c74f8..89f32ee6d 100644 --- a/cala-ledger/.sqlx/query-fd1773cb0943111b2424b64c1c8199b2ea77590bdbd3916beadfcb78471bedb7.json +++ b/cala-ledger/.sqlx/query-18b66d955afb396bee82661e6f0dcfc68f043e2992f67db202753f0124faefd8.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT id, sequence, payload, tracing_context, recorded_at\n FROM cala_persistent_outbox_events\n WHERE sequence > $1\n AND sequence <= $2\n ORDER BY sequence ASC", + "query": "\n SELECT id, sequence, payload, tracing_context, recorded_at, commit_xid\n FROM cala_persistent_outbox_events\n WHERE sequence > $1\n AND sequence <= $2\n ORDER BY sequence ASC", "describe": { "columns": [ { @@ -27,6 +27,11 @@ "ordinal": 4, "name": "recorded_at", "type_info": "Timestamptz" + }, + { + "ordinal": 5, + "name": "commit_xid", + "type_info": "Int8" } ], "parameters": { @@ -40,8 +45,9 @@ false, true, true, + false, false ] }, - "hash": "fd1773cb0943111b2424b64c1c8199b2ea77590bdbd3916beadfcb78471bedb7" + "hash": "18b66d955afb396bee82661e6f0dcfc68f043e2992f67db202753f0124faefd8" } diff --git a/cala-ledger/.sqlx/query-35ffe70c0d0b7c1d63e9cb42577e9f3359c7eb9c9a0a6230d0e3e75daaabf1d2.json b/cala-ledger/.sqlx/query-1a9c66dc201a67b5b3714d98e9f6ff8ac22071b70a47e63338653932656f3b2b.json similarity index 89% rename from cala-ledger/.sqlx/query-35ffe70c0d0b7c1d63e9cb42577e9f3359c7eb9c9a0a6230d0e3e75daaabf1d2.json rename to cala-ledger/.sqlx/query-1a9c66dc201a67b5b3714d98e9f6ff8ac22071b70a47e63338653932656f3b2b.json index 5f348e31b..7e338f56c 100644 --- a/cala-ledger/.sqlx/query-35ffe70c0d0b7c1d63e9cb42577e9f3359c7eb9c9a0a6230d0e3e75daaabf1d2.json +++ b/cala-ledger/.sqlx/query-1a9c66dc201a67b5b3714d98e9f6ff8ac22071b70a47e63338653932656f3b2b.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH heads AS (\n -- At most one row per input queue (the LATERAL's LIMIT 1\n -- against deduped input), so no DISTINCT is needed.\n SELECT p.id\n FROM UNNEST($1::text[]) AS q(queue_id)\n CROSS JOIN LATERAL (\n SELECT id FROM job_executions\n WHERE state = 'parked' AND queue_id = q.queue_id\n ORDER BY execute_at, id\n LIMIT 1\n ) p\n WHERE NOT EXISTS (\n SELECT 1 FROM job_executions a\n WHERE a.queue_id = q.queue_id AND a.state IN ('pending', 'running')\n )\n ), locked AS MATERIALIZED (\n -- Lock every head in (queue_id, id) order before the UPDATE\n -- below touches any of them -- the same global order\n -- `lock_queue_occupants` and `Self::apply`'s own `locked` CTE\n -- use. A bare `UPDATE ... FROM heads` has no ordering\n -- guarantee of its own (`heads`'s row order is not a lock\n -- order), so a multi-queue batch completion freeing several\n -- queues here could otherwise acquire in planner/`UNNEST`\n -- order and deadlock against a concurrent spawn's pin or\n -- swap touching the same rows in the opposite order.\n SELECT je.id FROM job_executions je\n WHERE je.id IN (SELECT id FROM heads)\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je SET state = 'pending'\n FROM locked l WHERE je.id = l.id\n RETURNING je.job_type, je.execute_at AS \"execute_at!\"\n ", + "query": "\n WITH heads AS (\n -- At most one row per input queue (the LATERAL's LIMIT 1\n -- against deduped input), so no DISTINCT is needed.\n SELECT p.id\n FROM UNNEST($1::text[]) AS q(queue_id)\n CROSS JOIN LATERAL (\n SELECT id FROM job_executions\n WHERE state = 'parked' AND queue_id = q.queue_id\n ORDER BY execute_at, id\n LIMIT 1\n ) p\n WHERE NOT EXISTS (\n SELECT 1 FROM job_executions a\n WHERE a.queue_id = q.queue_id AND a.state IN ('pending', 'running')\n )\n ), locked AS MATERIALIZED (\n -- Lock every head in (queue_id, id) order before the UPDATE\n -- below touches any of them -- the same global order\n -- `lock_queue_occupants` and `Self::apply`'s own `locked` CTE\n -- use. A bare `UPDATE ... FROM heads` has no ordering\n -- guarantee of its own (`heads`'s row order is not a lock\n -- order), so a multi-queue batch completion freeing several\n -- queues here could otherwise acquire in planner/`UNNEST`\n -- order and deadlock against a concurrent spawn's pin or\n -- swap touching the same rows in the opposite order.\n SELECT je.id FROM job_executions je\n WHERE je.id IN (SELECT id FROM heads)\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je SET state = 'pending'\n FROM locked l WHERE je.id = l.id AND je.state = 'parked'\n RETURNING je.job_type, je.execute_at AS \"execute_at?\"\n ", "describe": { "columns": [ { @@ -10,7 +10,7 @@ }, { "ordinal": 1, - "name": "execute_at!", + "name": "execute_at?", "type_info": "Timestamptz" } ], @@ -24,5 +24,5 @@ true ] }, - "hash": "35ffe70c0d0b7c1d63e9cb42577e9f3359c7eb9c9a0a6230d0e3e75daaabf1d2" + "hash": "1a9c66dc201a67b5b3714d98e9f6ff8ac22071b70a47e63338653932656f3b2b" } diff --git a/cala-ledger/.sqlx/query-bd38c42c1cca0e9b77744eb33bccabad662ce9f76958829910114fdf993dd170.json b/cala-ledger/.sqlx/query-25b23d03c8f4c17eaef798fa5fe60bb473e308127b8bf6803e545e4c06b1d9a3.json similarity index 89% rename from cala-ledger/.sqlx/query-bd38c42c1cca0e9b77744eb33bccabad662ce9f76958829910114fdf993dd170.json rename to cala-ledger/.sqlx/query-25b23d03c8f4c17eaef798fa5fe60bb473e308127b8bf6803e545e4c06b1d9a3.json index 28303063f..39f2a4ac5 100644 --- a/cala-ledger/.sqlx/query-bd38c42c1cca0e9b77744eb33bccabad662ce9f76958829910114fdf993dd170.json +++ b/cala-ledger/.sqlx/query-25b23d03c8f4c17eaef798fa5fe60bb473e308127b8bf6803e545e4c06b1d9a3.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH orphan_queues AS (\n SELECT DISTINCT p.queue_id\n FROM job_executions p\n WHERE p.state = 'parked'\n AND NOT EXISTS (\n SELECT 1 FROM job_executions a\n WHERE a.queue_id = p.queue_id AND a.state IN ('pending', 'running')\n )\n ), heads AS (\n SELECT h.id FROM orphan_queues oq\n CROSS JOIN LATERAL (\n SELECT id FROM job_executions\n WHERE state = 'parked' AND queue_id = oq.queue_id\n ORDER BY execute_at, id\n LIMIT 1\n ) h\n ), locked AS MATERIALIZED (\n SELECT je.id FROM job_executions je\n WHERE je.id IN (SELECT id FROM heads)\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je SET state = 'pending'\n FROM locked l WHERE je.id = l.id\n RETURNING je.job_type\n ", + "query": "\n WITH orphan_queues AS (\n SELECT DISTINCT p.queue_id\n FROM job_executions p\n WHERE p.state = 'parked'\n AND NOT EXISTS (\n SELECT 1 FROM job_executions a\n WHERE a.queue_id = p.queue_id AND a.state IN ('pending', 'running')\n )\n ), heads AS (\n SELECT h.id FROM orphan_queues oq\n CROSS JOIN LATERAL (\n SELECT id FROM job_executions\n WHERE state = 'parked' AND queue_id = oq.queue_id\n ORDER BY execute_at, id\n LIMIT 1\n ) h\n ), locked AS MATERIALIZED (\n SELECT je.id FROM job_executions je\n WHERE je.id IN (SELECT id FROM heads)\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je SET state = 'pending'\n FROM locked l WHERE je.id = l.id AND je.state = 'parked'\n RETURNING je.job_type\n ", "describe": { "columns": [ { @@ -16,5 +16,5 @@ false ] }, - "hash": "bd38c42c1cca0e9b77744eb33bccabad662ce9f76958829910114fdf993dd170" + "hash": "25b23d03c8f4c17eaef798fa5fe60bb473e308127b8bf6803e545e4c06b1d9a3" } diff --git a/cala-ledger/.sqlx/query-31384a875f0f689ff3e6140992df06ce14994f8a53a2a5c951c67d00b41e0f41.json b/cala-ledger/.sqlx/query-31384a875f0f689ff3e6140992df06ce14994f8a53a2a5c951c67d00b41e0f41.json deleted file mode 100644 index 7898fb647..000000000 --- a/cala-ledger/.sqlx/query-31384a875f0f689ff3e6140992df06ce14994f8a53a2a5c951c67d00b41e0f41.json +++ /dev/null @@ -1,24 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n -- `(queue_id, id)`-ordered lock first, for the same reason\n -- `reclaim_lost_jobs` takes one: this transaction goes on to call\n -- `PromoteHeadsHook::apply`, which locks in that order, over\n -- queues this reset just touched.\n WITH locked AS MATERIALIZED (\n SELECT je.id, u.execute_at FROM job_executions je\n JOIN UNNEST($1::uuid[], $2::timestamptz[]) AS u(id, execute_at)\n ON je.id = u.id\n WHERE je.state = 'running' AND je.poller_instance_id = $3\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE OF je\n )\n UPDATE job_executions je\n SET state = 'pending', poller_instance_id = NULL, execute_at = l.execute_at\n FROM locked l WHERE je.id = l.id\n RETURNING je.id AS \"id!: JobId\"\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!: JobId", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "UuidArray", - "TimestamptzArray", - "Uuid" - ] - }, - "nullable": [ - false - ] - }, - "hash": "31384a875f0f689ff3e6140992df06ce14994f8a53a2a5c951c67d00b41e0f41" -} diff --git a/cala-ledger/.sqlx/query-340045a3e3b2b59a47401a1a84f9539172c6a59e9b4b545fbe7f3196f2c91a16.json b/cala-ledger/.sqlx/query-340045a3e3b2b59a47401a1a84f9539172c6a59e9b4b545fbe7f3196f2c91a16.json deleted file mode 100644 index df35d999b..000000000 --- a/cala-ledger/.sqlx/query-340045a3e3b2b59a47401a1a84f9539172c6a59e9b4b545fbe7f3196f2c91a16.json +++ /dev/null @@ -1,24 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH deleted AS (\n DELETE FROM job_executions\n WHERE id = $1 AND poller_instance_id = $2\n RETURNING id, queue_id\n ), cleanup AS (\n DELETE FROM job_execution_states s USING deleted d\n WHERE s.id = d.id AND NOT $3::boolean\n )\n SELECT d.queue_id AS \"queue_id?\" FROM deleted d\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "queue_id?", - "type_info": "Varchar" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Uuid", - "Bool" - ] - }, - "nullable": [ - true - ] - }, - "hash": "340045a3e3b2b59a47401a1a84f9539172c6a59e9b4b545fbe7f3196f2c91a16" -} diff --git a/cala-ledger/.sqlx/query-44feb771a63f2a8c914513a865aab4b8fcd29e6b3c5aa49823784906e2b23c14.json b/cala-ledger/.sqlx/query-44feb771a63f2a8c914513a865aab4b8fcd29e6b3c5aa49823784906e2b23c14.json new file mode 100644 index 000000000..ea8d87d1e --- /dev/null +++ b/cala-ledger/.sqlx/query-44feb771a63f2a8c914513a865aab4b8fcd29e6b3c5aa49823784906e2b23c14.json @@ -0,0 +1,29 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT unique_key AS \"unique_key!\", id AS \"id: JobId\" FROM job_executions\n WHERE job_type = $1 AND unique_key = ANY($2)", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "unique_key!", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "id: JobId", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Text", + "TextArray" + ] + }, + "nullable": [ + true, + false + ] + }, + "hash": "44feb771a63f2a8c914513a865aab4b8fcd29e6b3c5aa49823784906e2b23c14" +} diff --git a/cala-ledger/.sqlx/query-51d00fc3a4fe9a2efc41a68cf70508343dbb97dcd3df85becfb2d9324ebea1c7.json b/cala-ledger/.sqlx/query-51d00fc3a4fe9a2efc41a68cf70508343dbb97dcd3df85becfb2d9324ebea1c7.json new file mode 100644 index 000000000..1f3af3707 --- /dev/null +++ b/cala-ledger/.sqlx/query-51d00fc3a4fe9a2efc41a68cf70508343dbb97dcd3df85becfb2d9324ebea1c7.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH to_touch AS MATERIALIZED (\n SELECT id FROM job_executions\n WHERE poller_instance_id = $2\n AND state = 'running'\n AND id = ANY($3)\n ORDER BY queue_id, id\n FOR NO KEY UPDATE SKIP LOCKED\n )\n UPDATE job_executions je\n SET alive_at = $1\n FROM to_touch t\n WHERE je.id = t.id\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Timestamptz", + "Uuid", + "UuidArray" + ] + }, + "nullable": [] + }, + "hash": "51d00fc3a4fe9a2efc41a68cf70508343dbb97dcd3df85becfb2d9324ebea1c7" +} diff --git a/cala-ledger/.sqlx/query-51eda6b786a3139160d8b6c06e4434f57c5bd08f43e030be61e4bd840d0323fa.json b/cala-ledger/.sqlx/query-51eda6b786a3139160d8b6c06e4434f57c5bd08f43e030be61e4bd840d0323fa.json deleted file mode 100644 index bfd09ad2a..000000000 --- a/cala-ledger/.sqlx/query-51eda6b786a3139160d8b6c06e4434f57c5bd08f43e030be61e4bd840d0323fa.json +++ /dev/null @@ -1,32 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH ins AS (\n INSERT INTO job_executions\n (id, job_type, queue_id, unique_key, execute_at, alive_at, created_at)\n VALUES ($1, $2, NULL, $3, $4, COALESCE($5, NOW()), COALESCE($5, NOW()))\n ON CONFLICT (job_type, unique_key) WHERE unique_key IS NOT NULL\n DO NOTHING\n RETURNING id\n )\n SELECT (SELECT id FROM ins) AS \"inserted?: JobId\",\n (SELECT id FROM job_executions\n WHERE job_type = $2 AND unique_key = $3) AS \"live?: JobId\"\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "inserted?: JobId", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "live?: JobId", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Varchar", - "Varchar", - "Timestamptz", - "Timestamptz" - ] - }, - "nullable": [ - null, - null - ] - }, - "hash": "51eda6b786a3139160d8b6c06e4434f57c5bd08f43e030be61e4bd840d0323fa" -} diff --git a/cala-ledger/.sqlx/query-525ff205aae7044cc1acac956d98cf6ee47b75cfb271e79b72607b939ea034a2.json b/cala-ledger/.sqlx/query-525ff205aae7044cc1acac956d98cf6ee47b75cfb271e79b72607b939ea034a2.json new file mode 100644 index 000000000..d6783b6eb --- /dev/null +++ b/cala-ledger/.sqlx/query-525ff205aae7044cc1acac956d98cf6ee47b75cfb271e79b72607b939ea034a2.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH locked AS MATERIALIZED (\n SELECT je.id, u.execute_at FROM job_executions je\n JOIN UNNEST($1::uuid[], $2::timestamptz[]) AS u(id, execute_at)\n ON je.id = u.id\n WHERE je.state = 'running' AND je.poller_instance_id = $3\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE OF je\n )\n UPDATE job_executions je\n SET state = 'pending', poller_instance_id = NULL, execute_at = l.execute_at\n FROM locked l\n WHERE je.id = l.id AND je.state = 'running' AND je.poller_instance_id = $3\n RETURNING je.id AS \"id!: JobId\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "UuidArray", + "TimestamptzArray", + "Uuid" + ] + }, + "nullable": [ + false + ] + }, + "hash": "525ff205aae7044cc1acac956d98cf6ee47b75cfb271e79b72607b939ea034a2" +} diff --git a/cala-ledger/.sqlx/query-eac483be72a79861524aebcd30b62948a4e87df694cd9cbe4a941879481756ae.json b/cala-ledger/.sqlx/query-588a243c4c09448d1abc2b8f21936247a55fc96787207d5a50b4a44314f8d97d.json similarity index 65% rename from cala-ledger/.sqlx/query-eac483be72a79861524aebcd30b62948a4e87df694cd9cbe4a941879481756ae.json rename to cala-ledger/.sqlx/query-588a243c4c09448d1abc2b8f21936247a55fc96787207d5a50b4a44314f8d97d.json index b3d8631e9..cb25d09e5 100644 --- a/cala-ledger/.sqlx/query-eac483be72a79861524aebcd30b62948a4e87df694cd9cbe4a941879481756ae.json +++ b/cala-ledger/.sqlx/query-588a243c4c09448d1abc2b8f21936247a55fc96787207d5a50b4a44314f8d97d.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH win AS (\n SELECT sequence, ROW_NUMBER() OVER (ORDER BY sequence) AS rn\n FROM cala_persistent_outbox_events\n WHERE sequence > $1\n AND sequence <= $1 + $2\n )\n SELECT e.sequence AS \"sequence!: i64\", e.id AS \"id!\", e.payload,\n e.tracing_context, e.recorded_at AS \"recorded_at!\"\n FROM cala_persistent_outbox_events e\n WHERE e.sequence > $1\n AND e.sequence <= $1 + $2\n AND e.sequence < (\n SELECT COALESCE(MIN(sequence), $1 + $2 + 1)\n FROM win\n WHERE sequence <> $1 + rn\n )\n ORDER BY e.sequence ASC", + "query": "\n WITH win AS (\n SELECT sequence, ROW_NUMBER() OVER (ORDER BY sequence) AS rn\n FROM cala_persistent_outbox_events\n WHERE sequence > $1\n AND sequence <= $1 + $2\n )\n SELECT e.sequence AS \"sequence!: i64\", e.id AS \"id!\", e.payload,\n e.tracing_context, e.recorded_at AS \"recorded_at!\", e.commit_xid\n FROM cala_persistent_outbox_events e\n WHERE e.sequence > $1\n AND e.sequence <= $1 + $2\n AND e.sequence < (\n SELECT COALESCE(MIN(sequence), $1 + $2 + 1)\n FROM win\n WHERE sequence <> $1 + rn\n )\n ORDER BY e.sequence ASC", "describe": { "columns": [ { @@ -27,6 +27,11 @@ "ordinal": 4, "name": "recorded_at!", "type_info": "Timestamptz" + }, + { + "ordinal": 5, + "name": "commit_xid", + "type_info": "Int8" } ], "parameters": { @@ -40,8 +45,9 @@ false, true, true, + false, false ] }, - "hash": "eac483be72a79861524aebcd30b62948a4e87df694cd9cbe4a941879481756ae" + "hash": "588a243c4c09448d1abc2b8f21936247a55fc96787207d5a50b4a44314f8d97d" } diff --git a/cala-ledger/.sqlx/query-5a71ac3f9107bd2bbd5cd5cd54b06b2b35d28d97bb7ac3c3cf022a697f6111a2.json b/cala-ledger/.sqlx/query-5a71ac3f9107bd2bbd5cd5cd54b06b2b35d28d97bb7ac3c3cf022a697f6111a2.json deleted file mode 100644 index cd09f5d0a..000000000 --- a/cala-ledger/.sqlx/query-5a71ac3f9107bd2bbd5cd5cd54b06b2b35d28d97bb7ac3c3cf022a697f6111a2.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE job_executions AS je\n SET state = 'pending', execute_at = u.execute_at, attempt_index = 1,\n poller_instance_id = NULL\n FROM UNNEST($1::uuid[], $2::timestamptz[]) AS u(id, execute_at)\n WHERE je.id = u.id AND je.poller_instance_id = $3\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "UuidArray", - "TimestamptzArray", - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "5a71ac3f9107bd2bbd5cd5cd54b06b2b35d28d97bb7ac3c3cf022a697f6111a2" -} diff --git a/cala-ledger/.sqlx/query-5cec2bcadeb25847daa2bb53b728a73432f4681c3e5d8da21702343f8b68b64c.json b/cala-ledger/.sqlx/query-5cec2bcadeb25847daa2bb53b728a73432f4681c3e5d8da21702343f8b68b64c.json deleted file mode 100644 index 7719cbc11..000000000 --- a/cala-ledger/.sqlx/query-5cec2bcadeb25847daa2bb53b728a73432f4681c3e5d8da21702343f8b68b64c.json +++ /dev/null @@ -1,38 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH input AS MATERIALIZED (\n SELECT * FROM UNNEST($1::uuid[], $2::text[], $3::text[], $4::timestamptz[])\n AS t(id, job_type, queue_id, execute_at)\n ORDER BY queue_id, id\n ), ins AS (\n INSERT INTO job_executions\n (id, job_type, queue_id, unique_key, state, attempt_index, execute_at, alive_at, created_at)\n SELECT id, job_type, queue_id, NULL, 'pending', 1, execute_at,\n COALESCE($5, NOW()), COALESCE($5, NOW())\n FROM input\n ON CONFLICT (queue_id) WHERE state IN ('pending','running') AND queue_id IS NOT NULL\n DO NOTHING\n RETURNING id, queue_id\n ), parked AS (\n INSERT INTO job_executions\n (id, job_type, queue_id, unique_key, state, attempt_index, execute_at, alive_at, created_at)\n SELECT i.id, i.job_type, i.queue_id, NULL, 'parked', 1, i.execute_at,\n COALESCE($5, NOW()), COALESCE($5, NOW())\n FROM input i\n WHERE i.id NOT IN (SELECT id FROM ins)\n RETURNING id, queue_id\n )\n SELECT r.id AS \"id!: JobId\", TRUE AS \"landed_pending!\", NULL::uuid AS \"occupant_id?\"\n FROM ins r\n UNION ALL\n SELECT p.id AS \"id!: JobId\", FALSE AS \"landed_pending!\",\n COALESCE(\n (SELECT w.id FROM ins w WHERE w.queue_id = p.queue_id),\n (SELECT o.id FROM job_executions o\n WHERE o.queue_id = p.queue_id AND o.state = 'pending')\n ) AS \"occupant_id?\"\n FROM parked p\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!: JobId", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "landed_pending!", - "type_info": "Bool" - }, - { - "ordinal": 2, - "name": "occupant_id?", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "UuidArray", - "TextArray", - "TextArray", - "TimestamptzArray", - "Timestamptz" - ] - }, - "nullable": [ - null, - null, - null - ] - }, - "hash": "5cec2bcadeb25847daa2bb53b728a73432f4681c3e5d8da21702343f8b68b64c" -} diff --git a/cala-ledger/.sqlx/query-606f573c3c6cd3898ce64ad6fe3bfb9b11aab3cb86b23465eaa41cee0f311079.json b/cala-ledger/.sqlx/query-606f573c3c6cd3898ce64ad6fe3bfb9b11aab3cb86b23465eaa41cee0f311079.json new file mode 100644 index 000000000..f315dd535 --- /dev/null +++ b/cala-ledger/.sqlx/query-606f573c3c6cd3898ce64ad6fe3bfb9b11aab3cb86b23465eaa41cee0f311079.json @@ -0,0 +1,65 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH limits AS (\n SELECT l.job_type, l.row_limit,\n LEAST(l.row_limit, $1::int4) * $7::int4 AS type_window_limit\n FROM UNNEST($4::text[], $6::int4[]) AS l(job_type, row_limit)\n WHERE l.row_limit > 0\n ),\n window_rows AS (\n SELECT d.id, d.execute_at, d.job_type\n FROM limits t\n CROSS JOIN LATERAL (\n SELECT je.id, je.execute_at, je.job_type\n FROM job_executions je\n WHERE je.state = 'pending'\n AND je.job_type = t.job_type\n AND je.execute_at <= $2::timestamptz\n ORDER BY je.execute_at, je.id\n LIMIT t.type_window_limit\n ) d\n ),\n ordered_candidates AS (\n SELECT id, execute_at, job_type,\n ROW_NUMBER() OVER (\n PARTITION BY job_type ORDER BY execute_at\n ) AS type_rn\n FROM window_rows\n ),\n locked AS (\n -- FOR UPDATE OF je: bare FOR UPDATE errors on a nullable join side.\n SELECT je.id, je.attempt_index, c.job_type, c.execute_at\n FROM ordered_candidates c\n JOIN job_executions je ON je.id = c.id\n ORDER BY c.type_rn ASC, c.execute_at ASC\n LIMIT $1\n FOR UPDATE OF je SKIP LOCKED\n ),\n selected_jobs AS (\n SELECT t.id, cp.execution_state_json AS data_json, t.attempt_index\n FROM (\n SELECT l.*,\n ROW_NUMBER() OVER (\n PARTITION BY l.job_type ORDER BY l.execute_at\n ) AS type_rn\n FROM locked l\n ) t\n JOIN limits lim ON lim.job_type = t.job_type\n LEFT JOIN job_execution_states cp ON cp.id = t.id\n WHERE t.type_rn <= lim.row_limit\n ),\n updated AS (\n UPDATE job_executions AS je\n SET state = 'running', alive_at = $5, execute_at = NULL, poller_instance_id = $3\n FROM selected_jobs\n WHERE je.id = selected_jobs.id\n AND je.state = 'pending'\n RETURNING je.id, selected_jobs.data_json, je.attempt_index, je.queue_id\n ),\n min_wait AS (\n SELECT MIN(execute_at) AS next_due_at\n FROM job_executions\n WHERE state = 'pending'\n AND job_type = ANY($4::text[] || $8::text[])\n AND execute_at > $2::timestamptz\n ),\n excluded_due AS (\n SELECT EXISTS (\n SELECT 1\n FROM UNNEST($8::text[]) AS et(job_type)\n CROSS JOIN LATERAL (\n SELECT 1 AS hit\n FROM job_executions je\n WHERE je.state = 'pending'\n AND je.job_type = et.job_type\n AND je.execute_at <= $2::timestamptz\n LIMIT 1\n ) probe\n ) AS excluded_due\n ),\n window_counts AS (\n SELECT job_type, COUNT(*) AS cnt FROM window_rows GROUP BY job_type\n ),\n poll_status AS (\n SELECT ((SELECT COUNT(*) FROM locked) >= $1\n OR (EXISTS (\n SELECT 1 FROM window_counts wc\n JOIN limits t ON t.job_type = wc.job_type\n WHERE wc.cnt >= t.type_window_limit\n )\n AND (SELECT COUNT(*) FROM ordered_candidates) > 0)) AS may_have_more\n )\n SELECT * FROM (\n SELECT\n u.id AS \"id?: JobId\",\n u.data_json AS \"data_json?: JsonValue\",\n u.attempt_index AS \"attempt_index?\",\n u.queue_id AS \"queue_id?\",\n NULL::TIMESTAMPTZ AS \"next_due_at?\",\n ps.may_have_more AS \"may_have_more!\",\n ed.excluded_due AS \"excluded_due!\"\n FROM updated u, poll_status ps, excluded_due ed\n UNION ALL\n SELECT\n NULL::UUID AS \"id?: JobId\",\n NULL::JSONB AS \"data_json?: JsonValue\",\n NULL::INT AS \"attempt_index?\",\n NULL::VARCHAR AS \"queue_id?\",\n mw.next_due_at AS \"next_due_at?\",\n ps.may_have_more AS \"may_have_more!\",\n ed.excluded_due AS \"excluded_due!\"\n FROM min_wait mw, poll_status ps, excluded_due ed\n ) AS result\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id?: JobId", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "data_json?: JsonValue", + "type_info": "Jsonb" + }, + { + "ordinal": 2, + "name": "attempt_index?", + "type_info": "Int4" + }, + { + "ordinal": 3, + "name": "queue_id?", + "type_info": "Varchar" + }, + { + "ordinal": 4, + "name": "next_due_at?", + "type_info": "Timestamptz" + }, + { + "ordinal": 5, + "name": "may_have_more!", + "type_info": "Bool" + }, + { + "ordinal": 6, + "name": "excluded_due!", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Int4", + "Timestamptz", + "Uuid", + "TextArray", + "Timestamptz", + "Int4Array", + "Int4", + "TextArray" + ] + }, + "nullable": [ + null, + null, + null, + null, + null, + null, + null + ] + }, + "hash": "606f573c3c6cd3898ce64ad6fe3bfb9b11aab3cb86b23465eaa41cee0f311079" +} diff --git a/cala-ledger/.sqlx/query-61f22dc336d52d7f3898e28088df3d08e6413d2e1dfd92af5d22f9ae9dfc8e82.json b/cala-ledger/.sqlx/query-61f22dc336d52d7f3898e28088df3d08e6413d2e1dfd92af5d22f9ae9dfc8e82.json deleted file mode 100644 index 9edcb6b06..000000000 --- a/cala-ledger/.sqlx/query-61f22dc336d52d7f3898e28088df3d08e6413d2e1dfd92af5d22f9ae9dfc8e82.json +++ /dev/null @@ -1,17 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE job_executions AS je\n SET state = 'pending', execute_at = u.execute_at,\n attempt_index = u.attempt_index, poller_instance_id = NULL\n FROM UNNEST($1::uuid[], $2::timestamptz[], $3::int4[])\n AS u(id, execute_at, attempt_index)\n WHERE je.id = u.id AND je.poller_instance_id = $4\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "UuidArray", - "TimestamptzArray", - "Int4Array", - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "61f22dc336d52d7f3898e28088df3d08e6413d2e1dfd92af5d22f9ae9dfc8e82" -} diff --git a/cala-ledger/.sqlx/query-828c03ee80336e13e64238f149cf80d842631584afb5f5ae32560355b246fce4.json b/cala-ledger/.sqlx/query-6f028dbcb7d30f9dc4a744936f6198082e2b9c48b06347736dd0c4e8b7649f91.json similarity index 80% rename from cala-ledger/.sqlx/query-828c03ee80336e13e64238f149cf80d842631584afb5f5ae32560355b246fce4.json rename to cala-ledger/.sqlx/query-6f028dbcb7d30f9dc4a744936f6198082e2b9c48b06347736dd0c4e8b7649f91.json index a0da51c07..a6187b293 100644 --- a/cala-ledger/.sqlx/query-828c03ee80336e13e64238f149cf80d842631584afb5f5ae32560355b246fce4.json +++ b/cala-ledger/.sqlx/query-6f028dbcb7d30f9dc4a744936f6198082e2b9c48b06347736dd0c4e8b7649f91.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n INSERT INTO cala_persistent_outbox_events (sequence)\n SELECT unnest($1::bigint[]) AS sequence\n ON CONFLICT (sequence) DO NOTHING\n RETURNING id, sequence AS \"sequence!: i64\", payload, tracing_context, recorded_at", + "query": "\n INSERT INTO cala_persistent_outbox_events (sequence)\n SELECT unnest($1::bigint[]) AS sequence\n ON CONFLICT (sequence) DO NOTHING\n RETURNING id, sequence AS \"sequence!: i64\", payload, tracing_context, recorded_at, commit_xid AS \"commit_xid!\"", "describe": { "columns": [ { @@ -27,6 +27,11 @@ "ordinal": 4, "name": "recorded_at", "type_info": "Timestamptz" + }, + { + "ordinal": 5, + "name": "commit_xid!", + "type_info": "Int8" } ], "parameters": { @@ -39,8 +44,9 @@ false, true, true, + false, false ] }, - "hash": "828c03ee80336e13e64238f149cf80d842631584afb5f5ae32560355b246fce4" + "hash": "6f028dbcb7d30f9dc4a744936f6198082e2b9c48b06347736dd0c4e8b7649f91" } diff --git a/cala-ledger/.sqlx/query-80bfbce31b7d3fb529766b245efa5f6bf7a8da2bee286f5851c3c5470ef80584.json b/cala-ledger/.sqlx/query-80bfbce31b7d3fb529766b245efa5f6bf7a8da2bee286f5851c3c5470ef80584.json new file mode 100644 index 000000000..c5209c953 --- /dev/null +++ b/cala-ledger/.sqlx/query-80bfbce31b7d3fb529766b245efa5f6bf7a8da2bee286f5851c3c5470ef80584.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH to_reschedule AS MATERIALIZED (\n SELECT je.id, u.execute_at\n FROM job_executions je\n JOIN UNNEST($1::uuid[], $2::timestamptz[]) AS u(id, execute_at)\n ON je.id = u.id\n WHERE je.poller_instance_id = $3\n ORDER BY je.queue_id, je.id\n FOR UPDATE\n )\n UPDATE job_executions AS je\n SET state = 'pending', execute_at = t.execute_at, attempt_index = 1,\n poller_instance_id = NULL\n FROM to_reschedule t\n WHERE je.id = t.id\n RETURNING je.id AS \"id!: JobId\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "UuidArray", + "TimestamptzArray", + "Uuid" + ] + }, + "nullable": [ + false + ] + }, + "hash": "80bfbce31b7d3fb529766b245efa5f6bf7a8da2bee286f5851c3c5470ef80584" +} diff --git a/cala-ledger/.sqlx/query-80ea63f5ddd05bbf7519cf6c46222a07e1c0eebdc0262b4dea9e63fe90468e0f.json b/cala-ledger/.sqlx/query-80ea63f5ddd05bbf7519cf6c46222a07e1c0eebdc0262b4dea9e63fe90468e0f.json new file mode 100644 index 000000000..91483488d --- /dev/null +++ b/cala-ledger/.sqlx/query-80ea63f5ddd05bbf7519cf6c46222a07e1c0eebdc0262b4dea9e63fe90468e0f.json @@ -0,0 +1,17 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH input AS (\n SELECT * FROM UNNEST($1::uuid[], $3::text[]) AS t(id, unique_key)\n ), pred AS (\n SELECT i.id AS new_id, p.pred_id\n FROM input i\n JOIN LATERAL (\n SELECT j.id AS pred_id\n FROM jobs j\n WHERE j.job_type = $2 AND j.unique_key = i.unique_key\n AND j.id != i.id\n ORDER BY j.created_at DESC, j.id DESC\n LIMIT 1\n ) p ON TRUE\n ), seeded AS (\n INSERT INTO job_execution_states (id, execution_state_json)\n SELECT pred.new_id, s.execution_state_json\n FROM pred\n JOIN job_execution_states s ON s.id = pred.pred_id\n WHERE $4::boolean\n RETURNING id\n )\n DELETE FROM job_execution_states s\n USING pred\n WHERE s.id = pred.pred_id\n ", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "UuidArray", + "Text", + "TextArray", + "Bool" + ] + }, + "nullable": [] + }, + "hash": "80ea63f5ddd05bbf7519cf6c46222a07e1c0eebdc0262b4dea9e63fe90468e0f" +} diff --git a/cala-ledger/.sqlx/query-91c71c304b86ba87df527714d44fbedd8a29b74b0c6e88da16fc0fbdaa5a3ff9.json b/cala-ledger/.sqlx/query-91c71c304b86ba87df527714d44fbedd8a29b74b0c6e88da16fc0fbdaa5a3ff9.json new file mode 100644 index 000000000..7fe6783c6 --- /dev/null +++ b/cala-ledger/.sqlx/query-91c71c304b86ba87df527714d44fbedd8a29b74b0c6e88da16fc0fbdaa5a3ff9.json @@ -0,0 +1,26 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT last_commit_seq AS \"last_commit_seq!: i64\",\n logged_through_sequence AS \"logged_through_sequence!: i64\"\n FROM cala_persistent_outbox_commit_log_state\n WHERE singleton", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "last_commit_seq!: i64", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "logged_through_sequence!: i64", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + false, + false + ] + }, + "hash": "91c71c304b86ba87df527714d44fbedd8a29b74b0c6e88da16fc0fbdaa5a3ff9" +} diff --git a/cala-ledger/.sqlx/query-932da1ca0fcd83172902956efaaa7dfed872f1aa8be8c117d0c1fbb02f4fa4aa.json b/cala-ledger/.sqlx/query-932da1ca0fcd83172902956efaaa7dfed872f1aa8be8c117d0c1fbb02f4fa4aa.json deleted file mode 100644 index 64ad15491..000000000 --- a/cala-ledger/.sqlx/query-932da1ca0fcd83172902956efaaa7dfed872f1aa8be8c117d0c1fbb02f4fa4aa.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE job_executions\n SET alive_at = $1\n WHERE poller_instance_id = $2\n AND state = 'running'\n AND id = ANY($3)\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Timestamptz", - "Uuid", - "UuidArray" - ] - }, - "nullable": [] - }, - "hash": "932da1ca0fcd83172902956efaaa7dfed872f1aa8be8c117d0c1fbb02f4fa4aa" -} diff --git a/cala-ledger/.sqlx/query-00b6ca945b532fe362dbd8eac3a08e168ebc55089279f47f29d8e9f78ce8fda5.json b/cala-ledger/.sqlx/query-9463b386cdfbeaa41e6a0e5b8b95ad6af995efc7f5c314d842740e903ee8fc9d.json similarity index 66% rename from cala-ledger/.sqlx/query-00b6ca945b532fe362dbd8eac3a08e168ebc55089279f47f29d8e9f78ce8fda5.json rename to cala-ledger/.sqlx/query-9463b386cdfbeaa41e6a0e5b8b95ad6af995efc7f5c314d842740e903ee8fc9d.json index 073f1c711..f75b3e6be 100644 --- a/cala-ledger/.sqlx/query-00b6ca945b532fe362dbd8eac3a08e168ebc55089279f47f29d8e9f78ce8fda5.json +++ b/cala-ledger/.sqlx/query-9463b386cdfbeaa41e6a0e5b8b95ad6af995efc7f5c314d842740e903ee8fc9d.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "WITH new_events AS (\n INSERT INTO cala_persistent_outbox_events (payload, tracing_context, recorded_at)\n SELECT unnest($1::jsonb[]) AS payload, $2::jsonb AS tracing_context, COALESCE($3::timestamptz, NOW()) AS recorded_at\n RETURNING id, sequence, recorded_at\n )\n SELECT ne.id AS \"id!\", ne.sequence AS \"sequence!\", ne.recorded_at AS \"recorded_at!\"\n FROM new_events ne\n ORDER BY ne.sequence", + "query": "WITH new_events AS (\n INSERT INTO cala_persistent_outbox_events (payload, tracing_context, recorded_at)\n SELECT unnest($1::jsonb[]) AS payload, $2::jsonb AS tracing_context, COALESCE($3::timestamptz, NOW()) AS recorded_at\n RETURNING id, sequence, recorded_at, commit_xid\n )\n SELECT ne.id AS \"id!\", ne.sequence AS \"sequence!\", ne.recorded_at AS \"recorded_at!\", ne.commit_xid AS \"commit_xid!\"\n FROM new_events ne\n ORDER BY ne.sequence", "describe": { "columns": [ { @@ -17,6 +17,11 @@ "ordinal": 2, "name": "recorded_at!", "type_info": "Timestamptz" + }, + { + "ordinal": 3, + "name": "commit_xid!", + "type_info": "Int8" } ], "parameters": { @@ -27,10 +32,11 @@ ] }, "nullable": [ + false, false, false, false ] }, - "hash": "00b6ca945b532fe362dbd8eac3a08e168ebc55089279f47f29d8e9f78ce8fda5" + "hash": "9463b386cdfbeaa41e6a0e5b8b95ad6af995efc7f5c314d842740e903ee8fc9d" } diff --git a/cala-ledger/.sqlx/query-9ab25bea85f68a606013136bf1d70be5173a6b992936215300ad0737c46bfeec.json b/cala-ledger/.sqlx/query-9ab25bea85f68a606013136bf1d70be5173a6b992936215300ad0737c46bfeec.json deleted file mode 100644 index 2838a67d9..000000000 --- a/cala-ledger/.sqlx/query-9ab25bea85f68a606013136bf1d70be5173a6b992936215300ad0737c46bfeec.json +++ /dev/null @@ -1,17 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE job_executions\n SET state = 'pending', execute_at = $2, attempt_index = $3, poller_instance_id = NULL\n WHERE id = $1 AND poller_instance_id = $4\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Uuid", - "Timestamptz", - "Int4", - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "9ab25bea85f68a606013136bf1d70be5173a6b992936215300ad0737c46bfeec" -} diff --git a/cala-ledger/.sqlx/query-9e4fa6286210c1c103985279291024c3a9fa3ea73e81f7561de13404ebe9a5d8.json b/cala-ledger/.sqlx/query-9e4fa6286210c1c103985279291024c3a9fa3ea73e81f7561de13404ebe9a5d8.json new file mode 100644 index 000000000..174c94e6a --- /dev/null +++ b/cala-ledger/.sqlx/query-9e4fa6286210c1c103985279291024c3a9fa3ea73e81f7561de13404ebe9a5d8.json @@ -0,0 +1,30 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH deleted AS (\n DELETE FROM job_executions\n WHERE id = $1 AND poller_instance_id = $2\n RETURNING id, queue_id\n ), cleanup AS (\n DELETE FROM job_execution_states s USING deleted d\n WHERE s.id = d.id AND NOT $3::boolean\n )\n SELECT id AS \"id!: JobId\", queue_id AS \"queue_id?\"\n FROM deleted\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "queue_id?", + "type_info": "Varchar" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Uuid", + "Bool" + ] + }, + "nullable": [ + false, + true + ] + }, + "hash": "9e4fa6286210c1c103985279291024c3a9fa3ea73e81f7561de13404ebe9a5d8" +} diff --git a/cala-ledger/.sqlx/query-992aef4c4cb3b91da4d7cb0572f1fbf7447c3d4fd92d0cdf83541a58afd8e570.json b/cala-ledger/.sqlx/query-aca7a9b6f0dd28404dd43921b767dfc3b514858c2d6e83d69e328b9cf57e56df.json similarity index 70% rename from cala-ledger/.sqlx/query-992aef4c4cb3b91da4d7cb0572f1fbf7447c3d4fd92d0cdf83541a58afd8e570.json rename to cala-ledger/.sqlx/query-aca7a9b6f0dd28404dd43921b767dfc3b514858c2d6e83d69e328b9cf57e56df.json index 95b1f39ba..a7e4863bb 100644 --- a/cala-ledger/.sqlx/query-992aef4c4cb3b91da4d7cb0572f1fbf7447c3d4fd92d0cdf83541a58afd8e570.json +++ b/cala-ledger/.sqlx/query-aca7a9b6f0dd28404dd43921b767dfc3b514858c2d6e83d69e328b9cf57e56df.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH candidates AS (\n SELECT je.id, je.queue_id, je.execute_at\n FROM job_executions je\n WHERE je.id = ANY($1) AND je.state = 'pending' AND je.queue_id IS NOT NULL\n ), swaps AS (\n SELECT c.id AS pending_id, sib.id AS parked_id\n FROM candidates c\n CROSS JOIN LATERAL (\n SELECT id, execute_at FROM job_executions\n WHERE state = 'parked' AND queue_id = c.queue_id\n ORDER BY execute_at, id\n LIMIT 1\n ) sib\n WHERE (sib.execute_at, sib.id) < (c.execute_at, c.id)\n ), locked AS MATERIALIZED (\n -- Take BOTH sides of every swap in (queue_id, id) order, in\n -- one pass, before either UPDATE below runs. `demote` reads\n -- from this (not from `swaps`) purely to force that\n -- dependency. See the doc comment for why the order is\n -- load-bearing; the strength is exactly what the two writes\n -- below would take anyway -- `state` is no index's key\n -- column -- so this changes lock ORDER only, never what\n -- conflicts with what. A swap's two rows always share one\n -- `queue_id` (the sibling lookup is `queue_id = c.queue_id`),\n -- so ordering by it groups each swap's pair together.\n SELECT je.id FROM job_executions je\n WHERE je.id IN (\n SELECT pending_id FROM swaps UNION SELECT parked_id FROM swaps\n )\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n ), demote AS (\n UPDATE job_executions SET state = 'parked'\n WHERE id IN (\n SELECT s.pending_id FROM swaps s JOIN locked l ON l.id = s.pending_id\n )\n RETURNING id\n )\n -- The promote UPDATE reads FROM `demote` (not `swaps`) so Postgres\n -- has a real data dependency forcing `demote` to run to completion\n -- first. Without it, this is two independent writes to the same\n -- table within one statement with no ordering guarantee between\n -- them, which can transiently make two rows active for one queue\n -- within the statement's own execution and violate\n -- `idx_job_executions_queue_active`.\n UPDATE job_executions je SET state = 'pending'\n FROM swaps s\n JOIN demote d ON d.id = s.pending_id\n WHERE je.id = s.parked_id\n RETURNING je.job_type, je.execute_at AS \"execute_at!\"\n ", + "query": "\n WITH candidates AS (\n SELECT je.id, je.queue_id, je.execute_at\n FROM job_executions je\n WHERE je.id = ANY($1) AND je.state = 'pending' AND je.queue_id IS NOT NULL\n ), swaps AS (\n SELECT c.id AS pending_id, sib.id AS parked_id\n FROM candidates c\n CROSS JOIN LATERAL (\n SELECT id, execute_at FROM job_executions\n WHERE state = 'parked' AND queue_id = c.queue_id\n ORDER BY execute_at, id\n LIMIT 1\n ) sib\n WHERE (sib.execute_at, sib.id) < (c.execute_at, c.id)\n ), locked AS MATERIALIZED (\n -- Take BOTH sides of every swap in (queue_id, id) order, in\n -- one pass, before either UPDATE below runs. `demote` reads\n -- from this (not from `swaps`) purely to force that\n -- dependency. See the doc comment for why the order is\n -- load-bearing; the strength is exactly what the two writes\n -- below would take anyway -- `state` is no index's key\n -- column -- so this changes lock ORDER only, never what\n -- conflicts with what. A swap's two rows always share one\n -- `queue_id` (the sibling lookup is `queue_id = c.queue_id`),\n -- so ordering by it groups each swap's pair together.\n SELECT je.id FROM job_executions je\n WHERE je.id IN (\n SELECT pending_id FROM swaps UNION SELECT parked_id FROM swaps\n )\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n ), demote AS (\n UPDATE job_executions SET state = 'parked'\n WHERE id IN (\n SELECT s.pending_id FROM swaps s JOIN locked l ON l.id = s.pending_id\n )\n AND state = 'pending'\n RETURNING id\n )\n -- The promote UPDATE reads FROM `demote` (not `swaps`) so Postgres\n -- has a real data dependency forcing `demote` to run to completion\n -- first. Without it, this is two independent writes to the same\n -- table within one statement with no ordering guarantee between\n -- them, which can transiently make two rows active for one queue\n -- within the statement's own execution and violate\n -- `idx_job_executions_queue_active`.\n UPDATE job_executions je SET state = 'pending'\n FROM swaps s\n JOIN demote d ON d.id = s.pending_id\n WHERE je.id = s.parked_id AND je.state = 'parked'\n RETURNING je.job_type, je.execute_at AS \"execute_at?\"\n ", "describe": { "columns": [ { @@ -10,7 +10,7 @@ }, { "ordinal": 1, - "name": "execute_at!", + "name": "execute_at?", "type_info": "Timestamptz" } ], @@ -24,5 +24,5 @@ true ] }, - "hash": "992aef4c4cb3b91da4d7cb0572f1fbf7447c3d4fd92d0cdf83541a58afd8e570" + "hash": "aca7a9b6f0dd28404dd43921b767dfc3b514858c2d6e83d69e328b9cf57e56df" } diff --git a/cala-ledger/.sqlx/query-bbbfd153dfa0453dc65549c83b34340b67b1347ec85793928f2658d1e7c63897.json b/cala-ledger/.sqlx/query-ad3164a1a7689fced3d6108d0eebad69d3c3abb79c1fe9127c13fc7bc3e7c303.json similarity index 69% rename from cala-ledger/.sqlx/query-bbbfd153dfa0453dc65549c83b34340b67b1347ec85793928f2658d1e7c63897.json rename to cala-ledger/.sqlx/query-ad3164a1a7689fced3d6108d0eebad69d3c3abb79c1fe9127c13fc7bc3e7c303.json index 3ec5165d8..97903d7e2 100644 --- a/cala-ledger/.sqlx/query-bbbfd153dfa0453dc65549c83b34340b67b1347ec85793928f2658d1e7c63897.json +++ b/cala-ledger/.sqlx/query-ad3164a1a7689fced3d6108d0eebad69d3c3abb79c1fe9127c13fc7bc3e7c303.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n SELECT sequence AS \"sequence!: i64\", id AS \"id!\", payload, tracing_context, recorded_at AS \"recorded_at!\"\n FROM cala_persistent_outbox_events\n WHERE sequence > $1\n AND sequence <= $1 + $2\n ORDER BY sequence ASC\n LIMIT $2", + "query": "\n SELECT sequence AS \"sequence!: i64\", id AS \"id!\", payload, tracing_context, recorded_at AS \"recorded_at!\", commit_xid\n FROM cala_persistent_outbox_events\n WHERE sequence > $1\n AND sequence <= $1 + $2\n ORDER BY sequence ASC\n LIMIT $2", "describe": { "columns": [ { @@ -27,6 +27,11 @@ "ordinal": 4, "name": "recorded_at!", "type_info": "Timestamptz" + }, + { + "ordinal": 5, + "name": "commit_xid", + "type_info": "Int8" } ], "parameters": { @@ -40,8 +45,9 @@ false, true, true, + false, false ] }, - "hash": "bbbfd153dfa0453dc65549c83b34340b67b1347ec85793928f2658d1e7c63897" + "hash": "ad3164a1a7689fced3d6108d0eebad69d3c3abb79c1fe9127c13fc7bc3e7c303" } diff --git a/cala-ledger/.sqlx/query-b01670f025421b5761f4907d81e1072abface895a3d1631aefa66773b7b7759c.json b/cala-ledger/.sqlx/query-b01670f025421b5761f4907d81e1072abface895a3d1631aefa66773b7b7759c.json deleted file mode 100644 index 78c0ac7e1..000000000 --- a/cala-ledger/.sqlx/query-b01670f025421b5761f4907d81e1072abface895a3d1631aefa66773b7b7759c.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n -- `(queue_id, id)`-ordered, like every other multi-row waiting locker\n -- of this table: shutdown releases an unbounded set of rows spread\n -- across arbitrary queues while peers are still spawning, completing\n -- and sweeping against the same ones.\n WITH locked AS MATERIALIZED (\n SELECT je.id FROM job_executions je\n WHERE je.poller_instance_id = $2 AND je.state = 'running'\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je\n SET state = 'pending',\n execute_at = $1,\n poller_instance_id = NULL\n FROM locked l WHERE je.id = l.id\n RETURNING je.id as \"id!: JobId\", je.attempt_index\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!: JobId", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "attempt_index", - "type_info": "Int4" - } - ], - "parameters": { - "Left": [ - "Timestamptz", - "Uuid" - ] - }, - "nullable": [ - false, - false - ] - }, - "hash": "b01670f025421b5761f4907d81e1072abface895a3d1631aefa66773b7b7759c" -} diff --git a/cala-ledger/.sqlx/query-b03d25d915a2be8e8270b8ea4d54d4f0e0c0613ed3eefa50a6295d02c3dbc907.json b/cala-ledger/.sqlx/query-b03d25d915a2be8e8270b8ea4d54d4f0e0c0613ed3eefa50a6295d02c3dbc907.json deleted file mode 100644 index 03817edc4..000000000 --- a/cala-ledger/.sqlx/query-b03d25d915a2be8e8270b8ea4d54d4f0e0c0613ed3eefa50a6295d02c3dbc907.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH to_delete AS MATERIALIZED (\n SELECT id FROM job_executions\n WHERE id = ANY($1) AND poller_instance_id = $2\n ORDER BY queue_id, id\n FOR UPDATE\n ), deleted AS (\n DELETE FROM job_executions je USING to_delete t WHERE je.id = t.id\n RETURNING je.id, je.queue_id\n ), cleanup AS (\n DELETE FROM job_execution_states s USING deleted d WHERE s.id = d.id\n )\n SELECT id AS \"id!: JobId\", queue_id AS \"queue_id?\"\n FROM deleted\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!: JobId", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "queue_id?", - "type_info": "Varchar" - } - ], - "parameters": { - "Left": [ - "UuidArray", - "Uuid" - ] - }, - "nullable": [ - false, - true - ] - }, - "hash": "b03d25d915a2be8e8270b8ea4d54d4f0e0c0613ed3eefa50a6295d02c3dbc907" -} diff --git a/cala-ledger/.sqlx/query-bd553029c1b59d0df10b2e681344152ca649b225518216413ad8aae57d1098c1.json b/cala-ledger/.sqlx/query-bd553029c1b59d0df10b2e681344152ca649b225518216413ad8aae57d1098c1.json deleted file mode 100644 index ccfcc44de..000000000 --- a/cala-ledger/.sqlx/query-bd553029c1b59d0df10b2e681344152ca649b225518216413ad8aae57d1098c1.json +++ /dev/null @@ -1,17 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH seeded AS (\n INSERT INTO job_execution_states (id, execution_state_json)\n SELECT $1, s.execution_state_json\n FROM job_execution_states s\n JOIN jobs j ON j.id = s.id\n WHERE j.job_type = $2 AND j.unique_key = $3 AND j.id != $1\n AND $4::boolean\n ORDER BY j.created_at DESC, j.id DESC\n LIMIT 1\n RETURNING id\n )\n DELETE FROM job_execution_states s\n USING jobs j\n WHERE s.id = j.id AND j.job_type = $2 AND j.unique_key = $3 AND j.id != $1\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Uuid", - "Text", - "Text", - "Bool" - ] - }, - "nullable": [] - }, - "hash": "bd553029c1b59d0df10b2e681344152ca649b225518216413ad8aae57d1098c1" -} diff --git a/cala-ledger/.sqlx/query-7089e2251160ce1d7266f49467204e83806992fa6f11588c25f26f5b90f23af9.json b/cala-ledger/.sqlx/query-bdfa672fbc2595d6eed0c82d542e1e564551d49b17045ad7a3c5c4006c6e4b08.json similarity index 50% rename from cala-ledger/.sqlx/query-7089e2251160ce1d7266f49467204e83806992fa6f11588c25f26f5b90f23af9.json rename to cala-ledger/.sqlx/query-bdfa672fbc2595d6eed0c82d542e1e564551d49b17045ad7a3c5c4006c6e4b08.json index f245fbb0f..2e43bc94f 100644 --- a/cala-ledger/.sqlx/query-7089e2251160ce1d7266f49467204e83806992fa6f11588c25f26f5b90f23af9.json +++ b/cala-ledger/.sqlx/query-bdfa672fbc2595d6eed0c82d542e1e564551d49b17045ad7a3c5c4006c6e4b08.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "WITH new_events AS (\n INSERT INTO cala_persistent_outbox_events (payload, tracing_context, recorded_at)\n SELECT unnest($1::jsonb[]) AS payload, $2::jsonb AS tracing_context, COALESCE($3::timestamptz, NOW()) AS recorded_at\n RETURNING id, sequence, recorded_at\n ),\n notified AS (\n SELECT pg_notify(\n 'cala_persistent_outbox_events',\n json_build_object('min_sequence', MIN(sequence), 'max_sequence', MAX(sequence))::TEXT\n )\n FROM new_events\n HAVING COUNT(*) > 0\n )\n SELECT ne.id AS \"id!\", ne.sequence AS \"sequence!\", ne.recorded_at AS \"recorded_at!\"\n FROM new_events ne\n LEFT JOIN notified ON TRUE\n ORDER BY ne.sequence", + "query": "WITH new_events AS (\n INSERT INTO cala_persistent_outbox_events (payload, tracing_context, recorded_at)\n SELECT unnest($1::jsonb[]) AS payload, $2::jsonb AS tracing_context, COALESCE($3::timestamptz, NOW()) AS recorded_at\n RETURNING id, sequence, recorded_at, commit_xid\n ),\n notified AS (\n SELECT pg_notify(\n 'cala_persistent_outbox_events',\n json_build_object('min_sequence', MIN(sequence), 'max_sequence', MAX(sequence))::TEXT\n )\n FROM new_events\n HAVING COUNT(*) > 0\n )\n SELECT ne.id AS \"id!\", ne.sequence AS \"sequence!\", ne.recorded_at AS \"recorded_at!\", ne.commit_xid AS \"commit_xid!\"\n FROM new_events ne\n LEFT JOIN notified ON TRUE\n ORDER BY ne.sequence", "describe": { "columns": [ { @@ -17,6 +17,11 @@ "ordinal": 2, "name": "recorded_at!", "type_info": "Timestamptz" + }, + { + "ordinal": 3, + "name": "commit_xid!", + "type_info": "Int8" } ], "parameters": { @@ -27,10 +32,11 @@ ] }, "nullable": [ + false, false, false, false ] }, - "hash": "7089e2251160ce1d7266f49467204e83806992fa6f11588c25f26f5b90f23af9" + "hash": "bdfa672fbc2595d6eed0c82d542e1e564551d49b17045ad7a3c5c4006c6e4b08" } diff --git a/cala-ledger/.sqlx/query-bfc3d07b931304ab6d1a4a290ff2da5bac04d17ecb2ca6590d273514a28e4156.json b/cala-ledger/.sqlx/query-bfc3d07b931304ab6d1a4a290ff2da5bac04d17ecb2ca6590d273514a28e4156.json new file mode 100644 index 000000000..049e86fb0 --- /dev/null +++ b/cala-ledger/.sqlx/query-bfc3d07b931304ab6d1a4a290ff2da5bac04d17ecb2ca6590d273514a28e4156.json @@ -0,0 +1,29 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH locked AS MATERIALIZED (\n SELECT je.id FROM job_executions je\n WHERE je.poller_instance_id = $2 AND je.state = 'running'\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je\n SET state = 'pending',\n execute_at = $1,\n poller_instance_id = NULL\n FROM locked l WHERE je.id = l.id\n RETURNING je.id as \"id!: JobId\", je.attempt_index\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "attempt_index", + "type_info": "Int4" + } + ], + "parameters": { + "Left": [ + "Timestamptz", + "Uuid" + ] + }, + "nullable": [ + false, + false + ] + }, + "hash": "bfc3d07b931304ab6d1a4a290ff2da5bac04d17ecb2ca6590d273514a28e4156" +} diff --git a/cala-ledger/.sqlx/query-cd605cbf07c5268c51f02f05d88f1685ee85f8a00888655cbb92ff35c5e3ba24.json b/cala-ledger/.sqlx/query-cd605cbf07c5268c51f02f05d88f1685ee85f8a00888655cbb92ff35c5e3ba24.json new file mode 100644 index 000000000..7c0e9344a --- /dev/null +++ b/cala-ledger/.sqlx/query-cd605cbf07c5268c51f02f05d88f1685ee85f8a00888655cbb92ff35c5e3ba24.json @@ -0,0 +1,25 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH to_retry AS MATERIALIZED (\n SELECT je.id, u.execute_at, u.attempt_index\n FROM job_executions je\n JOIN UNNEST($1::uuid[], $2::timestamptz[], $3::int4[])\n AS u(id, execute_at, attempt_index)\n ON je.id = u.id\n WHERE je.poller_instance_id = $4\n ORDER BY je.queue_id, je.id\n FOR UPDATE\n )\n UPDATE job_executions AS je\n SET state = 'pending', execute_at = t.execute_at,\n attempt_index = t.attempt_index, poller_instance_id = NULL\n FROM to_retry t\n WHERE je.id = t.id\n RETURNING je.id AS \"id!: JobId\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "UuidArray", + "TimestamptzArray", + "Int4Array", + "Uuid" + ] + }, + "nullable": [ + false + ] + }, + "hash": "cd605cbf07c5268c51f02f05d88f1685ee85f8a00888655cbb92ff35c5e3ba24" +} diff --git a/cala-ledger/.sqlx/query-d5a15e7d6001914866c9fd9149bd01169264da82273629d14afe509d1b6005e0.json b/cala-ledger/.sqlx/query-d5a15e7d6001914866c9fd9149bd01169264da82273629d14afe509d1b6005e0.json new file mode 100644 index 000000000..f7d90f527 --- /dev/null +++ b/cala-ledger/.sqlx/query-d5a15e7d6001914866c9fd9149bd01169264da82273629d14afe509d1b6005e0.json @@ -0,0 +1,65 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT l.commit_seq AS \"commit_seq!: i64\", l.group_last AS \"group_last!\",\n e.id AS \"id!\", e.sequence AS \"sequence!: i64\", e.payload,\n e.tracing_context, e.recorded_at AS \"recorded_at!\", e.commit_xid\n FROM cala_persistent_outbox_commit_log l\n JOIN cala_persistent_outbox_events e ON e.sequence = l.sequence\n WHERE l.commit_seq > $1 AND l.commit_seq <= $1 + $2\n ORDER BY l.commit_seq", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "commit_seq!: i64", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "group_last!", + "type_info": "Bool" + }, + { + "ordinal": 2, + "name": "id!", + "type_info": "Uuid" + }, + { + "ordinal": 3, + "name": "sequence!: i64", + "type_info": "Int8" + }, + { + "ordinal": 4, + "name": "payload", + "type_info": "Jsonb" + }, + { + "ordinal": 5, + "name": "tracing_context", + "type_info": "Jsonb" + }, + { + "ordinal": 6, + "name": "recorded_at!", + "type_info": "Timestamptz" + }, + { + "ordinal": 7, + "name": "commit_xid", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Int8", + "Int8" + ] + }, + "nullable": [ + false, + false, + false, + false, + true, + true, + false, + false + ] + }, + "hash": "d5a15e7d6001914866c9fd9149bd01169264da82273629d14afe509d1b6005e0" +} diff --git a/cala-ledger/.sqlx/query-d6964f14852c2b5d45efd02bde1fc59ba9158c3e1c361dd5841b4ad435fb9428.json b/cala-ledger/.sqlx/query-d6964f14852c2b5d45efd02bde1fc59ba9158c3e1c361dd5841b4ad435fb9428.json new file mode 100644 index 000000000..69f70d0e5 --- /dev/null +++ b/cala-ledger/.sqlx/query-d6964f14852c2b5d45efd02bde1fc59ba9158c3e1c361dd5841b4ad435fb9428.json @@ -0,0 +1,30 @@ +{ + "db_name": "PostgreSQL", + "query": "\n UPDATE job_executions je\n SET execute_at = LEAST(je.execute_at, t.target)\n FROM UNNEST($2::text[], $3::timestamptz[]) AS t(unique_key, target)\n WHERE je.job_type = $1\n AND je.unique_key = t.unique_key\n AND je.state = 'pending'\n AND je.attempt_index <= 1\n AND je.execute_at > t.target\n RETURNING je.unique_key AS \"unique_key!\", je.execute_at AS \"execute_at!\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "unique_key!", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "execute_at!", + "type_info": "Timestamptz" + } + ], + "parameters": { + "Left": [ + "Text", + "TextArray", + "TimestamptzArray" + ] + }, + "nullable": [ + true, + true + ] + }, + "hash": "d6964f14852c2b5d45efd02bde1fc59ba9158c3e1c361dd5841b4ad435fb9428" +} diff --git a/cala-ledger/.sqlx/query-da2a3f0d5bb15136393ada4a31886cd4e433c8d84930496e6f5487acc9f2fd93.json b/cala-ledger/.sqlx/query-da2a3f0d5bb15136393ada4a31886cd4e433c8d84930496e6f5487acc9f2fd93.json new file mode 100644 index 000000000..664ff3ffc --- /dev/null +++ b/cala-ledger/.sqlx/query-da2a3f0d5bb15136393ada4a31886cd4e433c8d84930496e6f5487acc9f2fd93.json @@ -0,0 +1,30 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH to_delete AS MATERIALIZED (\n SELECT id FROM job_executions\n WHERE id = ANY($1) AND poller_instance_id = $2\n ORDER BY queue_id, id\n FOR UPDATE\n ), deleted AS (\n DELETE FROM job_executions je USING to_delete t WHERE je.id = t.id\n RETURNING je.id, je.queue_id\n ), cleanup AS (\n DELETE FROM job_execution_states s USING deleted d\n WHERE s.id = d.id AND NOT $3::boolean\n )\n SELECT id AS \"id!: JobId\", queue_id AS \"queue_id?\"\n FROM deleted\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "queue_id?", + "type_info": "Varchar" + } + ], + "parameters": { + "Left": [ + "UuidArray", + "Uuid", + "Bool" + ] + }, + "nullable": [ + false, + true + ] + }, + "hash": "da2a3f0d5bb15136393ada4a31886cd4e433c8d84930496e6f5487acc9f2fd93" +} diff --git a/cala-ledger/.sqlx/query-f1125bc627825540d799741f15703f39401b3fb50755e188e6626e08c89e871c.json b/cala-ledger/.sqlx/query-e16b163bdb81c2afe37ccd396207859510d17971a4923bc4dd4cd95a1512748c.json similarity index 80% rename from cala-ledger/.sqlx/query-f1125bc627825540d799741f15703f39401b3fb50755e188e6626e08c89e871c.json rename to cala-ledger/.sqlx/query-e16b163bdb81c2afe37ccd396207859510d17971a4923bc4dd4cd95a1512748c.json index 53a1226c6..2bbeaa3e5 100644 --- a/cala-ledger/.sqlx/query-f1125bc627825540d799741f15703f39401b3fb50755e188e6626e08c89e871c.json +++ b/cala-ledger/.sqlx/query-e16b163bdb81c2afe37ccd396207859510d17971a4923bc4dd4cd95a1512748c.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "\n WITH locked AS MATERIALIZED (\n SELECT je.id FROM job_executions je\n WHERE je.state = 'running'\n AND je.alive_at < $1::timestamptz\n AND je.job_type = ANY($2)\n AND (je.poller_instance_id IS DISTINCT FROM $4 OR je.id <> ALL($5))\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je\n SET state = 'pending', execute_at = $3, attempt_index = attempt_index + 1, poller_instance_id = NULL\n FROM locked l WHERE je.id = l.id\n RETURNING je.id AS \"id!: JobId\", je.job_type AS \"job_type!: JobType\"\n ", + "query": "\n WITH locked AS MATERIALIZED (\n SELECT je.id FROM job_executions je\n WHERE je.state = 'running'\n AND je.alive_at < $1::timestamptz\n AND je.job_type = ANY($2)\n AND (je.poller_instance_id IS DISTINCT FROM $4 OR je.id <> ALL($5))\n ORDER BY je.queue_id, je.id\n FOR NO KEY UPDATE\n )\n UPDATE job_executions je\n SET state = 'pending', execute_at = $3, attempt_index = attempt_index + 1, poller_instance_id = NULL\n FROM locked l WHERE je.id = l.id\n RETURNING je.id AS \"id!: JobId\", je.job_type AS \"job_type!: JobType\",\n je.alive_at AS \"alive_at!\"\n ", "describe": { "columns": [ { @@ -12,6 +12,11 @@ "ordinal": 1, "name": "job_type!: JobType", "type_info": "Varchar" + }, + { + "ordinal": 2, + "name": "alive_at!", + "type_info": "Timestamptz" } ], "parameters": { @@ -24,9 +29,10 @@ ] }, "nullable": [ + false, false, false ] }, - "hash": "f1125bc627825540d799741f15703f39401b3fb50755e188e6626e08c89e871c" + "hash": "e16b163bdb81c2afe37ccd396207859510d17971a4923bc4dd4cd95a1512748c" } diff --git a/cala-ledger/.sqlx/query-e66100b21243cdecb6046f3748e9c7e90c6648c648d09a3f101edee12de354f2.json b/cala-ledger/.sqlx/query-e66100b21243cdecb6046f3748e9c7e90c6648c648d09a3f101edee12de354f2.json deleted file mode 100644 index 8a0d6b685..000000000 --- a/cala-ledger/.sqlx/query-e66100b21243cdecb6046f3748e9c7e90c6648c648d09a3f101edee12de354f2.json +++ /dev/null @@ -1,58 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n -- Claim admission, head-only. `state = 'pending'` contains ONLY\n -- already-claimable rows: a queue's blocked backlog is `parked`\n -- instead (promoted at completion time -- see\n -- `dispatcher.rs::delete_execution_in_op`), so queued and unqueued\n -- rows share one ordered scan with no anti-join and no per-queue\n -- LATERAL. See PERFORMANCE.md (\"Claim admission\") for the shape this\n -- replaced and the measurements behind this one.\n WITH limits AS (\n -- `type_window_limit` is each type's OWN admission budget for\n -- step 1, not a shared global one: a type's window is bounded by\n -- its own row_limit (never by another type's backlog), which is\n -- what stops one backlogged type from crowding a due row of\n -- another type out of the window entirely. Pinned by\n -- `capped_type_backlog_does_not_starve_another_type` and (for\n -- batched types) `claims_are_capped_by_free_batch_slots`.\n SELECT l.job_type, l.row_limit,\n LEAST(l.row_limit, $1::int4) * $7::int4 AS type_window_limit\n FROM UNNEST($4::text[], $6::int4[]) AS l(job_type, row_limit)\n WHERE l.row_limit > 0\n ),\n window_rows AS (\n -- One LATERAL probe per type, each bounded by ITS OWN budget\n -- ($1 x $7 rows capped at that type's own row_limit), never by\n -- how much is pending: cost is O(budget), flat in backlog --\n -- true only because `idx_job_executions_pending_execute_at`\n -- leads with `job_type`, so this probe is an index descent into\n -- that type's own slice rather than a filter-scan of every\n -- other type's pending rows too (see PERFORMANCE.md, \"Claim\n -- admission\"). Ordering is `(execute_at, id)` within that slice\n -- and the tiebreak is load-bearing for the same reason it\n -- always was -- a total order, so each type's window is a\n -- well-defined prefix rather than an arbitrary cut through a\n -- group of rows sharing a timestamp.\n SELECT d.id, d.execute_at, d.job_type\n FROM limits t\n CROSS JOIN LATERAL (\n SELECT je.id, je.execute_at, je.job_type\n FROM job_executions je\n WHERE je.state = 'pending'\n AND je.job_type = t.job_type\n AND je.execute_at <= $2::timestamptz\n ORDER BY je.execute_at, je.id\n LIMIT t.type_window_limit\n ) d\n ),\n ordered_candidates AS (\n -- Interleave types round-robin: every type's oldest candidate\n -- ranks ahead of any type's second, so the global LIMIT below\n -- cannot be consumed end-to-end by one backlogged type. Within a\n -- rank it is still oldest-first.\n SELECT id, execute_at, job_type,\n ROW_NUMBER() OVER (\n PARTITION BY job_type ORDER BY execute_at\n ) AS type_rn\n FROM window_rows\n ),\n locked AS (\n -- The join to job_executions sits BELOW the LIMIT so it runs\n -- lazily: only rows LockRows actually pulls get probed. The sort\n -- above is a blocking node, so the full candidate set is still\n -- materialised and SKIP LOCKED falls through a contended row\n -- exactly as before -- `headroom` still exists solely to give it\n -- somewhere to fall through to; see the constant's doc comment.\n --\n -- FOR UPDATE OF je: bare FOR UPDATE errors on a nullable join side.\n SELECT je.id, je.attempt_index, c.job_type, c.execute_at\n FROM ordered_candidates c\n JOIN job_executions je ON je.id = c.id\n ORDER BY c.type_rn ASC, c.execute_at ASC\n LIMIT $1\n FOR UPDATE OF je SKIP LOCKED\n ),\n selected_jobs AS (\n -- The budget is enforced HERE, on rows actually held: the window\n -- deliberately over-gathers (see $7) so there is something to\n -- fall through to when a peer holds a row. Rows over a type's cap\n -- are simply not claimed; their locks release at commit.\n -- execution_state_json is joined after the LIMIT, so it is\n -- fetched only for winners.\n SELECT t.id, cp.execution_state_json AS data_json, t.attempt_index\n FROM (\n SELECT l.*,\n ROW_NUMBER() OVER (\n PARTITION BY l.job_type ORDER BY l.execute_at\n ) AS type_rn\n FROM locked l\n ) t\n JOIN limits lim ON lim.job_type = t.job_type\n LEFT JOIN job_execution_states cp ON cp.id = t.id\n WHERE t.type_rn <= lim.row_limit\n ),\n updated AS (\n UPDATE job_executions AS je\n SET state = 'running', alive_at = $5, execute_at = NULL, poller_instance_id = $3\n FROM selected_jobs\n WHERE je.id = selected_jobs.id\n AND je.state = 'pending'\n RETURNING je.id, selected_jobs.data_json, je.attempt_index, je.queue_id\n ),\n min_wait AS (\n SELECT MIN(execute_at) AS next_due_at\n FROM job_executions\n WHERE state = 'pending'\n AND job_type = ANY($4)\n AND execute_at > $2::timestamptz\n ),\n window_counts AS (\n SELECT job_type, COUNT(*) AS cnt FROM window_rows GROUP BY job_type\n ),\n poll_status AS (\n -- Re-poll immediately only when this poll provably left claimable\n -- work behind: it filled its budget, or at least one type's OWN\n -- window came back full while still yielding at least one\n -- pollable candidate overall (rows past that type's window are\n -- unseen and already due). A window that came back short for\n -- every type means every claimable due row was examined, so\n -- `next_due_at` is the honest next deadline -- exact now, not\n -- merely a heuristic: every window row is by construction\n -- already a candidate. Blocked queues are\n -- covered by a wake rather than a spin: `delete_execution_in_op`/\n -- the orphan sweeper report `execution_ready` when a queue's\n -- head is promoted.\n SELECT ((SELECT COUNT(*) FROM locked) >= $1\n OR (EXISTS (\n SELECT 1 FROM window_counts wc\n JOIN limits t ON t.job_type = wc.job_type\n WHERE wc.cnt >= t.type_window_limit\n )\n AND (SELECT COUNT(*) FROM ordered_candidates) > 0)) AS may_have_more\n )\n SELECT * FROM (\n SELECT\n u.id AS \"id?: JobId\",\n u.data_json AS \"data_json?: JsonValue\",\n u.attempt_index AS \"attempt_index?\",\n u.queue_id AS \"queue_id?\",\n NULL::TIMESTAMPTZ AS \"next_due_at?\",\n ps.may_have_more AS \"may_have_more!\"\n FROM updated u, poll_status ps\n UNION ALL\n SELECT\n NULL::UUID AS \"id?: JobId\",\n NULL::JSONB AS \"data_json?: JsonValue\",\n NULL::INT AS \"attempt_index?\",\n NULL::VARCHAR AS \"queue_id?\",\n mw.next_due_at AS \"next_due_at?\",\n ps.may_have_more AS \"may_have_more!\"\n FROM min_wait mw, poll_status ps\n ) AS result\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id?: JobId", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "data_json?: JsonValue", - "type_info": "Jsonb" - }, - { - "ordinal": 2, - "name": "attempt_index?", - "type_info": "Int4" - }, - { - "ordinal": 3, - "name": "queue_id?", - "type_info": "Varchar" - }, - { - "ordinal": 4, - "name": "next_due_at?", - "type_info": "Timestamptz" - }, - { - "ordinal": 5, - "name": "may_have_more!", - "type_info": "Bool" - } - ], - "parameters": { - "Left": [ - "Int4", - "Timestamptz", - "Uuid", - "TextArray", - "Timestamptz", - "Int4Array", - "Int4" - ] - }, - "nullable": [ - null, - null, - null, - null, - null, - null - ] - }, - "hash": "e66100b21243cdecb6046f3748e9c7e90c6648c648d09a3f101edee12de354f2" -} diff --git a/cala-ledger/.sqlx/query-e67fb92be8782f30d872b37b7a8efdaae473009d3dc398aacff73a58eb0e0ad7.json b/cala-ledger/.sqlx/query-e67fb92be8782f30d872b37b7a8efdaae473009d3dc398aacff73a58eb0e0ad7.json deleted file mode 100644 index edf8a0bb8..000000000 --- a/cala-ledger/.sqlx/query-e67fb92be8782f30d872b37b7a8efdaae473009d3dc398aacff73a58eb0e0ad7.json +++ /dev/null @@ -1,16 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n UPDATE job_executions\n SET state = 'pending', execute_at = $2, attempt_index = 1, poller_instance_id = NULL\n WHERE id = $1 AND poller_instance_id = $3\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Uuid", - "Timestamptz", - "Uuid" - ] - }, - "nullable": [] - }, - "hash": "e67fb92be8782f30d872b37b7a8efdaae473009d3dc398aacff73a58eb0e0ad7" -} diff --git a/cala-ledger/.sqlx/query-ef771e1a68cc5b3722dcef68621e15fad82ebfca75ef553f9b3f0c165ccfb55c.json b/cala-ledger/.sqlx/query-ef771e1a68cc5b3722dcef68621e15fad82ebfca75ef553f9b3f0c165ccfb55c.json new file mode 100644 index 000000000..323eb32db --- /dev/null +++ b/cala-ledger/.sqlx/query-ef771e1a68cc5b3722dcef68621e15fad82ebfca75ef553f9b3f0c165ccfb55c.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH to_reschedule AS MATERIALIZED (\n SELECT je.id, u.execute_at\n FROM job_executions je\n JOIN UNNEST($1::uuid[], $2::timestamptz[]) AS u(id, execute_at)\n ON je.id = u.id\n WHERE je.poller_instance_id = $3\n ORDER BY je.queue_id, je.id\n FOR UPDATE\n )\n UPDATE job_executions AS je\n SET state = 'pending', execute_at = t.execute_at, poller_instance_id = NULL\n FROM to_reschedule t\n WHERE je.id = t.id\n RETURNING je.id AS \"id!: JobId\"\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "UuidArray", + "TimestamptzArray", + "Uuid" + ] + }, + "nullable": [ + false + ] + }, + "hash": "ef771e1a68cc5b3722dcef68621e15fad82ebfca75ef553f9b3f0c165ccfb55c" +} diff --git a/cala-ledger/.sqlx/query-f004ea481cacdb54848f6a506295bda267c3c163e2735aac61d4de4dcc8c6b0c.json b/cala-ledger/.sqlx/query-f004ea481cacdb54848f6a506295bda267c3c163e2735aac61d4de4dcc8c6b0c.json new file mode 100644 index 000000000..115ec6c56 --- /dev/null +++ b/cala-ledger/.sqlx/query-f004ea481cacdb54848f6a506295bda267c3c163e2735aac61d4de4dcc8c6b0c.json @@ -0,0 +1,39 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH raw AS (\n SELECT * FROM UNNEST($1::uuid[], $2::text[], $3::text[], $4::timestamptz[], $6::text[])\n AS t(id, job_type, queue_id, execute_at, unique_key)\n ), deduped AS MATERIALIZED (\n SELECT DISTINCT ON (job_type, COALESCE(unique_key, id::text)) *\n FROM raw\n ORDER BY job_type, COALESCE(unique_key, id::text), id\n ), input AS MATERIALIZED (\n SELECT * FROM deduped ORDER BY queue_id, id\n ), ins AS (\n INSERT INTO job_executions\n (id, job_type, queue_id, unique_key, state, attempt_index, execute_at, alive_at, created_at)\n SELECT id, job_type, queue_id, unique_key, 'pending', 1, execute_at,\n COALESCE($5, NOW()), COALESCE($5, NOW())\n FROM input\n ON CONFLICT (queue_id) WHERE state IN ('pending','running') AND queue_id IS NOT NULL\n DO NOTHING\n RETURNING id, queue_id\n ), parked AS (\n INSERT INTO job_executions\n (id, job_type, queue_id, unique_key, state, attempt_index, execute_at, alive_at, created_at)\n SELECT i.id, i.job_type, i.queue_id, i.unique_key, 'parked', 1, i.execute_at,\n COALESCE($5, NOW()), COALESCE($5, NOW())\n FROM input i\n WHERE i.id NOT IN (SELECT id FROM ins)\n RETURNING id, queue_id\n )\n SELECT r.id AS \"id!: JobId\", TRUE AS \"landed_pending!\", NULL::uuid AS \"occupant_id?\"\n FROM ins r\n UNION ALL\n SELECT p.id AS \"id!: JobId\", FALSE AS \"landed_pending!\",\n COALESCE(\n (SELECT w.id FROM ins w WHERE w.queue_id = p.queue_id),\n (SELECT o.id FROM job_executions o\n WHERE o.queue_id = p.queue_id AND o.state IN ('pending', 'running'))\n ) AS \"occupant_id?\"\n FROM parked p\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id!: JobId", + "type_info": "Uuid" + }, + { + "ordinal": 1, + "name": "landed_pending!", + "type_info": "Bool" + }, + { + "ordinal": 2, + "name": "occupant_id?", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "UuidArray", + "TextArray", + "TextArray", + "TimestamptzArray", + "Timestamptz", + "TextArray" + ] + }, + "nullable": [ + null, + null, + null + ] + }, + "hash": "f004ea481cacdb54848f6a506295bda267c3c163e2735aac61d4de4dcc8c6b0c" +} diff --git a/cala-ledger/.sqlx/query-f22cd97b0bf5369bd0c072c69a625b024b41060b447355a4ff9a8a244527a601.json b/cala-ledger/.sqlx/query-f22cd97b0bf5369bd0c072c69a625b024b41060b447355a4ff9a8a244527a601.json new file mode 100644 index 000000000..1919d9a83 --- /dev/null +++ b/cala-ledger/.sqlx/query-f22cd97b0bf5369bd0c072c69a625b024b41060b447355a4ff9a8a244527a601.json @@ -0,0 +1,71 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH st AS (\n SELECT last_commit_seq\n FROM cala_persistent_outbox_commit_log_state\n WHERE singleton AND logged_through_sequence < $2::bigint\n FOR UPDATE\n ),\n grp AS (\n SELECT sequence\n FROM cala_persistent_outbox_events\n WHERE commit_xid = $1::bigint AND payload IS NOT NULL\n ),\n numbered AS (\n SELECT g.sequence,\n st.last_commit_seq + ROW_NUMBER() OVER (ORDER BY g.sequence) AS commit_seq,\n (g.sequence = MAX(g.sequence) OVER ()) AS group_last\n FROM grp g, st\n ),\n ins AS (\n INSERT INTO cala_persistent_outbox_commit_log (commit_seq, sequence, group_last)\n SELECT commit_seq, sequence, group_last FROM numbered\n RETURNING commit_seq, sequence, group_last\n ),\n upd AS (\n UPDATE cala_persistent_outbox_commit_log_state s\n SET last_commit_seq = s.last_commit_seq + (SELECT COUNT(*) FROM ins),\n logged_through_sequence = $2::bigint\n WHERE s.singleton AND EXISTS (SELECT 1 FROM st)\n RETURNING s.last_commit_seq\n )\n SELECT (SELECT COALESCE(MAX(sequence), $2::bigint) FROM grp) AS \"group_max!: i64\",\n i.commit_seq AS \"commit_seq?: i64\",\n i.group_last AS \"group_last?\",\n e.id AS \"id?\",\n e.sequence AS \"sequence?: i64\",\n e.payload,\n e.tracing_context,\n e.recorded_at AS \"recorded_at?\",\n e.commit_xid AS \"commit_xid?\"\n FROM (SELECT 1) AS one\n LEFT JOIN ins i ON TRUE\n LEFT JOIN cala_persistent_outbox_events e ON e.sequence = i.sequence\n LEFT JOIN upd ON TRUE\n ORDER BY i.commit_seq NULLS FIRST", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "group_max!: i64", + "type_info": "Int8" + }, + { + "ordinal": 1, + "name": "commit_seq?: i64", + "type_info": "Int8" + }, + { + "ordinal": 2, + "name": "group_last?", + "type_info": "Bool" + }, + { + "ordinal": 3, + "name": "id?", + "type_info": "Uuid" + }, + { + "ordinal": 4, + "name": "sequence?: i64", + "type_info": "Int8" + }, + { + "ordinal": 5, + "name": "payload", + "type_info": "Jsonb" + }, + { + "ordinal": 6, + "name": "tracing_context", + "type_info": "Jsonb" + }, + { + "ordinal": 7, + "name": "recorded_at?", + "type_info": "Timestamptz" + }, + { + "ordinal": 8, + "name": "commit_xid?", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [ + "Int8", + "Int8" + ] + }, + "nullable": [ + null, + false, + false, + false, + false, + true, + true, + false, + false + ] + }, + "hash": "f22cd97b0bf5369bd0c072c69a625b024b41060b447355a4ff9a8a244527a601" +} diff --git a/cala-ledger/.sqlx/query-ff58bbfb397e2bbe3108ef363d3cff19f4760d88c2934249dfd2a0285be3f3e5.json b/cala-ledger/.sqlx/query-ff58bbfb397e2bbe3108ef363d3cff19f4760d88c2934249dfd2a0285be3f3e5.json new file mode 100644 index 000000000..1004c5188 --- /dev/null +++ b/cala-ledger/.sqlx/query-ff58bbfb397e2bbe3108ef363d3cff19f4760d88c2934249dfd2a0285be3f3e5.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH ordered AS MATERIALIZED (\n SELECT DISTINCT key FROM UNNEST($1::text[]) AS t(key) ORDER BY key\n )\n SELECT pg_advisory_xact_lock($2, hashtext(concat($3::text, ':', key)))\n FROM ordered\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "pg_advisory_xact_lock", + "type_info": "Void" + } + ], + "parameters": { + "Left": [ + "TextArray", + "Int4", + "Text" + ] + }, + "nullable": [ + null + ] + }, + "hash": "ff58bbfb397e2bbe3108ef363d3cff19f4760d88c2934249dfd2a0285be3f3e5" +} diff --git a/cala-ledger/.sqlx/query-ffeaefb95f238936233fd691b2923c9a01ad5840700c7ba784f0cd7d8864cfc3.json b/cala-ledger/.sqlx/query-ffeaefb95f238936233fd691b2923c9a01ad5840700c7ba784f0cd7d8864cfc3.json deleted file mode 100644 index 37fc82921..000000000 --- a/cala-ledger/.sqlx/query-ffeaefb95f238936233fd691b2923c9a01ad5840700c7ba784f0cd7d8864cfc3.json +++ /dev/null @@ -1,29 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH to_delete AS MATERIALIZED (\n SELECT id FROM job_executions\n WHERE id = ANY($1) AND poller_instance_id = $2\n ORDER BY queue_id, id\n FOR UPDATE\n ), deleted AS (\n DELETE FROM job_executions je USING to_delete t WHERE je.id = t.id\n RETURNING je.id, je.queue_id\n ), cleanup AS (\n DELETE FROM job_execution_states s USING deleted d WHERE s.id = d.id\n )\n SELECT id AS \"id!: JobId\", queue_id AS \"queue_id?\"\n FROM deleted\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id!: JobId", - "type_info": "Uuid" - }, - { - "ordinal": 1, - "name": "queue_id?", - "type_info": "Varchar" - } - ], - "parameters": { - "Left": [ - "UuidArray", - "Uuid" - ] - }, - "nullable": [ - false, - true - ] - }, - "hash": "ffeaefb95f238936233fd691b2923c9a01ad5840700c7ba784f0cd7d8864cfc3" -} diff --git a/cala-ledger/migrations/20251204130226_cala_obix_setup.sql b/cala-ledger/migrations/20251204130226_cala_obix_setup.sql index a9fa76448..406f681d4 100644 --- a/cala-ledger/migrations/20251204130226_cala_obix_setup.sql +++ b/cala-ledger/migrations/20251204130226_cala_obix_setup.sql @@ -3,17 +3,28 @@ -- Partitioned by `sequence` with a DEFAULT catch-all partition, so an insert -- can never fail to route. The primary key is `sequence` (a partitioned -- table's PK must include the partition key); `id` is a plain column. --- Partitions are pre-created ahead of the head by the obix partition --- maintainer job (registered by the application). +-- Partitions are pre-created ahead of the head by the maintainer job in +-- `src/out/partition`. +-- `commit_xid` holds the writer's top-level transaction id. Rows sharing one +-- value committed together. `pg_current_xact_id()` returns `xid8`, which is +-- 64-bit and does not wrap; there is no direct cast to bigint, hence the text +-- hop. CREATE TABLE cala_persistent_outbox_events ( id UUID NOT NULL DEFAULT gen_random_uuid(), sequence BIGSERIAL, payload JSONB, tracing_context JSONB, recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + commit_xid BIGINT NOT NULL DEFAULT pg_current_xact_id()::text::bigint, PRIMARY KEY (sequence) ) PARTITION BY RANGE (sequence); +-- Backs the sequencer's per-group lookup: every row of one `commit_xid`. +-- Created on the parent; Postgres cascades it to existing and future +-- partitions. +CREATE INDEX cala_idx_persistent_outbox_events_commit_xid + ON cala_persistent_outbox_events (commit_xid); + -- Initial partition. Its range MUST equal DEFAULT_PARTITION_WIDTH (a fixed -- constant) so maintainer-created partitions tile onto it without overlapping. -- Storage params are set per-partition (not inherited via PARTITION OF). @@ -30,6 +41,63 @@ CREATE TABLE cala_persistent_outbox_events_p0 PARTITION OF cala_persistent_outbo CREATE TABLE cala_persistent_outbox_events_default PARTITION OF cala_persistent_outbox_events DEFAULT; +-- Commit-ordered delivery lane: the materialised commit order, appended by +-- the sequencer (`src/out/persistent/sequencer.rs`). A group is appended +-- whole when its lowest member is reached in insert order, members by +-- `sequence` within it, and the resulting `commit_seq` is dense. +-- +-- Positions only; payloads are reached by joining on `sequence`. Retention +-- must therefore drop log partitions before or with the event partitions +-- they reference, never after. +CREATE TABLE cala_persistent_outbox_commit_log ( + commit_seq BIGINT NOT NULL, + sequence BIGINT NOT NULL, + group_last BOOLEAN NOT NULL, + PRIMARY KEY (commit_seq) +) PARTITION BY RANGE (commit_seq); + +-- Range equal to DEFAULT_PARTITION_WIDTH, as for the events table, so +-- maintainer-created partitions tile onto it without overlapping. +CREATE TABLE cala_persistent_outbox_commit_log_p0 PARTITION OF cala_persistent_outbox_commit_log + FOR VALUES FROM (0) TO (2000000) + WITH (autovacuum_vacuum_insert_scale_factor = 0.0, + autovacuum_vacuum_insert_threshold = 50000, + autovacuum_freeze_min_age = 0, + fillfactor = 100); + +CREATE TABLE cala_persistent_outbox_commit_log_default + PARTITION OF cala_persistent_outbox_commit_log DEFAULT; + +-- Backs the sequencer's restart read: the sequences above +-- `logged_through_sequence` that are already logged. Not UNIQUE — a +-- partitioned table's UNIQUE must include the partition key, which here is +-- `commit_seq`. +CREATE INDEX cala_idx_persistent_outbox_commit_log_sequence + ON cala_persistent_outbox_commit_log (sequence); + +-- Sequencer state. The two watermarks count different things: +-- `last_commit_seq` is a `commit_seq` (a position in the commit log), +-- `logged_through_sequence` is an insert `sequence` (a position in +-- `persistent_outbox_events`). Appending a group of three advances the +-- first by three and sets the second to that group's lowest member, so +-- neither tracks the other. +-- +-- Every payload-bearing event at or below `logged_through_sequence` is in +-- the log; events above it may be too, when a group straddles it. An append +-- locks this row and is conditional on `logged_through_sequence` being below +-- the appending sequence, which is what makes concurrent sequencers append +-- each group once. +-- +-- `singleton` admits one row and no other: the CHECK allows only TRUE and +-- the primary key allows only one of it. +CREATE TABLE cala_persistent_outbox_commit_log_state ( + singleton BOOLEAN PRIMARY KEY DEFAULT TRUE CHECK (singleton), + last_commit_seq BIGINT NOT NULL, + logged_through_sequence BIGINT NOT NULL +); +INSERT INTO cala_persistent_outbox_commit_log_state (last_commit_seq, logged_through_sequence) +VALUES (0, 0) ON CONFLICT (singleton) DO NOTHING; + -- Ephemeral outbox events CREATE TABLE cala_ephemeral_outbox_events ( event_type VARCHAR NOT NULL UNIQUE, @@ -81,19 +149,38 @@ CREATE INDEX cala_idx_inbox_events_status ON cala_inbox_events(status) WHERE status IN ('pending', 'processing', 'failed'); -- Keyed-subscriber subscriptions: one row per (subscriber_type, key) --- identity. Row presence IS the subscription: absence means cancelled. --- Identity and terms only (key, wake keys, instance config, birth --- frontier); execution and progress live in the job crate's own tables. +-- identity. -- --- `wake_keys` are a liveness signal, not a delivery filter: they decide --- whom to wake when a member has passivated, matched by set overlap. --- Never empty (rejected at subscribe time). +-- Row presence IS the subscription: absence means cancelled. This table +-- holds identity and terms only (key, wake keys, instance config, birth +-- frontier); execution and progress (liveness, generations, attempts, +-- watermark) live entirely in the job crate's own tables, addressed by +-- (subscriber_type, key) through job's keyed-job machinery. Readers here +-- must never join against job-crate tables to decide wakes (schema +-- boundary). -- +-- `wake_keys` are a liveness signal, NOT a delivery filter: a live member +-- reads the whole stream from its own cursor and decides per event in its +-- own handler. These keys only decide whom to *wake* when a member has +-- passivated. Matching is set-overlap on both sides — an event classifies +-- to a set of wake keys, a subscription declares the set it watches, and an +-- intersection respawns it. Never empty (rejected at subscribe time): an +-- empty set overlaps nothing, so such a row could never be woken again once +-- it passivated. +-- +-- A set rather than a scalar because a subscription's identity is its +-- `key`, so watching several partitions of the stream cannot be expressed +-- as extra rows. -- `checkpoint` mirrors the member's durable cursor, written in the same --- transaction as the job's own checkpoint, so the waker can find members --- drifting out of the in-memory event cache without joining across the --- schema boundary. It is a lower bound: a stale value costs at worst a --- spurious, idempotent wake. +-- transaction as the job's own checkpoint. The authoritative cursor still +-- lives in the job crate's execution state; this is a copy obix owns so the +-- waker can ask "who has fallen far enough behind the in-memory event cache +-- that waking them now would save a paged cold read from disk" without +-- joining across the schema boundary into job-crate tables. +-- +-- It is a lower bound, never an over-estimate: a member that dies without +-- passivating leaves it behind its true position, which costs at worst a +-- spurious wake (idempotent, resolves to the live holder or an empty run). CREATE TABLE cala_subscriptions ( subscriber_type VARCHAR NOT NULL, key VARCHAR NOT NULL, @@ -105,9 +192,14 @@ CREATE TABLE cala_subscriptions ( PRIMARY KEY (subscriber_type, key) ); --- Backs the waker's catch-up scan: the members nearest the eviction cliff --- are woken first. +-- Backs the waker's catch-up scan: `WHERE checkpoint < $1 ORDER BY +-- checkpoint ASC LIMIT $2` across every subscriber type, so the members +-- nearest the eviction cliff are the ones woken first and the per-pass +-- limit bounds the wake rate. CREATE INDEX cala_idx_subscriptions_checkpoint ON cala_subscriptions (checkpoint); --- Backs the waker's flush-time lookup on `wake_keys && $2::varchar[]`. +-- Backs the waker's flush-time lookup: `WHERE subscriber_type = $1 AND +-- wake_keys && $2::varchar[]`. The cast is load-bearing — Postgres has no +-- implicit varchar[]/text[] cast for `&&`, and sqlx infers text[] for an +-- untyped array parameter. CREATE INDEX cala_idx_subscriptions_wake_keys ON cala_subscriptions USING GIN (wake_keys); diff --git a/cala-ledger/src/ec_rollup.rs b/cala-ledger/src/ec_rollup.rs index 2ed73c020..f36aebb1b 100644 --- a/cala-ledger/src/ec_rollup.rs +++ b/cala-ledger/src/ec_rollup.rs @@ -46,6 +46,7 @@ use chrono::{DateTime, NaiveDate, Utc}; use std::collections::{HashMap, HashSet}; +use std::sync::Arc; use job::{JobType, Jobs}; use obix::{ @@ -204,7 +205,7 @@ impl SingletonSubscriber for EcBalanceRollupHandler { async fn handle_persistent<'inv>( &self, ctx: EventCtx<'inv, Self::Batch>, - event: &PersistentOutboxEvent, + event: &Arc>, ) -> Result, Box> { match &event.payload { Some(OutboxEventPayload::TransactionCreated { transaction }) => { From 9817393b00080feacd8e5e7dcc07335b257d34b9 Mon Sep 17 00:00:00 2001 From: bodymindarts Date: Wed, 16 Sep 2026 09:19:40 +0200 Subject: [PATCH 2/2] perf(ec_rollup): collect shared events instead of copying payloads MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit obix 0.11.0 delivers each persistent event as the Arc it decoded once, so a collect_with fold can retain a refcount instead of copying the payload out. The rollup batch now holds those Arcs: - EntryCreated: was a full EntryValues clone per entry per batch (Decimal, Currency, description, metadata JSON). Stragglers whose transaction flushed in an earlier landing were cloned and then dropped, which is pure waste. - TransactionCreated: was a Vec allocation per transaction. The Copy scalars are still copied out; entry_ids is read back through the event. Nothing downstream needed owned entries — the applier stores &'a EntryValues in SnapshotOrEntry and only reads them — so EcRollupTxn now carries Vec<&'a EntryValues>, borrowed either from the batch's events or from the flush's fetched map. Both outlive the apply call. Snapshots::from_ec_entries takes an IntoIterator of &EntryValues; &Vec still satisfies it, so its existing callers are unchanged. into_rollup_txns becomes rollup_txns(&self, &fetched). It reads the maps instead of draining them; entry ids are unique to one transaction, so get and remove are equivalent here. 104 tests pass (lib, ec_streaming_rollup, effective_balance). Co-Authored-By: Claude Opus 5 --- cala-ledger/src/balance/effective/mod.rs | 6 +- cala-ledger/src/balance/mod.rs | 15 +-- cala-ledger/src/balance/snapshot.rs | 6 +- cala-ledger/src/ec_rollup.rs | 117 ++++++++++++++++++----- 4 files changed, 106 insertions(+), 38 deletions(-) diff --git a/cala-ledger/src/balance/effective/mod.rs b/cala-ledger/src/balance/effective/mod.rs index fee30e316..333829107 100644 --- a/cala-ledger/src/balance/effective/mod.rs +++ b/cala-ledger/src/balance/effective/mod.rs @@ -309,7 +309,7 @@ impl EffectiveBalances { &self, op: &mut impl es_entity::AtomicOperation, journal_id: JournalId, - txns: &[EcRollupTxn], + txns: &[EcRollupTxn<'_>], ec_mappings: &HashMap>, ec_leaves: &HashSet, ) -> Result<(), BalanceError> { @@ -328,7 +328,7 @@ impl EffectiveBalances { // whose own entries are all later, for no benefit. let mut earliest: HashMap<(AccountId, Currency), NaiveDate> = HashMap::new(); for tx in txns { - for entry in tx.entries.iter() { + for entry in tx.entries.iter().copied() { for target in targets(&entry.account_id) { earliest .entry((target, entry.currency)) @@ -360,7 +360,7 @@ impl EffectiveBalances { } for (tx_index, tx) in txns.iter().enumerate() { - for entry in tx.entries.iter() { + for entry in tx.entries.iter().copied() { for target in targets(&entry.account_id) { if let Some(data) = all_data.get_mut(&(target, entry.currency)) { data.push(tx.effective, tx_index, tx.created_at, entry); diff --git a/cala-ledger/src/balance/mod.rs b/cala-ledger/src/balance/mod.rs index 3c378688d..f5e0938b8 100644 --- a/cala-ledger/src/balance/mod.rs +++ b/cala-ledger/src/balance/mod.rs @@ -65,11 +65,14 @@ pub(crate) use snapshot::*; /// One committed transaction's contribution to a streaming-rollup batch /// (see [`Balances::apply_ec_rollup_in_op`]). -pub(crate) struct EcRollupTxn { +pub(crate) struct EcRollupTxn<'a> { pub journal_id: JournalId, pub effective: NaiveDate, pub created_at: DateTime, - pub entries: Vec, + /// Borrowed for the length of the flush: stream-collected entries point + /// into the shared events the outbox decoded once, DB-fetched ones into + /// the flush's `fetched` map. The applier only ever reads them. + pub entries: Vec<&'a EntryValues>, } #[derive(Clone)] @@ -280,9 +283,9 @@ impl Balances { pub(crate) async fn apply_ec_rollup_in_op( &self, op: &mut impl es_entity::AtomicOperation, - txns: Vec, + txns: Vec>, ) -> Result<(), BalanceError> { - let mut groups: Vec<(JournalId, Vec)> = Vec::new(); + let mut groups: Vec<(JournalId, Vec>)> = Vec::new(); for tx in txns { match groups.iter_mut().find(|(j, _)| *j == tx.journal_id) { Some((_, group)) => group.push(tx), @@ -300,7 +303,7 @@ impl Balances { &self, op: &mut impl es_entity::AtomicOperation, journal_id: JournalId, - group: Vec, + group: Vec>, ) -> Result<(), BalanceError> { let member_account_ids: Vec = group .iter() @@ -348,7 +351,7 @@ impl Balances { let new_balances = Snapshots::from_ec_entries( tx.created_at, current_balances.clone(), - &tx.entries, + tx.entries.iter().copied(), &ec_mappings, &ec_leaves, ); diff --git a/cala-ledger/src/balance/snapshot.rs b/cala-ledger/src/balance/snapshot.rs index b9baadef3..b525f9eba 100644 --- a/cala-ledger/src/balance/snapshot.rs +++ b/cala-ledger/src/balance/snapshot.rs @@ -147,15 +147,15 @@ impl Snapshots { name = "cala_ledger.balances.from_ec_entries", skip_all )] - pub(crate) fn from_ec_entries( + pub(crate) fn from_ec_entries<'a>( time: DateTime, current_balances: HashMap<(AccountId, Currency), Option>, - entries: &[EntryValues], + entries: impl IntoIterator, ec_mappings: &HashMap>, ec_leaves: &HashSet, ) -> Vec { let mut fold = SnapshotFold::new(time, current_balances); - for entry in entries.iter() { + for entry in entries { for set in ec_mappings .get(&entry.account_id) .into_iter() diff --git a/cala-ledger/src/ec_rollup.rs b/cala-ledger/src/ec_rollup.rs index f36aebb1b..bc0cc7822 100644 --- a/cala-ledger/src/ec_rollup.rs +++ b/cala-ledger/src/ec_rollup.rs @@ -101,15 +101,41 @@ pub(crate) async fn register_ec_balance_rollup( } /// A transaction pulled from a `TransactionCreated` event, carrying just -/// what the rollup needs. `entry_ids` is the complete expected entry set, -/// which is what makes stream-collected entries verifiable (see -/// [`EcRollupBatch`]). +/// what the rollup needs. [`entry_ids`](Self::entry_ids) is the complete +/// expected entry set, which is what makes stream-collected entries +/// verifiable (see [`EcRollupBatch`]). +/// +/// The scalars are copied out (they are `Copy`); the id list is instead +/// read back through `event`, the shared event the outbox decoded once, +/// so collecting a transaction costs a refcount rather than a `Vec` +/// allocation per transaction. struct PendingTx { id: TransactionId, journal_id: JournalId, effective: NaiveDate, created_at: DateTime, - entry_ids: Vec, + event: Arc>, +} + +impl PendingTx { + fn entry_ids(&self) -> &[EntryId] { + match &self.event.payload { + Some(OutboxEventPayload::TransactionCreated { transaction }) => &transaction.entry_ids, + _ => unreachable!( + "PendingTx is only built from a TransactionCreated event in handle_persistent" + ), + } + } +} + +/// The entry carried by a collected `EntryCreated` event. +fn entry_of(event: &PersistentOutboxEvent) -> &EntryValues { + match &event.payload { + Some(OutboxEventPayload::EntryCreated { entry }) => entry, + _ => unreachable!( + "EcRollupBatch::entries only ever holds EntryCreated events, pushed in handle_persistent" + ), + } } /// One batch landing's accumulator. @@ -128,7 +154,7 @@ struct PendingTx { #[derive(Default)] struct EcRollupBatch { txns: Vec, - entries: HashMap>, + entries: HashMap>>>, } impl EcRollupBatch { @@ -136,11 +162,11 @@ impl EcRollupBatch { self.txns.push(tx); } - fn push_entry(&mut self, entry: EntryValues) { + fn push_entry(&mut self, event: Arc>) { self.entries - .entry(entry.transaction_id) + .entry(entry_of(&event).transaction_id) .or_default() - .push(entry); + .push(event); } /// Entry ids that were *not* collected from the stream in this landing @@ -153,9 +179,9 @@ impl EcRollupBatch { let collected: HashSet = self .entries .get(&tx.id) - .map(|entries| entries.iter().map(|e| e.id).collect()) + .map(|events| events.iter().map(|e| entry_of(e).id).collect()) .unwrap_or_default(); - tx.entry_ids + tx.entry_ids() .iter() .copied() .filter(move |id| !collected.contains(id)) @@ -166,17 +192,27 @@ impl EcRollupBatch { /// Assemble the applier's input in landing order: each transaction's /// stream-collected entries, topped up from the DB-`fetched` map where /// the group straddled a landing boundary, sorted by entry sequence. - fn into_rollup_txns(self, mut fetched: HashMap) -> Vec { - let EcRollupBatch { txns, mut entries } = self; - txns.into_iter() + /// + /// Entries are borrowed, never copied: stream-collected ones out of the + /// events this batch holds, fetched ones out of `fetched`. Both outlive + /// the applier call in [`flush`](EcBalanceRollupHandler::flush), and + /// stragglers left unused cost nothing. + fn rollup_txns<'a>(&'a self, fetched: &'a HashMap) -> Vec> { + self.txns + .iter() .map(|tx| { - let mut entry_values = entries.remove(&tx.id).unwrap_or_default(); - if entry_values.len() != tx.entry_ids.len() { + let entry_ids = tx.entry_ids(); + let mut entry_values: Vec<&EntryValues> = self + .entries + .get(&tx.id) + .map(|events| events.iter().map(|e| entry_of(e)).collect()) + .unwrap_or_default(); + if entry_values.len() != entry_ids.len() { entry_values.extend( - tx.entry_ids + entry_ids .iter() - .filter_map(|id| fetched.remove(id)) - .map(Entry::into_values), + .filter_map(|id| fetched.get(id)) + .map(Entry::values), ); } entry_values.sort_by_key(|e| e.sequence); @@ -214,13 +250,13 @@ impl SingletonSubscriber for EcBalanceRollupHandler { journal_id: transaction.journal_id, effective: transaction.effective, created_at: transaction.created_at, - entry_ids: transaction.entry_ids.clone(), + event: Arc::clone(event), }; Ok(ctx.collect_with(|batch| batch.push_tx(tx))) } - Some(OutboxEventPayload::EntryCreated { entry }) => { - let entry = entry.clone(); - Ok(ctx.collect_with(|batch| batch.push_entry(entry))) + Some(OutboxEventPayload::EntryCreated { .. }) => { + let event = Arc::clone(event); + Ok(ctx.collect_with(|batch| batch.push_entry(event))) } _ => Ok(ctx.skip()), } @@ -244,7 +280,7 @@ impl SingletonSubscriber for EcBalanceRollupHandler { self.entries.find_all_in_op(&mut *op, &missing_ids).await? }; - let rollup_txns = batch.into_rollup_txns(fetched); + let rollup_txns = batch.rollup_txns(&fetched); self.balances.apply_ec_rollup_in_op(op, rollup_txns).await?; Ok(()) } @@ -266,6 +302,19 @@ mod __fuzz { entry_ids: Vec, } + /// Wrap a payload in the shared event shape the batch now collects. + /// The envelope fields are inert here — only the payload is read. + fn event(payload: OutboxEventPayload) -> Arc> { + Arc::new(PersistentOutboxEvent { + id: obix::out::OutboxEventId::new(), + sequence: obix::EventSequence::from(0u64), + payload: Some(payload), + tracing_context: None, + recorded_at: Utc::now(), + commit_group: obix::CommitGroupId::from(0i64), + }) + } + pub fn fuzz_batch(data: &[u8]) { let parts: Vec<&[u8]> = data.split(|&b| b == 0xFF).collect(); if parts.len() < 2 { @@ -285,15 +334,31 @@ mod __fuzz { journal_id: t.journal_id, effective: t.effective, created_at: t.created_at, - entry_ids: t.entry_ids.clone(), + event: event(OutboxEventPayload::TransactionCreated { + transaction: cala_types::transaction::TransactionValues { + id: t.id, + journal_id: t.journal_id, + effective: t.effective, + created_at: t.created_at, + entry_ids: t.entry_ids.clone(), + // Inert: the batch only reads the fields above. + version: 1, + modified_at: t.created_at, + tx_template_id: crate::primitives::TxTemplateId::new(), + correlation_id: String::new(), + external_id: None, + description: None, + metadata: None, + }, + }), }); } for e in &entries { - batch.push_entry(e.clone()); + batch.push_entry(event(OutboxEventPayload::EntryCreated { entry: e.clone() })); } let _missing = batch.missing_entry_ids(); - let _rollup = batch.into_rollup_txns(HashMap::::new()); + let _rollup = batch.rollup_txns(&HashMap::::new()); } }