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
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
//! Rust SDK for RLMesh model-environment evaluation workflows.
//!
//! RLMesh connects a model to an environment over gRPC. This crate is the Rust
//! facade over the transport, wire, and runtime crates. Python users should
//! start with the `rlmesh` Python package; use this crate when serving or
//! driving environments and models from Rust.
//!
//! # The two roles
//!
//! Most deployments have one environment server and one model worker.
//!
//! - **Serve an environment.** Implement [`Env`] for one environment and host
//! it with [`EnvServer`]. Implement [`VectorEnv`] and use
//! [`VectorEnvServer`] only for an explicit local batching fast path.
//!
//! - **Drive or serve a model.** Implement [`ModelHandler`], then run it against
//! a remote environment with [`ModelWorker::run_local`] or serve it as an
//! endpoint with [`ModelWorker::serve`].
//!
//! Use [`RemoteEnv`] when you want to step an environment server directly.
//!
//! # Bind-first servers
//!
//! [`EnvServer::bind`] and [`ModelWorker::bind_async`] reserve the socket before
//! serving and return the resolved address, including OS-assigned port 0. This
//! avoids bind-drop-rebind races and poll-connect loops. Use the one-shot
//! `serve`/`serve_async` methods when you do not need that address first.
//!
//! # Errors
//!
//! Fallible operations return [`Result`] (alias for `Result<T, `[`Error`]`>`).
//! [`Error`] separates transport/server faults from two domain failures:
//! [`Error::Environment`] (carrying an [`ErrorCode`]) and [`Error::Model`] (a
//! failure your [`ModelHandler`] raised). Both carry an `is_recoverable` flag
//! surfaced by [`Error::is_recoverable`].
//!
//! # Implementing the traits
//!
//! [`Env`], [`VectorEnv`], and [`ModelHandler`] are `async-trait` traits, so
//! every impl carries `#[rlmesh::async_trait]`: this crate re-exports the macro
//! as [`async_trait`](macro@async_trait) (also in the [`prelude`]), already at
//! the version the traits were desugared with, so nothing beyond `rlmesh` needs
//! to be in your `Cargo.toml`.
//!
//! # Example: serve an environment
//!
//! ```no_run
//! use rlmesh::prelude::*;
//!
//! struct MyEnv {
//! observation_space: SpaceSpec,
//! action_space: SpaceSpec,
//! contract: EnvContract,
//! }
//!
//! #[rlmesh::async_trait]
//! impl Env for MyEnv {
//! fn observation_space(&self) -> &SpaceSpec { &self.observation_space }
//! fn action_space(&self) -> &SpaceSpec { &self.action_space }
//! fn env_contract(&self) -> &EnvContract { &self.contract }
//!
//! // Env methods use the two-arg std::result::Result form.
//! async fn reset(&mut self, _req: ResetRequest)
//! -> Result<ResetResult, EnvRuntimeError>
//! {
//! Ok(ResetResult::default())
//! }
//! async fn step(&mut self, _req: StepRequest)
//! -> Result<StepResult, EnvRuntimeError>
//! {
//! Ok(StepResult::default())
//! }
//! async fn render(&mut self, _req: RenderRequest)
//! -> Result<RenderResult, EnvRuntimeError>
//! {
//! Ok(RenderResult::default())
//! }
//! async fn close(&mut self, _req: CloseRequest)
//! -> Result<CloseResult, EnvRuntimeError>
//! {
//! Ok(CloseResult::default())
//! }
//! }
//!
//! # async fn run(env: MyEnv) -> rlmesh::Result<()> {
//! // Bind first when the caller needs the resolved address.
//! let bound = EnvServer::new(env).bind(BindAddress::parse("tcp://127.0.0.1:0")?).await?;
//! println!("listening on {}", bound.local_addr());
//! bound.serve().await
//! # }
//! ```
//!
//! # Example: drive a model against that environment
//!
//! ```no_run
//! use rlmesh::prelude::*;
//!
//! struct MyModel;
//!
//! #[rlmesh::async_trait]
//! impl ModelHandler for MyModel {
//! async fn predict(&mut self, _obs: ModelObservation)
//! -> rlmesh::Result<Vec<SpaceValue>>
//! {
//! // Read `_obs.decoded_lanes()`, run your policy, return one action per lane.
//! Ok(vec![SpaceValue::Discrete(0)])
//! }
//! }
//!
//! # async fn run() -> rlmesh::Result<()> {
//! // Drive a running env server for 100 episodes.
//! let report = ModelWorker::new(MyModel)
//! .run_local_async(RunLocalOptions::parse("tcp://127.0.0.1:50051")?.for_episodes(100))
//! .await?;
//! println!("ran {} steps", report.total_steps);
//! Ok(())
//! # }
//! ```
pub use remove_stale_socket;
pub use ;
pub use async_trait;
pub use ;
pub use ;
pub use ;
pub use ;
pub use telemetry;
pub use ;
pub use ;
pub use ServeOptions;
pub use ;
pub use CancellationToken;
/// Mint one fresh routing/episode id (UUIDv7 — time-ordered, sortable by creation,
/// never repeats). The single id-format home for this crate's id authorities (the
/// direct env client and the in-process / remote model paths). The runtime driver
/// mints its own ids in `rlmesh-runtime` (it does not depend on this crate).
pub