Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions sdk/eventhubs/azure_messaging_eventhubs/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@
- Fixed a deadlock when a CBS failure during management-client creation started connection recovery. ([#4728](https://github.com/Azure/azure-sdk-for-rust/issues/4728))
- Closed a stale-resource window in connection recovery. A `ReconnectConnection` recovery that fired while a slow-path attach (authorize, session begin, or sender/receiver link attach) was in flight could cache a resource bound to the just-dropped connection; the next operation on that resource failed (unauthorized / detached / closed) and triggered a second, redundant recovery cycle. A recovery generation counter now tags each cached resource, and a slow path that completes across a recovery discards its result and re-attaches against the new connection instead of caching the stale one. The authorizer's token cache is mutable (a background task refreshes tokens) so it cannot use the same one-shot cell as the connection caches; both of its writers, `authorize_path` and the refresh task, instead re-check the generation under the same lock that recovery's clear takes, and a recovery brackets its invalidation with a generation bump on each side, which leaves the counter odd for as long as the recovery runs, so a slow path that overlaps a recovery at either end also discards rather than caching a resource bound to the connection that recovery is dropping. A token refresh pass that a recovery discards now applies the same backoff floor as a failed pass, so a recovery storm cannot turn the refresh loop into an uncapped stream of credential and CBS calls. The per-path / per-partition concurrency is preserved. ([#4454](https://github.com/Azure/azure-sdk-for-rust/issues/4454))
- `InMemoryCheckpointStore` now rotates the ETag and refreshes `last_modified_time` when an existing ownership is renewed, matching the create path and the production `BlobCheckpointStore`. Previously the renewal path reinserted the caller's record verbatim, leaving a stale ETag and timestamp; that divergence from the real store could mask bugs in code that relies on ETag rotation for optimistic concurrency. ([#4594](https://github.com/Azure/azure-sdk-for-rust/issues/4594))
- `ConsumerClient::open_receiver_on_partition` no longer logs `Receiver attached on partition.` before an attach happens. The call creates the receiver and does no network I/O. The AMQP link attaches on the first poll of `EventReceiver::stream_events()`, which is where the service reports an unknown consumer group or an unknown partition id. The documentation on `open_receiver_on_partition`, `EventProcessor::run`, and `PartitionClient::stream_events` now states this. ([#5094](https://github.com/Azure/azure-sdk-for-rust/issues/5094))

### Other Changes

Expand Down
157 changes: 149 additions & 8 deletions sdk/eventhubs/azure_messaging_eventhubs/src/consumer/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -228,9 +228,15 @@ impl ConsumerClient {
})
}

/// Attaches a message receiver to a specific partition of the Event Hub.
/// Creates a message receiver for a specific partition of the Event Hub.
///
/// This function establishes a connection to the specified partition of the Event Hubs instance and returns a MessageReceiver which can be used to receive messages from it.
/// This function opens no AMQP link and does no network I/O. It builds an
/// [`EventReceiver`] from `options` and returns it. The AMQP link attaches on
/// the first poll of [`EventReceiver::stream_events`].
///
/// An unknown consumer group and an unknown partition id are reported from that
/// first poll, not from this call. A caller must poll the stream before it can
/// trust the receiver.
///
/// # Arguments
///
Expand All @@ -239,7 +245,7 @@ impl ConsumerClient {
///
/// # Returns
///
/// A MessageReceiver which can be used to receive messages from the partition.
/// An [`EventReceiver`] which can be used to receive messages from the partition.
///
/// Note that by default, a message receiver will receive events starting from the latest event in the partition (in
/// other words, it will receive new events only). To receive events from another location within the partition you can
Expand Down Expand Up @@ -344,7 +350,7 @@ impl ConsumerClient {
consumer_group = %self.consumer_group,
eventhub = %self.eventhub,
source_url = %source_url,
"Receiver attached on partition."
"Created receiver on partition. The AMQP link attaches on the first stream_events() poll."
);
Ok(EventReceiver::new(
self.recoverable_connection.clone(),
Expand Down Expand Up @@ -883,15 +889,15 @@ pub mod builders {
#[cfg(test)]
pub(crate) mod tests {
use crate::{
common::tests::force_errors, models::EventData, ConsumerClient, EventDataBatchOptions,
ProducerClient, Result, StartLocation, StartPosition,
common::tests::force_errors, error::ErrorKind, models::EventData, ConsumerClient,
EventDataBatchOptions, ProducerClient, Result, StartLocation, StartPosition,
};
use azure_core::{sleep::sleep, time::Duration};
use azure_core_amqp::{error::AmqpErrorKind, AmqpError, AmqpTransport};
use azure_core_test::{recorded, TestContext};
use azure_core_test::{credentials::MockCredential, recorded, TestContext};
use futures::stream::StreamExt;
use std::{
sync::Arc,
sync::{Arc, Mutex},
time::{SystemTime, UNIX_EPOCH},
};

Expand Down Expand Up @@ -1337,4 +1343,139 @@ pub(crate) mod tests {
})
.await
}
/// Collects the formatted output of a `tracing` subscriber, so a test can
/// read the records that the code under test made.
#[derive(Clone, Default)]
struct LogBuffer(Arc<Mutex<Vec<u8>>>);

impl LogBuffer {
fn contents(&self) -> String {
let buffer = self.0.lock().expect("the log buffer lock is poisoned");
String::from_utf8_lossy(&buffer).to_string()
}
}

impl std::io::Write for LogBuffer {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
let mut buffer = self.0.lock().expect("the log buffer lock is poisoned");
buffer.extend_from_slice(buf);
Ok(buf.len())
}

fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}

impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for LogBuffer {
type Writer = LogBuffer;

fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}

/// Installs a log capture and returns it with its guard. The dispatcher is
/// thread local, so the test must stay on one thread. Drop the guard
/// before you read `contents()`.
fn capture_logs() -> (LogBuffer, tracing::subscriber::DefaultGuard) {
let buffer = LogBuffer::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(buffer.clone())
.with_max_level(tracing::Level::TRACE)
.with_ansi(false)
.finish();
let guard = tracing::subscriber::set_default(subscriber);
(buffer, guard)
}

fn unconnected_consumer() -> ConsumerClient {
ConsumerClient::new_unconnected(
"example.servicebus.windows.net",
"test-eventhub",
Arc::new(MockCredential),
)
.expect("the client must build")
}

fn consumer_with_armed_attach_error() -> ConsumerClient {
let consumer = unconnected_consumer();
consumer
.recoverable_connection()
.force_attach_error(AmqpError::with_message("attach failed"))
.expect("the attach error must arm");
consumer
}

#[test]
fn open_receiver_on_partition_logs_a_deferred_attach() {
let consumer = unconnected_consumer();
let (buffer, guard) = capture_logs();
futures::executor::block_on(consumer.open_receiver_on_partition("0".to_string(), None))
.expect("the open must succeed without a connection");
drop(guard);
let logs = buffer.contents();

assert!(
logs.contains(
"Created receiver on partition. The AMQP link attaches on the first stream_events() poll."
),
"the log must say the link attaches on the first poll, got: {logs}"
);
assert!(
!logs.contains("Receiver attached on partition."),
"the log must not claim an attach happened, got: {logs}"
);
}

// Pins the lazy contract. An eager
// `recoverable_connection.get_receiver(...)` put before the
// `Ok(EventReceiver::new(...))` return makes the armed attach error leave
// the open, and this test fails.
#[tokio::test]
async fn open_receiver_on_partition_defers_the_attach_to_the_first_poll() {
let consumer = consumer_with_armed_attach_error();
let receiver = consumer
.open_receiver_on_partition("0".to_string(), None)
.await
.expect("the open must not attach, so the armed attach error must not surface here");

let mut stream = std::pin::pin!(receiver.stream_events());
let error = stream
.next()
.await
.expect("the stream yields the armed attach failure")
.expect_err("the armed attach error must surface on the first poll");
assert!(
matches!(error.kind, ErrorKind::AmqpError(_)),
"expected AmqpError, got {:?}",
error.kind
);
}

#[tokio::test]
async fn open_receiver_on_partition_never_logs_an_attach_that_failed() {
let consumer = consumer_with_armed_attach_error();
let (buffer, guard) = capture_logs();
let receiver = consumer
.open_receiver_on_partition("0".to_string(), None)
.await
.expect("the open must succeed without a connection");

let mut stream = std::pin::pin!(receiver.stream_events());
let _ = stream.next().await;
drop(guard);
let logs = buffer.contents();

// The first assertion anchors the capture. An empty capture would make
// the second assertion pass for the wrong reason.
assert!(
logs.contains("Opening receiver on partition."),
"the capture must hold the events of the call, got: {logs}"
);
assert!(
!logs.contains("Attached receiver on partition."),
"no attach was made, so nothing may record one, got: {logs}"
);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,10 @@ impl PartitionClient {
/// This method returns a stream of `ReceivedEventData` wrapped in a `Result`.
/// The stream yields events as they are received from the partition.
///
/// The partition's AMQP link attaches on the first poll of this stream. An unknown
/// consumer group and an unknown partition id are reported here, not from
/// [`EventProcessor::run`](crate::EventProcessor::run).
///
/// # Returns
/// A stream of `Result<ReceivedEventData>` representing the received events.
pub fn stream_events(&self) -> impl Stream<Item = Result<ReceivedEventData>> + '_ {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -223,6 +223,11 @@ impl EventProcessor {
/// to manage the ownership of partitions and distribute the load
/// among consumers.
/// The event processor will run until it is stopped or interrupted.
///
/// Each partition receiver attaches its AMQP link on the first poll of
/// [`PartitionClient::stream_events`](crate::processor::PartitionClient::stream_events).
/// This method does not report an invalid consumer group. The first poll of the
/// partition stream reports it.
/// # Errors
/// Returns an error if the event processor fails to start.
/// # Examples
Expand Down Expand Up @@ -430,13 +435,14 @@ impl EventProcessor {
));
}

// Since we can only have a single EventReceiver on a partition, we don't actually attempt to create the receiver until
let start_position = self.get_start_position(&partition_id, checkpoints);
debug!(
partition_id = %partition_id,
start_position = ?start_position,
"Start position for partition."
);
// The AMQP link for this partition attaches on the first poll of
// `stream_events()`, so this call only builds the receiver.
let receiver = self
.consumer_client
.open_receiver_on_partition(
Expand Down