use restate_sdk::prelude::*;
use std::time::Duration;
struct PeriodicTask;
const ACTIVE: &str = "active";
#[object]
impl PeriodicTask {
#[handler]
async fn start(&self, context: ObjectContext<'_>) -> Result<(), TerminalError> {
if context
.get::<bool>(ACTIVE)
.await?
.is_some_and(|enabled| enabled)
{
return Ok(());
}
PeriodicTask::schedule_next(&context);
context.set(ACTIVE, true);
Ok(())
}
#[handler]
async fn stop(&self, context: ObjectContext<'_>) -> Result<(), TerminalError> {
context.clear(ACTIVE);
Ok(())
}
#[handler]
async fn run(&self, context: ObjectContext<'_>) -> Result<(), TerminalError> {
if context.get::<bool>(ACTIVE).await?.is_none() {
return Ok(());
}
println!("Triggered the periodic task!");
PeriodicTask::schedule_next(&context);
Ok(())
}
}
impl PeriodicTask {
fn schedule_next(context: &ObjectContext<'_>) {
context
.object_client::<PeriodicTaskClient>(context.key())
.run()
.send_after(Duration::from_secs(10));
}
}
#[tokio::main]
async fn main() {
tracing_subscriber::fmt::init();
HttpServer::new(Endpoint::builder().bind(PeriodicTask).build())
.listen_and_serve("0.0.0.0:9080".parse().unwrap())
.await;
}