rill_server/actors/server/
actor.rs1use super::link;
2use crate::actors::router::Router;
3use crate::config::ServerConfig;
4use anyhow::Error;
5use async_trait::async_trait;
6use meio::{Actor, Context, Eliminated, IdOf, InteractionHandler, InterruptedBy, StartedBy};
7use meio_connect::server::{HttpServer, HttpServerLink};
8use rill_client::actors::broadcaster::Broadcaster;
9
10pub struct RillServer {
12 server_config: ServerConfig,
13 public_server: Option<HttpServerLink>,
14 private_server: Option<HttpServerLink>,
15}
16
17#[derive(Debug, Clone, PartialEq, Eq, Hash)]
18pub enum Group {
19 Tuning,
20 Broadcaster,
21 HttpServer,
22 Endpoints,
23}
24
25impl Actor for RillServer {
26 type GroupBy = Group;
27}
28
29impl RillServer {
30 pub fn new(server_config: Option<ServerConfig>) -> Self {
32 Self {
33 server_config: server_config.unwrap_or_default(),
34 public_server: None,
35 private_server: None,
36 }
37 }
38}
39
40#[async_trait]
41impl<T: Actor> StartedBy<T> for RillServer {
42 async fn handle(&mut self, ctx: &mut Context<Self>) -> Result<(), Error> {
43 ctx.termination_sequence(vec![
44 Group::Tuning,
45 Group::Broadcaster,
46 Group::HttpServer,
47 Group::Endpoints,
48 ]);
49
50 let extern_addr = format!("{}:{}", self.server_config.server_address(), 9090).parse()?;
52 let extern_http_server_actor = HttpServer::new(extern_addr);
53 let extern_http_server = ctx.spawn_actor(extern_http_server_actor, Group::HttpServer);
54 self.public_server = Some(extern_http_server.link());
55
56 let inner_addr = format!("127.0.0.1:{}", 0).parse()?;
57 let inner_http_server_actor = HttpServer::new(inner_addr);
58 let inner_http_server = ctx.spawn_actor(inner_http_server_actor, Group::HttpServer);
59 self.private_server = Some(inner_http_server.link());
60
61 let exporter_actor = Broadcaster::new();
62 let exporter = ctx.spawn_actor(exporter_actor, Group::Broadcaster);
63
64 let server_actor = Router::new(
65 inner_http_server.link(),
66 extern_http_server.link(),
67 exporter.link(),
68 );
69 let _router = ctx.spawn_actor(server_actor, Group::Endpoints);
70
71 Ok(())
72 }
73}
74
75#[async_trait]
76impl<T: Actor> InterruptedBy<T> for RillServer {
77 async fn handle(&mut self, ctx: &mut Context<Self>) -> Result<(), Error> {
78 ctx.shutdown();
79 Ok(())
80 }
81}
82
83#[async_trait]
84impl Eliminated<Broadcaster> for RillServer {
85 async fn handle(
86 &mut self,
87 _id: IdOf<Broadcaster>,
88 _ctx: &mut Context<Self>,
89 ) -> Result<(), Error> {
90 log::info!("Broadcaster finished");
91 Ok(())
92 }
93}
94
95#[async_trait]
96impl Eliminated<HttpServer> for RillServer {
97 async fn handle(
98 &mut self,
99 _id: IdOf<HttpServer>,
100 _ctx: &mut Context<Self>,
101 ) -> Result<(), Error> {
102 log::info!("HttpServer finished");
103 Ok(())
104 }
105}
106
107#[async_trait]
108impl Eliminated<Router> for RillServer {
109 async fn handle(&mut self, _id: IdOf<Router>, _ctx: &mut Context<Self>) -> Result<(), Error> {
110 log::info!("Router finished");
111 Ok(())
112 }
113}
114
115#[async_trait]
116impl InteractionHandler<link::WaitPublicEndpoint> for RillServer {
117 async fn handle(
118 &mut self,
119 _msg: link::WaitPublicEndpoint,
120 _ctx: &mut Context<Self>,
121 ) -> Result<HttpServerLink, Error> {
122 self.public_server
124 .clone()
125 .ok_or_else(|| Error::msg("no public server"))
126 }
127}
128
129#[async_trait]
130impl InteractionHandler<link::WaitPrivateEndpoint> for RillServer {
131 async fn handle(
132 &mut self,
133 _msg: link::WaitPrivateEndpoint,
134 _ctx: &mut Context<Self>,
135 ) -> Result<HttpServerLink, Error> {
136 self.private_server
138 .clone()
139 .ok_or_else(|| Error::msg("no private server"))
140 }
141}