development_workflow/
development_workflow.rs1use std::error::Error;
5use std::sync::atomic::{AtomicUsize, Ordering};
6use std::time::Duration;
7
8use rutis::{
9 BoxFuture, CordisError, Ctx, Effect, FiberState, FiberView, Plugin, PluginFactory, TypeKey,
10};
11use tokio::sync::{mpsc, oneshot};
12
13type BoxError = Box<dyn Error + Send + Sync>;
14
15struct Backend(String);
17#[derive(Default)]
18struct Completed(AtomicUsize);
19struct Job(String, oneshot::Sender<String>);
20
21struct Indexer(mpsc::Sender<Job>);
23impl Indexer {
24 async fn index(&self, document: &str) -> Result<String, BoxError> {
25 let (reply, result) = oneshot::channel();
26 self.0
27 .send(Job(document.to_owned(), reply))
28 .await
29 .map_err(|_| "indexer generation has stopped")?;
30 result.await.map_err(|_| "indexing was cancelled".into())
31 }
32}
33
34struct BackendPlugin(String);
35impl Plugin for BackendPlugin {
36 fn name(&self) -> &str {
37 "backend"
38 }
39
40 fn validate(&self) -> Result<(), CordisError> {
41 if self.0.trim().is_empty() {
42 return Err(CordisError::Validation {
43 issues: vec!["backend label must not be empty".into()],
44 });
45 }
46 Ok(())
47 }
48
49 fn apply<'a>(&'a self, ctx: &'a Ctx) -> BoxFuture<'a, Result<Effect, CordisError>> {
50 Box::pin(async move {
51 ctx.provide(Backend(self.0.clone()))?;
52 Ok(Effect::Done)
53 })
54 }
55}
56
57struct BackendFactory;
58impl PluginFactory<String> for BackendFactory {
59 fn name(&self) -> &str {
60 "backend"
61 }
62
63 fn build(&self, label: &String) -> Result<Box<dyn Plugin>, CordisError> {
64 Ok(Box::new(BackendPlugin(label.clone())))
66 }
67}
68
69struct IndexerPlugin(Vec<TypeKey>);
70impl Plugin for IndexerPlugin {
71 fn name(&self) -> &str {
72 "indexer"
73 }
74
75 fn injects(&self) -> &[TypeKey] {
76 &self.0
77 }
78
79 fn apply<'a>(&'a self, ctx: &'a Ctx) -> BoxFuture<'a, Result<Effect, CordisError>> {
80 Box::pin(async move {
81 let backend = ctx
82 .get::<Backend>()
83 .ok_or_else(|| CordisError::ServiceNotFound("Backend".into()))?;
84 let completed = ctx
85 .get::<Completed>()
86 .ok_or_else(|| CordisError::ServiceNotFound("Completed".into()))?;
87 let (sender, mut receiver) = mpsc::channel::<Job>(8);
88 let token = ctx.cancellation_token();
89 let stop = token.clone();
90 let runtime = ctx.handle().clone();
91
92 ctx.effect(move || {
95 let task = runtime.spawn(async move {
96 loop {
97 tokio::select! {
98 biased;
99 _ = token.cancelled() => break,
100 job = receiver.recv() => {
101 let Some(Job(document, reply)) = job else { break };
102 let value = format!("{}: {document}", backend.0);
104 completed.0.fetch_add(1, Ordering::SeqCst);
105 let _ = reply.send(value);
106 }
107 }
108 }
109 });
111 Effect::AsyncDisposer(Box::new(move || {
112 Box::pin(async move {
113 stop.cancel();
114 task.await
115 .map_err(|error| CordisError::PluginFailed(Box::new(error)))?;
116 Ok(())
117 })
118 }))
119 })?;
120 ctx.provide(Indexer(sender))?;
121 Ok(Effect::Done)
122 })
123 }
124}
125
126async fn wait_active(view: &FiberView) -> Result<(), BoxError> {
129 tokio::time::timeout(Duration::from_secs(5), async {
130 let mut state = view.watch();
131 loop {
132 let snapshot = state.borrow_and_update().clone();
133 match snapshot.state {
134 FiberState::Active => return Ok::<(), BoxError>(()),
135 FiberState::Failed | FiberState::Disposed => {
136 return Err(format!("{} did not start: {snapshot:?}", view.name()).into());
137 }
138 _ => state.changed().await?,
139 }
140 }
141 })
142 .await??;
143 Ok(())
144}
145
146async fn exercise(root: &Ctx) -> Result<(), BoxError> {
147 root.provide(Completed::default())?;
149 let indexer = root.plugin(IndexerPlugin(vec![
150 TypeKey::of::<Backend>(),
151 TypeKey::of::<Completed>(),
152 ]));
153 (&indexer).await?;
154 assert_eq!(indexer.state().state, FiberState::Pending);
155
156 let backend = root.plugin_with(BackendFactory, String::from("v1"));
157 wait_active(&indexer).await?;
158 let old = root.get::<Indexer>().ok_or("missing Indexer")?;
159 assert_eq!(old.index("first.txt").await?, "v1: first.txt");
160 println!("v1 indexed first.txt");
161
162 let generation = indexer.state().generation;
164 assert!(backend.update(String::new()).await.is_err());
165 assert_eq!(indexer.state().generation, generation);
166
167 backend.update(String::from("v2")).await?;
169 wait_active(&indexer).await?;
170 assert!(indexer.state().generation > generation);
171 assert!(old.index("stale.txt").await.is_err());
172 let current = root.get::<Indexer>().ok_or("missing reloaded Indexer")?;
173 assert_eq!(current.index("second.txt").await?, "v2: second.txt");
174 assert_eq!(
175 root.get::<Completed>()
176 .ok_or("missing Completed")?
177 .0
178 .load(Ordering::SeqCst),
179 2
180 );
181 println!("v2 indexed second.txt; completed = 2; old handle rejected");
182 Ok(())
183}
184
185#[tokio::main]
186async fn main() -> Result<(), BoxError> {
187 let root = Ctx::root()?;
188 let outcome = exercise(&root).await;
189 let closed = root.shutdown().await;
191 outcome?;
192 closed?;
193 println!("shutdown complete");
194 Ok(())
195}