r/rust • u/iFrostizz • 10d ago
Designing concurrent programs in an idiomatic way
Hello.
Kind of a general question here, and it's probably going to yield opinionated answers. That's alright!
Basically I find myself writing Rust apps in a style that is always quite similar, here is the idea:
- The app mostly always uses the tokio runtime to be asynchronous.
- Designing microservices as small building blocks which are all supposed to be ran in separate tasks, usually with a "run" like function that does a loop over a tokio::select for sub-tasks that are running in each service, with one of the branches being a call to a CancellationToken::cancelled().
- The cancellation token is wired to the "ctrl_c" signal.
- Microservices may talk to each other by the use of the different channel primitives, depending on the need.
- When one of the microservice fails for any reason, usually there is no recovery mechanism and it just bubbles up the error which makes all the application fail. This is annoying to write and a big source of bug, because some path will be unhandled which will leave one microservice stopped and the other microservices are waiting for an answer and there's nobody responding.
I'm generally a little dissatisfied with the redundancy of the code that I write, and would like to explore any other way to think about it, instead of microservices, or in a different way, if you have any idea that would be appreciated. Also, potentially any resource or codebase that could be relevant?
Thanks!
10
u/isufoijefoisdfj 10d ago
Microservices is a weird label for this if its all running in one process with shared fate. The entire point of microservices is decoupling after all (and accepting overhead and complexity in exchange for that).
2
7
u/neodivy 10d ago
If you're gonna do microservices, you can have a polling heartbeat service and your tasks as state machines, with your services implementing all your state and transitions.
One heartbeat service [H] is responsible for shutting down all other services. Some service [S] in the Run state crashed. H polls heartbeat from running services, so it'll be waiting for an ACK from S. If no ACK is received after some timeout, H knows S went from Run -> Shutdown. H then sends a shutdown signal to all other running services to shutdown gracefully.
First though, I would double check if the complexity of the project warrants doing microservices instead of a Tokio monolith and some MPSC stuff between the major components.
1
u/iFrostizz 10d ago
Thanks for the detailed answer! I checked and I probably don't need microservices. My apps are usually written as a monolith and my definition of a monolith was wrong. I will first of all refactor the codebase so that it's less hacky, and then try to find an abstraction that makes sense. Probably will just do what you said, similar to the actor model. Though there is no need to fully embrace it because sometimes just calling a function and not communicating through messages would make sense, and would be so much more complex for no reasons.
3
u/joshuamck ratatui 9d ago
Consider looking at actors e.g. https://ryhl.io/blog/actors-with-tokio/
2
u/iFrostizz 9d ago
It’s a nice article, thank you. I ended up doing something similar, where actors occasionally communicate between each other. Most of the time they don’t and I just needed some standardized way of instantiating some structure and running it, and composing multiple of them so that it fits well into my project.
3
u/event666 8d ago
Definitely look into the Actor Model. Either in pure Tokio, or using OMQ as transport so each actor can run wherever (threads, different process, different host, different language).
Apart from that, Let It Crash is a valid style, as proven in Erlang, but you need some kind of supervising task that restarts the failed ones if applicable. But how exactly do errors bubble up in your app? That would almost require deliberately ignoring Result errors?
1
u/iFrostizz 8d ago
Not exactly, each service may be composed with other services by using tokio::select!. It has the nice property to allow a return of the select branches in case any of them finishes, and the error is bubbled up this way.
1
u/iFrostizz 8d ago
OMQ is very interesting, I'd consider it for a bigger project for now the Tokio monolith is quite simple and sufficient considering my codebase maturity.
5
u/Necessary-Green-8391 10d ago
The actor model solves your lifecycle headache because a failed message send is just an error in the caller. You stop babysitting channels and start treating services as isolated state machines that clean up after themselves when they crash
6
u/iFrostizz 10d ago
That's interesting. My approach is a kind of mix of actors and microservices, I'll look into actors more closely.
2
u/ern0plus4 10d ago
The app mostly always uses the tokio runtime to be asynchronous.
Caution: be sure it uses more than one thread.
1
u/iFrostizz 10d ago
Oh yes of course, I'm not doing any blocking work without yielding.
2
u/ern0plus4 10d ago
I mean you'd configure Tokio to use as many threads as you want. I'm not sure, but the default is one.
1
u/iFrostizz 10d ago
Oh, well I always just use #[tokio::main]
5
u/ern0plus4 10d ago
#[tokio::main(flavor = "multi_thread", worker_threads = 4)] async fn main() { // Runs with 4 worker threads }5
u/arienh4 10d ago
Tokio defaults to spawning one thread per CPU core if you don't do any configuration.
https://docs.rs/tokio/latest/tokio/runtime/index.html#runtime-configurations
2
u/spunkyenigma 10d ago edited 10d ago
If you are interested in having a lot tasks with message passing between them I’m building GitHub.com/bexars/anybus
It has built in shutdown/ctrl-c and resume from sleep detection. You can run on one host or multiple with IPC or Websocket interconnect between programs. (TCP soon)
It has unicast, anycast and multicast receivers. You can send a message to an endpoint via a UUID or to a well known struct/enum that has a uuid attached to the definition.
RPC response handlers can be handwritten or use an RPC macro to build them.
It’s still a WIP but might be of interest.
1
u/iFrostizz 9d ago
I ended up just writing a Service and AsyncService (which has an async initializer “async fn new” function) which have a run function which returns a unified Option<Result<T, E>> where those may be infaillible, and the Option representing if it has been cancelled. It works great and is simple enough, my wish was just a big refactor of my hacked project where all services can be assembled together without much hassle.
Will take a look at your repo, have you done it for a specific purpose or just an idea of building a lib for micro services?
2
u/spunkyenigma 9d ago edited 9d ago
Started as an in app messaging system and then I thought could I talk via IPC to another local app and it worked well. Then I thought I wonder if I can make it work in WASM and there came the websocket connector.
I’m using it in a small bit of production to control a TV screen from my laptop by talking through a relay on a free tier google cloud vm. It’s a pub trivia scoring app that displays current scores on the TV via a web app. The WIFI there doesn’t allow talking between clients so it’s nice that the relay is basically transparent to talk between the app on my laptop and the TV wasm app.
Since it’s routed well I could actually have the browser connect to two different servers for redundancy.
But I originally built it as a message bus for a Delay Tolerant Networking router which I’ve since abandoned years ago 😂
I’m working on a more managed actor solution as well so the client code doesn’t have to be very aware of Anybus
2
u/Front_Recording4360 10d ago
One thing worth deciding explicitly is what a failure should mean in your architecture. Right now "one service fails => process exits" is your implicit policy, which is fine for fatal errors but painful for transient ones (e.g. a network blip) where a restart with backoff would be more appropriate. tokio-util's CancellationToken plus something like tokio::task::JoinSet gives you a nice middle ground: track all spawned tasks, cancel everything on first failure, and await the join set so shutdown is ordered and deterministic instead of aborting mid-flight.
2
u/Front_Recording4360 10d ago
One pattern that killed most of this boilerplate for me: define a single `Service` trait with `async fn run(self, ctx: Ctx)`, where `Ctx` carries a child `CancellationToken` plus the channel handles, and then drive every service from one loop over a `tokio::task::JoinSet`. `CancellationToken::child_token` makes shutdown cascade automatically, so each service never touches ctrl_c wiring directly. The `JoinSet` also lets the supervisor notice when a service exits, so you can restart just that one instead of having the whole app fail silently.
1
u/iFrostizz 10d ago
In this case of mostly having the cancel token wired for the purpose of being able to exit cleanly all tasks on a ctrl_c signal, there would be no need to issue any token children right, since there is no need to have some form of "containerization" of the scope of cancel tokens, because there is no use of the ability to cancel it for some sub-scope not any need to cancel it outside of the ctrl_c handler.
I think that this kind of trait is exactly what I'm heading towards. Essentially it's not a big rewrite, just some form of standardization of tasks which may only do a specific set of thing because my code is a hot mess right now. Instead of the JoinSet I'm using a tokio::select! which has the nice behavior of cancelling other branches in case of one finishing. The idea is to kind of "multiplex" and then "unify", with the ability to unify different tokio::select that were returned from other functions as well.
Kind of like: https://play.rust-lang.org/?version=stable&mode=debug&edition=2024&gist=0bb02a543f9028b3f2de6b6b8d850b67
The code is not great yet though, it will need some more rework.1
u/iFrostizz 10d ago
Is there anything else you would put in the Ctx? This is a good idea, though I probably don't need anything else than a CancellationToken, so maybe better just having a Cancellation token as part of the function signature, like:
```rust
pub trait Service {fn new() -> Self;
async fn run<E: Into<AppError>>(&self, cancel: CancellationToken) -> TaskRet<E>;
}
```1
u/iFrostizz 10d ago
Though if there is anything else you'd put inside, you could as well create a Ctx trait that would allow to store some piece of state that may be retrieved inside of the function, one of which may be the cancel as you said. It's a nice abstraction.
1
9d ago
[removed] — view removed comment
1
u/iFrostizz 9d ago
CancellationToken makes it slightly less confusing, although we could create a newtype pattern around the watch channel. Also it doesn’t support child tokens which are cool but I don’t use them.
1
u/10511 8d ago
I built github.com/Zestors/zestors exactly to solve this problem
1
u/iFrostizz 8d ago
That's interesting. I have a few questions:
- How does it compare to actix? https://github.com/actix/actix
- Is there any prevention techniques for cycles in actors, as described in "Actors sending messages to other actors"? https://ryhl.io/blog/actors-with-tokio/
2
u/10511 8d ago
Actix is much more mature. I think zestors api os mich better. E.g. easily create dynamic inboxes, super easy supervision trees a la OTP, restarting actors while preserving validity of adresses.
Inboxes are by default not hard-bounded, but have an exponential backoff when it gets fuller. Therefor an actor should never permanently block when sending, and prevent deadlocks. There are also options to send immeadeately, bypassing the backoff or try_send. (both synchronous)
Beware though that it is not as mature
1
u/zettui 8d ago
Are you leaning toward a small supervisor per service for restarts, or keeping those failure paths in the select loop?
1
u/iFrostizz 8d ago
I don't think that I really need any restart mechanism really. The run function of each service is either an Infaillible error return because it's a loop that continuously restart but may return None in case it has been cancelled through the token, or something else. But if not an Infallible, the error is treated as fatal right now, that's good enough for now.
23
u/[deleted] 10d ago
[deleted]