1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
use tracing::warn;
use wasmrs::{PayloadError, RawPayload};
use wasmrs_guest::{FluxChannel, Observer};

use crate::Packet;

pub struct Output<T>
where
  T: serde::Serialize,
{
  channel: FluxChannel<RawPayload, PayloadError>,
  name: String,
  _phantom: std::marker::PhantomData<T>,
}

impl std::fmt::Debug for Output<()> {
  fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
    f.debug_struct("Output").field("name", &self.name).finish()
  }
}

impl<T> Output<T>
where
  T: serde::Serialize,
{
  pub fn new(name: impl AsRef<str>, channel: FluxChannel<RawPayload, PayloadError>) -> Self {
    Self {
      channel,
      name: name.as_ref().to_owned(),
      _phantom: Default::default(),
    }
  }

  pub fn send(&mut self, value: &T) {
    self.send_raw(Packet::encode(&self.name, value));
  }

  pub fn send_raw(&mut self, value: Packet) {
    if let Err(e) = self.channel.send_result(value.into()) {
      warn!(
        port = self.name,
        error = %e,
        "failed sending packet on output channel, this is a bug"
      );
    };
  }

  pub fn open_bracket(&mut self) {
    self.send_raw(Packet::open_bracket(&self.name));
  }

  pub fn close_bracket(&mut self) {
    self.send_raw(Packet::close_bracket(&self.name));
  }

  pub fn done(&mut self) {
    self.send_raw(Packet::done(&self.name));
  }

  pub fn error(&mut self, err: impl AsRef<str>) {
    self.send_raw(Packet::err(&self.name, err));
  }
}