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 read_history(
367 &self,
368 request: Request<generated::ReadHistoryRequest>,
369 ) -> Result<Response<generated::ReadHistoryResponse>, Status> {
370 plain_rpcs::read_history(&self.state, request).await
371 }
372
373 async fn create_schedule(
374 &self,
375 request: Request<generated::CreateScheduleRequest>,
376 ) -> Result<Response<generated::CreateScheduleResponse>, Status> {
377 plain_rpcs::create_schedule(&self.state, request).await
378 }
379
380 async fn update_schedule(
381 &self,
382 request: Request<generated::UpdateScheduleRequest>,
383 ) -> Result<Response<generated::UpdateScheduleResponse>, Status> {
384 plain_rpcs::update_schedule(&self.state, request).await
385 }
386
387 async fn pause_schedule(
388 &self,
389 request: Request<generated::ScheduleIdRequest>,
390 ) -> Result<Response<generated::PauseScheduleResponse>, Status> {
391 plain_rpcs::pause_schedule(&self.state, request).await
392 }
393
394 async fn resume_schedule(
395 &self,
396 request: Request<generated::ScheduleIdRequest>,
397 ) -> Result<Response<generated::ResumeScheduleResponse>, Status> {
398 plain_rpcs::resume_schedule(&self.state, request).await
399 }
400
401 async fn delete_schedule(
402 &self,
403 request: Request<generated::ScheduleIdRequest>,
404 ) -> Result<Response<generated::DeleteScheduleResponse>, Status> {
405 plain_rpcs::delete_schedule(&self.state, request).await
406 }
407
408 async fn list_schedules(
409 &self,
410 request: Request<generated::ListSchedulesRequest>,
411 ) -> Result<Response<generated::ListSchedulesResponse>, Status> {
412 plain_rpcs::list_schedules(&self.state, request).await
413 }
414
415 async fn describe_schedule(
416 &self,
417 request: Request<generated::ScheduleIdRequest>,
418 ) -> Result<Response<generated::DescribeScheduleResponse>, Status> {
419 plain_rpcs::describe_schedule(&self.state, request).await
420 }
421
422 async fn mint_namespace(
427 &self,
428 request: Request<generated::MintNamespaceRequest>,
429 ) -> Result<Response<generated::MintNamespaceResponse>, Status> {
430 self.mint_namespace_here(request).await
431 }
432}
433
434#[cfg(test)]
435mod tests {
436 use std::{net::SocketAddr, sync::Arc};
437
438 use aion::EngineBuilder;
439 use aion_core::{Event, EventEnvelope, Payload, WorkflowId, WorkflowStatus};
440 use aion_proto::{
441 ProtoWireError, WireError, WireErrorCode,
442 convert::{decode_core_value, encode_core_value},
443 generated::workflow_service_server::WorkflowService,
444 };
445 use aion_store::{
446 EventStore, InMemoryStore, WriteToken,
447 visibility::{VisibilityRecord, VisibilityStore},
448 };
449 use chrono::Utc;
450 use prost::Message;
451 use serde_json::json;
452 use tonic::{Code, Request};
453
454 use super::convert::{decode_envelope, encode_envelope, encode_payload};
455 use super::*;
456 use crate::{
457 NamespaceResolver,
458 config::{
459 AuthConfig, AuthoringConfig, DeployConfig, ListenConfig, MetricsConfig,
460 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, RuntimeConfig,
461 WebSocketConfig, WorkerConfig,
462 },
463 };
464
465 const NAMESPACE: &str = "tenant-a";
466 const TOKEN: &str = "test-token";
467
468 async fn server_state(
473 resolver: NamespaceResolver,
474 runtime: RuntimeConfig,
475 ) -> Result<ServerState, Box<dyn std::error::Error>> {
476 #[cfg(feature = "auth")]
477 {
478 let url = crate::auth::test_support::serve_jwks()?;
479 let refresh = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
480 let cache = crate::auth::JwksCache::new(url, refresh).await?;
481 Ok(ServerState::from_parts_with_jwks(resolver, runtime, cache))
482 }
483 #[cfg(not(feature = "auth"))]
484 {
485 tokio::task::yield_now().await;
487 Ok(ServerState::from_parts(resolver, runtime))
488 }
489 }
490
491 #[tokio::test]
492 async fn in_process_tonic_start_and_list_use_shared_handlers()
493 -> Result<(), Box<dyn std::error::Error>> {
494 let backing = Arc::new(InMemoryStore::default());
495 let store: Arc<dyn EventStore> = backing.clone();
496 let visibility_store: Arc<dyn VisibilityStore> = backing;
497 let engine = Arc::new(
498 EngineBuilder::new()
499 .store_arc(Arc::clone(&store))
500 .visibility_store_arc(Arc::clone(&visibility_store))
501 .scheduler_threads(1)
502 .build()
503 .await?,
504 );
505 store
506 .append(
507 WriteToken::recorder(),
508 &workflow_id(),
509 &[started_event()?],
510 0,
511 )
512 .await?;
513 visibility_store
514 .record_visibility(VisibilityRecord {
515 workflow_id: workflow_id(),
516 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(2)),
517 workflow_type: String::from("fixture"),
518 status: WorkflowStatus::Running,
519 start_time: Utc::now(),
520 close_time: None,
521 failed_step: None,
522 failure_reason: None,
523 search_attributes: std::collections::HashMap::from([(
524 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
525 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
526 )]),
527 })
528 .await?;
529 let resolver = NamespaceResolver::from_config(
530 crate::config::NamespaceConfig {
531 mode: NamespaceMode::SharedEngine,
532 },
533 engine,
534 );
535 let state = server_state(resolver.clone(), runtime_config()).await?;
536 let service = WorkflowGrpcService::new(state);
537
538 let mut start = Request::new(generated::StartWorkflowRequest {
539 namespace: NAMESPACE.to_owned(),
540 workflow_type: "missing-workflow".to_owned(),
541 input: Some(encode_payload(proto_payload()?)),
542 routing_key: None,
543 task_queue: None,
544 display_name: None,
545 });
546 apply_metadata(start.metadata_mut())?;
547 let start_error = service.start_workflow(start).await;
548 let status = start_error
549 .err()
550 .ok_or_else(|| WireError::backend("expected error"))?;
551 assert_eq!(status.code(), Code::NotFound);
552 let detail = ProtoWireError::decode(status.details())?;
553 assert_eq!(detail.error_type.as_deref(), Some("WorkflowTypeNotFound"));
554 assert!(detail.message.contains("missing-workflow"));
555
556 let list_filter = encode_core_value(
557 NAMESPACE,
558 None,
559 &aion_store::visibility::ListWorkflowsFilter {
560 workflow_type: Some(String::from("fixture")),
561 status: Some(WorkflowStatus::Running),
562 ..aion_store::visibility::ListWorkflowsFilter::default()
563 },
564 )?;
565 let mut list = Request::new(generated::ListWorkflowsRequest {
566 namespace: NAMESPACE.to_owned(),
567 filter: Some(encode_envelope(list_filter)),
568 });
569 apply_metadata(list.metadata_mut())?;
570 let response = service.list_workflows(list).await?.into_inner();
571
572 assert_eq!(response.summaries.len(), 1);
573 let summary = response
574 .summaries
575 .into_iter()
576 .next()
577 .map(decode_envelope)
578 .map(|envelope| decode_core_value::<aion_store::visibility::WorkflowSummary>(&envelope))
579 .transpose()?
580 .ok_or_else(|| WireError::backend("summary missing"))?;
581 assert_eq!(summary.workflow_id, workflow_id());
582 assert_eq!(
588 resolver
589 .verify_workflow_ownership(NAMESPACE, &workflow_id())
590 .await
591 .err()
592 .map(|error| error.to_wire_error().code),
593 Some(WireErrorCode::NotFound)
594 );
595 Ok(())
596 }
597
598 async fn rename_fixture(
609 register_display_name: bool,
610 ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
611 rename_fixture_with_terminal(register_display_name, true).await
612 }
613
614 async fn rename_fixture_with_terminal(
618 register_display_name: bool,
619 terminal: bool,
620 ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
621 let backing = Arc::new(InMemoryStore::default());
622 let store: Arc<dyn EventStore> = backing.clone();
623 let visibility_store: Arc<dyn VisibilityStore> = backing;
624 let mut schema = aion_core::SearchAttributeSchema::new();
625 schema.register(
626 crate::namespace::NAMESPACE_ATTRIBUTE,
627 aion_core::SearchAttributeType::String,
628 )?;
629 if register_display_name {
630 schema.register(
631 crate::namespace::DISPLAY_NAME_ATTRIBUTE,
632 aion_core::SearchAttributeType::String,
633 )?;
634 }
635 let engine = Arc::new(
636 EngineBuilder::new()
637 .store_arc(Arc::clone(&store))
638 .visibility_store_arc(Arc::clone(&visibility_store))
639 .search_attribute_schema(schema)
640 .scheduler_threads(1)
641 .build()
642 .await?,
643 );
644 let namespace_event = |seq: u64| Event::SearchAttributesUpdated {
645 envelope: EventEnvelope {
646 seq,
647 recorded_at: Utc::now(),
648 workflow_id: workflow_id(),
649 },
650 workflow_id: workflow_id(),
651 attributes: std::collections::HashMap::from([(
652 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
653 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
654 )]),
655 };
656 let mut events = vec![started_event()?, namespace_event(2)];
657 if terminal {
658 events.push(Event::WorkflowCompleted {
659 envelope: EventEnvelope {
660 seq: 3,
661 recorded_at: Utc::now(),
662 workflow_id: workflow_id(),
663 },
664 result: payload()?,
665 });
666 }
667 store
668 .append(WriteToken::recorder(), &workflow_id(), &events, 0)
669 .await?;
670 let resolver = NamespaceResolver::from_config(
671 crate::config::NamespaceConfig {
672 mode: NamespaceMode::SharedEngine,
673 },
674 engine,
675 );
676 let state = server_state(resolver, runtime_config()).await?;
677 Ok((WorkflowGrpcService::new(state), store))
678 }
679
680 #[tokio::test]
684 async fn in_process_tonic_rename_records_the_name_and_supersedes_it()
685 -> Result<(), Box<dyn std::error::Error>> {
686 let (service, store) = rename_fixture(true).await?;
687
688 let rename = |name: &str| {
689 let mut request = Request::new(generated::RenameRequest {
690 namespace: NAMESPACE.to_owned(),
691 workflow_id: Some(generated::WorkflowId {
692 uuid: workflow_id().to_string(),
693 }),
694 run_id: None,
696 display_name: name.to_owned(),
697 });
698 apply_metadata(request.metadata_mut()).map(|()| request)
699 };
700
701 let response = service.rename(rename(" Nightly settlement ")?).await?;
703 let response = response.into_inner();
704 assert_eq!(response.display_name, "Nightly settlement");
705 assert_eq!(
706 response.run_id.map(|id| id.uuid),
707 Some(aion_core::RunId::new(uuid::Uuid::from_u128(1)).to_string())
708 );
709
710 service
711 .rename(rename("Nightly settlement (rerun)")?)
712 .await?;
713
714 let history = store.read_history(&workflow_id()).await?;
715 assert_eq!(
716 aion_core::display_name(&history).as_deref(),
717 Some("Nightly settlement (rerun)"),
718 "the latest recorded name wins"
719 );
720 let names: Vec<_> = history
721 .iter()
722 .filter_map(|event| match event {
723 Event::SearchAttributesUpdated { attributes, .. } => attributes
724 .get(aion_core::DISPLAY_NAME_ATTRIBUTE)
725 .and_then(|value| match value {
726 aion_core::SearchAttributeValue::String(name) => Some(name.clone()),
727 _ => None,
728 }),
729 _ => None,
730 })
731 .collect();
732 assert_eq!(
733 names,
734 vec![
735 String::from("Nightly settlement"),
736 String::from("Nightly settlement (rerun)")
737 ],
738 "history keeps every name the run has worn"
739 );
740 Ok(())
741 }
742
743 #[tokio::test]
746 async fn in_process_tonic_rename_refuses_a_blank_name() -> Result<(), Box<dyn std::error::Error>>
747 {
748 let (service, store) = rename_fixture(false).await?;
749 let before = store.read_history(&workflow_id()).await?.len();
750
751 let mut request = Request::new(generated::RenameRequest {
752 namespace: NAMESPACE.to_owned(),
753 workflow_id: Some(generated::WorkflowId {
754 uuid: workflow_id().to_string(),
755 }),
756 run_id: None,
757 display_name: String::from(" "),
758 });
759 apply_metadata(request.metadata_mut())?;
760 let status = service
761 .rename(request)
762 .await
763 .err()
764 .ok_or_else(|| WireError::backend("expected a blank-name refusal"))?;
765 assert_eq!(status.code(), Code::InvalidArgument);
766 assert_eq!(
767 store.read_history(&workflow_id()).await?.len(),
768 before,
769 "a refused rename appends nothing"
770 );
771 Ok(())
772 }
773
774 #[tokio::test]
780 async fn in_process_tonic_rename_of_a_non_resident_running_run_is_failed_precondition()
781 -> Result<(), Box<dyn std::error::Error>> {
782 let (service, store) = rename_fixture_with_terminal(true, false).await?;
783 let before = store.read_history(&workflow_id()).await?.len();
784
785 let mut request = Request::new(generated::RenameRequest {
786 namespace: NAMESPACE.to_owned(),
787 workflow_id: Some(generated::WorkflowId {
788 uuid: workflow_id().to_string(),
789 }),
790 run_id: None,
791 display_name: String::from("Nightly settlement"),
792 });
793 apply_metadata(request.metadata_mut())?;
794 let status = service
795 .rename(request)
796 .await
797 .err()
798 .ok_or_else(|| WireError::backend("expected a residency refusal"))?;
799
800 assert_eq!(status.code(), Code::FailedPrecondition);
801 let detail = ProtoWireError::decode(status.details())?;
802 assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
803 assert!(
804 detail.message.contains("not resident"),
805 "the refusal must name the reason: {}",
806 detail.message
807 );
808 assert_eq!(
809 store.read_history(&workflow_id()).await?.len(),
810 before,
811 "a refused rename appends nothing"
812 );
813 Ok(())
814 }
815
816 #[tokio::test]
820 async fn in_process_tonic_reopen_completed_is_failed_precondition_invalid_state()
821 -> Result<(), Box<dyn std::error::Error>> {
822 let backing = Arc::new(InMemoryStore::default());
823 let store: Arc<dyn EventStore> = backing.clone();
824 let visibility_store: Arc<dyn VisibilityStore> = backing;
825 let engine = Arc::new(
826 EngineBuilder::new()
827 .store_arc(Arc::clone(&store))
828 .visibility_store_arc(Arc::clone(&visibility_store))
829 .scheduler_threads(1)
830 .build()
831 .await?,
832 );
833 store
837 .append(
838 WriteToken::recorder(),
839 &workflow_id(),
840 &[
841 started_event()?,
842 Event::SearchAttributesUpdated {
843 envelope: EventEnvelope {
844 seq: 2,
845 recorded_at: Utc::now(),
846 workflow_id: workflow_id(),
847 },
848 workflow_id: workflow_id(),
849 attributes: std::collections::HashMap::from([(
850 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
851 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
852 )]),
853 },
854 Event::WorkflowCompleted {
855 envelope: EventEnvelope {
856 seq: 3,
857 recorded_at: Utc::now(),
858 workflow_id: workflow_id(),
859 },
860 result: payload()?,
861 },
862 ],
863 0,
864 )
865 .await?;
866 let resolver = NamespaceResolver::from_config(
867 crate::config::NamespaceConfig {
868 mode: NamespaceMode::SharedEngine,
869 },
870 engine,
871 );
872 let state = server_state(resolver, runtime_config()).await?;
873 let service = WorkflowGrpcService::new(state);
874
875 let mut reopen = Request::new(generated::ReopenRequest {
876 namespace: NAMESPACE.to_owned(),
877 workflow_id: Some(generated::WorkflowId {
878 uuid: workflow_id().to_string(),
879 }),
880 run_id: None,
881 });
882 apply_metadata(reopen.metadata_mut())?;
883 let status = service
884 .reopen(reopen)
885 .await
886 .err()
887 .ok_or_else(|| WireError::backend("expected a reopen precondition error"))?;
888 assert_eq!(status.code(), Code::FailedPrecondition);
889 let detail = ProtoWireError::decode(status.details())?;
890 assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
891 assert_eq!(
892 detail.code,
893 aion_proto::ProtoWireErrorCode::InvalidState as i32
894 );
895 Ok(())
896 }
897
898 fn apply_metadata(
899 metadata: &mut tonic::metadata::MetadataMap,
900 ) -> Result<(), Box<dyn std::error::Error>> {
901 #[cfg(feature = "auth")]
905 let bearer = crate::auth::test_support::mint_token("alice", NAMESPACE)?;
906 #[cfg(not(feature = "auth"))]
907 let bearer = TOKEN.to_owned();
908 metadata.insert("authorization", format!("Bearer {bearer}").parse()?);
909 metadata.insert("x-aion-subject", "alice".parse()?);
910 metadata.insert("x-aion-namespaces", NAMESPACE.parse()?);
911 Ok(())
912 }
913
914 fn runtime_config() -> RuntimeConfig {
918 RuntimeConfig {
919 listen: ListenConfig {
920 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
921 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
922 },
923 tls: None,
924 auth: AuthConfig {
925 enabled: true,
926 jwks_url: Some(TOKEN.to_owned()),
927 jwks_refresh_seconds: 300,
928 },
929 ops_console: OpsConsoleConfig {
930 source: OpsConsoleAssetSource::Embedded,
931 },
932 namespace: NamespaceConfig {
933 mode: NamespaceMode::SharedEngine,
934 },
935 worker: WorkerConfig {
936 heartbeat_window: std::time::Duration::from_secs(30),
937 ..WorkerConfig::default()
938 },
939 websocket: WebSocketConfig {
940 outbound_buffer_bound: 32,
941 event_broadcast_capacity: Some(64),
942 cluster_broadcast_capacity: Some(64),
943 },
944 workflow_packages: Vec::new(),
945 deploy: DeployConfig::default(),
946 authoring: AuthoringConfig::default(),
947 dev: crate::config::DevConfig::default(),
948 outbox: crate::config::OutboxConfig::default(),
949 observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
950 mcp: crate::config::ResolvedMcpConfig::default(),
951 scheduler_threads: 1,
952 jit_threshold: None,
953 query_timeout: Some(std::time::Duration::from_secs(10)),
954 default_namespace: "default".to_owned(),
955 auto_create: crate::config::AutoCreate::Open,
956 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
957 drain_timeout: std::time::Duration::from_secs(30),
958 metrics: MetricsConfig { enabled: true },
959 owned_shards: Vec::new(),
960 cors_allowed_origins: Vec::new(),
961 }
962 }
963
964 fn started_event() -> Result<Event, aion_core::PayloadError> {
965 Ok(Event::WorkflowStarted {
966 envelope: EventEnvelope {
967 seq: 1,
968 recorded_at: Utc::now(),
969 workflow_id: workflow_id(),
970 },
971 workflow_type: "fixture".to_owned(),
972 input: payload()?,
973 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
974 parent_run_id: None,
975 parent_workflow_id: None,
976 package_version: aion_core::PackageVersion::new("a".repeat(64)),
977 })
978 }
979
980 fn proto_payload() -> Result<aion_proto::ProtoPayload, aion_core::PayloadError> {
981 Ok(payload()?.into())
982 }
983
984 fn payload() -> Result<Payload, aion_core::PayloadError> {
985 Payload::from_json(&json!({ "fixture": "input" }))
986 }
987
988 fn workflow_id() -> WorkflowId {
989 WorkflowId::new(uuid::Uuid::from_u128(1))
990 }
991}