|
| 1 | +//! AMQP 1.0 binding example |
| 2 | +//! |
| 3 | +//! You need a running AMQP 1.0 broker to try out this example. |
| 4 | +//! With docker: docker run -it --rm -e ARTEMIS_USERNAME=guest -e ARTEMIS_PASSWORD=guest -p 5672:5672 vromero/activemq-artemis |
| 5 | +
|
| 6 | +use cloudevents::{binding::fe2o3_amqp::EventMessage, message::BinaryDeserializer, Event, EventBuilderV10, EventBuilder}; |
| 7 | +use fe2o3_amqp::{Connection, Sender, Receiver, types::messaging::Message, Session}; |
| 8 | +use serde_json::json; |
| 9 | + |
| 10 | +type BoxError = Box<dyn std::error::Error>; |
| 11 | +type Result<T> = std::result::Result<T, BoxError>; |
| 12 | + |
| 13 | +async fn send_event(sender: &mut Sender, i: usize) -> Result<()> { |
| 14 | + let event = EventBuilderV10::new() |
| 15 | + .id(i.to_string()) |
| 16 | + .ty("example.test") |
| 17 | + .source("localhost") |
| 18 | + .data("application/json", json!({"hello": "world"})) |
| 19 | + .build()?; |
| 20 | + let event_message = EventMessage::from_binary_event(event)?; |
| 21 | + let message = Message::from(event_message); |
| 22 | + sender.send(message).await? |
| 23 | + .accepted_or("not accepted")?; |
| 24 | + Ok(()) |
| 25 | +} |
| 26 | + |
| 27 | +async fn recv_event(receiver: &mut Receiver) -> Result<Event> { |
| 28 | + use fe2o3_amqp::types::primitives::Value; |
| 29 | + |
| 30 | + let delivery = receiver.recv::<Value>().await?; |
| 31 | + receiver.accept(&delivery).await?; |
| 32 | + |
| 33 | + let event_message = EventMessage::from(delivery.into_message()); |
| 34 | + let event = event_message.into_event()?; |
| 35 | + Ok(event) |
| 36 | +} |
| 37 | + |
| 38 | +#[tokio::main] |
| 39 | +async fn main() { |
| 40 | + let mut connection = |
| 41 | + Connection::open("cloudevents-sdk-rust", "amqp://guest:guest@localhost:5672") |
| 42 | + .await |
| 43 | + .unwrap(); |
| 44 | + let mut session = Session::begin(&mut connection).await.unwrap(); |
| 45 | + let mut sender = Sender::attach(&mut session, "sender", "q1").await.unwrap(); |
| 46 | + let mut receiver = Receiver::attach(&mut session, "receiver", "q1").await.unwrap(); |
| 47 | + |
| 48 | + send_event(&mut sender, 1).await.unwrap(); |
| 49 | + let event = recv_event(&mut receiver).await.unwrap(); |
| 50 | + println!("{:?}", event); |
| 51 | + |
| 52 | + sender.close().await.unwrap(); |
| 53 | + receiver.close().await.unwrap(); |
| 54 | + session.end().await.unwrap(); |
| 55 | + connection.close().await.unwrap(); |
| 56 | +} |
0 commit comments