pub struct Client { /* private fields */ }Implementations§
Source§impl Client
impl Client
Sourcepub async fn open_cancellation_scope_on_claim(
&self,
task: &WorkflowTask,
sequence: u64,
parent_scope_id: &str,
shield_parent: bool,
) -> Result<CancellationScopeOpenReceipt>
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
impl Client
Sourcepub async fn acknowledge_activity_cancellation(
&self,
task_id: &str,
activity_attempt_id: &str,
lease_owner: &str,
request_id: &str,
) -> Result<Value>
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.
Sourcepub async fn activity_task_status(
&self,
task_id: &str,
activity_attempt_id: &str,
lease_owner: &str,
) -> Result<Value>
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.
Sourcepub async fn poll_cooperative_workflow_task(
&self,
worker_id: &str,
task_queue: &str,
timeout: Duration,
) -> Result<CooperativeWorkflowTaskPoll>
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.
Sourcepub async fn heartbeat_workflow_task(
&self,
task: &WorkflowTask,
original: Option<&CancellationRequest>,
) -> Result<WorkflowTaskHeartbeat>
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.
Sourcepub async fn deliver_workflow_cancellation(
&self,
task: &WorkflowTask,
delivery: &CancellationDelivery,
) -> Result<CancellationDeliveryReply>
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.
Sourcepub async fn refresh_workflow_cancellation_history(
&self,
task: &WorkflowTask,
observation: &CancellationRequest,
) -> Result<Vec<HistoryEvent>>
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.
Sourcepub async fn request_workflow_cancellation(
&self,
workflow_id: &str,
options: CooperativeCancellationOptions,
) -> Result<WorkflowCancellationRequest>
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.
Sourcepub async fn request_workflow_run_cancellation(
&self,
workflow_id: &str,
run_id: &str,
options: CooperativeCancellationOptions,
) -> Result<WorkflowCancellationRequest>
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
impl Client
Sourcepub async fn create_worker_session(
&self,
worker_id: &str,
options: &WorkerSessionOptions,
) -> Result<Value>
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§impl Client
impl Client
pub fn new(base_url: impl Into<String>) -> Result<Self>
pub fn builder(base_url: impl Into<String>) -> ClientBuilder
pub async fn health(&self) -> Result<Value>
pub async fn cluster_info(&self) -> Result<Value>
pub async fn start_workflow<T: Serialize>( &self, workflow_type: &str, task_queue: &str, workflow_id: &str, input: T, ) -> Result<WorkflowHandle>
Sourcepub async fn start_workflow_with_options<T: Serialize>(
&self,
workflow_type: &str,
task_queue: &str,
workflow_id: &str,
options: WorkflowStartOptions,
input: T,
) -> Result<WorkflowHandle>
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.
pub async fn signal_workflow<T: Serialize>( &self, workflow_id: &str, signal_name: &str, input: T, ) -> Result<Value>
Sourcepub async fn append_message_stream<T: Serialize>(
&self,
workflow_id: &str,
stream_name: &str,
message_id: &str,
input: T,
) -> Result<Value>
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.
Sourcepub async fn signal_workflow_run<T: Serialize>(
&self,
workflow_id: &str,
run_id: &str,
signal_name: &str,
input: T,
) -> Result<Value>
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.
Sourcepub async fn cancel_workflow(
&self,
workflow_id: &str,
options: WorkflowCommandOptions,
) -> Result<WorkflowCommandResult>
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.
Sourcepub async fn cancel_workflow_run(
&self,
workflow_id: &str,
run_id: &str,
options: WorkflowCommandOptions,
) -> Result<WorkflowCommandResult>
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.
Sourcepub async fn terminate_workflow(
&self,
workflow_id: &str,
options: WorkflowCommandOptions,
) -> Result<WorkflowCommandResult>
pub async fn terminate_workflow( &self, workflow_id: &str, options: WorkflowCommandOptions, ) -> Result<WorkflowCommandResult>
Forcefully terminate the current run for an instance.
Sourcepub async fn terminate_workflow_run(
&self,
workflow_id: &str,
run_id: &str,
options: WorkflowCommandOptions,
) -> Result<WorkflowCommandResult>
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.
Sourcepub async fn redrive_workflow_run(
&self,
workflow_id: &str,
failed_run_id: &str,
request_id: Option<&str>,
) -> Result<WorkflowRedriveResult>
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.
Sourcepub async fn query_workflow<T: Serialize>(
&self,
workflow_id: &str,
query_name: &str,
input: T,
) -> Result<Value>
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.
Sourcepub async fn query_workflow_run<T: Serialize>(
&self,
workflow_id: &str,
run_id: &str,
query_name: &str,
input: T,
) -> Result<Value>
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.
Sourcepub async fn query_workflow_avro_value<T: Serialize>(
&self,
workflow_id: &str,
query_name: &str,
input: T,
) -> Result<AvroValue>
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.
Sourcepub async fn query_workflow_run_avro_value<T: Serialize>(
&self,
workflow_id: &str,
run_id: &str,
query_name: &str,
input: T,
) -> Result<AvroValue>
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.
Sourcepub async fn update_workflow<T: Serialize>(
&self,
workflow_id: &str,
update_name: &str,
input: T,
request_id: Option<&str>,
) -> Result<Value>
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.
Sourcepub async fn update_workflow_avro_value<T: Serialize>(
&self,
workflow_id: &str,
update_name: &str,
input: T,
request_id: Option<&str>,
) -> Result<AvroValue>
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.
pub async fn describe_workflow( &self, workflow_id: &str, ) -> Result<WorkflowDescription>
Sourcepub async fn describe_workflow_run(
&self,
workflow_id: &str,
run_id: &str,
) -> Result<WorkflowDescription>
pub async fn describe_workflow_run( &self, workflow_id: &str, run_id: &str, ) -> Result<WorkflowDescription>
Describe one selected run, including historical terminal runs.
Sourcepub async fn list_workflow_streams(
&self,
workflow_id: &str,
run_id: &str,
) -> Result<Vec<WorkflowStreamDescription>>
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.
Sourcepub async fn describe_workflow_stream(
&self,
workflow_id: &str,
run_id: &str,
stream_name: &str,
) -> Result<WorkflowStreamDescription>
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.
Sourcepub 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>
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.
Sourcepub async fn append_workflow_stream(
&self,
workflow_id: &str,
run_id: &str,
stream_name: &str,
items: &[WorkflowStreamAppendItem],
max_pending_items: Option<u64>,
) -> Result<WorkflowStreamAppendResult>
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.
Sourcepub 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>
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.
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>
Sourcepub 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>
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.
Sourcepub 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>
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.
Sourcepub async fn deregister_worker_registration(
&self,
worker_id: &str,
) -> Result<WorkerDeregistrationEnvelope>
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.
Sourcepub async fn poll_query_task(
&self,
worker_id: &str,
task_queue: &str,
timeout: Duration,
) -> Result<Option<QueryTask>>
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.
Sourcepub async fn poll_query_task_response(
&self,
worker_id: &str,
task_queue: &str,
timeout: Duration,
) -> Result<PollQueryTaskResponse>
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.
Sourcepub 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>
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.
Sourcepub 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>
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.
pub async fn heartbeat_worker( &self, worker_id: &str, workflow_available: usize, activity_available: usize, ) -> Result<Value>
pub async fn poll_workflow_task( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<Option<WorkflowTask>>
pub async fn poll_workflow_task_response( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<PollWorkflowTaskResponse>
pub async fn complete_workflow_task( &self, task_id: &str, lease_owner: &str, workflow_task_attempt: u64, commands: Vec<Value>, ) -> Result<Value>
pub async fn fail_workflow_task( &self, task_id: &str, lease_owner: &str, workflow_task_attempt: u64, message: impl Into<String>, ) -> Result<Value>
pub async fn poll_activity_task( &self, worker_id: &str, task_queue: &str, timeout: Duration, ) -> Result<Option<ActivityTask>>
Sourcepub async fn poll_activity_task_response(
&self,
worker_id: &str,
task_queue: &str,
timeout: Duration,
) -> Result<PollActivityTaskResponse>
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.
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>
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>
Sourcepub async fn heartbeat_activity_task<T: Serialize>(
&self,
task_id: &str,
activity_attempt_id: &str,
lease_owner: &str,
details: T,
) -> Result<ActivityHeartbeatResponse>
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.