use crate::host::error::*;
use crate::host::output_sink::*;
use crate::host::initialisation_context::*;
use crate::host::scene_context::*;
use crate::host::scene_message::*;
use crate::host::serialization::*;
use crate::host::stream_target::*;
use futures::prelude::*;
use std::marker::{PhantomData};
use serde::*;
#[derive(Clone)]
#[derive(Serialize, Deserialize)]
pub struct Subscribe<TMessageType: SceneMessage>(StreamTarget, PhantomData<TMessageType>);
impl<TMessageType: SceneMessage> SceneMessage for Subscribe<TMessageType> {
fn initialise(_: &impl SceneInitialisationContext) {
#[cfg(feature="json")]
install_serializable_type(|msg: TMessageType| msg.to_json(), |json| TMessageType::from_json(json)).unwrap();
}
#[inline]
fn message_type_name() -> String { format!("subscribe::{}", TMessageType::message_type_name()) }
}
impl<TMessageType: SceneMessage> Subscribe<TMessageType> {
#[inline]
pub fn with_target(target: StreamTarget) -> Self {
Subscribe(target, PhantomData)
}
#[inline]
pub fn target(&self) -> StreamTarget {
self.0.clone()
}
}
#[inline]
pub fn subscribe<TMessageType: SceneMessage>(target: impl Into<StreamTarget>) -> Subscribe<TMessageType> {
Subscribe::with_target(target.into())
}
pub struct EventSubscribers<TEventMessage>
where
TEventMessage: 'static + SceneMessage,
{
receivers: Vec<OutputSink<TEventMessage>>,
next_receiver: usize,
}
impl<TEventMessage> EventSubscribers<TEventMessage>
where
TEventMessage: 'static + SceneMessage,
{
pub fn new() -> Self {
EventSubscribers {
receivers: vec![],
next_receiver: 0,
}
}
pub fn subscribe(&mut self, context: &SceneContext, target: impl Into<StreamTarget>) {
let target = target.into();
self.receivers.retain(|sink| sink.is_attached());
let output_sink = context.send(target);
let output_sink = if let Ok(output_sink) = output_sink { output_sink } else { return; };
self.receivers.push(output_sink);
}
pub fn add_target(&mut self, output_sink: OutputSink<TEventMessage>) {
self.receivers.push(output_sink)
}
pub async fn send_round_robin(&mut self, message: TEventMessage) -> Result<(), TEventMessage> {
let mut message = message;
loop {
if self.receivers.is_empty() {
break Err(message);
}
self.next_receiver += 1;
if self.next_receiver >= self.receivers.len() {
self.next_receiver = 0;
}
match self.receivers[self.next_receiver].send(message).await {
Ok(()) => { break Ok(()); }
Err(SceneSendError::CouldNotConnect(_)) |
Err(SceneSendError::TargetProgramEndedBeforeReady) |
Err(SceneSendError::ErrorAfterDeserialization) |
Err(SceneSendError::CannotReEnterTargetProgram) => {
self.receivers.remove(self.next_receiver);
break Ok(());
}
Err(SceneSendError::StreamClosed(returned_message)) |
Err(SceneSendError::CannotAcceptMoreInputUntilSceneIsIdle(returned_message)) |
Err(SceneSendError::TargetProgramEnded(returned_message)) |
Err(SceneSendError::CannotDeserialize(returned_message, _)) |
Err(SceneSendError::CannotSerialize(returned_message, _)) |
Err(SceneSendError::NoConnection(returned_message)) |
Err(SceneSendError::StreamDisconnected(returned_message)) => {
self.receivers.remove(self.next_receiver);
if self.next_receiver > 0 {
self.next_receiver -= 1;
} else if !self.receivers.is_empty() {
self.next_receiver = self.receivers.len() - 1;
}
message = returned_message;
}
}
}
}
}
impl<TEventMessage> EventSubscribers<TEventMessage>
where
TEventMessage: 'static + Clone + SceneMessage,
{
pub async fn send(&mut self, message: TEventMessage) -> bool {
self.receivers.retain(|sink| sink.is_attached());
let senders = self.receivers.iter_mut()
.enumerate()
.map(|(idx, sender)| sender.send(message.clone()).map(move |result| (idx, result)))
.collect::<Vec<_>>();
let mut results = future::join_all(senders).await;
let mut sent_successfully = false;
results.sort_by(|(a, _), (b, _)| a.cmp(b));
for (idx, result) in results.into_iter().rev() {
if result.is_err() {
self.receivers.remove(idx);
} else {
sent_successfully = true;
}
}
sent_successfully
}
}