From 776958703c77990f847120880440934061c0813e Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Thu, 16 Oct 2025 13:28:49 +0530 Subject: [PATCH 01/12] feat: migrate to job crate --- Cargo.lock | 372 +++++++++++++++++------- Cargo.toml | 1 + migrations/20250904065521_job_setup.sql | 56 ++++ src/app/error.rs | 4 + src/app/mod.rs | 17 ++ src/job/mod.rs | 50 +--- src/job_svc/mod.rs | 20 ++ src/job_svc/populate_outbox.rs | 129 ++++++++ src/lib.rs | 1 + 9 files changed, 494 insertions(+), 156 deletions(-) create mode 100644 migrations/20250904065521_job_setup.sql create mode 100644 src/job_svc/mod.rs create mode 100644 src/job_svc/populate_outbox.rs diff --git a/Cargo.lock b/Cargo.lock index c4dfd12f..f4d8507e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,21 +2,6 @@ # It is not intended for manual editing. version = 4 -[[package]] -name = "addr2line" -version = "0.24.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dfbe277e56a376000877090da837660b4427aad530e3028d44e0bffe4f89a1c1" -dependencies = [ - "gimli", -] - -[[package]] -name = "adler2" -version = "2.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "512761e0bb2578dd7380c6baaa0f4ce03e84f95e960231d1dec8bf4d7d6e2627" - [[package]] name = "aead" version = "0.5.2" @@ -180,9 +165,9 @@ dependencies = [ [[package]] name = "async-trait" -version = "0.1.88" +version = "0.1.89" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e539d3fca749fcee5236ab05e93a52867dd549cc157c8cb7f99595f3cedffdb5" +checksum = "9035ad2d096bed7955a320ee7e2230574d28fd3c3a0f186cbea1ff3c7eed5dbb" dependencies = [ "proc-macro2", "quote", @@ -311,21 +296,6 @@ dependencies = [ "tower-service", ] -[[package]] -name = "backtrace" -version = "0.3.74" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8d82cb332cdfaed17ae235a638438ac4d4839913cc2af585c3c6746e8f8bee1a" -dependencies = [ - "addr2line", - "cfg-if", - "libc", - "miniz_oxide", - "object", - "rustc-demangle", - "windows-targets 0.52.6", -] - [[package]] name = "base64" version = "0.13.1" @@ -364,7 +334,7 @@ dependencies = [ "js-sys", "log", "miniscript", - "rand", + "rand 0.8.5", "serde", "serde_json", "sled", @@ -546,6 +516,7 @@ dependencies = [ "fedimint-tonic-lnd", "futures", "hex", + "job", "miniscript", "opentelemetry", "opentelemetry-otlp", @@ -553,7 +524,7 @@ dependencies = [ "prost 0.12.6", "prost-wkt-types", "protobuf-src", - "rand", + "rand 0.8.5", "regex", "reqwest", "reqwest-middleware", @@ -716,7 +687,7 @@ dependencies = [ "num-traits", "serde", "wasm-bindgen", - "windows-link", + "windows-link 0.1.1", ] [[package]] @@ -870,7 +841,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1bfb12502f3fc46cca1bb51ac28df9d618d813cdc3d2f25b9fe775a34af26bb3" dependencies = [ "generic-array", - "rand_core", + "rand_core 0.6.4", "typenum", ] @@ -1031,6 +1002,12 @@ version = "0.15.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1aaf95b3e5c8f23aa320147307562d361db0ae0d51242340f558153b4eb2439b" +[[package]] +name = "dyn-clone" +version = "1.0.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d0881ea181b1df73ff77ffaaf9c7544ecc11e82fba9b5f27b262a3c73a332555" + [[package]] name = "either" version = "1.13.0" @@ -1373,12 +1350,6 @@ dependencies = [ "wasi 0.14.2+wasi-0.2.4", ] -[[package]] -name = "gimli" -version = "0.31.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "07e28edb80900c19c28f1072f2e8aeca7fa06b23cd4169cefe1af5aa3260783f" - [[package]] name = "glob" version = "0.3.1" @@ -1598,7 +1569,7 @@ dependencies = [ "httpdate", "itoa", "pin-project-lite", - "socket2", + "socket2 0.5.7", "tokio", "tower-service", "tracing", @@ -1696,7 +1667,7 @@ dependencies = [ "http-body 1.0.1", "hyper 1.5.1", "pin-project-lite", - "socket2", + "socket2 0.5.7", "tokio", "tower-service", "tracing", @@ -1877,7 +1848,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d0acd33ff0285af998aaf9b57342af478078f53492322fafc47450e09397e0e9" dependencies = [ "bitmaps", - "rand_core", + "rand_core 0.6.4", "rand_xoshiro", "serde", "sized-chunks", @@ -1970,6 +1941,28 @@ version = "1.0.14" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d75a2a4b1b190afb6f5425f10f6a8f959d2ea0b9c2b1d79553551850539e4674" +[[package]] +name = "job" +version = "0.1.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fbc97547bab9416fb9e5490474afef9f452611dbff8a478cf9dc3cf5bb093bdb" +dependencies = [ + "async-trait", + "chrono", + "derive_builder", + "es-entity", + "futures", + "rand 0.9.2", + "serde", + "serde_json", + "serde_with", + "sqlx", + "thiserror 2.0.12", + "tokio", + "tracing", + "uuid", +] + [[package]] name = "js-sys" version = "0.3.77" @@ -2043,9 +2036,9 @@ dependencies = [ [[package]] name = "libc" -version = "0.2.165" +version = "0.2.177" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fcb4d3d38eab6c5239a362fa8bae48c03baf980a6e7079f063942d563ef3533e" +checksum = "2874a2af47a2325c2001a6e6fad9b16a53b802102b528163885171cf92b15976" [[package]] name = "libm" @@ -2139,15 +2132,6 @@ dependencies = [ "serde", ] -[[package]] -name = "miniz_oxide" -version = "0.8.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e2d80299ef12ff69b16a84bb182e3b9df68b5a91574d3d4fa6e41b65deec4df1" -dependencies = [ - "adler2", -] - [[package]] name = "mio" version = "1.0.2" @@ -2193,7 +2177,7 @@ dependencies = [ "num-integer", "num-iter", "num-traits", - "rand", + "rand 0.8.5", "smallvec", "zeroize", ] @@ -2234,15 +2218,6 @@ dependencies = [ "libm", ] -[[package]] -name = "object" -version = "0.36.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "aedf0a2d09c573ed1d8d85b30c119153926a2b36dce0ab28322c09a117a4683e" -dependencies = [ - "memchr", -] - [[package]] name = "once_cell" version = "1.20.2" @@ -2329,7 +2304,7 @@ dependencies = [ "once_cell", "opentelemetry", "percent-encoding", - "rand", + "rand 0.8.5", "serde_json", "thiserror 1.0.69", "tokio", @@ -2724,7 +2699,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls 0.23.18", - "socket2", + "socket2 0.5.7", "thiserror 2.0.12", "tokio", "tracing", @@ -2738,7 +2713,7 @@ checksum = "a2fe5ef3495d7d2e377ff17b1a8ce2ee2ec2a18cde8b6ad6619d65d0701c135d" dependencies = [ "bytes", "getrandom 0.2.15", - "rand", + "rand 0.8.5", "ring", "rustc-hash", "rustls 0.23.18", @@ -2759,7 +2734,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2", + "socket2 0.5.7", "tracing", "windows-sys 0.59.0", ] @@ -2792,8 +2767,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" dependencies = [ "libc", - "rand_chacha", - "rand_core", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + +[[package]] +name = "rand" +version = "0.9.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6db2770f06117d490610c7488547d543617b21bfa07796d7a12f6f1bd53850d1" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.3", ] [[package]] @@ -2803,7 +2788,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.6.4", +] + +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.3", ] [[package]] @@ -2815,13 +2810,22 @@ dependencies = [ "getrandom 0.2.15", ] +[[package]] +name = "rand_core" +version = "0.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" +dependencies = [ + "getrandom 0.3.3", +] + [[package]] name = "rand_xoshiro" version = "0.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6f97cdb2a36ed4183de61b2f824cc45c9f1037f28afe0a322e9fff4c108b5aaa" dependencies = [ - "rand_core", + "rand_core 0.6.4", ] [[package]] @@ -2842,6 +2846,26 @@ dependencies = [ "bitflags 2.6.0", ] +[[package]] +name = "ref-cast" +version = "1.0.25" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f354300ae66f76f1c85c5f84693f0ce81d747e2c3f21a45fef496d89c960bf7d" +dependencies = [ + "ref-cast-impl", +] + +[[package]] +name = "ref-cast-impl" +version = "1.0.25" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b7186006dcb21920990093f30e3dea63b7d6e977bf1256be20c3563a5db070da" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.104", +] + [[package]] name = "regex" version = "1.11.1" @@ -2968,7 +2992,7 @@ checksum = "493b4243e32d6eedd29f9a398896e35c6943a123b55eec97dcaee98310d25810" dependencies = [ "anyhow", "chrono", - "rand", + "rand 0.8.5", ] [[package]] @@ -3027,7 +3051,7 @@ dependencies = [ "num-traits", "pkcs1", "pkcs8", - "rand_core", + "rand_core 0.6.4", "signature", "spki", "subtle", @@ -3044,7 +3068,7 @@ dependencies = [ "borsh", "bytes", "num-traits", - "rand", + "rand 0.8.5", "rkyv", "serde", "serde_json", @@ -3060,12 +3084,6 @@ dependencies = [ "rust_decimal", ] -[[package]] -name = "rustc-demangle" -version = "0.1.24" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "719b953e2095829ee67db738b3bfa9fa368c94900df327b3f07fe6e794d2fe1f" - [[package]] name = "rustc-hash" version = "2.0.0" @@ -3199,6 +3217,30 @@ dependencies = [ "sdd", ] +[[package]] +name = "schemars" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4cd191f9397d57d581cddd31014772520aa448f65ef991055d7f61582c65165f" +dependencies = [ + "dyn-clone", + "ref-cast", + "serde", + "serde_json", +] + +[[package]] +name = "schemars" +version = "1.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "82d20c4491bc164fa2f6c5d44565947a52ad80b9505d8e36f8d54c27c739fcd0" +dependencies = [ + "dyn-clone", + "ref-cast", + "serde", + "serde_json", +] + [[package]] name = "scopeguard" version = "1.2.0" @@ -3234,7 +3276,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "25996b82292a7a57ed3508f052cfff8640d38d32018784acd714758b43da9c8f" dependencies = [ "bitcoin_hashes", - "rand", + "rand 0.8.5", "secp256k1-sys", "serde", ] @@ -3250,18 +3292,28 @@ dependencies = [ [[package]] name = "serde" -version = "1.0.219" +version = "1.0.228" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5f0e2c6ed6606019b4e29e69dbaba95b11854410e5347d525002456dbbb786b6" +checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e" +dependencies = [ + "serde_core", + "serde_derive", +] + +[[package]] +name = "serde_core" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad" dependencies = [ "serde_derive", ] [[package]] name = "serde_derive" -version = "1.0.219" +version = "1.0.228" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5b0276cf7f2c73365f7157c8123c21cd9a50fbbd844757af28ca1f5925fc2a00" +checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" dependencies = [ "proc-macro2", "quote", @@ -3270,14 +3322,15 @@ dependencies = [ [[package]] name = "serde_json" -version = "1.0.140" +version = "1.0.145" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "20068b6e96dc6c9bd23e01df8827e6c7e1f2fddd43c21810382803c136b99373" +checksum = "402a6f66d8c709116cf22f558eab210f5a50187f702eb4d7e5ef38d9a7f1c79c" dependencies = [ "itoa", "memchr", "ryu", "serde", + "serde_core", ] [[package]] @@ -3294,17 +3347,18 @@ dependencies = [ [[package]] name = "serde_with" -version = "3.11.0" +version = "3.15.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8e28bdad6db2b8340e449f7108f020b3b092e8583a9e3fb82713e1d4e71fe817" +checksum = "6093cd8c01b25262b84927e0f7151692158fab02d961e04c979d3903eba7ecc5" dependencies = [ "base64 0.22.1", "chrono", "hex", "indexmap 1.9.3", "indexmap 2.6.0", - "serde", - "serde_derive", + "schemars 0.9.0", + "schemars 1.0.4", + "serde_core", "serde_json", "serde_with_macros", "time", @@ -3312,11 +3366,11 @@ dependencies = [ [[package]] name = "serde_with_macros" -version = "3.11.0" +version = "3.15.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9d846214a9854ef724f3da161b426242d8de7c1fc7de2f89bb1efcb154dca79d" +checksum = "a7e6c180db0816026a61afa1cff5344fb7ebded7e4d3062772179f2501481c27" dependencies = [ - "darling 0.20.11", + "darling 0.21.3", "proc-macro2", "quote", "syn 2.0.104", @@ -3423,7 +3477,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de" dependencies = [ "digest", - "rand_core", + "rand_core 0.6.4", ] [[package]] @@ -3492,6 +3546,16 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "socket2" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "17129e116933cf371d018bb80ae557e889637989d8638274fb25622827b03881" +dependencies = [ + "libc", + "windows-sys 0.60.2", +] + [[package]] name = "spin" version = "0.9.8" @@ -3680,7 +3744,7 @@ dependencies = [ "memchr", "once_cell", "percent-encoding", - "rand", + "rand 0.8.5", "rsa", "rust_decimal", "serde", @@ -3721,7 +3785,7 @@ dependencies = [ "md-5", "memchr", "once_cell", - "rand", + "rand 0.8.5", "rust_decimal", "serde", "serde_json", @@ -4014,20 +4078,19 @@ checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" [[package]] name = "tokio" -version = "1.41.1" +version = "1.48.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "22cfb5bee7a6a52939ca9224d6ac897bb669134078daa8735560897f69de4d33" +checksum = "ff360e02eab121e0bc37a2d3b4d4dc622e6eda3a8e5253d5435ecf5bd4c68408" dependencies = [ - "backtrace", "bytes", "libc", "mio", "parking_lot 0.12.3", "pin-project-lite", "signal-hook-registry", - "socket2", + "socket2 0.6.1", "tokio-macros", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -4042,9 +4105,9 @@ dependencies = [ [[package]] name = "tokio-macros" -version = "2.4.0" +version = "2.6.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "693d596312e88961bc67d7f1f97af8a70227d9f90c31bba5806eec004978d752" +checksum = "af407857209536a95c8e56f8231ef2c2e2aff839b22e07a1ffcbc617e9db9fa5" dependencies = [ "proc-macro2", "quote", @@ -4192,7 +4255,7 @@ dependencies = [ "percent-encoding", "pin-project", "prost 0.13.3", - "socket2", + "socket2 0.5.7", "tokio", "tokio-stream", "tower 0.4.13", @@ -4251,7 +4314,7 @@ dependencies = [ "indexmap 1.9.3", "pin-project", "pin-project-lite", - "rand", + "rand 0.8.5", "slab", "tokio", "tokio-util", @@ -4770,6 +4833,12 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "76840935b766e1b0a05c0066835fb9ec80071d4c09a16f6bd5f7e655e3c14c38" +[[package]] +name = "windows-link" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" + [[package]] name = "windows-registry" version = "0.2.0" @@ -4827,6 +4896,24 @@ dependencies = [ "windows-targets 0.52.6", ] +[[package]] +name = "windows-sys" +version = "0.60.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2f500e4d28234f72040990ec9d39e3a6b950f9f22d3dba18416c35882612bcb" +dependencies = [ + "windows-targets 0.53.5", +] + +[[package]] +name = "windows-sys" +version = "0.61.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae137229bcbd6cdf0f7b80a31df61766145077ddf49416a728b02cb3921ff3fc" +dependencies = [ + "windows-link 0.2.1", +] + [[package]] name = "windows-targets" version = "0.48.5" @@ -4851,13 +4938,30 @@ dependencies = [ "windows_aarch64_gnullvm 0.52.6", "windows_aarch64_msvc 0.52.6", "windows_i686_gnu 0.52.6", - "windows_i686_gnullvm", + "windows_i686_gnullvm 0.52.6", "windows_i686_msvc 0.52.6", "windows_x86_64_gnu 0.52.6", "windows_x86_64_gnullvm 0.52.6", "windows_x86_64_msvc 0.52.6", ] +[[package]] +name = "windows-targets" +version = "0.53.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4945f9f551b88e0d65f3db0bc25c33b8acea4d9e41163edf90dcd0b19f9069f3" +dependencies = [ + "windows-link 0.2.1", + "windows_aarch64_gnullvm 0.53.1", + "windows_aarch64_msvc 0.53.1", + "windows_i686_gnu 0.53.1", + "windows_i686_gnullvm 0.53.1", + "windows_i686_msvc 0.53.1", + "windows_x86_64_gnu 0.53.1", + "windows_x86_64_gnullvm 0.53.1", + "windows_x86_64_msvc 0.53.1", +] + [[package]] name = "windows_aarch64_gnullvm" version = "0.48.5" @@ -4870,6 +4974,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32a4622180e7a0ec044bb555404c800bc9fd9ec262ec147edd5989ccd0c02cd3" +[[package]] +name = "windows_aarch64_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53" + [[package]] name = "windows_aarch64_msvc" version = "0.48.5" @@ -4882,6 +4992,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09ec2a7bb152e2252b53fa7803150007879548bc709c039df7627cabbd05d469" +[[package]] +name = "windows_aarch64_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006" + [[package]] name = "windows_i686_gnu" version = "0.48.5" @@ -4894,12 +5010,24 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8e9b5ad5ab802e97eb8e295ac6720e509ee4c243f69d781394014ebfe8bbfa0b" +[[package]] +name = "windows_i686_gnu" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "960e6da069d81e09becb0ca57a65220ddff016ff2d6af6a223cf372a506593a3" + [[package]] name = "windows_i686_gnullvm" version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0eee52d38c090b3caa76c563b86c3a4bd71ef1a819287c19d586d7334ae8ed66" +[[package]] +name = "windows_i686_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c" + [[package]] name = "windows_i686_msvc" version = "0.48.5" @@ -4912,6 +5040,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "240948bc05c5e7c6dabba28bf89d89ffce3e303022809e73deaefe4f6ec56c66" +[[package]] +name = "windows_i686_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e7ac75179f18232fe9c285163565a57ef8d3c89254a30685b57d83a38d326c2" + [[package]] name = "windows_x86_64_gnu" version = "0.48.5" @@ -4924,6 +5058,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "147a5c80aabfbf0c7d901cb5895d1de30ef2907eb21fbbab29ca94c5b08b1a78" +[[package]] +name = "windows_x86_64_gnu" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9c3842cdd74a865a8066ab39c8a7a473c0778a3f29370b5fd6b4b9aa7df4a499" + [[package]] name = "windows_x86_64_gnullvm" version = "0.48.5" @@ -4936,6 +5076,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "24d5b23dc417412679681396f2b49f3de8c1473deb516bd34410872eff51ed0d" +[[package]] +name = "windows_x86_64_gnullvm" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0ffa179e2d07eee8ad8f57493436566c7cc30ac536a3379fdf008f47f6bb7ae1" + [[package]] name = "windows_x86_64_msvc" version = "0.48.5" @@ -4948,6 +5094,12 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" +[[package]] +name = "windows_x86_64_msvc" +version = "0.53.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650" + [[package]] name = "winnow" version = "0.6.20" diff --git a/Cargo.toml b/Cargo.toml index 39da81d8..b7fc86db 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,6 +10,7 @@ fail-on-warnings = [] [dependencies] es-entity = "0.9.0" sqlx-ledger = { version = "0.11.5", features = ["otel"] } +job_crate = { package = "job", version = "0.1.12" } anyhow = "1.0.82" bitcoincore-rpc = "0.17.0" diff --git a/migrations/20250904065521_job_setup.sql b/migrations/20250904065521_job_setup.sql new file mode 100644 index 00000000..d5a70a90 --- /dev/null +++ b/migrations/20250904065521_job_setup.sql @@ -0,0 +1,56 @@ +CREATE TABLE jobs ( + id UUID PRIMARY KEY, + unique_per_type BOOLEAN NOT NULL, + job_type VARCHAR NOT NULL, + created_at TIMESTAMPTZ NOT NULL DEFAULT NOW() +); +CREATE UNIQUE INDEX idx_unique_job_type ON jobs (job_type) WHERE unique_per_type = TRUE; + +CREATE TABLE job_events ( + id UUID NOT NULL REFERENCES jobs(id), + sequence INT NOT NULL, + event_type VARCHAR NOT NULL, + event JSONB NOT NULL, + context JSONB DEFAULT NULL, + recorded_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), + UNIQUE(id, sequence) +); + +CREATE TYPE JobExecutionState AS ENUM ('pending', 'running'); + +CREATE TABLE job_executions ( + id UUID REFERENCES jobs(id) NOT NULL UNIQUE, + job_type VARCHAR NOT NULL, + attempt_index INT NOT NULL DEFAULT 1, + state JobExecutionState NOT NULL DEFAULT 'pending', + execution_state_json JSONB, + execute_at TIMESTAMPTZ, + alive_at TIMESTAMPTZ NOT NULL, + created_at TIMESTAMPTZ NOT NULL +); + +CREATE OR REPLACE FUNCTION notify_job_execution_insert() RETURNS TRIGGER AS $$ +BEGIN + PERFORM pg_notify('job_execution', ''); + RETURN NULL; +END; +$$ LANGUAGE plpgsql; + +CREATE OR REPLACE FUNCTION notify_job_execution_update() RETURNS TRIGGER AS $$ +BEGIN + IF NEW.execute_at IS DISTINCT FROM OLD.execute_at THEN + PERFORM pg_notify('job_execution', ''); + END IF; + RETURN NULL; +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER job_executions_notify_insert_trigger +AFTER INSERT ON job_executions +FOR EACH STATEMENT +EXECUTE FUNCTION notify_job_execution_insert(); + +CREATE TRIGGER job_executions_notify_update_trigger +AFTER UPDATE ON job_executions +FOR EACH STATEMENT +EXECUTE FUNCTION notify_job_execution_update(); diff --git a/src/app/error.rs b/src/app/error.rs index 9ba387d7..a231deb9 100644 --- a/src/app/error.rs +++ b/src/app/error.rs @@ -1,6 +1,8 @@ use chacha20poly1305; use thiserror::Error; +use job_crate::error::JobError as JobSvcError; + use crate::{ address::error::AddressError, batch::error::BatchError, @@ -44,6 +46,8 @@ pub enum ApplicationError { #[error("{0}")] JobError(#[from] JobError), #[error("{0}")] + JobSvcError(#[from] JobSvcError), + #[error("{0}")] OutboxError(#[from] OutboxError), #[error("{0}")] UtxoError(#[from] UtxoError), diff --git a/src/app/mod.rs b/src/app/mod.rs index 1b4cc206..8a28a7a7 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -9,6 +9,8 @@ use std::collections::HashMap; pub use config::*; use error::*; +use job_crate::{Jobs, JobSvcConfig}; + use crate::{ account::balance::AccountBalanceSummary, address::*, @@ -17,6 +19,7 @@ use crate::{ descriptor::*, fees::{self, *}, job, + job_svc::*, ledger::*, outbox::*, payout::*, @@ -48,6 +51,7 @@ pub struct App { batch_inclusion: BatchInclusion, pool: sqlx::PgPool, config: AppConfig, + jobs: Jobs, } impl App { @@ -68,6 +72,16 @@ impl App { ) .await?; let fees_client = FeesClient::new(config.fees.clone()); + + let job_svc_config = JobSvcConfig::builder().pool(pool.clone()).build().expect("couldn't build job_svc_config"); + let mut jobs = Jobs::init(job_svc_config).await?; + JobSvc::init( + &jobs, + outbox.clone(), + ledger.clone(), + ); + jobs.start_poll().await?; + let runner = job::start_job_runner( &pool, outbox.clone(), @@ -84,6 +98,7 @@ impl App { config.blockchain.clone(), config.signer_encryption.clone(), fees_client.clone(), + jobs.clone(), ) .await?; Self::spawn_sync_all_wallets(pool.clone(), config.jobs.sync_all_wallets_delay).await?; @@ -97,6 +112,7 @@ impl App { config.jobs.respawn_all_outbox_handlers_delay, ) .await?; + let app = Self { outbox, profiles: Profiles::new(&pool), @@ -115,6 +131,7 @@ impl App { batch_inclusion, config, _runner: runner, + jobs, }; if let Some(deprecrated_encryption_key) = app.config.deprecated_encryption_key.as_ref() { app.rotate_encryption_key(deprecrated_encryption_key) diff --git a/src/job/mod.rs b/src/job/mod.rs index a7072ef0..1ae44ebd 100644 --- a/src/job/mod.rs +++ b/src/job/mod.rs @@ -3,7 +3,6 @@ mod batch_signing; mod batch_wallet_accounting; mod config; mod executor; -mod populate_outbox; mod sync_wallet; pub mod error; @@ -26,7 +25,6 @@ use batch_wallet_accounting::BatchWalletAccountingData; use error::JobError; pub use executor::JobExecutionError; use executor::JobExecutor; -use populate_outbox::PopulateOutboxData; use process_payout_queue::ProcessPayoutQueueData; use sync_wallet::SyncWalletData; @@ -51,6 +49,7 @@ pub async fn start_job_runner( blockchain_cfg: BlockchainConfig, signer_encryption_config: SignerEncryptionConfig, fees_client: FeesClient, + jobs: job_crate::Jobs, ) -> Result { let mut registry = JobRegistry::new(&[ sync_all_wallets, @@ -62,7 +61,6 @@ pub async fn start_job_runner( batch_signing, batch_broadcasting, respawn_all_outbox_handlers, - populate_outbox, ]); registry.set_context(config); registry.set_context(blockchain_cfg); @@ -78,6 +76,7 @@ pub async fn start_job_runner( registry.set_context(addresses); registry.set_context(signer_encryption_config); registry.set_context(fees_client); + registry.set_context(jobs); Ok(registry.runner(pool).set_keep_alive(false).run().await?) } @@ -139,25 +138,6 @@ async fn process_all_payout_queues( Ok(()) } -#[job(name = "populate_outbox")] -async fn populate_outbox( - mut current_job: CurrentJob, - outbox: Outbox, - ledger: Ledger, -) -> Result<(), JobError> { - JobExecutor::builder(&mut current_job) - .max_retry_delay(std::time::Duration::from_secs(20)) - .build() - .expect("couldn't build JobExecutor") - .execute(|data| async move { - let data: PopulateOutboxData = data.expect("no PopulateOutboxData available"); - let data = populate_outbox::execute(data, outbox, ledger).await?; - Ok::<_, JobError>(data) - }) - .await?; - Ok(()) -} - #[job(name = "respawn_all_outbox_handlers")] async fn respawn_all_outbox_handlers( mut current_job: CurrentJob, @@ -165,6 +145,7 @@ async fn respawn_all_outbox_handlers( respawn_all_outbox_handlers_delay: delay, .. }: JobsConfig, + jobs: job_crate::Jobs, ) -> Result<(), JobError> { let pool = current_job.pool().clone(); let accounts = Accounts::new(&pool); @@ -173,7 +154,7 @@ async fn respawn_all_outbox_handlers( .expect("couldn't build JobExecutor") .execute(|_| async move { for account in accounts.list().await? { - let _ = spawn_outbox_handler(&pool, account).await; + let _ = crate::job_svc::spawn_outbox_handler(&jobs, account).await; } Ok::<(), JobError>(()) }) @@ -568,29 +549,6 @@ async fn spawn_batch_broadcasting( } } -#[instrument(name = "job.spawn_outbox_handler", skip_all)] -pub async fn spawn_outbox_handler(pool: &sqlx::PgPool, account: Account) -> Result<(), JobError> { - let data = PopulateOutboxData { - account_id: account.id, - journal_id: account.journal_id(), - tracing_data: crate::tracing::extract_tracing_data(), - }; - match JobBuilder::new_with_id(Uuid::from(data.journal_id), "populate_outbox") - .set_channel_name("populate_outbox") - .set_channel_args(&format!("account_id:{}", data.account_id)) - .set_json(&data) - .expect("Couldn't set json") - .spawn(pool) - .await - { - Err(sqlx::Error::Database(err)) if err.message().contains("duplicate key") => Ok(()), - Err(e) => { - crate::tracing::insert_error_fields(tracing::Level::ERROR, &e); - Err(e.into()) - } - Ok(_) => Ok(()), - } -} #[instrument(name = "job.spawn_respawn_all_outbox_handlers", skip_all, fields(error, error.level, error.message), err)] pub async fn spawn_respawn_all_outbox_handlers( pool: &sqlx::PgPool, diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs new file mode 100644 index 00000000..0ac5ca62 --- /dev/null +++ b/src/job_svc/mod.rs @@ -0,0 +1,20 @@ +mod populate_outbox; + +use job_crate::Jobs; + +use crate::{ledger::Ledger, outbox::Outbox}; +use crate::job_svc::populate_outbox::PopulateOutboxJobInit; + +pub use populate_outbox::spawn_outbox_handler; + +pub struct JobSvc; + +impl JobSvc { + pub fn init( + jobs: &Jobs, + outbox: Outbox, + ledger: Ledger, + ) { + jobs.add_initializer(PopulateOutboxJobInit::new(outbox, ledger)); + } +} diff --git a/src/job_svc/populate_outbox.rs b/src/job_svc/populate_outbox.rs new file mode 100644 index 00000000..d1eeadb9 --- /dev/null +++ b/src/job_svc/populate_outbox.rs @@ -0,0 +1,129 @@ +use async_trait::async_trait; +use futures::StreamExt; +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; +use tracing::instrument; + +use job_crate::{ + error::JobError as JobSvcError, + Job, + JobConfig, + JobId, + JobInitializer, + JobRunner, + Jobs, + JobType, + RetrySettings, + CurrentJob, + JobCompletion +}; + +use crate::{ + account::Account, + ledger::Ledger, + outbox::Outbox, + primitives::{AccountId, LedgerJournalId}, +}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct PopulateOutboxJobConfig { + pub account_id: AccountId, + pub journal_id: LedgerJournalId, + #[serde(flatten)] + pub tracing_data: HashMap, +} + +impl JobConfig for PopulateOutboxJobConfig { + type Initializer = PopulateOutboxJobInit; +} + +pub struct PopulateOutboxJobInit { + outbox: Outbox, + ledger: Ledger, +} + +impl PopulateOutboxJobInit { + pub fn new(outbox: Outbox, ledger: Ledger) -> Self { + Self { outbox, ledger } + } +} + +impl JobInitializer for PopulateOutboxJobInit { + fn job_type() -> JobType + where + Self: Sized, + { + JobType::new("populate_outbox") + } + + fn init(&self, job: &Job) -> Result, Box> { + let config: PopulateOutboxJobConfig = job.config()?; + Ok(Box::new(PopulateOutboxJobRunner { + config, + outbox: self.outbox.clone(), + ledger: self.ledger.clone(), + })) + } + + fn retry_on_error_settings() -> RetrySettings + where + Self: Sized, + { + RetrySettings::repeat_indefinitely() + } +} + +pub struct PopulateOutboxJobRunner { + config: PopulateOutboxJobConfig, + outbox: Outbox, + ledger: Ledger, +} + +#[async_trait] +impl JobRunner for PopulateOutboxJobRunner { + async fn run( + &self, + _current_job: CurrentJob, + ) -> Result> { + let mut stream = self + .ledger + .journal_events( + self.config.journal_id, + self.outbox + .last_ledger_event_id(self.config.account_id) + .await?, + ) + .await?; + + while let Some(event) = stream.next().await { + self.outbox + .handle_journal_event(event?, tracing::Span::current()) + .await?; + } + + Ok(JobCompletion::Complete) + } +} + +#[instrument(name = "job.spawn_outbox_handler", skip_all)] +pub async fn spawn_outbox_handler( + jobs: &Jobs, + account: Account, +) -> Result<(), JobSvcError> { + let config = PopulateOutboxJobConfig { + account_id: account.id, + journal_id: account.journal_id(), + tracing_data: crate::tracing::extract_tracing_data(), + }; + + let job_id = JobId::from(uuid::Uuid::from(config.journal_id)); + + match jobs.create_and_spawn(job_id, config).await { + Ok(_) => Ok(()), + Err(JobSvcError::DuplicateUniqueJobType) => Ok(()), + Err(e) => { + crate::tracing::insert_error_fields(tracing::Level::ERROR, &e); + Err(e) + } + } +} diff --git a/src/lib.rs b/src/lib.rs index d05ebcbd..be9c3614 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -14,6 +14,7 @@ pub mod descriptor; mod dev_constants; pub mod fees; mod job; +pub mod job_svc; pub mod ledger; mod outbox; pub mod payout; From 0afc813ae26b3cb35fcef80af0774936ed1d057c Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Thu, 16 Oct 2025 13:29:59 +0530 Subject: [PATCH 02/12] chore: update job migration --- migrations/20250904065521_job_setup.sql | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/migrations/20250904065521_job_setup.sql b/migrations/20250904065521_job_setup.sql index d5a70a90..d5dd16e8 100644 --- a/migrations/20250904065521_job_setup.sql +++ b/migrations/20250904065521_job_setup.sql @@ -21,6 +21,7 @@ CREATE TYPE JobExecutionState AS ENUM ('pending', 'running'); CREATE TABLE job_executions ( id UUID REFERENCES jobs(id) NOT NULL UNIQUE, job_type VARCHAR NOT NULL, + poller_instance_id UUID, attempt_index INT NOT NULL DEFAULT 1, state JobExecutionState NOT NULL DEFAULT 'pending', execution_state_json JSONB, @@ -29,6 +30,10 @@ CREATE TABLE job_executions ( created_at TIMESTAMPTZ NOT NULL ); +CREATE INDEX idx_job_executions_poller_instance + ON job_executions(poller_instance_id) + WHERE state = 'running'; + CREATE OR REPLACE FUNCTION notify_job_execution_insert() RETURNS TRIGGER AS $$ BEGIN PERFORM pg_notify('job_execution', ''); From 82a1380566aecedbf1aced301ff5432348a9ec20 Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Thu, 16 Oct 2025 13:52:08 +0530 Subject: [PATCH 03/12] chore: run cargo fmt --- src/app/mod.rs | 13 ++++++------- src/job_svc/mod.rs | 8 ++------ src/job_svc/populate_outbox.rs | 20 ++++---------------- 3 files changed, 12 insertions(+), 29 deletions(-) diff --git a/src/app/mod.rs b/src/app/mod.rs index 8a28a7a7..cc7738bc 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -9,7 +9,7 @@ use std::collections::HashMap; pub use config::*; use error::*; -use job_crate::{Jobs, JobSvcConfig}; +use job_crate::{JobSvcConfig, Jobs}; use crate::{ account::balance::AccountBalanceSummary, @@ -73,13 +73,12 @@ impl App { .await?; let fees_client = FeesClient::new(config.fees.clone()); - let job_svc_config = JobSvcConfig::builder().pool(pool.clone()).build().expect("couldn't build job_svc_config"); + let job_svc_config = JobSvcConfig::builder() + .pool(pool.clone()) + .build() + .expect("couldn't build job_svc_config"); let mut jobs = Jobs::init(job_svc_config).await?; - JobSvc::init( - &jobs, - outbox.clone(), - ledger.clone(), - ); + JobSvc::init(&jobs, outbox.clone(), ledger.clone()); jobs.start_poll().await?; let runner = job::start_job_runner( diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index 0ac5ca62..9e6ba83d 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -2,19 +2,15 @@ mod populate_outbox; use job_crate::Jobs; -use crate::{ledger::Ledger, outbox::Outbox}; use crate::job_svc::populate_outbox::PopulateOutboxJobInit; +use crate::{ledger::Ledger, outbox::Outbox}; pub use populate_outbox::spawn_outbox_handler; pub struct JobSvc; impl JobSvc { - pub fn init( - jobs: &Jobs, - outbox: Outbox, - ledger: Ledger, - ) { + pub fn init(jobs: &Jobs, outbox: Outbox, ledger: Ledger) { jobs.add_initializer(PopulateOutboxJobInit::new(outbox, ledger)); } } diff --git a/src/job_svc/populate_outbox.rs b/src/job_svc/populate_outbox.rs index d1eeadb9..fe69a383 100644 --- a/src/job_svc/populate_outbox.rs +++ b/src/job_svc/populate_outbox.rs @@ -5,17 +5,8 @@ use std::collections::HashMap; use tracing::instrument; use job_crate::{ - error::JobError as JobSvcError, - Job, - JobConfig, - JobId, - JobInitializer, - JobRunner, - Jobs, - JobType, - RetrySettings, - CurrentJob, - JobCompletion + error::JobError as JobSvcError, CurrentJob, Job, JobCompletion, JobConfig, JobId, + JobInitializer, JobRunner, JobType, Jobs, RetrySettings, }; use crate::{ @@ -94,7 +85,7 @@ impl JobRunner for PopulateOutboxJobRunner { .await?, ) .await?; - + while let Some(event) = stream.next().await { self.outbox .handle_journal_event(event?, tracing::Span::current()) @@ -106,10 +97,7 @@ impl JobRunner for PopulateOutboxJobRunner { } #[instrument(name = "job.spawn_outbox_handler", skip_all)] -pub async fn spawn_outbox_handler( - jobs: &Jobs, - account: Account, -) -> Result<(), JobSvcError> { +pub async fn spawn_outbox_handler(jobs: &Jobs, account: Account) -> Result<(), JobSvcError> { let config = PopulateOutboxJobConfig { account_id: account.id, journal_id: account.journal_id(), From ea251f9bf1d5059f77e67b79eff19a3b6368e97a Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Mon, 20 Oct 2025 15:10:08 +0530 Subject: [PATCH 04/12] refactor: job svc --- Cargo.lock | 90 +++++++++++++++++++++++++++------- Cargo.toml | 2 +- src/app/error.rs | 3 +- src/app/mod.rs | 40 +-------------- src/job/mod.rs | 51 +------------------ src/job_svc/mod.rs | 47 ++++++++++++++++-- src/job_svc/populate_outbox.rs | 24 +-------- 7 files changed, 118 insertions(+), 139 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index f4d8507e..0e3ae12e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -518,9 +518,9 @@ dependencies = [ "hex", "job", "miniscript", - "opentelemetry", + "opentelemetry 0.27.0", "opentelemetry-otlp", - "opentelemetry_sdk", + "opentelemetry_sdk 0.27.0", "prost 0.12.6", "prost-wkt-types", "protobuf-src", @@ -548,7 +548,7 @@ dependencies = [ "tonic-build 0.11.0", "tonic-health", "tracing", - "tracing-opentelemetry", + "tracing-opentelemetry 0.28.0", "tracing-subscriber", "url", "uuid", @@ -1073,27 +1073,31 @@ dependencies = [ [[package]] name = "es-entity" -version = "0.9.0" +version = "0.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "51574a2e447ca840cf2d80b610b8c3e88088dc8477c5d3ba24e20c8f6f0d05af" +checksum = "e936315f0f604d9b34f67684f01963fd6eb816cbcfe5ee2ec770fd7d0b61ee86" dependencies = [ "chrono", "derive_builder", "es-entity-macros", "im", + "opentelemetry 0.30.0", + "opentelemetry_sdk 0.30.0", "pin-project", "serde", "serde_json", "sqlx", "thiserror 2.0.12", + "tracing", + "tracing-opentelemetry 0.31.0", "uuid", ] [[package]] name = "es-entity-macros" -version = "0.9.0" +version = "0.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9beed351da1b6a67b0e82a815aee153bcf2187bf5f4702e6d7e6a48d38aa6f45" +checksum = "c7130dc1d0205b2dc2928a41336b103d8510440b9b0bb9a344c3a0cbf59765a8" dependencies = [ "convert_case", "darling 0.21.3", @@ -1943,9 +1947,9 @@ checksum = "d75a2a4b1b190afb6f5425f10f6a8f959d2ea0b9c2b1d79553551850539e4674" [[package]] name = "job" -version = "0.1.12" +version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fbc97547bab9416fb9e5490474afef9f452611dbff8a478cf9dc3cf5bb093bdb" +checksum = "0e12be50af95af1b818da9de22c96183d96b4b54c672faa69c101905172f598b" dependencies = [ "async-trait", "chrono", @@ -2244,6 +2248,20 @@ dependencies = [ "thiserror 1.0.69", ] +[[package]] +name = "opentelemetry" +version = "0.30.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "aaf416e4cb72756655126f7dd7bb0af49c674f4c1b9903e80c009e0c37e552e6" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "pin-project-lite", + "thiserror 2.0.12", + "tracing", +] + [[package]] name = "opentelemetry-http" version = "0.27.0" @@ -2253,7 +2271,7 @@ dependencies = [ "async-trait", "bytes", "http 1.1.0", - "opentelemetry", + "opentelemetry 0.27.0", "reqwest", ] @@ -2266,10 +2284,10 @@ dependencies = [ "async-trait", "futures-core", "http 1.1.0", - "opentelemetry", + "opentelemetry 0.27.0", "opentelemetry-http", "opentelemetry-proto", - "opentelemetry_sdk", + "opentelemetry_sdk 0.27.0", "prost 0.13.3", "reqwest", "thiserror 1.0.69", @@ -2284,8 +2302,8 @@ version = "0.27.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6e05acbfada5ec79023c85368af14abd0b307c015e9064d249b2a950ef459a6" dependencies = [ - "opentelemetry", - "opentelemetry_sdk", + "opentelemetry 0.27.0", + "opentelemetry_sdk 0.27.0", "prost 0.13.3", "tonic 0.12.3", ] @@ -2302,7 +2320,7 @@ dependencies = [ "futures-util", "glob", "once_cell", - "opentelemetry", + "opentelemetry 0.27.0", "percent-encoding", "rand 0.8.5", "serde_json", @@ -2312,6 +2330,24 @@ dependencies = [ "tracing", ] +[[package]] +name = "opentelemetry_sdk" +version = "0.30.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "11f644aa9e5e31d11896e024305d7e3c98a88884d9f8919dbf37a9991bc47a4b" +dependencies = [ + "futures-channel", + "futures-executor", + "futures-util", + "opentelemetry 0.30.0", + "percent-encoding", + "rand 0.9.2", + "serde_json", + "thiserror 2.0.12", + "tokio", + "tokio-stream", +] + [[package]] name = "parking" version = "2.2.1" @@ -3636,7 +3672,7 @@ dependencies = [ "cached", "chrono", "derive_builder", - "opentelemetry", + "opentelemetry 0.27.0", "rust_decimal", "rusty-money", "serde", @@ -3646,7 +3682,7 @@ dependencies = [ "thiserror 1.0.69", "tokio", "tracing", - "tracing-opentelemetry", + "tracing-opentelemetry 0.28.0", "uuid", ] @@ -4401,8 +4437,8 @@ checksum = "97a971f6058498b5c0f1affa23e7ea202057a7301dbff68e968b2d578bcbd053" dependencies = [ "js-sys", "once_cell", - "opentelemetry", - "opentelemetry_sdk", + "opentelemetry 0.27.0", + "opentelemetry_sdk 0.27.0", "smallvec", "tracing", "tracing-core", @@ -4411,6 +4447,22 @@ dependencies = [ "web-time", ] +[[package]] +name = "tracing-opentelemetry" +version = "0.31.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ddcf5959f39507d0d04d6413119c04f33b623f4f951ebcbdddddfad2d0623a9c" +dependencies = [ + "js-sys", + "once_cell", + "opentelemetry 0.30.0", + "opentelemetry_sdk 0.30.0", + "tracing", + "tracing-core", + "tracing-subscriber", + "web-time", +] + [[package]] name = "tracing-serde" version = "0.2.0" diff --git a/Cargo.toml b/Cargo.toml index b7fc86db..441ed10f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,7 +10,7 @@ fail-on-warnings = [] [dependencies] es-entity = "0.9.0" sqlx-ledger = { version = "0.11.5", features = ["otel"] } -job_crate = { package = "job", version = "0.1.12" } +job_crate = { package = "job", version = "0.2.1" } anyhow = "1.0.82" bitcoincore-rpc = "0.17.0" diff --git a/src/app/error.rs b/src/app/error.rs index a231deb9..9942eb90 100644 --- a/src/app/error.rs +++ b/src/app/error.rs @@ -1,8 +1,6 @@ use chacha20poly1305; use thiserror::Error; -use job_crate::error::JobError as JobSvcError; - use crate::{ address::error::AddressError, batch::error::BatchError, @@ -11,6 +9,7 @@ use crate::{ descriptor::error::DescriptorError, fees::error::FeeEstimationError, job::error::JobError, + job_svc::JobSvcError, ledger::error::LedgerError, outbox::error::OutboxError, payout::error::PayoutError, diff --git a/src/app/mod.rs b/src/app/mod.rs index cc7738bc..dddbbfea 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -9,8 +9,6 @@ use std::collections::HashMap; pub use config::*; use error::*; -use job_crate::{JobSvcConfig, Jobs}; - use crate::{ account::balance::AccountBalanceSummary, address::*, @@ -51,7 +49,6 @@ pub struct App { batch_inclusion: BatchInclusion, pool: sqlx::PgPool, config: AppConfig, - jobs: Jobs, } impl App { @@ -73,13 +70,7 @@ impl App { .await?; let fees_client = FeesClient::new(config.fees.clone()); - let job_svc_config = JobSvcConfig::builder() - .pool(pool.clone()) - .build() - .expect("couldn't build job_svc_config"); - let mut jobs = Jobs::init(job_svc_config).await?; - JobSvc::init(&jobs, outbox.clone(), ledger.clone()); - jobs.start_poll().await?; + let job_svc = JobSvc::init(pool.clone(), outbox.clone(), ledger.clone()).await?; let runner = job::start_job_runner( &pool, @@ -97,7 +88,6 @@ impl App { config.blockchain.clone(), config.signer_encryption.clone(), fees_client.clone(), - jobs.clone(), ) .await?; Self::spawn_sync_all_wallets(pool.clone(), config.jobs.sync_all_wallets_delay).await?; @@ -106,11 +96,6 @@ impl App { config.jobs.process_all_payout_queues_delay, ) .await?; - Self::spawn_respawn_all_outbox_handlers( - pool.clone(), - config.jobs.respawn_all_outbox_handlers_delay, - ) - .await?; let app = Self { outbox, @@ -130,7 +115,6 @@ impl App { batch_inclusion, config, _runner: runner, - jobs, }; if let Some(deprecrated_encryption_key) = app.config.deprecated_encryption_key.as_ref() { app.rotate_encryption_key(deprecrated_encryption_key) @@ -1097,26 +1081,4 @@ impl App { Ok(()) } - #[instrument( - name = "app.spawn_respawn_all_outbox_handlers", - level = "trace", - skip_all, - err - )] - async fn spawn_respawn_all_outbox_handlers( - pool: sqlx::PgPool, - delay: std::time::Duration, - ) -> Result<(), ApplicationError> { - tokio::spawn(async move { - loop { - let _ = job::spawn_respawn_all_outbox_handlers( - &pool, - std::time::Duration::from_secs(1), - ) - .await; - tokio::time::sleep(delay).await; - } - }); - Ok(()) - } } diff --git a/src/job/mod.rs b/src/job/mod.rs index 1ae44ebd..5c822e78 100644 --- a/src/job/mod.rs +++ b/src/job/mod.rs @@ -15,7 +15,7 @@ use tracing::instrument; use uuid::{uuid, Uuid}; use crate::{ - account::*, address::Addresses, app::BlockchainConfig, batch::*, fees::FeesClient, + address::Addresses, app::BlockchainConfig, batch::*, fees::FeesClient, ledger::Ledger, outbox::*, payout::*, payout_queue::*, primitives::*, signing_session::*, utxo::Utxos, wallet::*, xpub::*, }; @@ -30,7 +30,6 @@ use sync_wallet::SyncWalletData; const SYNC_ALL_WALLETS_ID: Uuid = uuid!("00000000-0000-0000-0000-000000000001"); const PROCESS_ALL_PAYOUT_QUEUES_ID: Uuid = uuid!("00000000-0000-0000-0000-000000000002"); -const RESPAWN_ALL_OUTBOX_ID: Uuid = uuid!("00000000-0000-0000-0000-000000000003"); #[allow(clippy::too_many_arguments)] pub async fn start_job_runner( @@ -49,7 +48,6 @@ pub async fn start_job_runner( blockchain_cfg: BlockchainConfig, signer_encryption_config: SignerEncryptionConfig, fees_client: FeesClient, - jobs: job_crate::Jobs, ) -> Result { let mut registry = JobRegistry::new(&[ sync_all_wallets, @@ -60,7 +58,6 @@ pub async fn start_job_runner( batch_wallet_accounting, batch_signing, batch_broadcasting, - respawn_all_outbox_handlers, ]); registry.set_context(config); registry.set_context(blockchain_cfg); @@ -76,7 +73,6 @@ pub async fn start_job_runner( registry.set_context(addresses); registry.set_context(signer_encryption_config); registry.set_context(fees_client); - registry.set_context(jobs); Ok(registry.runner(pool).set_keep_alive(false).run().await?) } @@ -138,31 +134,6 @@ async fn process_all_payout_queues( Ok(()) } -#[job(name = "respawn_all_outbox_handlers")] -async fn respawn_all_outbox_handlers( - mut current_job: CurrentJob, - JobsConfig { - respawn_all_outbox_handlers_delay: delay, - .. - }: JobsConfig, - jobs: job_crate::Jobs, -) -> Result<(), JobError> { - let pool = current_job.pool().clone(); - let accounts = Accounts::new(&pool); - JobExecutor::builder(&mut current_job) - .build() - .expect("couldn't build JobExecutor") - .execute(|_| async move { - for account in accounts.list().await? { - let _ = crate::job_svc::spawn_outbox_handler(&jobs, account).await; - } - Ok::<(), JobError>(()) - }) - .await?; - spawn_respawn_all_outbox_handlers(current_job.pool(), delay).await?; - Ok(()) -} - #[job(name = "sync_wallet")] #[allow(clippy::too_many_arguments)] async fn sync_wallet( @@ -549,26 +520,6 @@ async fn spawn_batch_broadcasting( } } -#[instrument(name = "job.spawn_respawn_all_outbox_handlers", skip_all, fields(error, error.level, error.message), err)] -pub async fn spawn_respawn_all_outbox_handlers( - pool: &sqlx::PgPool, - duration: std::time::Duration, -) -> Result<(), JobError> { - match JobBuilder::new_with_id(RESPAWN_ALL_OUTBOX_ID, "respawn_all_outbox_handlers") - .set_channel_name("respawn_all_outbox_handlers") - .set_delay(duration) - .spawn(pool) - .await - { - Err(sqlx::Error::Database(err)) if err.message().contains("duplicate key") => Ok(()), - Err(e) => { - crate::tracing::insert_error_fields(tracing::Level::ERROR, &e); - Err(e.into()) - } - Ok(_) => Ok(()), - } -} - fn schedule_payout_queue_channel_arg(payout_queue_id: PayoutQueueId) -> String { format!("payout_queue_id:{payout_queue_id}") } diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index 9e6ba83d..bb5e76ea 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -1,16 +1,53 @@ +pub mod error; mod populate_outbox; -use job_crate::Jobs; +use job_crate::{error::JobError as JobCrateError, JobId, JobSvcConfig, Jobs}; +use tracing::instrument; use crate::job_svc::populate_outbox::PopulateOutboxJobInit; -use crate::{ledger::Ledger, outbox::Outbox}; +use crate::{account::Account, ledger::Ledger, outbox::Outbox}; -pub use populate_outbox::spawn_outbox_handler; +pub use error::JobSvcError; +pub use populate_outbox::PopulateOutboxJobConfig; -pub struct JobSvc; +#[derive(Clone)] +pub struct JobSvc { + jobs: Jobs, +} impl JobSvc { - pub fn init(jobs: &Jobs, outbox: Outbox, ledger: Ledger) { + pub async fn init( + pool: sqlx::PgPool, + outbox: Outbox, + ledger: Ledger, + ) -> Result { + let job_svc_config = JobSvcConfig::builder() + .pool(pool) + .build() + .map_err(|e| JobSvcError::ConfigBuild(e.to_string()))?; + + let mut jobs = Jobs::init(job_svc_config).await?; jobs.add_initializer(PopulateOutboxJobInit::new(outbox, ledger)); + jobs.start_poll().await?; + + Ok(Self { jobs }) + } + + #[instrument(name = "job_svc.spawn_outbox_handler_in_op", skip_all)] + pub async fn spawn_outbox_handler_in_op( + &self, + op: &mut impl es_entity::AtomicOperation, + account: Account, + ) -> Result<(), JobSvcError> { + let config = PopulateOutboxJobConfig { + account_id: account.id, + journal_id: account.journal_id(), + tracing_data: crate::tracing::extract_tracing_data(), + }; + + let job_id = JobId::from(uuid::Uuid::from(config.journal_id)); + + self.jobs.create_and_spawn_in_op(op, job_id, config).await? + Ok(()) } } diff --git a/src/job_svc/populate_outbox.rs b/src/job_svc/populate_outbox.rs index fe69a383..0457176c 100644 --- a/src/job_svc/populate_outbox.rs +++ b/src/job_svc/populate_outbox.rs @@ -2,15 +2,12 @@ use async_trait::async_trait; use futures::StreamExt; use serde::{Deserialize, Serialize}; use std::collections::HashMap; -use tracing::instrument; use job_crate::{ - error::JobError as JobSvcError, CurrentJob, Job, JobCompletion, JobConfig, JobId, - JobInitializer, JobRunner, JobType, Jobs, RetrySettings, + CurrentJob, Job, JobCompletion, JobConfig, JobInitializer, JobRunner, JobType, RetrySettings, }; use crate::{ - account::Account, ledger::Ledger, outbox::Outbox, primitives::{AccountId, LedgerJournalId}, @@ -96,22 +93,3 @@ impl JobRunner for PopulateOutboxJobRunner { } } -#[instrument(name = "job.spawn_outbox_handler", skip_all)] -pub async fn spawn_outbox_handler(jobs: &Jobs, account: Account) -> Result<(), JobSvcError> { - let config = PopulateOutboxJobConfig { - account_id: account.id, - journal_id: account.journal_id(), - tracing_data: crate::tracing::extract_tracing_data(), - }; - - let job_id = JobId::from(uuid::Uuid::from(config.journal_id)); - - match jobs.create_and_spawn(job_id, config).await { - Ok(_) => Ok(()), - Err(JobSvcError::DuplicateUniqueJobType) => Ok(()), - Err(e) => { - crate::tracing::insert_error_fields(tracing::Level::ERROR, &e); - Err(e) - } - } -} From 061d5b032019413b1ed80ce56813e49ec99e311b Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Tue, 21 Oct 2025 13:11:09 +0530 Subject: [PATCH 05/12] refactor: borrow job svc from app for adminapp --- src/admin/app.rs | 22 ++++++++++++++++++---- src/admin/error.rs | 6 ++++-- src/admin/mod.rs | 8 +++++--- src/app/mod.rs | 7 ++++++- src/cli/mod.rs | 12 ++++++++++-- src/job/mod.rs | 6 +++--- src/job_svc/mod.rs | 10 +++++++--- src/job_svc/populate_outbox.rs | 1 - 8 files changed, 53 insertions(+), 19 deletions(-) diff --git a/src/admin/app.rs b/src/admin/app.rs index ec3d7f42..5367b072 100644 --- a/src/admin/app.rs +++ b/src/admin/app.rs @@ -1,7 +1,9 @@ use tracing::instrument; use super::{error::*, keys::*}; -use crate::{account::*, dev_constants, ledger::Ledger, primitives::bitcoin, profile::*}; +use crate::{ + account::*, dev_constants, job_svc::JobSvc, ledger::Ledger, primitives::bitcoin, profile::*, +}; const BOOTSTRAP_KEY_NAME: &str = "admin_bootstrap_key"; @@ -11,16 +13,18 @@ pub struct AdminApp { profiles: Profiles, ledger: Ledger, network: bitcoin::Network, + job_svc: JobSvc, } impl AdminApp { - pub fn new(pool: sqlx::PgPool, network: bitcoin::Network) -> Self { + pub fn new(pool: sqlx::PgPool, network: bitcoin::Network, job_svc: JobSvc) -> Self { Self { keys: AdminApiKeys::new(&pool), accounts: Accounts::new(&pool), profiles: Profiles::new(&pool), ledger: Ledger::new(&pool), network, + job_svc, } } } @@ -43,7 +47,7 @@ impl AdminApp { .await?; let new_profile = NewProfile::builder() .account_id(account.id) - .name(account.name) + .name(account.name.clone()) .build() .expect("Couldn't build NewProfile"); let profile = self.profiles.create_in_op(&mut op, new_profile).await?; @@ -51,6 +55,11 @@ impl AdminApp { .profiles .create_key_for_profile_in_op(&mut op, profile, true) .await?; + + self.job_svc + .spawn_outbox_handler_in_op(&mut op, account) + .await?; + op.commit().await?; Ok((admin_key, profile_key)) } @@ -81,7 +90,7 @@ impl AdminApp { .await?; let new_profile = NewProfile::builder() .account_id(account.id) - .name(account.name) + .name(account.name.clone()) .build() .expect("Couldn't build NewProfile"); let profile = self.profiles.create_in_op(&mut op, new_profile).await?; @@ -89,6 +98,11 @@ impl AdminApp { .profiles .create_key_for_profile_in_op(&mut op, profile, false) .await?; + + self.job_svc + .spawn_outbox_handler_in_op(&mut op, account) + .await?; + op.commit().await?; Ok(key) } diff --git a/src/admin/error.rs b/src/admin/error.rs index 71394487..d26d2902 100644 --- a/src/admin/error.rs +++ b/src/admin/error.rs @@ -1,8 +1,8 @@ use thiserror::Error; use crate::{ - account::error::AccountError, app::error::ApplicationError, ledger::error::LedgerError, - profile::error::ProfileError, + account::error::AccountError, app::error::ApplicationError, job_svc::JobSvcError, + ledger::error::LedgerError, profile::error::ProfileError, }; #[allow(clippy::large_enum_variant)] @@ -22,6 +22,8 @@ pub enum AdminApiError { ProfileError(#[from] ProfileError), #[error("{0}")] LedgerError(#[from] LedgerError), + #[error("{0}")] + JobSvcError(#[from] JobSvcError), #[error("AdminApiError - DevBootstrapError: {0}")] DevBootstrapError(#[from] anyhow::Error), } diff --git a/src/admin/mod.rs b/src/admin/mod.rs index 4603e077..8fc7caea 100644 --- a/src/admin/mod.rs +++ b/src/admin/mod.rs @@ -4,7 +4,7 @@ pub mod error; mod keys; mod server; -use crate::{dev_constants, primitives::bitcoin, token_store}; +use crate::{dev_constants, job_svc::JobSvc, primitives::bitcoin, token_store}; pub use app::*; pub use config::*; @@ -17,8 +17,9 @@ pub async fn run_dev( config: AdminApiConfig, network: bitcoin::Network, bria_home: String, + job_svc: JobSvc, ) -> Result<(), AdminApiError> { - let app = AdminApp::new(pool, network); + let app = AdminApp::new(pool, network, job_svc); let (admin_key, profile_key) = app.dev_bootstrap().await?; token_store::store_admin_token(&bria_home, &admin_key.key)?; println!("Admin API key"); @@ -44,8 +45,9 @@ pub async fn run( pool: sqlx::PgPool, config: AdminApiConfig, network: bitcoin::Network, + job_svc: JobSvc, ) -> Result<(), AdminApiError> { - let app = AdminApp::new(pool, network); + let app = AdminApp::new(pool, network, job_svc); server::start(config, app).await?; Ok(()) } diff --git a/src/app/mod.rs b/src/app/mod.rs index dddbbfea..eef62127 100644 --- a/src/app/mod.rs +++ b/src/app/mod.rs @@ -33,6 +33,7 @@ use crate::{ #[allow(dead_code)] pub struct App { _runner: JobRunnerHandle, + job_svc: JobSvc, outbox: Outbox, profiles: Profiles, xpubs: XPubs, @@ -98,6 +99,7 @@ impl App { .await?; let app = Self { + job_svc, outbox, profiles: Profiles::new(&pool), xpubs, @@ -127,6 +129,10 @@ impl App { self.config.blockchain.network } + pub fn job_svc(&self) -> &JobSvc { + &self.job_svc + } + #[instrument(name = "app.authenticate", skip_all, err)] pub async fn authenticate(&self, key: &str) -> Result { let profile = self.profiles.find_by_key(key).await?; @@ -1080,5 +1086,4 @@ impl App { }); Ok(()) } - } diff --git a/src/cli/mod.rs b/src/cli/mod.rs index 6d926acb..431896b1 100644 --- a/src/cli/mod.rs +++ b/src/cli/mod.rs @@ -1071,23 +1071,31 @@ async fn run_cmd( let mut handles = Vec::new(); let pool = init_pool(&db).await?; + // Initialize App first (creates JobSvc internally with its infrastructure) + let main_app = crate::app::App::run(pool.clone(), app.clone()).await?; + + // Extract JobSvc from App to pass to AdminApp + let admin_job_svc = main_app.job_svc().clone(); + let admin_send = send.clone(); let admin_pool = pool.clone(); let network = app.blockchain.network; handles.push(tokio::spawn(async move { let _ = admin_send.try_send(if dev { - super::admin::run_dev(admin_pool, admin, network, bria_home) + super::admin::run_dev(admin_pool, admin, network, bria_home, admin_job_svc) .await .context("Admin server error") } else { - super::admin::run(admin_pool, admin, network) + super::admin::run(admin_pool, admin, network, admin_job_svc) .await .context("Admin server error") }); })); + let api_send = send.clone(); handles.push(tokio::spawn(async move { let _ = api_send.try_send(if dev { + // API runs with the main_app (which already has JobSvc) super::api::run_dev(pool, api, app, dev_xpub, dev_derivation) .await .context("Api server error") diff --git a/src/job/mod.rs b/src/job/mod.rs index 5c822e78..f0284ab3 100644 --- a/src/job/mod.rs +++ b/src/job/mod.rs @@ -15,9 +15,9 @@ use tracing::instrument; use uuid::{uuid, Uuid}; use crate::{ - address::Addresses, app::BlockchainConfig, batch::*, fees::FeesClient, - ledger::Ledger, outbox::*, payout::*, payout_queue::*, primitives::*, signing_session::*, - utxo::Utxos, wallet::*, xpub::*, + address::Addresses, app::BlockchainConfig, batch::*, fees::FeesClient, ledger::Ledger, + outbox::*, payout::*, payout_queue::*, primitives::*, signing_session::*, utxo::Utxos, + wallet::*, xpub::*, }; use batch_broadcasting::BatchBroadcastingData; use batch_signing::BatchSigningData; diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index bb5e76ea..f418c649 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -1,7 +1,7 @@ pub mod error; mod populate_outbox; -use job_crate::{error::JobError as JobCrateError, JobId, JobSvcConfig, Jobs}; +use job_crate::{JobId, JobSvcConfig, Jobs}; use tracing::instrument; use crate::job_svc::populate_outbox::PopulateOutboxJobInit; @@ -25,7 +25,7 @@ impl JobSvc { .pool(pool) .build() .map_err(|e| JobSvcError::ConfigBuild(e.to_string()))?; - + let mut jobs = Jobs::init(job_svc_config).await?; jobs.add_initializer(PopulateOutboxJobInit::new(outbox, ledger)); jobs.start_poll().await?; @@ -33,6 +33,10 @@ impl JobSvc { Ok(Self { jobs }) } + pub fn jobs(&self) -> &Jobs { + &self.jobs + } + #[instrument(name = "job_svc.spawn_outbox_handler_in_op", skip_all)] pub async fn spawn_outbox_handler_in_op( &self, @@ -47,7 +51,7 @@ impl JobSvc { let job_id = JobId::from(uuid::Uuid::from(config.journal_id)); - self.jobs.create_and_spawn_in_op(op, job_id, config).await? + self.jobs.create_and_spawn_in_op(op, job_id, config).await?; Ok(()) } } diff --git a/src/job_svc/populate_outbox.rs b/src/job_svc/populate_outbox.rs index 0457176c..a76a9e84 100644 --- a/src/job_svc/populate_outbox.rs +++ b/src/job_svc/populate_outbox.rs @@ -92,4 +92,3 @@ impl JobRunner for PopulateOutboxJobRunner { Ok(JobCompletion::Complete) } } - From f459d9f4ff95676139cc8cb7225b328876db294b Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Tue, 21 Oct 2025 13:14:43 +0530 Subject: [PATCH 06/12] fix: commit job svc error.rs --- src/job_svc/error.rs | 9 +++++++++ 1 file changed, 9 insertions(+) create mode 100644 src/job_svc/error.rs diff --git a/src/job_svc/error.rs b/src/job_svc/error.rs new file mode 100644 index 00000000..95e1f149 --- /dev/null +++ b/src/job_svc/error.rs @@ -0,0 +1,9 @@ +use thiserror::Error; + +#[derive(Error, Debug)] +pub enum JobSvcError { + #[error("JobSvcError - ConfigBuild: {0}")] + ConfigBuild(String), + #[error("JobSvcError - JobCrateError: {0}")] + JobCrateError(#[from] job_crate::error::JobError), +} From c9de08be65152ce92c8cab620e87541f13f36709 Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Tue, 21 Oct 2025 15:16:02 +0530 Subject: [PATCH 07/12] fix: integration test --- src/job_svc/error.rs | 4 ++++ src/job_svc/mod.rs | 21 +++++++++++++++++++++ tests/helpers.rs | 7 +++++-- 3 files changed, 30 insertions(+), 2 deletions(-) diff --git a/src/job_svc/error.rs b/src/job_svc/error.rs index 95e1f149..5b5d9f6a 100644 --- a/src/job_svc/error.rs +++ b/src/job_svc/error.rs @@ -6,4 +6,8 @@ pub enum JobSvcError { ConfigBuild(String), #[error("JobSvcError - JobCrateError: {0}")] JobCrateError(#[from] job_crate::error::JobError), + #[error("JobSvcError - LedgerError: {0}")] + Ledger(#[from] crate::ledger::error::LedgerError), + #[error("JobSvcError - OutboxError: {0}")] + Outbox(#[from] crate::outbox::error::OutboxError), } diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index f418c649..bc6a98d0 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -33,6 +33,27 @@ impl JobSvc { Ok(Self { jobs }) } + /// Initialize JobSvc for testing - creates its own infrastructure + pub async fn init_for_test(pool: sqlx::PgPool) -> Result { + use crate::{ + address::Addresses, batch_inclusion::BatchInclusion, payout::Payouts, + payout_queue::PayoutQueues, + }; + + let ledger = Ledger::init(&pool).await?; + let addresses = Addresses::new(&pool); + let payouts = Payouts::new(&pool); + let payout_queues = PayoutQueues::new(&pool); + let batch_inclusion = BatchInclusion::new(pool.clone(), payout_queues); + let outbox = Outbox::init( + &pool, + crate::outbox::Augmenter::new(&addresses, &payouts, &batch_inclusion), + ) + .await?; + + Self::init(pool, outbox, ledger).await + } + pub fn jobs(&self) -> &Jobs { &self.jobs } diff --git a/tests/helpers.rs b/tests/helpers.rs index 5fcc4eb6..c091f8de 100644 --- a/tests/helpers.rs +++ b/tests/helpers.rs @@ -14,7 +14,7 @@ use bdk::{ miniscript::Segwitv0, }; use bitcoincore_rpc::{Client as BitcoindClient, RpcApi}; -use bria::{admin::*, primitives::*, profile::*, xpub::*}; +use bria::{admin::*, job_svc::JobSvc, primitives::*, profile::*, xpub::*}; use rand::distributions::{Alphanumeric, DistString}; pub async fn init_pool() -> anyhow::Result { @@ -32,7 +32,10 @@ pub async fn create_test_account(pool: &sqlx::PgPool) -> anyhow::Result "TEST_{}", Alphanumeric.sample_string(&mut rand::thread_rng(), 32) ); - let app = AdminApp::new(pool.clone(), bitcoin::Network::Regtest); + + let job_svc = JobSvc::init_for_test(pool.clone()).await?; + + let app = AdminApp::new(pool.clone(), bitcoin::Network::Regtest, job_svc); let profile_key = app.create_account(name.clone()).await?; Ok(Profiles::new(pool).find_by_key(&profile_key.key).await?) From 7034e4e5ff6f424ed7fcb9666b0ead7a96c116a0 Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Tue, 21 Oct 2025 16:33:45 +0530 Subject: [PATCH 08/12] fix: return anyhow from test func --- src/job_svc/error.rs | 4 ---- src/job_svc/mod.rs | 4 ++-- 2 files changed, 2 insertions(+), 6 deletions(-) diff --git a/src/job_svc/error.rs b/src/job_svc/error.rs index 5b5d9f6a..95e1f149 100644 --- a/src/job_svc/error.rs +++ b/src/job_svc/error.rs @@ -6,8 +6,4 @@ pub enum JobSvcError { ConfigBuild(String), #[error("JobSvcError - JobCrateError: {0}")] JobCrateError(#[from] job_crate::error::JobError), - #[error("JobSvcError - LedgerError: {0}")] - Ledger(#[from] crate::ledger::error::LedgerError), - #[error("JobSvcError - OutboxError: {0}")] - Outbox(#[from] crate::outbox::error::OutboxError), } diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index bc6a98d0..e4c02aae 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -34,7 +34,7 @@ impl JobSvc { } /// Initialize JobSvc for testing - creates its own infrastructure - pub async fn init_for_test(pool: sqlx::PgPool) -> Result { + pub async fn init_for_test(pool: sqlx::PgPool) -> anyhow::Result { use crate::{ address::Addresses, batch_inclusion::BatchInclusion, payout::Payouts, payout_queue::PayoutQueues, @@ -51,7 +51,7 @@ impl JobSvc { ) .await?; - Self::init(pool, outbox, ledger).await + Self::init(pool, outbox, ledger).await.map_err(Into::into) } pub fn jobs(&self) -> &Jobs { From f37ef780dce26746d3696830b1ccd610ab1973e2 Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Tue, 21 Oct 2025 20:31:08 +0530 Subject: [PATCH 09/12] refactor: spawn outbox handler and admin app test --- src/admin/app.rs | 4 ++-- src/cli/mod.rs | 3 --- src/job_svc/error.rs | 2 -- src/job_svc/mod.rs | 35 +++++++---------------------------- src/lib.rs | 4 ++-- tests/helpers.rs | 10 +++++++++- 6 files changed, 20 insertions(+), 38 deletions(-) diff --git a/src/admin/app.rs b/src/admin/app.rs index 5367b072..fe5e4e19 100644 --- a/src/admin/app.rs +++ b/src/admin/app.rs @@ -57,7 +57,7 @@ impl AdminApp { .await?; self.job_svc - .spawn_outbox_handler_in_op(&mut op, account) + .spawn_outbox_handler_in_op(&mut op, account.id, account.journal_id()) .await?; op.commit().await?; @@ -100,7 +100,7 @@ impl AdminApp { .await?; self.job_svc - .spawn_outbox_handler_in_op(&mut op, account) + .spawn_outbox_handler_in_op(&mut op, account.id, account.journal_id()) .await?; op.commit().await?; diff --git a/src/cli/mod.rs b/src/cli/mod.rs index 431896b1..ee1e7e5f 100644 --- a/src/cli/mod.rs +++ b/src/cli/mod.rs @@ -1071,10 +1071,8 @@ async fn run_cmd( let mut handles = Vec::new(); let pool = init_pool(&db).await?; - // Initialize App first (creates JobSvc internally with its infrastructure) let main_app = crate::app::App::run(pool.clone(), app.clone()).await?; - // Extract JobSvc from App to pass to AdminApp let admin_job_svc = main_app.job_svc().clone(); let admin_send = send.clone(); @@ -1095,7 +1093,6 @@ async fn run_cmd( let api_send = send.clone(); handles.push(tokio::spawn(async move { let _ = api_send.try_send(if dev { - // API runs with the main_app (which already has JobSvc) super::api::run_dev(pool, api, app, dev_xpub, dev_derivation) .await .context("Api server error") diff --git a/src/job_svc/error.rs b/src/job_svc/error.rs index 95e1f149..db0ce363 100644 --- a/src/job_svc/error.rs +++ b/src/job_svc/error.rs @@ -2,8 +2,6 @@ use thiserror::Error; #[derive(Error, Debug)] pub enum JobSvcError { - #[error("JobSvcError - ConfigBuild: {0}")] - ConfigBuild(String), #[error("JobSvcError - JobCrateError: {0}")] JobCrateError(#[from] job_crate::error::JobError), } diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index e4c02aae..87bccd8a 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -4,11 +4,10 @@ mod populate_outbox; use job_crate::{JobId, JobSvcConfig, Jobs}; use tracing::instrument; -use crate::job_svc::populate_outbox::PopulateOutboxJobInit; -use crate::{account::Account, ledger::Ledger, outbox::Outbox}; +use crate::job_svc::populate_outbox::{PopulateOutboxJobInit, PopulateOutboxJobConfig}; +use crate::{ledger::Ledger, outbox::Outbox, primitives::{AccountId, LedgerJournalId}}; pub use error::JobSvcError; -pub use populate_outbox::PopulateOutboxJobConfig; #[derive(Clone)] pub struct JobSvc { @@ -24,7 +23,7 @@ impl JobSvc { let job_svc_config = JobSvcConfig::builder() .pool(pool) .build() - .map_err(|e| JobSvcError::ConfigBuild(e.to_string()))?; + .expect("Couldn't build JobSvcConfig"); let mut jobs = Jobs::init(job_svc_config).await?; jobs.add_initializer(PopulateOutboxJobInit::new(outbox, ledger)); @@ -33,27 +32,6 @@ impl JobSvc { Ok(Self { jobs }) } - /// Initialize JobSvc for testing - creates its own infrastructure - pub async fn init_for_test(pool: sqlx::PgPool) -> anyhow::Result { - use crate::{ - address::Addresses, batch_inclusion::BatchInclusion, payout::Payouts, - payout_queue::PayoutQueues, - }; - - let ledger = Ledger::init(&pool).await?; - let addresses = Addresses::new(&pool); - let payouts = Payouts::new(&pool); - let payout_queues = PayoutQueues::new(&pool); - let batch_inclusion = BatchInclusion::new(pool.clone(), payout_queues); - let outbox = Outbox::init( - &pool, - crate::outbox::Augmenter::new(&addresses, &payouts, &batch_inclusion), - ) - .await?; - - Self::init(pool, outbox, ledger).await.map_err(Into::into) - } - pub fn jobs(&self) -> &Jobs { &self.jobs } @@ -62,11 +40,12 @@ impl JobSvc { pub async fn spawn_outbox_handler_in_op( &self, op: &mut impl es_entity::AtomicOperation, - account: Account, + account_id: AccountId, + journal_id: LedgerJournalId, ) -> Result<(), JobSvcError> { let config = PopulateOutboxJobConfig { - account_id: account.id, - journal_id: account.journal_id(), + account_id, + journal_id, tracing_data: crate::tracing::extract_tracing_data(), }; diff --git a/src/lib.rs b/src/lib.rs index be9c3614..fe305426 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -7,7 +7,7 @@ pub mod admin; mod api; pub mod app; pub mod batch; -mod batch_inclusion; +pub mod batch_inclusion; pub mod bdk; pub mod cli; pub mod descriptor; @@ -16,7 +16,7 @@ pub mod fees; mod job; pub mod job_svc; pub mod ledger; -mod outbox; +pub mod outbox; pub mod payout; pub mod payout_queue; pub mod primitives; diff --git a/tests/helpers.rs b/tests/helpers.rs index c091f8de..f0c9c160 100644 --- a/tests/helpers.rs +++ b/tests/helpers.rs @@ -17,6 +17,7 @@ use bitcoincore_rpc::{Client as BitcoindClient, RpcApi}; use bria::{admin::*, job_svc::JobSvc, primitives::*, profile::*, xpub::*}; use rand::distributions::{Alphanumeric, DistString}; +use bria::{address::Addresses, batch_inclusion::BatchInclusion, payout::Payouts, payout_queue::PayoutQueues, outbox::{Outbox, Augmenter}, ledger::Ledger}; pub async fn init_pool() -> anyhow::Result { let pg_host = std::env::var("PG_HOST").unwrap_or("localhost".to_string()); let pg_con = format!("postgres://user:password@{pg_host}:5432/pg"); @@ -33,7 +34,14 @@ pub async fn create_test_account(pool: &sqlx::PgPool) -> anyhow::Result Alphanumeric.sample_string(&mut rand::thread_rng(), 32) ); - let job_svc = JobSvc::init_for_test(pool.clone()).await?; + let addresses = Addresses::new(pool); + let payouts = Payouts::new(pool); + let payout_queues = PayoutQueues::new(pool); + let batch_inclusion = BatchInclusion::new(pool.clone(), payout_queues); + let augmenter = Augmenter::new(&addresses, &payouts, &batch_inclusion); + let outbox = Outbox::init(pool, augmenter).await?; + let ledger = Ledger::init(&pool.clone()).await?; + let job_svc = JobSvc::init(pool.clone(), outbox, ledger).await?; let app = AdminApp::new(pool.clone(), bitcoin::Network::Regtest, job_svc); From 8fc770db9ffa3422f7b4fc68f1c3543e8a639c8e Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Tue, 21 Oct 2025 20:34:07 +0530 Subject: [PATCH 10/12] chore: run cargo fmt --- src/job_svc/mod.rs | 8 ++++++-- tests/helpers.rs | 9 ++++++++- 2 files changed, 14 insertions(+), 3 deletions(-) diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index 87bccd8a..e38ec384 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -4,8 +4,12 @@ mod populate_outbox; use job_crate::{JobId, JobSvcConfig, Jobs}; use tracing::instrument; -use crate::job_svc::populate_outbox::{PopulateOutboxJobInit, PopulateOutboxJobConfig}; -use crate::{ledger::Ledger, outbox::Outbox, primitives::{AccountId, LedgerJournalId}}; +use crate::job_svc::populate_outbox::{PopulateOutboxJobConfig, PopulateOutboxJobInit}; +use crate::{ + ledger::Ledger, + outbox::Outbox, + primitives::{AccountId, LedgerJournalId}, +}; pub use error::JobSvcError; diff --git a/tests/helpers.rs b/tests/helpers.rs index f0c9c160..85dddc9c 100644 --- a/tests/helpers.rs +++ b/tests/helpers.rs @@ -17,7 +17,14 @@ use bitcoincore_rpc::{Client as BitcoindClient, RpcApi}; use bria::{admin::*, job_svc::JobSvc, primitives::*, profile::*, xpub::*}; use rand::distributions::{Alphanumeric, DistString}; -use bria::{address::Addresses, batch_inclusion::BatchInclusion, payout::Payouts, payout_queue::PayoutQueues, outbox::{Outbox, Augmenter}, ledger::Ledger}; +use bria::{ + address::Addresses, + batch_inclusion::BatchInclusion, + ledger::Ledger, + outbox::{Augmenter, Outbox}, + payout::Payouts, + payout_queue::PayoutQueues, +}; pub async fn init_pool() -> anyhow::Result { let pg_host = std::env::var("PG_HOST").unwrap_or("localhost".to_string()); let pg_con = format!("postgres://user:password@{pg_host}:5432/pg"); From e53cf596904179bd42552719fb21d9c9b115e9b0 Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Wed, 22 Oct 2025 13:26:22 +0530 Subject: [PATCH 11/12] chore: remove tracing data attr from job config --- src/job_svc/mod.rs | 1 - src/job_svc/populate_outbox.rs | 3 --- 2 files changed, 4 deletions(-) diff --git a/src/job_svc/mod.rs b/src/job_svc/mod.rs index e38ec384..c4bd25b3 100644 --- a/src/job_svc/mod.rs +++ b/src/job_svc/mod.rs @@ -50,7 +50,6 @@ impl JobSvc { let config = PopulateOutboxJobConfig { account_id, journal_id, - tracing_data: crate::tracing::extract_tracing_data(), }; let job_id = JobId::from(uuid::Uuid::from(config.journal_id)); diff --git a/src/job_svc/populate_outbox.rs b/src/job_svc/populate_outbox.rs index a76a9e84..c1aa49d9 100644 --- a/src/job_svc/populate_outbox.rs +++ b/src/job_svc/populate_outbox.rs @@ -1,7 +1,6 @@ use async_trait::async_trait; use futures::StreamExt; use serde::{Deserialize, Serialize}; -use std::collections::HashMap; use job_crate::{ CurrentJob, Job, JobCompletion, JobConfig, JobInitializer, JobRunner, JobType, RetrySettings, @@ -17,8 +16,6 @@ use crate::{ pub struct PopulateOutboxJobConfig { pub account_id: AccountId, pub journal_id: LedgerJournalId, - #[serde(flatten)] - pub tracing_data: HashMap, } impl JobConfig for PopulateOutboxJobConfig { From 4c8606badfdf7ce89c7b8012a25d068783897597 Mon Sep 17 00:00:00 2001 From: Kartik Shah Date: Wed, 22 Oct 2025 17:46:18 +0530 Subject: [PATCH 12/12] chore: remove populate outbox from legacy job svc --- src/job/populate_outbox.rs | 36 ------------------------------------ 1 file changed, 36 deletions(-) delete mode 100644 src/job/populate_outbox.rs diff --git a/src/job/populate_outbox.rs b/src/job/populate_outbox.rs deleted file mode 100644 index 209d8382..00000000 --- a/src/job/populate_outbox.rs +++ /dev/null @@ -1,36 +0,0 @@ -use futures::StreamExt; -use serde::{Deserialize, Serialize}; -use tracing::instrument; - -use super::error::JobError; -use crate::{ledger::*, outbox::*, primitives::*}; - -use std::collections::HashMap; - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct PopulateOutboxData { - pub(super) account_id: AccountId, - pub(super) journal_id: LedgerJournalId, - #[serde(flatten)] - pub(super) tracing_data: HashMap, -} - -#[instrument("job.populate_outbox", skip(outbox, ledger))] -pub async fn execute( - data: PopulateOutboxData, - outbox: Outbox, - ledger: Ledger, -) -> Result { - let mut stream = ledger - .journal_events( - data.journal_id, - outbox.last_ledger_event_id(data.account_id).await?, - ) - .await?; - while let Some(event) = stream.next().await { - outbox - .handle_journal_event(event?, tracing::Span::current()) - .await?; - } - Ok(data) -}