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.
6/// Proto/generated wire conversions. Crate-visible because the HTTP edge
7/// reuses the rename conversions when it forwards to a shard owner (#211).
8pub(crate) mod convert;
9/// The owner side of namespace-mint routing (`MintNamespace`).
10mod mint_resolve;
11/// The pure-delegation RPC bodies (no routing, no forwarding, no cfg arms).
12mod plain_rpcs;
13/// Cluster request routing/resolution (steered start + forward-or-local).
14mod routing_resolve;
15/// Wire-error-to-tonic-Status mapping.
16mod status;
17
18pub(crate) use auth::caller_from_metadata;
19pub(crate) use status::{status_from_wire_error, status_with_code};
20
21use aion_proto::generated::{self, workflow_service_server::WorkflowServiceServer};
22use tonic::{Request, Response, Status};
23
24use crate::namespace::MintCredentials;
25use crate::routing::{ForwardReply, ForwardRequest};
26use crate::{CallerIdentity, ServerState, api::handlers};
27use convert::decode_workflow_id;
28use convert::{
29    decode_cancel_request, decode_pause_request, decode_query_request, decode_rename_request,
30    decode_reopen_request, decode_resume_request, decode_signal_request, decode_start_request,
31    encode_cancel_response, encode_pause_response, encode_query_response, encode_rename_response,
32    encode_reopen_response, encode_resume_response, encode_signal_response, encode_start_response,
33};
34use routing_resolve::{RouteResolution, StartResolution};
35
36/// Cloneable tonic implementation for workflow management.
37#[derive(Clone)]
38pub struct WorkflowGrpcService {
39    state: ServerState,
40}
41
42impl WorkflowGrpcService {
43    /// Build a tonic workflow service from shared server state.
44    #[must_use]
45    pub const fn new(state: ServerState) -> Self {
46        Self { state }
47    }
48
49    async fn caller<T>(&self, request: &Request<T>) -> Result<CallerIdentity, Status> {
50        caller_from_metadata(request.metadata(), &self.state).await
51    }
52}
53
54/// Construct the generated tonic server wrapper.
55#[must_use]
56pub fn workflow_service(state: ServerState) -> WorkflowServiceServer<WorkflowGrpcService> {
57    WorkflowServiceServer::new(WorkflowGrpcService::new(state))
58}
59
60#[tonic::async_trait]
61impl generated::workflow_service_server::WorkflowService for WorkflowGrpcService {
62    async fn start_workflow(
63        &self,
64        request: Request<generated::StartWorkflowRequest>,
65    ) -> Result<Response<generated::StartWorkflowResponse>, Status> {
66        if self.state.drain_state().is_draining() {
67            return Err(Status::unavailable(
68                "server is draining and not accepting new workflow starts",
69            ));
70        }
71        let caller = self.caller(&request).await?;
72        // R-4 steered start over R-1 remint: a non-empty routing key steers the
73        // start to its shard owner (forwarding when remote); otherwise the R-1
74        // remint places it locally. `placement` is `None` for single-node /
75        // non-clustered boots (and own-all scope), so the engine mints as usual —
76        // default path unchanged.
77        let (metadata, _ext, inner) = request.into_parts();
78        let placement = match self.resolve_start(&inner, &metadata).await {
79            StartResolution::Reject(status) => return Err(status),
80            StartResolution::Reply(reply) => return Ok(Response::new(reply)),
81            StartResolution::Local(placement) => placement,
82        };
83        // Minted-on-use safety net (Phase 1 S6): mint/gate the authorized
84        // namespace before the engine start, sharing the SAME policy the
85        // worker-registration seam applies. Auth-scoped (`guard.scope` runs
86        // inside the handler); placement is the orthogonal steered-start id.
87        // Carry the caller's own credentials onto any forwarded mint, so
88        // the namespace shard's owner authorizes THIS caller rather than an
89        // anonymous peer (the same verbatim-metadata discipline the R-3
90        // request forwarder already applies to a forwarded start).
91        let minter = self
92            .state
93            .namespace_minter()
94            .with_caller_credentials(MintCredentials::from_grpc_metadata(&metadata));
95        let response = handlers::start_with_placement(
96            self.state.namespace_guard(),
97            &caller,
98            decode_start_request(inner),
99            placement,
100            Some(&minter),
101        )
102        .await
103        .map_err(status_from_wire_error)?;
104        return Ok(Response::new(encode_start_response(response)));
105    }
106
107    async fn signal(
108        &self,
109        request: Request<generated::SignalRequest>,
110    ) -> Result<Response<generated::SignalResponse>, Status> {
111        let caller = self.caller(&request).await?;
112        let (metadata, _ext, inner) = request.into_parts();
113        let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
114        match self
115            .resolve_route(
116                workflow_id,
117                &metadata,
118                ForwardRequest::Signal(inner.clone()),
119            )
120            .await
121        {
122            RouteResolution::Reject(status) => return Err(status),
123            RouteResolution::Reply(ForwardReply::Signal(reply)) => {
124                return Ok(Response::new(reply));
125            }
126            RouteResolution::Reply(_) => {
127                return Err(Status::internal("forwarder returned a mismatched reply"));
128            }
129            RouteResolution::Local => {
130                let response = handlers::signal(
131                    self.state.namespace_guard(),
132                    &caller,
133                    decode_signal_request(inner),
134                )
135                .await
136                .map_err(status_from_wire_error)?;
137                return Ok(Response::new(encode_signal_response(response)));
138            }
139        }
140    }
141
142    async fn query(
143        &self,
144        request: Request<generated::QueryRequest>,
145    ) -> Result<Response<generated::QueryResponse>, Status> {
146        let caller = self.caller(&request).await?;
147        let (metadata, _ext, inner) = request.into_parts();
148        let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
149        match self
150            .resolve_route(workflow_id, &metadata, ForwardRequest::Query(inner.clone()))
151            .await
152        {
153            RouteResolution::Reject(status) => return Err(status),
154            RouteResolution::Reply(ForwardReply::Query(reply)) => {
155                return Ok(Response::new(reply));
156            }
157            RouteResolution::Reply(_) => {
158                return Err(Status::internal("forwarder returned a mismatched reply"));
159            }
160            RouteResolution::Local => {
161                let response = handlers::query(
162                    self.state.namespace_guard(),
163                    &caller,
164                    decode_query_request(inner),
165                )
166                .await
167                .map_err(status_from_wire_error)?;
168                return Ok(Response::new(encode_query_response(response)));
169            }
170        }
171    }
172
173    async fn cancel(
174        &self,
175        request: Request<generated::CancelRequest>,
176    ) -> Result<Response<generated::CancelResponse>, Status> {
177        let caller = self.caller(&request).await?;
178        let (metadata, _ext, inner) = request.into_parts();
179        let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
180        match self
181            .resolve_route(
182                workflow_id,
183                &metadata,
184                ForwardRequest::Cancel(inner.clone()),
185            )
186            .await
187        {
188            RouteResolution::Reject(status) => return Err(status),
189            RouteResolution::Reply(ForwardReply::Cancel(reply)) => {
190                return Ok(Response::new(reply));
191            }
192            RouteResolution::Reply(_) => {
193                return Err(Status::internal("forwarder returned a mismatched reply"));
194            }
195            RouteResolution::Local => {
196                let response = handlers::cancel(
197                    &self.state,
198                    self.state.namespace_guard(),
199                    &caller,
200                    decode_cancel_request(inner),
201                )
202                .await
203                .map_err(status_from_wire_error)?;
204                return Ok(Response::new(encode_cancel_response(response)));
205            }
206        }
207    }
208
209    async fn reopen(
210        &self,
211        request: Request<generated::ReopenRequest>,
212    ) -> Result<Response<generated::ReopenResponse>, Status> {
213        let caller = self.caller(&request).await?;
214        let (metadata, _ext, inner) = request.into_parts();
215        let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
216        match self
217            .resolve_route(
218                workflow_id,
219                &metadata,
220                ForwardRequest::Reopen(inner.clone()),
221            )
222            .await
223        {
224            RouteResolution::Reject(status) => return Err(status),
225            RouteResolution::Reply(ForwardReply::Reopen(reply)) => {
226                return Ok(Response::new(reply));
227            }
228            RouteResolution::Reply(_) => {
229                return Err(Status::internal("forwarder returned a mismatched reply"));
230            }
231            RouteResolution::Local => {
232                let response = handlers::reopen(
233                    self.state.namespace_guard(),
234                    &caller,
235                    decode_reopen_request(inner),
236                )
237                .await
238                .map_err(status_from_wire_error)?;
239                return Ok(Response::new(encode_reopen_response(response)));
240            }
241        }
242    }
243
244    async fn pause(
245        &self,
246        request: Request<generated::PauseRequest>,
247    ) -> Result<Response<generated::PauseResponse>, Status> {
248        let caller = self.caller(&request).await?;
249        let (metadata, _ext, inner) = request.into_parts();
250        let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
251        match self
252            .resolve_route(workflow_id, &metadata, ForwardRequest::Pause(inner.clone()))
253            .await
254        {
255            RouteResolution::Reject(status) => return Err(status),
256            RouteResolution::Reply(ForwardReply::Pause(reply)) => {
257                return Ok(Response::new(reply));
258            }
259            RouteResolution::Reply(_) => {
260                return Err(Status::internal("forwarder returned a mismatched reply"));
261            }
262            RouteResolution::Local => {
263                let response = handlers::pause(
264                    self.state.namespace_guard(),
265                    &caller,
266                    decode_pause_request(inner),
267                )
268                .await
269                .map_err(status_from_wire_error)?;
270                return Ok(Response::new(encode_pause_response(response)));
271            }
272        }
273    }
274
275    async fn resume(
276        &self,
277        request: Request<generated::ResumeRequest>,
278    ) -> Result<Response<generated::ResumeResponse>, Status> {
279        let caller = self.caller(&request).await?;
280        let (metadata, _ext, inner) = request.into_parts();
281        let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
282        match self
283            .resolve_route(
284                workflow_id,
285                &metadata,
286                ForwardRequest::Resume(inner.clone()),
287            )
288            .await
289        {
290            RouteResolution::Reject(status) => return Err(status),
291            RouteResolution::Reply(ForwardReply::Resume(reply)) => {
292                return Ok(Response::new(reply));
293            }
294            RouteResolution::Reply(_) => {
295                return Err(Status::internal("forwarder returned a mismatched reply"));
296            }
297            RouteResolution::Local => {
298                let response = handlers::resume(
299                    self.state.namespace_guard(),
300                    &caller,
301                    decode_resume_request(inner),
302                )
303                .await
304                .map_err(status_from_wire_error)?;
305                return Ok(Response::new(encode_resume_response(response)));
306            }
307        }
308    }
309
310    async fn rename(
311        &self,
312        request: Request<generated::RenameRequest>,
313    ) -> Result<Response<generated::RenameResponse>, Status> {
314        let caller = self.caller(&request).await?;
315        let (metadata, _ext, inner) = request.into_parts();
316        let workflow_id = inner.workflow_id.clone().map(decode_workflow_id);
317        match self
318            .resolve_route(
319                workflow_id,
320                &metadata,
321                ForwardRequest::Rename(inner.clone()),
322            )
323            .await
324        {
325            RouteResolution::Reject(status) => return Err(status),
326            RouteResolution::Reply(ForwardReply::Rename(reply)) => {
327                return Ok(Response::new(reply));
328            }
329            RouteResolution::Reply(_) => {
330                return Err(Status::internal("forwarder returned a mismatched reply"));
331            }
332            RouteResolution::Local => {
333                let response = handlers::rename(
334                    self.state.namespace_guard(),
335                    &caller,
336                    decode_rename_request(inner),
337                )
338                .await
339                .map_err(status_from_wire_error)?;
340                return Ok(Response::new(encode_rename_response(response)));
341            }
342        }
343    }
344
345    async fn list_workflows(
346        &self,
347        request: Request<generated::ListWorkflowsRequest>,
348    ) -> Result<Response<generated::ListWorkflowsResponse>, Status> {
349        plain_rpcs::list_workflows(&self.state, request).await
350    }
351
352    async fn count_workflows(
353        &self,
354        request: Request<generated::CountWorkflowsRequest>,
355    ) -> Result<Response<generated::CountWorkflowsResponse>, Status> {
356        plain_rpcs::count_workflows(&self.state, request).await
357    }
358
359    async fn describe_workflow(
360        &self,
361        request: Request<generated::DescribeWorkflowRequest>,
362    ) -> Result<Response<generated::DescribeWorkflowResponse>, Status> {
363        plain_rpcs::describe_workflow(&self.state, request).await
364    }
365
366    async fn create_schedule(
367        &self,
368        request: Request<generated::CreateScheduleRequest>,
369    ) -> Result<Response<generated::CreateScheduleResponse>, Status> {
370        plain_rpcs::create_schedule(&self.state, request).await
371    }
372
373    async fn update_schedule(
374        &self,
375        request: Request<generated::UpdateScheduleRequest>,
376    ) -> Result<Response<generated::UpdateScheduleResponse>, Status> {
377        plain_rpcs::update_schedule(&self.state, request).await
378    }
379
380    async fn pause_schedule(
381        &self,
382        request: Request<generated::ScheduleIdRequest>,
383    ) -> Result<Response<generated::PauseScheduleResponse>, Status> {
384        plain_rpcs::pause_schedule(&self.state, request).await
385    }
386
387    async fn resume_schedule(
388        &self,
389        request: Request<generated::ScheduleIdRequest>,
390    ) -> Result<Response<generated::ResumeScheduleResponse>, Status> {
391        plain_rpcs::resume_schedule(&self.state, request).await
392    }
393
394    async fn delete_schedule(
395        &self,
396        request: Request<generated::ScheduleIdRequest>,
397    ) -> Result<Response<generated::DeleteScheduleResponse>, Status> {
398        plain_rpcs::delete_schedule(&self.state, request).await
399    }
400
401    async fn list_schedules(
402        &self,
403        request: Request<generated::ListSchedulesRequest>,
404    ) -> Result<Response<generated::ListSchedulesResponse>, Status> {
405        plain_rpcs::list_schedules(&self.state, request).await
406    }
407
408    async fn describe_schedule(
409        &self,
410        request: Request<generated::ScheduleIdRequest>,
411    ) -> Result<Response<generated::DescribeScheduleResponse>, Status> {
412        plain_rpcs::describe_schedule(&self.state, request).await
413    }
414
415    /// Mint (or gate) an already-authorized namespace set on THIS node — the
416    /// owner side of namespace-mint routing. See
417    /// [`WorkflowGrpcService::mint_namespace_here`] for the authorization and
418    /// policy-parity contract this delegates to.
419    async fn mint_namespace(
420        &self,
421        request: Request<generated::MintNamespaceRequest>,
422    ) -> Result<Response<generated::MintNamespaceResponse>, Status> {
423        self.mint_namespace_here(request).await
424    }
425}
426
427#[cfg(test)]
428mod tests {
429    use std::{net::SocketAddr, sync::Arc};
430
431    use aion::EngineBuilder;
432    use aion_core::{Event, EventEnvelope, Payload, WorkflowId, WorkflowStatus};
433    use aion_proto::{
434        ProtoWireError, WireError, WireErrorCode,
435        convert::{decode_core_value, encode_core_value},
436        generated::workflow_service_server::WorkflowService,
437    };
438    use aion_store::{
439        EventStore, InMemoryStore, WriteToken,
440        visibility::{VisibilityRecord, VisibilityStore},
441    };
442    use chrono::Utc;
443    use prost::Message;
444    use serde_json::json;
445    use tonic::{Code, Request};
446
447    use super::convert::{decode_envelope, encode_envelope, encode_payload};
448    use super::*;
449    use crate::{
450        NamespaceResolver,
451        config::{
452            AuthConfig, AuthoringConfig, DeployConfig, ListenConfig, MetricsConfig,
453            NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, RuntimeConfig,
454            WebSocketConfig, WorkerConfig,
455        },
456    };
457
458    const NAMESPACE: &str = "tenant-a";
459    const TOKEN: &str = "test-token";
460
461    /// Server state whose bearer validation matches the compiled auth path:
462    /// under `feature = "auth"` a real [`crate::auth::JwksCache`] is fetched
463    /// from a live fixture JWKS endpoint; otherwise the development token path
464    /// needs no cache.
465    async fn server_state(
466        resolver: NamespaceResolver,
467        runtime: RuntimeConfig,
468    ) -> Result<ServerState, Box<dyn std::error::Error>> {
469        #[cfg(feature = "auth")]
470        {
471            let url = crate::auth::test_support::serve_jwks()?;
472            let refresh = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
473            let cache = crate::auth::JwksCache::new(url, refresh).await?;
474            Ok(ServerState::from_parts_with_jwks(resolver, runtime, cache))
475        }
476        #[cfg(not(feature = "auth"))]
477        {
478            // Yield to preserve the async signature the auth arm requires.
479            tokio::task::yield_now().await;
480            Ok(ServerState::from_parts(resolver, runtime))
481        }
482    }
483
484    #[tokio::test]
485    async fn in_process_tonic_start_and_list_use_shared_handlers()
486    -> Result<(), Box<dyn std::error::Error>> {
487        let backing = Arc::new(InMemoryStore::default());
488        let store: Arc<dyn EventStore> = backing.clone();
489        let visibility_store: Arc<dyn VisibilityStore> = backing;
490        let engine = Arc::new(
491            EngineBuilder::new()
492                .store_arc(Arc::clone(&store))
493                .visibility_store_arc(Arc::clone(&visibility_store))
494                .scheduler_threads(1)
495                .build()
496                .await?,
497        );
498        store
499            .append(
500                WriteToken::recorder(),
501                &workflow_id(),
502                &[started_event()?],
503                0,
504            )
505            .await?;
506        visibility_store
507            .record_visibility(VisibilityRecord {
508                workflow_id: workflow_id(),
509                run_id: aion_core::RunId::new(uuid::Uuid::from_u128(2)),
510                workflow_type: String::from("fixture"),
511                status: WorkflowStatus::Running,
512                start_time: Utc::now(),
513                close_time: None,
514                failed_step: None,
515                failure_reason: None,
516                search_attributes: std::collections::HashMap::from([(
517                    crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
518                    aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
519                )]),
520            })
521            .await?;
522        let resolver = NamespaceResolver::from_config(
523            crate::config::NamespaceConfig {
524                mode: NamespaceMode::SharedEngine,
525            },
526            engine,
527        );
528        let state = server_state(resolver.clone(), runtime_config()).await?;
529        let service = WorkflowGrpcService::new(state);
530
531        let mut start = Request::new(generated::StartWorkflowRequest {
532            namespace: NAMESPACE.to_owned(),
533            workflow_type: "missing-workflow".to_owned(),
534            input: Some(encode_payload(proto_payload()?)),
535            routing_key: None,
536            task_queue: None,
537            display_name: None,
538        });
539        apply_metadata(start.metadata_mut())?;
540        let start_error = service.start_workflow(start).await;
541        let status = start_error
542            .err()
543            .ok_or_else(|| WireError::backend("expected error"))?;
544        assert_eq!(status.code(), Code::NotFound);
545        let detail = ProtoWireError::decode(status.details())?;
546        assert_eq!(detail.error_type.as_deref(), Some("WorkflowTypeNotFound"));
547        assert!(detail.message.contains("missing-workflow"));
548
549        let list_filter = encode_core_value(
550            NAMESPACE,
551            None,
552            &aion_store::visibility::ListWorkflowsFilter {
553                workflow_type: Some(String::from("fixture")),
554                status: Some(WorkflowStatus::Running),
555                ..aion_store::visibility::ListWorkflowsFilter::default()
556            },
557        )?;
558        let mut list = Request::new(generated::ListWorkflowsRequest {
559            namespace: NAMESPACE.to_owned(),
560            filter: Some(encode_envelope(list_filter)),
561        });
562        apply_metadata(list.metadata_mut())?;
563        let response = service.list_workflows(list).await?.into_inner();
564
565        assert_eq!(response.summaries.len(), 1);
566        let summary = response
567            .summaries
568            .into_iter()
569            .next()
570            .map(decode_envelope)
571            .map(|envelope| decode_core_value::<aion_store::visibility::WorkflowSummary>(&envelope))
572            .transpose()?
573            .ok_or_else(|| WireError::backend("summary missing"))?;
574        assert_eq!(summary.workflow_id, workflow_id());
575        // The seeded history records no namespace attribute, so durable
576        // ownership verification must reject targeted access with NotFound:
577        // a missing ownership attribute is indistinguishable from a
578        // nonexistent workflow (anti-existence-leak), and NamespaceDenied is
579        // reserved for callers without a grant for the requested namespace.
580        assert_eq!(
581            resolver
582                .verify_workflow_ownership(NAMESPACE, &workflow_id())
583                .await
584                .err()
585                .map(|error| error.to_wire_error().code),
586            Some(WireErrorCode::NotFound)
587        );
588        Ok(())
589    }
590
591    /// A gRPC service over a TERMINAL run that records its own namespace.
592    ///
593    /// Terminal is the durable posture in which no live recorder can exist, so
594    /// the engine's non-resident append path is legitimately open (a `Running`
595    /// run with no registered handle is refused instead — see the engine's
596    /// residency gate).
597    ///
598    /// `register_display_name` selects whether the engine's search-attribute
599    /// schema knows `aion.display_name`; the blank-name arm does not need it,
600    /// because the name is refused before any append is attempted.
601    async fn rename_fixture(
602        register_display_name: bool,
603    ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
604        rename_fixture_with_terminal(register_display_name, true).await
605    }
606
607    /// The same fixture, with the run's terminality selectable, so a test can
608    /// also exercise the engine's residency refusal for a durably-`Running`
609    /// run that no registry handle owns.
610    async fn rename_fixture_with_terminal(
611        register_display_name: bool,
612        terminal: bool,
613    ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
614        let backing = Arc::new(InMemoryStore::default());
615        let store: Arc<dyn EventStore> = backing.clone();
616        let visibility_store: Arc<dyn VisibilityStore> = backing;
617        let mut schema = aion_core::SearchAttributeSchema::new();
618        schema.register(
619            crate::namespace::NAMESPACE_ATTRIBUTE,
620            aion_core::SearchAttributeType::String,
621        )?;
622        if register_display_name {
623            schema.register(
624                crate::namespace::DISPLAY_NAME_ATTRIBUTE,
625                aion_core::SearchAttributeType::String,
626            )?;
627        }
628        let engine = Arc::new(
629            EngineBuilder::new()
630                .store_arc(Arc::clone(&store))
631                .visibility_store_arc(Arc::clone(&visibility_store))
632                .search_attribute_schema(schema)
633                .scheduler_threads(1)
634                .build()
635                .await?,
636        );
637        let namespace_event = |seq: u64| Event::SearchAttributesUpdated {
638            envelope: EventEnvelope {
639                seq,
640                recorded_at: Utc::now(),
641                workflow_id: workflow_id(),
642            },
643            workflow_id: workflow_id(),
644            attributes: std::collections::HashMap::from([(
645                crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
646                aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
647            )]),
648        };
649        let mut events = vec![started_event()?, namespace_event(2)];
650        if terminal {
651            events.push(Event::WorkflowCompleted {
652                envelope: EventEnvelope {
653                    seq: 3,
654                    recorded_at: Utc::now(),
655                    workflow_id: workflow_id(),
656                },
657                result: payload()?,
658            });
659        }
660        store
661            .append(WriteToken::recorder(), &workflow_id(), &events, 0)
662            .await?;
663        let resolver = NamespaceResolver::from_config(
664            crate::config::NamespaceConfig {
665                mode: NamespaceMode::SharedEngine,
666            },
667            engine,
668        );
669        let state = server_state(resolver, runtime_config()).await?;
670        Ok((WorkflowGrpcService::new(state), store))
671    }
672
673    /// #211: a rename over the gRPC service records the name durably and
674    /// returns it as recorded (trimmed), and a SECOND rename supersedes it
675    /// while history keeps both — the whole wire path, guard included.
676    #[tokio::test]
677    async fn in_process_tonic_rename_records_the_name_and_supersedes_it()
678    -> Result<(), Box<dyn std::error::Error>> {
679        let (service, store) = rename_fixture(true).await?;
680
681        let rename = |name: &str| {
682            let mut request = Request::new(generated::RenameRequest {
683                namespace: NAMESPACE.to_owned(),
684                workflow_id: Some(generated::WorkflowId {
685                    uuid: workflow_id().to_string(),
686                }),
687                // Absent run id: the rename lands on the LATEST run.
688                run_id: None,
689                display_name: name.to_owned(),
690            });
691            apply_metadata(request.metadata_mut()).map(|()| request)
692        };
693
694        // Surrounding whitespace is trimmed, and the response says so.
695        let response = service.rename(rename("  Nightly settlement  ")?).await?;
696        let response = response.into_inner();
697        assert_eq!(response.display_name, "Nightly settlement");
698        assert_eq!(
699            response.run_id.map(|id| id.uuid),
700            Some(aion_core::RunId::new(uuid::Uuid::from_u128(1)).to_string())
701        );
702
703        service
704            .rename(rename("Nightly settlement (rerun)")?)
705            .await?;
706
707        let history = store.read_history(&workflow_id()).await?;
708        assert_eq!(
709            aion_core::display_name(&history).as_deref(),
710            Some("Nightly settlement (rerun)"),
711            "the latest recorded name wins"
712        );
713        let names: Vec<_> = history
714            .iter()
715            .filter_map(|event| match event {
716                Event::SearchAttributesUpdated { attributes, .. } => attributes
717                    .get(aion_core::DISPLAY_NAME_ATTRIBUTE)
718                    .and_then(|value| match value {
719                        aion_core::SearchAttributeValue::String(name) => Some(name.clone()),
720                        _ => None,
721                    }),
722                _ => None,
723            })
724            .collect();
725        assert_eq!(
726            names,
727            vec![
728                String::from("Nightly settlement"),
729                String::from("Nightly settlement (rerun)")
730            ],
731            "history keeps every name the run has worn"
732        );
733        Ok(())
734    }
735
736    /// #211: a rename carrying a blank name is refused as tonic
737    /// `InvalidArgument` before anything is appended.
738    #[tokio::test]
739    async fn in_process_tonic_rename_refuses_a_blank_name() -> Result<(), Box<dyn std::error::Error>>
740    {
741        let (service, store) = rename_fixture(false).await?;
742        let before = store.read_history(&workflow_id()).await?.len();
743
744        let mut request = Request::new(generated::RenameRequest {
745            namespace: NAMESPACE.to_owned(),
746            workflow_id: Some(generated::WorkflowId {
747                uuid: workflow_id().to_string(),
748            }),
749            run_id: None,
750            display_name: String::from("   "),
751        });
752        apply_metadata(request.metadata_mut())?;
753        let status = service
754            .rename(request)
755            .await
756            .err()
757            .ok_or_else(|| WireError::backend("expected a blank-name refusal"))?;
758        assert_eq!(status.code(), Code::InvalidArgument);
759        assert_eq!(
760            store.read_history(&workflow_id()).await?.len(),
761            before,
762            "a refused rename appends nothing"
763        );
764        Ok(())
765    }
766
767    /// #211 single writer, over the wire: renaming a durably-`Running` run that
768    /// no registry handle owns is refused as tonic `FailedPrecondition`
769    /// carrying the typed `InvalidState` detail — never served through a second
770    /// writer appending behind the run's live recorder. This is the wire
771    /// contract `docs/operations/API.md` states for the route.
772    #[tokio::test]
773    async fn in_process_tonic_rename_of_a_non_resident_running_run_is_failed_precondition()
774    -> Result<(), Box<dyn std::error::Error>> {
775        let (service, store) = rename_fixture_with_terminal(true, false).await?;
776        let before = store.read_history(&workflow_id()).await?.len();
777
778        let mut request = Request::new(generated::RenameRequest {
779            namespace: NAMESPACE.to_owned(),
780            workflow_id: Some(generated::WorkflowId {
781                uuid: workflow_id().to_string(),
782            }),
783            run_id: None,
784            display_name: String::from("Nightly settlement"),
785        });
786        apply_metadata(request.metadata_mut())?;
787        let status = service
788            .rename(request)
789            .await
790            .err()
791            .ok_or_else(|| WireError::backend("expected a residency refusal"))?;
792
793        assert_eq!(status.code(), Code::FailedPrecondition);
794        let detail = ProtoWireError::decode(status.details())?;
795        assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
796        assert!(
797            detail.message.contains("not resident"),
798            "the refusal must name the reason: {}",
799            detail.message
800        );
801        assert_eq!(
802            store.read_history(&workflow_id()).await?.len(),
803            before,
804            "a refused rename appends nothing"
805        );
806        Ok(())
807    }
808
809    /// Reopening a terminal-Completed workflow over the gRPC service surfaces
810    /// the engine's `InvalidState` precondition as tonic `FailedPrecondition`
811    /// carrying the typed `InvalidState` detail (AO-007 C35/C38).
812    #[tokio::test]
813    async fn in_process_tonic_reopen_completed_is_failed_precondition_invalid_state()
814    -> Result<(), Box<dyn std::error::Error>> {
815        let backing = Arc::new(InMemoryStore::default());
816        let store: Arc<dyn EventStore> = backing.clone();
817        let visibility_store: Arc<dyn VisibilityStore> = backing;
818        let engine = Arc::new(
819            EngineBuilder::new()
820                .store_arc(Arc::clone(&store))
821                .visibility_store_arc(Arc::clone(&visibility_store))
822                .scheduler_threads(1)
823                .build()
824                .await?,
825        );
826        // A terminal-Completed run whose history records its namespace, so the
827        // guard's durable-ownership verification passes and the request reaches
828        // the engine reopen op (which rejects Completed with InvalidState).
829        store
830            .append(
831                WriteToken::recorder(),
832                &workflow_id(),
833                &[
834                    started_event()?,
835                    Event::SearchAttributesUpdated {
836                        envelope: EventEnvelope {
837                            seq: 2,
838                            recorded_at: Utc::now(),
839                            workflow_id: workflow_id(),
840                        },
841                        workflow_id: workflow_id(),
842                        attributes: std::collections::HashMap::from([(
843                            crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
844                            aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
845                        )]),
846                    },
847                    Event::WorkflowCompleted {
848                        envelope: EventEnvelope {
849                            seq: 3,
850                            recorded_at: Utc::now(),
851                            workflow_id: workflow_id(),
852                        },
853                        result: payload()?,
854                    },
855                ],
856                0,
857            )
858            .await?;
859        let resolver = NamespaceResolver::from_config(
860            crate::config::NamespaceConfig {
861                mode: NamespaceMode::SharedEngine,
862            },
863            engine,
864        );
865        let state = server_state(resolver, runtime_config()).await?;
866        let service = WorkflowGrpcService::new(state);
867
868        let mut reopen = Request::new(generated::ReopenRequest {
869            namespace: NAMESPACE.to_owned(),
870            workflow_id: Some(generated::WorkflowId {
871                uuid: workflow_id().to_string(),
872            }),
873            run_id: None,
874        });
875        apply_metadata(reopen.metadata_mut())?;
876        let status = service
877            .reopen(reopen)
878            .await
879            .err()
880            .ok_or_else(|| WireError::backend("expected a reopen precondition error"))?;
881        assert_eq!(status.code(), Code::FailedPrecondition);
882        let detail = ProtoWireError::decode(status.details())?;
883        assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
884        assert_eq!(
885            detail.code,
886            aion_proto::ProtoWireErrorCode::InvalidState as i32
887        );
888        Ok(())
889    }
890
891    fn apply_metadata(
892        metadata: &mut tonic::metadata::MetadataMap,
893    ) -> Result<(), Box<dyn std::error::Error>> {
894        // Bearer credential accepted by the compiled authentication path: a
895        // JWT minted against the fixture JWKS under `feature = "auth"`, the
896        // development shared-secret token otherwise.
897        #[cfg(feature = "auth")]
898        let bearer = crate::auth::test_support::mint_token("alice", NAMESPACE)?;
899        #[cfg(not(feature = "auth"))]
900        let bearer = TOKEN.to_owned();
901        metadata.insert("authorization", format!("Bearer {bearer}").parse()?);
902        metadata.insert("x-aion-subject", "alice".parse()?);
903        metadata.insert("x-aion-namespaces", NAMESPACE.parse()?);
904        Ok(())
905    }
906
907    /// Test runtime settings with authentication enabled; under
908    /// `feature = "auth"` validation runs against the [`server_state`]-injected
909    /// JWKS cache, so the configured dev-secret `jwks_url` is never fetched.
910    fn runtime_config() -> RuntimeConfig {
911        RuntimeConfig {
912            listen: ListenConfig {
913                grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
914                http: SocketAddr::from(([127, 0, 0, 1], 8080)),
915            },
916            tls: None,
917            auth: AuthConfig {
918                enabled: true,
919                jwks_url: Some(TOKEN.to_owned()),
920                jwks_refresh_seconds: 300,
921            },
922            ops_console: OpsConsoleConfig {
923                source: OpsConsoleAssetSource::Embedded,
924            },
925            namespace: NamespaceConfig {
926                mode: NamespaceMode::SharedEngine,
927            },
928            worker: WorkerConfig {
929                heartbeat_window: std::time::Duration::from_secs(30),
930                ..WorkerConfig::default()
931            },
932            websocket: WebSocketConfig {
933                outbound_buffer_bound: 32,
934                event_broadcast_capacity: Some(64),
935                cluster_broadcast_capacity: Some(64),
936            },
937            workflow_packages: Vec::new(),
938            deploy: DeployConfig::default(),
939            authoring: AuthoringConfig::default(),
940            dev: crate::config::DevConfig::default(),
941            outbox: crate::config::OutboxConfig::default(),
942            observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
943            mcp: crate::config::ResolvedMcpConfig::default(),
944            scheduler_threads: 1,
945            jit_threshold: None,
946            query_timeout: Some(std::time::Duration::from_secs(10)),
947            default_namespace: "default".to_owned(),
948            auto_create: crate::config::AutoCreate::Open,
949            max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
950            drain_timeout: std::time::Duration::from_secs(30),
951            metrics: MetricsConfig { enabled: true },
952            owned_shards: Vec::new(),
953            cors_allowed_origins: Vec::new(),
954        }
955    }
956
957    fn started_event() -> Result<Event, aion_core::PayloadError> {
958        Ok(Event::WorkflowStarted {
959            envelope: EventEnvelope {
960                seq: 1,
961                recorded_at: Utc::now(),
962                workflow_id: workflow_id(),
963            },
964            workflow_type: "fixture".to_owned(),
965            input: payload()?,
966            run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
967            parent_run_id: None,
968            parent_workflow_id: None,
969            package_version: aion_core::PackageVersion::new("a".repeat(64)),
970        })
971    }
972
973    fn proto_payload() -> Result<aion_proto::ProtoPayload, aion_core::PayloadError> {
974        Ok(payload()?.into())
975    }
976
977    fn payload() -> Result<Payload, aion_core::PayloadError> {
978        Payload::from_json(&json!({ "fixture": "input" }))
979    }
980
981    fn workflow_id() -> WorkflowId {
982        WorkflowId::new(uuid::Uuid::from_u128(1))
983    }
984}