Skip to main content

aion_server/api/grpc/
mod.rs

1//! tonic workflow service adapter.
2
3/// Caller-identity extraction from gRPC request metadata.
4mod auth;
5/// Codec conversions between the generated wire messages and `aion-proto` types.
6mod convert;
7/// Cluster request routing/resolution (steered start + forward-or-local).
8#[cfg(feature = "haematite-backend")]
9mod routing_resolve;
10/// Wire-error-to-tonic-Status mapping.
11mod 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/// Cloneable tonic implementation for workflow management.
38#[derive(Clone)]
39pub struct WorkflowGrpcService {
40    state: ServerState,
41}
42
43impl WorkflowGrpcService {
44    /// Build a tonic workflow service from shared server state.
45    #[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/// Construct the generated tonic server wrapper.
56#[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        // R-4 steered start over R-1 remint: a non-empty routing key steers the
74        // start to its shard owner (forwarding when remote); otherwise the R-1
75        // remint places it locally. `placement` is `None` for single-node /
76        // non-clustered boots (and own-all scope), so the engine mints as usual —
77        // default path unchanged. Without the cluster backend there is no
78        // steering or placement at all.
79        #[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            // Minted-on-use safety net (Phase 1 S6): mint/gate the authorized
88            // namespace before the engine start, sharing the SAME policy the
89            // worker-registration seam applies. Auth-scoped (`guard.scope` runs
90            // inside the handler); placement is the orthogonal steered-start id.
91            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            // Minted-on-use safety net (Phase 1 S6): mint/gate the authorized
107            // namespace before the engine start, sharing the SAME policy the
108            // worker-registration seam applies.
109            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    /// Server state whose bearer validation matches the compiled auth path:
601    /// under `feature = "auth"` a real [`crate::auth::JwksCache`] is fetched
602    /// from a live fixture JWKS endpoint; otherwise the development token path
603    /// needs no cache.
604    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            // Yield to preserve the async signature required by the auth-feature branch.
618            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        // The seeded history records no namespace attribute, so durable
714        // ownership verification must reject targeted access with NotFound:
715        // a missing ownership attribute is indistinguishable from a
716        // nonexistent workflow (anti-existence-leak), and NamespaceDenied is
717        // reserved for callers without a grant for the requested namespace.
718        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    /// Reopening a terminal-Completed workflow over the gRPC service surfaces
730    /// the engine's `InvalidState` precondition as tonic `FailedPrecondition`
731    /// carrying the typed `InvalidState` detail (AO-007 C35/C38).
732    #[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        // A terminal-Completed run whose history records its namespace, so the
747        // guard's durable-ownership verification passes and the request reaches
748        // the engine reopen op (which rejects Completed with InvalidState).
749        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        // Bearer credential accepted by the compiled authentication path: a
815        // JWT minted against the fixture JWKS under `feature = "auth"`, the
816        // development shared-secret token otherwise.
817        #[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    /// Test runtime settings with authentication enabled; under
828    /// `feature = "auth"` validation runs against the [`server_state`]-injected
829    /// JWKS cache, so the configured dev-secret `jwks_url` is never fetched.
830    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_millis(30_000),
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_millis(10_000)),
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}