use std::time::Duration;
use actix::actors::resolver::{Connect, Resolver};
use actix::prelude::*;
use actix::io::{FramedWrite, WriteHandler};
use backoff::backoff::Backoff;
use backoff::ExponentialBackoff;
use log::{error, info, trace, warn};
use tokio::io::WriteHalf;
use tokio::net::TcpStream;
use tokio_util::codec::{FramedRead, LinesCodec, LinesCodecError};
#[derive(Message, Clone)]
#[rtype(result = "()")]
pub struct OGNMessage {
pub raw: String,
}
pub struct OGNActor {
recipient: Recipient<OGNMessage>,
backoff: ExponentialBackoff,
writer: Option<FramedWrite<String, WriteHalf<TcpStream>, LinesCodec>>,
}
impl OGNActor {
pub fn new(recipient: Recipient<OGNMessage>) -> OGNActor {
let mut backoff = ExponentialBackoff::default();
backoff.max_elapsed_time = None;
OGNActor { recipient, backoff, writer: None }
}
fn schedule_keepalive(ctx: &mut Context<Self>) {
ctx.run_later(Duration::from_secs(30), |act, ctx| {
info!("Sending keepalive to OGN server");
if let Some(ref mut writer) = act.writer {
writer.write("# keep alive".to_string());
}
OGNActor::schedule_keepalive(ctx);
});
}
}
impl Actor for OGNActor {
type Context = Context<Self>;
fn started(&mut self, ctx: &mut Self::Context) {
info!("Connecting to OGN server...");
Resolver::from_registry()
.send(Connect::host("aprs.glidernet.org:10152"))
.into_actor(self)
.map(|res, act, ctx| match res {
Ok(Ok(stream)) => {
info!("Connected to OGN server");
act.backoff.reset();
let (r, w) = tokio::io::split(stream);
let mut writer = FramedWrite::new(w, LinesCodec::new(), ctx);
let login_message = {
let username = "test";
let password = "-1";
let app_name = option_env!("CARGO_PKG_NAME").unwrap_or("unknown");
let app_version = option_env!("CARGO_PKG_VERSION").unwrap_or("0.0.0");
format!(
"user {} pass {} vers {} {}",
username,
password,
app_name,
app_version,
)
};
writer.write(login_message);
act.writer = Some(writer);
ctx.add_stream(FramedRead::new(r, LinesCodec::new()));
OGNActor::schedule_keepalive(ctx);
}
Ok(Err(err)) => {
error!("Can not connect to OGN server: {}", err);
if let Some(timeout) = act.backoff.next_backoff() {
ctx.run_later(timeout, |_, ctx| ctx.stop());
} else {
ctx.stop();
}
}
Err(err) => {
error!("Can not connect to OGN server: {}", err);
if let Some(timeout) = act.backoff.next_backoff() {
ctx.run_later(timeout, |_, ctx| ctx.stop());
} else {
ctx.stop();
}
}
})
.wait(ctx);
}
fn stopped(&mut self, _: &mut Self::Context) {
info!("Disconnected from OGN server");
}
}
impl Supervised for OGNActor {
fn restarting(&mut self, _: &mut Self::Context) {
info!("Restarting OGN client...");
self.writer.take();
}
}
impl WriteHandler<LinesCodecError> for OGNActor {
fn error(&mut self, err: LinesCodecError, _: &mut Self::Context) -> Running {
warn!("OGN connection dropped: error: {}", err);
Running::Stop
}
}
impl StreamHandler<Result<String, LinesCodecError>> for OGNActor {
fn handle(&mut self, line: Result<String, LinesCodecError>, _: &mut Self::Context) {
if let Ok(line) = line {
trace!("{}", line);
if !line.starts_with('#') {
if let Err(error) = self.recipient.do_send(OGNMessage { raw: line }) {
warn!("do_send failed: {}", error);
}
}
}
}
}