pub struct Queue {
pub connection_pool: Pool<ConnectionManager<PgConnection>>,
}Expand description
An async queue that can be used to enqueue tasks.
It uses a PostgreSQL storage. It must be connected to perform any operation.
To connect a Queue to the PostgreSQL database call the get_connection method.
A Queue can be created with the TypedBuilder.
// Set DATABASE_URL enviroment variable if you would like to try this function.
pub fn connection_pool(pool_size: u32) -> r2d2::Pool<r2d2::ConnectionManager<PgConnection>> {
let database_url = env::var("DATABASE_URL").expect("DATABASE_URL must be set");
let manager = r2d2::ConnectionManager::<PgConnection>::new(database_url);
r2d2::Pool::builder()
.max_size(pool_size)
.build(manager)
.unwrap()
}
let queue = Queue::builder().connection_pool(connection_pool(3)).build();Fields§
§connection_pool: Pool<ConnectionManager<PgConnection>>Implementations§
Source§impl Queue
impl Queue
Sourcepub fn get_connection(&self) -> Result<PoolConnection, QueueError>
pub fn get_connection(&self) -> Result<PoolConnection, QueueError>
Connect to the db if not connected
pub fn schedule_task_query( connection: &mut PgConnection, params: &dyn Runnable, ) -> Result<Task, QueueError>
pub fn insert_query( connection: &mut PgConnection, params: &dyn Runnable, scheduled_at: DateTime<Utc>, ) -> Result<Task, QueueError>
pub fn fetch_task_query( connection: &mut PgConnection, task_type: String, ) -> Option<Task>
pub fn fetch_and_touch_query( connection: &mut PgConnection, task_type: String, ) -> Result<Option<Task>, QueueError>
pub fn find_task_by_id_query( connection: &mut PgConnection, id: &Uuid, ) -> Option<Task>
pub fn remove_all_tasks_query( connection: &mut PgConnection, ) -> Result<usize, QueueError>
pub fn remove_all_scheduled_tasks_query( connection: &mut PgConnection, ) -> Result<usize, QueueError>
pub fn remove_tasks_of_type_query( connection: &mut PgConnection, task_type: &str, ) -> Result<usize, QueueError>
pub fn remove_task_by_metadata_query( connection: &mut PgConnection, task: &dyn Runnable, ) -> Result<usize, QueueError>
pub fn remove_task_query( connection: &mut PgConnection, id: &Uuid, ) -> Result<usize, QueueError>
pub fn update_task_state_query( connection: &mut PgConnection, task: &Task, state: FangTaskState, ) -> Result<Task, QueueError>
pub fn fail_task_query( connection: &mut PgConnection, task: &Task, error: &str, ) -> Result<Task, QueueError>
pub fn schedule_retry_query( connection: &mut PgConnection, task: &Task, backoff_seconds: u32, error: &str, ) -> Result<Task, QueueError>
Trait Implementations§
Source§impl Queueable for Queue
impl Queueable for Queue
Source§fn remove_task_by_metadata(
&self,
task: &dyn Runnable,
) -> Result<usize, QueueError>
fn remove_task_by_metadata( &self, task: &dyn Runnable, ) -> Result<usize, QueueError>
To use this function task has to be uniq. uniq() has to return true. If task is not uniq this function will not do anything.
Source§fn fetch_and_touch_task(
&self,
task_type: String,
) -> Result<Option<Task>, QueueError>
fn fetch_and_touch_task( &self, task_type: String, ) -> Result<Option<Task>, QueueError>
This method should retrieve one task of the
task_type type. After fetching it should update the state
of the task to FangTaskState::InProgress.Source§fn insert_task(&self, params: &dyn Runnable) -> Result<Task, QueueError>
fn insert_task(&self, params: &dyn Runnable) -> Result<Task, QueueError>
Enqueue a task to the queue, The task will be executed as soon as possible by the worker of the same type
created by an
WorkerPool.Source§fn schedule_task(&self, params: &dyn Runnable) -> Result<Task, QueueError>
fn schedule_task(&self, params: &dyn Runnable) -> Result<Task, QueueError>
Schedule a task.
Source§fn remove_all_scheduled_tasks(&self) -> Result<usize, QueueError>
fn remove_all_scheduled_tasks(&self) -> Result<usize, QueueError>
Remove all tasks that are scheduled in the future.
Source§fn remove_all_tasks(&self) -> Result<usize, QueueError>
fn remove_all_tasks(&self) -> Result<usize, QueueError>
The method will remove all tasks from the queue
Source§fn remove_tasks_of_type(&self, task_type: &str) -> Result<usize, QueueError>
fn remove_tasks_of_type(&self, task_type: &str) -> Result<usize, QueueError>
Removes all tasks that have the specified
task_type.Source§fn remove_task(&self, id: &Uuid) -> Result<usize, QueueError>
fn remove_task(&self, id: &Uuid) -> Result<usize, QueueError>
Remove a task by its id.
Source§fn update_task_state(
&self,
task: &Task,
state: FangTaskState,
) -> Result<Task, QueueError>
fn update_task_state( &self, task: &Task, state: FangTaskState, ) -> Result<Task, QueueError>
Update the state field of the specified task
See the
FangTaskState enum for possible states.Source§fn fail_task(&self, task: &Task, error: &str) -> Result<Task, QueueError>
fn fail_task(&self, task: &Task, error: &str) -> Result<Task, QueueError>
Update the state of a task to
FangTaskState::Failed and set an error_message.fn find_task_by_id(&self, id: &Uuid) -> Option<Task>
fn schedule_retry( &self, task: &Task, backoff_seconds: u32, error: &str, ) -> Result<Task, QueueError>
Auto Trait Implementations§
impl !RefUnwindSafe for Queue
impl !UnwindSafe for Queue
impl Freeze for Queue
impl Send for Queue
impl Sync for Queue
impl Unpin for Queue
impl UnsafeUnpin for Queue
Blanket Implementations§
Source§impl<T> AggregateExpressionMethods for T
impl<T> AggregateExpressionMethods for T
Source§fn aggregate_distinct(self) -> Self::Outputwhere
Self: DistinctDsl,
fn aggregate_distinct(self) -> Self::Outputwhere
Self: DistinctDsl,
DISTINCT modifier for aggregate functions Read moreSource§fn aggregate_all(self) -> Self::Outputwhere
Self: AllDsl,
fn aggregate_all(self) -> Self::Outputwhere
Self: AllDsl,
ALL modifier for aggregate functions Read moreSource§fn aggregate_filter<P>(self, f: P) -> Self::Output
fn aggregate_filter<P>(self, f: P) -> Self::Output
Add an aggregate function filter Read more
Source§fn aggregate_order<O>(self, o: O) -> Self::Outputwhere
Self: OrderAggregateDsl<O>,
fn aggregate_order<O>(self, o: O) -> Self::Outputwhere
Self: OrderAggregateDsl<O>,
Add an aggregate function order Read more
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Downcast for Twhere
T: Any,
impl<T> Downcast for Twhere
T: Any,
Source§fn into_any(self: Box<T>) -> Box<dyn Any>
fn into_any(self: Box<T>) -> Box<dyn Any>
Converts
Box<dyn Trait> (where Trait: Downcast) to Box<dyn Any>, which can then be
downcast into Box<dyn ConcreteType> where ConcreteType implements Trait.Source§fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
fn into_any_rc(self: Rc<T>) -> Rc<dyn Any>
Converts
Rc<Trait> (where Trait: Downcast) to Rc<Any>, which can then be further
downcast into Rc<ConcreteType> where ConcreteType implements Trait.Source§fn as_any(&self) -> &(dyn Any + 'static)
fn as_any(&self) -> &(dyn Any + 'static)
Converts
&Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &Any’s vtable from &Trait’s.Source§fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)
Converts
&mut Trait (where Trait: Downcast) to &Any. This is needed since Rust cannot
generate &mut Any’s vtable from &mut Trait’s.Source§impl<T> DowncastSend for T
impl<T> DowncastSend for T
Source§impl<T> DowncastSync for T
impl<T> DowncastSync for T
Source§impl<T> IntoSql for T
impl<T> IntoSql for T
Source§fn into_sql<T>(self) -> Self::Expression
fn into_sql<T>(self) -> Self::Expression
Convert
self to an expression for Diesel’s query builder. Read moreSource§fn as_sql<'a, T>(&'a self) -> <&'a Self as AsExpression<T>>::Expression
fn as_sql<'a, T>(&'a self) -> <&'a Self as AsExpression<T>>::Expression
Convert
&self to an expression for Diesel’s query builder. Read moreSource§impl<T> WindowExpressionMethods for T
impl<T> WindowExpressionMethods for T
Source§fn over(self) -> Self::Outputwhere
Self: OverDsl,
fn over(self) -> Self::Outputwhere
Self: OverDsl,
Turn a function call into a window function call Read more
Source§fn window_filter<P>(self, f: P) -> Self::Output
fn window_filter<P>(self, f: P) -> Self::Output
Add a filter to the current window function Read more
Source§fn partition_by<E>(self, expr: E) -> Self::Outputwhere
Self: PartitionByDsl<E>,
fn partition_by<E>(self, expr: E) -> Self::Outputwhere
Self: PartitionByDsl<E>,
Add a partition clause to the current window function Read more
Source§fn window_order<E>(self, expr: E) -> Self::Outputwhere
Self: OrderWindowDsl<E>,
fn window_order<E>(self, expr: E) -> Self::Outputwhere
Self: OrderWindowDsl<E>,
Add a order clause to the current window function Read more