use super::handler_registry::TheHandlerRegistry;
use super::AxonServerHandle;
use crate::axon_server::event::event_store_client::EventStoreClient;
use crate::axon_server::event::{Event, EventWithToken, GetEventsRequest};
use crate::axon_utils::WorkerControl;
use crate::intellij_work_around::Debuggable;
use anyhow::Result;
use async_stream::stream;
use futures_core::stream::Stream;
use log::{debug, info};
use tokio::select;
use tokio::sync::mpsc::{channel, Receiver, Sender};
#[derive(Debug)]
struct AxonEventProcessed {
message_identifier: String,
}
#[tonic::async_trait]
pub trait TokenStore {
async fn store_token(&self, token: i64);
async fn retrieve_token(&self) -> Result<i64>;
}
pub async fn event_processor<Q: TokenStore + Send + Sync + Clone>(
axon_server_handle: AxonServerHandle,
query_model: Q,
event_handler_registry: TheHandlerRegistry<Q, Event, Option<Q>>,
worker_control: WorkerControl,
) -> Result<()> {
let WorkerControl {
control_channel,
label,
} = &worker_control;
debug!("Event processor: start: {:?}", label);
let conn = axon_server_handle.conn.clone();
let mut client = EventStoreClient::new(conn);
let (tx, rx): (Sender<AxonEventProcessed>, Receiver<AxonEventProcessed>) = channel(10);
let initial_token = query_model.retrieve_token().await.unwrap_or(-1) + 1;
debug!("Initial token: {:?}", initial_token);
let outbound = create_output_stream(label.clone(), axon_server_handle, initial_token, rx);
debug!("Event Processor: calling open_stream");
let response = select! {
response_result = client.list_events(outbound) => response_result?,
_command = control_channel.recv() => {
info!("Event processor stopped while waiting for open stream: {:?}", label);
return Ok(())
}
};
debug!("Stream response: {:?}: {:?}", label, response);
let mut events = response.into_inner();
loop {
debug!("Waiting for event: {:?}", label);
let event_with_token = select! {
event = events.message() => event?,
_command = control_channel.recv() => {
info!("Event processor stopped: {:?}", label);
return Ok(())
}
};
debug!(
"Event with token: {:?}: {:?}",
label,
event_with_token.as_ref().map(|e| Debuggable::from(e))
);
if let Some(EventWithToken {
event: Some(event),
token,
..
}) = event_with_token
{
let message_identifier = event.message_identifier.clone();
if let Event {
payload: Some(serialized_object),
..
} = event.clone()
{
let event_type = serialized_object.r#type;
let mut event_handler_option = event_handler_registry.handlers.get(&event_type);
if event_handler_option.is_none() {
for (regex, event_handler) in &event_handler_registry.category_handlers {
if regex.is_match(&event_type) {
event_handler_option = Some(event_handler);
break;
}
}
}
if let Some(event_handler) = event_handler_option {
(event_handler)
.handle(serialized_object.data, event, query_model.clone())
.await?;
}
}
query_model.store_token(token).await;
tx.send(AxonEventProcessed { message_identifier }).await?;
}
}
}
fn create_output_stream(
label: String,
axon_server_handle: AxonServerHandle,
initial_token: i64,
mut rx: Receiver<AxonEventProcessed>,
) -> impl Stream<Item = GetEventsRequest> {
stream! {
debug!("Event Processor: stream: start: {:?}: {:?}", &label, rx);
let permits_batch_size: i64 = 3;
let mut permits = permits_batch_size * 2;
let processor = format!("Event processor: {:?}", &label);
let mut request = GetEventsRequest {
tracking_token: initial_token,
number_of_permits: permits,
client_id: axon_server_handle.client_id.clone(),
component_name: axon_server_handle.display_name.clone(),
processor,
blacklist: Vec::new(),
force_read_from_leader: false,
};
yield request.clone();
request.number_of_permits = permits_batch_size;
while let Some(axon_event_processed) = rx.recv().await {
debug!("Event processed: {:?}: {:?}", &label, axon_event_processed.message_identifier);
permits -= 1;
if permits <= permits_batch_size {
debug!("Event Processor: stream: send more flow-control permits: {:?}: amount: {:?}", &label, permits_batch_size);
yield request.clone();
permits += permits_batch_size;
}
debug!("Event Processor: stream: flow-control permits: {:?}: balance: {:?}", &label, permits);
}
}
}