|
| 1 | +//! Provides a mock implmentation of the aggregator actor. |
| 2 | +use async_trait::async_trait; |
| 3 | +use ceramic_actor::{Actor, Handler, Message}; |
| 4 | +use mockall::mock; |
| 5 | +use prometheus_client::registry::Registry; |
| 6 | + |
| 7 | +use crate::metrics::Metrics; |
| 8 | + |
| 9 | +use super::{ |
| 10 | + NewEventStatesMsg, Resolver, ResolverActor, ResolverEnvelope, ResolverHandle, StreamStateMsg, |
| 11 | + SubscribeSinceMsg, |
| 12 | +}; |
| 13 | + |
| 14 | +mock! { |
| 15 | + // mockall does not support multiple methods on the struct with the same name. |
| 16 | + // This arises when implementing multiple traits that have methods with the same name as is |
| 17 | + // the case with the [`ceramic_actor::Handler`] trait. |
| 18 | + // |
| 19 | + // We add a layer of indirection to get around this limitation. |
| 20 | + pub Resolver { |
| 21 | + #[allow(missing_docs)] |
| 22 | + pub fn handle_subscribe_since( |
| 23 | + &mut self, |
| 24 | + message: SubscribeSinceMsg, |
| 25 | + ) -> <SubscribeSinceMsg as Message>::Result; |
| 26 | + #[allow(missing_docs)] |
| 27 | + pub fn handle_new_event_states( |
| 28 | + &mut self, |
| 29 | + message: NewEventStatesMsg, |
| 30 | + ) -> <NewEventStatesMsg as Message>::Result; |
| 31 | + #[allow(missing_docs)] |
| 32 | + pub fn handle_stream_state( |
| 33 | + &mut self, |
| 34 | + message: StreamStateMsg, |
| 35 | + ) -> <StreamStateMsg |
| 36 | + as Message>::Result; |
| 37 | + } |
| 38 | +} |
| 39 | + |
| 40 | +#[async_trait] |
| 41 | +impl Handler<SubscribeSinceMsg> for MockResolver { |
| 42 | + async fn handle( |
| 43 | + &mut self, |
| 44 | + message: SubscribeSinceMsg, |
| 45 | + ) -> <SubscribeSinceMsg as Message>::Result { |
| 46 | + self.handle_subscribe_since(message) |
| 47 | + } |
| 48 | +} |
| 49 | + |
| 50 | +#[async_trait] |
| 51 | +impl Handler<NewEventStatesMsg> for MockResolver { |
| 52 | + async fn handle( |
| 53 | + &mut self, |
| 54 | + message: NewEventStatesMsg, |
| 55 | + ) -> <NewEventStatesMsg as Message>::Result { |
| 56 | + self.handle_new_event_states(message) |
| 57 | + } |
| 58 | +} |
| 59 | + |
| 60 | +#[async_trait] |
| 61 | +impl Handler<StreamStateMsg> for MockResolver { |
| 62 | + async fn handle(&mut self, message: StreamStateMsg) -> <StreamStateMsg as Message>::Result { |
| 63 | + self.handle_stream_state(message) |
| 64 | + } |
| 65 | +} |
| 66 | + |
| 67 | +impl Actor for MockResolver { |
| 68 | + type Envelope = ResolverEnvelope; |
| 69 | +} |
| 70 | +impl ResolverActor for MockResolver {} |
| 71 | + |
| 72 | +impl MockResolver { |
| 73 | + /// Spawn a mock aggregator actor. |
| 74 | + pub fn spawn(mock_actor: MockResolver) -> ResolverHandle { |
| 75 | + let metrics = Metrics::register(&mut Registry::default()); |
| 76 | + let (handle, _task_handle) = |
| 77 | + Resolver::spawn(1_000, mock_actor, metrics, std::future::pending()); |
| 78 | + handle |
| 79 | + } |
| 80 | +} |
0 commit comments