use std::collections::HashMap;
use futures_channel::{mpsc, oneshot};
use futures_core::Future;
use futures_util::{stream::FuturesUnordered, StreamExt};
use js_sys::Array;
use serde::Serialize;
use crate::{port::Port, Dispatcher, MessageHeader};
pub type Outgoing<Response> = (Response, Array, Array);
pub enum ExecuteResult<Response> {
Response(Option<Outgoing<Response>>),
StreamComplete,
}
pub type StreamMessage<Response> = (u32, Option<Outgoing<Response>>);
pub trait Service {
type Response;
fn execute(
&self,
sequence: u32,
abort_rx: oneshot::Receiver<()>,
payload: Vec<u8>,
js_args: Array,
stream_tx: mpsc::UnboundedSender<StreamMessage<Self::Response>>,
) -> impl Future<Output = (u32, ExecuteResult<Self::Response>)>;
}
pub type Request = (u32, Vec<u8>, Array);
pub(crate) async fn task<S>(
service: S,
port: Port,
mut dispatcher: Dispatcher,
mut requests_rx: mpsc::UnboundedReceiver<Request>,
mut aborts_rx: mpsc::UnboundedReceiver<u32>,
) where
S: Service + 'static,
S::Response: Serialize,
{
let (stream_tx, mut stream_rx) = mpsc::unbounded();
let mut running: HashMap<u32, oneshot::Sender<()>> = HashMap::new();
let mut executions: FuturesUnordered<_> = FuturesUnordered::new();
loop {
futures_util::select! {
_ = dispatcher => {}
request = requests_rx.next() => {
let (sequence, payload, js_args) = request.expect("web_rpc: the request channel closed");
let (abort_tx, abort_rx) = oneshot::channel();
running.insert(sequence, abort_tx);
executions.push(service.execute(sequence, abort_rx, payload, js_args, stream_tx.clone()));
},
abort = aborts_rx.next() => {
if let Some(sequence) = abort {
if let Some(abort_tx) = running.remove(&sequence) {
let _ = abort_tx.send(());
}
}
},
message = stream_rx.next() => {
if let Some((sequence, message)) = message {
match message {
Some((item, post_args, transfer_args)) => {
crate::post_message(&port, MessageHeader::StreamItem(sequence), &item, &post_args, &transfer_args);
}
None => {
running.remove(&sequence);
crate::post_header(&port, MessageHeader::StreamEnd(sequence));
}
}
}
},
execution = executions.next() => {
if let Some((sequence, result)) = execution {
match result {
ExecuteResult::Response(response) => {
if running.remove(&sequence).is_some() {
if let Some((response, post_args, transfer_args)) = response {
crate::post_message(&port, MessageHeader::Response(sequence), &response, &post_args, &transfer_args);
}
}
}
ExecuteResult::StreamComplete => {}
}
}
}
}
}
}