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