pub struct Publisher { /* private fields */ }Expand description
A concurrently usable publisher. Dropping it releases its registration.
Implementations§
Source§impl Publisher
impl Publisher
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 9)
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 28)
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_send(&self, message: &[u8]) -> Result<()>
pub fn try_send(&self, message: &[u8]) -> Result<()>
Returns Error::Full when there is insufficient space or recovery closes admission.
Examples found in repository?
examples/benchmark.rs (line 17)
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 45)
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_send_batch(&self, messages: &[&[u8]]) -> Result<usize>
pub fn try_send_batch(&self, messages: &[&[u8]]) -> Result<usize>
Publishes an ordered prefix, amortizing publisher admission across a batch. Returns the committed prefix length. A short count (including zero) means full, recovery, or a mid-batch error; retry the unsent suffix to observe a persistent error. An error before any commit is returned immediately.
Trait Implementations§
Auto Trait Implementations§
impl Freeze for Publisher
impl RefUnwindSafe for Publisher
impl Send for Publisher
impl Sync for Publisher
impl Unpin for Publisher
impl UnsafeUnpin for Publisher
impl UnwindSafe for Publisher
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