January 2019
Intermediate to advanced
520 pages
14h 32m
English
We used the basic_consume method of a Channel to create a Stream that returns Delivery objects from a queue. Since we want to attach that Stream to the actor, we have to implement StreamHandler for the QueueActor type:
impl<T: QueueHandler> StreamHandler<Delivery, LapinError> for QueueActor<T> { fn handle(&mut self, item: Delivery, ctx: &mut Context<Self>) { debug!("Message received!"); let fut = self .channel .basic_ack(item.delivery_tag, false) .map_err(drop); ctx.spawn(wrap_future(fut)); match self.process_message(item, ctx) { Ok(pair) => { if let Some((corr_id, data)) = pair { self.send_message(corr_id, data, ctx); } } Err(err) => { warn!("Message processing error: {}", err); } } } }
Our StreamHandler implementation ...
Read now
Unlock full access