1use crate::agent::AgentOutput;
6use crate::app::{AgentToolRef, RequestContext};
7use crate::codec::host_service::{HostServiceChannel, connect_host_service, plain_channel};
8use crate::codec::workflow::{
9 from_wire_get_workflow_provider_run_events_response,
10 from_wire_get_workflow_provider_run_output_response,
11 from_wire_list_workflow_provider_definitions_response,
12 from_wire_list_workflow_provider_runs_response, from_wire_signal_workflow_run_response,
13 from_wire_workflow_definition, from_wire_workflow_event, from_wire_workflow_run,
14 to_wire_apply_workflow_provider_definition_request,
15 to_wire_cancel_workflow_provider_run_request,
16 to_wire_delete_workflow_provider_definition_request,
17 to_wire_deliver_workflow_provider_event_request,
18 to_wire_get_workflow_provider_definition_request,
19 to_wire_get_workflow_provider_run_events_request,
20 to_wire_get_workflow_provider_run_output_request, to_wire_get_workflow_provider_run_request,
21 to_wire_list_workflow_provider_definitions_request,
22 to_wire_list_workflow_provider_runs_request,
23 to_wire_set_workflow_provider_activation_paused_request,
24 to_wire_set_workflow_provider_definition_paused_request,
25 to_wire_signal_or_start_workflow_provider_run_request,
26 to_wire_signal_workflow_provider_run_request, to_wire_start_workflow_provider_run_request,
27};
28use crate::generated::v1;
29use crate::rpc_support::GestaltError;
30
31pub type WorkflowRunStatus = i32;
33
34pub mod workflow_run_status {
36 pub const WORKFLOW_RUN_STATUS_UNSPECIFIED: i32 = 0;
38 pub const WORKFLOW_RUN_STATUS_PENDING: i32 = 1;
40 pub const WORKFLOW_RUN_STATUS_RUNNING: i32 = 2;
42 pub const WORKFLOW_RUN_STATUS_SUCCEEDED: i32 = 3;
44 pub const WORKFLOW_RUN_STATUS_FAILED: i32 = 4;
46 pub const WORKFLOW_RUN_STATUS_CANCELED: i32 = 5;
48}
49
50pub type WorkflowStepStatus = i32;
52
53pub mod workflow_step_status {
55 pub const WORKFLOW_STEP_STATUS_UNSPECIFIED: i32 = 0;
57 pub const WORKFLOW_STEP_STATUS_PENDING: i32 = 1;
59 pub const WORKFLOW_STEP_STATUS_RUNNING: i32 = 2;
61 pub const WORKFLOW_STEP_STATUS_SKIPPED: i32 = 3;
63 pub const WORKFLOW_STEP_STATUS_SUCCEEDED: i32 = 4;
65 pub const WORKFLOW_STEP_STATUS_FAILED: i32 = 5;
67 pub const WORKFLOW_STEP_STATUS_UNKNOWN: i32 = 6;
69}
70
71#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
73#[serde(rename_all = "camelCase")]
74pub struct ApplyWorkflowProviderDefinitionRequest {
75 pub provider: String,
77 pub spec: Option<WorkflowDefinitionSpec>,
79 pub idempotency_key: String,
81 pub context: Option<RequestContext>,
83}
84
85#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
87#[serde(rename_all = "camelCase")]
88pub struct BoundWorkflowTarget {
89 pub steps: Vec<WorkflowStep>,
91}
92
93#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
95#[serde(rename_all = "camelCase")]
96pub struct CancelWorkflowProviderRunRequest {
97 pub run_id: String,
99 pub reason: String,
101 pub context: Option<RequestContext>,
103 pub provider: String,
105}
106
107#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
109#[serde(rename_all = "camelCase")]
110pub struct DeleteWorkflowProviderDefinitionRequest {
111 pub definition_id: String,
113 pub context: Option<RequestContext>,
115 pub provider: String,
117}
118
119#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
121#[serde(rename_all = "camelCase")]
122pub struct DeliverWorkflowProviderEventRequest {
123 pub provider: String,
125 pub event: Option<WorkflowEvent>,
127 pub context: Option<RequestContext>,
129}
130
131#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
133#[serde(rename_all = "camelCase")]
134pub struct GetWorkflowProviderDefinitionRequest {
135 pub definition_id: String,
137 pub context: Option<RequestContext>,
139 pub provider: String,
141}
142
143#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
145#[serde(rename_all = "camelCase")]
146pub struct GetWorkflowProviderRunEventsRequest {
147 pub run_id: String,
149 pub context: Option<RequestContext>,
151 pub provider: String,
153}
154
155#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
157#[serde(rename_all = "camelCase")]
158pub struct GetWorkflowProviderRunEventsResponse {
159 pub events: Vec<WorkflowRunEvent>,
161}
162
163#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
165#[serde(rename_all = "camelCase")]
166pub struct GetWorkflowProviderRunOutputRequest {
167 pub run_id: String,
169 pub context: Option<RequestContext>,
171 pub provider: String,
173}
174
175#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
177#[serde(rename_all = "camelCase")]
178pub struct GetWorkflowProviderRunOutputResponse {
179 pub output: Option<serde_json::Value>,
181}
182
183#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
185#[serde(rename_all = "camelCase")]
186pub struct GetWorkflowProviderRunRequest {
187 pub run_id: String,
189 pub context: Option<RequestContext>,
191 pub provider: String,
193}
194
195#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
197#[serde(rename_all = "camelCase")]
198pub struct ListWorkflowProviderDefinitionsRequest {
199 pub context: Option<RequestContext>,
201 pub provider: String,
203}
204
205#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
207#[serde(rename_all = "camelCase")]
208pub struct ListWorkflowProviderDefinitionsResponse {
209 pub definitions: Vec<WorkflowDefinition>,
211}
212
213#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
215#[serde(rename_all = "camelCase")]
216pub struct ListWorkflowProviderRunsRequest {
217 pub page_size: i32,
219 pub page_token: String,
221 pub status: WorkflowRunStatus,
223 pub target_app: String,
225 pub context: Option<RequestContext>,
227 pub provider: String,
229}
230
231#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
233#[serde(rename_all = "camelCase")]
234pub struct ListWorkflowProviderRunsResponse {
235 pub runs: Vec<WorkflowRun>,
237 pub next_page_token: String,
239}
240
241#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
243#[serde(rename_all = "camelCase")]
244pub struct SetWorkflowProviderActivationPausedRequest {
245 pub definition_id: String,
247 pub activation_id: String,
249 pub paused: bool,
251 pub context: Option<RequestContext>,
253 pub provider: String,
255}
256
257#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
259#[serde(rename_all = "camelCase")]
260pub struct SetWorkflowProviderDefinitionPausedRequest {
261 pub definition_id: String,
263 pub paused: bool,
265 pub context: Option<RequestContext>,
267 pub provider: String,
269}
270
271#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
273#[serde(rename_all = "camelCase")]
274pub struct SignalOrStartWorkflowProviderRunRequest {
275 pub workflow_key: String,
277 pub idempotency_key: String,
279 pub signal: Option<WorkflowSignal>,
281 pub provider: String,
283 pub definition_id: String,
285 pub input: Option<serde_json::Map<String, serde_json::Value>>,
287 pub expected_definition_generation: i64,
289 pub context: Option<RequestContext>,
291}
292
293#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
295#[serde(rename_all = "camelCase")]
296pub struct SignalWorkflowProviderRunRequest {
297 pub run_id: String,
299 pub signal: Option<WorkflowSignal>,
301 pub context: Option<RequestContext>,
303 pub provider: String,
305}
306
307#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
309#[serde(rename_all = "camelCase")]
310pub struct SignalWorkflowRunResponse {
311 pub run: Option<WorkflowRun>,
313 pub signal: Option<WorkflowSignal>,
315 pub started_run: bool,
317 pub workflow_key: String,
319}
320
321#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
323#[serde(rename_all = "camelCase")]
324pub struct StartWorkflowProviderRunRequest {
325 pub idempotency_key: String,
327 pub workflow_key: String,
329 pub provider: String,
331 pub definition_id: String,
333 pub input: Option<serde_json::Map<String, serde_json::Value>>,
335 pub expected_definition_generation: i64,
337 pub context: Option<RequestContext>,
339}
340
341#[allow(clippy::enum_variant_names, clippy::large_enum_variant)]
343#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
344pub enum WorkflowActivationTrigger {
345 Schedule(WorkflowScheduleActivation),
347 Event(WorkflowEventActivation),
349}
350
351#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
353#[serde(rename_all = "camelCase")]
354pub struct WorkflowActivation {
355 pub id: String,
357 pub input: Option<WorkflowValue>,
359 pub paused: bool,
361 pub trigger: Option<WorkflowActivationTrigger>,
363}
364
365#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
367#[serde(rename_all = "camelCase")]
368pub struct WorkflowAgentMessage {
369 pub role: String,
371 pub text: Option<WorkflowText>,
373 pub metadata: Option<serde_json::Map<String, serde_json::Value>>,
375}
376
377#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
379#[serde(rename_all = "camelCase")]
380pub struct WorkflowArray {
381 pub values: Vec<WorkflowValue>,
383}
384
385#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
387#[serde(rename_all = "camelCase")]
388pub struct WorkflowDefinition {
389 pub id: String,
391 pub generation: i64,
393 pub target: Option<BoundWorkflowTarget>,
395 pub activations: Vec<WorkflowActivation>,
397 pub paused: bool,
399 #[serde(with = "crate::serde_time")]
400 pub created_at: Option<std::time::SystemTime>,
402 #[serde(with = "crate::serde_time")]
403 pub updated_at: Option<std::time::SystemTime>,
405 pub provider: String,
407 pub run_as: String,
409}
410
411#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
413#[serde(rename_all = "camelCase")]
414pub struct WorkflowDefinitionSpec {
415 pub id: String,
417 pub target: Option<BoundWorkflowTarget>,
419 pub activations: Vec<WorkflowActivation>,
421 pub paused: bool,
423 pub run_as: String,
425}
426
427#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
429#[serde(rename_all = "camelCase")]
430pub struct WorkflowEvent {
431 pub id: String,
433 pub source: String,
435 pub spec_version: String,
437 pub r#type: String,
439 pub subject: String,
441 #[serde(with = "crate::serde_time")]
442 pub time: Option<std::time::SystemTime>,
444 pub datacontenttype: String,
446 pub data: Option<serde_json::Map<String, serde_json::Value>>,
448 pub extensions: std::collections::BTreeMap<String, serde_json::Value>,
450}
451
452#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
454#[serde(rename_all = "camelCase")]
455pub struct WorkflowEventActivation {
456 pub r#match: Option<WorkflowEventMatch>,
458}
459
460#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
462#[serde(rename_all = "camelCase")]
463pub struct WorkflowEventMatch {
464 pub r#type: String,
466 pub source: String,
468 pub subject: String,
470}
471
472#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
474#[serde(rename_all = "camelCase")]
475pub struct WorkflowEventTriggerInvocation {
476 pub activation_id: String,
478 pub event: Option<WorkflowEvent>,
480}
481
482#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
484#[serde(rename_all = "camelCase")]
485pub struct WorkflowManualTrigger {}
486
487#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
489#[serde(rename_all = "camelCase")]
490pub struct WorkflowObject {
491 pub fields: std::collections::BTreeMap<String, WorkflowValue>,
493}
494
495#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
497#[serde(rename_all = "camelCase")]
498pub struct WorkflowPathSource {
499 pub path: String,
501}
502
503#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
505#[serde(rename_all = "camelCase")]
506pub struct WorkflowRun {
507 pub id: String,
509 pub status: WorkflowRunStatus,
511 pub target: Option<BoundWorkflowTarget>,
513 pub trigger: Option<WorkflowRunTrigger>,
515 #[serde(with = "crate::serde_time")]
516 pub created_at: Option<std::time::SystemTime>,
518 #[serde(with = "crate::serde_time")]
519 pub started_at: Option<std::time::SystemTime>,
521 #[serde(with = "crate::serde_time")]
522 pub completed_at: Option<std::time::SystemTime>,
524 pub status_message: String,
526 pub output: Option<serde_json::Value>,
528 pub workflow_key: String,
530 pub provider: String,
532 pub definition_id: String,
534 pub run_as: String,
536 pub input: Option<serde_json::Map<String, serde_json::Value>>,
538 pub definition_generation: i64,
540 pub current_step_id: String,
542 pub steps: Vec<WorkflowStepExecution>,
544}
545
546#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
548#[serde(rename_all = "camelCase")]
549pub struct WorkflowRunEvent {
550 pub id: String,
552 pub run_id: String,
554 pub step_id: String,
556 pub r#type: String,
558 pub data: Option<serde_json::Map<String, serde_json::Value>>,
560 #[serde(with = "crate::serde_time")]
561 pub created_at: Option<std::time::SystemTime>,
563}
564
565#[allow(clippy::enum_variant_names, clippy::large_enum_variant)]
567#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
568pub enum WorkflowRunTriggerKind {
569 Manual(WorkflowManualTrigger),
571 Schedule(WorkflowScheduleTrigger),
573 Event(WorkflowEventTriggerInvocation),
575}
576
577#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
579#[serde(rename_all = "camelCase")]
580pub struct WorkflowRunTrigger {
581 pub kind: Option<WorkflowRunTriggerKind>,
583}
584
585#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
587#[serde(rename_all = "camelCase")]
588pub struct WorkflowScheduleActivation {
589 pub cron: String,
591 pub timezone: String,
593}
594
595#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
597#[serde(rename_all = "camelCase")]
598pub struct WorkflowScheduleTrigger {
599 pub activation_id: String,
601 #[serde(with = "crate::serde_time")]
602 pub scheduled_for: Option<std::time::SystemTime>,
604}
605
606#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
608#[serde(rename_all = "camelCase")]
609pub struct WorkflowSignal {
610 pub id: String,
612 pub name: String,
614 pub payload: Option<serde_json::Map<String, serde_json::Value>>,
616 pub metadata: Option<serde_json::Map<String, serde_json::Value>>,
618 #[serde(with = "crate::serde_time")]
619 pub created_at: Option<std::time::SystemTime>,
621 pub idempotency_key: String,
623 pub sequence: i64,
625}
626
627#[allow(clippy::enum_variant_names, clippy::large_enum_variant)]
629#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
630pub enum WorkflowStepAction {
631 App(WorkflowStepAppCall),
633 Agent(WorkflowStepAgentTurn),
635}
636
637#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
639#[serde(rename_all = "camelCase")]
640pub struct WorkflowStep {
641 pub id: String,
643 pub inputs: std::collections::BTreeMap<String, WorkflowValue>,
645 pub when: Option<WorkflowStepWhen>,
647 pub timeout_seconds: i32,
649 pub metadata: Option<serde_json::Map<String, serde_json::Value>>,
651 pub action: Option<WorkflowStepAction>,
653}
654
655#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
657#[serde(rename_all = "camelCase")]
658pub struct WorkflowStepAgentTurn {
659 pub provider: String,
661 pub model: String,
663 pub session_key: String,
665 pub prompt: Option<WorkflowText>,
667 pub messages: Vec<WorkflowAgentMessage>,
669 pub tools: Vec<AgentToolRef>,
671 pub output: Option<AgentOutput>,
673 pub model_options: Option<serde_json::Map<String, serde_json::Value>>,
675}
676
677#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
679#[serde(rename_all = "camelCase")]
680pub struct WorkflowStepAppCall {
681 pub name: String,
683 pub operation: String,
685 pub input: Option<WorkflowValue>,
687 pub connection: String,
689 pub instance: String,
691 pub credential_mode: String,
693}
694
695#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
697#[serde(rename_all = "camelCase")]
698pub struct WorkflowStepAttempt {
699 pub id: String,
701 pub status: WorkflowStepStatus,
703 pub idempotency_key: String,
705 pub input: Option<serde_json::Value>,
707 pub output: Option<serde_json::Value>,
709 pub status_message: String,
711 #[serde(with = "crate::serde_time")]
712 pub started_at: Option<std::time::SystemTime>,
714 #[serde(with = "crate::serde_time")]
715 pub completed_at: Option<std::time::SystemTime>,
717}
718
719#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
721#[serde(rename_all = "camelCase")]
722pub struct WorkflowStepExecution {
723 pub step_id: String,
725 pub status: WorkflowStepStatus,
727 pub attempts: Vec<WorkflowStepAttempt>,
729 pub input: Option<serde_json::Value>,
731 pub output: Option<serde_json::Value>,
733 pub status_message: String,
735 pub skip_reason: String,
737 #[serde(with = "crate::serde_time")]
738 pub started_at: Option<std::time::SystemTime>,
740 #[serde(with = "crate::serde_time")]
741 pub completed_at: Option<std::time::SystemTime>,
743}
744
745#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
747#[serde(rename_all = "camelCase")]
748pub struct WorkflowStepInputSource {
749 pub step_id: String,
751 pub path: String,
753}
754
755#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
757#[serde(rename_all = "camelCase")]
758pub struct WorkflowStepOutputSource {
759 pub step_id: String,
761 pub path: String,
763}
764
765#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
767#[serde(rename_all = "camelCase")]
768pub struct WorkflowStepWhen {
769 pub value: Option<WorkflowValue>,
771 pub equals: Option<serde_json::Value>,
773}
774
775#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
777#[serde(rename_all = "camelCase")]
778pub struct WorkflowText {
779 pub template: String,
781}
782
783#[allow(clippy::enum_variant_names, clippy::large_enum_variant)]
785#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
786pub enum WorkflowValueKind {
787 Literal(serde_json::Value),
789 Object(WorkflowObject),
791 Array(WorkflowArray),
793 Template(WorkflowText),
795 Input(WorkflowPathSource),
797 Signal(WorkflowPathSource),
799 StepOutput(WorkflowStepOutputSource),
801 StepInput(WorkflowStepInputSource),
803}
804
805#[derive(Clone, Debug, Default, PartialEq, serde::Serialize, serde::Deserialize)]
807#[serde(rename_all = "camelCase")]
808pub struct WorkflowValue {
809 pub kind: Option<WorkflowValueKind>,
811}
812
813pub struct Workflow {
815 inner: v1::workflow_client::WorkflowClient<HostServiceChannel>,
816 timeout: Option<std::time::Duration>,
817 context: Option<RequestContext>,
818}
819
820impl Workflow {
821 pub fn new(channel: tonic::transport::Channel) -> Self {
823 Self {
824 inner: v1::workflow_client::WorkflowClient::new(plain_channel(channel)),
825 timeout: None,
826 context: None,
827 }
828 }
829
830 pub fn with_timeout(mut self, timeout: std::time::Duration) -> Self {
833 self.timeout = Some(timeout);
834 self
835 }
836
837 pub fn with_context(mut self, context: RequestContext) -> Self {
840 self.context = Some(context);
841 self
842 }
843
844 pub async fn connect() -> Result<Self, GestaltError> {
846 Self::connect_named("").await
847 }
848
849 pub async fn connect_named(name: &str) -> Result<Self, GestaltError> {
851 Ok(Self {
852 inner: v1::workflow_client::WorkflowClient::new(
853 connect_host_service("workflow", name).await?,
854 ),
855 timeout: None,
856 context: None,
857 })
858 }
859
860 pub async fn apply_definition(
862 &mut self,
863 provider: String,
864 idempotency_key: String,
865 spec: Option<WorkflowDefinitionSpec>,
866 ) -> Result<WorkflowDefinition, GestaltError> {
867 let request = ApplyWorkflowProviderDefinitionRequest {
868 provider,
869 idempotency_key,
870 spec,
871 context: self.context.clone(),
872 };
873 let mut tonic_request =
874 tonic::Request::new(to_wire_apply_workflow_provider_definition_request(request));
875 if let Some(timeout) = self.timeout {
876 tonic_request.set_timeout(timeout);
877 }
878 let response = self.inner.apply_definition(tonic_request).await?;
879 Ok(from_wire_workflow_definition(response.into_inner()))
880 }
881
882 pub async fn apply_definition_raw(
884 &mut self,
885 request: ApplyWorkflowProviderDefinitionRequest,
886 ) -> Result<WorkflowDefinition, GestaltError> {
887 let mut request = request;
888 if request.context.is_none() {
889 request.context = self.context.clone();
890 }
891 let mut tonic_request =
892 tonic::Request::new(to_wire_apply_workflow_provider_definition_request(request));
893 if let Some(timeout) = self.timeout {
894 tonic_request.set_timeout(timeout);
895 }
896 let response = self.inner.apply_definition(tonic_request).await?;
897 Ok(from_wire_workflow_definition(response.into_inner()))
898 }
899
900 pub async fn get_definition(
902 &mut self,
903 provider: String,
904 definition_id: String,
905 ) -> Result<WorkflowDefinition, GestaltError> {
906 let request = GetWorkflowProviderDefinitionRequest {
907 provider,
908 definition_id,
909 context: self.context.clone(),
910 };
911 let mut tonic_request =
912 tonic::Request::new(to_wire_get_workflow_provider_definition_request(request));
913 if let Some(timeout) = self.timeout {
914 tonic_request.set_timeout(timeout);
915 }
916 let response = self.inner.get_definition(tonic_request).await?;
917 Ok(from_wire_workflow_definition(response.into_inner()))
918 }
919
920 pub async fn get_definition_raw(
922 &mut self,
923 request: GetWorkflowProviderDefinitionRequest,
924 ) -> Result<WorkflowDefinition, GestaltError> {
925 let mut request = request;
926 if request.context.is_none() {
927 request.context = self.context.clone();
928 }
929 let mut tonic_request =
930 tonic::Request::new(to_wire_get_workflow_provider_definition_request(request));
931 if let Some(timeout) = self.timeout {
932 tonic_request.set_timeout(timeout);
933 }
934 let response = self.inner.get_definition(tonic_request).await?;
935 Ok(from_wire_workflow_definition(response.into_inner()))
936 }
937
938 pub async fn list_definitions(
940 &mut self,
941 provider: String,
942 ) -> Result<Vec<WorkflowDefinition>, GestaltError> {
943 let request = ListWorkflowProviderDefinitionsRequest {
944 provider,
945 context: self.context.clone(),
946 };
947 let mut tonic_request =
948 tonic::Request::new(to_wire_list_workflow_provider_definitions_request(request));
949 if let Some(timeout) = self.timeout {
950 tonic_request.set_timeout(timeout);
951 }
952 let response = from_wire_list_workflow_provider_definitions_response(
953 self.inner
954 .list_definitions(tonic_request)
955 .await?
956 .into_inner(),
957 );
958 Ok(response.definitions)
959 }
960
961 pub async fn list_definitions_raw(
963 &mut self,
964 request: ListWorkflowProviderDefinitionsRequest,
965 ) -> Result<ListWorkflowProviderDefinitionsResponse, GestaltError> {
966 let mut request = request;
967 if request.context.is_none() {
968 request.context = self.context.clone();
969 }
970 let mut tonic_request =
971 tonic::Request::new(to_wire_list_workflow_provider_definitions_request(request));
972 if let Some(timeout) = self.timeout {
973 tonic_request.set_timeout(timeout);
974 }
975 let response = self.inner.list_definitions(tonic_request).await?;
976 Ok(from_wire_list_workflow_provider_definitions_response(
977 response.into_inner(),
978 ))
979 }
980
981 pub async fn set_definition_paused(
983 &mut self,
984 provider: String,
985 definition_id: String,
986 paused: bool,
987 ) -> Result<WorkflowDefinition, GestaltError> {
988 let request = SetWorkflowProviderDefinitionPausedRequest {
989 provider,
990 definition_id,
991 paused,
992 context: self.context.clone(),
993 };
994 let mut tonic_request = tonic::Request::new(
995 to_wire_set_workflow_provider_definition_paused_request(request),
996 );
997 if let Some(timeout) = self.timeout {
998 tonic_request.set_timeout(timeout);
999 }
1000 let response = self.inner.set_definition_paused(tonic_request).await?;
1001 Ok(from_wire_workflow_definition(response.into_inner()))
1002 }
1003
1004 pub async fn set_definition_paused_raw(
1006 &mut self,
1007 request: SetWorkflowProviderDefinitionPausedRequest,
1008 ) -> Result<WorkflowDefinition, GestaltError> {
1009 let mut request = request;
1010 if request.context.is_none() {
1011 request.context = self.context.clone();
1012 }
1013 let mut tonic_request = tonic::Request::new(
1014 to_wire_set_workflow_provider_definition_paused_request(request),
1015 );
1016 if let Some(timeout) = self.timeout {
1017 tonic_request.set_timeout(timeout);
1018 }
1019 let response = self.inner.set_definition_paused(tonic_request).await?;
1020 Ok(from_wire_workflow_definition(response.into_inner()))
1021 }
1022
1023 pub async fn set_activation_paused(
1025 &mut self,
1026 provider: String,
1027 definition_id: String,
1028 activation_id: String,
1029 paused: bool,
1030 ) -> Result<WorkflowDefinition, GestaltError> {
1031 let request = SetWorkflowProviderActivationPausedRequest {
1032 provider,
1033 definition_id,
1034 activation_id,
1035 paused,
1036 context: self.context.clone(),
1037 };
1038 let mut tonic_request = tonic::Request::new(
1039 to_wire_set_workflow_provider_activation_paused_request(request),
1040 );
1041 if let Some(timeout) = self.timeout {
1042 tonic_request.set_timeout(timeout);
1043 }
1044 let response = self.inner.set_activation_paused(tonic_request).await?;
1045 Ok(from_wire_workflow_definition(response.into_inner()))
1046 }
1047
1048 pub async fn set_activation_paused_raw(
1050 &mut self,
1051 request: SetWorkflowProviderActivationPausedRequest,
1052 ) -> Result<WorkflowDefinition, GestaltError> {
1053 let mut request = request;
1054 if request.context.is_none() {
1055 request.context = self.context.clone();
1056 }
1057 let mut tonic_request = tonic::Request::new(
1058 to_wire_set_workflow_provider_activation_paused_request(request),
1059 );
1060 if let Some(timeout) = self.timeout {
1061 tonic_request.set_timeout(timeout);
1062 }
1063 let response = self.inner.set_activation_paused(tonic_request).await?;
1064 Ok(from_wire_workflow_definition(response.into_inner()))
1065 }
1066
1067 pub async fn delete_definition(
1069 &mut self,
1070 provider: String,
1071 definition_id: String,
1072 ) -> Result<(), GestaltError> {
1073 let request = DeleteWorkflowProviderDefinitionRequest {
1074 provider,
1075 definition_id,
1076 context: self.context.clone(),
1077 };
1078 let mut tonic_request =
1079 tonic::Request::new(to_wire_delete_workflow_provider_definition_request(request));
1080 if let Some(timeout) = self.timeout {
1081 tonic_request.set_timeout(timeout);
1082 }
1083 self.inner.delete_definition(tonic_request).await?;
1084 Ok(())
1085 }
1086
1087 pub async fn delete_definition_raw(
1089 &mut self,
1090 request: DeleteWorkflowProviderDefinitionRequest,
1091 ) -> Result<(), GestaltError> {
1092 let mut request = request;
1093 if request.context.is_none() {
1094 request.context = self.context.clone();
1095 }
1096 let mut tonic_request =
1097 tonic::Request::new(to_wire_delete_workflow_provider_definition_request(request));
1098 if let Some(timeout) = self.timeout {
1099 tonic_request.set_timeout(timeout);
1100 }
1101 self.inner.delete_definition(tonic_request).await?;
1102 Ok(())
1103 }
1104
1105 pub async fn start_run(
1107 &mut self,
1108 idempotency_key: String,
1109 workflow_key: String,
1110 provider: String,
1111 definition_id: String,
1112 expected_definition_generation: i64,
1113 input: Option<serde_json::Map<String, serde_json::Value>>,
1114 ) -> Result<WorkflowRun, GestaltError> {
1115 let request = StartWorkflowProviderRunRequest {
1116 idempotency_key,
1117 workflow_key,
1118 provider,
1119 definition_id,
1120 expected_definition_generation,
1121 input,
1122 context: self.context.clone(),
1123 };
1124 let mut tonic_request =
1125 tonic::Request::new(to_wire_start_workflow_provider_run_request(request));
1126 if let Some(timeout) = self.timeout {
1127 tonic_request.set_timeout(timeout);
1128 }
1129 let response = self.inner.start_run(tonic_request).await?;
1130 Ok(from_wire_workflow_run(response.into_inner()))
1131 }
1132
1133 pub async fn start_run_raw(
1135 &mut self,
1136 request: StartWorkflowProviderRunRequest,
1137 ) -> Result<WorkflowRun, GestaltError> {
1138 let mut request = request;
1139 if request.context.is_none() {
1140 request.context = self.context.clone();
1141 }
1142 let mut tonic_request =
1143 tonic::Request::new(to_wire_start_workflow_provider_run_request(request));
1144 if let Some(timeout) = self.timeout {
1145 tonic_request.set_timeout(timeout);
1146 }
1147 let response = self.inner.start_run(tonic_request).await?;
1148 Ok(from_wire_workflow_run(response.into_inner()))
1149 }
1150
1151 pub async fn list_runs(
1153 &mut self,
1154 provider: String,
1155 page_size: i32,
1156 page_token: String,
1157 status: WorkflowRunStatus,
1158 target_app: String,
1159 ) -> Result<ListWorkflowProviderRunsResponse, GestaltError> {
1160 let request = ListWorkflowProviderRunsRequest {
1161 provider,
1162 page_size,
1163 page_token,
1164 status,
1165 target_app,
1166 context: self.context.clone(),
1167 };
1168 let mut tonic_request =
1169 tonic::Request::new(to_wire_list_workflow_provider_runs_request(request));
1170 if let Some(timeout) = self.timeout {
1171 tonic_request.set_timeout(timeout);
1172 }
1173 let response = self.inner.list_runs(tonic_request).await?;
1174 Ok(from_wire_list_workflow_provider_runs_response(
1175 response.into_inner(),
1176 ))
1177 }
1178
1179 pub async fn list_runs_raw(
1181 &mut self,
1182 request: ListWorkflowProviderRunsRequest,
1183 ) -> Result<ListWorkflowProviderRunsResponse, GestaltError> {
1184 let mut request = request;
1185 if request.context.is_none() {
1186 request.context = self.context.clone();
1187 }
1188 let mut tonic_request =
1189 tonic::Request::new(to_wire_list_workflow_provider_runs_request(request));
1190 if let Some(timeout) = self.timeout {
1191 tonic_request.set_timeout(timeout);
1192 }
1193 let response = self.inner.list_runs(tonic_request).await?;
1194 Ok(from_wire_list_workflow_provider_runs_response(
1195 response.into_inner(),
1196 ))
1197 }
1198
1199 pub async fn get_run(
1201 &mut self,
1202 provider: String,
1203 run_id: String,
1204 ) -> Result<WorkflowRun, GestaltError> {
1205 let request = GetWorkflowProviderRunRequest {
1206 provider,
1207 run_id,
1208 context: self.context.clone(),
1209 };
1210 let mut tonic_request =
1211 tonic::Request::new(to_wire_get_workflow_provider_run_request(request));
1212 if let Some(timeout) = self.timeout {
1213 tonic_request.set_timeout(timeout);
1214 }
1215 let response = self.inner.get_run(tonic_request).await?;
1216 Ok(from_wire_workflow_run(response.into_inner()))
1217 }
1218
1219 pub async fn get_run_raw(
1221 &mut self,
1222 request: GetWorkflowProviderRunRequest,
1223 ) -> Result<WorkflowRun, GestaltError> {
1224 let mut request = request;
1225 if request.context.is_none() {
1226 request.context = self.context.clone();
1227 }
1228 let mut tonic_request =
1229 tonic::Request::new(to_wire_get_workflow_provider_run_request(request));
1230 if let Some(timeout) = self.timeout {
1231 tonic_request.set_timeout(timeout);
1232 }
1233 let response = self.inner.get_run(tonic_request).await?;
1234 Ok(from_wire_workflow_run(response.into_inner()))
1235 }
1236
1237 pub async fn get_run_events(
1239 &mut self,
1240 provider: String,
1241 run_id: String,
1242 ) -> Result<Vec<WorkflowRunEvent>, GestaltError> {
1243 let request = GetWorkflowProviderRunEventsRequest {
1244 provider,
1245 run_id,
1246 context: self.context.clone(),
1247 };
1248 let mut tonic_request =
1249 tonic::Request::new(to_wire_get_workflow_provider_run_events_request(request));
1250 if let Some(timeout) = self.timeout {
1251 tonic_request.set_timeout(timeout);
1252 }
1253 let response = from_wire_get_workflow_provider_run_events_response(
1254 self.inner.get_run_events(tonic_request).await?.into_inner(),
1255 );
1256 Ok(response.events)
1257 }
1258
1259 pub async fn get_run_events_raw(
1261 &mut self,
1262 request: GetWorkflowProviderRunEventsRequest,
1263 ) -> Result<GetWorkflowProviderRunEventsResponse, GestaltError> {
1264 let mut request = request;
1265 if request.context.is_none() {
1266 request.context = self.context.clone();
1267 }
1268 let mut tonic_request =
1269 tonic::Request::new(to_wire_get_workflow_provider_run_events_request(request));
1270 if let Some(timeout) = self.timeout {
1271 tonic_request.set_timeout(timeout);
1272 }
1273 let response = self.inner.get_run_events(tonic_request).await?;
1274 Ok(from_wire_get_workflow_provider_run_events_response(
1275 response.into_inner(),
1276 ))
1277 }
1278
1279 pub async fn get_run_output(
1281 &mut self,
1282 provider: String,
1283 run_id: String,
1284 ) -> Result<Option<serde_json::Value>, GestaltError> {
1285 let request = GetWorkflowProviderRunOutputRequest {
1286 provider,
1287 run_id,
1288 context: self.context.clone(),
1289 };
1290 let mut tonic_request =
1291 tonic::Request::new(to_wire_get_workflow_provider_run_output_request(request));
1292 if let Some(timeout) = self.timeout {
1293 tonic_request.set_timeout(timeout);
1294 }
1295 let response = from_wire_get_workflow_provider_run_output_response(
1296 self.inner.get_run_output(tonic_request).await?.into_inner(),
1297 );
1298 Ok(response.output)
1299 }
1300
1301 pub async fn get_run_output_raw(
1303 &mut self,
1304 request: GetWorkflowProviderRunOutputRequest,
1305 ) -> Result<GetWorkflowProviderRunOutputResponse, GestaltError> {
1306 let mut request = request;
1307 if request.context.is_none() {
1308 request.context = self.context.clone();
1309 }
1310 let mut tonic_request =
1311 tonic::Request::new(to_wire_get_workflow_provider_run_output_request(request));
1312 if let Some(timeout) = self.timeout {
1313 tonic_request.set_timeout(timeout);
1314 }
1315 let response = self.inner.get_run_output(tonic_request).await?;
1316 Ok(from_wire_get_workflow_provider_run_output_response(
1317 response.into_inner(),
1318 ))
1319 }
1320
1321 pub async fn cancel_run(
1323 &mut self,
1324 provider: String,
1325 run_id: String,
1326 reason: String,
1327 ) -> Result<WorkflowRun, GestaltError> {
1328 let request = CancelWorkflowProviderRunRequest {
1329 provider,
1330 run_id,
1331 reason,
1332 context: self.context.clone(),
1333 };
1334 let mut tonic_request =
1335 tonic::Request::new(to_wire_cancel_workflow_provider_run_request(request));
1336 if let Some(timeout) = self.timeout {
1337 tonic_request.set_timeout(timeout);
1338 }
1339 let response = self.inner.cancel_run(tonic_request).await?;
1340 Ok(from_wire_workflow_run(response.into_inner()))
1341 }
1342
1343 pub async fn cancel_run_raw(
1345 &mut self,
1346 request: CancelWorkflowProviderRunRequest,
1347 ) -> Result<WorkflowRun, GestaltError> {
1348 let mut request = request;
1349 if request.context.is_none() {
1350 request.context = self.context.clone();
1351 }
1352 let mut tonic_request =
1353 tonic::Request::new(to_wire_cancel_workflow_provider_run_request(request));
1354 if let Some(timeout) = self.timeout {
1355 tonic_request.set_timeout(timeout);
1356 }
1357 let response = self.inner.cancel_run(tonic_request).await?;
1358 Ok(from_wire_workflow_run(response.into_inner()))
1359 }
1360
1361 pub async fn signal_run(
1363 &mut self,
1364 provider: String,
1365 run_id: String,
1366 signal: Option<WorkflowSignal>,
1367 ) -> Result<SignalWorkflowRunResponse, GestaltError> {
1368 let request = SignalWorkflowProviderRunRequest {
1369 provider,
1370 run_id,
1371 signal,
1372 context: self.context.clone(),
1373 };
1374 let mut tonic_request =
1375 tonic::Request::new(to_wire_signal_workflow_provider_run_request(request));
1376 if let Some(timeout) = self.timeout {
1377 tonic_request.set_timeout(timeout);
1378 }
1379 let response = self.inner.signal_run(tonic_request).await?;
1380 Ok(from_wire_signal_workflow_run_response(
1381 response.into_inner(),
1382 ))
1383 }
1384
1385 pub async fn signal_run_raw(
1387 &mut self,
1388 request: SignalWorkflowProviderRunRequest,
1389 ) -> Result<SignalWorkflowRunResponse, GestaltError> {
1390 let mut request = request;
1391 if request.context.is_none() {
1392 request.context = self.context.clone();
1393 }
1394 let mut tonic_request =
1395 tonic::Request::new(to_wire_signal_workflow_provider_run_request(request));
1396 if let Some(timeout) = self.timeout {
1397 tonic_request.set_timeout(timeout);
1398 }
1399 let response = self.inner.signal_run(tonic_request).await?;
1400 Ok(from_wire_signal_workflow_run_response(
1401 response.into_inner(),
1402 ))
1403 }
1404
1405 #[allow(clippy::too_many_arguments)]
1407 pub async fn signal_or_start_run(
1408 &mut self,
1409 workflow_key: String,
1410 idempotency_key: String,
1411 provider: String,
1412 definition_id: String,
1413 expected_definition_generation: i64,
1414 signal: Option<WorkflowSignal>,
1415 input: Option<serde_json::Map<String, serde_json::Value>>,
1416 ) -> Result<SignalWorkflowRunResponse, GestaltError> {
1417 let request = SignalOrStartWorkflowProviderRunRequest {
1418 workflow_key,
1419 idempotency_key,
1420 provider,
1421 definition_id,
1422 expected_definition_generation,
1423 signal,
1424 input,
1425 context: self.context.clone(),
1426 };
1427 let mut tonic_request = tonic::Request::new(
1428 to_wire_signal_or_start_workflow_provider_run_request(request),
1429 );
1430 if let Some(timeout) = self.timeout {
1431 tonic_request.set_timeout(timeout);
1432 }
1433 let response = self.inner.signal_or_start_run(tonic_request).await?;
1434 Ok(from_wire_signal_workflow_run_response(
1435 response.into_inner(),
1436 ))
1437 }
1438
1439 pub async fn signal_or_start_run_raw(
1441 &mut self,
1442 request: SignalOrStartWorkflowProviderRunRequest,
1443 ) -> Result<SignalWorkflowRunResponse, GestaltError> {
1444 let mut request = request;
1445 if request.context.is_none() {
1446 request.context = self.context.clone();
1447 }
1448 let mut tonic_request = tonic::Request::new(
1449 to_wire_signal_or_start_workflow_provider_run_request(request),
1450 );
1451 if let Some(timeout) = self.timeout {
1452 tonic_request.set_timeout(timeout);
1453 }
1454 let response = self.inner.signal_or_start_run(tonic_request).await?;
1455 Ok(from_wire_signal_workflow_run_response(
1456 response.into_inner(),
1457 ))
1458 }
1459
1460 pub async fn deliver_event(
1462 &mut self,
1463 provider: String,
1464 event: Option<WorkflowEvent>,
1465 ) -> Result<WorkflowEvent, GestaltError> {
1466 let request = DeliverWorkflowProviderEventRequest {
1467 provider,
1468 event,
1469 context: self.context.clone(),
1470 };
1471 let mut tonic_request =
1472 tonic::Request::new(to_wire_deliver_workflow_provider_event_request(request));
1473 if let Some(timeout) = self.timeout {
1474 tonic_request.set_timeout(timeout);
1475 }
1476 let response = self.inner.deliver_event(tonic_request).await?;
1477 Ok(from_wire_workflow_event(response.into_inner()))
1478 }
1479
1480 pub async fn deliver_event_raw(
1482 &mut self,
1483 request: DeliverWorkflowProviderEventRequest,
1484 ) -> Result<WorkflowEvent, GestaltError> {
1485 let mut request = request;
1486 if request.context.is_none() {
1487 request.context = self.context.clone();
1488 }
1489 let mut tonic_request =
1490 tonic::Request::new(to_wire_deliver_workflow_provider_event_request(request));
1491 if let Some(timeout) = self.timeout {
1492 tonic_request.set_timeout(timeout);
1493 }
1494 let response = self.inner.deliver_event(tonic_request).await?;
1495 Ok(from_wire_workflow_event(response.into_inner()))
1496 }
1497}