tuwunel_api/client/admin/misc/
scheduled_tasks.rs1use axum::extract::State;
2use itertools::Itertools;
3use ruma::{MilliSecondsSinceUnixEpoch, UInt};
4use synapse_admin_api::scheduled_tasks::list::v1::{
5 Request, Response, ScheduledTask, TaskStatus,
6};
7use tuwunel_core::Result;
8use tuwunel_service::tasks::{Status, TaskInfo};
9
10use crate::{Ruma, client::admin::require_admin};
11
12pub(crate) async fn admin_scheduled_tasks_route(
17 State(services): State<crate::State>,
18 body: Ruma<Request>,
19) -> Result<Response> {
20 require_admin(&services, body.sender_user()).await?;
21
22 let scheduled_tasks = services
23 .tasks
24 .list()
25 .into_iter()
26 .filter(|task| matches_filters(task, &body))
27 .sorted_by_key(|task| task.timestamp_ms)
28 .map(scheduled_task)
29 .collect();
30
31 Ok(Response { scheduled_tasks })
32}
33
34fn matches_filters(task: &TaskInfo, req: &Request) -> bool {
35 req.action_name
36 .as_deref()
37 .is_none_or(|name| task.action == name)
38 && req
39 .resource_id
40 .as_deref()
41 .is_none_or(|id| task.resource_id.as_str() == id)
42 && req
43 .status
44 .as_ref()
45 .is_none_or(|status| task_status(task.status) == *status)
46 && req
47 .max_timestamp
48 .is_none_or(|max| task.timestamp_ms <= u64::from(max))
49}
50
51fn scheduled_task(task: TaskInfo) -> ScheduledTask {
52 ScheduledTask {
53 id: task.id.to_string(),
54 action: task.action.to_owned(),
55 status: task_status(task.status),
56 timestamp_ms: MilliSecondsSinceUnixEpoch(
57 UInt::try_from(task.timestamp_ms).unwrap_or_default(),
58 ),
59 resource_id: Some(task.resource_id),
60 result: task
61 .result
62 .and_then(|value| serde_json::from_value(value).ok()),
63 error: task.error,
64 }
65}
66
67fn task_status(status: Status) -> TaskStatus {
68 match status {
69 | Status::Scheduled => TaskStatus::Scheduled,
70 | Status::Active => TaskStatus::Active,
71 | Status::Complete => TaskStatus::Complete,
72 | Status::Failed => TaskStatus::Failed,
73 }
74}
75
76#[cfg(test)]
77mod tests {
78 use serde_json::json;
79 use tuwunel_service::tasks::Status;
80
81 use super::task_status;
82
83 #[test]
84 fn maps_each_tracked_status() {
85 let status = |status| serde_json::to_value(task_status(status)).unwrap();
86
87 assert_eq!(status(Status::Scheduled), json!("scheduled"));
88 assert_eq!(status(Status::Active), json!("active"));
89 assert_eq!(status(Status::Complete), json!("complete"));
90 assert_eq!(status(Status::Failed), json!("failed"));
91 }
92}