use std::time::Duration;
use async_trait::async_trait;
use futures::{StreamExt, pin_mut};
use tari_service_framework::{
ServiceInitializationError,
ServiceInitializer,
ServiceInitializerContext,
reply_channel,
reply_channel::SenderService,
};
use tari_shutdown::ShutdownSignal;
use tokio::time::sleep;
use tower::Service;
pub struct ServiceB {
response_msg: String,
request_stream: Option<reply_channel::Receiver<String, String>>,
shutdown_signal: Option<ShutdownSignal>,
}
impl ServiceB {
pub fn new(
response_msg: String,
request_stream: reply_channel::Receiver<String, String>,
shutdown_signal: ShutdownSignal,
) -> Self {
Self {
response_msg,
request_stream: Some(request_stream),
shutdown_signal: Some(shutdown_signal),
}
}
pub async fn run(mut self) {
println!("Starting Service B");
let mut shutdown_signal = self
.shutdown_signal
.take()
.expect("Service B initialized without shutdown signal");
let request_stream = self
.request_stream
.take()
.expect("Service B initialized without request_stream")
.fuse();
pin_mut!(request_stream);
loop {
tokio::select! {
request_context = request_stream.select_next_some() => {
println!("Handling Service B API Request");
let (request, reply_tx) = request_context.split();
let mut response = self.response_msg.clone();
response.push_str(request.clone().as_str());
let _resp = reply_tx.send(response);
},
_ = shutdown_signal.wait() => {
println!("Service B shutting down because the shutdown signal was received");
break;
}
}
}
println!("Service B is shutdown");
}
}
#[derive(Clone)]
pub struct ServiceBHandle {
request_tx: SenderService<String, String>,
}
impl ServiceBHandle {
pub fn new(request_tx: SenderService<String, String>) -> Self {
Self { request_tx }
}
pub async fn send_msg(&mut self, msg: String) -> String {
self.request_tx.call(msg).await.unwrap()
}
}
pub struct ServiceBInitializer {
response_msg: String,
}
impl ServiceBInitializer {
pub fn new(response_msg: String) -> Self {
Self { response_msg }
}
}
#[async_trait]
impl ServiceInitializer for ServiceBInitializer {
async fn initialize(&mut self, context: ServiceInitializerContext) -> Result<(), ServiceInitializationError> {
let (sender, receiver) = reply_channel::unbounded();
let service_b_handle = ServiceBHandle::new(sender);
println!("Service B is going to wait to register its handle");
println!("Service B is registering its handle now");
context.register_handle(service_b_handle);
let response_msg = self.response_msg.clone();
println!("Service B initialized waiting on Handles Future to complete");
context.spawn_when_ready(move |handles| async move {
println!("Service B got the handles");
let service = ServiceB::new(response_msg, receiver, handles.get_shutdown_signal());
service.run().await;
println!("Service B has shutdown and initializer spawned task is now ending");
});
sleep(Duration::from_secs(10)).await;
Ok(())
}
}