Skip to main content

rill_server/actors/server/
actor.rs

1use 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
10/// Embedded node.
11pub 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    /// Create a new instance of an embedded node.
31    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        // Starting all basic actors
51        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        // `public_server` always available here since it's attached in `StartedBy` handler
123        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        // `private_server` always available here since it's attached in `StartedBy` handler
137        self.private_server
138            .clone()
139            .ok_or_else(|| Error::msg("no private server"))
140    }
141}