Skip to main content

Queue

Struct Queue 

Source
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

Source

pub fn builder() -> QueueBuilder<((),)>

Create a builder for building Queue. On the builder, call .connection_pool(...) to set the values of the fields. Finally, call .build() to create the instance of Queue.

Source§

impl Queue

Source

pub fn get_connection(&self) -> Result<PoolConnection, QueueError>

Connect to the db if not connected

Source

pub fn schedule_task_query( connection: &mut PgConnection, params: &dyn Runnable, ) -> Result<Task, QueueError>

Source

pub fn insert_query( connection: &mut PgConnection, params: &dyn Runnable, scheduled_at: DateTime<Utc>, ) -> Result<Task, QueueError>

Source

pub fn fetch_task_query( connection: &mut PgConnection, task_type: String, ) -> Option<Task>

Source

pub fn fetch_and_touch_query( connection: &mut PgConnection, task_type: String, ) -> Result<Option<Task>, QueueError>

Source

pub fn find_task_by_id_query( connection: &mut PgConnection, id: &Uuid, ) -> Option<Task>

Source

pub fn remove_all_tasks_query( connection: &mut PgConnection, ) -> Result<usize, QueueError>

Source

pub fn remove_all_scheduled_tasks_query( connection: &mut PgConnection, ) -> Result<usize, QueueError>

Source

pub fn remove_tasks_of_type_query( connection: &mut PgConnection, task_type: &str, ) -> Result<usize, QueueError>

Source

pub fn remove_task_by_metadata_query( connection: &mut PgConnection, task: &dyn Runnable, ) -> Result<usize, QueueError>

Source

pub fn remove_task_query( connection: &mut PgConnection, id: &Uuid, ) -> Result<usize, QueueError>

Source

pub fn update_task_state_query( connection: &mut PgConnection, task: &Task, state: FangTaskState, ) -> Result<Task, QueueError>

Source

pub fn fail_task_query( connection: &mut PgConnection, task: &Task, error: &str, ) -> Result<Task, QueueError>

Source

pub fn schedule_retry_query( connection: &mut PgConnection, task: &Task, backoff_seconds: u32, error: &str, ) -> Result<Task, QueueError>

Trait Implementations§

Source§

impl Clone for Queue

Source§

fn clone(&self) -> Queue

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Queueable for Queue

Source§

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>

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>

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>

Schedule a task.
Source§

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>

The method will remove all tasks from the queue
Source§

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>

Remove a task by its id.
Source§

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>

Update the state of a task to FangTaskState::Failed and set an error_message.
Source§

fn find_task_by_id(&self, id: &Uuid) -> Option<Task>

Source§

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

Source§

fn aggregate_distinct(self) -> Self::Output
where Self: DistinctDsl,

DISTINCT modifier for aggregate functions Read more
Source§

fn aggregate_all(self) -> Self::Output
where Self: AllDsl,

ALL modifier for aggregate functions Read more
Source§

fn aggregate_filter<P>(self, f: P) -> Self::Output
where P: AsExpression<Bool>, Self: FilterDsl<<P as AsExpression<Bool>>::Expression>,

Add an aggregate function filter Read more
Source§

fn aggregate_order<O>(self, o: O) -> Self::Output
where Self: OrderAggregateDsl<O>,

Add an aggregate function order Read more
Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> Downcast for T
where T: Any,

Source§

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>

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)

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)

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
where T: Any + Send,

Source§

fn into_any_send(self: Box<T>) -> Box<dyn Any + Send>

Converts Box<Trait> (where Trait: DowncastSend) to Box<dyn Any + Send>, which can then be downcast into Box<ConcreteType> where ConcreteType implements Trait.
Source§

impl<T> DowncastSync for T
where T: Any + Send + Sync,

Source§

fn into_any_sync(self: Box<T>) -> Box<dyn Any + Send + Sync>

Converts Box<Trait> (where Trait: DowncastSync) to Box<dyn Any + Send + Sync>, which can then be downcast into Box<ConcreteType> where ConcreteType implements Trait.
Source§

fn into_any_arc(self: Arc<T>) -> Arc<dyn Any + Send + Sync>

Converts Arc<Trait> (where Trait: DowncastSync) to Arc<Any>, which can then be downcast into Arc<ConcreteType> where ConcreteType implements Trait.
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoSql for T

Source§

fn into_sql<T>(self) -> Self::Expression

Convert self to an expression for Diesel’s query builder. Read more
Source§

fn as_sql<'a, T>(&'a self) -> <&'a Self as AsExpression<T>>::Expression
where &'a Self: AsExpression<T>, T: SqlType + TypedExpressionType,

Convert &self to an expression for Diesel’s query builder. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WindowExpressionMethods for T

Source§

fn over(self) -> Self::Output
where Self: OverDsl,

Turn a function call into a window function call Read more
Source§

fn window_filter<P>(self, f: P) -> Self::Output
where P: AsExpression<Bool>, Self: FilterDsl<<P as AsExpression<Bool>>::Expression>,

Add a filter to the current window function Read more
Source§

fn partition_by<E>(self, expr: E) -> Self::Output
where Self: PartitionByDsl<E>,

Add a partition clause to the current window function Read more
Source§

fn window_order<E>(self, expr: E) -> Self::Output
where Self: OrderWindowDsl<E>,

Add a order clause to the current window function Read more
Source§

fn frame_by<E>(self, expr: E) -> Self::Output
where Self: FrameDsl<E>,

Add a frame clause to the current window function Read more