Struct Client

Source
pub struct Client { /* private fields */ }

Implementations§

Source§

impl Client

Source

pub async fn open_cancellation_scope_on_claim( &self, task: &WorkflowTask, sequence: u64, parent_scope_id: &str, shield_parent: bool, ) -> Result<CancellationScopeOpenReceipt>

Open a scope and prove its canonical identity on the original claim.

This explicit protocol 1.20 operation keeps one five-second budget for lost acknowledgement recovery and all history pages. No lease, worker capability or cancellation execution authority is added by this proof.

Source§

impl Client

Source

pub async fn acknowledge_activity_cancellation( &self, task_id: &str, activity_attempt_id: &str, lease_owner: &str, request_id: &str, ) -> Result<Value>

Report a stopped and dropped remote callback under its original claim.

Call only after callback execution has ended. This explicit protocol 1.20 operation preserves the original cancellation request and grants no lease or result publication authority. The request has one five-second budget.

Source

pub async fn activity_task_status( &self, task_id: &str, activity_attempt_id: &str, lease_owner: &str, ) -> Result<Value>

Observe an exact activity attempt without renewing its lease or progress.

This explicit worker-protocol 1.20 operation has one five-second budget. A pending workflow request alone does not cancel an activity. Its reply reports the runtime’s current continuation authority after delivery.

Source

pub async fn poll_cooperative_workflow_task( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<CooperativeWorkflowTaskPoll>

Acquire a claim with protocol 1.20 and retain its pending observation.

The Server must have admitted a capable worker registration. This method registers no capability and does not change ordinary Worker polling. It never substitutes a worker ID for a missing lease owner or invents a task attempt. Poll retries reuse one request ID. Poll and history loading share one long-poll budget plus five seconds, with at most 128 pages.

Source

pub async fn heartbeat_workflow_task( &self, task: &WorkflowTask, original: Option<&CancellationRequest>, ) -> Result<WorkflowTaskHeartbeat>

Renew the exact selected workflow-task lease using worker protocol 1.20.

The reply must acknowledge the actual task, owner and attempt. Renewal neither commits cancellation delivery nor extends its cleanup deadline.

Source

pub async fn deliver_workflow_cancellation( &self, task: &WorkflowTask, delivery: &CancellationDelivery, ) -> Result<CancellationDeliveryReply>

Commit cooperative delivery on the exact selected workflow-task claim.

This explicit worker-protocol 1.20 operation renews no lease and advertises no worker capability. The Server requires a previously admitted capable claim. Reload canonical history before exposing cancellation to workflow code, including when an acknowledgment is lost. A cancellation wait explicitly releases the claim and must return the worker to polling.

Source

pub async fn refresh_workflow_cancellation_history( &self, task: &WorkflowTask, observation: &CancellationRequest, ) -> Result<Vec<HistoryEvent>>

Reload canonical history with the Server-issued token and exact claim.

The refresh has one five-second budget, at most 128 pages, and the existing SDK page-size limit. Invalid, repeated or non-progressing tokens and pages fail closed. This returns fresh history without modifying the caller’s previous task snapshot or its original request identity.

Source

pub async fn request_workflow_cancellation( &self, workflow_id: &str, options: CooperativeCancellationOptions, ) -> Result<WorkflowCancellationRequest>

Request bounded workflow-authored cleanup on a capable runtime.

This is separate from terminal cancellation. Discovery must advertise cooperative support and compatible protocol 1.20. The Server also refuses an active workflow claim that cannot deliver cooperation. Repeated requests retain its original request identity and deadline.

Source

pub async fn request_workflow_run_cancellation( &self, workflow_id: &str, run_id: &str, options: CooperativeCancellationOptions, ) -> Result<WorkflowCancellationRequest>

Request cleanup only if the selected run remains current.

Source§

impl Client

Source

pub async fn create_worker_session( &self, worker_id: &str, options: &WorkerSessionOptions, ) -> Result<Value>

Create, reuse or reacquire a session for this registered worker.

Server owns admission, requirements, capacity and holder authority. An admitted reacquisition requires rebuilding worker-local resources.

Source

pub async fn renew_worker_session( &self, worker_id: &str, session_id: &str, lease_seconds: u64, ) -> Result<Value>

Renew only the current session holder’s lease. This does not extend TTL.

Source

pub async fn close_worker_session( &self, worker_id: &str, session_id: &str, reason: &str, ) -> Result<Value>

Close one holder’s session. A closed identity cannot be reacquired.

Source§

impl Client

Source

pub fn new(base_url: impl Into<String>) -> Result<Self>

Source

pub fn builder(base_url: impl Into<String>) -> ClientBuilder

Source

pub async fn health(&self) -> Result<Value>

Source

pub async fn cluster_info(&self) -> Result<Value>

Source

pub async fn start_workflow<T: Serialize>( &self, workflow_type: &str, task_queue: &str, workflow_id: &str, input: T, ) -> Result<WorkflowHandle>

Source

pub async fn start_workflow_with_options<T: Serialize>( &self, workflow_type: &str, task_queue: &str, workflow_id: &str, options: WorkflowStartOptions, input: T, ) -> Result<WorkflowHandle>

Start a workflow with explicit server-enforced execution and run deadlines.

Source

pub async fn signal_workflow<T: Serialize>( &self, workflow_id: &str, signal_name: &str, input: T, ) -> Result<Value>

Source

pub async fn append_message_stream<T: Serialize>( &self, workflow_id: &str, stream_name: &str, message_id: &str, input: T, ) -> Result<Value>

Append one idempotently identified item to an instance-scoped input stream.

Source

pub async fn signal_workflow_run<T: Serialize>( &self, workflow_id: &str, run_id: &str, signal_name: &str, input: T, ) -> Result<Value>

Signal only if run_id is still the current run for this instance.

Source

pub async fn cancel_workflow( &self, workflow_id: &str, options: WorkflowCommandOptions, ) -> Result<WorkflowCommandResult>

Close the current run as cancelled immediately.

Server revokes open tasks and timers without resuming workflow code for cleanup. Use Client::request_workflow_cancellation for a separate cooperative request on a supporting runtime and opted-in worker.

Source

pub async fn cancel_workflow_run( &self, workflow_id: &str, run_id: &str, options: WorkflowCommandOptions, ) -> Result<WorkflowCommandResult>

Close the selected run as cancelled, only if run_id is still current.

Source

pub async fn terminate_workflow( &self, workflow_id: &str, options: WorkflowCommandOptions, ) -> Result<WorkflowCommandResult>

Forcefully terminate the current run for an instance.

Source

pub async fn terminate_workflow_run( &self, workflow_id: &str, run_id: &str, options: WorkflowCommandOptions, ) -> Result<WorkflowCommandResult>

Forcefully terminate only if run_id is still current.

Source

pub async fn redrive_workflow_run( &self, workflow_id: &str, failed_run_id: &str, request_id: Option<&str>, ) -> Result<WorkflowRedriveResult>

Continue a failed run from its recorded activity failure boundary.

Source

pub async fn query_workflow<T: Serialize>( &self, workflow_id: &str, query_name: &str, input: T, ) -> Result<Value>

Execute a named, read-only query against a running or completed workflow.

Arguments and results use the platform payload envelope. Server and worker rejections are returned as Error::QueryFailed with a stable reason, HTTP status, and original response body.

Source

pub async fn query_workflow_run<T: Serialize>( &self, workflow_id: &str, run_id: &str, query_name: &str, input: T, ) -> Result<Value>

Query only if run_id is still current, preventing accidental retargeting.

Source

pub async fn query_workflow_avro_value<T: Serialize>( &self, workflow_id: &str, query_name: &str, input: T, ) -> Result<AvroValue>

Query a workflow and return the lossless fixed Avro Value result.

Source

pub async fn query_workflow_run_avro_value<T: Serialize>( &self, workflow_id: &str, run_id: &str, query_name: &str, input: T, ) -> Result<AvroValue>

Query a selected run and return the lossless fixed Avro Value result.

Source

pub async fn update_workflow<T: Serialize>( &self, workflow_id: &str, update_name: &str, input: T, request_id: Option<&str>, ) -> Result<Value>

Send a synchronous update using fixed Avro Value arguments.

Source

pub async fn update_workflow_avro_value<T: Serialize>( &self, workflow_id: &str, update_name: &str, input: T, request_id: Option<&str>, ) -> Result<AvroValue>

Send a synchronous update and retain a bytes-capable Avro result.

Source

pub async fn describe_workflow( &self, workflow_id: &str, ) -> Result<WorkflowDescription>

Source

pub async fn describe_workflow_run( &self, workflow_id: &str, run_id: &str, ) -> Result<WorkflowDescription>

Describe one selected run, including historical terminal runs.

Source

pub async fn list_workflow_streams( &self, workflow_id: &str, run_id: &str, ) -> Result<Vec<WorkflowStreamDescription>>

List the run-scoped output streams already opened by a workflow.

Source

pub async fn describe_workflow_stream( &self, workflow_id: &str, run_id: &str, stream_name: &str, ) -> Result<WorkflowStreamDescription>

Describe stream lifecycle, offsets, pending count, and terminal error.

Source

pub async fn subscribe_workflow_stream( &self, workflow_id: &str, run_id: &str, stream_name: &str, from_offset: u64, max_items: usize, wait: Duration, ) -> Result<WorkflowStreamPage>

Read one bounded page beginning at a zero-based offset.

Delivery is at least once: persist next_offset only after processing the page. The future is cancellation-safe; dropping it cancels the in-flight request. Long polling is capped at 60 seconds by the SDK and service contract.

Source

pub async fn append_workflow_stream( &self, workflow_id: &str, run_id: &str, stream_name: &str, items: &[WorkflowStreamAppendItem], max_pending_items: Option<u64>, ) -> Result<WorkflowStreamAppendResult>

Append inline Avro envelopes or opaque external payload references.

Source

pub async fn close_workflow_stream( &self, workflow_id: &str, run_id: &str, stream_name: &str, error_reason: Option<&str>, retention_seconds: Option<u64>, ) -> Result<WorkflowStreamDescription>

Close a stream, or mark it errored when error_reason is supplied.

Source

pub async fn register_worker( &self, worker_id: &str, task_queue: &str, supported_workflow_types: Vec<String>, supported_activity_types: Vec<String>, max_concurrent_workflow_tasks: usize, max_concurrent_activity_tasks: usize, ) -> Result<RegisterWorkerResponse>

Source

pub async fn register_worker_with_capabilities( &self, worker_id: &str, task_queue: &str, supported_workflow_types: Vec<String>, supported_activity_types: Vec<String>, max_concurrent_workflow_tasks: usize, max_concurrent_activity_tasks: usize, capabilities: Vec<String>, ) -> Result<RegisterWorkerResponse>

Register a worker and explicitly advertise additive worker capabilities.

Source

pub async fn register_worker_with_command_contracts( &self, worker_id: &str, task_queue: &str, supported_workflow_types: Vec<String>, supported_activity_types: Vec<String>, max_concurrent_workflow_tasks: usize, max_concurrent_activity_tasks: usize, capabilities: Vec<String>, workflow_command_contracts: Value, ) -> Result<RegisterWorkerResponse>

Register a worker and advertise its named query and update handlers.

This Rust SDK cannot execute synchronous pre-accept update validation, so a workflow contract with a non-empty or malformed update_validators declaration returns Error::UnsupportedUpdateValidators before registration transport.

Source

pub async fn deregister_worker_registration( &self, worker_id: &str, ) -> Result<WorkerDeregistrationEnvelope>

Gracefully remove one worker’s registration through the worker plane.

This operation is separate from operator-facing worker management. It uses worker-protocol authentication and returns the server’s lease recovery result.

Source

pub async fn poll_query_task( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<Option<QueryTask>>

Long-poll for an ephemeral, read-only workflow query task.

Source

pub async fn poll_query_task_response( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<PollQueryTaskResponse>

Poll a query task while preserving server stop and drain metadata.

Source

pub async fn complete_query_task<T: Serialize>( &self, query_task_id: &str, lease_owner: &str, query_task_attempt: u64, result: T, codec: &str, ) -> Result<Value>

Complete a query task without appending workflow history.

Source

pub async fn fail_query_task( &self, query_task_id: &str, lease_owner: &str, query_task_attempt: u64, message: impl Into<String>, reason: impl Into<String>, failure_type: impl Into<String>, ) -> Result<Value>

Report a stable machine-readable query-task failure.

Source

pub async fn heartbeat_worker( &self, worker_id: &str, workflow_available: usize, activity_available: usize, ) -> Result<Value>

Source

pub async fn poll_workflow_task( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<Option<WorkflowTask>>

Source

pub async fn poll_workflow_task_response( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<PollWorkflowTaskResponse>

Source

pub async fn complete_workflow_task( &self, task_id: &str, lease_owner: &str, workflow_task_attempt: u64, commands: Vec<Value>, ) -> Result<Value>

Source

pub async fn fail_workflow_task( &self, task_id: &str, lease_owner: &str, workflow_task_attempt: u64, message: impl Into<String>, ) -> Result<Value>

Source

pub async fn poll_activity_task( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<Option<ActivityTask>>

Source

pub async fn poll_activity_task_response( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<PollActivityTaskResponse>

Poll an activity task while preserving server stop and drain metadata.

Source

pub async fn complete_activity_task<T: Serialize>( &self, task_id: &str, activity_attempt_id: &str, lease_owner: &str, result: T, codec: &str, ) -> Result<Value>

Source

pub async fn fail_activity_task( &self, task_id: &str, activity_attempt_id: &str, lease_owner: &str, message: impl Into<String>, non_retryable: bool, ) -> Result<Value>

Source

pub async fn heartbeat_activity_task<T: Serialize>( &self, task_id: &str, activity_attempt_id: &str, lease_owner: &str, details: T, ) -> Result<ActivityHeartbeatResponse>

Record a flat map of scalar/null activity progress metadata.

Heartbeat details use the Server’s bounded JSON progress contract. They are not workflow payload envelopes.

Trait Implementations§

Source§

impl Clone for Client

Source§

fn clone(&self) -> Client

Returns a copy of the value. Read more
1.0.0 · Source§

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

Performs copy-assignment from source. Read more
Source§

impl Debug for Client

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

§

impl Freeze for Client

§

impl !RefUnwindSafe for Client

§

impl Send for Client

§

impl Sync for Client

§

impl Unpin for Client

§

impl !UnwindSafe for Client

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<T> CloneToUninit for T
where T: Clone,

Source§

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

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

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

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

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

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.

§

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

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: 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
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: 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> 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 = 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.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

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
§

fn with_current_subscriber(self) -> WithDispatch<Self>

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

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

§

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

§

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