use std::sync::Arc;
use parking_lot::Mutex;
use wabot_addon_async_in_memory::InMemoryJobRepository;
use wabot_core::injection::Container;
use wabot_feature_async::{
job, register_async_runtime, register_job_repository, AsyncError, CommandData,
CommandHandlerEntry, CommandRegistry, CronHandlerEntry, Job, JobRepository, JobRunner,
};
pub struct AsyncHarness {
container: Container,
repository: Arc<InMemoryJobRepository>,
registry: Arc<CommandRegistry>,
runner: Arc<JobRunner>,
ran: Mutex<Vec<String>>,
}
impl AsyncHarness {
pub fn builder() -> AsyncHarnessBuilder {
AsyncHarnessBuilder {
container: None,
commands: Vec::new(),
crons: Vec::new(),
}
}
pub fn container(&self) -> &Container {
&self.container
}
pub async fn execute<C: CommandData + serde::Serialize>(&self, command: &C) -> FinishedJob {
let payload = serde_json::to_value(command).expect("a serializable command");
self.execute_named(C::COMMAND_NAME, payload).await
}
pub async fn execute_named(
&self,
command_name: &str,
payload: serde_json::Value,
) -> FinishedJob {
assert!(
self.registry.get(command_name).is_some(),
"AsyncHarness: no handler registered for command '{command_name}'. \
Registered: {:?}",
self.registry.command_names()
);
let job = job::new_job(job::JobData {
base: Default::default(),
command_name: command_name.to_string(),
command_data: payload,
scheduled_at: Some(chrono::Utc::now().timestamp_millis()),
started_at: None,
success_at: None,
failed_at: None,
retry_delays_seconds: self
.registry
.options_for(command_name)
.and_then(|o| o.retry_delays_seconds),
intent_number: None,
error: None,
acceptable_running_time_seconds: None,
stuck_retry_attempts: None,
dedup_key: None,
actor: wabot_core::audit::audit_actor(),
request_id: wabot_core::log_context::request_id(),
});
self.repository.create(&job).await.expect("stored");
self.ran.lock().push(job.id().to_string());
let result = self.runner.run(self.container.clone(), job.clone()).await;
let stored = self
.repository
.find(job.id())
.await
.expect("a readable repository")
.expect("the job it just ran");
FinishedJob {
job: stored,
run_error: result.err(),
}
}
pub async fn run_cron(&self, command_name: &str) -> FinishedJob {
self.execute_named(command_name, serde_json::Value::Null)
.await
}
pub async fn jobs(&self) -> Vec<Job> {
let ids = self.ran.lock().clone();
let mut jobs = Vec::with_capacity(ids.len());
for id in ids {
if let Ok(Some(job)) = self.repository.find(&id).await {
jobs.push(job);
}
}
jobs
}
}
pub struct FinishedJob {
pub job: Job,
pub run_error: Option<AsyncError>,
}
impl FinishedJob {
pub fn succeeded(&self) -> bool {
job::was_success(&self.job)
}
pub fn error(&self) -> Option<String> {
self.job.data().error.as_ref().map(|e| e.message.clone())
}
pub fn attempts(&self) -> u32 {
self.job.data().intent_number.unwrap_or(0)
}
pub fn retry_at_ms(&self) -> Option<i64> {
self.job
.data()
.failed_at
.is_none()
.then(|| self.job.data().scheduled_at)
.flatten()
}
pub fn assert_succeeded(&self) -> &Self {
assert!(
self.succeeded(),
"expected the job to succeed, but it {}",
match self.error() {
Some(message) => format!("failed: {message}"),
None => "did not finish".to_string(),
}
);
self
}
}
pub struct AsyncHarnessBuilder {
container: Option<Container>,
commands: Vec<CommandHandlerEntry>,
crons: Vec<CronHandlerEntry>,
}
impl AsyncHarnessBuilder {
pub fn container(mut self, container: Container) -> Self {
self.container = Some(container);
self
}
pub fn command(mut self, entry: CommandHandlerEntry) -> Self {
self.commands.push(entry);
self
}
pub fn cron(mut self, entry: CronHandlerEntry) -> Self {
self.crons.push(entry);
self
}
pub fn build(self) -> AsyncHarness {
let container = self.container.unwrap_or_default();
let repository = Arc::new(InMemoryJobRepository::new());
register_job_repository(&container, repository.clone());
register_async_runtime(&container);
let registry: Arc<CommandRegistry> = container.resolve();
for entry in self.commands {
registry.register(entry);
}
for entry in self.crons {
registry.register(wabot_feature_async::cron_command_entry(&entry));
}
let runner = Arc::new(JobRunner::new(repository.clone(), registry.clone()));
AsyncHarness {
container,
repository,
registry,
runner,
ran: Mutex::new(Vec::new()),
}
}
}