pub struct Ctx(/* private fields */);Expand description
上下文 = Arc<CtxInner>(所有权模型:Clone 廉价;isolate/plugin 返回共享内核的新 Ctx)。
Implementations§
Source§impl Ctx
impl Ctx
Sourcepub fn instance(&self) -> InstanceId
pub fn instance(&self) -> InstanceId
Identity of the fiber owning this context. Isolated contexts retain it.
Examples found in repository?
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 }Sourcepub fn handle(&self) -> &Handle
pub fn handle(&self) -> &Handle
Examples found in repository?
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 }Sourcepub fn error_sink(&self) -> ErrorSink
pub fn error_sink(&self) -> ErrorSink
错误路由:插件运行期(apply 之外)的异步错误经此上报,不崩 root。 供外部插件在 turn 边界兜底(如 session 落盘失败),与框架内部 路由同一 sink——可观测,不静默。
Sourcepub fn take_cleanup_errors(&self) -> Vec<Arc<CordisError>>
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.
Sourcepub fn root() -> Result<Ctx, CordisError>
pub fn root() -> Result<Ctx, CordisError>
自动路径:Handle::try_current() 失败返回明确错误,绝不隐式建 runtime(D8)。
Examples found in repository?
More examples
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}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}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}Sourcepub fn root_with_sink(handle: Handle, sink: ErrorSink) -> Ctx
pub fn root_with_sink(handle: Handle, sink: ErrorSink) -> Ctx
注入构造 + 自定义 ErrorSink。
Sourcepub fn root_view(&self) -> Option<FiberView>
pub fn root_view(&self) -> Option<FiberView>
root fiber 句柄(root dispose 清子树 / root restart,§五 root_restart)。
最终 shutdown 并释放所有 FiberView 后返回 None。
Sourcepub fn diagnostics(&self) -> RuntimeDiagnostics
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.
Sourcepub fn shutdown(&self) -> BoxFuture<'static, Result<(), Arc<CordisError>>>
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?
More examples
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}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}Sourcepub fn shutdown_with_timeout(
&self,
limit: Duration,
) -> BoxFuture<'static, Result<(), DisposeWaitError>>
pub fn shutdown_with_timeout( &self, limit: Duration, ) -> BoxFuture<'static, Result<(), DisposeWaitError>>
限制等待最终关闭的时间。超时后关闭仍在后台进行,重复调用
shutdown() 可继续 join 同一结果。
Sourcepub fn events(&self) -> &EventBus
pub fn events(&self) -> &EventBus
事件总线(全局唯一;事件分发不跨 isolate 过滤,D29)。
Examples found in repository?
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
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}Sourcepub fn isolate(&self, key: impl Into<TypeKey>, label: &str) -> Ctx
pub fn isolate(&self, key: impl Into<TypeKey>, label: &str) -> Ctx
isolate 作用域(支柱 3):按 ServiceKey 隔离,同 label 合并(TS 语义,D21); 返回的 Ctx 保留原 fiber 所有权(D28)。
Sourcepub fn get<T: Send + Sync + 'static>(&self) -> Option<Arc<T>>
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?
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}Sourcepub fn get_as<T: ?Sized + Send + Sync + 'static>(
&self,
key: impl Into<TypeKey>,
) -> Option<Arc<T>>
pub fn get_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, ) -> Option<Arc<T>>
带显式 key 的读取(多实例,shaku Keyed 模式)。
Sourcepub fn require<T: Send + Sync + 'static>(
&self,
) -> Result<Arc<T>, ServiceReadError>
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.
Sourcepub fn require_as<T: ?Sized + Send + Sync + 'static>(
&self,
key: impl Into<TypeKey>,
) -> Result<Arc<T>, ServiceReadError>
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.
Sourcepub fn provide<T: Send + Sync + 'static>(
&self,
value: T,
) -> Result<Disposer, CordisError>
pub fn provide<T: Send + Sync + 'static>( &self, value: T, ) -> Result<Disposer, CordisError>
值语义注册便捷入口(D13)。
Examples found in repository?
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
Sourcepub fn provide_as<T: ?Sized + Send + Sync + 'static>(
&self,
key: impl Into<TypeKey>,
value: Arc<T>,
) -> Result<Disposer, CordisError>
pub fn provide_as<T: ?Sized + Send + Sync + 'static>( &self, key: impl Into<TypeKey>, value: Arc<T>, ) -> Result<Disposer, CordisError>
trait 对象 / 共享实例注册入口。
Sourcepub 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>
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 门控保留)。
Sourcepub fn provide_mut_as<T: ?Sized + Send + Sync + 'static>(
&self,
key: impl Into<TypeKey>,
value: Arc<T>,
) -> Result<(Disposer, ServiceWriter<T>), CordisError>
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.
Sourcepub fn effect(
&self,
f: impl FnOnce() -> Effect,
) -> Result<Disposer, CordisError>
pub fn effect( &self, f: impl FnOnce() -> Effect, ) -> Result<Disposer, CordisError>
注册清理效应(D23):f 立即执行,返回的清理在卸载时 LIFO 执行。
fiber 已 Disposed/Unloading 时返回 InactiveEffect(§四:重入报错)。
Examples found in repository?
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
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 }Sourcepub fn effect_named(
&self,
label: impl Into<String>,
f: impl FnOnce() -> Effect,
) -> Result<Disposer, CordisError>
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.
Sourcepub fn plugin(&self, p: impl Plugin) -> FiberView
pub fn plugin(&self, p: impl Plugin) -> FiberView
装载插件(支柱 1)。返回 FiberView;级联卸载:child dispose 注册为 parent fiber 的 effect(D28:child plugin 自动归 parent fiber 所有)。
Examples found in repository?
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
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}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}Sourcepub fn plugin_with<C: Send + Sync + 'static>(
&self,
factory: impl PluginFactory<C>,
config: C,
) -> FiberView
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?
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}Sourcepub fn plugin_from<C: Send + Sync + 'static>(
&self,
build: impl Fn(&C) -> Result<Box<dyn Plugin>, CordisError> + Send + Sync + 'static,
config: C,
) -> FiberView
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。
Sourcepub fn cancellation_token(&self) -> CancellationToken
pub fn cancellation_token(&self) -> CancellationToken
当前 fiber 代的取消 token(D27:每代独立 token,卸载第②步取消)。
Examples found in repository?
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§impl Ctx
impl Ctx
Sourcepub 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>
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.
Sourcepub 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>
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.