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/// The owner side of namespace-mint routing (`MintNamespace`).
8mod mint_resolve;
9/// Cluster request routing/resolution (steered start + forward-or-local).
10#[cfg(feature = "haematite-backend")]
11mod routing_resolve;
12/// Wire-error-to-tonic-Status mapping.
13mod 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/// Cloneable tonic implementation for workflow management.
41#[derive(Clone)]
42pub struct WorkflowGrpcService {
43    state: ServerState,
44}
45
46impl WorkflowGrpcService {
47    /// Build a tonic workflow service from shared server state.
48    #[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/// Construct the generated tonic server wrapper.
59#[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        // R-4 steered start over R-1 remint: a non-empty routing key steers the
77        // start to its shard owner (forwarding when remote); otherwise the R-1
78        // remint places it locally. `placement` is `None` for single-node /
79        // non-clustered boots (and own-all scope), so the engine mints as usual —
80        // default path unchanged. Without the cluster backend there is no
81        // steering or placement at all.
82        #[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            // Minted-on-use safety net (Phase 1 S6): mint/gate the authorized
91            // namespace before the engine start, sharing the SAME policy the
92            // worker-registration seam applies. Auth-scoped (`guard.scope` runs
93            // inside the handler); placement is the orthogonal steered-start id.
94            // Carry the caller's own credentials onto any forwarded mint, so
95            // the namespace shard's owner authorizes THIS caller rather than an
96            // anonymous peer (the same verbatim-metadata discipline the R-3
97            // request forwarder already applies to a forwarded start).
98            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            // Minted-on-use safety net (Phase 1 S6): mint/gate the authorized
117            // namespace before the engine start, sharing the SAME policy the
118            // worker-registration seam applies.
119            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    /// Mint (or gate) an already-authorized namespace set on THIS node — the
576    /// owner side of namespace-mint routing. See
577    /// [`WorkflowGrpcService::mint_namespace_here`] for the authorization and
578    /// policy-parity contract this delegates to.
579    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    /// Server state whose bearer validation matches the compiled auth path:
622    /// under `feature = "auth"` a real [`crate::auth::JwksCache`] is fetched
623    /// from a live fixture JWKS endpoint; otherwise the development token path
624    /// needs no cache.
625    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            // Yield to preserve the async signature required by the auth-feature branch.
639            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        // The seeded history records no namespace attribute, so durable
735        // ownership verification must reject targeted access with NotFound:
736        // a missing ownership attribute is indistinguishable from a
737        // nonexistent workflow (anti-existence-leak), and NamespaceDenied is
738        // reserved for callers without a grant for the requested namespace.
739        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    /// Reopening a terminal-Completed workflow over the gRPC service surfaces
751    /// the engine's `InvalidState` precondition as tonic `FailedPrecondition`
752    /// carrying the typed `InvalidState` detail (AO-007 C35/C38).
753    #[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        // A terminal-Completed run whose history records its namespace, so the
768        // guard's durable-ownership verification passes and the request reaches
769        // the engine reopen op (which rejects Completed with InvalidState).
770        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        // Bearer credential accepted by the compiled authentication path: a
836        // JWT minted against the fixture JWKS under `feature = "auth"`, the
837        // development shared-secret token otherwise.
838        #[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    /// Test runtime settings with authentication enabled; under
849    /// `feature = "auth"` validation runs against the [`server_state`]-injected
850    /// JWKS cache, so the configured dev-secret `jwks_url` is never fetched.
851    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}