pub struct Subscriber { /* private fields */ }Expand description
A subscriber. Multiple subscribers compete for messages; delivery is not broadcast.
Implementations§
Source§impl Subscriber
impl Subscriber
Sourcepub fn open(options: &Options) -> Result<Self>
pub fn open(options: &Options) -> Result<Self>
Creates or joins a queue with the supplied identity and capacity.
§Errors
Rejects invalid options, capacity mismatches, exhausted registrations, publisher limits (publishers only), and operating system failures.
Examples found in repository?
examples/benchmark.rs (line 10)
5fn main() {
6 const ITERATIONS: usize = 1_000_000;
7 for size in [3, 50, 1024] {
8 let options = Options::new(format!("b{}x{size}", std::process::id()), 1 << 20);
9 let publisher = Publisher::open(&options).unwrap();
10 let subscriber = Subscriber::open(&options).unwrap();
11 let message = vec![42; size];
12 let mut received = vec![0; size];
13 let mut samples = Vec::new();
14 for round in 0..12 {
15 let start = Instant::now();
16 for _ in 0..ITERATIONS {
17 publisher.try_send(black_box(&message)).unwrap();
18 assert_eq!(
19 subscriber.try_recv_into(black_box(&mut received)).unwrap(),
20 Some(size)
21 );
22 black_box(&received);
23 }
24 if round >= 4 {
25 samples.push(start.elapsed().as_nanos() as f64 / ITERATIONS as f64);
26 }
27 }
28 let mean = samples.iter().sum::<f64>() / samples.len() as f64;
29 let deviation = (samples.iter().map(|x| (x - mean).powi(2)).sum::<f64>()
30 / (samples.len() - 1) as f64)
31 .sqrt();
32 println!(
33 "{size} bytes: {mean:.2} ns/roundtrip, stddev {deviation:.2}, samples {samples:?}"
34 );
35 }
36}More examples
examples/interop.rs (line 30)
14fn main() {
15 let args = std::env::args().collect::<Vec<_>>();
16 let options = Options::new(
17 &args[2],
18 std::env::var("INTEROP_CAPACITY")
19 .map(|v| v.parse().unwrap())
20 .unwrap_or(4096),
21 )
22 .with_path(&args[3]);
23 let count = args[4].parse::<usize>().unwrap();
24 if args[1].starts_with("hold-") {
25 let _publisher;
26 let _subscriber;
27 if args[1] == "hold-publisher" {
28 _publisher = Publisher::open(&options).unwrap();
29 } else {
30 _subscriber = Subscriber::open(&options).unwrap();
31 }
32 println!("READY");
33 io::stdout().flush().unwrap();
34 io::stdin().read_line(&mut String::new()).unwrap();
35 } else if args[1] == "publish" {
36 let publisher = Publisher::open(&options).unwrap();
37 let start = args.get(5).map_or(0, |s| s.parse::<usize>().unwrap());
38 if args.len() > 5 {
39 println!("READY");
40 io::stdout().flush().unwrap();
41 io::stdin().read_line(&mut String::new()).unwrap();
42 }
43 for i in start..start + count {
44 let data = message(i);
45 while let Err(error) = publisher.try_send(&data) {
46 assert!(error.is_full(), "{error}");
47 std::thread::yield_now();
48 }
49 }
50 } else {
51 let subscriber = Subscriber::open(&options).unwrap();
52 println!("READY");
53 if args[1] == "collect" {
54 loop {
55 let data = subscriber
56 .recv_timeout(Duration::from_secs(30))
57 .unwrap()
58 .unwrap();
59 if data.is_empty() {
60 break;
61 }
62 let id = u64::from_le_bytes(data[..8].try_into().unwrap()) as usize;
63 assert_eq!(data, message(id));
64 println!("{id}");
65 }
66 return;
67 }
68 for i in 0..count {
69 assert_eq!(
70 subscriber
71 .recv_timeout(Duration::from_secs(30))
72 .unwrap()
73 .unwrap(),
74 message(i)
75 );
76 }
77 }
78}Sourcepub fn try_recv(&self) -> Result<Option<Vec<u8>>>
pub fn try_recv(&self) -> Result<Option<Vec<u8>>>
Copies and consumes a ready message, allocating a result vector.
Sourcepub fn try_recv_into(&self, buffer: &mut [u8]) -> Result<Option<usize>>
pub fn try_recv_into(&self, buffer: &mut [u8]) -> Result<Option<usize>>
Copies into caller-owned storage. An undersized buffer truncates and consumes the message, matching the .NET v3 API. The return value is bytes copied.
Examples found in repository?
examples/benchmark.rs (line 19)
5fn main() {
6 const ITERATIONS: usize = 1_000_000;
7 for size in [3, 50, 1024] {
8 let options = Options::new(format!("b{}x{size}", std::process::id()), 1 << 20);
9 let publisher = Publisher::open(&options).unwrap();
10 let subscriber = Subscriber::open(&options).unwrap();
11 let message = vec![42; size];
12 let mut received = vec![0; size];
13 let mut samples = Vec::new();
14 for round in 0..12 {
15 let start = Instant::now();
16 for _ in 0..ITERATIONS {
17 publisher.try_send(black_box(&message)).unwrap();
18 assert_eq!(
19 subscriber.try_recv_into(black_box(&mut received)).unwrap(),
20 Some(size)
21 );
22 black_box(&received);
23 }
24 if round >= 4 {
25 samples.push(start.elapsed().as_nanos() as f64 / ITERATIONS as f64);
26 }
27 }
28 let mean = samples.iter().sum::<f64>() / samples.len() as f64;
29 let deviation = (samples.iter().map(|x| (x - mean).powi(2)).sum::<f64>()
30 / (samples.len() - 1) as f64)
31 .sqrt();
32 println!(
33 "{size} bytes: {mean:.2} ns/roundtrip, stddev {deviation:.2}, samples {samples:?}"
34 );
35 }
36}Sourcepub fn recv_timeout(&self, timeout: Duration) -> Result<Option<Vec<u8>>>
pub fn recv_timeout(&self, timeout: Duration) -> Result<Option<Vec<u8>>>
Waits for a message, or returns None after the timeout. A zero timeout performs one attempt. Missed notifications retain the five-millisecond retry.
Examples found in repository?
examples/interop.rs (line 56)
14fn main() {
15 let args = std::env::args().collect::<Vec<_>>();
16 let options = Options::new(
17 &args[2],
18 std::env::var("INTEROP_CAPACITY")
19 .map(|v| v.parse().unwrap())
20 .unwrap_or(4096),
21 )
22 .with_path(&args[3]);
23 let count = args[4].parse::<usize>().unwrap();
24 if args[1].starts_with("hold-") {
25 let _publisher;
26 let _subscriber;
27 if args[1] == "hold-publisher" {
28 _publisher = Publisher::open(&options).unwrap();
29 } else {
30 _subscriber = Subscriber::open(&options).unwrap();
31 }
32 println!("READY");
33 io::stdout().flush().unwrap();
34 io::stdin().read_line(&mut String::new()).unwrap();
35 } else if args[1] == "publish" {
36 let publisher = Publisher::open(&options).unwrap();
37 let start = args.get(5).map_or(0, |s| s.parse::<usize>().unwrap());
38 if args.len() > 5 {
39 println!("READY");
40 io::stdout().flush().unwrap();
41 io::stdin().read_line(&mut String::new()).unwrap();
42 }
43 for i in start..start + count {
44 let data = message(i);
45 while let Err(error) = publisher.try_send(&data) {
46 assert!(error.is_full(), "{error}");
47 std::thread::yield_now();
48 }
49 }
50 } else {
51 let subscriber = Subscriber::open(&options).unwrap();
52 println!("READY");
53 if args[1] == "collect" {
54 loop {
55 let data = subscriber
56 .recv_timeout(Duration::from_secs(30))
57 .unwrap()
58 .unwrap();
59 if data.is_empty() {
60 break;
61 }
62 let id = u64::from_le_bytes(data[..8].try_into().unwrap()) as usize;
63 assert_eq!(data, message(id));
64 println!("{id}");
65 }
66 return;
67 }
68 for i in 0..count {
69 assert_eq!(
70 subscriber
71 .recv_timeout(Duration::from_secs(30))
72 .unwrap()
73 .unwrap(),
74 message(i)
75 );
76 }
77 }
78}Trait Implementations§
Source§impl Debug for Subscriber
impl Debug for Subscriber
impl Sync for Subscriber
Auto Trait Implementations§
impl !Freeze for Subscriber
impl !RefUnwindSafe for Subscriber
impl Send for Subscriber
impl Unpin for Subscriber
impl UnsafeUnpin for Subscriber
impl UnwindSafe for Subscriber
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more