1mod auth;
5pub(crate) mod convert;
9mod mint_resolve;
11mod plain_rpcs;
13mod routing_resolve;
15mod status;
17
18pub(crate) use auth::caller_from_metadata;
19pub(crate) use status::{status_from_wire_error, status_with_code};
20
21use aion_proto::generated::{self, workflow_service_server::WorkflowServiceServer};
22use tonic::{Request, Response, Status};
23
24use crate::namespace::MintCredentials;
25use crate::routing::{ForwardReply, ForwardRequest};
26use crate::{CallerIdentity, ServerState, api::handlers};
27use convert::decode_workflow_id;
28use convert::{
29 decode_cancel_request, decode_pause_request, decode_query_request, decode_rename_request,
30 decode_reopen_request, decode_resume_request, decode_signal_request, decode_start_request,
31 encode_cancel_response, encode_pause_response, encode_query_response, encode_rename_response,
32 encode_reopen_response, encode_resume_response, encode_signal_response, encode_start_response,
33};
34use routing_resolve::{RouteResolution, StartResolution};
35
36#[derive(Clone)]
38pub struct WorkflowGrpcService {
39 state: ServerState,
40}
41
42impl WorkflowGrpcService {
43 #[must_use]
45 pub const fn new(state: ServerState) -> Self {
46 Self { state }
47 }
48
49 async fn caller<T>(&self, request: &Request<T>) -> Result<CallerIdentity, Status> {
50 caller_from_metadata(request.metadata(), &self.state).await
51 }
52}
53
54#[must_use]
56pub fn workflow_service(state: ServerState) -> WorkflowServiceServer<WorkflowGrpcService> {
57 WorkflowServiceServer::new(WorkflowGrpcService::new(state))
58}
59
60#[tonic::async_trait]
61impl generated::workflow_service_server::WorkflowService for WorkflowGrpcService {
62 async fn start_workflow(
63 &self,
64 request: Request<generated::StartWorkflowRequest>,
65 ) -> Result<Response<generated::StartWorkflowResponse>, Status> {
66 if self.state.drain_state().is_draining() {
67 return Err(Status::unavailable(
68 "server is draining and not accepting new workflow starts",
69 ));
70 }
71 let caller = self.caller(&request).await?;
72 let (metadata, _ext, inner) = request.into_parts();
78 let placement = match self.resolve_start(&inner, &metadata).await {
79 StartResolution::Reject(status) => return Err(status),
80 StartResolution::Reply(reply) => return Ok(Response::new(reply)),
81 StartResolution::Local(placement) => placement,
82 };
83 let minter = self
92 .state
93 .namespace_minter()
94 .with_caller_credentials(MintCredentials::from_grpc_metadata(&metadata));
95 let response = handlers::start_with_placement(
96 self.state.namespace_guard(),
97 &caller,
98 decode_start_request(inner),
99 placement,
100 Some(&minter),
101 )
102 .await
103 .map_err(status_from_wire_error)?;
104 return Ok(Response::new(encode_start_response(response)));
105 }
106
107 async fn signal(
108 &self,
109 request: Request<generated::SignalRequest>,
110 ) -> Result<Response<generated::SignalResponse>, Status> {
111 let caller = self.caller(&request).await?;
112 let (metadata, _ext, inner) = request.into_parts();
113 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
114 match self
115 .resolve_route(
116 workflow_id,
117 &metadata,
118 ForwardRequest::Signal(inner.clone()),
119 )
120 .await
121 {
122 RouteResolution::Reject(status) => return Err(status),
123 RouteResolution::Reply(ForwardReply::Signal(reply)) => {
124 return Ok(Response::new(reply));
125 }
126 RouteResolution::Reply(_) => {
127 return Err(Status::internal("forwarder returned a mismatched reply"));
128 }
129 RouteResolution::Local => {
130 let response = handlers::signal(
131 self.state.namespace_guard(),
132 &caller,
133 decode_signal_request(inner),
134 )
135 .await
136 .map_err(status_from_wire_error)?;
137 return Ok(Response::new(encode_signal_response(response)));
138 }
139 }
140 }
141
142 async fn query(
143 &self,
144 request: Request<generated::QueryRequest>,
145 ) -> Result<Response<generated::QueryResponse>, Status> {
146 let caller = self.caller(&request).await?;
147 let (metadata, _ext, inner) = request.into_parts();
148 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
149 match self
150 .resolve_route(workflow_id, &metadata, ForwardRequest::Query(inner.clone()))
151 .await
152 {
153 RouteResolution::Reject(status) => return Err(status),
154 RouteResolution::Reply(ForwardReply::Query(reply)) => {
155 return Ok(Response::new(reply));
156 }
157 RouteResolution::Reply(_) => {
158 return Err(Status::internal("forwarder returned a mismatched reply"));
159 }
160 RouteResolution::Local => {
161 let response = handlers::query(
162 self.state.namespace_guard(),
163 &caller,
164 decode_query_request(inner),
165 )
166 .await
167 .map_err(status_from_wire_error)?;
168 return Ok(Response::new(encode_query_response(response)));
169 }
170 }
171 }
172
173 async fn cancel(
174 &self,
175 request: Request<generated::CancelRequest>,
176 ) -> Result<Response<generated::CancelResponse>, Status> {
177 let caller = self.caller(&request).await?;
178 let (metadata, _ext, inner) = request.into_parts();
179 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
180 match self
181 .resolve_route(
182 workflow_id,
183 &metadata,
184 ForwardRequest::Cancel(inner.clone()),
185 )
186 .await
187 {
188 RouteResolution::Reject(status) => return Err(status),
189 RouteResolution::Reply(ForwardReply::Cancel(reply)) => {
190 return Ok(Response::new(reply));
191 }
192 RouteResolution::Reply(_) => {
193 return Err(Status::internal("forwarder returned a mismatched reply"));
194 }
195 RouteResolution::Local => {
196 let response = handlers::cancel(
197 &self.state,
198 self.state.namespace_guard(),
199 &caller,
200 decode_cancel_request(inner),
201 )
202 .await
203 .map_err(status_from_wire_error)?;
204 return Ok(Response::new(encode_cancel_response(response)));
205 }
206 }
207 }
208
209 async fn reopen(
210 &self,
211 request: Request<generated::ReopenRequest>,
212 ) -> Result<Response<generated::ReopenResponse>, Status> {
213 let caller = self.caller(&request).await?;
214 let (metadata, _ext, inner) = request.into_parts();
215 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
216 match self
217 .resolve_route(
218 workflow_id,
219 &metadata,
220 ForwardRequest::Reopen(inner.clone()),
221 )
222 .await
223 {
224 RouteResolution::Reject(status) => return Err(status),
225 RouteResolution::Reply(ForwardReply::Reopen(reply)) => {
226 return Ok(Response::new(reply));
227 }
228 RouteResolution::Reply(_) => {
229 return Err(Status::internal("forwarder returned a mismatched reply"));
230 }
231 RouteResolution::Local => {
232 let response = handlers::reopen(
233 self.state.namespace_guard(),
234 &caller,
235 decode_reopen_request(inner),
236 )
237 .await
238 .map_err(status_from_wire_error)?;
239 return Ok(Response::new(encode_reopen_response(response)));
240 }
241 }
242 }
243
244 async fn pause(
245 &self,
246 request: Request<generated::PauseRequest>,
247 ) -> Result<Response<generated::PauseResponse>, Status> {
248 let caller = self.caller(&request).await?;
249 let (metadata, _ext, inner) = request.into_parts();
250 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
251 match self
252 .resolve_route(workflow_id, &metadata, ForwardRequest::Pause(inner.clone()))
253 .await
254 {
255 RouteResolution::Reject(status) => return Err(status),
256 RouteResolution::Reply(ForwardReply::Pause(reply)) => {
257 return Ok(Response::new(reply));
258 }
259 RouteResolution::Reply(_) => {
260 return Err(Status::internal("forwarder returned a mismatched reply"));
261 }
262 RouteResolution::Local => {
263 let response = handlers::pause(
264 self.state.namespace_guard(),
265 &caller,
266 decode_pause_request(inner),
267 )
268 .await
269 .map_err(status_from_wire_error)?;
270 return Ok(Response::new(encode_pause_response(response)));
271 }
272 }
273 }
274
275 async fn resume(
276 &self,
277 request: Request<generated::ResumeRequest>,
278 ) -> Result<Response<generated::ResumeResponse>, Status> {
279 let caller = self.caller(&request).await?;
280 let (metadata, _ext, inner) = request.into_parts();
281 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
282 match self
283 .resolve_route(
284 workflow_id,
285 &metadata,
286 ForwardRequest::Resume(inner.clone()),
287 )
288 .await
289 {
290 RouteResolution::Reject(status) => return Err(status),
291 RouteResolution::Reply(ForwardReply::Resume(reply)) => {
292 return Ok(Response::new(reply));
293 }
294 RouteResolution::Reply(_) => {
295 return Err(Status::internal("forwarder returned a mismatched reply"));
296 }
297 RouteResolution::Local => {
298 let response = handlers::resume(
299 self.state.namespace_guard(),
300 &caller,
301 decode_resume_request(inner),
302 )
303 .await
304 .map_err(status_from_wire_error)?;
305 return Ok(Response::new(encode_resume_response(response)));
306 }
307 }
308 }
309
310 async fn rename(
311 &self,
312 request: Request<generated::RenameRequest>,
313 ) -> Result<Response<generated::RenameResponse>, Status> {
314 let caller = self.caller(&request).await?;
315 let (metadata, _ext, inner) = request.into_parts();
316 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
317 match self
318 .resolve_route(
319 workflow_id,
320 &metadata,
321 ForwardRequest::Rename(inner.clone()),
322 )
323 .await
324 {
325 RouteResolution::Reject(status) => return Err(status),
326 RouteResolution::Reply(ForwardReply::Rename(reply)) => {
327 return Ok(Response::new(reply));
328 }
329 RouteResolution::Reply(_) => {
330 return Err(Status::internal("forwarder returned a mismatched reply"));
331 }
332 RouteResolution::Local => {
333 let response = handlers::rename(
334 self.state.namespace_guard(),
335 &caller,
336 decode_rename_request(inner),
337 )
338 .await
339 .map_err(status_from_wire_error)?;
340 return Ok(Response::new(encode_rename_response(response)));
341 }
342 }
343 }
344
345 async fn list_workflows(
346 &self,
347 request: Request<generated::ListWorkflowsRequest>,
348 ) -> Result<Response<generated::ListWorkflowsResponse>, Status> {
349 plain_rpcs::list_workflows(&self.state, request).await
350 }
351
352 async fn count_workflows(
353 &self,
354 request: Request<generated::CountWorkflowsRequest>,
355 ) -> Result<Response<generated::CountWorkflowsResponse>, Status> {
356 plain_rpcs::count_workflows(&self.state, request).await
357 }
358
359 async fn describe_workflow(
360 &self,
361 request: Request<generated::DescribeWorkflowRequest>,
362 ) -> Result<Response<generated::DescribeWorkflowResponse>, Status> {
363 plain_rpcs::describe_workflow(&self.state, request).await
364 }
365
366 async fn create_schedule(
367 &self,
368 request: Request<generated::CreateScheduleRequest>,
369 ) -> Result<Response<generated::CreateScheduleResponse>, Status> {
370 plain_rpcs::create_schedule(&self.state, request).await
371 }
372
373 async fn update_schedule(
374 &self,
375 request: Request<generated::UpdateScheduleRequest>,
376 ) -> Result<Response<generated::UpdateScheduleResponse>, Status> {
377 plain_rpcs::update_schedule(&self.state, request).await
378 }
379
380 async fn pause_schedule(
381 &self,
382 request: Request<generated::ScheduleIdRequest>,
383 ) -> Result<Response<generated::PauseScheduleResponse>, Status> {
384 plain_rpcs::pause_schedule(&self.state, request).await
385 }
386
387 async fn resume_schedule(
388 &self,
389 request: Request<generated::ScheduleIdRequest>,
390 ) -> Result<Response<generated::ResumeScheduleResponse>, Status> {
391 plain_rpcs::resume_schedule(&self.state, request).await
392 }
393
394 async fn delete_schedule(
395 &self,
396 request: Request<generated::ScheduleIdRequest>,
397 ) -> Result<Response<generated::DeleteScheduleResponse>, Status> {
398 plain_rpcs::delete_schedule(&self.state, request).await
399 }
400
401 async fn list_schedules(
402 &self,
403 request: Request<generated::ListSchedulesRequest>,
404 ) -> Result<Response<generated::ListSchedulesResponse>, Status> {
405 plain_rpcs::list_schedules(&self.state, request).await
406 }
407
408 async fn describe_schedule(
409 &self,
410 request: Request<generated::ScheduleIdRequest>,
411 ) -> Result<Response<generated::DescribeScheduleResponse>, Status> {
412 plain_rpcs::describe_schedule(&self.state, request).await
413 }
414
415 async fn mint_namespace(
420 &self,
421 request: Request<generated::MintNamespaceRequest>,
422 ) -> Result<Response<generated::MintNamespaceResponse>, Status> {
423 self.mint_namespace_here(request).await
424 }
425}
426
427#[cfg(test)]
428mod tests {
429 use std::{net::SocketAddr, sync::Arc};
430
431 use aion::EngineBuilder;
432 use aion_core::{Event, EventEnvelope, Payload, WorkflowId, WorkflowStatus};
433 use aion_proto::{
434 ProtoWireError, WireError, WireErrorCode,
435 convert::{decode_core_value, encode_core_value},
436 generated::workflow_service_server::WorkflowService,
437 };
438 use aion_store::{
439 EventStore, InMemoryStore, WriteToken,
440 visibility::{VisibilityRecord, VisibilityStore},
441 };
442 use chrono::Utc;
443 use prost::Message;
444 use serde_json::json;
445 use tonic::{Code, Request};
446
447 use super::convert::{decode_envelope, encode_envelope, encode_payload};
448 use super::*;
449 use crate::{
450 NamespaceResolver,
451 config::{
452 AuthConfig, AuthoringConfig, DeployConfig, ListenConfig, MetricsConfig,
453 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, RuntimeConfig,
454 WebSocketConfig, WorkerConfig,
455 },
456 };
457
458 const NAMESPACE: &str = "tenant-a";
459 const TOKEN: &str = "test-token";
460
461 async fn server_state(
466 resolver: NamespaceResolver,
467 runtime: RuntimeConfig,
468 ) -> Result<ServerState, Box<dyn std::error::Error>> {
469 #[cfg(feature = "auth")]
470 {
471 let url = crate::auth::test_support::serve_jwks()?;
472 let refresh = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
473 let cache = crate::auth::JwksCache::new(url, refresh).await?;
474 Ok(ServerState::from_parts_with_jwks(resolver, runtime, cache))
475 }
476 #[cfg(not(feature = "auth"))]
477 {
478 tokio::task::yield_now().await;
480 Ok(ServerState::from_parts(resolver, runtime))
481 }
482 }
483
484 #[tokio::test]
485 async fn in_process_tonic_start_and_list_use_shared_handlers()
486 -> Result<(), Box<dyn std::error::Error>> {
487 let backing = Arc::new(InMemoryStore::default());
488 let store: Arc<dyn EventStore> = backing.clone();
489 let visibility_store: Arc<dyn VisibilityStore> = backing;
490 let engine = Arc::new(
491 EngineBuilder::new()
492 .store_arc(Arc::clone(&store))
493 .visibility_store_arc(Arc::clone(&visibility_store))
494 .scheduler_threads(1)
495 .build()
496 .await?,
497 );
498 store
499 .append(
500 WriteToken::recorder(),
501 &workflow_id(),
502 &[started_event()?],
503 0,
504 )
505 .await?;
506 visibility_store
507 .record_visibility(VisibilityRecord {
508 workflow_id: workflow_id(),
509 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(2)),
510 workflow_type: String::from("fixture"),
511 status: WorkflowStatus::Running,
512 start_time: Utc::now(),
513 close_time: None,
514 failed_step: None,
515 failure_reason: None,
516 search_attributes: std::collections::HashMap::from([(
517 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
518 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
519 )]),
520 })
521 .await?;
522 let resolver = NamespaceResolver::from_config(
523 crate::config::NamespaceConfig {
524 mode: NamespaceMode::SharedEngine,
525 },
526 engine,
527 );
528 let state = server_state(resolver.clone(), runtime_config()).await?;
529 let service = WorkflowGrpcService::new(state);
530
531 let mut start = Request::new(generated::StartWorkflowRequest {
532 namespace: NAMESPACE.to_owned(),
533 workflow_type: "missing-workflow".to_owned(),
534 input: Some(encode_payload(proto_payload()?)),
535 routing_key: None,
536 task_queue: None,
537 display_name: None,
538 });
539 apply_metadata(start.metadata_mut())?;
540 let start_error = service.start_workflow(start).await;
541 let status = start_error
542 .err()
543 .ok_or_else(|| WireError::backend("expected error"))?;
544 assert_eq!(status.code(), Code::NotFound);
545 let detail = ProtoWireError::decode(status.details())?;
546 assert_eq!(detail.error_type.as_deref(), Some("WorkflowTypeNotFound"));
547 assert!(detail.message.contains("missing-workflow"));
548
549 let list_filter = encode_core_value(
550 NAMESPACE,
551 None,
552 &aion_store::visibility::ListWorkflowsFilter {
553 workflow_type: Some(String::from("fixture")),
554 status: Some(WorkflowStatus::Running),
555 ..aion_store::visibility::ListWorkflowsFilter::default()
556 },
557 )?;
558 let mut list = Request::new(generated::ListWorkflowsRequest {
559 namespace: NAMESPACE.to_owned(),
560 filter: Some(encode_envelope(list_filter)),
561 });
562 apply_metadata(list.metadata_mut())?;
563 let response = service.list_workflows(list).await?.into_inner();
564
565 assert_eq!(response.summaries.len(), 1);
566 let summary = response
567 .summaries
568 .into_iter()
569 .next()
570 .map(decode_envelope)
571 .map(|envelope| decode_core_value::<aion_store::visibility::WorkflowSummary>(&envelope))
572 .transpose()?
573 .ok_or_else(|| WireError::backend("summary missing"))?;
574 assert_eq!(summary.workflow_id, workflow_id());
575 assert_eq!(
581 resolver
582 .verify_workflow_ownership(NAMESPACE, &workflow_id())
583 .await
584 .err()
585 .map(|error| error.to_wire_error().code),
586 Some(WireErrorCode::NotFound)
587 );
588 Ok(())
589 }
590
591 async fn rename_fixture(
602 register_display_name: bool,
603 ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
604 rename_fixture_with_terminal(register_display_name, true).await
605 }
606
607 async fn rename_fixture_with_terminal(
611 register_display_name: bool,
612 terminal: bool,
613 ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
614 let backing = Arc::new(InMemoryStore::default());
615 let store: Arc<dyn EventStore> = backing.clone();
616 let visibility_store: Arc<dyn VisibilityStore> = backing;
617 let mut schema = aion_core::SearchAttributeSchema::new();
618 schema.register(
619 crate::namespace::NAMESPACE_ATTRIBUTE,
620 aion_core::SearchAttributeType::String,
621 )?;
622 if register_display_name {
623 schema.register(
624 crate::namespace::DISPLAY_NAME_ATTRIBUTE,
625 aion_core::SearchAttributeType::String,
626 )?;
627 }
628 let engine = Arc::new(
629 EngineBuilder::new()
630 .store_arc(Arc::clone(&store))
631 .visibility_store_arc(Arc::clone(&visibility_store))
632 .search_attribute_schema(schema)
633 .scheduler_threads(1)
634 .build()
635 .await?,
636 );
637 let namespace_event = |seq: u64| Event::SearchAttributesUpdated {
638 envelope: EventEnvelope {
639 seq,
640 recorded_at: Utc::now(),
641 workflow_id: workflow_id(),
642 },
643 workflow_id: workflow_id(),
644 attributes: std::collections::HashMap::from([(
645 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
646 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
647 )]),
648 };
649 let mut events = vec![started_event()?, namespace_event(2)];
650 if terminal {
651 events.push(Event::WorkflowCompleted {
652 envelope: EventEnvelope {
653 seq: 3,
654 recorded_at: Utc::now(),
655 workflow_id: workflow_id(),
656 },
657 result: payload()?,
658 });
659 }
660 store
661 .append(WriteToken::recorder(), &workflow_id(), &events, 0)
662 .await?;
663 let resolver = NamespaceResolver::from_config(
664 crate::config::NamespaceConfig {
665 mode: NamespaceMode::SharedEngine,
666 },
667 engine,
668 );
669 let state = server_state(resolver, runtime_config()).await?;
670 Ok((WorkflowGrpcService::new(state), store))
671 }
672
673 #[tokio::test]
677 async fn in_process_tonic_rename_records_the_name_and_supersedes_it()
678 -> Result<(), Box<dyn std::error::Error>> {
679 let (service, store) = rename_fixture(true).await?;
680
681 let rename = |name: &str| {
682 let mut request = Request::new(generated::RenameRequest {
683 namespace: NAMESPACE.to_owned(),
684 workflow_id: Some(generated::WorkflowId {
685 uuid: workflow_id().to_string(),
686 }),
687 run_id: None,
689 display_name: name.to_owned(),
690 });
691 apply_metadata(request.metadata_mut()).map(|()| request)
692 };
693
694 let response = service.rename(rename(" Nightly settlement ")?).await?;
696 let response = response.into_inner();
697 assert_eq!(response.display_name, "Nightly settlement");
698 assert_eq!(
699 response.run_id.map(|id| id.uuid),
700 Some(aion_core::RunId::new(uuid::Uuid::from_u128(1)).to_string())
701 );
702
703 service
704 .rename(rename("Nightly settlement (rerun)")?)
705 .await?;
706
707 let history = store.read_history(&workflow_id()).await?;
708 assert_eq!(
709 aion_core::display_name(&history).as_deref(),
710 Some("Nightly settlement (rerun)"),
711 "the latest recorded name wins"
712 );
713 let names: Vec<_> = history
714 .iter()
715 .filter_map(|event| match event {
716 Event::SearchAttributesUpdated { attributes, .. } => attributes
717 .get(aion_core::DISPLAY_NAME_ATTRIBUTE)
718 .and_then(|value| match value {
719 aion_core::SearchAttributeValue::String(name) => Some(name.clone()),
720 _ => None,
721 }),
722 _ => None,
723 })
724 .collect();
725 assert_eq!(
726 names,
727 vec![
728 String::from("Nightly settlement"),
729 String::from("Nightly settlement (rerun)")
730 ],
731 "history keeps every name the run has worn"
732 );
733 Ok(())
734 }
735
736 #[tokio::test]
739 async fn in_process_tonic_rename_refuses_a_blank_name() -> Result<(), Box<dyn std::error::Error>>
740 {
741 let (service, store) = rename_fixture(false).await?;
742 let before = store.read_history(&workflow_id()).await?.len();
743
744 let mut request = Request::new(generated::RenameRequest {
745 namespace: NAMESPACE.to_owned(),
746 workflow_id: Some(generated::WorkflowId {
747 uuid: workflow_id().to_string(),
748 }),
749 run_id: None,
750 display_name: String::from(" "),
751 });
752 apply_metadata(request.metadata_mut())?;
753 let status = service
754 .rename(request)
755 .await
756 .err()
757 .ok_or_else(|| WireError::backend("expected a blank-name refusal"))?;
758 assert_eq!(status.code(), Code::InvalidArgument);
759 assert_eq!(
760 store.read_history(&workflow_id()).await?.len(),
761 before,
762 "a refused rename appends nothing"
763 );
764 Ok(())
765 }
766
767 #[tokio::test]
773 async fn in_process_tonic_rename_of_a_non_resident_running_run_is_failed_precondition()
774 -> Result<(), Box<dyn std::error::Error>> {
775 let (service, store) = rename_fixture_with_terminal(true, false).await?;
776 let before = store.read_history(&workflow_id()).await?.len();
777
778 let mut request = Request::new(generated::RenameRequest {
779 namespace: NAMESPACE.to_owned(),
780 workflow_id: Some(generated::WorkflowId {
781 uuid: workflow_id().to_string(),
782 }),
783 run_id: None,
784 display_name: String::from("Nightly settlement"),
785 });
786 apply_metadata(request.metadata_mut())?;
787 let status = service
788 .rename(request)
789 .await
790 .err()
791 .ok_or_else(|| WireError::backend("expected a residency refusal"))?;
792
793 assert_eq!(status.code(), Code::FailedPrecondition);
794 let detail = ProtoWireError::decode(status.details())?;
795 assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
796 assert!(
797 detail.message.contains("not resident"),
798 "the refusal must name the reason: {}",
799 detail.message
800 );
801 assert_eq!(
802 store.read_history(&workflow_id()).await?.len(),
803 before,
804 "a refused rename appends nothing"
805 );
806 Ok(())
807 }
808
809 #[tokio::test]
813 async fn in_process_tonic_reopen_completed_is_failed_precondition_invalid_state()
814 -> Result<(), Box<dyn std::error::Error>> {
815 let backing = Arc::new(InMemoryStore::default());
816 let store: Arc<dyn EventStore> = backing.clone();
817 let visibility_store: Arc<dyn VisibilityStore> = backing;
818 let engine = Arc::new(
819 EngineBuilder::new()
820 .store_arc(Arc::clone(&store))
821 .visibility_store_arc(Arc::clone(&visibility_store))
822 .scheduler_threads(1)
823 .build()
824 .await?,
825 );
826 store
830 .append(
831 WriteToken::recorder(),
832 &workflow_id(),
833 &[
834 started_event()?,
835 Event::SearchAttributesUpdated {
836 envelope: EventEnvelope {
837 seq: 2,
838 recorded_at: Utc::now(),
839 workflow_id: workflow_id(),
840 },
841 workflow_id: workflow_id(),
842 attributes: std::collections::HashMap::from([(
843 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
844 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
845 )]),
846 },
847 Event::WorkflowCompleted {
848 envelope: EventEnvelope {
849 seq: 3,
850 recorded_at: Utc::now(),
851 workflow_id: workflow_id(),
852 },
853 result: payload()?,
854 },
855 ],
856 0,
857 )
858 .await?;
859 let resolver = NamespaceResolver::from_config(
860 crate::config::NamespaceConfig {
861 mode: NamespaceMode::SharedEngine,
862 },
863 engine,
864 );
865 let state = server_state(resolver, runtime_config()).await?;
866 let service = WorkflowGrpcService::new(state);
867
868 let mut reopen = Request::new(generated::ReopenRequest {
869 namespace: NAMESPACE.to_owned(),
870 workflow_id: Some(generated::WorkflowId {
871 uuid: workflow_id().to_string(),
872 }),
873 run_id: None,
874 });
875 apply_metadata(reopen.metadata_mut())?;
876 let status = service
877 .reopen(reopen)
878 .await
879 .err()
880 .ok_or_else(|| WireError::backend("expected a reopen precondition error"))?;
881 assert_eq!(status.code(), Code::FailedPrecondition);
882 let detail = ProtoWireError::decode(status.details())?;
883 assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
884 assert_eq!(
885 detail.code,
886 aion_proto::ProtoWireErrorCode::InvalidState as i32
887 );
888 Ok(())
889 }
890
891 fn apply_metadata(
892 metadata: &mut tonic::metadata::MetadataMap,
893 ) -> Result<(), Box<dyn std::error::Error>> {
894 #[cfg(feature = "auth")]
898 let bearer = crate::auth::test_support::mint_token("alice", NAMESPACE)?;
899 #[cfg(not(feature = "auth"))]
900 let bearer = TOKEN.to_owned();
901 metadata.insert("authorization", format!("Bearer {bearer}").parse()?);
902 metadata.insert("x-aion-subject", "alice".parse()?);
903 metadata.insert("x-aion-namespaces", NAMESPACE.parse()?);
904 Ok(())
905 }
906
907 fn runtime_config() -> RuntimeConfig {
911 RuntimeConfig {
912 listen: ListenConfig {
913 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
914 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
915 },
916 tls: None,
917 auth: AuthConfig {
918 enabled: true,
919 jwks_url: Some(TOKEN.to_owned()),
920 jwks_refresh_seconds: 300,
921 },
922 ops_console: OpsConsoleConfig {
923 source: OpsConsoleAssetSource::Embedded,
924 },
925 namespace: NamespaceConfig {
926 mode: NamespaceMode::SharedEngine,
927 },
928 worker: WorkerConfig {
929 heartbeat_window: std::time::Duration::from_secs(30),
930 ..WorkerConfig::default()
931 },
932 websocket: WebSocketConfig {
933 outbound_buffer_bound: 32,
934 event_broadcast_capacity: Some(64),
935 cluster_broadcast_capacity: Some(64),
936 },
937 workflow_packages: Vec::new(),
938 deploy: DeployConfig::default(),
939 authoring: AuthoringConfig::default(),
940 dev: crate::config::DevConfig::default(),
941 outbox: crate::config::OutboxConfig::default(),
942 observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
943 mcp: crate::config::ResolvedMcpConfig::default(),
944 scheduler_threads: 1,
945 jit_threshold: None,
946 query_timeout: Some(std::time::Duration::from_secs(10)),
947 default_namespace: "default".to_owned(),
948 auto_create: crate::config::AutoCreate::Open,
949 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
950 drain_timeout: std::time::Duration::from_secs(30),
951 metrics: MetricsConfig { enabled: true },
952 owned_shards: Vec::new(),
953 cors_allowed_origins: Vec::new(),
954 }
955 }
956
957 fn started_event() -> Result<Event, aion_core::PayloadError> {
958 Ok(Event::WorkflowStarted {
959 envelope: EventEnvelope {
960 seq: 1,
961 recorded_at: Utc::now(),
962 workflow_id: workflow_id(),
963 },
964 workflow_type: "fixture".to_owned(),
965 input: payload()?,
966 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
967 parent_run_id: None,
968 parent_workflow_id: None,
969 package_version: aion_core::PackageVersion::new("a".repeat(64)),
970 })
971 }
972
973 fn proto_payload() -> Result<aion_proto::ProtoPayload, aion_core::PayloadError> {
974 Ok(payload()?.into())
975 }
976
977 fn payload() -> Result<Payload, aion_core::PayloadError> {
978 Payload::from_json(&json!({ "fixture": "input" }))
979 }
980
981 fn workflow_id() -> WorkflowId {
982 WorkflowId::new(uuid::Uuid::from_u128(1))
983 }
984}