Skip to main content

Ctx

Struct Ctx 

Source
pub struct Ctx(/* private fields */);
Expand description

上下文 = Arc<CtxInner>(所有权模型:Clone 廉价;isolate/plugin 返回共享内核的新 Ctx)。

Implementations§

Source§

impl Ctx

Source

pub fn instance(&self) -> InstanceId

Identity of the fiber owning this context. Isolated contexts retain it.

Examples found in repository?
examples/listener_ctx_ownership.rs (line 28)
22    fn call<'a>(
23        &'a self,
24        emitter: &'a Ctx,
25        _event: &'a Tick,
26    ) -> BoxFuture<'a, Result<Option<()>, CordisError>> {
27        Box::pin(async move {
28            assert_ne!(emitter.instance(), self.owner.instance());
29            let cleaned = self.cleaned.clone();
30            self.owner.effect(move || {
31                Effect::Disposer(Box::new(move || {
32                    cleaned.fetch_add(1, Ordering::SeqCst);
33                    Ok(())
34                }))
35            })?;
36            Ok(None)
37        })
38    }
Source

pub fn handle(&self) -> &Handle

Examples found in repository?
examples/development_workflow.rs (line 90)
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    }
Source

pub fn error_sink(&self) -> ErrorSink

错误路由:插件运行期(apply 之外)的异步错误经此上报,不崩 root。 供外部插件在 turn 边界兜底(如 session 落盘失败),与框架内部 路由同一 sink——可观测,不静默。

Source

pub fn take_cleanup_errors(&self) -> Vec<Arc<CordisError>>

Take completed, early-disposed effect errors owned by this fiber. Disposer::dispose() returns an error immediately without sending it to the error sink. Each error is returned here once; taken errors no longer appear in a later unload result or restart sink notification. Errors from effects still draining remain available to a subsequent call or to the unload that owns their completion.

Source

pub fn root() -> Result<Ctx, CordisError>

自动路径:Handle::try_current() 失败返回明确错误,绝不隐式建 runtime(D8)。

Examples found in repository?
examples/development_workflow.rs (line 187)
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}
More examples
Hide additional examples
examples/listener_ctx_ownership.rs (line 65)
64async fn main() -> Result<(), Box<dyn std::error::Error>> {
65    let root = Ctx::root()?;
66    let cleaned = Arc::new(AtomicUsize::new(0));
67    let plugin = root.plugin(Register(cleaned.clone()));
68    (&plugin).await?;
69    root.events()
70        .serial(&root, &rutis::EventKey::of(), &Tick)
71        .await?;
72    plugin.shutdown().await?;
73    assert_eq!(cleaned.load(Ordering::SeqCst), 1);
74    root.shutdown().await?;
75    Ok(())
76}
examples/sync_decision.rs (line 37)
36async fn main() -> Result<(), Box<dyn std::error::Error>> {
37    let root = Ctx::root()?;
38    root.events().on_waterfall_sync(
39        &root,
40        &BEFORE_SAVE,
41        |_: &Ctx, _: &BeforeSave, next: SyncNext<'_, BeforeSave>| {
42            // Fixed input: rewrite the downstream result as it returns.
43            Ok(next.call()?.min(100))
44        },
45    )?;
46    let state = Mutex::new(10);
47    assert_eq!(save(&root, &state, 120)?, 100);
48    println!("saved {}", state.lock().unwrap());
49    root.shutdown().await?;
50    Ok(())
51}
examples/quickstart.rs (line 73)
72async fn main() -> Result<(), Box<dyn std::error::Error>> {
73    let ctx = Ctx::root()?;
74
75    // 1. 先装 consumer:依赖未到,它停在 Pending(依赖门控)。
76    let listener = ctx.plugin(Listener {
77        deps: vec![TypeKey::of::<Greeting>()],
78    });
79
80    // 2. provider 到位 → 门控放行,consumer 自动装载。
81    let v1 = ctx.plugin(Greeter { version: 1 });
82    wait_active(&listener).await; // apply 已打印第一行
83
84    // 3. 换 provider:旧的卸载、服务摘除,consumer 被驱逐回 Pending;
85    //    新 provider 提供同类型服务,consumer 自动重载——打印第二行。
86    //    全程没有碰过 consumer。
87    v1.dispose().await?;
88    let _v2 = ctx.plugin(Greeter { version: 2 });
89    wait_active(&listener).await;
90    println!("done: consumer reloaded without touching it");
91    Ok(())
92}
Source

pub fn root_with(handle: Handle) -> Ctx

注入构造(优先路径,D8)。

Source

pub fn root_with_sink(handle: Handle, sink: ErrorSink) -> Ctx

注入构造 + 自定义 ErrorSink。

Source

pub fn root_view(&self) -> Option<FiberView>

root fiber 句柄(root dispose 清子树 / root restart,§五 root_restart)。 最终 shutdown 并释放所有 FiberView 后返回 None。

Source

pub fn diagnostics(&self) -> RuntimeDiagnostics

Read-only, best-effort snapshot of live fibers and services.

This scans fibers and bindings under separate locks. Concurrent lifecycle changes may therefore mix states or generations from different moments, even within one plugin’s state, dependency, and binding entries. The result is useful for diagnosis, not an atomic transaction or a change stream. Dependency checks are never called here: their status is the last result recorded by normal gate resolution, or CheckPending if that binding has not been checked yet. Reading does not call plugin name, injects, or other user code and does not advance a fiber.

Source

pub fn shutdown(&self) -> BoxFuture<'static, Result<(), Arc<CordisError>>>

最终关闭 root。重复或并发调用共享同一完成结果;现有 dispose 仍可 restart。关闭会拒绝新注册并等待准入关闭前完成挂载的子树 及其清理。创建已开始但挂载被关闭拒绝的子 fiber 不会进入 apply, 会独立关闭;其返回的 FiberView::shutdown() 可等待该 driver 退出。 root 的完成结果不包含这种未挂载子 fiber。 不协作的代码可能使等待无限延长,可用 shutdown_with_timeout 限制等待。

Examples found in repository?
examples/development_workflow.rs (line 190)
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}
More examples
Hide additional examples
examples/listener_ctx_ownership.rs (line 74)
64async fn main() -> Result<(), Box<dyn std::error::Error>> {
65    let root = Ctx::root()?;
66    let cleaned = Arc::new(AtomicUsize::new(0));
67    let plugin = root.plugin(Register(cleaned.clone()));
68    (&plugin).await?;
69    root.events()
70        .serial(&root, &rutis::EventKey::of(), &Tick)
71        .await?;
72    plugin.shutdown().await?;
73    assert_eq!(cleaned.load(Ordering::SeqCst), 1);
74    root.shutdown().await?;
75    Ok(())
76}
examples/sync_decision.rs (line 49)
36async fn main() -> Result<(), Box<dyn std::error::Error>> {
37    let root = Ctx::root()?;
38    root.events().on_waterfall_sync(
39        &root,
40        &BEFORE_SAVE,
41        |_: &Ctx, _: &BeforeSave, next: SyncNext<'_, BeforeSave>| {
42            // Fixed input: rewrite the downstream result as it returns.
43            Ok(next.call()?.min(100))
44        },
45    )?;
46    let state = Mutex::new(10);
47    assert_eq!(save(&root, &state, 120)?, 100);
48    println!("saved {}", state.lock().unwrap());
49    root.shutdown().await?;
50    Ok(())
51}
Source

pub fn shutdown_with_timeout( &self, limit: Duration, ) -> BoxFuture<'static, Result<(), DisposeWaitError>>

限制等待最终关闭的时间。超时后关闭仍在后台进行,重复调用 shutdown() 可继续 join 同一结果。

Source

pub fn events(&self) -> &EventBus

事件总线(全局唯一;事件分发不跨 isolate 过滤,D29)。

Examples found in repository?
examples/listener_ctx_ownership.rs (line 50)
48    fn apply<'a>(&'a self, ctx: &'a Ctx) -> BoxFuture<'a, Result<Effect, CordisError>> {
49        Box::pin(async move {
50            ctx.events().on::<Tick>(
51                ctx,
52                &rutis::EventKey::of(),
53                OwnedListener {
54                    owner: ctx.clone(),
55                    cleaned: self.0.clone(),
56                },
57            )?;
58            Ok(Effect::Done)
59        })
60    }
61}
62
63#[tokio::main]
64async fn main() -> Result<(), Box<dyn std::error::Error>> {
65    let root = Ctx::root()?;
66    let cleaned = Arc::new(AtomicUsize::new(0));
67    let plugin = root.plugin(Register(cleaned.clone()));
68    (&plugin).await?;
69    root.events()
70        .serial(&root, &rutis::EventKey::of(), &Tick)
71        .await?;
72    plugin.shutdown().await?;
73    assert_eq!(cleaned.load(Ordering::SeqCst), 1);
74    root.shutdown().await?;
75    Ok(())
76}
More examples
Hide additional examples
examples/sync_decision.rs (line 21)
18fn save(ctx: &Ctx, state: &Mutex<u64>, candidate: u64) -> Result<u64, CordisError> {
19    let mut current = state.lock().unwrap();
20    let proposed =
21        ctx.events()
22            .waterfall_sync(ctx, &BEFORE_SAVE, &BeforeSave { candidate }, |_, event| {
23                Ok((*current).max(event.candidate))
24            })?;
25    // Middleware only proposes a value; the caller validates before writing.
26    if proposed < *current {
27        return Err(CordisError::Validation {
28            issues: vec!["saved value must not decrease".into()],
29        });
30    }
31    *current = proposed;
32    Ok(proposed)
33}
34
35#[tokio::main]
36async fn main() -> Result<(), Box<dyn std::error::Error>> {
37    let root = Ctx::root()?;
38    root.events().on_waterfall_sync(
39        &root,
40        &BEFORE_SAVE,
41        |_: &Ctx, _: &BeforeSave, next: SyncNext<'_, BeforeSave>| {
42            // Fixed input: rewrite the downstream result as it returns.
43            Ok(next.call()?.min(100))
44        },
45    )?;
46    let state = Mutex::new(10);
47    assert_eq!(save(&root, &state, 120)?, 100);
48    println!("saved {}", state.lock().unwrap());
49    root.shutdown().await?;
50    Ok(())
51}
Source

pub fn isolate(&self, key: impl Into<TypeKey>, label: &str) -> Ctx

isolate 作用域(支柱 3):按 ServiceKey 隔离,同 label 合并(TS 语义,D21); 返回的 Ctx 保留原 fiber 所有权(D28)。

Source

pub fn get<T: Send + Sync + 'static>(&self) -> Option<Arc<T>>

类型键读取(显式定位器,D13):沿父链解析作用域; provider 非 Active 时不可见,但其子树内自访问除外(清理期自访问,§四)。 访问方自身失活(Unloading/Disposed)时同样不可见——TS inactive context 语义(reflect.spec ‘service inject leak’ 的语言无关内核;provider 子树 内自访问豁免,与清理期自访问同一条规则)。

Examples found in repository?
examples/development_workflow.rs (line 82)
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}
Source

pub fn get_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, ) -> Option<Arc<T>>

带显式 key 的读取(多实例,shaku Keyed 模式)。

Source

pub fn require<T: Send + Sync + 'static>( &self, ) -> Result<Arc<T>, ServiceReadError>

Read a declared service, or return a call-site error. Unlike get, this enforces the caller’s dependency declaration at every invocation.

Examples found in repository?
examples/quickstart.rs (line 52)
50    fn apply<'a>(&'a self, ctx: &'a Ctx) -> BoxFuture<'a, Result<Effect, CordisError>> {
51        Box::pin(async move {
52            let greeting = ctx.require::<Greeting>()?.0.clone();
53            println!("[listener] loaded: {greeting}");
54            Ok(Effect::Done)
55        })
56    }
Source

pub fn require_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, ) -> Result<Arc<T>, ServiceReadError>

Strict keyed read. A declaration on an ancestor fiber is usable only when that ancestor sees the same isolate scope as this context.

Source

pub fn provide<T: Send + Sync + 'static>( &self, value: T, ) -> Result<Disposer, CordisError>

值语义注册便捷入口(D13)。

Examples found in repository?
examples/development_workflow.rs (line 51)
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}
More examples
Hide additional examples
examples/quickstart.rs (line 32)
29    fn apply<'a>(&'a self, ctx: &'a Ctx) -> BoxFuture<'a, Result<Effect, CordisError>> {
30        let greeting = Greeting(format!("hello from greeter v{}", self.version));
31        Box::pin(async move {
32            ctx.provide(greeting)?;
33            Ok(Effect::Done)
34        })
35    }
Source

pub fn provide_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, value: Arc<T>, ) -> Result<Disposer, CordisError>

trait 对象 / 共享实例注册入口。

Source

pub fn provide_as_with_check<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, value: Arc<T>, check: impl Fn() -> bool + Send + Sync + 'static, ) -> Result<Disposer, CordisError>

带 check() 谓词的注册(§四:check 门控保留)。

Source

pub fn provide_mut_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, value: Arc<T>, ) -> Result<(Disposer, ServiceWriter<T>), CordisError>

Register a provider-owned mutable binding and return a generation-bound writer. A write replaces the registered Arc, not existing Arc snapshots.

Source

pub fn effect( &self, f: impl FnOnce() -> Effect, ) -> Result<Disposer, CordisError>

注册清理效应(D23):f 立即执行,返回的清理在卸载时 LIFO 执行。 fiber 已 Disposed/Unloading 时返回 InactiveEffect(§四:重入报错)。

Examples found in repository?
examples/listener_ctx_ownership.rs (lines 30-35)
22    fn call<'a>(
23        &'a self,
24        emitter: &'a Ctx,
25        _event: &'a Tick,
26    ) -> BoxFuture<'a, Result<Option<()>, CordisError>> {
27        Box::pin(async move {
28            assert_ne!(emitter.instance(), self.owner.instance());
29            let cleaned = self.cleaned.clone();
30            self.owner.effect(move || {
31                Effect::Disposer(Box::new(move || {
32                    cleaned.fetch_add(1, Ordering::SeqCst);
33                    Ok(())
34                }))
35            })?;
36            Ok(None)
37        })
38    }
More examples
Hide additional examples
examples/development_workflow.rs (lines 94-119)
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    }
Source

pub fn effect_named( &self, label: impl Into<String>, f: impl FnOnce() -> Effect, ) -> Result<Disposer, CordisError>

Register a cleanup with a label visible through FiberView::effects.

Source

pub fn plugin(&self, p: impl Plugin) -> FiberView

装载插件(支柱 1)。返回 FiberView;级联卸载:child dispose 注册为 parent fiber 的 effect(D28:child plugin 自动归 parent fiber 所有)。

Examples found in repository?
examples/listener_ctx_ownership.rs (line 67)
64async fn main() -> Result<(), Box<dyn std::error::Error>> {
65    let root = Ctx::root()?;
66    let cleaned = Arc::new(AtomicUsize::new(0));
67    let plugin = root.plugin(Register(cleaned.clone()));
68    (&plugin).await?;
69    root.events()
70        .serial(&root, &rutis::EventKey::of(), &Tick)
71        .await?;
72    plugin.shutdown().await?;
73    assert_eq!(cleaned.load(Ordering::SeqCst), 1);
74    root.shutdown().await?;
75    Ok(())
76}
More examples
Hide additional examples
examples/quickstart.rs (lines 76-78)
72async fn main() -> Result<(), Box<dyn std::error::Error>> {
73    let ctx = Ctx::root()?;
74
75    // 1. 先装 consumer:依赖未到,它停在 Pending(依赖门控)。
76    let listener = ctx.plugin(Listener {
77        deps: vec![TypeKey::of::<Greeting>()],
78    });
79
80    // 2. provider 到位 → 门控放行,consumer 自动装载。
81    let v1 = ctx.plugin(Greeter { version: 1 });
82    wait_active(&listener).await; // apply 已打印第一行
83
84    // 3. 换 provider:旧的卸载、服务摘除,consumer 被驱逐回 Pending;
85    //    新 provider 提供同类型服务,consumer 自动重载——打印第二行。
86    //    全程没有碰过 consumer。
87    v1.dispose().await?;
88    let _v2 = ctx.plugin(Greeter { version: 2 });
89    wait_active(&listener).await;
90    println!("done: consumer reloaded without touching it");
91    Ok(())
92}
examples/development_workflow.rs (lines 149-152)
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}
Source

pub fn plugin_with<C: Send + Sync + 'static>( &self, factory: impl PluginFactory<C>, config: C, ) -> FiberView

工厂模式装载(D32:配置热更新):依赖门控声明在注册时取自工厂并固定, 每代装载用当前 config 构造实例;返回的 FiberView 可 update(config)。

Examples found in repository?
examples/development_workflow.rs (line 156)
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}
Source

pub fn plugin_from<C: Send + Sync + 'static>( &self, build: impl Fn(&C) -> Result<Box<dyn Plugin>, CordisError> + Send + Sync + 'static, config: C, ) -> FiberView

工厂模式装载的闭包便捷形态(D32):零依赖声明的单方法工厂。 需要声明 injects/validate_config 时实现 PluginFactory。

Source

pub fn cancellation_token(&self) -> CancellationToken

当前 fiber 代的取消 token(D27:每代独立 token,卸载第②步取消)。

Examples found in repository?
examples/development_workflow.rs (line 88)
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    }
Source

pub fn cancelled(&self) -> impl Future<Output = ()> + Send + 'static

等待当前 fiber 代被取消(协作取消;不观察则 dispose 无限等待,D27 限制)。

Source

pub fn refresh(&self)

触发依赖重查(check() 谓词结果变更等场景)。

Source§

impl Ctx

Source

pub fn intercept_require_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, hook: impl Fn(Arc<T>) -> ServiceIntercept<T> + Send + Sync + 'static, ) -> Result<Disposer, CordisError>

Intercept only strict reads after existing declaration, visibility and readiness checks. Explicit get_as remains an unhooked locator.

Source

pub fn intercept_set_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, hook: impl Fn(Arc<T>) -> ServiceIntercept<T> + Send + Sync + 'static, ) -> Result<Disposer, CordisError>

Intercept a provider’s update of a binding made with provide_mut_as.

Trait Implementations§

Source§

impl Clone for Ctx

Source§

fn clone(&self) -> Self

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more

Auto Trait Implementations§

§

impl !RefUnwindSafe for Ctx

§

impl !UnwindSafe for Ctx

§

impl Freeze for Ctx

§

impl Send for Ctx

§

impl Sync for Ctx

§

impl Unpin for Ctx

§

impl UnsafeUnpin for Ctx

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.