apalis_core/backend/
expose.rs1use std::str::FromStr;
2
3use crate::{
4 backend::{Backend, TaskSink, WireFormatBackend},
5 task::{Task, status::Status},
6};
7
8const DEFAULT_PAGE_SIZE: u32 = 10;
9pub trait Expose<Args, Kind> {}
11
12impl<B, Args, Kind> Expose<Args, Kind> for B where
13 B: Backend
14 + Metrics
15 + ListWorkers
16 + ListQueues
17 + ListAllTasks
18 + ListTasks
19 + TaskSink<Args, Kind>
20{
21}
22
23pub trait ListQueues: Backend {
25 fn list_queues(&self) -> impl Future<Output = Result<Vec<QueueInfo>, Self::Error>> + Send;
27}
28
29pub trait ListWorkers: Backend {
31 fn list_workers(&self) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send;
33
34 fn list_all_workers(
36 &self,
37 ) -> impl Future<Output = Result<Vec<RunningWorker>, Self::Error>> + Send;
38}
39pub trait ListTasks: WireFormatBackend + Backend {
41 #[allow(clippy::type_complexity)]
43 fn list_tasks(
44 &self,
45 filter: &Filter,
46 ) -> impl Future<Output = Result<Vec<Task<Self::Compact>>, Self::Error>> + Send;
47}
48
49pub trait ListAllTasks: WireFormatBackend + Backend {
51 #[allow(clippy::type_complexity)]
53 fn list_all_tasks(
54 &self,
55 filter: &Filter,
56 ) -> impl Future<Output = Result<Vec<Task<Self::Compact>>, Self::Error>> + Send;
57}
58
59pub trait Metrics: Backend {
61 fn global(&self) -> impl Future<Output = Result<Vec<Statistic>, Self::Error>> + Send;
63
64 fn fetch_by_queue(&self) -> impl Future<Output = Result<Vec<Statistic>, Self::Error>> + Send;
66}
67
68#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
70#[derive(Debug, Clone)]
71pub struct QueueInfo {
72 pub name: String,
74 pub stats: Vec<Statistic>,
76 pub workers: Vec<String>,
78 pub activity: Vec<usize>,
80}
81
82#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
84#[derive(Debug, Clone)]
85pub struct RunningWorker {
86 pub id: String,
88 pub queue: String,
90 pub backend: String,
92 pub started_at: u64,
94 pub last_heartbeat: u64,
96 pub layers: String,
98}
99
100#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
101#[derive(Debug, Clone)]
102pub struct Filter {
104 #[cfg_attr(feature = "serde", serde(default))]
106 pub status: Option<Status>,
107 #[cfg_attr(feature = "serde", serde(default = "default_page"))]
108 pub page: u32,
110 #[cfg_attr(feature = "serde", serde(default))]
112 pub page_size: Option<u32>,
113}
114
115impl Filter {
116 #[must_use]
118 pub fn offset(&self) -> u32 {
119 (self.page - 1) * self.page_size.unwrap_or(DEFAULT_PAGE_SIZE)
120 }
121
122 #[must_use]
124 pub fn limit(&self) -> u32 {
125 self.page_size.unwrap_or(DEFAULT_PAGE_SIZE)
126 }
127}
128
129#[cfg(feature = "serde")]
130fn default_page() -> u32 {
131 1
132}
133#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
135#[derive(Debug, Clone)]
136pub struct Statistic {
137 pub title: String,
139 pub stat_type: StatType,
141 pub value: String,
143 pub priority: Option<u64>,
145}
146#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
148#[derive(Debug, Clone, PartialEq, Eq, Default)]
149#[non_exhaustive]
150pub enum StatType {
151 Timestamp,
153 #[default]
155 Number,
156 Decimal,
158 Percentage,
160}
161
162#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
164#[non_exhaustive]
165pub enum ParseStatTypeError {
166 #[error("invalid stat type: `{0}`")]
168 InvalidStatType(String),
169}
170
171impl FromStr for StatType {
172 type Err = ParseStatTypeError;
173
174 fn from_str(s: &str) -> Result<Self, Self::Err> {
175 match s {
176 "Timestamp" => Ok(Self::Timestamp),
177 "Decimal" => Ok(Self::Decimal),
178 "Percentage" => Ok(Self::Percentage),
179 "Number" => Ok(Self::Number),
180 _ => Err(ParseStatTypeError::InvalidStatType(s.to_owned())),
181 }
182 }
183}