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 retire_workloop(
213 &self,
214 request: Request<generated::RetireWorkloopRequest>,
215 ) -> Result<Response<generated::RetireWorkloopResponse>, Status> {
216 let caller = self.caller(&request).await?;
217 let (metadata, _ext, inner) = request.into_parts();
218 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
219 match self
220 .resolve_route(
221 workflow_id,
222 &metadata,
223 ForwardRequest::RetireWorkloop(inner.clone()),
224 )
225 .await
226 {
227 RouteResolution::Reject(status) => Err(status),
228 RouteResolution::Reply(ForwardReply::RetireWorkloop(reply)) => Ok(Response::new(reply)),
229 RouteResolution::Reply(_) => {
230 Err(Status::internal("forwarder returned a mismatched reply"))
231 }
232 RouteResolution::Local => {
233 let id = inner
234 .workflow_id
235 .clone()
236 .map(decode_workflow_id)
237 .map(|id| id.uuid)
238 .ok_or_else(|| Status::invalid_argument("workflow id is required"))?;
239 let (_, reason) = handlers::retire_workloop(
240 self.state.namespace_guard(),
241 &caller,
242 inner.namespace,
243 id,
244 inner.reason,
245 )
246 .await
247 .map_err(status_from_wire_error)?;
248 Ok(Response::new(generated::RetireWorkloopResponse { reason }))
249 }
250 }
251 }
252
253 async fn reopen(
254 &self,
255 request: Request<generated::ReopenRequest>,
256 ) -> Result<Response<generated::ReopenResponse>, Status> {
257 let caller = self.caller(&request).await?;
258 let (metadata, _ext, inner) = request.into_parts();
259 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
260 match self
261 .resolve_route(
262 workflow_id,
263 &metadata,
264 ForwardRequest::Reopen(inner.clone()),
265 )
266 .await
267 {
268 RouteResolution::Reject(status) => return Err(status),
269 RouteResolution::Reply(ForwardReply::Reopen(reply)) => {
270 return Ok(Response::new(reply));
271 }
272 RouteResolution::Reply(_) => {
273 return Err(Status::internal("forwarder returned a mismatched reply"));
274 }
275 RouteResolution::Local => {
276 let response = handlers::reopen(
277 self.state.namespace_guard(),
278 &caller,
279 decode_reopen_request(inner),
280 )
281 .await
282 .map_err(status_from_wire_error)?;
283 return Ok(Response::new(encode_reopen_response(response)));
284 }
285 }
286 }
287
288 async fn pause(
289 &self,
290 request: Request<generated::PauseRequest>,
291 ) -> Result<Response<generated::PauseResponse>, Status> {
292 let caller = self.caller(&request).await?;
293 let (metadata, _ext, inner) = request.into_parts();
294 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
295 match self
296 .resolve_route(workflow_id, &metadata, ForwardRequest::Pause(inner.clone()))
297 .await
298 {
299 RouteResolution::Reject(status) => return Err(status),
300 RouteResolution::Reply(ForwardReply::Pause(reply)) => {
301 return Ok(Response::new(reply));
302 }
303 RouteResolution::Reply(_) => {
304 return Err(Status::internal("forwarder returned a mismatched reply"));
305 }
306 RouteResolution::Local => {
307 let response = handlers::pause(
308 self.state.namespace_guard(),
309 &caller,
310 decode_pause_request(inner),
311 )
312 .await
313 .map_err(status_from_wire_error)?;
314 return Ok(Response::new(encode_pause_response(response)));
315 }
316 }
317 }
318
319 async fn resume(
320 &self,
321 request: Request<generated::ResumeRequest>,
322 ) -> Result<Response<generated::ResumeResponse>, Status> {
323 let caller = self.caller(&request).await?;
324 let (metadata, _ext, inner) = request.into_parts();
325 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
326 match self
327 .resolve_route(
328 workflow_id,
329 &metadata,
330 ForwardRequest::Resume(inner.clone()),
331 )
332 .await
333 {
334 RouteResolution::Reject(status) => return Err(status),
335 RouteResolution::Reply(ForwardReply::Resume(reply)) => {
336 return Ok(Response::new(reply));
337 }
338 RouteResolution::Reply(_) => {
339 return Err(Status::internal("forwarder returned a mismatched reply"));
340 }
341 RouteResolution::Local => {
342 let response = handlers::resume(
343 self.state.namespace_guard(),
344 &caller,
345 decode_resume_request(inner),
346 )
347 .await
348 .map_err(status_from_wire_error)?;
349 return Ok(Response::new(encode_resume_response(response)));
350 }
351 }
352 }
353
354 async fn rename(
355 &self,
356 request: Request<generated::RenameRequest>,
357 ) -> Result<Response<generated::RenameResponse>, Status> {
358 let caller = self.caller(&request).await?;
359 let (metadata, _ext, inner) = request.into_parts();
360 let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
361 match self
362 .resolve_route(
363 workflow_id,
364 &metadata,
365 ForwardRequest::Rename(inner.clone()),
366 )
367 .await
368 {
369 RouteResolution::Reject(status) => return Err(status),
370 RouteResolution::Reply(ForwardReply::Rename(reply)) => {
371 return Ok(Response::new(reply));
372 }
373 RouteResolution::Reply(_) => {
374 return Err(Status::internal("forwarder returned a mismatched reply"));
375 }
376 RouteResolution::Local => {
377 let response = handlers::rename(
378 self.state.namespace_guard(),
379 &caller,
380 decode_rename_request(inner),
381 )
382 .await
383 .map_err(status_from_wire_error)?;
384 return Ok(Response::new(encode_rename_response(response)));
385 }
386 }
387 }
388
389 async fn list_workflows(
390 &self,
391 request: Request<generated::ListWorkflowsRequest>,
392 ) -> Result<Response<generated::ListWorkflowsResponse>, Status> {
393 plain_rpcs::list_workflows(&self.state, request).await
394 }
395
396 async fn count_workflows(
397 &self,
398 request: Request<generated::CountWorkflowsRequest>,
399 ) -> Result<Response<generated::CountWorkflowsResponse>, Status> {
400 plain_rpcs::count_workflows(&self.state, request).await
401 }
402
403 async fn describe_workflow(
404 &self,
405 request: Request<generated::DescribeWorkflowRequest>,
406 ) -> Result<Response<generated::DescribeWorkflowResponse>, Status> {
407 plain_rpcs::describe_workflow(&self.state, request).await
408 }
409
410 async fn read_history(
411 &self,
412 request: Request<generated::ReadHistoryRequest>,
413 ) -> Result<Response<generated::ReadHistoryResponse>, Status> {
414 plain_rpcs::read_history(&self.state, request).await
415 }
416
417 async fn create_schedule(
418 &self,
419 request: Request<generated::CreateScheduleRequest>,
420 ) -> Result<Response<generated::CreateScheduleResponse>, Status> {
421 plain_rpcs::create_schedule(&self.state, request).await
422 }
423
424 async fn update_schedule(
425 &self,
426 request: Request<generated::UpdateScheduleRequest>,
427 ) -> Result<Response<generated::UpdateScheduleResponse>, Status> {
428 plain_rpcs::update_schedule(&self.state, request).await
429 }
430
431 async fn pause_schedule(
432 &self,
433 request: Request<generated::ScheduleIdRequest>,
434 ) -> Result<Response<generated::PauseScheduleResponse>, Status> {
435 plain_rpcs::pause_schedule(&self.state, request).await
436 }
437
438 async fn resume_schedule(
439 &self,
440 request: Request<generated::ScheduleIdRequest>,
441 ) -> Result<Response<generated::ResumeScheduleResponse>, Status> {
442 plain_rpcs::resume_schedule(&self.state, request).await
443 }
444
445 async fn delete_schedule(
446 &self,
447 request: Request<generated::ScheduleIdRequest>,
448 ) -> Result<Response<generated::DeleteScheduleResponse>, Status> {
449 plain_rpcs::delete_schedule(&self.state, request).await
450 }
451
452 async fn list_schedules(
453 &self,
454 request: Request<generated::ListSchedulesRequest>,
455 ) -> Result<Response<generated::ListSchedulesResponse>, Status> {
456 plain_rpcs::list_schedules(&self.state, request).await
457 }
458
459 async fn describe_schedule(
460 &self,
461 request: Request<generated::ScheduleIdRequest>,
462 ) -> Result<Response<generated::DescribeScheduleResponse>, Status> {
463 plain_rpcs::describe_schedule(&self.state, request).await
464 }
465
466 async fn mint_namespace(
471 &self,
472 request: Request<generated::MintNamespaceRequest>,
473 ) -> Result<Response<generated::MintNamespaceResponse>, Status> {
474 self.mint_namespace_here(request).await
475 }
476}
477
478#[cfg(test)]
479mod tests {
480 use std::{net::SocketAddr, sync::Arc};
481
482 use aion::EngineBuilder;
483 use aion_core::{Event, EventEnvelope, Payload, WorkflowId, WorkflowStatus};
484 use aion_proto::{
485 ProtoWireError, WireError, WireErrorCode,
486 convert::{decode_core_value, encode_core_value},
487 generated::workflow_service_server::WorkflowService,
488 };
489 use aion_store::{
490 EventStore, InMemoryStore, WriteToken,
491 visibility::{VisibilityRecord, VisibilityStore},
492 };
493 use chrono::Utc;
494 use prost::Message;
495 use serde_json::json;
496 use tonic::{Code, Request};
497
498 use super::convert::{decode_envelope, encode_envelope, encode_payload};
499 use super::*;
500 use crate::{
501 NamespaceResolver,
502 config::{
503 AuthConfig, AuthoringConfig, DeployConfig, ListenConfig, MetricsConfig,
504 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, RuntimeConfig,
505 WebSocketConfig, WorkerConfig,
506 },
507 };
508
509 const NAMESPACE: &str = "tenant-a";
510 const TOKEN: &str = "test-token";
511
512 async fn server_state(
517 resolver: NamespaceResolver,
518 runtime: RuntimeConfig,
519 ) -> Result<ServerState, Box<dyn std::error::Error>> {
520 #[cfg(feature = "auth")]
521 {
522 let url = crate::auth::test_support::serve_jwks()?;
523 let refresh = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
524 let cache = crate::auth::JwksCache::new(url, refresh).await?;
525 Ok(ServerState::from_parts_with_jwks(resolver, runtime, cache))
526 }
527 #[cfg(not(feature = "auth"))]
528 {
529 tokio::task::yield_now().await;
531 Ok(ServerState::from_parts(resolver, runtime))
532 }
533 }
534
535 #[tokio::test]
536 async fn in_process_tonic_start_and_list_use_shared_handlers()
537 -> Result<(), Box<dyn std::error::Error>> {
538 let backing = Arc::new(InMemoryStore::default());
539 let store: Arc<dyn EventStore> = backing.clone();
540 let visibility_store: Arc<dyn VisibilityStore> = backing;
541 let engine = Arc::new(
542 EngineBuilder::new()
543 .store_arc(Arc::clone(&store))
544 .visibility_store_arc(Arc::clone(&visibility_store))
545 .scheduler_threads(1)
546 .build()
547 .await?,
548 );
549 store
550 .append(
551 WriteToken::recorder(),
552 &workflow_id(),
553 &[started_event()?],
554 0,
555 )
556 .await?;
557 visibility_store
558 .record_visibility(VisibilityRecord {
559 workflow_id: workflow_id(),
560 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(2)),
561 workflow_type: String::from("fixture"),
562 status: WorkflowStatus::Running,
563 start_time: Utc::now(),
564 close_time: None,
565 failed_step: None,
566 failure_reason: None,
567 search_attributes: std::collections::HashMap::from([(
568 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
569 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
570 )]),
571 })
572 .await?;
573 let resolver = NamespaceResolver::from_config(
574 crate::config::NamespaceConfig {
575 mode: NamespaceMode::SharedEngine,
576 },
577 engine,
578 );
579 let state = server_state(resolver.clone(), runtime_config()).await?;
580 let service = WorkflowGrpcService::new(state);
581
582 let mut start = Request::new(generated::StartWorkflowRequest {
583 namespace: NAMESPACE.to_owned(),
584 workflow_type: "missing-workflow".to_owned(),
585 input: Some(encode_payload(proto_payload()?)),
586 routing_key: None,
587 task_queue: None,
588 display_name: None,
589 });
590 apply_metadata(start.metadata_mut())?;
591 let start_error = service.start_workflow(start).await;
592 let status = start_error
593 .err()
594 .ok_or_else(|| WireError::backend("expected error"))?;
595 assert_eq!(status.code(), Code::NotFound);
596 let detail = ProtoWireError::decode(status.details())?;
597 assert_eq!(detail.error_type.as_deref(), Some("WorkflowTypeNotFound"));
598 assert!(detail.message.contains("missing-workflow"));
599
600 let list_filter = encode_core_value(
601 NAMESPACE,
602 None,
603 &aion_store::visibility::ListWorkflowsFilter {
604 workflow_type: Some(String::from("fixture")),
605 status: Some(WorkflowStatus::Running),
606 ..aion_store::visibility::ListWorkflowsFilter::default()
607 },
608 )?;
609 let mut list = Request::new(generated::ListWorkflowsRequest {
610 namespace: NAMESPACE.to_owned(),
611 filter: Some(encode_envelope(list_filter)),
612 });
613 apply_metadata(list.metadata_mut())?;
614 let response = service.list_workflows(list).await?.into_inner();
615
616 assert_eq!(response.summaries.len(), 1);
617 let summary = response
618 .summaries
619 .into_iter()
620 .next()
621 .map(decode_envelope)
622 .map(|envelope| decode_core_value::<aion_store::visibility::WorkflowSummary>(&envelope))
623 .transpose()?
624 .ok_or_else(|| WireError::backend("summary missing"))?;
625 assert_eq!(summary.workflow_id, workflow_id());
626 assert_eq!(
632 resolver
633 .verify_workflow_ownership(NAMESPACE, &workflow_id())
634 .await
635 .err()
636 .map(|error| error.to_wire_error().code),
637 Some(WireErrorCode::NotFound)
638 );
639 Ok(())
640 }
641
642 async fn rename_fixture(
653 register_display_name: bool,
654 ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
655 rename_fixture_with_terminal(register_display_name, true).await
656 }
657
658 async fn rename_fixture_with_terminal(
662 register_display_name: bool,
663 terminal: bool,
664 ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
665 let backing = Arc::new(InMemoryStore::default());
666 let store: Arc<dyn EventStore> = backing.clone();
667 let visibility_store: Arc<dyn VisibilityStore> = backing;
668 let mut schema = aion_core::SearchAttributeSchema::new();
669 schema.register(
670 crate::namespace::NAMESPACE_ATTRIBUTE,
671 aion_core::SearchAttributeType::String,
672 )?;
673 if register_display_name {
674 schema.register(
675 crate::namespace::DISPLAY_NAME_ATTRIBUTE,
676 aion_core::SearchAttributeType::String,
677 )?;
678 }
679 let engine = Arc::new(
680 EngineBuilder::new()
681 .store_arc(Arc::clone(&store))
682 .visibility_store_arc(Arc::clone(&visibility_store))
683 .search_attribute_schema(schema)
684 .scheduler_threads(1)
685 .build()
686 .await?,
687 );
688 let namespace_event = |seq: u64| Event::SearchAttributesUpdated {
689 envelope: EventEnvelope {
690 seq,
691 recorded_at: Utc::now(),
692 workflow_id: workflow_id(),
693 },
694 workflow_id: workflow_id(),
695 attributes: std::collections::HashMap::from([(
696 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
697 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
698 )]),
699 };
700 let mut events = vec![started_event()?, namespace_event(2)];
701 if terminal {
702 events.push(Event::WorkflowCompleted {
703 envelope: EventEnvelope {
704 seq: 3,
705 recorded_at: Utc::now(),
706 workflow_id: workflow_id(),
707 },
708 result: payload()?,
709 });
710 }
711 store
712 .append(WriteToken::recorder(), &workflow_id(), &events, 0)
713 .await?;
714 let resolver = NamespaceResolver::from_config(
715 crate::config::NamespaceConfig {
716 mode: NamespaceMode::SharedEngine,
717 },
718 engine,
719 );
720 let state = server_state(resolver, runtime_config()).await?;
721 Ok((WorkflowGrpcService::new(state), store))
722 }
723
724 #[tokio::test]
728 async fn in_process_tonic_rename_records_the_name_and_supersedes_it()
729 -> Result<(), Box<dyn std::error::Error>> {
730 let (service, store) = rename_fixture(true).await?;
731
732 let rename = |name: &str| {
733 let mut request = Request::new(generated::RenameRequest {
734 namespace: NAMESPACE.to_owned(),
735 workflow_id: Some(generated::WorkflowId {
736 uuid: workflow_id().to_string(),
737 }),
738 run_id: None,
740 display_name: name.to_owned(),
741 });
742 apply_metadata(request.metadata_mut()).map(|()| request)
743 };
744
745 let response = service.rename(rename(" Nightly settlement ")?).await?;
747 let response = response.into_inner();
748 assert_eq!(response.display_name, "Nightly settlement");
749 assert_eq!(
750 response.run_id.map(|id| id.uuid),
751 Some(aion_core::RunId::new(uuid::Uuid::from_u128(1)).to_string())
752 );
753
754 service
755 .rename(rename("Nightly settlement (rerun)")?)
756 .await?;
757
758 let history = store.read_history(&workflow_id()).await?;
759 assert_eq!(
760 aion_core::display_name(&history).as_deref(),
761 Some("Nightly settlement (rerun)"),
762 "the latest recorded name wins"
763 );
764 let names: Vec<_> = history
765 .iter()
766 .filter_map(|event| match event {
767 Event::SearchAttributesUpdated { attributes, .. } => attributes
768 .get(aion_core::DISPLAY_NAME_ATTRIBUTE)
769 .and_then(|value| match value {
770 aion_core::SearchAttributeValue::String(name) => Some(name.clone()),
771 _ => None,
772 }),
773 _ => None,
774 })
775 .collect();
776 assert_eq!(
777 names,
778 vec![
779 String::from("Nightly settlement"),
780 String::from("Nightly settlement (rerun)")
781 ],
782 "history keeps every name the run has worn"
783 );
784 Ok(())
785 }
786
787 #[tokio::test]
790 async fn in_process_tonic_rename_refuses_a_blank_name() -> Result<(), Box<dyn std::error::Error>>
791 {
792 let (service, store) = rename_fixture(false).await?;
793 let before = store.read_history(&workflow_id()).await?.len();
794
795 let mut request = Request::new(generated::RenameRequest {
796 namespace: NAMESPACE.to_owned(),
797 workflow_id: Some(generated::WorkflowId {
798 uuid: workflow_id().to_string(),
799 }),
800 run_id: None,
801 display_name: String::from(" "),
802 });
803 apply_metadata(request.metadata_mut())?;
804 let status = service
805 .rename(request)
806 .await
807 .err()
808 .ok_or_else(|| WireError::backend("expected a blank-name refusal"))?;
809 assert_eq!(status.code(), Code::InvalidArgument);
810 assert_eq!(
811 store.read_history(&workflow_id()).await?.len(),
812 before,
813 "a refused rename appends nothing"
814 );
815 Ok(())
816 }
817
818 #[tokio::test]
824 async fn in_process_tonic_rename_of_a_non_resident_running_run_is_failed_precondition()
825 -> Result<(), Box<dyn std::error::Error>> {
826 let (service, store) = rename_fixture_with_terminal(true, false).await?;
827 let before = store.read_history(&workflow_id()).await?.len();
828
829 let mut request = Request::new(generated::RenameRequest {
830 namespace: NAMESPACE.to_owned(),
831 workflow_id: Some(generated::WorkflowId {
832 uuid: workflow_id().to_string(),
833 }),
834 run_id: None,
835 display_name: String::from("Nightly settlement"),
836 });
837 apply_metadata(request.metadata_mut())?;
838 let status = service
839 .rename(request)
840 .await
841 .err()
842 .ok_or_else(|| WireError::backend("expected a residency refusal"))?;
843
844 assert_eq!(status.code(), Code::FailedPrecondition);
845 let detail = ProtoWireError::decode(status.details())?;
846 assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
847 assert!(
848 detail.message.contains("not resident"),
849 "the refusal must name the reason: {}",
850 detail.message
851 );
852 assert_eq!(
853 store.read_history(&workflow_id()).await?.len(),
854 before,
855 "a refused rename appends nothing"
856 );
857 Ok(())
858 }
859
860 #[tokio::test]
864 async fn in_process_tonic_reopen_completed_is_failed_precondition_invalid_state()
865 -> Result<(), Box<dyn std::error::Error>> {
866 let backing = Arc::new(InMemoryStore::default());
867 let store: Arc<dyn EventStore> = backing.clone();
868 let visibility_store: Arc<dyn VisibilityStore> = backing;
869 let engine = Arc::new(
870 EngineBuilder::new()
871 .store_arc(Arc::clone(&store))
872 .visibility_store_arc(Arc::clone(&visibility_store))
873 .scheduler_threads(1)
874 .build()
875 .await?,
876 );
877 store
881 .append(
882 WriteToken::recorder(),
883 &workflow_id(),
884 &[
885 started_event()?,
886 Event::SearchAttributesUpdated {
887 envelope: EventEnvelope {
888 seq: 2,
889 recorded_at: Utc::now(),
890 workflow_id: workflow_id(),
891 },
892 workflow_id: workflow_id(),
893 attributes: std::collections::HashMap::from([(
894 crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
895 aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
896 )]),
897 },
898 Event::WorkflowCompleted {
899 envelope: EventEnvelope {
900 seq: 3,
901 recorded_at: Utc::now(),
902 workflow_id: workflow_id(),
903 },
904 result: payload()?,
905 },
906 ],
907 0,
908 )
909 .await?;
910 let resolver = NamespaceResolver::from_config(
911 crate::config::NamespaceConfig {
912 mode: NamespaceMode::SharedEngine,
913 },
914 engine,
915 );
916 let state = server_state(resolver, runtime_config()).await?;
917 let service = WorkflowGrpcService::new(state);
918
919 let mut reopen = Request::new(generated::ReopenRequest {
920 namespace: NAMESPACE.to_owned(),
921 workflow_id: Some(generated::WorkflowId {
922 uuid: workflow_id().to_string(),
923 }),
924 run_id: None,
925 });
926 apply_metadata(reopen.metadata_mut())?;
927 let status = service
928 .reopen(reopen)
929 .await
930 .err()
931 .ok_or_else(|| WireError::backend("expected a reopen precondition error"))?;
932 assert_eq!(status.code(), Code::FailedPrecondition);
933 let detail = ProtoWireError::decode(status.details())?;
934 assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
935 assert_eq!(
936 detail.code,
937 aion_proto::ProtoWireErrorCode::InvalidState as i32
938 );
939 Ok(())
940 }
941
942 fn apply_metadata(
943 metadata: &mut tonic::metadata::MetadataMap,
944 ) -> Result<(), Box<dyn std::error::Error>> {
945 #[cfg(feature = "auth")]
949 let bearer = crate::auth::test_support::mint_token("alice", NAMESPACE)?;
950 #[cfg(not(feature = "auth"))]
951 let bearer = TOKEN.to_owned();
952 metadata.insert("authorization", format!("Bearer {bearer}").parse()?);
953 metadata.insert("x-aion-subject", "alice".parse()?);
954 metadata.insert("x-aion-namespaces", NAMESPACE.parse()?);
955 Ok(())
956 }
957
958 fn runtime_config() -> RuntimeConfig {
962 RuntimeConfig {
963 listen: ListenConfig {
964 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
965 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
966 },
967 tls: None,
968 auth: AuthConfig {
969 enabled: true,
970 jwks_url: Some(TOKEN.to_owned()),
971 jwks_refresh_seconds: 300,
972 },
973 ops_console: OpsConsoleConfig {
974 source: OpsConsoleAssetSource::Embedded,
975 },
976 namespace: NamespaceConfig {
977 mode: NamespaceMode::SharedEngine,
978 },
979 worker: WorkerConfig {
980 heartbeat_window: std::time::Duration::from_secs(30),
981 ..WorkerConfig::default()
982 },
983 websocket: WebSocketConfig {
984 outbound_buffer_bound: 32,
985 event_broadcast_capacity: Some(64),
986 cluster_broadcast_capacity: Some(64),
987 },
988 workflow_packages: Vec::new(),
989 deploy: DeployConfig::default(),
990 authoring: AuthoringConfig::default(),
991 dev: crate::config::DevConfig::default(),
992 outbox: crate::config::OutboxConfig::default(),
993 observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
994 mcp: crate::config::ResolvedMcpConfig::default(),
995 scheduler_threads: 1,
996 jit_threshold: None,
997 query_timeout: Some(std::time::Duration::from_secs(10)),
998 workloop_sweep_interval: Some(std::time::Duration::from_millis(50)),
999 default_namespace: "default".to_owned(),
1000 auto_create: crate::config::AutoCreate::Open,
1001 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
1002 drain_timeout: std::time::Duration::from_secs(30),
1003 metrics: MetricsConfig { enabled: true },
1004 owned_shards: Vec::new(),
1005 cors_allowed_origins: Vec::new(),
1006 }
1007 }
1008
1009 fn started_event() -> Result<Event, aion_core::PayloadError> {
1010 Ok(Event::WorkflowStarted {
1011 envelope: EventEnvelope {
1012 seq: 1,
1013 recorded_at: Utc::now(),
1014 workflow_id: workflow_id(),
1015 },
1016 workflow_type: "fixture".to_owned(),
1017 input: payload()?,
1018 run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
1019 parent_run_id: None,
1020 parent_workflow_id: None,
1021 package_version: aion_core::PackageVersion::new("a".repeat(64)),
1022 })
1023 }
1024
1025 fn proto_payload() -> Result<aion_proto::ProtoPayload, aion_core::PayloadError> {
1026 Ok(payload()?.into())
1027 }
1028
1029 fn payload() -> Result<Payload, aion_core::PayloadError> {
1030 Payload::from_json(&json!({ "fixture": "input" }))
1031 }
1032
1033 fn workflow_id() -> WorkflowId {
1034 WorkflowId::new(uuid::Uuid::from_u128(1))
1035 }
1036}