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