diff --git a/Cargo.lock b/Cargo.lock index 4914d55cbef..3af73852114 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -405,13 +405,13 @@ checksum = "f4d270c927b83dbee3c6f5fb5340b38c6019452eeb6ea097f40148ec164fafe1" dependencies = [ "async-trait", "azure_core 1.1.0", - "fe2o3-amqp", - "fe2o3-amqp-cbs", - "fe2o3-amqp-ext", - "fe2o3-amqp-management", - "fe2o3-amqp-types", + "fe2o3-amqp 0.14.0", + "fe2o3-amqp-cbs 0.14.0", + "fe2o3-amqp-ext 0.14.0", + "fe2o3-amqp-management 0.14.0", + "fe2o3-amqp-types 0.14.0", "serde", - "serde_amqp", + "serde_amqp 0.14.1", "serde_bytes", "tokio", "tracing", @@ -425,16 +425,20 @@ version = "1.2.0-beta.1" dependencies = [ "async-trait", "azure_core 1.2.0-beta.1", - "fe2o3-amqp", - "fe2o3-amqp-cbs", - "fe2o3-amqp-ext", - "fe2o3-amqp-management", - "fe2o3-amqp-types", + "fe2o3-amqp 0.16.0", + "fe2o3-amqp-cbs 0.16.0", + "fe2o3-amqp-ext 0.16.0", + "fe2o3-amqp-management 0.16.0", + "fe2o3-amqp-types 0.16.0", + "fe2o3-amqp-ws", + "rustls", + "rustls-platform-verifier", "serde", - "serde_amqp", + "serde_amqp 0.16.0", "serde_bytes", "serde_json", "tokio", + "tokio-rustls", "tracing", "tracing-subscriber", "typespec 1.2.0-beta.1", @@ -477,18 +481,6 @@ dependencies = [ "tracing-subscriber", ] -[[package]] -name = "azure_core_opentelemetry" -version = "1.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b71bd721089367eee3399ab7820052421ddb35916327d70d14667303617dfd21" -dependencies = [ - "azure_core 1.1.0", - "opentelemetry 0.31.0", - "opentelemetry_sdk 0.31.0", - "tracing", -] - [[package]] name = "azure_core_opentelemetry" version = "1.1.0-beta.1" @@ -497,8 +489,8 @@ dependencies = [ "azure_core_test 0.2.0", "azure_core_test_macros 0.2.0", "azure_identity 1.1.0-beta.1", - "opentelemetry 0.32.0", - "opentelemetry_sdk 0.32.1", + "opentelemetry", + "opentelemetry_sdk", "tokio", "tracing", "tracing-subscriber", @@ -591,8 +583,8 @@ dependencies = [ "base64", "clap", "futures", - "opentelemetry 0.32.0", - "opentelemetry_sdk 0.32.1", + "opentelemetry", + "opentelemetry_sdk", "pin-project", "reqwest", "serde", @@ -710,10 +702,10 @@ dependencies = [ "azure_identity 1.0.0", "clap", "futures", - "opentelemetry 0.32.0", + "opentelemetry", "opentelemetry-otlp", "opentelemetry-stdout", - "opentelemetry_sdk 0.32.1", + "opentelemetry_sdk", "rand 0.10.2", "serde", "serde_json", @@ -799,16 +791,16 @@ dependencies = [ "async-lock", "async-stream", "async-trait", - "azure_core 1.1.0", - "azure_core_amqp 1.1.0", - "azure_core_test 0.2.0 (registry+https://github.com/rust-lang/crates.io-index)", - "azure_identity 1.0.0", + "azure_core 1.2.0-beta.1", + "azure_core_amqp 1.2.0-beta.1", + "azure_core_test 0.2.0", + "azure_identity 1.1.0-beta.1", "azure_messaging_eventhubs", "azure_messaging_eventhubs_checkpointstore_blob", - "azure_storage_blob 1.0.0", + "azure_storage_blob", "base64", "criterion", - "fe2o3-amqp", + "fe2o3-amqp 0.16.0", "futures", "hmac", "include-file", @@ -827,17 +819,17 @@ name = "azure_messaging_eventhubs_checkpointstore_blob" version = "0.9.0" dependencies = [ "async-trait", - "azure_core 1.1.0", - "azure_core_opentelemetry 1.0.0", - "azure_core_test 0.2.0 (registry+https://github.com/rust-lang/crates.io-index)", - "azure_identity 1.0.0", + "azure_core 1.2.0-beta.1", + "azure_core_opentelemetry", + "azure_core_test 0.2.0", + "azure_identity 1.1.0-beta.1", "azure_messaging_eventhubs", - "azure_storage_blob 1.0.0", + "azure_storage_blob", "futures", - "opentelemetry 0.32.0", + "opentelemetry", "opentelemetry-appender-tracing", "opentelemetry-stdout", - "opentelemetry_sdk 0.32.1", + "opentelemetry_sdk", "serde", "serde_json", "time", @@ -945,25 +937,6 @@ dependencies = [ "tokio", ] -[[package]] -name = "azure_storage_blob" -version = "1.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1756febbcca86c862ef718b983b505d08bd65a9bc984a915b0a16af4a4c3fe5b" -dependencies = [ - "async-stream", - "async-trait", - "azure_core 1.1.0", - "bytes", - "futures", - "percent-encoding", - "pin-project", - "serde", - "serde_json", - "time", - "tokio", -] - [[package]] name = "azure_storage_blob" version = "1.1.0-beta.2" @@ -1539,6 +1512,12 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "data-encoding" +version = "2.11.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06" + [[package]] name = "der" version = "0.8.1" @@ -1653,14 +1632,14 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a579ef4f1fb186f04bcdc9caf0c335adedebe879227c96d56876d473aa3d20a" dependencies = [ "bytes", - "fe2o3-amqp-types", + "fe2o3-amqp-types 0.14.0", "futures-util", "getrandom 0.3.4", "native-tls", "parking_lot", "pin-project-lite", "serde", - "serde_amqp", + "serde_amqp 0.14.1", "serde_bytes", "slab", "thiserror", @@ -1668,10 +1647,38 @@ dependencies = [ "tokio-native-tls", "tokio-stream", "tokio-util", + "url", + "uuid", + "wasmtimer", +] + +[[package]] +name = "fe2o3-amqp" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c1bf4f8851ad8c71afaa1aa164e0caa46b9b9282c78d9ba690b3ac3fa02ec9c6" +dependencies = [ + "bytes", + "fe2o3-amqp-types 0.16.0", + "futures-util", + "getrandom 0.4.3", + "parking_lot", + "pin-project-lite", + "rustls", + "serde", + "serde_amqp 0.16.0", + "serde_bytes", + "slab", + "thiserror", + "tokio", + "tokio-rustls", + "tokio-stream", + "tokio-util", "tracing", "url", "uuid", "wasmtimer", + "webpki-roots", ] [[package]] @@ -1680,8 +1687,19 @@ version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6cae904b214ffa3c9bae26e4129d300fe79189d2ef70503071fb25ff9127531e" dependencies = [ - "fe2o3-amqp", - "fe2o3-amqp-management", + "fe2o3-amqp 0.14.0", + "fe2o3-amqp-management 0.14.0", + "trait-variant", +] + +[[package]] +name = "fe2o3-amqp-cbs" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cc7df6781f8e68441039978dfccfb207abcb21c12cecc06d7012944c98b5af50" +dependencies = [ + "fe2o3-amqp 0.16.0", + "fe2o3-amqp-management 0.16.0", "trait-variant", ] @@ -1691,8 +1709,18 @@ version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6362c13b91a80dca77360eecdfbe85425e6dff1523f16c3a79fac541c47cf27d" dependencies = [ - "fe2o3-amqp-types", - "serde_amqp", + "fe2o3-amqp-types 0.14.0", + "serde_amqp 0.14.1", +] + +[[package]] +name = "fe2o3-amqp-ext" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5ab270e0b2e2a991a7b5809df3d250f24066e3ad3b9f68c31c91a695c36b9435" +dependencies = [ + "fe2o3-amqp-types 0.16.0", + "serde_amqp 0.16.0", ] [[package]] @@ -1701,8 +1729,20 @@ version = "0.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0582084762bdf022540c37868a0808e9f54dbcc51fe56f6212da59c167569cda" dependencies = [ - "fe2o3-amqp", - "fe2o3-amqp-types", + "fe2o3-amqp 0.14.0", + "fe2o3-amqp-types 0.14.0", + "serde", + "thiserror", +] + +[[package]] +name = "fe2o3-amqp-management" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b1fe67bcc1831af0046d2b56615845b887d3e89a5827990dceb259b82316385a" +dependencies = [ + "fe2o3-amqp 0.16.0", + "fe2o3-amqp-types 0.16.0", "serde", "thiserror", ] @@ -1715,11 +1755,44 @@ checksum = "8bcc8d13ed13fbb2fb664a6df114bcc32f8ca85c9cb6b89d4e7576c47f583706" dependencies = [ "ordered-float", "serde", - "serde_amqp", + "serde_amqp 0.14.1", + "serde_bytes", + "serde_repr", +] + +[[package]] +name = "fe2o3-amqp-types" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "63a759b57e923ff56dde0935a479be557b8681095e190f0cdf2a2777e8060cfd" +dependencies = [ + "ordered-float", + "serde", + "serde_amqp 0.16.0", "serde_bytes", "serde_repr", ] +[[package]] +name = "fe2o3-amqp-ws" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d97c7c11da9b00c1703fee21d5b1ca129ade1ae0d44d4b9f20eecdde29dbbe7e" +dependencies = [ + "bytes", + "futures-util", + "getrandom 0.4.3", + "http", + "js-sys", + "pin-project-lite", + "thiserror", + "tokio", + "tokio-tungstenite", + "tungstenite", + "wasm-bindgen", + "web-sys", +] + [[package]] name = "filetime" version = "0.2.29" @@ -2667,20 +2740,6 @@ dependencies = [ "vcpkg", ] -[[package]] -name = "opentelemetry" -version = "0.31.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b84bcd6ae87133e903af7ef497404dda70c60d0ea14895fc8a5e6722754fc2a0" -dependencies = [ - "futures-core", - "futures-sink", - "js-sys", - "pin-project-lite", - "thiserror", - "tracing", -] - [[package]] name = "opentelemetry" version = "0.32.0" @@ -2701,7 +2760,7 @@ version = "0.32.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2c0080f0dc1d7c786f467cd85a4e395fcab11ee852004f39a29a18ab7c25d837" dependencies = [ - "opentelemetry 0.32.0", + "opentelemetry", "tracing", "tracing-core", "tracing-subscriber", @@ -2714,9 +2773,9 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9966929966d17620d7c316c643ba62631826e10021409357772d5eea84f62c35" dependencies = [ "http", - "opentelemetry 0.32.0", + "opentelemetry", "opentelemetry-proto", - "opentelemetry_sdk 0.32.1", + "opentelemetry_sdk", "prost 0.14.4", "thiserror", "tokio", @@ -2730,8 +2789,8 @@ version = "0.32.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "56d658ba1faf63f7b9c492cfbe6e0ec365440a16132d3270c1065f7b33f1b638" dependencies = [ - "opentelemetry 0.32.0", - "opentelemetry_sdk 0.32.1", + "opentelemetry", + "opentelemetry_sdk", "prost 0.14.4", "tonic 0.14.6", "tonic-prost", @@ -2744,23 +2803,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a1b1c6a247d79091f0062a5f4bd058589525cf987a8d4c169440d9c1be72f0ad" dependencies = [ "chrono", - "opentelemetry 0.32.0", - "opentelemetry_sdk 0.32.1", -] - -[[package]] -name = "opentelemetry_sdk" -version = "0.31.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e14ae4f5991976fd48df6d843de219ca6d31b01daaab2dad5af2badeded372bd" -dependencies = [ - "futures-channel", - "futures-executor", - "futures-util", - "opentelemetry 0.31.0", - "percent-encoding", - "rand 0.9.5", - "thiserror", + "opentelemetry", + "opentelemetry_sdk", ] [[package]] @@ -2772,7 +2816,7 @@ dependencies = [ "futures-channel", "futures-executor", "futures-util", - "opentelemetry 0.32.0", + "opentelemetry", "percent-encoding", "portable-atomic", "rand 0.9.5", @@ -3610,11 +3654,27 @@ dependencies = [ "uuid", ] +[[package]] +name = "serde_amqp" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "110918b753cd22eaeb30632c7b73cc548b2a6258ec07a2bd5ea8ebf6365f8d69" +dependencies = [ + "bytes", + "indexmap 2.14.0", + "ordered-float", + "serde", + "serde_amqp_derive", + "serde_bytes", + "thiserror", + "uuid", +] + [[package]] name = "serde_amqp_derive" -version = "0.3.0" +version = "0.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "22da57ecf44834259b4416250608e11da620750be91305bf6ae5d398954ddc6d" +checksum = "b1e9b8826519d5a00c5de47e74ee76001a50276de6716a91fb40efc70a3c95fa" dependencies = [ "convert_case", "darling", @@ -3743,6 +3803,17 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "sha1" +version = "0.10.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8" +dependencies = [ + "cfg-if", + "cpufeatures 0.2.17", + "digest", +] + [[package]] name = "sha2" version = "0.10.9" @@ -3941,7 +4012,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.3", "once_cell", "rustix", "windows-sys 0.61.2", @@ -4112,6 +4183,22 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-tungstenite" +version = "0.26.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a9daff607c6d2bf6c16fd681ccb7eecc83e4e2cdc1ca067ffaadfca5de7f084" +dependencies = [ + "futures-util", + "log", + "rustls", + "rustls-native-certs", + "rustls-pki-types", + "tokio", + "tokio-rustls", + "tungstenite", +] + [[package]] name = "tokio-util" version = "0.7.18" @@ -4418,6 +4505,25 @@ version = "0.2.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" +[[package]] +name = "tungstenite" +version = "0.26.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4793cb5e56680ecbb1d843515b23b6de9a75eb04b66643e256a396d43be33c13" +dependencies = [ + "bytes", + "data-encoding", + "http", + "httparse", + "log", + "rand 0.9.5", + "rustls", + "rustls-pki-types", + "sha1", + "thiserror", + "utf-8", +] + [[package]] name = "typed-path" version = "0.12.3" @@ -4439,7 +4545,6 @@ dependencies = [ "base64", "bytes", "futures", - "quick-xml", "serde", "serde_json", "url", @@ -4596,6 +4701,12 @@ dependencies = [ "serde", ] +[[package]] +name = "utf-8" +version = "0.7.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09cc8ee72d2a9becf2f2febe0205bbed8fc6615b7cb429ad062dc7b7ddd036a9" + [[package]] name = "utf8-zero" version = "0.8.1" diff --git a/Cargo.toml b/Cargo.toml index 4a5b4892f2a..6600d3006c7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -104,11 +104,12 @@ criterion = { version = "0.8", features = ["async_tokio"] } crossbeam = { version = "0.8", default-features = false } crossbeam-epoch = { version = "0.9", default-features = false } dyn-clone = "1.0" -fe2o3-amqp = { version = "0.14", features = ["uuid"] } -fe2o3-amqp-ext = { version = "0.14" } -fe2o3-amqp-management = { version = "0.14" } -fe2o3-amqp-cbs = { version = "0.14" } -fe2o3-amqp-types = { version = "0.14" } +fe2o3-amqp = { version = "0.16", features = ["uuid"] } +fe2o3-amqp-ws = { version = "0.16" } +fe2o3-amqp-ext = { version = "0.16" } +fe2o3-amqp-management = { version = "0.16" } +fe2o3-amqp-cbs = { version = "0.16" } +fe2o3-amqp-types = { version = "0.16" } flate2 = "1.1.9" futures = "0.3" getrandom = { version = "0.4" } @@ -136,8 +137,10 @@ reqwest = { version = "0.13.2", features = [ ], default-features = false } rust_decimal = "1.40.0" rustc_version = "0.4" +rustls = "0.23" +rustls-platform-verifier = "0.7" serde = { version = "1.0", features = ["derive"] } -serde_amqp = { version = "0.14", features = ["uuid"] } +serde_amqp = { version = "0.16", features = ["uuid"] } serde_bytes = { version = "0.11" } serde_json = "1.0.149" serde_test = "1" @@ -156,6 +159,7 @@ tokio = { version = "1.49", default-features = false, features = [ "time", ] } tokio-metrics = "0.4" +tokio-rustls = { version = "0.26", default-features = false } tracing = "0.1.44" tracing-subscriber = "0.3" url = "2.5" diff --git a/eng/dict/crates.txt b/eng/dict/crates.txt index 53900c0621d..5da45bf457b 100644 --- a/eng/dict/crates.txt +++ b/eng/dict/crates.txt @@ -54,11 +54,13 @@ fe2o3_amqp_cbs fe2o3_amqp_ext fe2o3_amqp_management fe2o3_amqp_types +fe2o3_amqp_ws fe2o3-amqp fe2o3-amqp-cbs fe2o3-amqp-ext fe2o3-amqp-management fe2o3-amqp-types +fe2o3-amqp-ws flate2 futures getrandom @@ -113,6 +115,8 @@ time tokio tokio_metrics tokio-metrics +tokio_tungstenite +tokio-tungstenite tracing tracing_subscriber tracing-subscriber diff --git a/sdk/core/azure_core_amqp/CHANGELOG.md b/sdk/core/azure_core_amqp/CHANGELOG.md index 5f4bba0d71d..e8f6113858c 100644 --- a/sdk/core/azure_core_amqp/CHANGELOG.md +++ b/sdk/core/azure_core_amqp/CHANGELOG.md @@ -4,8 +4,15 @@ ### Features Added +- Added `AmqpTransport` and an `AmqpConnectionOptions::transport` field to select the connection transport. `AmqpTransport::WebSocket` tunnels AMQP over secure WebSockets (`wss://`, port 443) for networks that block the native AMQP ports. +- Added the `fe2o3_amqp_ws` and `fe2o3_amqp_ws_rustls` features. `fe2o3_amqp_ws` is the base feature and turns on the WebSocket transport code, and `fe2o3_amqp_ws_rustls` adds the TLS stack that the rest of `sdk/core` uses, rustls with the aws-lc-rs provider. This is the shape that `azure_core` uses for HTTP, where `reqwest` is the base and `reqwest_rustls` adds the stack. The `default` feature selects both. To build the transport on another TLS stack, turn off the default features, name `fe2o3_amqp_ws`, and take a direct dependency on `fe2o3-amqp-ws` with the stack you want; Cargo unifies the features. One stack must be selected somewhere in the graph, and a build that selects none still compiles and reports the missing stack when the connection opens. A build without `fe2o3_amqp_ws` still accepts `AmqpTransport::WebSocket`, and the connection then returns an error when it opens. +- Added the `fe2o3_amqp_rustls` feature, which adds the TLS stack for AMQP framed directly on TCP (`amqps://`, port 5671) on top of the `fe2o3_amqp` base feature. It is rustls with the aws-lc-rs provider, the stack that the rest of `sdk/core` uses, and `default` selects it in place of the `fe2o3-amqp/native-tls` entry that `default` named before. The connection supplies a TLS connector built on `rustls-platform-verifier`, so the handshake reads the trust store of the operating system and trusts the same roots as `reqwest/rustls` on the HTTP side. The default connector of `fe2o3-amqp` instead fills its root store from `webpki-roots`, a compiled-in copy of the Mozilla root set, which would drop the roots that an operator installs in the operating system. There is no native-tls feature on this crate, because `fe2o3-amqp` accepts one TLS stack only, and two features that cannot both be on would break `--all-features`. ([#4189](https://github.com/Azure/azure-sdk-for-rust/issues/4189)) + ### Breaking Changes +- Added the `transport` field to `AmqpConnectionOptions`. The struct is not `#[non_exhaustive]`, so an existing struct literal that names every field no longer compiles. Add `..Default::default()` to the initializer, and put `#[allow(clippy::needless_update)]` on it. The workspace allows the `constructible_struct_adds_field` semver lint for this reason, so `cargo semver-checks` does not report the addition. +- The `default` feature now selects `fe2o3_amqp_rustls`, so AMQP framed directly on TCP (`amqps://`, port 5671) runs on rustls with the aws-lc-rs provider where it ran on native-tls. Both stacks read the trust store of the operating system, so a broker behind a private or an enterprise certificate authority keeps working. The stacks read that store through different platform APIs, and a deployment that tunes native-tls directly, such as one that sets OpenSSL environment variables, can still see a difference. To keep native-tls, turn off the default features, name `fe2o3_amqp`, and take a direct dependency on `fe2o3-amqp` with its `native-tls` feature. ([#4189](https://github.com/Azure/azure-sdk-for-rust/issues/4189)) + ### Bugs Fixed - Link properties set through `AmqpReceiverOptions::properties` and `AmqpSenderOptions::properties` now reach the Attach frame. They were discarded before the link attached. @@ -13,6 +20,8 @@ ### Other Changes +- Updated the `fe2o3-amqp` family of dependencies from 0.14 to 0.16. The rustls backend of 0.16 is built on aws-lc-rs, where 0.14 was built on `ring`, which `deny.toml` bans. This is what makes `fe2o3_amqp_rustls` possible. The update needed no source change. + ## 1.1.0 (2026-07-09) ### Features Added diff --git a/sdk/core/azure_core_amqp/Cargo.toml b/sdk/core/azure_core_amqp/Cargo.toml index c50a8fc8428..86fcebd684b 100644 --- a/sdk/core/azure_core_amqp/Cargo.toml +++ b/sdk/core/azure_core_amqp/Cargo.toml @@ -21,14 +21,18 @@ edition.workspace = true async-trait.workspace = true azure_core = { path = "../azure_core", version = "1.2.0-beta.1", default-features = false } fe2o3-amqp = { workspace = true, optional = true } +fe2o3-amqp-ws = { workspace = true, optional = true } fe2o3-amqp-cbs = { workspace = true, optional = true } fe2o3-amqp-ext = { workspace = true, optional = true } fe2o3-amqp-management = { workspace = true, optional = true } fe2o3-amqp-types = { workspace = true, optional = true } +rustls = { workspace = true, optional = true } +rustls-platform-verifier = { workspace = true, optional = true } serde.workspace = true serde_amqp = { workspace = true, optional = true } serde_bytes = { workspace = true, optional = true } tokio.workspace = true +tokio-rustls = { workspace = true, optional = true } tracing.workspace = true typespec = { path = "../typespec", version = "1.2.0-beta.1" } typespec_macros = { path = "../typespec_macros", version = "1.1.0-beta.1" } @@ -38,7 +42,12 @@ serde_json.workspace = true tracing-subscriber = { workspace = true, features = ["env-filter"] } [features] -default = ["fe2o3_amqp", "fe2o3-amqp/native-tls"] +default = [ + "fe2o3_amqp", + "fe2o3_amqp_rustls", + "fe2o3_amqp_ws", + "fe2o3_amqp_ws_rustls", +] ffi = [] test = [] fe2o3_amqp = [ @@ -51,9 +60,100 @@ fe2o3_amqp = [ "serde_bytes", "azure_core/tokio", ] +# The TLS stack for AMQP framed directly on TCP (`amqps://`, port 5671). +# `fe2o3_amqp` is the base, and this feature adds the stack, in the same shape +# that `azure_core` uses for HTTP, where `reqwest` is the base and +# `reqwest_rustls` adds the stack. It is rustls with the aws-lc-rs provider, and +# `default` selects it. +# +# The connection builds its own TLS connector on `rustls-platform-verifier`. +# The default connector of `fe2o3-amqp` fills its root store from +# `webpki-roots`, which holds a compiled-in copy of the Mozilla root set and +# ignores the trust store of the operating system. The platform verifier reads +# the trust store of the operating system, so this feature trusts the same roots +# as `reqwest/rustls` on the HTTP side, and a broker behind a private or an +# enterprise certificate authority keeps working. `webpki-roots` still arrives +# as a dependency, because the `rustls` feature of `fe2o3-amqp` always names it, +# but nothing reads it. +# +# The direct `rustls` dependency selects the crypto provider. +# `ClientConfig::builder()` panics when no process-level default is installed +# and the crate features name no single provider. The default features of +# `rustls` give aws-lc-rs, std, and tls12, and `deny.toml` bans `ring`, so this +# repository cannot reach that panic. An application that unifies a second +# provider into the graph makes the choice ambiguous again, and it must call +# `CryptoProvider::install_default()` before it opens a connection. The +# WebSocket feature below carries the same condition. +# +# This feature needs `fe2o3-amqp` 0.16 or later. The rustls backend of 0.14 was +# built on `ring`, which `deny.toml` bans (#4189). +# +# There is no matching native-tls feature, because `fe2o3-amqp` accepts one TLS +# stack only and reports `TlsConnectorNotFound` at run time when both are on. +# Two features that cannot both be on would break `--all-features`. To use +# another stack, turn off the default features, name `fe2o3_amqp`, and take a +# direct dependency on `fe2o3-amqp` with the stack you want. Cargo unifies the +# features, and nothing pulls rustls in: +# +# azure_core_amqp = { version = "...", default-features = false, features = [ +# "fe2o3_amqp", +# ] } +# fe2o3-amqp = { version = "0.16", features = ["native-tls"] } +# +# `samples/list_blobs_native_tls` shows the same pattern for `reqwest`. +fe2o3_amqp_rustls = [ + "fe2o3_amqp", + "fe2o3-amqp/rustls", + "dep:rustls", + "dep:rustls-platform-verifier", + "dep:tokio-rustls", +] +# AMQP over WebSockets. `fe2o3_amqp_ws` is the base feature: it turns on the +# transport code and the `fe2o3-amqp-ws` dependency, and it names no TLS stack. +# `fe2o3_amqp_ws_rustls` adds the stack that the rest of `sdk/core` uses, rustls +# with the aws-lc-rs provider, and `default` selects it. +# +# To build the transport on another TLS stack, turn off the default features, +# name `fe2o3_amqp_ws`, and take a direct dependency on the backend crate with +# the stack you want. Cargo unifies the features, so the transport then uses +# your stack and nothing pulls rustls in: +# +# azure_core_amqp = { version = "...", default-features = false, features = [ +# "fe2o3_amqp", "fe2o3_amqp_ws", +# ] } +# fe2o3-amqp = { version = "0.16", features = ["native-tls"] } +# fe2o3-amqp-ws = { version = "0.16", features = ["native-tls"] } +# +# `samples/list_blobs_native_tls` shows the same pattern for `reqwest`. +# +# One stack must be selected somewhere in the graph. `fe2o3_amqp_ws` on its own +# builds, and the connection reports `TlsFeatureNotEnabled` when it opens. +# `reqwest` behaves the same way. +fe2o3_amqp_ws = ["fe2o3_amqp", "dep:fe2o3-amqp-ws"] +# The direct `rustls` dependency is only here to select the crypto provider. +# This crate names no `rustls` type. `reqwest/rustls` selects aws-lc-rs itself, +# so `typespec_client_core` needs no such dependency. The `tokio-tungstenite` +# chain under `fe2o3-amqp-ws` is different. It takes rustls with +# `default-features = false` and names no provider, so `ClientConfig::builder()` +# panics when no process-level default is installed. The default features of +# `rustls` give aws-lc-rs, std, and tls12. +# +# `fe2o3_amqp_rustls` selects the same provider, so the dependency is redundant +# when `default` is on. This feature has to stand on its own, because a consumer +# can turn the default features off and name only the WebSocket transport. +fe2o3_amqp_ws_rustls = [ + "fe2o3_amqp_ws", + "fe2o3-amqp-ws/rustls-tls-native-roots", + "dep:rustls", +] [lints] workspace = true [package.metadata.docs.rs] -features = ["fe2o3_amqp"] +features = [ + "fe2o3_amqp", + "fe2o3_amqp_rustls", + "fe2o3_amqp_ws", + "fe2o3_amqp_ws_rustls", +] diff --git a/sdk/core/azure_core_amqp/src/connection.rs b/sdk/core/azure_core_amqp/src/connection.rs index ba15a6329c8..d77674affc6 100644 --- a/sdk/core/azure_core_amqp/src/connection.rs +++ b/sdk/core/azure_core_amqp/src/connection.rs @@ -14,7 +14,46 @@ type ConnectionImplementation = super::fe2o3::connection::Fe2o3AmqpConnection; #[cfg(not(feature = "fe2o3_amqp"))] type ConnectionImplementation = super::noop::NoopAmqpConnection; +/// The transport used to carry the AMQP protocol. +/// +/// AMQP is normally framed directly over a TCP/TLS socket, but some network +/// environments (for example corporate firewalls) only permit outbound +/// connections on port 443. In those cases the AMQP frames can be tunneled +/// over a WebSocket connection instead. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub enum AmqpTransport { + /// AMQP framing over a TCP/TLS socket (port 5671). This is the default. + #[default] + Tcp, + /// AMQP framing tunneled over secure WebSockets (`wss://`, port 443). + /// + /// This variant needs the `fe2o3_amqp_ws` feature, which the `default` + /// feature selects. A build without it can still select this variant, but + /// the connection then returns an error when it opens. + /// + /// `fe2o3_amqp_ws` names no TLS stack. `default` adds + /// `fe2o3_amqp_ws_rustls`, which selects rustls with the aws-lc-rs + /// provider. To use another stack, turn off the default features, name + /// `fe2o3_amqp_ws`, and take a direct dependency on `fe2o3-amqp-ws` with + /// the stack you want. A build that selects no stack also compiles, and + /// the connection reports the missing stack when it opens. + WebSocket, +} + /// Options for configuring an AMQP connection. +/// +/// Build it from [`Default`] and set only the fields you need, so a later field +/// addition does not break the call site: +/// +/// ``` +/// use azure_core_amqp::{AmqpConnectionOptions, AmqpTransport}; +/// +/// #[allow(clippy::needless_update)] +/// let options = AmqpConnectionOptions { +/// transport: Some(AmqpTransport::WebSocket), +/// ..Default::default() +/// }; +/// ``` #[derive(Debug, Default, Clone)] pub struct AmqpConnectionOptions { /// Maximum frame size for the connection in bytes. @@ -36,11 +75,17 @@ pub struct AmqpConnectionOptions { /// Buffer size for the connection. pub buffer_size: Option, /// Custom endpoint for the connection. Used to connect to a local AMQP proxy server. + /// + /// The host and an explicit port both carry into the address that the + /// connection dials. Under [`AmqpTransport::WebSocket`] the port carries into + /// the `wss://` address, so name the port that the proxy accepts WebSockets + /// on, and leave the port out to dial the default port 443. The .NET Azure SDK + /// treats `CustomEndpointAddress` the same way. pub custom_endpoint: Option, + /// The transport used to carry the AMQP protocol. Defaults to [`AmqpTransport::Tcp`]. + pub transport: Option, } -impl AmqpConnectionOptions {} - /// Trait defining the asynchronous APIs for AMQP connection operations. #[async_trait::async_trait] pub trait AmqpConnectionApis { @@ -250,6 +295,7 @@ mod tests { .collect(), ), buffer_size: Some(1024), + transport: Some(AmqpTransport::WebSocket), }; assert_eq!(connection_options.max_frame_size, Some(1024)); @@ -281,6 +327,7 @@ mod tests { connection_options.custom_endpoint, Some(Url::parse("http://localhost:8080").unwrap()) ); + assert_eq!(connection_options.transport, Some(AmqpTransport::WebSocket)); } // On macOS, there is a periodic issue where loopback TCP connections fail. diff --git a/sdk/core/azure_core_amqp/src/fe2o3/connection.rs b/sdk/core/azure_core_amqp/src/fe2o3/connection.rs index a4a71c59bcc..67b545a1294 100644 --- a/sdk/core/azure_core_amqp/src/fe2o3/connection.rs +++ b/sdk/core/azure_core_amqp/src/fe2o3/connection.rs @@ -1,8 +1,10 @@ // Copyright (c) Microsoft Corporation. All Rights reserved // Licensed under the MIT license. +#[cfg(feature = "fe2o3_amqp_ws")] +use crate::fe2o3::error::Fe2o3WebSocketError; use crate::{ - connection::{AmqpConnectionApis, AmqpConnectionOptions}, + connection::{AmqpConnectionApis, AmqpConnectionOptions, AmqpTransport}, error::{AmqpErrorKind, Result}, fe2o3::error::{Fe2o3ConnectionError, Fe2o3ConnectionOpenError, Fe2o3TransportError}, value::{AmqpOrderedMap, AmqpSymbol, AmqpValue}, @@ -44,6 +46,51 @@ impl Drop for Fe2o3AmqpConnection { } } +// cspell:ignore servicebus + +/// The well-known path that Service Bus and Event Hubs expose for the AMQP +/// WebSocket binding. Matches the suffix used by the other Azure SDKs. +#[cfg(feature = "fe2o3_amqp_ws")] +const WEBSOCKET_PATH: &str = "/$servicebus/websocket/"; + +/// Builds the secure WebSocket (`wss://`) address used to tunnel AMQP for the +/// given connection target. The target is the AMQP service URL (or a custom +/// endpoint proxy). Its scheme and path are discarded: only the host and an +/// explicit port (if any) are carried over, since AMQP-over-WebSockets always +/// uses TLS and a fixed binding path. When no port is present the default +/// `wss` port (443) is used. +#[cfg(feature = "fe2o3_amqp_ws")] +fn websocket_address(target: &Url) -> Result { + let host = target + .host_str() + .ok_or_else(|| AmqpError::with_message("AMQP connection URL is missing a host."))?; + let authority = match target.port() { + Some(port) => format!("{host}:{port}"), + None => host.to_string(), + }; + Ok(format!("wss://{authority}{WEBSOCKET_PATH}")) +} + +/// Builds the TLS connector for AMQP framed directly on TCP. +/// +/// The default connector of `fe2o3-amqp` fills its root store from +/// `webpki-roots`, a compiled-in copy of the Mozilla root set, and it ignores +/// the trust store of the operating system. This connector uses the platform +/// verifier instead, so the TCP transport trusts the same roots as the HTTP +/// stack of `azure_core`, and a broker behind a private or an enterprise +/// certificate authority keeps working. +#[cfg(feature = "fe2o3_amqp_rustls")] +fn platform_verifier_connector() -> Result { + use rustls_platform_verifier::ConfigVerifierExt as _; + + let config = rustls::ClientConfig::with_platform_verifier().map_err(|e| { + AmqpError::with_message(format!("Could not build the AMQP TLS configuration: {e}")) + })?; + Ok(tokio_rustls::TlsConnector::from(std::sync::Arc::new( + config, + ))) +} + #[async_trait::async_trait] impl AmqpConnectionApis for Fe2o3AmqpConnection { async fn open( @@ -111,18 +158,91 @@ impl AmqpConnectionApis for Fe2o3AmqpConnection { builder = builder.buffer_size(buffer_size); } - if let Some(custom_endpoint) = options.custom_endpoint { - endpoint = custom_endpoint; - builder = builder.hostname(url.host_str()); - } + let handle = match options.transport.unwrap_or_default() { + AmqpTransport::Tcp => { + // `custom_endpoint` redirects the socket to a proxy while the + // AMQP `hostname` stays the real service host. + if let Some(custom_endpoint) = options.custom_endpoint { + endpoint = custom_endpoint; + builder = builder.hostname(url.host_str()); + } - self.connection - .set(Mutex::new( - builder - .open(endpoint) + // Supply the connector so that the handshake uses the trust + // store of the operating system. See + // `platform_verifier_connector`. Without a connector, + // `fe2o3-amqp` falls back to its `webpki-roots` default. + #[cfg(feature = "fe2o3_amqp_rustls")] + { + builder + .rustls_connector(platform_verifier_connector()?) + .open(endpoint) + .await + .map_err(|e| AmqpError::from(Fe2o3ConnectionOpenError(e)))? + } + + // Another TLS stack, selected through a direct dependency on + // `fe2o3-amqp`, keeps the default connector of that stack. + #[cfg(not(feature = "fe2o3_amqp_rustls"))] + { + builder + .open(endpoint) + .await + .map_err(|e| AmqpError::from(Fe2o3ConnectionOpenError(e)))? + } + } + AmqpTransport::WebSocket => { + // A build without the transport code still accepts the variant, + // so report the missing feature rather than fail to compile. + #[cfg(not(feature = "fe2o3_amqp_ws"))] + { + Err(AmqpError::with_message( + "The WebSocket transport needs the `fe2o3_amqp_ws` feature of \ + azure_core_amqp.", + ))?; + unreachable!() + } + + // Tunnel AMQP over a secure WebSocket (port 443) for networks + // that block the native AMQP ports. The socket connects to the + // websocket address (or the custom endpoint proxy, if set), + // while the AMQP `hostname` remains the real service host. + // `open_with_stream` does not derive the hostname from a URL, + // so it must be set explicitly. + // + // `connect_with_config` takes no connector, so `fe2o3-amqp-ws` + // uses whichever TLS stack its own features select. The + // `fe2o3_amqp_ws_rustls` feature of this crate selects rustls, + // and a direct dependency on `fe2o3-amqp-ws` in the application + // can select another stack instead. + // + // `connect_tls_with_config` does the same thing, but it sits + // behind the TLS features of `fe2o3-amqp-ws`, so a call to it + // would keep `fe2o3_amqp_ws` from building on its own. + // `connect_with_config` carries no such gate, and it reports + // `TlsFeatureNotEnabled` when no stack is selected. + #[cfg(feature = "fe2o3_amqp_ws")] + { + let ws_target = options.custom_endpoint.as_ref().unwrap_or(&url); + let ws_address = websocket_address(ws_target)?; + debug!("Opening AMQP-over-WebSockets connection to {ws_address}."); + let ws_stream = fe2o3_amqp_ws::WebSocketStream::connect_with_config( + &ws_address, + None, + false, + ) .await - .map_err(|e| AmqpError::from(Fe2o3ConnectionOpenError(e)))?, - )) + .map_err(|e| AmqpError::from(Fe2o3WebSocketError(e)))?; + builder + .hostname(url.host_str()) + .open_with_stream(ws_stream) + .await + .map_err(|e| AmqpError::from(Fe2o3ConnectionOpenError(e)))? + } + } + }; + + self.connection + .set(Mutex::new(handle)) .map_err(|_| Self::connection_already_set())?; Ok(()) } @@ -221,3 +341,67 @@ impl From for AmqpError { } } } + +#[cfg(all(test, feature = "fe2o3_amqp_rustls"))] +mod rustls_tests { + use super::*; + + #[test] + fn platform_verifier_connector_builds() { + // `ClientConfig::builder()` panics when the process has no default + // crypto provider, and the platform verifier reports an error when it + // cannot read the trust store of the operating system. Both faults + // would otherwise appear only when a connection opens. + assert!(platform_verifier_connector().is_ok()); + } +} + +#[cfg(all(test, feature = "fe2o3_amqp_ws"))] +mod tests { + use super::*; + + #[test] + fn websocket_address_uses_default_port_for_service_url() { + // The Event Hubs connection URL has no explicit port; the wss default + // (443) is implied and the binding path is appended. + let url = Url::parse("amqps://my-namespace.servicebus.windows.net/my-eventhub").unwrap(); + assert_eq!( + websocket_address(&url).unwrap(), + "wss://my-namespace.servicebus.windows.net/$servicebus/websocket/" + ); + } + + #[test] + fn websocket_address_preserves_explicit_port() { + // A custom endpoint (e.g. a local proxy) may carry an explicit port, + // which must be preserved in the websocket address. + let proxy = Url::parse("amqps://localhost:8081/").unwrap(); + assert_eq!( + websocket_address(&proxy).unwrap(), + "wss://localhost:8081/$servicebus/websocket/" + ); + } + + #[test] + fn websocket_address_preserves_an_amqp_port() { + // A custom endpoint that names an AMQP port keeps it. The .NET Azure SDK + // carries `CustomEndpointAddress` into the WebSocket address the same way, + // so a proxy that accepts WebSockets on 5671 stays reachable. A caller who + // wants port 443 leaves the port out. + let proxy = Url::parse("amqps://proxy.example.com:5671/").unwrap(); + assert_eq!( + websocket_address(&proxy).unwrap(), + "wss://proxy.example.com:5671/$servicebus/websocket/" + ); + } + + #[test] + fn websocket_address_keeps_brackets_around_ipv6_host() { + // `Url::host_str` keeps the brackets around an IPv6 literal, so the + // authority stays valid when the host and the port are joined. + let proxy = Url::parse("amqps://[::1]:8081/").unwrap(); + let address = websocket_address(&proxy).unwrap(); + assert_eq!(address, "wss://[::1]:8081/$servicebus/websocket/"); + assert!(Url::parse(&address).is_ok()); + } +} diff --git a/sdk/core/azure_core_amqp/src/fe2o3/error.rs b/sdk/core/azure_core_amqp/src/fe2o3/error.rs index a5a57524933..0a74307d625 100644 --- a/sdk/core/azure_core_amqp/src/fe2o3/error.rs +++ b/sdk/core/azure_core_amqp/src/fe2o3/error.rs @@ -46,6 +46,15 @@ impl From for Fe2o3TransportError { } } +#[cfg(feature = "fe2o3_amqp_ws")] +pub(crate) struct Fe2o3WebSocketError(pub fe2o3_amqp_ws::Error); +#[cfg(feature = "fe2o3_amqp_ws")] +impl From for Fe2o3WebSocketError { + fn from(e: fe2o3_amqp_ws::Error) -> Self { + Fe2o3WebSocketError(e) + } +} + // Specializations of From for common AMQP types. impl From<&fe2o3_amqp_types::definitions::ErrorCondition> for AmqpErrorCondition { fn from(e: &fe2o3_amqp_types::definitions::ErrorCondition) -> Self { @@ -157,6 +166,17 @@ impl From for AmqpError { } } +#[cfg(feature = "fe2o3_amqp_ws")] +impl From for AmqpError { + fn from(e: Fe2o3WebSocketError) -> Self { + // The websocket establishment error wraps the underlying WebSocket and + // I/O failures that occur before the AMQP protocol handshake begins. + // There is no AMQP error condition for these, so they map to a + // transport-layer error. + AmqpError::from(AmqpErrorKind::TransportImplementationError(Box::new(e.0))) + } +} + impl From for AmqpError { fn from(e: Fe2o3TransportError) -> Self { match e.0 { diff --git a/sdk/core/azure_core_amqp/src/lib.rs b/sdk/core/azure_core_amqp/src/lib.rs index e691eb96e72..391f0489188 100644 --- a/sdk/core/azure_core_amqp/src/lib.rs +++ b/sdk/core/azure_core_amqp/src/lib.rs @@ -24,7 +24,7 @@ mod simple_value; mod value; pub use cbs::{AmqpClaimsBasedSecurity, AmqpClaimsBasedSecurityApis}; -pub use connection::{AmqpConnection, AmqpConnectionApis, AmqpConnectionOptions}; +pub use connection::{AmqpConnection, AmqpConnectionApis, AmqpConnectionOptions, AmqpTransport}; pub use error::*; pub use management::{AmqpManagement, AmqpManagementApis}; pub use messaging::{AmqpDelivery, AmqpDeliveryApis, AmqpMessage, AmqpSource, AmqpTarget}; diff --git a/sdk/eventhubs/azure_messaging_eventhubs/CHANGELOG.md b/sdk/eventhubs/azure_messaging_eventhubs/CHANGELOG.md index 8575377c254..af383a77018 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/CHANGELOG.md +++ b/sdk/eventhubs/azure_messaging_eventhubs/CHANGELOG.md @@ -5,12 +5,15 @@ ### Features Added - Added connection-string authentication. `ProducerClientBuilder` and `ConsumerClientBuilder` now have an `open_with_connection_string` method that authenticates with a Shared Access Signature parsed from an Event Hubs connection string (`Endpoint=sb://...;SharedAccessKeyName=...;SharedAccessKey=...`, optionally with `EntityPath`, or a pre-formed `SharedAccessSignature`). The connection-string parser is exposed publicly as `ConnectionString`. This reaches parity with the other Azure SDKs for development and test scenarios; Microsoft Entra ID via `open` with a `TokenCredential` remains the recommended path for production. The parser rejects empty required values and empty Event Hub names up front, and a pre-formed `SharedAccessSignature` reports its own `se` as the token expiry (rather than a rolling client-side window); because such a token cannot be renewed, the connection's token refresher detects the non-advancing expiry and leaves the broker to enforce it. ([#3459](https://github.com/Azure/azure-sdk-for-rust/issues/3459)) +- Added a `with_transport` builder method on `ProducerClient` and `ConsumerClient`, which takes the `AmqpTransport` of `azure_core_amqp` (re-exported as `models::AmqpTransport`). `AmqpTransport::WebSocket` tunnels AMQP over secure WebSockets (`wss://`, port 443), allowing clients to connect from networks that block the native AMQP ports (5671/5672). This matches the transport option offered by the .NET, Java, and Python Azure SDKs. The `EventProcessor` inherits the transport from the `ConsumerClient` passed to `build`, so it runs over WebSockets when that client selects them. ([#3601](https://github.com/Azure/azure-sdk-for-rust/issues/3601)) +- Added the `fe2o3_amqp`, `fe2o3_amqp_rustls`, `fe2o3_amqp_ws`, and `fe2o3_amqp_ws_rustls` features, which forward the matching features of `azure_core_amqp`. The `default` feature selects the AMQP backend and the rustls stack with the aws-lc-rs provider for both the TCP and the WebSocket transport. That is the stack that the rest of `sdk/core` uses. See Breaking Changes for the effect on the TCP transport, which ran on native-tls before. To build on another stack, turn off the default features, name the base features, and take a direct dependency on `fe2o3-amqp` and `fe2o3-amqp-ws` with the stack you want; Cargo unifies the features. `samples/list_blobs_native_tls` shows the same pattern for `reqwest`. - The `EventProcessor` now opens every partition receiver with AMQP epoch (owner level) `0` and surfaces broker-initiated displacement as the new `EventHubsError::ConsumerDisconnected` error kind. When a second `EventProcessor` instance claims a partition this instance is currently holding, the broker disconnects this instance's receiver and the consumer's `stream_events()` resolves with `ConsumerDisconnected`. This matches the behavior of `EventProcessorClient` in the .NET and Java Azure SDKs. Consumers should pattern-match on `ErrorKind::ConsumerDisconnected` to detect a stolen partition and re-acquire a client via `next_partition_client()`. - Added `EventHubsError::ConsumerDisconnected(Option)` error variant. - Added the `ErrorKind::InvalidBatchSize { requested, max_allowed }` error variant. `create_batch` reports it when `EventDataBatchOptions::max_size_in_bytes` is zero or is larger than the maximum the sender link allows, so a caller can branch on the kind instead of the message. This matches the `ArgumentOutOfRangeException` that .NET raises and the typed error that Go returns for the same input. ### Breaking Changes +- The `default` feature now selects `fe2o3_amqp_rustls`, so AMQP framed directly on TCP (`amqps://`, port 5671) runs on rustls with the aws-lc-rs provider where it ran on native-tls. The two stacks trust different root certificates. `fe2o3-amqp` builds its rustls root store from `webpki-roots`, which carries the Mozilla root set, and native-tls read the root store of the operating system. A namespace behind a private or an enterprise certificate authority can fail the handshake after this change, even when the operating system trusts that authority. To keep native-tls, turn off the default features, name `fe2o3_amqp`, and take a direct dependency on `fe2o3-amqp` with its `native-tls` feature. ([#4189](https://github.com/Azure/azure-sdk-for-rust/issues/4189)) - On the receive path, the `amqp:link:stolen` AMQP condition is no longer auto-retried. A receiver displaced by a higher-or-equal-epoch attacher now surfaces the error (translated to `EventHubsError::ConsumerDisconnected` by `EventReceiver::stream_events`) instead of silently re-attaching. Sender, CBS, and management operations retain the historical retry-on-stolen behavior. ### Bugs Fixed diff --git a/sdk/eventhubs/azure_messaging_eventhubs/Cargo.toml b/sdk/eventhubs/azure_messaging_eventhubs/Cargo.toml index 0d7ee4773e4..78cc24eca26 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/Cargo.toml +++ b/sdk/eventhubs/azure_messaging_eventhubs/Cargo.toml @@ -20,8 +20,14 @@ edition.workspace = true async-lock.workspace = true async-stream.workspace = true async-trait.workspace = true -azure_core = { workspace = true, default-features = false } -azure_core_amqp.workspace = true +azure_core = { path = "../../core/azure_core", version = "1.2.0-beta.1", default-features = false } +# Unreleased: AmqpTransport for AMQP-over-WebSockets (issue #3601). Path plus +# version per the dependency policy in AGENTS.md. +# `default-features = false` so that the `default` feature below is what selects +# the TLS stack, in the same shape as `azure_core` above. With the default +# features on here, a consumer that turns them off could not drop the rustls +# stack to bring their own. +azure_core_amqp = { path = "../../core/azure_core_amqp", version = "1.2.0-beta.1", default-features = false } base64.workspace = true futures.workspace = true hmac.workspace = true @@ -35,16 +41,22 @@ tracing.workspace = true rustc_version.workspace = true [dev-dependencies] -azure_core_amqp = { workspace = true, features = ["test"] } -azure_core_test = { workspace = true, features = ["tracing"] } -azure_identity.workspace = true -azure_messaging_eventhubs = { path = ".", features = [ +azure_core_amqp = { path = "../../core/azure_core_amqp", default-features = false, features = [ + "test", +] } +azure_core_test = { path = "../../core/azure_core_test", features = [ + "tracing", +] } +azure_identity = { path = "../../identity/azure_identity" } +azure_messaging_eventhubs = { path = ".", default-features = false, features = [ "in_memory_checkpoint_store", ] } # Path-only so `cargo package` works: this crate is a dependency of the # checkpoint store crate, and the migration guide's blob sample compiles here. azure_messaging_eventhubs_checkpointstore_blob = { path = "../azure_messaging_eventhubs_checkpointstore_blob" } -azure_storage_blob.workspace = true +# Path-only, for the same reason as the checkpoint store crate above: the +# migration guide's blob sample must resolve the same azure_core as this crate. +azure_storage_blob = { path = "../../storage/azure_storage_blob" } criterion.workspace = true fe2o3-amqp = { workspace = true, features = ["tracing"] } include-file.workspace = true @@ -54,6 +66,30 @@ tracing-subscriber = { workspace = true, features = ["env-filter", "fmt"] } [features] in_memory_checkpoint_store = [] default = ["azure_core_amqp/default"] +# The AMQP backend and the AMQP-over-WebSockets transport, forwarded from +# `azure_core_amqp`. The `default` feature selects the backend and the rustls +# stack with the aws-lc-rs provider for both the TCP and the WebSocket +# transport. That is the stack that the rest of `sdk/core` uses. +# +# There is no native-tls feature, because `fe2o3-amqp` accepts one TLS stack +# only. To build on another stack, turn off the default features, name the base +# features, and take a direct dependency on the backend crates with the stack +# you want. Cargo unifies the features: +# +# azure_messaging_eventhubs = { version = "...", default-features = false, features = [ +# "fe2o3_amqp", "fe2o3_amqp_ws", +# ] } +# fe2o3-amqp = { version = "0.16", features = ["native-tls"] } +# fe2o3-amqp-ws = { version = "0.16", features = ["native-tls"] } +# +# `samples/list_blobs_native_tls` shows the same pattern for `reqwest`. +fe2o3_amqp = ["azure_core_amqp/fe2o3_amqp"] +fe2o3_amqp_rustls = ["fe2o3_amqp", "azure_core_amqp/fe2o3_amqp_rustls"] +fe2o3_amqp_ws = ["fe2o3_amqp", "azure_core_amqp/fe2o3_amqp_ws"] +fe2o3_amqp_ws_rustls = [ + "fe2o3_amqp_ws", + "azure_core_amqp/fe2o3_amqp_ws_rustls", +] [[bench]] name = "benchmarks" diff --git a/sdk/eventhubs/azure_messaging_eventhubs/examples/eventhubs_websocket_transport.rs b/sdk/eventhubs/azure_messaging_eventhubs/examples/eventhubs_websocket_transport.rs new file mode 100644 index 00000000000..e3b754226bd --- /dev/null +++ b/sdk/eventhubs/azure_messaging_eventhubs/examples/eventhubs_websocket_transport.rs @@ -0,0 +1,98 @@ +// Copyright (c) Microsoft Corporation. +// Licensed under the MIT license. + +//! This sample demonstrates AMQP-over-WebSockets transport. It opens both a +//! [`ProducerClient`] and a [`ConsumerClient`] with +//! [`AmqpTransport::WebSocket`], sends a uniquely tagged event, and reads it +//! back. The clients talk to the broker over `wss://` on port 443 instead of +//! AMQP on port 5671, which is useful when a firewall blocks 5671. +//! +//! Environment: +//! EVENTHUBS_CONNECTION_STRING required, e.g. +//! `Endpoint=sb://.servicebus.windows.net/;SharedAccessKeyName=;SharedAccessKey=` +//! EVENTHUB_NAME required only if the connection string has no `EntityPath` + +use azure_core::{time::Duration, Uuid}; +use azure_messaging_eventhubs::{ + models::AmqpTransport, ConsumerClient, OpenReceiverOptions, ProducerClient, SendEventOptions, + StartLocation, StartPosition, +}; +use futures::StreamExt; + +#[tokio::main] +async fn main() -> Result<(), Box> { + let connection_string = std::env::var("EVENTHUBS_CONNECTION_STRING")?; + // `None` when the connection string already carries an `EntityPath`. + let eventhub_name = std::env::var("EVENTHUB_NAME").ok(); + + let producer = ProducerClient::builder() + .with_transport(AmqpTransport::WebSocket) + .open_with_connection_string(&connection_string, eventhub_name.as_deref()) + .await?; + println!("Opened producer over WebSockets."); + + // Pick a partition and capture its current tail, so we read only events we + // enqueue after this point. + let properties = producer.get_eventhub_properties().await?; + let partition_id = properties.partition_ids[0].clone(); + let before = producer.get_partition_properties(&partition_id).await?; + let start_sequence = before.last_enqueued_sequence_number; + + let marker = Uuid::new_v4().to_string(); + producer + .send_event( + marker.clone(), + Some(SendEventOptions { + partition_id: Some(partition_id.clone()), + }), + ) + .await?; + println!("Sent event with marker {marker} to partition {partition_id}."); + + let consumer = ConsumerClient::builder() + .with_transport(AmqpTransport::WebSocket) + .open_with_connection_string(&connection_string, eventhub_name.as_deref()) + .await?; + println!("Opened consumer over WebSockets."); + + let receiver = consumer + .open_receiver_on_partition( + partition_id.clone(), + Some(OpenReceiverOptions { + start_position: Some(StartPosition { + location: StartLocation::SequenceNumber(start_sequence), + inclusive: false, + }), + receive_timeout: Some(Duration::seconds(30)), + ..Default::default() + }), + ) + .await?; + + let mut found = false; + { + let mut stream = receiver.stream_events(); + while let Some(event) = stream.next().await { + let event = event?; + if event.event_data().body() == Some(marker.as_bytes()) { + found = true; + println!("Received the marker event back over WebSockets."); + break; + } + } + // `stream` borrows `receiver`; drop it (end of block) before closing. + } + + // `consumer.close()` detaches the receiver link too, so the order of these + // two calls is free (#4931). Close the receiver first to show the intent. + receiver.close().await?; + consumer.close().await?; + producer.close().await?; + + if found { + println!("PASS: AMQP-over-WebSockets validated end to end."); + Ok(()) + } else { + Err("FAIL: did not receive the marker event within the timeout".into()) + } +} diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/common/authorizer.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/common/authorizer.rs index ed9e6ad31c4..109fd463ee2 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/common/authorizer.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/common/authorizer.rs @@ -740,6 +740,7 @@ mod tests { url, None, None, + azure_core_amqp::AmqpTransport::default(), mock_credential.clone(), Default::default(), None, @@ -806,6 +807,7 @@ mod tests { url, None, None, + azure_core_amqp::AmqpTransport::default(), mock_credential.clone(), Default::default(), None, @@ -876,6 +878,7 @@ mod tests { host.clone(), None, None, + azure_core_amqp::AmqpTransport::default(), mock_credential.clone(), Default::default(), None, @@ -1007,6 +1010,7 @@ mod tests { url.clone(), None, None, + Default::default(), credential.clone(), Default::default(), None, @@ -1103,6 +1107,7 @@ mod tests { url.clone(), None, None, + azure_core_amqp::AmqpTransport::default(), credential.clone(), Default::default(), None, @@ -1208,6 +1213,7 @@ mod tests { url.clone(), None, None, + azure_core_amqp::AmqpTransport::default(), credential.clone(), Default::default(), None, @@ -1340,6 +1346,7 @@ mod tests { url.clone(), None, None, + azure_core_amqp::AmqpTransport::default(), credential.clone(), Default::default(), None, diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/common/recoverable/connection.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/common/recoverable/connection.rs index 678516feeb9..64760d7efa2 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/common/recoverable/connection.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/common/recoverable/connection.rs @@ -25,7 +25,7 @@ use azure_core_amqp::{ AmqpClaimsBasedSecurity, AmqpConnection, AmqpConnectionApis, AmqpConnectionOptions, AmqpError, AmqpManagement, AmqpManagementApis, AmqpReceiver, AmqpReceiverApis, AmqpReceiverOptions, AmqpSender, AmqpSenderApis, AmqpSession, AmqpSessionApis, AmqpSessionOptions, AmqpSource, - AmqpSymbol, + AmqpSymbol, AmqpTransport, }; #[cfg(test)] use std::sync::Mutex; @@ -79,6 +79,7 @@ pub(crate) struct RecoverableConnection { pub(super) url: Url, application_id: Option, custom_endpoint: Option, + transport: AmqpTransport, // The management client is a single cached instance, held in a `OnceCell` // for the same reason the per-path caches are: the expensive build (connect // + session begin + CBS authorize + link attach) must not run while a lock @@ -283,6 +284,7 @@ impl RecoverableConnection { url: Url, application_id: Option, custom_endpoint: Option, + transport: AmqpTransport, credential: Arc, retry_options: RetryOptions, cbs_token_type: Option<&'static str>, @@ -299,6 +301,7 @@ impl RecoverableConnection { application_id, connection_name, custom_endpoint, + transport, retry_options, cbs_lock: AsyncMutex::new(()), connections: AsyncMutex::new(None), @@ -808,6 +811,28 @@ impl RecoverableConnection { .cell } + /// Builds the options handed to [`AmqpConnection::open`]. Kept separate from + /// `create_connection` so the wiring can be asserted without a broker. + fn connection_options(&self) -> AmqpConnectionOptions { + AmqpConnectionOptions { + properties: Some( + vec![ + ("user-agent", get_user_agent(&self.application_id)), + ("version", get_package_version()), + ("platform", get_platform_info()), + ("product", get_package_name()), + ] + .into_iter() + .map(|(k, v)| (AmqpSymbol::from(k), AmqpValue::from(v))) + .collect(), + ), + desired_capabilities: Some(vec![GEODR_REPLICATION_CAPABILITY.into()]), + custom_endpoint: self.custom_endpoint.clone(), + transport: Some(self.transport), + ..Default::default() + } + } + #[instrument( level = "debug", skip_all, @@ -829,22 +854,7 @@ impl RecoverableConnection { .open( self.connection_name.clone(), self.url.clone(), - Some(AmqpConnectionOptions { - properties: Some( - vec![ - ("user-agent", get_user_agent(&self.application_id)), - ("version", get_package_version()), - ("platform", get_platform_info()), - ("product", get_package_name()), - ] - .into_iter() - .map(|(k, v)| (AmqpSymbol::from(k), AmqpValue::from(v))) - .collect(), - ), - desired_capabilities: Some(vec![GEODR_REPLICATION_CAPABILITY.into()]), - custom_endpoint: self.custom_endpoint.clone(), - ..Default::default() - }), + Some(self.connection_options()), ) .await?; info!( @@ -1515,6 +1525,7 @@ mod tests { Url::parse("amqps://example.com").unwrap(), None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1550,6 +1561,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1573,6 +1585,7 @@ mod tests { url, Some(app_id.clone()), None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1593,6 +1606,7 @@ mod tests { url.clone(), None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1613,6 +1627,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1667,6 +1682,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1696,6 +1712,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1740,6 +1757,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1785,6 +1803,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1840,6 +1859,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1943,6 +1963,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -1992,6 +2013,7 @@ mod tests { url, None, Some(custom_endpoint.clone()), + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -2000,6 +2022,51 @@ mod tests { assert_eq!(connection_manager.custom_endpoint, Some(custom_endpoint)); } + // The transport selected on a client builder (and, transitively, on an + // EventProcessor's injected ConsumerClient) must reach the connection so it + // is applied when the AMQP connection is opened. This verifies the field is + // stored on the RecoverableConnection. + #[test] + fn constructor_with_websocket_transport() { + let url = Url::parse("amqps://example.com").unwrap(); + let connection_manager = RecoverableConnection::new( + url, + None, + None, + AmqpTransport::WebSocket, + Arc::new(MockCredential), + Default::default(), + None, + ); + + assert_eq!(connection_manager.transport, AmqpTransport::WebSocket); + } + + // The stored transport must also reach the options handed to + // `AmqpConnection::open`. Asserting on the constructor alone would still + // pass if `create_connection` dropped the `with_transport` call. + #[test] + fn connection_options_carry_the_transport() { + let url = Url::parse("amqps://example.com").unwrap(); + let custom_endpoint = Url::parse("amqps://proxy.example.com:8081").unwrap(); + for transport in [AmqpTransport::Tcp, AmqpTransport::WebSocket] { + let connection_manager = RecoverableConnection::new( + url.clone(), + None, + Some(custom_endpoint.clone()), + transport, + Arc::new(MockCredential), + Default::default(), + None, + ); + + let options = connection_manager.connection_options(); + assert_eq!(options.transport, Some(transport)); + assert_eq!(options.custom_endpoint, Some(custom_endpoint.clone())); + assert!(options.properties.is_some()); + } + } + #[test] fn test_should_retry_amqp_error() { use azure_core_amqp::AmqpDescribedError; @@ -2363,6 +2430,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -2439,6 +2507,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, @@ -2488,6 +2557,7 @@ mod tests { url, None, None, + AmqpTransport::default(), Arc::new(MockCredential), Default::default(), None, diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/event_receiver.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/event_receiver.rs index c933e59ec51..1d7a3664df2 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/event_receiver.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/event_receiver.rs @@ -385,6 +385,7 @@ mod tests { Url::parse("amqps://example.servicebus.windows.net").unwrap(), None, None, + Default::default(), Arc::new(azure_core_test::credentials::MockCredential), Default::default(), None, diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/mod.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/mod.rs index 2e89e763b42..205e7092e90 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/mod.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/consumer/mod.rs @@ -16,7 +16,7 @@ use azure_core::{credentials::TokenCredential, http::Url, time::Duration, Uuid}; use azure_core_amqp::AmqpError; use azure_core_amqp::{ message::AmqpSourceFilter, AmqpDescribed, AmqpOrderedMap, AmqpReceiverOptions, AmqpSource, - AmqpSymbol, AmqpValue, ReceiverCreditMode, + AmqpSymbol, AmqpTransport, AmqpValue, ReceiverCreditMode, }; pub use event_receiver::EventReceiver; use std::{ @@ -45,6 +45,7 @@ struct ConsumerClientOptions { retry_options: Option, custom_endpoint: Option, cbs_token_type: Option<&'static str>, + transport: AmqpTransport, } impl ConsumerClient { @@ -98,6 +99,7 @@ impl ConsumerClient { url.clone(), options.application_id, options.custom_endpoint, + options.transport, credential, retry_options, options.cbs_token_type, @@ -195,6 +197,7 @@ impl ConsumerClient { retry_options: None, custom_endpoint: None, cbs_token_type: None, + transport: AmqpTransport::default(), }, ) } @@ -610,6 +613,7 @@ pub mod builders { }, Result, }; + use azure_core_amqp::AmqpTransport; use std::sync::Arc; /// A builder for creating a [`ConsumerClient`]. @@ -637,6 +641,7 @@ pub mod builders { instance_id: Option, retry_options: Option, custom_endpoint: Option, + transport: Option, } impl ConsumerClientBuilder { @@ -704,11 +709,37 @@ pub mod builders { /// Note: The custom endpoint option allows a customer to specify an AMQP proxy /// which will be used to forward requests to the actual Event Hub instance. /// + /// An explicit port on the endpoint carries into the address that the client + /// dials. Under [`AmqpTransport::WebSocket`] that is the `wss://` address, so + /// name the port that the proxy accepts WebSockets on, and leave the port out + /// to dial the default port 443. + /// pub fn with_custom_endpoint(mut self, endpoint: String) -> Self { self.custom_endpoint = Some(endpoint); self } + /// Sets the transport used to communicate with the Event Hub. + /// + /// # Arguments + /// * `transport` - The transport to use. Defaults to + /// [`AmqpTransport::Tcp`]. Use [`AmqpTransport::WebSocket`] to + /// tunnel AMQP over WebSockets (port 443) when the native AMQP + /// ports are blocked. + /// + /// # Returns + /// The updated [`ConsumerClientBuilder`]. + pub fn with_transport(mut self, transport: AmqpTransport) -> Self { + self.transport = Some(transport); + self + } + + /// Returns the AMQP transport this builder opens the connection with. + /// Shared by every `open` path so they cannot drift apart. + pub(crate) fn transport(&self) -> AmqpTransport { + self.transport.unwrap_or_default() + } + /// Opens a connection to the Event Hub. /// /// This method establishes a connection to the Event Hubs instance associated @@ -750,6 +781,7 @@ pub mod builders { eventhub_name: String, credential: Arc, ) -> Result { + let transport = self.transport(); let custom_endpoint = match self.custom_endpoint { Some(endpoint) => Some(Url::parse(&endpoint).map_err(azure_core::Error::from)?), None => None, @@ -766,6 +798,7 @@ pub mod builders { retry_options: self.retry_options, custom_endpoint, cbs_token_type: None, + transport, }, )?; consumer.ensure_connection().await?; @@ -814,6 +847,7 @@ pub mod builders { connection_string: &str, eventhub: Option<&str>, ) -> Result { + let transport = self.transport(); let connection_string: ConnectionString = connection_string.parse()?; let eventhub = resolve_eventhub(&connection_string, eventhub)?; let credential = Arc::new(SasCredential::from_connection_string( @@ -837,6 +871,7 @@ pub mod builders { retry_options: self.retry_options, custom_endpoint, cbs_token_type: Some(SAS_TOKEN_TYPE), + transport, }, )?; consumer.ensure_connection().await?; @@ -852,13 +887,33 @@ pub(crate) mod tests { ProducerClient, Result, StartLocation, StartPosition, }; use azure_core::{sleep::sleep, time::Duration}; - use azure_core_amqp::{error::AmqpErrorKind, AmqpError}; + use azure_core_amqp::{error::AmqpErrorKind, AmqpError, AmqpTransport}; use azure_core_test::{recorded, TestContext}; use futures::stream::StreamExt; use std::{ sync::Arc, time::{SystemTime, UNIX_EPOCH}, }; + + // Every `open` path on the builder reads the transport through one helper, + // so this covers the plumbing that the connection-string path shares. + #[test] + fn builder_reads_the_transport_through_one_helper() { + assert_eq!( + ConsumerClient::builder() + .with_transport(AmqpTransport::WebSocket) + .transport(), + AmqpTransport::WebSocket + ); + assert_eq!( + ConsumerClient::builder() + .with_transport(AmqpTransport::Tcp) + .transport(), + AmqpTransport::Tcp + ); + // An unset transport keeps the TCP default. + assert_eq!(ConsumerClient::builder().transport(), AmqpTransport::Tcp); + } use tracing::info; // static INIT_LOGGING: std::sync::Once = std::sync::Once::new(); diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/event_processor/processor.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/event_processor/processor.rs index 6e0a6ba00da..97102d90971 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/event_processor/processor.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/event_processor/processor.rs @@ -760,6 +760,45 @@ pub mod builders { /// Builds the event processor with the specified consumer client and checkpoint store. /// Returns a `Result` containing the constructed `EventProcessor`. + /// + /// # Connection options (including transport) + /// + /// The event processor does not open its own connection. It processes + /// partitions using the [`ConsumerClient`] passed here, and every + /// per-partition receiver reuses that client's connection. Connection-level + /// options, such as the transport, a custom endpoint, retry options, and the + /// application id, are therefore configured on the [`ConsumerClient`] before + /// it is passed to `build`. + /// + /// To run the processor over AMQP-over-WebSockets (port 443, useful when the + /// native AMQP ports are blocked), select the transport on the consumer + /// client with + /// [`ConsumerClientBuilder::with_transport`](crate::builders::ConsumerClientBuilder::with_transport): + /// + /// ```no_run + /// use azure_messaging_eventhubs::{EventProcessor, CheckpointStore, ConsumerClient}; + /// use azure_messaging_eventhubs::models::AmqpTransport; + /// use std::sync::Arc; + /// + /// async fn create_processor(checkpoint_store: Arc) -> Result<(), Box> { + /// use azure_identity::DeveloperToolsCredential; + /// + /// let eventhub_namespace = std::env::var("EVENTHUBS_HOST")?; + /// let eventhub_name = std::env::var("EVENTHUB_NAME")?; + /// let consumer = ConsumerClient::builder() + /// .with_transport(AmqpTransport::WebSocket) + /// .open( + /// &eventhub_namespace, + /// eventhub_name, + /// DeveloperToolsCredential::new(None)?.clone(), + /// ) + /// .await?; + /// let processor = EventProcessor::builder() + /// .build(consumer, checkpoint_store.clone()) + /// .await?; + /// Ok(()) + /// } + /// ``` pub async fn build( self, consumer_client: ConsumerClient, diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/models/mod.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/models/mod.rs index 92b28ad90c4..319846fb0b1 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/models/mod.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/models/mod.rs @@ -10,6 +10,32 @@ pub use azure_core_amqp::AmqpMessage; /// An AMQP Value. pub use azure_core_amqp::AmqpValue; +/// The transport that carries the AMQP protocol. +/// +/// Event Hubs is normally reached with AMQP framed directly over a TCP/TLS +/// socket ([`AmqpTransport::Tcp`], port 5671). Some networks (for example +/// corporate firewalls) only permit outbound connections on port 443; there, +/// [`AmqpTransport::WebSocket`] tunnels AMQP over secure WebSockets instead. +/// Select it with `with_transport` on a client builder. +/// +/// # Examples +/// +/// ```no_run +/// use azure_messaging_eventhubs::{ProducerClient, models::AmqpTransport}; +/// use azure_identity::DeveloperToolsCredential; +/// +/// #[tokio::main] +/// async fn main() -> Result<(), Box> { +/// let credential = DeveloperToolsCredential::new(None)?; +/// let producer = ProducerClient::builder() +/// .with_transport(AmqpTransport::WebSocket) +/// .open("my_namespace.servicebus.windows.net", "my_eventhub", credential) +/// .await?; +/// Ok(()) +/// } +/// ``` +pub use azure_core_amqp::AmqpTransport; + /// An AMQP Simple Value. /// /// An AMQP Simple Value is a primitive type in AMQP 1.0. diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/producer/batch.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/producer/batch.rs index a664749abfe..2d7bcdc4491 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/producer/batch.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/producer/batch.rs @@ -392,6 +392,7 @@ pub struct EventDataBatchOptions { mod tests { use super::*; use crate::RetryOptions; + use azure_core_amqp::AmqpTransport; use azure_core_test::credentials::MockCredential; use std::sync::Arc; @@ -428,6 +429,7 @@ mod tests { RetryOptions::default(), None, None, + AmqpTransport::default(), ) } diff --git a/sdk/eventhubs/azure_messaging_eventhubs/src/producer/mod.rs b/sdk/eventhubs/azure_messaging_eventhubs/src/producer/mod.rs index 7e8e14c8b9f..81f9f2a188d 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs/src/producer/mod.rs +++ b/sdk/eventhubs/azure_messaging_eventhubs/src/producer/mod.rs @@ -17,6 +17,7 @@ use azure_core::{ }; use azure_core_amqp::{ error::AmqpErrorKind, AmqpError, AmqpSendOptions, AmqpSendOutcome, AmqpSenderApis, + AmqpTransport, }; use batch::{EventDataBatch, EventDataBatchOptions}; use std::{fmt::Debug, sync::Arc}; @@ -92,6 +93,7 @@ impl From for SendMessageOptions { } impl ProducerClient { + #[allow(clippy::too_many_arguments, reason = "private API")] pub(crate) fn new( endpoint: Url, eventhub: String, @@ -100,12 +102,14 @@ impl ProducerClient { retry_options: RetryOptions, custom_endpoint: Option, cbs_token_type: Option<&'static str>, + transport: AmqpTransport, ) -> Self { Self { connection: RecoverableConnection::new( endpoint.clone(), application_id, custom_endpoint, + transport, credential, retry_options, cbs_token_type, @@ -561,6 +565,7 @@ pub mod builders { Result, RetryOptions, }; use azure_core::{http::Url, Error}; + use azure_core_amqp::AmqpTransport; use std::sync::Arc; /// A builder for creating a [`ProducerClient`]. @@ -590,6 +595,9 @@ pub mod builders { /// The custom endpoint for the Event Hub. custom_endpoint: Option, + + /// The transport used to communicate with the Event Hub. + transport: Option, } impl ProducerClientBuilder { @@ -640,11 +648,37 @@ pub mod builders { /// Note: The custom endpoint option allows a customer to specify an AMQP proxy /// which will be used to forward requests to the actual Event Hub instance. /// + /// An explicit port on the endpoint carries into the address that the client + /// dials. Under [`AmqpTransport::WebSocket`] that is the `wss://` address, so + /// name the port that the proxy accepts WebSockets on, and leave the port out + /// to dial the default port 443. + /// pub fn with_custom_endpoint(mut self, endpoint: String) -> Self { self.custom_endpoint = Some(endpoint); self } + /// Sets the transport used to communicate with the Event Hub. + /// + /// # Arguments + /// * `transport` - The transport to use. Defaults to + /// [`AmqpTransport::Tcp`]. Use [`AmqpTransport::WebSocket`] to + /// tunnel AMQP over WebSockets (port 443) when the native AMQP + /// ports are blocked. + /// + /// # Returns + /// The updated [`ProducerClientBuilder`]. + pub fn with_transport(mut self, transport: AmqpTransport) -> Self { + self.transport = Some(transport); + self + } + + /// Returns the AMQP transport this builder opens the connection with. + /// Shared by every `open` path so they cannot drift apart. + pub(crate) fn transport(&self) -> AmqpTransport { + self.transport.unwrap_or_default() + } + /// Opens the connection to the Event Hub. /// /// # Arguments @@ -661,6 +695,7 @@ pub mod builders { eventhub: &str, credential: Arc, ) -> Result { + let transport = self.transport(); let url = format!("amqps://{}/{}", fully_qualified_namespace, eventhub); let url = Url::parse(&url).map_err(azure_core::Error::from)?; @@ -677,6 +712,7 @@ pub mod builders { self.retry_options.unwrap_or_default(), custom_endpoint, None, + transport, ); // Open a connection to the Event Hub to ensure that the client is ready to send messages. @@ -726,6 +762,7 @@ pub mod builders { connection_string: &str, eventhub: Option<&str>, ) -> Result { + let transport = self.transport(); let connection_string: ConnectionString = connection_string.parse()?; let eventhub = resolve_eventhub(&connection_string, eventhub)?; let credential = Arc::new(SasCredential::from_connection_string( @@ -752,6 +789,7 @@ pub mod builders { self.retry_options.unwrap_or_default(), custom_endpoint, Some(SAS_TOKEN_TYPE), + transport, ); client.ensure_connection().await?; @@ -765,10 +803,30 @@ mod tests { use crate::common::tests::force_errors; use crate::{models::EventData, EventDataBatchOptions, ProducerClient, Result}; use azure_core::time::Duration; - use azure_core_amqp::error::AmqpErrorKind; + use azure_core_amqp::{error::AmqpErrorKind, AmqpTransport}; use azure_core_test::{recorded, TestContext}; use std::sync::Arc; + // Every `open` path on the builder reads the transport through one helper, + // so this covers the plumbing that the connection-string path shares. + #[test] + fn builder_reads_the_transport_through_one_helper() { + assert_eq!( + ProducerClient::builder() + .with_transport(AmqpTransport::WebSocket) + .transport(), + AmqpTransport::WebSocket + ); + assert_eq!( + ProducerClient::builder() + .with_transport(AmqpTransport::Tcp) + .transport(), + AmqpTransport::Tcp + ); + // An unset transport keeps the TCP default. + assert_eq!(ProducerClient::builder().transport(), AmqpTransport::Tcp); + } + #[recorded::test(live)] async fn force_errors_send_batch_link_error(ctx: TestContext) -> Result<()> { const EVENTHUB_PARTITION: &str = "1"; diff --git a/sdk/eventhubs/azure_messaging_eventhubs_checkpointstore_blob/Cargo.toml b/sdk/eventhubs/azure_messaging_eventhubs_checkpointstore_blob/Cargo.toml index 881adf4414e..606c407ebef 100644 --- a/sdk/eventhubs/azure_messaging_eventhubs_checkpointstore_blob/Cargo.toml +++ b/sdk/eventhubs/azure_messaging_eventhubs_checkpointstore_blob/Cargo.toml @@ -20,9 +20,12 @@ rust-version.workspace = true [dependencies] async-trait.workspace = true -azure_core.workspace = true +azure_core = { path = "../../core/azure_core", version = "1.2.0-beta.1", default-features = false } azure_messaging_eventhubs = { path = "../azure_messaging_eventhubs", version = "0.15.0" } -azure_storage_blob.workspace = true +# Unreleased: must resolve the same azure_core as azure_messaging_eventhubs, which +# needs AmqpTransport (issue #3601). The registry build of azure_storage_blob pins +# the previous azure_core, which puts two copies of the crate in the graph. +azure_storage_blob = { path = "../../storage/azure_storage_blob", version = "1.1.0-beta.2" } futures.workspace = true serde = { workspace = true, features = ["derive"] } serde_json.workspace = true @@ -30,9 +33,11 @@ time = { workspace = true, features = ["serde"] } tracing = { workspace = true } [dev-dependencies] -azure_core_opentelemetry.workspace = true -azure_core_test = { workspace = true, features = ["tracing"] } -azure_identity.workspace = true +azure_core_opentelemetry = { path = "../../core/azure_core_opentelemetry" } +azure_core_test = { path = "../../core/azure_core_test", features = [ + "tracing", +] } +azure_identity = { path = "../../identity/azure_identity" } opentelemetry.workspace = true opentelemetry-appender-tracing.workspace = true opentelemetry-stdout.workspace = true