1use std::collections::HashMap;
2use std::sync::Arc;
3
4use mm1_address::address::Address;
5use mm1_common::errors::chain::ExactTypeDisplayChainExt;
6use mm1_common::types::AnyError;
7use mm1_core::tap::{MessageTap, NoopTap};
8use mm1_core::tracing::TraceId;
9use mm1_runnable::local::{self, BoxedRunnable};
10use tokio::runtime::{Handle, Runtime};
11use tokio::sync::mpsc;
12use tracing::{error, instrument, trace};
13
14use crate::actor_key::ActorKey;
15use crate::config::{EffectiveActorConfig, Mm1NodeConfig, Valid};
16use crate::init::InitActorArgs;
17use crate::runtime::container::{Container, ContainerArgs, ContainerError};
18use crate::runtime::context;
19use crate::runtime::rt_api::{RequestAddressError, RtApi};
20
21#[derive(derive_more::Debug)]
22pub struct Rt {
23 #[allow(unused)]
24 config: Valid<Mm1NodeConfig>,
25 rt_default: Runtime,
26 rt_named: HashMap<String, Runtime>,
27 tx_actor_failure: mpsc::UnboundedSender<(Address, AnyError)>,
28 #[debug(skip)]
29 message_taps: HashMap<String, Arc<dyn MessageTap>>,
30}
31
32#[derive(Debug, thiserror::Error)]
33pub enum RtCreateError {
34 #[error("runtime config error: {}", _0)]
35 RuntimeConfigError(crate::config::ValidationError<Mm1NodeConfig>),
36 #[error("runtime init error: {}", _0)]
37 RuntimeInitError(#[source] std::io::Error),
38}
39
40#[derive(Debug, thiserror::Error)]
41pub enum RtRunError {
42 #[error("request address error: {}", _0)]
43 RequestAddressError(
44 #[allow(private_interfaces)]
45 #[source]
46 RequestAddressError,
47 ),
48 #[error("container error: {}", _0)]
49 #[allow(private_interfaces)]
50 ContainerError(#[source] ContainerError),
51}
52
53impl Rt {
54 pub fn create(config: Mm1NodeConfig) -> Result<Self, RtCreateError> {
55 let config = config
56 .validate()
57 .map_err(RtCreateError::RuntimeConfigError)?;
58 let (rt_default, rt_named) = config
59 .build_runtimes()
60 .map_err(RtCreateError::RuntimeInitError)?;
61
62 let (tx_actor_failure, _rx_actor_failure) = mpsc::unbounded_channel();
63 Ok(Self {
64 config,
65 rt_default,
66 rt_named,
67 tx_actor_failure,
68 message_taps: Default::default(),
69 })
70 }
71
72 pub fn with_actor_failure_sink(
73 self,
74 tx_actor_failure: mpsc::UnboundedSender<(Address, AnyError)>,
75 ) -> Self {
76 Self {
77 tx_actor_failure,
78 ..self
79 }
80 }
81
82 pub fn with_tap(mut self, key: impl Into<String>, tap: Arc<dyn MessageTap>) -> Self {
83 self.message_taps.insert(key.into(), tap);
84 self
85 }
86
87 pub fn run(&self, main_actor: BoxedRunnable<context::ActorContext>) -> Result<(), RtRunError> {
94 let config = self.config.clone();
95 let rt_default = self.rt_default.handle().to_owned();
96 let rt_named = self
97 .rt_named
98 .iter()
99 .map(|(k, v)| (k.to_owned(), v.handle().to_owned()))
100 .collect::<HashMap<_, _>>();
101
102 let rt_handle = rt_default.clone();
103
104 let init_actor_args = InitActorArgs {
105 local_subnet_auto: config.local_subnet_address_auto(),
106 local_subnets_bind: config.local_subnet_addresses_bind().collect(),
107 #[cfg(feature = "multinode")]
108 multinode_inbound: config.multinode_inbound().collect(),
109 #[cfg(feature = "multinode")]
110 multinode_outbound: config.multinode_outbound().collect(),
111 };
112
113 rt_handle.block_on(run_inner(
114 config.clone(),
115 ActorKey::root(),
116 crate::init::init_actor_config(),
117 rt_default,
118 rt_named,
119 local::boxed_from_fn((crate::init::run, (main_actor, init_actor_args))),
120 &self.message_taps,
121 self.tx_actor_failure.clone(),
122 ))
123 }
124}
125
126#[allow(clippy::too_many_arguments)]
127#[instrument(skip_all, fields(func = main_actor.func_name()))]
128async fn run_inner(
129 config: Valid<Mm1NodeConfig>,
130 actor_key: ActorKey,
131 actor_config: impl EffectiveActorConfig,
132 rt_default: Handle,
133 rt_named: HashMap<String, Handle>,
134 main_actor: BoxedRunnable<context::ActorContext>,
135 message_taps: &HashMap<String, Arc<dyn MessageTap>>,
136 tx_actor_failure: mpsc::UnboundedSender<(Address, AnyError)>,
137) -> Result<(), RtRunError> {
138 let rt_api = RtApi::create(
139 config.local_subnet_address_auto(),
140 rt_default,
141 rt_named,
142 Arc::new(NoopTap),
143 message_taps.clone(),
144 );
145
146 let subnet_lease = rt_api
147 .request_address(actor_config.netmask())
148 .await
149 .map_err(RtRunError::RequestAddressError)?;
150 let message_tap_key = actor_config.message_tap_key();
151 let message_tap = rt_api.message_tap(message_tap_key);
152
153 let args = ContainerArgs {
154 ack_to: None,
155 link_to: Default::default(),
156 actor_key,
157 trace_id: TraceId::random(),
158 subnet_lease,
159 rt_api: rt_api.clone(),
160 rt_config: Arc::new(config),
161 message_tap,
162 tx_actor_failure,
163 };
164 trace!("creating and running container...");
165 let container = Container::create(args, main_actor).map_err(RtRunError::ContainerError)?;
166 let exit_result = container.run().await;
167
168 rt_api.shutdown().await;
169
170 let exit_reason = exit_result.map_err(RtRunError::ContainerError)?;
171
172 trace!(reason = ?exit_reason.as_ref().map_err(|e| e.as_display_chain()), "container exited");
173
174 if let Err(failure) = exit_reason {
175 error!(
176 error = %failure.as_display_chain(),
177 "main-actor failure"
178 );
179 }
180
181 Ok(())
182}