Skip to main content

development_workflow/
development_workflow.rs

1//! 开发指南配套:配置更新 → 消费者重载 → 旧队列停止 → 跨代计数保留。
2//! Run: cargo run -p rutis --example development_workflow
3
4use 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
15// 契约:Backend 是本代实现;Completed 是独立于 worker 代际的应用状态。
16struct Backend(String);
17#[derive(Default)]
18struct Completed(AtomicUsize);
19struct Job(String, oneshot::Sender<String>);
20
21// 一个真实的有界队列服务。旧代接收端关闭后,旧句柄明确返回错误。
22struct 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        // 纯构造:预检查和正式装载都会调用这里。
65        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            // 在 effect 工厂内启动,启动后立即交回清理责任。
93            // 先登记 worker,再提供服务:清理时先撤销服务,再 join worker。
94            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                                // 示例用字符串处理代替真正的异步索引 I/O。
103                                let value = format!("{}: {document}", backend.0);
104                                completed.0.fetch_add(1, Ordering::SeqCst);
105                                let _ = reply.send(value);
106                            }
107                        }
108                    }
109                    // 本例选择取消未处理请求:drop receiver 及排队的 reply。
110                });
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
126// 等待 Active,并处理 Failed / Disposed / 观察通道关闭 / 超时。
127// Pending 也能 settle,所以不能只使用 (&view).await 代替此检查。
128async 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 级状态;Indexer 重载不清空它。
148    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    // 预检查失败:现有实例和代数保持不变。
163    let generation = indexer.state().generation;
164    assert!(backend.update(String::new()).await.is_err());
165    assert_eq!(indexer.state().generation, generation);
166
167    // 正式更新只操作 provider;框架会重载声明依赖的消费者。
168    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    // 即使演示中的业务步骤失败,也先尝试关闭 root。
190    let closed = root.shutdown().await;
191    outcome?;
192    closed?;
193    println!("shutdown complete");
194    Ok(())
195}