use crate::modules::input::token;
use crate::{
RuntimeError,
futures::{
net::stream::{RecvTask, SendTask},
task::{Task, sealed::Step},
},
};
use std::sync::Arc;
#[derive(Default)]
pub(crate) enum Stage {
#[default]
Connecting,
Sending(SendTask),
Reading(RecvTask),
}
pub(crate) trait Sends {
fn send_all(&self, data: Arc<[u8]>) -> SendTask;
}
pub(crate) fn advance<C, S>(
connect: &mut C,
stage: &mut Stage,
data: &Arc<[u8]>,
reactor_id: i32,
task_id: usize,
) -> Step<Result<Vec<u8>, RuntimeError>>
where
C: Task<Output = Result<S, RuntimeError>>,
S: Sends,
{
loop {
match &mut *stage {
Stage::Connecting => match connect.step(token(), reactor_id, task_id) {
Step::Done(Ok(conn)) => {
*stage = Stage::Sending(conn.send_all(Arc::clone(data)));
}
Step::Done(Err(error)) => return Step::Done(Err(error)),
Step::Park(park) => return Step::Park(park),
},
Stage::Sending(send) => match send.step(token(), reactor_id, task_id) {
Step::Done(Ok(_)) => {
let read = RecvTask::to_end(send.source().clone());
*stage = Stage::Reading(read);
}
Step::Done(Err(error)) => return Step::Done(Err(error)),
Step::Park(park) => return Step::Park(park),
},
Stage::Reading(read) => return read.step(token(), reactor_id, task_id),
}
}
}