Skip to main content

Task

Enum Task 

Source
pub enum Task {
    Sync(SyncFn),
    Async(AsyncFn),
    SyncIter(SyncIterFn),
    AsyncStream(AsyncStreamFn),
    SyncBatch(SyncBatchFn),
    AsyncBatch(AsyncBatchFn),
    SyncIterBatch(SyncIterBatchFn),
    AsyncStreamBatch(AsyncStreamBatchFn),
}
Expand description

A single reusable unit of work in a cognee pipeline.

VariantExecutionInputOutput
Task::Syncblockingsingle valuesingle value
Task::Asyncnon-blockingsingle valuesingle value
Task::SyncIterblockingsingle valuelazy iterator
Task::AsyncStreamnon-blockingsingle valueasync stream
Task::SyncBatchblockingslice of valuessingle value
Task::AsyncBatchnon-blockingslice of valuessingle value
Task::SyncIterBatchblockingslice of valueslazy iterator
Task::AsyncStreamBatchnon-blockingslice of valuesasync stream

Single-value variants are called once per item. Batch variants receive a &[Box<dyn Value>] slice of items accumulated up to the task’s batch_size. The pipeline executor detects which kind the next task is and routes accordingly.

Variants§

§

Sync(SyncFn)

§

Async(AsyncFn)

§

SyncIter(SyncIterFn)

§

AsyncStream(AsyncStreamFn)

§

SyncBatch(SyncBatchFn)

§

AsyncBatch(AsyncBatchFn)

§

SyncIterBatch(SyncIterBatchFn)

§

AsyncStreamBatch(AsyncStreamBatchFn)

Implementations§

Source§

impl Task

Source

pub fn is_batch(&self) -> bool

Returns true if this task accepts a batch slice rather than a single value.

Source

pub fn python_task_type(&self) -> &'static str

Python-compat label used in the ${task_type} Task Started/ Completed/Errored analytics event names.

Mirrors Python’s tasks/task.py:194-207 inspect.isasyncgenfunction / iscoroutinefunction branch:

Rust variantPython label
Task::Sync, Task::SyncBatch"Function"
Task::Async, Task::AsyncBatch"Coroutine"
Task::SyncIter, Task::SyncIterBatch"Generator"
Task::AsyncStream, Task::AsyncStreamBatch"Async Generator"

The match is intentionally exhaustive (no wildcard arm) so that adding a new Task::* variant fails the build until the analytics mapping is decided.

Source§

impl Task

Source

pub fn sync<F>(f: F) -> Self
where F: Fn(Arc<dyn Value>, Arc<TaskContext>) -> Result<Arc<dyn Value>, TaskError> + Send + Sync + 'static,

Create a Task::Sync from a raw closure.

Source

pub fn async_fn<F>(f: F) -> Self
where F: Fn(Arc<dyn Value>, Arc<TaskContext>) -> BoxFuture<'static, Result<Arc<dyn Value>, TaskError>> + Send + Sync + 'static,

Create a Task::Async from a raw closure returning a BoxFuture.

Source

pub fn sync_iter<F>(f: F) -> Self
where F: Fn(Arc<dyn Value>, Arc<TaskContext>) -> Result<ValueIter, TaskError> + Send + Sync + 'static,

Create a Task::SyncIter from a raw closure returning a ValueIter.

Source

pub fn async_stream<F>(f: F) -> Self
where F: Fn(Arc<dyn Value>, Arc<TaskContext>) -> Result<ValueStream, TaskError> + Send + Sync + 'static,

Create a Task::AsyncStream from a raw closure returning a ValueStream.

Source

pub fn sync_batch<F>(f: F) -> Self
where F: for<'a> Fn(&'a [Box<dyn Value>], Arc<TaskContext>) -> Result<Arc<dyn Value>, TaskError> + Send + Sync + 'static,

Create a Task::SyncBatch from a raw closure.

Source

pub fn async_batch<F>(f: F) -> Self
where F: for<'a> Fn(&'a [Box<dyn Value>], Arc<TaskContext>) -> BoxFuture<'static, Result<Arc<dyn Value>, TaskError>> + Send + Sync + 'static,

Create a Task::AsyncBatch from a raw closure returning a BoxFuture.

Source

pub fn sync_iter_batch<F>(f: F) -> Self
where F: for<'a> Fn(&'a [Box<dyn Value>], Arc<TaskContext>) -> Result<ValueIter, TaskError> + Send + Sync + 'static,

Create a Task::SyncIterBatch from a raw closure returning a ValueIter.

Source

pub fn async_stream_batch<F>(f: F) -> Self
where F: for<'a> Fn(&'a [Box<dyn Value>], Arc<TaskContext>) -> Result<ValueStream, TaskError> + Send + Sync + 'static,

Create a Task::AsyncStreamBatch from a raw closure returning a ValueStream.

Source

pub fn sync_typed<I, O, F>(f: F) -> Self
where I: Value, O: Value, F: Fn(&I, Arc<TaskContext>) -> Result<Box<O>, TaskError> + Send + Sync + 'static,

Create a Task::Sync from a typed closure.

Task::sync_typed(|input: &MyInput, ctx| {
    Ok(Box::new(process(input)))
})
Source

pub fn async_fn_typed<I, O, F>(f: F) -> Self
where I: Value, O: Value, F: Fn(&I, Arc<TaskContext>) -> BoxFuture<'static, Result<Box<O>, TaskError>> + Send + Sync + 'static,

Create a Task::Async from a typed closure returning a 'static future.

Data needed inside the async block must be copied/cloned before it:

Task::async_fn_typed(|input: &MyInput, ctx| {
    let id = input.id;  // copy before async block
    Box::pin(async move {
        Ok(Box::new(fetch(id).await?))
    })
})
Source

pub fn sync_iter_typed<I, O, F, Iter>(f: F) -> Self
where I: Value, O: Value, F: Fn(&I, Arc<TaskContext>) -> Result<Iter, TaskError> + Send + Sync + 'static, Iter: Iterator<Item = Box<O>> + Send + 'static,

Create a Task::SyncIter from a typed closure returning a concrete iterator. The iterator must be 'static (may not borrow the input).

Task::sync_iter_typed(|input: &Document, ctx| {
    let chunks = split(input.text.clone());
    Ok(chunks.into_iter().map(Box::new))
})
Source

pub fn async_stream_typed<I, O, F, S>(f: F) -> Self
where I: Value, O: Value, F: Fn(&I, Arc<TaskContext>) -> Result<S, TaskError> + Send + Sync + 'static, S: Stream<Item = Box<O>> + Send + 'static,

Create a Task::AsyncStream from a typed closure returning a concrete stream. The stream must be 'static.

Task::async_stream_typed(|input: &DatasetId, ctx| {
    let id = *input;
    Ok(stream_chunks(id))
})
Source

pub fn sync_batch_typed<I, O, F>(f: F) -> Self
where I: Value, O: Value, F: for<'a> Fn(&'a [&'a I], Arc<TaskContext>) -> Result<Box<O>, TaskError> + Send + Sync + 'static,

Create a Task::SyncBatch from a typed closure receiving &[&I].

Task::sync_batch_typed(|chunks: &[&DocumentChunk], ctx| {
    Ok(Box::new(embed_all(chunks)))
})
Source

pub fn async_batch_typed<I, O, F>(f: F) -> Self
where I: Value, O: Value, F: for<'a> Fn(&'a [&'a I], Arc<TaskContext>) -> BoxFuture<'static, Result<Box<O>, TaskError>> + Send + Sync + 'static,

Create a Task::AsyncBatch from a typed closure returning a 'static future.

Data needed inside the async block must be copied/cloned before it:

Task::async_batch_typed(|chunks: &[&DocumentChunk], ctx| {
    let texts: Vec<String> = chunks.iter().map(|c| c.text.clone()).collect();
    Box::pin(async move {
        Ok(Box::new(embed_batch(texts).await?))
    })
})
Source

pub fn sync_iter_batch_typed<I, O, F, Iter>(f: F) -> Self
where I: Value, O: Value, F: for<'a> Fn(&'a [&'a I], Arc<TaskContext>) -> Result<Iter, TaskError> + Send + Sync + 'static, Iter: Iterator<Item = Box<O>> + Send + 'static,

Create a Task::SyncIterBatch from a typed closure returning a concrete iterator.

Source

pub fn async_stream_batch_typed<I, O, F, S>(f: F) -> Self
where I: Value, O: Value, F: for<'a> Fn(&'a [&'a I], Arc<TaskContext>) -> Result<S, TaskError> + Send + Sync + 'static, S: Stream<Item = Box<O>> + Send + 'static,

Create a Task::AsyncStreamBatch from a typed closure returning a concrete stream.

Source

pub fn call(&self, input: Arc<dyn Value>, ctx: Arc<TaskContext>) -> TaskCall

Call this task with a single input value.

Panics if called on a batch variant — use Task::call_batch for those. The task is Fn, so &self suffices — the same task object handles every input in a fan-out scenario and every retry attempt.

Source

pub fn parallel(tasks: Vec<Task>) -> Self

Build a task that runs multiple sub-tasks concurrently on the same input.

Semantics (matching Python run_tasks_parallel):

  • Each sub-task receives Arc::clone(&input) and the shared context.
  • All sub-tasks run concurrently via futures::future::join_all.
  • If any sub-task fails, the whole parallel task fails with that error.
  • On success, returns the result of the last sub-task (by position).

Only single-value (Sync / Async) sub-tasks are supported. Iter/stream sub-tasks inside a parallel group don’t have well-defined “last result” semantics and will panic at call time.

Source

pub fn call_batch( &self, items: &[Box<dyn Value>], ctx: Arc<TaskContext>, ) -> TaskCall

Call this batch task with a slice of accumulated values.

Panics if called on a single-value variant — use Task::call for those.

Trait Implementations§

Source§

impl From<Task> for TaskInfo

Source§

fn from(task: Task) -> Self

Converts to this type from the input type.
Source§

impl<I: Value, O: Value> From<TypedTask<I, O>> for Task

Source§

fn from(typed: TypedTask<I, O>) -> Self

Erase I and O, producing the type-erased Task.

Delegates to the corresponding Task::sync_typed / Task::async_fn_typed / … constructor, reusing their downcast logic.

Auto Trait Implementations§

§

impl !RefUnwindSafe for Task

§

impl !UnwindSafe for Task

§

impl Freeze for Task

§

impl Send for Task

§

impl Sync for Task

§

impl Unpin for Task

§

impl UnsafeUnpin for Task

Blanket Implementations§

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> Conv for T

Source§

fn conv<T>(self) -> T
where Self: Into<T>,

Converts self into T using Into<T>. Read more
Source§

impl<T> FmtForward for T

Source§

fn fmt_binary(self) -> FmtBinary<Self>
where Self: Binary,

Causes self to use its Binary implementation when Debug-formatted.
Source§

fn fmt_display(self) -> FmtDisplay<Self>
where Self: Display,

Causes self to use its Display implementation when Debug-formatted.
Source§

fn fmt_lower_exp(self) -> FmtLowerExp<Self>
where Self: LowerExp,

Causes self to use its LowerExp implementation when Debug-formatted.
Source§

fn fmt_lower_hex(self) -> FmtLowerHex<Self>
where Self: LowerHex,

Causes self to use its LowerHex implementation when Debug-formatted.
Source§

fn fmt_octal(self) -> FmtOctal<Self>
where Self: Octal,

Causes self to use its Octal implementation when Debug-formatted.
Source§

fn fmt_pointer(self) -> FmtPointer<Self>
where Self: Pointer,

Causes self to use its Pointer implementation when Debug-formatted.
Source§

fn fmt_upper_exp(self) -> FmtUpperExp<Self>
where Self: UpperExp,

Causes self to use its UpperExp implementation when Debug-formatted.
Source§

fn fmt_upper_hex(self) -> FmtUpperHex<Self>
where Self: UpperHex,

Causes self to use its UpperHex implementation when Debug-formatted.
Source§

fn fmt_list(self) -> FmtList<Self>
where &'a Self: for<'a> IntoIterator,

Formats each item in a sequence. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
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> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

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

Source§

fn pipe<R>(self, func: impl FnOnce(Self) -> R) -> R
where Self: Sized,

Pipes by value. This is generally the method you want to use. Read more
Source§

fn pipe_ref<'a, R>(&'a self, func: impl FnOnce(&'a Self) -> R) -> R
where R: 'a,

Borrows self and passes that borrow into the pipe function. Read more
Source§

fn pipe_ref_mut<'a, R>(&'a mut self, func: impl FnOnce(&'a mut Self) -> R) -> R
where R: 'a,

Mutably borrows self and passes that borrow into the pipe function. Read more
Source§

fn pipe_borrow<'a, B, R>(&'a self, func: impl FnOnce(&'a B) -> R) -> R
where Self: Borrow<B>, B: 'a + ?Sized, R: 'a,

Borrows self, then passes self.borrow() into the pipe function. Read more
Source§

fn pipe_borrow_mut<'a, B, R>( &'a mut self, func: impl FnOnce(&'a mut B) -> R, ) -> R
where Self: BorrowMut<B>, B: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.borrow_mut() into the pipe function. Read more
Source§

fn pipe_as_ref<'a, U, R>(&'a self, func: impl FnOnce(&'a U) -> R) -> R
where Self: AsRef<U>, U: 'a + ?Sized, R: 'a,

Borrows self, then passes self.as_ref() into the pipe function.
Source§

fn pipe_as_mut<'a, U, R>(&'a mut self, func: impl FnOnce(&'a mut U) -> R) -> R
where Self: AsMut<U>, U: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.as_mut() into the pipe function.
Source§

fn pipe_deref<'a, T, R>(&'a self, func: impl FnOnce(&'a T) -> R) -> R
where Self: Deref<Target = T>, T: 'a + ?Sized, R: 'a,

Borrows self, then passes self.deref() into the pipe function.
Source§

fn pipe_deref_mut<'a, T, R>( &'a mut self, func: impl FnOnce(&'a mut T) -> R, ) -> R
where Self: DerefMut<Target = T> + Deref, T: 'a + ?Sized, R: 'a,

Mutably borrows self, then passes self.deref_mut() into the pipe function.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

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

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> Tap for T

Source§

fn tap(self, func: impl FnOnce(&Self)) -> Self

Immutable access to a value. Read more
Source§

fn tap_mut(self, func: impl FnOnce(&mut Self)) -> Self

Mutable access to a value. Read more
Source§

fn tap_borrow<B>(self, func: impl FnOnce(&B)) -> Self
where Self: Borrow<B>, B: ?Sized,

Immutable access to the Borrow<B> of a value. Read more
Source§

fn tap_borrow_mut<B>(self, func: impl FnOnce(&mut B)) -> Self
where Self: BorrowMut<B>, B: ?Sized,

Mutable access to the BorrowMut<B> of a value. Read more
Source§

fn tap_ref<R>(self, func: impl FnOnce(&R)) -> Self
where Self: AsRef<R>, R: ?Sized,

Immutable access to the AsRef<R> view of a value. Read more
Source§

fn tap_ref_mut<R>(self, func: impl FnOnce(&mut R)) -> Self
where Self: AsMut<R>, R: ?Sized,

Mutable access to the AsMut<R> view of a value. Read more
Source§

fn tap_deref<T>(self, func: impl FnOnce(&T)) -> Self
where Self: Deref<Target = T>, T: ?Sized,

Immutable access to the Deref::Target of a value. Read more
Source§

fn tap_deref_mut<T>(self, func: impl FnOnce(&mut T)) -> Self
where Self: DerefMut<Target = T> + Deref, T: ?Sized,

Mutable access to the Deref::Target of a value. Read more
Source§

fn tap_dbg(self, func: impl FnOnce(&Self)) -> Self

Calls .tap() only in debug builds, and is erased in release builds.
Source§

fn tap_mut_dbg(self, func: impl FnOnce(&mut Self)) -> Self

Calls .tap_mut() only in debug builds, and is erased in release builds.
Source§

fn tap_borrow_dbg<B>(self, func: impl FnOnce(&B)) -> Self
where Self: Borrow<B>, B: ?Sized,

Calls .tap_borrow() only in debug builds, and is erased in release builds.
Source§

fn tap_borrow_mut_dbg<B>(self, func: impl FnOnce(&mut B)) -> Self
where Self: BorrowMut<B>, B: ?Sized,

Calls .tap_borrow_mut() only in debug builds, and is erased in release builds.
Source§

fn tap_ref_dbg<R>(self, func: impl FnOnce(&R)) -> Self
where Self: AsRef<R>, R: ?Sized,

Calls .tap_ref() only in debug builds, and is erased in release builds.
Source§

fn tap_ref_mut_dbg<R>(self, func: impl FnOnce(&mut R)) -> Self
where Self: AsMut<R>, R: ?Sized,

Calls .tap_ref_mut() only in debug builds, and is erased in release builds.
Source§

fn tap_deref_dbg<T>(self, func: impl FnOnce(&T)) -> Self
where Self: Deref<Target = T>, T: ?Sized,

Calls .tap_deref() only in debug builds, and is erased in release builds.
Source§

fn tap_deref_mut_dbg<T>(self, func: impl FnOnce(&mut T)) -> Self
where Self: DerefMut<Target = T> + Deref, T: ?Sized,

Calls .tap_deref_mut() only in debug builds, and is erased in release builds.
Source§

impl<T> TryConv for T

Source§

fn try_conv<T>(self) -> Result<T, Self::Error>
where Self: TryInto<T>,

Attempts to convert self into T using TryInto<T>. Read more
Source§

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

Source§

type Error = Infallible

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

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

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<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> Value for T
where T: Any + Send + Sync + 'static,

Source§

fn as_any(&self) -> &(dyn Any + 'static)

Source§

fn as_any_mut(&mut self) -> &mut (dyn Any + 'static)

Source§

fn into_any(self: Box<T>) -> Box<dyn Any>

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more