Skip to main content

TaskContext

Struct TaskContext 

Source
pub struct TaskContext { /* private fields */ }
Available on crate feature graph only.
Expand description

Runtime context passed to #[entrypoint] and #[task] functions.

Provides access to workflow state, checkpointing, interrupt/resume, and progress streaming. Each task function receives a mutable reference to TaskContext enabling state reads, writes, event emission, and interrupt requests.

§Example

use adk_graph::functional::TaskContext;

#[task]
async fn my_step(ctx: &mut TaskContext) -> Result<Value> {
    // Read state
    let count: i64 = ctx.get("counter").unwrap_or(0);

    // Write state
    ctx.set("counter", serde_json::json!(count + 1));

    // Emit progress
    ctx.emit(StreamEvent::custom("my_step", "progress", serde_json::json!({"count": count + 1})));

    Ok(serde_json::json!({"new_count": count + 1}))
}

Implementations§

Source§

impl TaskContext

Source

pub fn new( thread_id: String, state: HashMap<String, Value>, checkpointer: Arc<dyn Checkpointer>, event_tx: Sender<StreamEvent>, execution_log: Arc<RwLock<ExecutionLog>>, cancel_token: CancellationToken, schema: Option<StateSchema>, ) -> TaskContext

Available on crate feature functional only.

Create a new TaskContext.

Typically constructed by the macro-generated entrypoint, not by user code.

Source

pub fn state(&self) -> &HashMap<String, Value>

Available on crate feature functional only.

Get the current workflow state (read-only).

Source

pub fn get<T>(&self, key: &str) -> Option<T>

Available on crate feature functional only.

Get a typed value from state.

Attempts to deserialize the value stored at key into type T. Returns None if the key does not exist or deserialization fails.

§Example
let count: Option<i64> = ctx.get("counter");
Source

pub fn set(&mut self, key: &str, value: impl Into<Value>)

Available on crate feature functional only.

Set a value in state.

If a StateSchema is configured, the update is applied using the appropriate reducer for the key. Otherwise the value is set directly (overwrite semantics).

§Example
ctx.set("counter", serde_json::json!(42));
Source

pub fn emit(&self, event: StreamEvent)

Available on crate feature functional only.

Emit a progress event to stream listeners.

Events are broadcast to all registered receivers. If no listeners are active the event is silently dropped.

§Example
ctx.emit(StreamEvent::custom("my_task", "progress", json!({"pct": 50})));
Source

pub async fn interrupt<T>(&self, message: &str) -> Result<T, GraphError>

Available on crate feature functional only.

Interrupt execution and wait for external input.

Persists the current state as an interrupt checkpoint, emits an interrupted event, and suspends execution. When the workflow is resumed with an interrupt value, the value is deserialized into T and returned.

§Errors

Returns FunctionalError::InterruptTypeMismatch if the resume value cannot be deserialized into T.

Returns FunctionalError::CheckpointFailed if persisting the interrupt checkpoint fails.

§Example
let approval: bool = ctx.interrupt("Please approve this action").await?;
Source

pub fn with_resume_values( self, resume_values: HashMap<String, Value>, ) -> TaskContext

Available on crate feature functional only.

Supplies values for interrupt sites, keyed by continuation key.

Re-invoke the entrypoint with these set to resume: an interrupt whose key is present returns the deserialized value instead of suspending.

§Example
// First run suspends and reports its continuation key.
let Err(e) = run(ctx).await else { unreachable!() };

// Second run supplies the value under that key.
let ctx = ctx.with_resume_values(HashMap::from([
    ("interrupt-1".to_string(), serde_json::json!({ "approved": true })),
]));
let output = run(ctx).await?;
Source

pub fn resume_values(&self) -> &HashMap<String, Value>

Available on crate feature functional only.

The values available to interrupt sites in this context.

Source

pub fn thread_id(&self) -> &str

Available on crate feature functional only.

Get the thread identifier for this context.

Source

pub fn cancel_token(&self) -> &CancellationToken

Available on crate feature functional only.

Get a reference to the cancellation token.

Source

pub fn is_cancelled(&self) -> bool

Available on crate feature functional only.

Check if the workflow has been cancelled.

Source

pub async fn current_step(&self) -> usize

Available on crate feature functional only.

Get the current step number from the execution log.

Source

pub fn with_schema_validator( self, validator: StateSchemaValidator, ) -> TaskContext

Available on crate feature functional only.

Set a StateSchemaValidator for this context.

When set, the validator is used to validate initial state at workflow start and task output before applying reducers.

Source

pub fn schema_validator(&self) -> Option<&StateSchemaValidator>

Available on crate feature functional only.

Get the schema validator, if configured.

Source

pub fn validate_state(&self) -> Result<(), FunctionalError>

Available on crate feature functional only.

Validate the current state against the schema validator.

Called at workflow start to validate initial state.

§Errors

Returns FunctionalError::SchemaValidation if validation fails.

Source

pub fn validate_task_output( &self, output: &HashMap<String, Value>, ) -> Result<(), FunctionalError>

Available on crate feature functional only.

Validate task output against the schema validator.

Called after a task produces output, before applying reducers.

§Errors

Returns FunctionalError::SchemaValidation if validation fails.

Source

pub fn iteration_key(&mut self, task_name: &str) -> String

Available on crate feature functional only.

Generate a unique checkpoint key for a task inside a loop.

Each call to this method for the same task_name increments the iteration counter, producing keys like "step_a::iter_0", "step_a::iter_1", etc. Keys are deterministic from task name and iteration index.

§Example
for item in items {
    let key = ctx.iteration_key("process_item");
    // key = "process_item::iter_0", "process_item::iter_1", ...
}
Source

pub fn current_iteration(&self, task_name: &str) -> Option<usize>

Available on crate feature functional only.

Get the current iteration index for a task without incrementing.

Returns None if the task has not been called in a loop yet.

Source

pub fn reset_iteration(&mut self, task_name: &str)

Available on crate feature functional only.

Reset the iteration counter for a task.

Useful when re-entering a loop (e.g., nested loops or retry).

Source

pub fn reset_all_iterations(&mut self)

Available on crate feature functional only.

Reset all iteration counters.

Source

pub fn route_to(&mut self, targets: &[&str])

Available on crate feature functional only.

Records the task names this task chose, for its caller to read.

Nothing in this crate acts on the value. The functional API runs your own control flow: #[entrypoint] calls your function once, and the awaits inside it are the order of execution. So a task states its choice here and the surrounding code reads it with Self::take_pending_route and calls what it names.

For routing the framework performs, build the workflow with StateGraph::add_conditional_edges, which resolves a route key against declared targets and dispatches to them.

§Example
use std::collections::HashMap;
use std::sync::Arc;

use adk_graph::checkpoint::{Checkpointer, MemoryCheckpointer};
use adk_graph::functional::{ExecutionLog, TaskContext};
use tokio::sync::{RwLock, broadcast};
use tokio_util::sync::CancellationToken;

let (event_tx, _rx) = broadcast::channel(8);
let mut ctx = TaskContext::new(
    "thread-1".to_string(),
    HashMap::new(),
    Arc::new(MemoryCheckpointer::new()) as Arc<dyn Checkpointer>,
    event_tx,
    Arc::new(RwLock::new(ExecutionLog::new())),
    CancellationToken::new(),
    None,
);

// A task states which branches it chose.
ctx.route_to(&["process_a", "process_b"]);

// The surrounding code reads the choice and calls those tasks itself.
let chosen = ctx.take_pending_route().expect("a route was set");
assert_eq!(chosen, vec!["process_a".to_string(), "process_b".to_string()]);

// Reading it clears it, so the next task starts with no choice pending.
assert_eq!(ctx.take_pending_route(), None);
Source

pub fn take_pending_route(&mut self) -> Option<Vec<String>>

Available on crate feature functional only.

Takes the task names recorded by Self::route_to, clearing them.

Returns None when no task set a route. Call this from the code that sequences your tasks; nothing in this crate calls it for you.

Auto Trait Implementations§

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> MaybeSend for T
where T: Send,

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> 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 = !

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> 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