diff --git a/roles/Cargo.lock b/roles/Cargo.lock index 4263188502..f145aca4e0 100644 --- a/roles/Cargo.lock +++ b/roles/Cargo.lock @@ -58,7 +58,7 @@ version = "0.7.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "891477e0c6a8957309ee5c45a6368af3ae14bb510732d2684ffa19af310920f9" dependencies = [ - "getrandom 0.2.15", + "getrandom 0.2.16", "once_cell", "version_check", ] @@ -75,6 +75,15 @@ dependencies = [ "zerocopy 0.7.35", ] +[[package]] +name = "aho-corasick" +version = "1.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e60d3430d3a69478ad0993f19238d2df97c507009a52b3c10addcd7f6bcb916" +dependencies = [ + "memchr", +] + [[package]] name = "allocator-api2" version = "0.2.21" @@ -133,9 +142,9 @@ dependencies = [ [[package]] name = "anyhow" -version = "1.0.97" +version = "1.0.98" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "dcfed56ad506cb2c684a14971b8861fdc3baaaae314b9e5f9bb532cbe3ba7a4f" +checksum = "e16d2d3311acee920a9eb8d33b8cbc1787ce4a264e85f964c2404b969bdcd487" [[package]] name = "arraydeque" @@ -174,14 +183,15 @@ dependencies = [ [[package]] name = "async-executor" -version = "1.13.1" +version = "1.13.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "30ca9a001c1e8ba5149f91a74362376cc6bc5b919d92d988668657bd570bdcec" +checksum = "bb812ffb58524bdd10860d7d974e2f01cc0950c2438a74ee5ec2e2280c6c4ffa" dependencies = [ "async-task", "concurrent-queue", "fastrand", "futures-lite", + "pin-project-lite", "slab", ] @@ -249,7 +259,7 @@ checksum = "3b43422f69d8ff38f95f1b2bb76517c91589a924d1559a0e935d7c8ce0274c11" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -292,7 +302,7 @@ checksum = "e539d3fca749fcee5236ab05e93a52867dd549cc157c8cb7f99595f3cedffdb5" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -537,11 +547,51 @@ version = "1.10.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d71b6127be86fdcfddb610f7182ac57211d4b18a3e9c82eb2d17662f2227ad6a" +[[package]] +name = "capnp" +version = "0.20.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "053b81915c2ce1629b8fb964f578b18cb39b23ef9d5b24120d0dfc959569a1d9" +dependencies = [ + "embedded-io", +] + +[[package]] +name = "capnp-futures" +version = "0.20.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b70b0d44372d42654e3efac38c1643c7b0f9d3a9e9b72b635f942ff3f17e891" +dependencies = [ + "capnp", + "futures-channel", + "futures-util", +] + +[[package]] +name = "capnp-rpc" +version = "0.20.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d5a945dd7eac211c30763aa1dbf86ed8e58129d01442d4d2d516facfdb859a1e" +dependencies = [ + "capnp", + "capnp-futures", + "futures", +] + +[[package]] +name = "capnpc" +version = "0.20.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1aa3d5f01e69ed11656d2c7c47bf34327ea9bfb5c85c7de787fcd7b6c5e45b61" +dependencies = [ + "capnp", +] + [[package]] name = "cc" -version = "1.2.17" +version = "1.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1fcb57c740ae1daf453ae85f16e37396f672b039e00d9d866e07ddb24e328e3a" +checksum = "8691782945451c1c383942c4874dbe63814f61cb57ef773cda2972682b7bb3c0" dependencies = [ "shlex", ] @@ -589,9 +639,9 @@ dependencies = [ [[package]] name = "clap" -version = "4.5.34" +version = "4.5.37" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e958897981290da2a852763fe9cdb89cd36977a5d729023127095fa94d95e2ff" +checksum = "eccb054f56cbd38340b380d4a8e69ef1f02f1af43db2f0cc817a4774d80ae071" dependencies = [ "clap_builder", "clap_derive", @@ -599,9 +649,9 @@ dependencies = [ [[package]] name = "clap_builder" -version = "4.5.34" +version = "4.5.37" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "83b0f35019843db2160b5bb19ae09b4e6411ac33fc6a712003c33e03090e2489" +checksum = "efd9466fac8543255d3b1fcad4762c5e116ffe808c8a3043d4263cd4fd4862a2" dependencies = [ "anstream", "anstyle", @@ -618,7 +668,7 @@ dependencies = [ "heck", "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -705,7 +755,7 @@ version = "0.1.16" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f9d839f2a20b0aee515dc581a6172f2321f96cab76c1a38a4c584a194955390e" dependencies = [ - "getrandom 0.2.15", + "getrandom 0.2.16", "once_cell", "tiny-keccak", ] @@ -863,6 +913,12 @@ dependencies = [ "const-random", ] +[[package]] +name = "embedded-io" +version = "0.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "edd0f118536f44f5ccd48bcb8b111bdc3de888b58c74639dfb034a357d0f206d" + [[package]] name = "encoding_rs" version = "0.8.35" @@ -880,9 +936,9 @@ checksum = "877a4ace8713b0bcf2a4e7eec82529c029f1d0619886d18145fea96c3ffe5c0f" [[package]] name = "errno" -version = "0.3.10" +version = "0.3.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "33d852cb9b869c2a9b3df2f71a3074817f01e1844f839a144f5fcef059a4eb5d" +checksum = "976dd42dc7e85965fe702eb8164f21f450704bdde31faefd6471dba214cb594e" dependencies = [ "libc", "windows-sys 0.59.0", @@ -1049,7 +1105,7 @@ checksum = "162ee34ebcb7c64a8abebc059ce0fee27c2262618d7b60ed8faf72fef13c3650" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -1094,9 +1150,9 @@ dependencies = [ [[package]] name = "getrandom" -version = "0.2.15" +version = "0.2.16" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c4567c8db10ae91089c99af84c68c38da3ec2f087c3f82960bcdbf3656b6f4d7" +checksum = "335ff9f135e4384c8150d6f27c6daed433577f86b4750418338c01a1a2528592" dependencies = [ "cfg-if", "libc", @@ -1145,9 +1201,9 @@ dependencies = [ [[package]] name = "h2" -version = "0.4.8" +version = "0.4.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5017294ff4bb30944501348f6f8e42e6ad28f42c8bbef7a74029aff064a4e3c2" +checksum = "75249d144030531f8dee69fe9cea04d3edf809a017ae445e2abdff6629e86633" dependencies = [ "atomic-waker", "bytes", @@ -1184,9 +1240,9 @@ dependencies = [ [[package]] name = "hashbrown" -version = "0.15.2" +version = "0.15.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bf151400ff0baff5465007dd2f3e717f3fe502074ca563069ce3a6629d07b289" +checksum = "84b26c544d002229e640969970a2e74021aadf6e2f96372b9c58eff97de08eb3" [[package]] name = "hashlink" @@ -1314,9 +1370,9 @@ dependencies = [ [[package]] name = "hyper-util" -version = "0.1.10" +version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df2dcfbe0677734ab2f3ffa7fa7bfd4706bfdc1ef393f2ee30184aed67e631b4" +checksum = "497bbc33a26fdd4af9ed9c70d63f61cf56a938375fbb32df34db9b1cd6d643f2" dependencies = [ "bytes", "futures-channel", @@ -1324,6 +1380,7 @@ dependencies = [ "http", "http-body", "hyper", + "libc", "pin-project-lite", "socket2", "tokio", @@ -1348,17 +1405,17 @@ checksum = "a0eb5a3343abf848c0984fe4604b2b105da9539376e24fc0a3b0007411ae4fd9" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] name = "indexmap" -version = "2.8.0" +version = "2.9.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3954d50fe15b02142bf25d3b8bdadb634ec3948f103d04ffe3031bc8fe9d7058" +checksum = "cea70ddb795996207ad57735b50c5982d8844f38ba9ee5f1aedcfb708a2aa11e" dependencies = [ "equivalent", - "hashbrown 0.15.2", + "hashbrown 0.15.3", ] [[package]] @@ -1390,7 +1447,7 @@ dependencies = [ "network_helpers_sv2", "once_cell", "pool_sv2", - "rand 0.9.0", + "rand 0.9.1", "roles_logic_sv2", "stratum-common", "tar", @@ -1535,9 +1592,9 @@ checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" [[package]] name = "libc" -version = "0.2.171" +version = "0.2.172" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c19937216e9d3aa9956d9bb8dfc0b0c8beb6058fc4f7a4dc4d850edf86a237d6" +checksum = "d750af042f7ef4f724306de029d18836c26c1765a54a6a3f094cbd23a7267ffa" [[package]] name = "libredox" @@ -1558,9 +1615,9 @@ checksum = "d26c52dbd32dccf2d10cac7725f8eae5296885fb5703b261f7d0a0739ec807ab" [[package]] name = "linux-raw-sys" -version = "0.9.3" +version = "0.9.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fe7db12097d22ec582439daf8618b8fdd1a7bef6270e9af3b1ebcd30893cf413" +checksum = "cd945864f07fe9f5371a27ad7b52a172b4b499999f1d97574c9fa68373937e12" [[package]] name = "lock_api" @@ -1609,7 +1666,7 @@ dependencies = [ "primitive-types", "rand 0.8.5", "roles_logic_sv2", - "sha2 0.10.8", + "sha2 0.10.9", "stratum-common", "tokio", "tracing", @@ -1667,21 +1724,20 @@ dependencies = [ [[package]] name = "miniz_oxide" -version = "0.8.5" +version = "0.8.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8e3e04debbb59698c15bacbb6d93584a8c0ca9cc3213cb423d31f760d8843ce5" +checksum = "3be647b768db090acb35d5ec5db2b0e1f1de11133ca123b9eacf5137868f892a" dependencies = [ "adler2", ] [[package]] name = "minreq" -version = "2.13.2" +version = "2.13.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "da0c420feb01b9fb5061f8c8f452534361dd783756dcf38ec45191ce55e7a161" +checksum = "f0d2aaba477837b46ec1289588180fabfccf0c3b1d1a0c6b1866240cd6cd5ce9" dependencies = [ "log", - "once_cell", "rustls", "rustls-webpki", "serde", @@ -1845,7 +1901,7 @@ dependencies = [ "proc-macro-crate", "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -1914,7 +1970,7 @@ dependencies = [ "pest_meta", "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -1925,7 +1981,7 @@ checksum = "7f9f832470494906d1fca5329f8ab5791cc60beb230c74815dff541cbd2b5ca0" dependencies = [ "once_cell", "pest", - "sha2 0.10.8", + "sha2 0.10.9", ] [[package]] @@ -2021,7 +2077,7 @@ version = "0.2.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" dependencies = [ - "zerocopy 0.8.24", + "zerocopy 0.8.25", ] [[package]] @@ -2046,9 +2102,9 @@ dependencies = [ [[package]] name = "proc-macro2" -version = "1.0.94" +version = "1.0.95" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a31971752e70b8b2686d7e46ec17fb38dad4051d94024c88df49b667caea9c84" +checksum = "02b3e5e68a3a1a02aad3ec490a98007cbc13c37cbe84a3cd7b8e406d76e7f778" dependencies = [ "unicode-ident", ] @@ -2087,13 +2143,12 @@ dependencies = [ [[package]] name = "rand" -version = "0.9.0" +version = "0.9.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3779b94aeb87e8bd4e834cee3650289ee9e0d5677f976ecdb6d219e5f4f6cd94" +checksum = "9fbfd9d094a40bf3ae768db9361049ace4c0e04a4fd6b359518bd7b73a73dd97" dependencies = [ "rand_chacha 0.9.0", "rand_core 0.9.3", - "zerocopy 0.8.24", ] [[package]] @@ -2122,7 +2177,7 @@ version = "0.6.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" dependencies = [ - "getrandom 0.2.15", + "getrandom 0.2.16", ] [[package]] @@ -2136,13 +2191,42 @@ dependencies = [ [[package]] name = "redox_syscall" -version = "0.5.10" +version = "0.5.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0b8c0c260b63a8219631167be35e6a988e9554dbd323f8bd08439c8ed1302bd1" +checksum = "928fca9cf2aa042393a8325b9ead81d2f0df4cb12e1e24cef072922ccd99c5af" dependencies = [ "bitflags", ] +[[package]] +name = "regex" +version = "1.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b544ef1b4eac5dc2db33ea63606ae9ffcfac26c1416a2806ae0bf5f56b201191" +dependencies = [ + "aho-corasick", + "memchr", + "regex-automata", + "regex-syntax", +] + +[[package]] +name = "regex-automata" +version = "0.4.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "809e8dc61f6de73b46c85f4c96486310fe304c434cfa43669d7b40f711150908" +dependencies = [ + "aho-corasick", + "memchr", + "regex-syntax", +] + +[[package]] +name = "regex-syntax" +version = "0.8.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b15c43186be67a4fd63bee50d0303afffcef381492ebe2c5d87f324e1b8815c" + [[package]] name = "ring" version = "0.17.14" @@ -2151,7 +2235,7 @@ checksum = "a4689e6c2294d81e88dc6261c768b63bc4fcdb852be6d1352498b114f61383b7" dependencies = [ "cc", "cfg-if", - "getrandom 0.2.15", + "getrandom 0.2.16", "libc", "untrusted", "windows-sys 0.52.0", @@ -2239,14 +2323,14 @@ dependencies = [ [[package]] name = "rustix" -version = "1.0.3" +version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e56a18552996ac8d29ecc3b190b4fdbb2d91ca4ec396de7bbffaf43f3d637e96" +checksum = "c71e83d6afe7ff64890ec6b71d6a69bb8a610ab78ce364b3352876bb4c801266" dependencies = [ "bitflags", "errno", "libc", - "linux-raw-sys 0.9.3", + "linux-raw-sys 0.9.4", "windows-sys 0.59.0", ] @@ -2357,7 +2441,7 @@ checksum = "5b0276cf7f2c73365f7157c8123c21cd9a50fbbd844757af28ca1f5925fc2a00" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -2396,9 +2480,9 @@ dependencies = [ [[package]] name = "sha2" -version = "0.10.8" +version = "0.10.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "793db75ad2bcafc3ffa7c68b215fee268f537982cd901d132f89c6343f3a3dc8" +checksum = "a7507d819769d01a365ab707794a4084392c824f54a7a6a7862f8c3d0892b283" dependencies = [ "cfg-if", "cpufeatures", @@ -2422,9 +2506,9 @@ checksum = "0fda2ff0d084019ba4d7c6f371c95d8fd75ce3524c3cb8fb653a3023f6323e64" [[package]] name = "signal-hook-registry" -version = "1.4.2" +version = "1.4.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a9e9e0b4211b72e7b8b6e85c807d36c212bdb33ea8587f7569562a84df5465b1" +checksum = "9203b8055f63a2a00e2f593bb0510367fe707d7ff1e5c872de2f537b339e5410" dependencies = [ "libc", ] @@ -2446,15 +2530,15 @@ dependencies = [ [[package]] name = "smallvec" -version = "1.14.0" +version = "1.15.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7fcf8323ef1faaee30a44a340193b1ac6814fd9b7b4e88e9d4519a3e4abe1cfd" +checksum = "8917285742e9f3e1683f0a9c4e6b57960b7314d0b08d30d1ecd426713ee2eee9" [[package]] name = "socket2" -version = "0.5.8" +version = "0.5.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c970269d99b64e60ec3bd6ad27270092a5394c4e309314b18ae3fe575695fbe8" +checksum = "4f5fd57c80058a56cf5c777ab8a126398ece8e442983605d280a44ce79d0edef" dependencies = [ "libc", "windows-sys 0.52.0", @@ -2512,9 +2596,9 @@ dependencies = [ [[package]] name = "syn" -version = "2.0.100" +version = "2.0.101" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b09a44accad81e1ba1cd74a32461ba89dee89095ba17b32f5d03683b1b1fc2a0" +checksum = "8ce2b7fc941b3a24138a0a7cf8e858bfc6a992e7978a068a5c760deb0ed43caf" dependencies = [ "proc-macro2", "quote", @@ -2545,10 +2629,35 @@ checksum = "7437ac7763b9b123ccf33c338a5cc1bac6f69b45a136c19bdd8a65e3916435bf" dependencies = [ "fastrand", "once_cell", - "rustix 1.0.3", + "rustix 1.0.7", "windows-sys 0.59.0", ] +[[package]] +name = "template-provider-role" +version = "0.1.0" +dependencies = [ + "async-channel 1.9.0", + "binary_sv2", + "buffer_sv2", + "capnp", + "capnp-rpc", + "capnpc", + "clap", + "codec_sv2", + "futures", + "key-utils", + "network_helpers_sv2", + "noise_sv2", + "regex", + "roles_logic_sv2", + "stratum-common", + "tokio", + "tokio-util", + "tracing", + "tracing-subscriber", +] + [[package]] name = "template_distribution_sv2" version = "3.0.0" @@ -2574,7 +2683,7 @@ checksum = "7f7cf42b4507d8ea322120659672cf1b9dbb93f8f2d4ecfd6e51350ff5b17a1d" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -2598,9 +2707,9 @@ dependencies = [ [[package]] name = "tokio" -version = "1.44.1" +version = "1.44.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f382da615b842244d4b8738c82ed1275e6c5dd90c459a30941cd07080b06c91a" +checksum = "e6b88822cbe49de4185e3a4cbf8321dd487cf5fe0c5c65695fef6346371e9c48" dependencies = [ "backtrace", "bytes", @@ -2623,17 +2732,18 @@ checksum = "6e06d43f1345a3bcd39f6a56dbb7dcab2ba47e68e8ac134855e7e2bdbaf8cab8" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] name = "tokio-util" -version = "0.7.14" +version = "0.7.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6b9590b93e6fcc1739458317cccd391ad3955e2bde8913edf6f95f9e65a8f034" +checksum = "66a539a9ad6d5d281510d5bd368c973d636c02dbf8a67300bfb6b950696ad7df" dependencies = [ "bytes", "futures-core", + "futures-io", "futures-sink", "pin-project-lite", "tokio", @@ -2641,9 +2751,9 @@ dependencies = [ [[package]] name = "toml" -version = "0.8.20" +version = "0.8.22" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cd87a5cdd6ffab733b2f74bc4fd7ee5fff6634124999ac278c35fc78c6120148" +checksum = "05ae329d1f08c4d17a59bed7ff5b5a769d062e64a62d34a3261b219e62cd5aae" dependencies = [ "serde", "serde_spanned", @@ -2653,26 +2763,33 @@ dependencies = [ [[package]] name = "toml_datetime" -version = "0.6.8" +version = "0.6.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0dd7358ecb8fc2f8d014bf86f6f638ce72ba252a2c3a2572f2a795f1d23efb41" +checksum = "3da5db5a963e24bc68be8b17b6fa82814bb22ee8660f192bb182771d498f09a3" dependencies = [ "serde", ] [[package]] name = "toml_edit" -version = "0.22.24" +version = "0.22.26" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "17b4795ff5edd201c7cd6dca065ae59972ce77d1b80fa0a84d94950ece7d1474" +checksum = "310068873db2c5b3e7659d2cc35d21855dbafa50d1ce336397c666e3cb08137e" dependencies = [ "indexmap", "serde", "serde_spanned", "toml_datetime", + "toml_write", "winnow", ] +[[package]] +name = "toml_write" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bfb942dfe1d8e29a7ee7fcbde5bd2b9a25fb89aa70caea2eba3bee836ff41076" + [[package]] name = "tower-service" version = "0.3.3" @@ -2698,7 +2815,7 @@ checksum = "395ae124c09f9e6918a2310af6038fba074bcf474ac352496d5910dd59a2226d" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] @@ -2757,7 +2874,7 @@ dependencies = [ "roles_logic_sv2", "serde", "serde_json", - "sha2 0.10.8", + "sha2 0.10.9", "stratum-common", "sv1_api", "tokio", @@ -2900,7 +3017,7 @@ dependencies = [ "log", "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", "wasm-bindgen-shared", ] @@ -2935,7 +3052,7 @@ checksum = "8ae87ea40c9f689fc23f209965b6fb8a99ad69aeeb0231408be24920604395de" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", "wasm-bindgen-backend", "wasm-bindgen-shared", ] @@ -3080,9 +3197,9 @@ checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" [[package]] name = "winnow" -version = "0.7.4" +version = "0.7.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0e97b544156e9bebe1a0ffbc03484fc1ffe3100cbce3ffb17eac35f7cdd7ab36" +checksum = "d9fb597c990f03753e08d3c29efbfcf2019a003b4bf4ba19225c158e1549f0f3" dependencies = [ "memchr", ] @@ -3127,11 +3244,11 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.24" +version = "0.8.25" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2586fea28e186957ef732a5f8b3be2da217d65c5969d4b1e17f973ebbe876879" +checksum = "a1702d9583232ddb9174e01bb7c15a2ab8fb1bc6f227aa1233858c351a3ba0cb" dependencies = [ - "zerocopy-derive 0.8.24", + "zerocopy-derive 0.8.25", ] [[package]] @@ -3142,18 +3259,18 @@ checksum = "fa4f8080344d4671fb4e831a13ad1e68092748387dfc4f55e356242fae12ce3e" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] name = "zerocopy-derive" -version = "0.8.24" +version = "0.8.25" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a996a8f63c5c4448cd959ac1bab0aaa3306ccfd060472f85943ee0750f0169be" +checksum = "28a6e20d751156648aa063f3800b706ee209a32c0b4d9f24be3d980b01be55ef" dependencies = [ "proc-macro2", "quote", - "syn 2.0.100", + "syn 2.0.101", ] [[package]] diff --git a/roles/Cargo.toml b/roles/Cargo.toml index 4208ef2989..74c7f875b8 100644 --- a/roles/Cargo.toml +++ b/roles/Cargo.toml @@ -10,7 +10,7 @@ members = [ "translator", "jd-client", "jd-server" -] +, "template-provider-role"] [profile.dev] # Required by super_safe_lock diff --git a/roles/template-provider-role/Cargo.toml b/roles/template-provider-role/Cargo.toml new file mode 100644 index 0000000000..2f01a37b10 --- /dev/null +++ b/roles/template-provider-role/Cargo.toml @@ -0,0 +1,35 @@ +[package] +name = "template-provider-role" +version = "0.1.0" +edition = "2021" + +# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html + + +[lib] +name = "template_provider_sv2" +path = "src/lib.rs" + +[dependencies] +clap = { version = "4.5.37", features = ["derive"] } +binary_sv2 = { path = "../../protocols/v2/binary-sv2" } +buffer_sv2 = { path = "../../utils/buffer" } +codec_sv2 = { path = "../../protocols/v2/codec-sv2", features = ["noise_sv2"] } +network_helpers_sv2 = { path = "../roles-utils/network-helpers" } +noise_sv2 = { path = "../../protocols/v2/noise-sv2" } +roles_logic_sv2 = { path = "../../protocols/v2/roles-logic-sv2" } +async-channel = "1.5.1" +tracing = "0.1.41" +stratum-common = { version = "2.0.0", path = "../../common", features = ["constants"] } +tracing-subscriber = "0.3.19" +tokio = { version = "1.44.2", features = ["full"] } +key-utils = { path = "../../utils/key-utils" } +capnp = "0.20.0" +capnp-rpc = "0.20.0" +tokio-util = { version = "0.7.15", features = ["compat"] } +futures = "0.3.31" + +[build-dependencies] +capnpc = "0.20.0" +regex = "1.11.1" + diff --git a/roles/template-provider-role/build.rs b/roles/template-provider-role/build.rs new file mode 100644 index 0000000000..1a99eb1cbe --- /dev/null +++ b/roles/template-provider-role/build.rs @@ -0,0 +1,43 @@ +use std::{ + env, + path::{Path, PathBuf}, +}; + +fn main() { + println!("cargo:rerun-if-changed=capnp"); + + let out_dir = PathBuf::from(env::var("OUT_DIR").expect("Missing OUT_DIR")); + let capnp_dir = Path::new("capnp"); + + let schemas = [ + "common.capnp", + "echo.capnp", + "init.capnp", + "mining.capnp", + "proxy.capnp", + ]; + let mut cmd = capnpc::CompilerCommand::new(); + cmd.src_prefix(capnp_dir).output_path(&out_dir); + for schema in &schemas { + cmd.file(capnp_dir.join(schema)); + } + cmd.run().expect("capnpc compilation failed"); + + // Too much hassle look into it later. + // let re = Regex::new(r"crate::(\w+_capnp)").unwrap(); + + // for schema in &schemas { + // let module_name = schema.strip_suffix(".capnp").unwrap(); + // let file_name = format!("{}_capnp.rs", module_name); + // let generated_file = out_dir.join(&file_name); + + // let content = fs::read_to_string(&generated_file) + // .unwrap_or_else(|e| panic!("Failed to read {}: {}", generated_file.display(), e)); + + // // Patch all references like `crate::proxy_capnp::...` → `crate::capnp::proxy_capnp::...` + // let patched = re.replace_all(&content, "crate::capnp::$1"); + + // fs::write(&generated_file, patched.as_ref()) + // .unwrap_or_else(|e| panic!("Failed to write {}: {}", generated_file.display(), e)); + // } +} diff --git a/roles/template-provider-role/capnp/common.capnp b/roles/template-provider-role/capnp/common.capnp new file mode 100644 index 0000000000..f6bed24855 --- /dev/null +++ b/roles/template-provider-role/capnp/common.capnp @@ -0,0 +1,18 @@ +# Copyright (c) 2024 The Bitcoin Core developers +# Distributed under the MIT software license, see the accompanying +# file COPYING or http://www.opensource.org/licenses/mit-license.php. + +@0xcd2c6232cb484a28; + +using Cxx = import "/capnp/c++.capnp"; +$Cxx.namespace("ipc::capnp::messages"); + +# using Proxy = import "/mp/proxy.capnp"; +# $Proxy.includeTypes("ipc/capnp/common-types.h"); + +using Proxy = import "proxy.capnp"; + +struct BlockRef $Proxy.wrap("interfaces::BlockRef") { + hash @0 :Data; + height @1 :Int32; +} diff --git a/roles/template-provider-role/capnp/echo.capnp b/roles/template-provider-role/capnp/echo.capnp new file mode 100644 index 0000000000..5c531d265c --- /dev/null +++ b/roles/template-provider-role/capnp/echo.capnp @@ -0,0 +1,18 @@ +# Copyright (c) 2021 The Bitcoin Core developers +# Distributed under the MIT software license, see the accompanying +# file COPYING or http://www.opensource.org/licenses/mit-license.php. + +@0x888b4f7f51e691f7; + +using Cxx = import "/capnp/c++.capnp"; +$Cxx.namespace("ipc::capnp::messages"); + +# using Proxy = import "/mp/proxy.capnp"; +# $Proxy.include("interfaces/echo.h"); +# $Proxy.includeTypes("ipc/capnp/echo-types.h"); +using Proxy = import "proxy.capnp"; + +interface Echo $Proxy.wrap("interfaces::Echo") { + destroy @0 (context :Proxy.Context) -> (); + echo @1 (context :Proxy.Context, echo: Text) -> (result :Text); +} diff --git a/roles/template-provider-role/capnp/init.capnp b/roles/template-provider-role/capnp/init.capnp new file mode 100644 index 0000000000..8dc628ac67 --- /dev/null +++ b/roles/template-provider-role/capnp/init.capnp @@ -0,0 +1,24 @@ +# Copyright (c) 2021 The Bitcoin Core developers +# Distributed under the MIT software license, see the accompanying +# file COPYING or http://www.opensource.org/licenses/mit-license.php. + +@0xf2c5cfa319406aa6; + +using Cxx = import "/capnp/c++.capnp"; +$Cxx.namespace("ipc::capnp::messages"); + +# using Proxy = import "/mp/proxy.capnp"; +# $Proxy.include("interfaces/echo.h"); +# $Proxy.include("interfaces/init.h"); +# $Proxy.include("interfaces/mining.h"); +# $Proxy.includeTypes("ipc/capnp/init-types.h"); + +using Echo = import "echo.capnp"; +using Mining = import "mining.capnp"; +using Proxy = import "proxy.capnp"; + +interface Init $Proxy.wrap("interfaces::Init") { + construct @0 (threadMap: Proxy.ThreadMap) -> (threadMap :Proxy.ThreadMap); + makeEcho @1 (context :Proxy.Context) -> (result :Echo.Echo); + makeMining @2 (context :Proxy.Context) -> (result :Mining.Mining); +} diff --git a/roles/template-provider-role/capnp/mining.capnp b/roles/template-provider-role/capnp/mining.capnp new file mode 100644 index 0000000000..93052da7b6 --- /dev/null +++ b/roles/template-provider-role/capnp/mining.capnp @@ -0,0 +1,57 @@ +# Copyright (c) 2024 The Bitcoin Core developers +# Distributed under the MIT software license, see the accompanying +# file COPYING or http://www.opensource.org/licenses/mit-license.php. + +@0xc77d03df6a41b505; + +using Cxx = import "/capnp/c++.capnp"; +$Cxx.namespace("ipc::capnp::messages"); + +using Common = import "common.capnp"; +# using Proxy = import "/mp/proxy.capnp"; +# $Proxy.include("interfaces/mining.h"); +# $Proxy.includeTypes("ipc/capnp/mining-types.h"); +using Proxy = import "proxy.capnp"; + +interface Mining $Proxy.wrap("interfaces::Mining") { + isTestChain @0 (context :Proxy.Context) -> (result: Bool); + isInitialBlockDownload @1 (context :Proxy.Context) -> (result: Bool); + getTip @2 (context :Proxy.Context) -> (result: Common.BlockRef, hasResult: Bool); + waitTipChanged @3 (context :Proxy.Context, currentTip: Data, timeout: Float64) -> (result: Common.BlockRef); + createNewBlock @4 (options: BlockCreateOptions) -> (result: BlockTemplate); +} + +interface BlockTemplate $Proxy.wrap("interfaces::BlockTemplate") { + destroy @0 (context :Proxy.Context) -> (); + getBlockHeader @1 (context: Proxy.Context) -> (result: Data); + getBlock @2 (context: Proxy.Context) -> (result: Data); + getTxFees @3 (context: Proxy.Context) -> (result: List(Int64)); + getTxSigops @4 (context: Proxy.Context) -> (result: List(Int64)); + getCoinbaseTx @5 (context: Proxy.Context) -> (result: Data); + getCoinbaseCommitment @6 (context: Proxy.Context) -> (result: Data); + getWitnessCommitmentIndex @7 (context: Proxy.Context) -> (result: Int32); + getCoinbaseMerklePath @8 (context: Proxy.Context) -> (result: List(Data)); + submitSolution @9 (context: Proxy.Context, version: UInt32, timestamp: UInt32, nonce: UInt32, coinbase :Data) -> (result: Bool); + waitNext @10 (context: Proxy.Context, options: BlockWaitOptions) -> (result: BlockTemplate); +} + +struct BlockCreateOptions $Proxy.wrap("node::BlockCreateOptions") { + useMempool @0 :Bool $Proxy.name("use_mempool"); + blockReservedWeight @1 :UInt64 $Proxy.name("block_reserved_weight"); + coinbaseOutputMaxAdditionalSigops @2 :UInt64 $Proxy.name("coinbase_output_max_additional_sigops"); +} + +struct BlockWaitOptions $Proxy.wrap("node::BlockWaitOptions") { + timeout @0 : Float64 $Proxy.name("timeout"); + feeThreshold @1 : Int64 $Proxy.name("fee_threshold"); +} + +# Note: serialization of the BlockValidationState C++ type is somewhat fragile +# and using the struct can be awkward. It would be good if testBlockValidity +# method were changed to return validity information in a simpler format. +struct BlockValidationState { + mode @0 :Int32; + result @1 :Int32; + rejectReason @2 :Text; + debugMessage @3 :Text; +} diff --git a/roles/template-provider-role/capnp/proxy.capnp b/roles/template-provider-role/capnp/proxy.capnp new file mode 100644 index 0000000000..abd02e437f --- /dev/null +++ b/roles/template-provider-role/capnp/proxy.capnp @@ -0,0 +1,65 @@ +# Copyright (c) 2019 The Bitcoin Core developers +# Distributed under the MIT software license, see the accompanying +# file COPYING or http://www.opensource.org/licenses/mit-license.php. + +@0xcc316e3f71a040fb; + +using Cxx = import "/capnp/c++.capnp"; +$Cxx.namespace("mp"); + +annotation include(file): Text; +annotation includeTypes(file): Text; +# Extra include paths to add to generated files. + +annotation wrap(interface, struct): Text; +# Wrap capnp interface generating ProxyClient / ProxyServer C++ classes that +# forward calls to a C++ interface with same methods and parameters. Text +# string should be the name of the C++ interface. +# If applied to struct rather than an interface, this will generate a ProxyType +# struct with get methods for introspection and copying fields between C++ and +# capnp structs. + +annotation count(param, struct, interface): Int32; +# Indicate how many C++ method parameters there are corresponding to one capnp +# parameter (default is 1). If not 1, multiple C++ method arguments will be +# condensed into a single capnp parameter by the client and then expanded by +# the server by CustomReadField/CustomBuildField overloads which need to be +# provided separately. An example would be a capnp Text parameter initialized +# from C++ char* and size arguments. Can be 0 to fill an implicit capnp +# parameter from client or server side context. If annotation is applied to an +# interface or struct type it will apply to all parameters of that type. + +annotation exception(param): Text; +# Indicate that a result parameter corresponds to a C++ exception. Text string +# should be the name of a C++ exception type that the generated server class +# will catch and the client class will rethrow. + +annotation name(field, method): Text; +# Name of the C++ method or field corresponding to a capnp method or field. + +annotation skip(field): Void; +# Synonym for count(0). + +interface ThreadMap $count(0) { + # Interface letting clients control which thread a method call should + # execute on. Clients create and name threads and pass the thread handle as + # a call parameter. + makeThread @0 (name :Text) -> (result :Thread); +} + +interface Thread { + # Thread handle returned by makeThread corresponding to one server thread. + + getName @0 () -> (result: Text); +} + +struct Context $count(0) { + # Execution context passed as a parameter from the client class to the server class. + + thread @0 : Thread; + # Handle of the server thread the current method call should execute on. + + callbackThread @1 : Thread; + # Handle of the client thread that is calling the current method, and that + # any callbacks made by the server thread should be made on. +} diff --git a/roles/template-provider-role/src/error.rs b/roles/template-provider-role/src/error.rs new file mode 100644 index 0000000000..5c55192888 --- /dev/null +++ b/roles/template-provider-role/src/error.rs @@ -0,0 +1,110 @@ +use std::{ + convert::From, + fmt::Debug, + sync::{MutexGuard, PoisonError}, +}; + +use roles_logic_sv2::parsers::Mining; + +#[derive(std::fmt::Debug)] +pub enum TPError { + Io(std::io::Error), + ChannelSend(Box), + ChannelRecv(async_channel::RecvError), + BinarySv2(binary_sv2::Error), + Codec(codec_sv2::Error), + Noise(noise_sv2::Error), + RolesLogic(roles_logic_sv2::Error), + Framing(codec_sv2::framing_sv2::Error), + PoisonLock(String), + Custom(String), + Sv2ProtocolError((u32, Mining<'static>)), +} + +impl std::fmt::Display for TPError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + use TPError::*; + match self { + Io(ref e) => write!(f, "I/O error: `{:?}", e), + ChannelSend(ref e) => write!(f, "Channel send failed: `{:?}`", e), + ChannelRecv(ref e) => write!(f, "Channel recv failed: `{:?}`", e), + BinarySv2(ref e) => write!(f, "Binary SV2 error: `{:?}`", e), + Codec(ref e) => write!(f, "Codec SV2 error: `{:?}", e), + Framing(ref e) => write!(f, "Framing SV2 error: `{:?}`", e), + Noise(ref e) => write!(f, "Noise SV2 error: `{:?}", e), + RolesLogic(ref e) => write!(f, "Roles Logic SV2 error: `{:?}`", e), + PoisonLock(ref e) => write!(f, "Poison lock: {:?}", e), + Custom(ref e) => write!(f, "Custom SV2 error: `{:?}`", e), + Sv2ProtocolError(ref e) => { + write!(f, "Received Sv2 Protocol Error from upstream: `{:?}`", e) + } + } + } +} + +pub type TPResult = Result; + +impl From for TPError { + fn from(e: std::io::Error) -> TPError { + TPError::Io(e) + } +} + +impl From for TPError { + fn from(e: async_channel::RecvError) -> TPError { + TPError::ChannelRecv(e) + } +} + +impl From for TPError { + fn from(e: binary_sv2::Error) -> TPError { + TPError::BinarySv2(e) + } +} + +impl From for TPError { + fn from(e: codec_sv2::Error) -> TPError { + TPError::Codec(e) + } +} + +impl From for TPError { + fn from(e: noise_sv2::Error) -> TPError { + TPError::Noise(e) + } +} + +impl From for TPError { + fn from(e: roles_logic_sv2::Error) -> TPError { + TPError::RolesLogic(e) + } +} + +impl From> for TPError { + fn from(e: async_channel::SendError) -> TPError { + TPError::ChannelSend(Box::new(e)) + } +} + +impl From for TPError { + fn from(e: String) -> TPError { + TPError::Custom(e) + } +} +impl From for TPError { + fn from(e: codec_sv2::framing_sv2::Error) -> TPError { + TPError::Framing(e) + } +} + +impl From>> for TPError { + fn from(e: PoisonError>) -> TPError { + TPError::PoisonLock(e.to_string()) + } +} + +impl From<(u32, Mining<'static>)> for TPError { + fn from(e: (u32, Mining<'static>)) -> Self { + TPError::Sv2ProtocolError(e) + } +} diff --git a/roles/template-provider-role/src/lib.rs b/roles/template-provider-role/src/lib.rs new file mode 100644 index 0000000000..63fe023407 --- /dev/null +++ b/roles/template-provider-role/src/lib.rs @@ -0,0 +1,299 @@ +#![allow(warnings)] +use std::time::Duration; + +use async_channel::{Receiver, Sender}; +use codec_sv2::{Frame, HandshakeRole, StandardEitherFrame, StandardSv2Frame}; +use key_utils::{Secp256k1PublicKey, Secp256k1SecretKey}; +use network_helpers_sv2::noise_connection::Connection; +use noise_sv2::Responder; +use roles_logic_sv2::parsers::{ + message_type_to_name, AnyMessage, CommonMessages, IsSv2Message, + JobDeclaration::{ + AllocateMiningJobToken, AllocateMiningJobTokenSuccess, DeclareMiningJob, + DeclareMiningJobError, DeclareMiningJobSuccess, IdentifyTransactions, + IdentifyTransactionsSuccess, ProvideMissingTransactions, ProvideMissingTransactionsSuccess, + PushSolution, + }, + TemplateDistribution::{self, CoinbaseOutputConstraints}, +}; +use setup_connection::SetupConnectionHandler; +use tokio::{net::TcpListener, task::LocalSet}; +use tracing::info; + +mod error; +mod message_utils; +mod provider; +mod rpc_client; +mod setup_connection; + +pub type Message = AnyMessage<'static>; +pub type StdFrame = StandardSv2Frame; +pub type EitherFrame = StandardEitherFrame; + +pub mod common_capnp { + include!(concat!(env!("OUT_DIR"), "/common_capnp.rs")); +} +pub mod echo_capnp { + include!(concat!(env!("OUT_DIR"), "/echo_capnp.rs")); +} +pub mod init_capnp { + include!(concat!(env!("OUT_DIR"), "/init_capnp.rs")); +} +pub mod mining_capnp { + include!(concat!(env!("OUT_DIR"), "/mining_capnp.rs")); +} +pub mod proxy_capnp { + include!(concat!(env!("OUT_DIR"), "/proxy_capnp.rs")); +} + +use provider::Provider; +use rpc_client::error::RpcClientError; +use tracing::{debug, error}; + +#[derive(Debug, Clone)] +pub struct Config { + sv2bind: String, + sv2port: u16, + sv2interval: u64, + sv2feedelta: u64, + socket_path: String, + authority_public_key: Secp256k1PublicKey, + authority_secret_key: Secp256k1SecretKey, + cert_validity_sec: u64, +} + +impl Config { + pub fn new( + sv2bind: String, + sv2port: u16, + sv2interval: u64, + sv2feedelta: u64, + socket_path: String, + ) -> Self { + let authority_public_key = "9auqWEzQDVyd2oe1JVGFLMLHZtCo2FFqZwtKA5gd9xbuEu7PH72"; + let authority_secret_key = "mkDLTBBRxdBv998612qipDYoTK3YUrqLe8uWw7gu3iXbSrn2n"; + let cert_validity_sec = 3600; + Config { + sv2bind, + sv2port, + sv2interval, + sv2feedelta, + authority_public_key: authority_public_key + .parse() + .expect("Invalid authority public key"), + authority_secret_key: authority_secret_key + .parse() + .expect("Invalid authority secret key"), + cert_validity_sec, + socket_path, + } + } +} + +/// Main entry point to start the Template Provider service. +pub async fn start_template_provider(config: Config) { + info!("Starting Template Provider service..."); + let (prev_hash_sender, prev_hash_receiver) = async_channel::unbounded::(); + let (template_sender, template_receiver) = async_channel::unbounded::(); + let (tx_data_response_sender, tx_data_response_receiver) = + async_channel::unbounded::(); + let (message_sender, message_receiver) = async_channel::unbounded::(); + info!("Communication channels initialized."); + let network_config = config.clone(); + tokio::spawn(async move { + listen_for_connections( + network_config, + prev_hash_receiver, + template_receiver, + tx_data_response_receiver, + message_sender, + ) + .await; + error!("Network listener task finished unexpectedly."); + }); + + let local_set = LocalSet::new(); + info!("Running Provider initialization and loop within LocalSet."); + + local_set + .run_until(async move { + info!("Attempting to initialize Provider..."); + match Provider::new( + prev_hash_sender, + template_sender, + tx_data_response_sender, + message_receiver, + config.sv2feedelta, + config.sv2interval, + config.socket_path, + ) + .await + { + Ok(provider) => { + info!("Provider initialized successfully. Starting run loop."); + provider.run().await; + info!("Provider run loop finished."); + } + Err(RpcClientError::Connection(e)) => { + error!( + "Provider initialization failed: Could not connect to RPC socket: {}", + e + ); + } + Err(RpcClientError::Capnp(e)) => { + error!( + "Provider initialization failed: Cap'n Proto RPC error: {}", + e + ); + } + Err(e) => { + error!("Provider initialization failed"); + } + } + // If initialization fails, the application might effectively stop here. + // Consider more robust shutdown logic if needed. + error!("Provider task finished or failed to initialize."); + }) + .await; + + info!("Template Provider LocalSet finished execution."); +} + +async fn listen_for_connections( + config: Config, + prev_hash_receiver: Receiver, + template_receiver: Receiver, + tx_data_response_receiver: Receiver, + message_sender: Sender, +) { + let bind_addr = format!("{}:{}", config.sv2bind, config.sv2port); + let listener = match TcpListener::bind(&bind_addr).await { + Ok(l) => { + info!("SV2 Listener started on: {}", bind_addr); + l + } + Err(e) => { + error!("Failed to bind TCP listener to {}: {}", bind_addr, e); + return; + } + }; + + loop { + match listener.accept().await { + Ok((stream, address)) => { + info!("Accepted SV2 connection from: {}", address); + let conn_prev_hash_receiver = prev_hash_receiver.clone(); + let conn_template_receiver = template_receiver.clone(); + let conn_tx_data_receiver = tx_data_response_receiver.clone(); + let conn_message_sender = message_sender.clone(); + let conn_config = config.clone(); + + tokio::spawn(async move { + info!("Spawning connection handler for {}", address); + if let Err(e) = handle_connection( + conn_config, + stream, + address, + conn_prev_hash_receiver, + conn_template_receiver, + conn_tx_data_receiver, + conn_message_sender, + ) + .await + { + error!("Error handling connection from {}: {}", address, e); + } else { + info!("Connection handler finished gracefully for {}", address); + } + }); + } + Err(e) => { + error!("Error accepting connection: {}", e); + tokio::time::sleep(Duration::from_millis(100)).await; + } + } + } +} + +type ConnectionHandlerError = Box; + +async fn handle_connection( + config: Config, + stream: tokio::net::TcpStream, + address: std::net::SocketAddr, + prev_hash_receiver: Receiver, + template_receiver: Receiver, + tx_data_response_receiver: Receiver, + message_sender: Sender, +) -> Result<(), ConnectionHandlerError> { + let responder = Responder::from_authority_kp( + &config.authority_public_key.into_bytes(), + &config.authority_secret_key.into_bytes(), + std::time::Duration::from_secs(config.cert_validity_sec), + ) + .unwrap(); + let (mut receiver, mut sender) = Connection::new(stream, HandshakeRole::Responder(responder)) + .await + .unwrap(); + + _ = SetupConnectionHandler::setup(&mut receiver, &mut sender, address) + .await + .unwrap(); + + info!("SV2 SetupConnection successful with {}", address); + tokio::spawn(forward_messages( + sender.clone(), + prev_hash_receiver, + "PrevHash", + )); + tokio::spawn(forward_messages( + sender.clone(), + template_receiver, + "Template", + )); + tokio::spawn(forward_messages( + sender.clone(), + tx_data_response_receiver, + "TxDataResponse", + )); + + loop { + match receiver.recv().await { + Ok(msg) => { + if message_sender.send(msg).await.is_err() { + error!("Failed to forward message to Provider; channel closed. Closing connection handler for {}.", address); + break; + } + } + Err(e) => { + info!("Connection closed by peer {}", address); + return Err(e.into()); + } + } + } + Ok(()) +} + +async fn forward_messages( + sender: Sender, + receiver: Receiver, + message_type_name: &'static str, +) { + debug!( + "Starting forwarder task for {} messages.", + message_type_name + ); + while let Ok(msg) = receiver.recv().await { + if sender.send(msg).await.is_err() { + debug!( + "Failed to forward {} message; connection likely closed.", + message_type_name + ); + break; + } + } + debug!( + "Stopping forwarder task for {} messages (receiver closed or send failed).", + message_type_name + ); +} diff --git a/roles/template-provider-role/src/main.rs b/roles/template-provider-role/src/main.rs new file mode 100644 index 0000000000..cf71d50f06 --- /dev/null +++ b/roles/template-provider-role/src/main.rs @@ -0,0 +1,49 @@ +#![allow(warnings)] +use clap::Parser; +use template_provider_sv2::start_template_provider; +use tracing::{debug, info}; + +/// Handles SV2 communication with configurable network and timing parameters +#[derive(Parser, Debug)] +#[command(name = "sv2", version, about, long_about = None)] +pub struct Sv2CLI { + /// The IP address to bind for incoming SV2 connections + #[arg(long, default_value = "127.0.0.1")] + sv2bind: String, + + /// The port number for incoming SV2 connections + #[arg(long, default_value_t = 8442)] + sv2port: u16, + + /// Time interval (in seconds) between SV2 messages + #[arg(long, default_value_t = 10)] + sv2interval: u64, + + /// Fee delta to apply to SV2 jobs (in sats) + #[arg(long, default_value_t = 0)] + sv2feedelta: u64, + + #[arg(long, default_value = "/home/shourya/.bitcoin/testnet4/node.sock")] + unix_socket_path: String, +} + +impl From for template_provider_sv2::Config { + fn from(value: Sv2CLI) -> Self { + template_provider_sv2::Config::new( + value.sv2bind, + value.sv2port, + value.sv2interval, + value.sv2feedelta, + value.unix_socket_path, + ) + } +} + +#[tokio::main] +async fn main() { + tracing_subscriber::fmt::init(); + debug!("Parsing the CLI"); + let config = Sv2CLI::parse(); + info!("Starting the template provider instance"); + start_template_provider(config.into()).await; +} diff --git a/roles/template-provider-role/src/message_utils.rs b/roles/template-provider-role/src/message_utils.rs new file mode 100644 index 0000000000..538ece7347 --- /dev/null +++ b/roles/template-provider-role/src/message_utils.rs @@ -0,0 +1,121 @@ +use codec_sv2::Frame; +use roles_logic_sv2::parsers::{ + message_type_to_name, AnyMessage, CommonMessages, IsSv2Message, + JobDeclaration::{ + AllocateMiningJobToken, AllocateMiningJobTokenSuccess, DeclareMiningJob, + DeclareMiningJobError, DeclareMiningJobSuccess, IdentifyTransactions, + IdentifyTransactionsSuccess, ProvideMissingTransactions, ProvideMissingTransactionsSuccess, + PushSolution, + }, + TemplateDistribution::{self, CoinbaseOutputConstraints}, +}; + +use crate::EitherFrame; + +pub fn message_from_frame(frame: &mut EitherFrame) -> (u8, AnyMessage<'static>) { + match frame { + Frame::Sv2(frame) => { + if let Some(header) = frame.get_header() { + let message_type = header.msg_type(); + let mut payload = frame.payload().to_vec(); + let message: Result, _> = + (message_type, payload.as_mut_slice()).try_into(); + match message { + Ok(message) => { + let message = into_static(message); + (message_type, message) + } + _ => { + println!("Received frame with invalid payload or message type: {frame:?}"); + panic!(); + } + } + } else { + println!("Received frame with invalid header: {frame:?}"); + panic!(); + } + } + Frame::HandShake(f) => { + println!("Received unexpected handshake frame: {f:?}"); + panic!(); + } + } +} + +pub fn into_static(m: AnyMessage<'_>) -> AnyMessage<'static> { + match m { + AnyMessage::Mining(m) => AnyMessage::Mining(m.into_static()), + AnyMessage::Common(m) => match m { + CommonMessages::ChannelEndpointChanged(m) => { + AnyMessage::Common(CommonMessages::ChannelEndpointChanged(m.into_static())) + } + CommonMessages::SetupConnection(m) => { + AnyMessage::Common(CommonMessages::SetupConnection(m.into_static())) + } + CommonMessages::SetupConnectionError(m) => { + AnyMessage::Common(CommonMessages::SetupConnectionError(m.into_static())) + } + CommonMessages::SetupConnectionSuccess(m) => { + AnyMessage::Common(CommonMessages::SetupConnectionSuccess(m.into_static())) + } + CommonMessages::Reconnect(m) => { + AnyMessage::Common(CommonMessages::Reconnect(m.into_static())) + } + }, + AnyMessage::JobDeclaration(m) => match m { + AllocateMiningJobToken(m) => { + AnyMessage::JobDeclaration(AllocateMiningJobToken(m.into_static())) + } + AllocateMiningJobTokenSuccess(m) => { + AnyMessage::JobDeclaration(AllocateMiningJobTokenSuccess(m.into_static())) + } + DeclareMiningJob(m) => AnyMessage::JobDeclaration(DeclareMiningJob(m.into_static())), + DeclareMiningJobError(m) => { + AnyMessage::JobDeclaration(DeclareMiningJobError(m.into_static())) + } + DeclareMiningJobSuccess(m) => { + AnyMessage::JobDeclaration(DeclareMiningJobSuccess(m.into_static())) + } + IdentifyTransactions(m) => { + AnyMessage::JobDeclaration(IdentifyTransactions(m.into_static())) + } + IdentifyTransactionsSuccess(m) => { + AnyMessage::JobDeclaration(IdentifyTransactionsSuccess(m.into_static())) + } + ProvideMissingTransactions(m) => { + AnyMessage::JobDeclaration(ProvideMissingTransactions(m.into_static())) + } + ProvideMissingTransactionsSuccess(m) => { + AnyMessage::JobDeclaration(ProvideMissingTransactionsSuccess(m.into_static())) + } + PushSolution(m) => AnyMessage::JobDeclaration(PushSolution(m.into_static())), + }, + AnyMessage::TemplateDistribution(m) => match m { + CoinbaseOutputConstraints(m) => { + AnyMessage::TemplateDistribution(CoinbaseOutputConstraints(m.into_static())) + } + TemplateDistribution::NewTemplate(m) => { + AnyMessage::TemplateDistribution(TemplateDistribution::NewTemplate(m.into_static())) + } + TemplateDistribution::RequestTransactionData(m) => AnyMessage::TemplateDistribution( + TemplateDistribution::RequestTransactionData(m.into_static()), + ), + TemplateDistribution::RequestTransactionDataError(m) => { + AnyMessage::TemplateDistribution(TemplateDistribution::RequestTransactionDataError( + m.into_static(), + )) + } + TemplateDistribution::RequestTransactionDataSuccess(m) => { + AnyMessage::TemplateDistribution( + TemplateDistribution::RequestTransactionDataSuccess(m.into_static()), + ) + } + TemplateDistribution::SetNewPrevHash(m) => AnyMessage::TemplateDistribution( + TemplateDistribution::SetNewPrevHash(m.into_static()), + ), + TemplateDistribution::SubmitSolution(m) => AnyMessage::TemplateDistribution( + TemplateDistribution::SubmitSolution(m.into_static()), + ), + }, + } +} diff --git a/roles/template-provider-role/src/provider/mod.rs b/roles/template-provider-role/src/provider/mod.rs new file mode 100644 index 0000000000..e8af4c2fdb --- /dev/null +++ b/roles/template-provider-role/src/provider/mod.rs @@ -0,0 +1,506 @@ +use crate::{ + init_capnp, proxy_capnp, + rpc_client::{error::RpcClientError, BackendRpcClient, BlockTemplateClient}, +}; +use async_channel::{Receiver, Sender}; +use binary_sv2::{B016M, U256}; +use codec_sv2::Sv2Frame; +use roles_logic_sv2::{ + parsers::{ + AnyMessage, + TemplateDistribution::{CoinbaseOutputConstraints, RequestTransactionData, SubmitSolution}, + }, + template_distribution_sv2::{ + NewTemplate, RequestTransactionDataError, RequestTransactionDataSuccess, SetNewPrevHash, + }, +}; +use state::{ChainTip, ProviderState}; +use std::{ + cell::Cell, + collections::HashMap, + io::Cursor, + path::Path, + rc::Rc, + sync::atomic::{AtomicBool, AtomicU64}, + time::{Duration, Instant}, +}; +use stratum_common::bitcoin::{ + block::Header, + consensus, + consensus::{Decodable, Encodable}, + hashes::Hash, + Block, BlockHash, Transaction, +}; +use tokio::{task, time::interval}; +use tracing::{debug, error, info, warn}; + +use crate::{ + message_utils::{into_static, message_from_frame}, + EitherFrame, +}; +const BLOCK_RESERVED_WEIGHT: u64 = 3990000; +const USE_MEMPOOL_DEFAULT: bool = true; +const TEMPLATE_EVICTION_DURATION: Duration = Duration::from_secs(600); +mod state; + +pub struct Provider { + state: ProviderState, + rpc_client: BackendRpcClient, + prev_hash_sender: Sender, + template_sender: Sender, + tx_data_response_sender: Sender, + message_receiver: Receiver, + sv2interval: Duration, + sv2feedelta: u64, + active_template: Option, +} + +impl Provider { + pub async fn new( + prev_hash_sender: Sender, + template_sender: Sender, + tx_data_response_sender: Sender, + message_receiver: Receiver, + sv2feedelta: u64, + sv2interval: u64, + socket_path: String, + ) -> Result { + info!("Initializing Provider..."); + let rpc_client = BackendRpcClient::connect(std::path::Path::new(&socket_path)).await?; + info!("RPC client connected successfully."); + Ok(Provider { + state: ProviderState::new(), + rpc_client, + prev_hash_sender, + template_sender, + tx_data_response_sender, + message_receiver, + sv2interval: Duration::from_secs(sv2interval), + sv2feedelta, + active_template: None, + }) + } + + /// Runs the main event loop for the provider. + pub async fn run(mut self) { + info!("Starting Provider event loop."); + let mut template_timer = interval(self.sv2interval); + let mut block_checker = interval(Duration::from_secs(1)); + let mut template_eviction_timer = interval(Duration::from_secs(360)); + + loop { + let message_receiver_clone = self.message_receiver.clone(); + + tokio::select! { + biased; + maybe_msg = message_receiver_clone.recv() => { + match maybe_msg { + Ok(mut msg) => self.handle_upstream_message(&mut msg).await, + Err(e) => { + error!("Message receiver channel closed or error: {}. Shutting down provider loop.", e); + break; + } + } + }, + + _ = block_checker.tick() => { + if let Err(e) = self.check_for_tip_change().await { + error!("Error during tip change check"); + } + }, + + _ = template_timer.tick() => { + if let Err(e) = self.generate_periodic_template().await { + error!("Error during periodic template generation"); + } + }, + + _ = template_eviction_timer.tick() => { + self.evict_old_templates(); + }, + } + } + info!("Provider event loop finished."); + } + + /// Handles messages received from downstream. + async fn handle_upstream_message(&mut self, frame: &mut EitherFrame) { + let (msg_type, msg) = message_from_frame(frame); + debug!("Received message type {}: {:?}", msg_type, msg); + + match msg { + AnyMessage::TemplateDistribution(CoinbaseOutputConstraints(m)) => { + info!("Received CoinbaseOutputConstraints: {:?}", m); + let sigops = m.coinbase_output_max_additional_sigops as u64; + self.state.set_coinbase_constraints(sigops); + if let Err(e) = self + .rpc_client + .update_coinbase_constraints(sigops, BLOCK_RESERVED_WEIGHT) + .await + { + error!("Failed to update coinbase constraints via RPC"); + } else { + info!("Coinbase constraints options sent successfully via RPC."); + } + } + AnyMessage::TemplateDistribution(RequestTransactionData(m)) => { + info!( + "Received RequestTransactionData for template ID {}: {:?}", + m.template_id, m + ); + self.handle_request_transaction_data(m).await; + } + AnyMessage::TemplateDistribution(SubmitSolution(m)) => { + info!( + "Received SubmitSolution for template ID {}: {:?}", + m.template_id, m + ); + if let Some(template_client) = &self.active_template { + match template_client + .submit_solution( + m.version, + m.header_timestamp, + m.header_nonce, + &m.coinbase_tx.to_vec()[..], + ) + .await + { + Ok(true) => { + info!("Solution submitted successfully via RPC for template associated with current capability."); + self.active_template = None; + } + Ok(false) => { + warn!("Solution submission via RPC indicated logical failure (e.g., rejected) for template."); + self.active_template = None; + } + Err(e) => { + error!("Failed to submit solution via RPC"); + } + } + } else { + warn!("Received SubmitSolution but no active block template client is available (Template ID: {}). Solution ignored.", m.template_id); + } + } + _ => { + warn!("Received unhandled message type {}: {:?}", msg_type, msg); + } + } + } + + async fn check_for_tip_change(&mut self) -> Result<(), RpcClientError> { + if !self.state.is_coinbase_constraint_received() { + debug!("Coinbase constraints not yet received, skipping tip check."); + return Ok(()); + } + let maybe_tip = self.rpc_client.get_tip().await?; + let (current_tip_hash, tip_height) = match maybe_tip { + Some(tip) => tip, + None => { + warn!("Node reported no tip available during check."); + return Ok(()); + } + }; + + if self.state.last_tip.prevhash != current_tip_hash + || self.state.last_tip.height as u32 != tip_height + { + info!( + "Tip change detected! New tip: height={}, hash={}", + tip_height, current_tip_hash + ); + + self.generate_and_process_template(true).await?; + self.state.last_tip = ChainTip { + prevhash: current_tip_hash, + height: tip_height as i32, + }; + info!("Tip change processing complete and local state updated."); + } + Ok(()) + } + + async fn generate_periodic_template(&mut self) -> Result<(), RpcClientError> { + if !self.state.is_coinbase_constraint_received() || !self.state.is_prev_hash_received() { + debug!("Skipping periodic template generation: Coinbase constraints received: {}, PrevHash received: {}", + self.state.is_coinbase_constraint_received(), self.state.is_prev_hash_received()); + return Ok(()); + } + self.generate_and_process_template(false).await?; + debug!("Periodic template generation/processing attempt complete."); + Ok(()) + } + + /// Fetches new template data, updates state, and sends messages downstream. + async fn generate_and_process_template( + &mut self, + is_future: bool, + ) -> Result<(), RpcClientError> { + let template_id = self.state.peek_next_template_id(); + debug!( + template_id, + is_future, "Generating template data via RPC..." + ); + + let new_template_client = self + .rpc_client + .create_new_block( + USE_MEMPOOL_DEFAULT, + BLOCK_RESERVED_WEIGHT, + self.state.get_coinbase_max_additional_sigops(), + ) + .await?; + + let block = new_template_client.get_block().await?; + let merkle_path_bytes = new_template_client.get_coinbase_merkle_path().await?; + + let merkle_path: Vec = merkle_path_bytes + .iter() + .map(|e| U256::Owned(e.to_vec())) + .collect(); + let fee = block.txdata[0] + .output + .iter() + .fold(0, |acc, x| acc + x.value.to_sat()); + + let map_try_into_error = |field: &'static str| { + RpcClientError::InvalidData(format!("Data too large for {}", field)) + }; + + let coinbase_outputs_concat = block.txdata[0] + .output + .iter() + .map(|txout| txout.consensus_encode_to_vec()) + .collect::>>() + .concat(); + + let new_template_msg_data = NewTemplate { + template_id, + future_template: is_future, + version: block.header.version.to_consensus() as u32, + coinbase_tx_version: block.txdata[0].version.0 as u32, + coinbase_prefix: block.txdata[0].input[0] + .script_sig + .to_bytes() + .try_into() + .map_err(|_| map_try_into_error("coinbase_prefix"))?, + coinbase_tx_input_sequence: block.txdata[0].input[0].sequence.to_consensus_u32(), + coinbase_tx_value_remaining: fee, + coinbase_tx_outputs_count: block.txdata[0].output.len() as u32, + coinbase_tx_outputs: coinbase_outputs_concat + .try_into() + .map_err(|_| map_try_into_error("coinbase_tx_outputs"))?, + coinbase_tx_locktime: block.txdata[0].lock_time.to_consensus_u32(), + merkle_path: merkle_path + .try_into() + .map_err(|_| map_try_into_error("merkle_path"))?, + }; + + let prev_hash_msg_data = SetNewPrevHash { + template_id, + prev_hash: block + .header + .prev_blockhash + .to_raw_hash() + .to_byte_array() + .try_into() + .unwrap(), + header_timestamp: block.header.time, + n_bits: block.header.bits.to_consensus(), + target: block.header.target().to_be_bytes().try_into().unwrap(), + }; + + let new_template_msg = into_static(AnyMessage::TemplateDistribution( + roles_logic_sv2::parsers::TemplateDistribution::NewTemplate(new_template_msg_data), + )); + let prev_hash_msg = into_static(AnyMessage::TemplateDistribution( + roles_logic_sv2::parsers::TemplateDistribution::SetNewPrevHash(prev_hash_msg_data), + )); + + let should_send_template = is_future + || (fee > (self.state.current_fee - self.sv2feedelta)) + || self.state.current_fee == 0; + + if is_future || should_send_template { + let template_frame = + Sv2Frame::from_message(new_template_msg.clone(), 0x71, 0x0001, false).unwrap(); + + if let Err(e) = self.template_sender.send(template_frame.into()).await { + error!( + "(Template {}) Failed to send NewTemplate downstream: {}", + template_id, e + ); + task::spawn_local(async move { + if let Err(e) = new_template_client.destroy().await { + error!( + "(Template {}) Failed to destroy block template after send failure ", + template_id + ); + } + }); + } else { + info!( + "(Template {}) Sent NewTemplate (is_future={})", + template_id, is_future + ); + self.state.insert_template(template_id, block); + self.state.current_fee = fee; + + if let Some(old_client) = self.active_template.take() { + task::spawn_local(async move { + if let Err(e) = old_client.destroy().await { + error!("Failed to destroy previous block template"); + } else { + debug!("Destroyed previous block template."); + } + }); + } + self.active_template = Some(new_template_client); + self.state.get_next_template_id(); + if is_future { + let prevhash_frame = + Sv2Frame::from_message(prev_hash_msg.clone(), 0x72, 0x0001, false).unwrap(); + if let Err(e) = self.prev_hash_sender.send(prevhash_frame.into()).await { + error!( + "(Template {}) Failed to send SetNewPrevHash downstream: {}", + template_id, e + ); + } else { + debug!("(Template {}) Sent SetNewPrevHash", template_id); + self.state.set_prev_hash_received(true); + } + } + } + } else { + debug!("(Template {}) New template fee ({}) does not meet threshold (current: {}), skipping send.", template_id, fee, self.state.current_fee); + task::spawn_local(async move { + if let Err(e) = new_template_client.destroy().await { + error!( + "(Template {}) Failed to destroy unused block template", + template_id + ); + } else { + debug!( + "(Template {}) Destroyed unused block template.", + template_id + ); + } + }); + } + + Ok(()) + } + + /// Handles a RequestTransactionData message from downstream. + async fn handle_request_transaction_data( + &self, + m: roles_logic_sv2::template_distribution_sv2::RequestTransactionData, + ) { + let response_msg = match self.state.templates.get(&m.template_id) { + Some(tx_data) => { + if tx_data.block.header.prev_blockhash != self.state.last_tip.prevhash { + warn!( + "Responding with stale-template-id error for template {}", + m.template_id + ); + let error = RequestTransactionDataError { + template_id: m.template_id, + error_code: b"stale-template-id"[..].to_vec().try_into().unwrap(), + }; + AnyMessage::TemplateDistribution( + roles_logic_sv2::parsers::TemplateDistribution::RequestTransactionDataError( + error, + ), + ) + } else { + info!("Providing transaction data for template {}", m.template_id); + // Ensure witness data exists and is not empty before accessing index 0 + let excess_data = tx_data + .block + .txdata + .get(0) + .and_then(|tx| tx.input.get(0)) + .map(|input| input.witness.to_vec().concat()) + .unwrap_or_else(|| { + warn!( + "Coinbase transaction or input missing/malformed for template {}", + m.template_id + ); + Vec::new() + }); + + let transaction_list: Vec = tx_data.block.txdata[1..] + .iter() + .map(|tx| B016M::Owned(tx.consensus_encode_to_vec())) + .collect(); + + let success = RequestTransactionDataSuccess { + template_id: m.template_id, + excess_data: excess_data + .try_into() + .unwrap_or_else(|_| Vec::new().try_into().unwrap()), + transaction_list: transaction_list + .try_into() + .unwrap_or_else(|_| Vec::new().try_into().unwrap()), + }; + AnyMessage::TemplateDistribution(roles_logic_sv2::parsers::TemplateDistribution::RequestTransactionDataSuccess(success)) + } + } + None => { + warn!( + "Responding with template-id-not-found error for template {}", + m.template_id + ); + let error = RequestTransactionDataError { + template_id: m.template_id, + error_code: b"template-id-not-found".to_vec().try_into().unwrap(), + }; + AnyMessage::TemplateDistribution( + roles_logic_sv2::parsers::TemplateDistribution::RequestTransactionDataError( + error, + ), + ) + } + }; + + let response_frame = + Sv2Frame::from_message(into_static(response_msg), 0x75, 0x0001, false).unwrap(); + + if self + .tx_data_response_sender + .send(response_frame.into()) + .await + .is_err() + { + error!("Failed to send transaction data response for template {} downstream, channel closed?", m.template_id); + } + } + + fn evict_old_templates(&mut self) { + let before_count = self.state.templates.len(); + if before_count > 0 { + self.state + .templates + .retain(|_key, value| value.timeout.elapsed() <= TEMPLATE_EVICTION_DURATION); + let after_count = self.state.templates.len(); + if before_count != after_count { + info!( + "Evicted {} old templates ({} remaining).", + before_count - after_count, + after_count + ); + } + } + } +} + +trait ConsensusEncodeVec { + fn consensus_encode_to_vec(&self) -> Vec; +} +impl ConsensusEncodeVec for T { + fn consensus_encode_to_vec(&self) -> Vec { + let mut writer = Vec::new(); + self.consensus_encode(&mut writer) + .expect("Vec encode failed"); + writer + } +} diff --git a/roles/template-provider-role/src/provider/state.rs b/roles/template-provider-role/src/provider/state.rs new file mode 100644 index 0000000000..cc5beecc64 --- /dev/null +++ b/roles/template-provider-role/src/provider/state.rs @@ -0,0 +1,95 @@ +use std::{collections::HashMap, time::Instant}; + +use stratum_common::bitcoin::{hashes::Hash, Block, BlockHash}; +use tracing::debug; + +const DEFAULT_MAX_ADDITIONAL_SIGOPS: u64 = 0; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ChainTip { + pub prevhash: BlockHash, + pub height: i32, +} + +impl Default for ChainTip { + fn default() -> Self { + ChainTip { + prevhash: BlockHash::from_byte_array([0; 32]), + height: -1, + } + } +} + +#[derive(Debug)] +pub struct CoolDownTransactionData { + pub block: Block, + pub timeout: Instant, +} + +#[derive(Debug)] +pub struct ProviderState { + template_id_counter: u64, + pub last_tip: ChainTip, + is_prev_hash_received: bool, + pub templates: HashMap, + pub current_fee: u64, + is_coinbase_constraint_received: bool, + coinbase_max_additional_sigops: Option, +} + +impl ProviderState { + pub fn new() -> Self { + ProviderState { + template_id_counter: 0, + last_tip: ChainTip::default(), + is_prev_hash_received: false, + templates: HashMap::new(), + current_fee: 0, + is_coinbase_constraint_received: false, + coinbase_max_additional_sigops: None, + } + } + + pub fn get_next_template_id(&mut self) -> u64 { + let id = self.template_id_counter; + self.template_id_counter += 1; + id + } + + pub fn peek_next_template_id(&self) -> u64 { + self.template_id_counter + 1 + } + + pub fn set_prev_hash_received(&mut self, received: bool) { + self.is_prev_hash_received = received; + } + + pub fn is_prev_hash_received(&self) -> bool { + self.is_coinbase_constraint_received + } + + pub fn set_coinbase_constraints(&mut self, max_additional_sigops: u64) { + self.coinbase_max_additional_sigops = Some(max_additional_sigops); + self.is_coinbase_constraint_received = true; + } + + pub fn is_coinbase_constraint_received(&self) -> bool { + self.is_coinbase_constraint_received + } + + pub fn get_coinbase_max_additional_sigops(&self) -> u64 { + self.coinbase_max_additional_sigops + .unwrap_or(DEFAULT_MAX_ADDITIONAL_SIGOPS) + } + + pub fn insert_template(&mut self, template_id: u64, block: Block) { + debug!("Storing block data for template ID {}", template_id); + self.templates.insert( + template_id, + CoolDownTransactionData { + block, + timeout: Instant::now(), + }, + ); + } +} diff --git a/roles/template-provider-role/src/rpc_client/error.rs b/roles/template-provider-role/src/rpc_client/error.rs new file mode 100644 index 0000000000..d2dba146f5 --- /dev/null +++ b/roles/template-provider-role/src/rpc_client/error.rs @@ -0,0 +1,34 @@ +#[derive(Debug)] +pub enum RpcClientError { + Connection(std::io::Error), + Capnp(capnp::Error), + Schema(capnp::NotInSchema), + Encode(stratum_common::bitcoin::consensus::encode::Error), + InvalidData(String), + OperationFailed(String), + NotFound, +} + +impl From for RpcClientError { + fn from(value: std::io::Error) -> Self { + Self::Connection(value) + } +} + +impl From for RpcClientError { + fn from(value: capnp::Error) -> Self { + Self::Capnp(value) + } +} + +impl From for RpcClientError { + fn from(value: capnp::NotInSchema) -> Self { + Self::Schema(value) + } +} + +impl From for RpcClientError { + fn from(value: stratum_common::bitcoin::consensus::encode::Error) -> Self { + Self::Encode(value) + } +} diff --git a/roles/template-provider-role/src/rpc_client/mod.rs b/roles/template-provider-role/src/rpc_client/mod.rs new file mode 100644 index 0000000000..a1f2a8d1d2 --- /dev/null +++ b/roles/template-provider-role/src/rpc_client/mod.rs @@ -0,0 +1,233 @@ +use crate::{common_capnp, init_capnp, mining_capnp, proxy_capnp}; +use capnp::capability::Promise; +use capnp_rpc::{rpc_twoparty_capnp, twoparty, RpcSystem}; +use error::RpcClientError; +use futures::FutureExt; +use std::{io::Cursor, path::Path, sync::Arc}; +use stratum_common::bitcoin::{consensus::Decodable, hashes::Hash, Block, BlockHash}; + +use tokio::{net::UnixStream, task}; +use tokio_util::compat::*; +use tracing::{debug, error, info}; + +pub mod error; + +pub struct BackendRpcClient { + bootstrap_client: init_capnp::init::Client, + rpc_execution_thread: Arc, +} + +impl BackendRpcClient { + pub async fn connect(socket_path: &Path) -> Result { + info!("Connection RPC client to socket: {}", socket_path.display()); + let stream = UnixStream::connect(socket_path).await?; + let (reader, writer) = stream.into_split(); + let reader = reader.compat(); + let writer = writer.compat_write(); + + let network = Box::new(twoparty::VatNetwork::new( + reader, + writer, + rpc_twoparty_capnp::Side::Client, + Default::default(), + )); + + let mut rpc_system = RpcSystem::new(network, None); + let bootstrap_client: init_capnp::init::Client = + rpc_system.bootstrap(rpc_twoparty_capnp::Side::Server); + + task::spawn_local(rpc_system.map(|_| debug!("RPC system finished"))); + info!("RPC System spawned locally."); + + let construct_request = bootstrap_client.construct_request(); + let construct_response = construct_request.send().promise.await?; + + let thread_map = construct_response.get()?.get_thread_map()?; + debug!("Construct response received. Requesting execution thread handle..."); + + let thread_response = thread_map.make_thread_request().send().promise.await?; + + let rpc_execution_thread = thread_response.get()?.get_result()?; + + info!("RPC execution thread handle obtained successfully."); + + Ok(Self { + bootstrap_client, + rpc_execution_thread: Arc::new(rpc_execution_thread), + }) + } + + async fn get_mining_service(&self) -> Result { + debug!("Requesting Mining service capability..."); + let mut request = self.bootstrap_client.make_mining_request(); + request + .get() + .get_context()? + .set_thread(self.rpc_execution_thread.as_ref().clone()); + let response = request.send().promise.await?; + response.get()?.get_result().map_err(RpcClientError::Capnp) + } + + pub async fn get_tip(&self) -> Result, RpcClientError> { + let mining = self.get_mining_service().await?; + let mut request = mining.get_tip_request(); + request + .get() + .get_context()? + .set_thread(self.rpc_execution_thread.as_ref().clone()); + + let response = request.send().promise.await?; + + if !response.get()?.get_has_result() { + debug!("Node reported no current tip available"); + return Ok(None); + } + + let result_reader = response.get()?.get_result()?; + + let hash_bytes = result_reader.get_hash()?; + let height = result_reader.get_height(); + + let block_hash = BlockHash::from_slice(hash_bytes) + .map_err(|e| RpcClientError::InvalidData(format!("Invalid tip hash bytes: {}", e)))?; + + debug!("Tip received: height={}, hash={}", height, block_hash); + Ok(Some((block_hash, height as u32))) + } + + pub async fn update_coinbase_constraints( + &self, + max_additional_sigops: u64, + reserved_weight: u64, + ) -> Result<(), RpcClientError> { + info!( + max_additional_sigops, + reserved_weight, "Updating coinbase constraints via RPC options." + ); + + let mining = self.get_mining_service().await?; + let mut create_block_req = mining.create_new_block_request(); + + let mut options = create_block_req.get().init_options(); + options.set_block_reserved_weight(reserved_weight); + options.set_coinbase_output_max_additional_sigops(max_additional_sigops); + options.set_use_mempool(true); + + let _ = create_block_req.send().promise.await?; + Ok(()) + } + + pub async fn create_new_block( + &self, + use_mempool: bool, + reserved_weight: u64, + max_additional_sigops: u64, + ) -> Result { + let mining = self.get_mining_service().await?; + let mut request = mining.create_new_block_request(); + + let mut options = request.get().init_options(); + options.set_use_mempool(use_mempool); + options.set_block_reserved_weight(reserved_weight); + options.set_coinbase_output_max_additional_sigops(max_additional_sigops); + let response = request.send().promise.await?; + let template_client_capnp = response.get()?.get_result()?; + debug!("Received BlockTemplate capability."); + + Ok(BlockTemplateClient::new( + template_client_capnp, + self.rpc_execution_thread.clone(), + )) + } +} + +pub struct BlockTemplateClient { + template_capnp: mining_capnp::block_template::Client, + rpc_execution_thread: Arc, +} + +impl BlockTemplateClient { + fn new( + template_capnp: mining_capnp::block_template::Client, + rpc_execution_thread: Arc, + ) -> Self { + Self { + template_capnp, + rpc_execution_thread, + } + } + + pub async fn get_block(&self) -> Result { + debug!("Requesting block data for template"); + let mut request = self.template_capnp.get_block_request(); + request + .get() + .get_context()? + .set_thread(self.rpc_execution_thread.as_ref().clone()); + let response = request.send().promise.await?; + let block_bytes = response.get()?.get_result()?; + debug!("Received raw block data ({} bytes).", block_bytes.len()); + let mut cursor = Cursor::new(block_bytes); + Block::consensus_decode(&mut cursor).map_err(RpcClientError::Encode) + } + + pub async fn get_coinbase_merkle_path(&self) -> Result>, RpcClientError> { + debug!("Requesting coinbase merkle path for template..."); + let mut request = self.template_capnp.get_coinbase_merkle_path_request(); + request + .get() + .get_context()? + .set_thread(self.rpc_execution_thread.as_ref().clone()); + let response = request.send().promise.await?; + let list_reader = response.get()?.get_result()?; + debug!("Received merkle path ({} elements).", list_reader.len()); + list_reader + .iter() + .map(|result| result.map(|bytes| bytes.to_vec())) + .collect::, _>>() + .map_err(RpcClientError::Capnp) + } + + pub async fn submit_solution( + &self, + version: u32, + timestamp: u32, + nonce: u32, + coinbase_txn_bytes: &[u8], + ) -> Result { + let mut request = self.template_capnp.submit_solution_request(); + request.get().set_version(version); + request.get().set_timestamp(timestamp); + request.get().set_nonce(nonce); + request.get().set_coinbase(coinbase_txn_bytes); + + request + .get() + .get_context()? + .set_thread(self.rpc_execution_thread.as_ref().clone()); + + let response = request.send().promise.await?; + let success = response.get()?.get_result(); + info!("Solution submission result: {}", success); + if success { + Ok(true) + } else { + Err(RpcClientError::OperationFailed( + "SubmitSolution returned false".to_string(), + )) + } + } + + pub async fn destroy(self) -> Result<(), RpcClientError> { + debug!("Destroying BlockTemplate capability..."); + let mut request = self.template_capnp.destroy_request(); + request + .get() + .get_context()? + .set_thread(self.rpc_execution_thread.as_ref().clone()); + + request.send().promise.await?; + debug!("BlockTemplate destroy request sent."); + Ok(()) + } +} diff --git a/roles/template-provider-role/src/setup_connection.rs b/roles/template-provider-role/src/setup_connection.rs new file mode 100644 index 0000000000..30c6cd8b4c --- /dev/null +++ b/roles/template-provider-role/src/setup_connection.rs @@ -0,0 +1,82 @@ +use super::{ + error::{TPError, TPResult}, + EitherFrame, StdFrame, +}; +use async_channel::{Receiver, Sender}; +use roles_logic_sv2::{ + common_messages_sv2::{SetupConnection, SetupConnectionSuccess}, + errors::Error, + handlers::common::{ParseCommonMessagesFromDownstream, SendTo}, + parsers::{AnyMessage, CommonMessages}, + utils::Mutex, +}; +use std::{convert::TryInto, net::SocketAddr, sync::Arc}; +use tracing::{debug, error, info}; + +pub struct SetupConnectionHandler; + +impl SetupConnectionHandler { + pub async fn setup( + receiver: &mut Receiver, + sender: &mut Sender, + address: SocketAddr, + ) -> TPResult<()> { + let mut incoming: StdFrame = match receiver.recv().await { + Ok(EitherFrame::Sv2(s)) => { + debug!("Got sv2 message: {:?}", s); + s + } + Ok(EitherFrame::HandShake(s)) => { + error!( + "Got unexpected handshake message from upstream: {:?} at {}", + s, address + ); + panic!() + } + Err(e) => { + error!("Error receiving message: {:?}", e); + return Err(Error::NoDownstreamsConnected.into()); + } + }; + + let message_type = incoming + .get_header() + .ok_or_else(|| TPError::Custom(String::from("No header set")))? + .msg_type(); + let payload = incoming.payload(); + let response = ParseCommonMessagesFromDownstream::handle_message_common( + Arc::new(Mutex::new(SetupConnectionHandler)), + message_type, + payload, + )?; + + let message = response.into_message().ok_or(TPError::RolesLogic( + roles_logic_sv2::Error::NoDownstreamsConnected, + ))?; + + let sv2_frame: StdFrame = AnyMessage::Common(message.clone()).try_into()?; + let sv2_frame = sv2_frame.into(); + sender.send(sv2_frame).await?; + + Ok(()) + } +} + +impl ParseCommonMessagesFromDownstream for SetupConnectionHandler { + fn handle_setup_connection( + &mut self, + incoming: SetupConnection, + ) -> Result { + info!( + "Received `SetupConnection`: version={}, flags={:b}", + incoming.min_version, incoming.flags + ); + Ok(SendTo::RelayNewMessageToRemote( + Arc::new(Mutex::new(())), + CommonMessages::SetupConnectionSuccess(SetupConnectionSuccess { + flags: incoming.flags, + used_version: 2, + }), + )) + } +}