use std::error::Error;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use rutis::{
BoxFuture, CordisError, Ctx, Effect, FiberState, FiberView, Plugin, PluginFactory, TypeKey,
};
use tokio::sync::{mpsc, oneshot};
type BoxError = Box<dyn Error + Send + Sync>;
struct Backend(String);
#[derive(Default)]
struct Completed(AtomicUsize);
struct Job(String, oneshot::Sender<String>);
struct Indexer(mpsc::Sender<Job>);
impl Indexer {
async fn index(&self, document: &str) -> Result<String, BoxError> {
let (reply, result) = oneshot::channel();
self.0
.send(Job(document.to_owned(), reply))
.await
.map_err(|_| "indexer generation has stopped")?;
result.await.map_err(|_| "indexing was cancelled".into())
}
}
struct BackendPlugin(String);
impl Plugin for BackendPlugin {
fn name(&self) -> &str {
"backend"
}
fn validate(&self) -> Result<(), CordisError> {
if self.0.trim().is_empty() {
return Err(CordisError::Validation {
issues: vec!["backend label must not be empty".into()],
});
}
Ok(())
}
fn apply<'a>(&'a self, ctx: &'a Ctx) -> BoxFuture<'a, Result<Effect, CordisError>> {
Box::pin(async move {
ctx.provide(Backend(self.0.clone()))?;
Ok(Effect::Done)
})
}
}
struct BackendFactory;
impl PluginFactory<String> for BackendFactory {
fn name(&self) -> &str {
"backend"
}
fn build(&self, label: &String) -> Result<Box<dyn Plugin>, CordisError> {
Ok(Box::new(BackendPlugin(label.clone())))
}
}
struct IndexerPlugin(Vec<TypeKey>);
impl Plugin for IndexerPlugin {
fn name(&self) -> &str {
"indexer"
}
fn injects(&self) -> &[TypeKey] {
&self.0
}
fn apply<'a>(&'a self, ctx: &'a Ctx) -> BoxFuture<'a, Result<Effect, CordisError>> {
Box::pin(async move {
let backend = ctx
.get::<Backend>()
.ok_or_else(|| CordisError::ServiceNotFound("Backend".into()))?;
let completed = ctx
.get::<Completed>()
.ok_or_else(|| CordisError::ServiceNotFound("Completed".into()))?;
let (sender, mut receiver) = mpsc::channel::<Job>(8);
let token = ctx.cancellation_token();
let stop = token.clone();
let runtime = ctx.handle().clone();
ctx.effect(move || {
let task = runtime.spawn(async move {
loop {
tokio::select! {
biased;
_ = token.cancelled() => break,
job = receiver.recv() => {
let Some(Job(document, reply)) = job else { break };
let value = format!("{}: {document}", backend.0);
completed.0.fetch_add(1, Ordering::SeqCst);
let _ = reply.send(value);
}
}
}
});
Effect::AsyncDisposer(Box::new(move || {
Box::pin(async move {
stop.cancel();
task.await
.map_err(|error| CordisError::PluginFailed(Box::new(error)))?;
Ok(())
})
}))
})?;
ctx.provide(Indexer(sender))?;
Ok(Effect::Done)
})
}
}
async fn wait_active(view: &FiberView) -> Result<(), BoxError> {
tokio::time::timeout(Duration::from_secs(5), async {
let mut state = view.watch();
loop {
let snapshot = state.borrow_and_update().clone();
match snapshot.state {
FiberState::Active => return Ok::<(), BoxError>(()),
FiberState::Failed | FiberState::Disposed => {
return Err(format!("{} did not start: {snapshot:?}", view.name()).into());
}
_ => state.changed().await?,
}
}
})
.await??;
Ok(())
}
async fn exercise(root: &Ctx) -> Result<(), BoxError> {
root.provide(Completed::default())?;
let indexer = root.plugin(IndexerPlugin(vec![
TypeKey::of::<Backend>(),
TypeKey::of::<Completed>(),
]));
(&indexer).await?;
assert_eq!(indexer.state().state, FiberState::Pending);
let backend = root.plugin_with(BackendFactory, String::from("v1"));
wait_active(&indexer).await?;
let old = root.get::<Indexer>().ok_or("missing Indexer")?;
assert_eq!(old.index("first.txt").await?, "v1: first.txt");
println!("v1 indexed first.txt");
let generation = indexer.state().generation;
assert!(backend.update(String::new()).await.is_err());
assert_eq!(indexer.state().generation, generation);
backend.update(String::from("v2")).await?;
wait_active(&indexer).await?;
assert!(indexer.state().generation > generation);
assert!(old.index("stale.txt").await.is_err());
let current = root.get::<Indexer>().ok_or("missing reloaded Indexer")?;
assert_eq!(current.index("second.txt").await?, "v2: second.txt");
assert_eq!(
root.get::<Completed>()
.ok_or("missing Completed")?
.0
.load(Ordering::SeqCst),
2
);
println!("v2 indexed second.txt; completed = 2; old handle rejected");
Ok(())
}
#[tokio::main]
async fn main() -> Result<(), BoxError> {
let root = Ctx::root()?;
let outcome = exercise(&root).await;
let closed = root.shutdown().await;
outcome?;
closed?;
println!("shutdown complete");
Ok(())
}