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 read_history(
367        &self,
368        request: Request<generated::ReadHistoryRequest>,
369    ) -> Result<Response<generated::ReadHistoryResponse>, Status> {
370        plain_rpcs::read_history(&self.state, request).await
371    }
372
373    async fn create_schedule(
374        &self,
375        request: Request<generated::CreateScheduleRequest>,
376    ) -> Result<Response<generated::CreateScheduleResponse>, Status> {
377        plain_rpcs::create_schedule(&self.state, request).await
378    }
379
380    async fn update_schedule(
381        &self,
382        request: Request<generated::UpdateScheduleRequest>,
383    ) -> Result<Response<generated::UpdateScheduleResponse>, Status> {
384        plain_rpcs::update_schedule(&self.state, request).await
385    }
386
387    async fn pause_schedule(
388        &self,
389        request: Request<generated::ScheduleIdRequest>,
390    ) -> Result<Response<generated::PauseScheduleResponse>, Status> {
391        plain_rpcs::pause_schedule(&self.state, request).await
392    }
393
394    async fn resume_schedule(
395        &self,
396        request: Request<generated::ScheduleIdRequest>,
397    ) -> Result<Response<generated::ResumeScheduleResponse>, Status> {
398        plain_rpcs::resume_schedule(&self.state, request).await
399    }
400
401    async fn delete_schedule(
402        &self,
403        request: Request<generated::ScheduleIdRequest>,
404    ) -> Result<Response<generated::DeleteScheduleResponse>, Status> {
405        plain_rpcs::delete_schedule(&self.state, request).await
406    }
407
408    async fn list_schedules(
409        &self,
410        request: Request<generated::ListSchedulesRequest>,
411    ) -> Result<Response<generated::ListSchedulesResponse>, Status> {
412        plain_rpcs::list_schedules(&self.state, request).await
413    }
414
415    async fn describe_schedule(
416        &self,
417        request: Request<generated::ScheduleIdRequest>,
418    ) -> Result<Response<generated::DescribeScheduleResponse>, Status> {
419        plain_rpcs::describe_schedule(&self.state, request).await
420    }
421
422    /// Mint (or gate) an already-authorized namespace set on THIS node — the
423    /// owner side of namespace-mint routing. See
424    /// [`WorkflowGrpcService::mint_namespace_here`] for the authorization and
425    /// policy-parity contract this delegates to.
426    async fn mint_namespace(
427        &self,
428        request: Request<generated::MintNamespaceRequest>,
429    ) -> Result<Response<generated::MintNamespaceResponse>, Status> {
430        self.mint_namespace_here(request).await
431    }
432}
433
434#[cfg(test)]
435mod tests {
436    use std::{net::SocketAddr, sync::Arc};
437
438    use aion::EngineBuilder;
439    use aion_core::{Event, EventEnvelope, Payload, WorkflowId, WorkflowStatus};
440    use aion_proto::{
441        ProtoWireError, WireError, WireErrorCode,
442        convert::{decode_core_value, encode_core_value},
443        generated::workflow_service_server::WorkflowService,
444    };
445    use aion_store::{
446        EventStore, InMemoryStore, WriteToken,
447        visibility::{VisibilityRecord, VisibilityStore},
448    };
449    use chrono::Utc;
450    use prost::Message;
451    use serde_json::json;
452    use tonic::{Code, Request};
453
454    use super::convert::{decode_envelope, encode_envelope, encode_payload};
455    use super::*;
456    use crate::{
457        NamespaceResolver,
458        config::{
459            AuthConfig, AuthoringConfig, DeployConfig, ListenConfig, MetricsConfig,
460            NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, RuntimeConfig,
461            WebSocketConfig, WorkerConfig,
462        },
463    };
464
465    const NAMESPACE: &str = "tenant-a";
466    const TOKEN: &str = "test-token";
467
468    /// Server state whose bearer validation matches the compiled auth path:
469    /// under `feature = "auth"` a real [`crate::auth::JwksCache`] is fetched
470    /// from a live fixture JWKS endpoint; otherwise the development token path
471    /// needs no cache.
472    async fn server_state(
473        resolver: NamespaceResolver,
474        runtime: RuntimeConfig,
475    ) -> Result<ServerState, Box<dyn std::error::Error>> {
476        #[cfg(feature = "auth")]
477        {
478            let url = crate::auth::test_support::serve_jwks()?;
479            let refresh = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
480            let cache = crate::auth::JwksCache::new(url, refresh).await?;
481            Ok(ServerState::from_parts_with_jwks(resolver, runtime, cache))
482        }
483        #[cfg(not(feature = "auth"))]
484        {
485            // Yield to preserve the async signature the auth arm requires.
486            tokio::task::yield_now().await;
487            Ok(ServerState::from_parts(resolver, runtime))
488        }
489    }
490
491    #[tokio::test]
492    async fn in_process_tonic_start_and_list_use_shared_handlers()
493    -> Result<(), Box<dyn std::error::Error>> {
494        let backing = Arc::new(InMemoryStore::default());
495        let store: Arc<dyn EventStore> = backing.clone();
496        let visibility_store: Arc<dyn VisibilityStore> = backing;
497        let engine = Arc::new(
498            EngineBuilder::new()
499                .store_arc(Arc::clone(&store))
500                .visibility_store_arc(Arc::clone(&visibility_store))
501                .scheduler_threads(1)
502                .build()
503                .await?,
504        );
505        store
506            .append(
507                WriteToken::recorder(),
508                &workflow_id(),
509                &[started_event()?],
510                0,
511            )
512            .await?;
513        visibility_store
514            .record_visibility(VisibilityRecord {
515                workflow_id: workflow_id(),
516                run_id: aion_core::RunId::new(uuid::Uuid::from_u128(2)),
517                workflow_type: String::from("fixture"),
518                status: WorkflowStatus::Running,
519                start_time: Utc::now(),
520                close_time: None,
521                failed_step: None,
522                failure_reason: None,
523                search_attributes: std::collections::HashMap::from([(
524                    crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
525                    aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
526                )]),
527            })
528            .await?;
529        let resolver = NamespaceResolver::from_config(
530            crate::config::NamespaceConfig {
531                mode: NamespaceMode::SharedEngine,
532            },
533            engine,
534        );
535        let state = server_state(resolver.clone(), runtime_config()).await?;
536        let service = WorkflowGrpcService::new(state);
537
538        let mut start = Request::new(generated::StartWorkflowRequest {
539            namespace: NAMESPACE.to_owned(),
540            workflow_type: "missing-workflow".to_owned(),
541            input: Some(encode_payload(proto_payload()?)),
542            routing_key: None,
543            task_queue: None,
544            display_name: None,
545        });
546        apply_metadata(start.metadata_mut())?;
547        let start_error = service.start_workflow(start).await;
548        let status = start_error
549            .err()
550            .ok_or_else(|| WireError::backend("expected error"))?;
551        assert_eq!(status.code(), Code::NotFound);
552        let detail = ProtoWireError::decode(status.details())?;
553        assert_eq!(detail.error_type.as_deref(), Some("WorkflowTypeNotFound"));
554        assert!(detail.message.contains("missing-workflow"));
555
556        let list_filter = encode_core_value(
557            NAMESPACE,
558            None,
559            &aion_store::visibility::ListWorkflowsFilter {
560                workflow_type: Some(String::from("fixture")),
561                status: Some(WorkflowStatus::Running),
562                ..aion_store::visibility::ListWorkflowsFilter::default()
563            },
564        )?;
565        let mut list = Request::new(generated::ListWorkflowsRequest {
566            namespace: NAMESPACE.to_owned(),
567            filter: Some(encode_envelope(list_filter)),
568        });
569        apply_metadata(list.metadata_mut())?;
570        let response = service.list_workflows(list).await?.into_inner();
571
572        assert_eq!(response.summaries.len(), 1);
573        let summary = response
574            .summaries
575            .into_iter()
576            .next()
577            .map(decode_envelope)
578            .map(|envelope| decode_core_value::<aion_store::visibility::WorkflowSummary>(&envelope))
579            .transpose()?
580            .ok_or_else(|| WireError::backend("summary missing"))?;
581        assert_eq!(summary.workflow_id, workflow_id());
582        // The seeded history records no namespace attribute, so durable
583        // ownership verification must reject targeted access with NotFound:
584        // a missing ownership attribute is indistinguishable from a
585        // nonexistent workflow (anti-existence-leak), and NamespaceDenied is
586        // reserved for callers without a grant for the requested namespace.
587        assert_eq!(
588            resolver
589                .verify_workflow_ownership(NAMESPACE, &workflow_id())
590                .await
591                .err()
592                .map(|error| error.to_wire_error().code),
593            Some(WireErrorCode::NotFound)
594        );
595        Ok(())
596    }
597
598    /// A gRPC service over a TERMINAL run that records its own namespace.
599    ///
600    /// Terminal is the durable posture in which no live recorder can exist, so
601    /// the engine's non-resident append path is legitimately open (a `Running`
602    /// run with no registered handle is refused instead — see the engine's
603    /// residency gate).
604    ///
605    /// `register_display_name` selects whether the engine's search-attribute
606    /// schema knows `aion.display_name`; the blank-name arm does not need it,
607    /// because the name is refused before any append is attempted.
608    async fn rename_fixture(
609        register_display_name: bool,
610    ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
611        rename_fixture_with_terminal(register_display_name, true).await
612    }
613
614    /// The same fixture, with the run's terminality selectable, so a test can
615    /// also exercise the engine's residency refusal for a durably-`Running`
616    /// run that no registry handle owns.
617    async fn rename_fixture_with_terminal(
618        register_display_name: bool,
619        terminal: bool,
620    ) -> Result<(WorkflowGrpcService, Arc<dyn EventStore>), Box<dyn std::error::Error>> {
621        let backing = Arc::new(InMemoryStore::default());
622        let store: Arc<dyn EventStore> = backing.clone();
623        let visibility_store: Arc<dyn VisibilityStore> = backing;
624        let mut schema = aion_core::SearchAttributeSchema::new();
625        schema.register(
626            crate::namespace::NAMESPACE_ATTRIBUTE,
627            aion_core::SearchAttributeType::String,
628        )?;
629        if register_display_name {
630            schema.register(
631                crate::namespace::DISPLAY_NAME_ATTRIBUTE,
632                aion_core::SearchAttributeType::String,
633            )?;
634        }
635        let engine = Arc::new(
636            EngineBuilder::new()
637                .store_arc(Arc::clone(&store))
638                .visibility_store_arc(Arc::clone(&visibility_store))
639                .search_attribute_schema(schema)
640                .scheduler_threads(1)
641                .build()
642                .await?,
643        );
644        let namespace_event = |seq: u64| Event::SearchAttributesUpdated {
645            envelope: EventEnvelope {
646                seq,
647                recorded_at: Utc::now(),
648                workflow_id: workflow_id(),
649            },
650            workflow_id: workflow_id(),
651            attributes: std::collections::HashMap::from([(
652                crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
653                aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
654            )]),
655        };
656        let mut events = vec![started_event()?, namespace_event(2)];
657        if terminal {
658            events.push(Event::WorkflowCompleted {
659                envelope: EventEnvelope {
660                    seq: 3,
661                    recorded_at: Utc::now(),
662                    workflow_id: workflow_id(),
663                },
664                result: payload()?,
665            });
666        }
667        store
668            .append(WriteToken::recorder(), &workflow_id(), &events, 0)
669            .await?;
670        let resolver = NamespaceResolver::from_config(
671            crate::config::NamespaceConfig {
672                mode: NamespaceMode::SharedEngine,
673            },
674            engine,
675        );
676        let state = server_state(resolver, runtime_config()).await?;
677        Ok((WorkflowGrpcService::new(state), store))
678    }
679
680    /// #211: a rename over the gRPC service records the name durably and
681    /// returns it as recorded (trimmed), and a SECOND rename supersedes it
682    /// while history keeps both — the whole wire path, guard included.
683    #[tokio::test]
684    async fn in_process_tonic_rename_records_the_name_and_supersedes_it()
685    -> Result<(), Box<dyn std::error::Error>> {
686        let (service, store) = rename_fixture(true).await?;
687
688        let rename = |name: &str| {
689            let mut request = Request::new(generated::RenameRequest {
690                namespace: NAMESPACE.to_owned(),
691                workflow_id: Some(generated::WorkflowId {
692                    uuid: workflow_id().to_string(),
693                }),
694                // Absent run id: the rename lands on the LATEST run.
695                run_id: None,
696                display_name: name.to_owned(),
697            });
698            apply_metadata(request.metadata_mut()).map(|()| request)
699        };
700
701        // Surrounding whitespace is trimmed, and the response says so.
702        let response = service.rename(rename("  Nightly settlement  ")?).await?;
703        let response = response.into_inner();
704        assert_eq!(response.display_name, "Nightly settlement");
705        assert_eq!(
706            response.run_id.map(|id| id.uuid),
707            Some(aion_core::RunId::new(uuid::Uuid::from_u128(1)).to_string())
708        );
709
710        service
711            .rename(rename("Nightly settlement (rerun)")?)
712            .await?;
713
714        let history = store.read_history(&workflow_id()).await?;
715        assert_eq!(
716            aion_core::display_name(&history).as_deref(),
717            Some("Nightly settlement (rerun)"),
718            "the latest recorded name wins"
719        );
720        let names: Vec<_> = history
721            .iter()
722            .filter_map(|event| match event {
723                Event::SearchAttributesUpdated { attributes, .. } => attributes
724                    .get(aion_core::DISPLAY_NAME_ATTRIBUTE)
725                    .and_then(|value| match value {
726                        aion_core::SearchAttributeValue::String(name) => Some(name.clone()),
727                        _ => None,
728                    }),
729                _ => None,
730            })
731            .collect();
732        assert_eq!(
733            names,
734            vec![
735                String::from("Nightly settlement"),
736                String::from("Nightly settlement (rerun)")
737            ],
738            "history keeps every name the run has worn"
739        );
740        Ok(())
741    }
742
743    /// #211: a rename carrying a blank name is refused as tonic
744    /// `InvalidArgument` before anything is appended.
745    #[tokio::test]
746    async fn in_process_tonic_rename_refuses_a_blank_name() -> Result<(), Box<dyn std::error::Error>>
747    {
748        let (service, store) = rename_fixture(false).await?;
749        let before = store.read_history(&workflow_id()).await?.len();
750
751        let mut request = Request::new(generated::RenameRequest {
752            namespace: NAMESPACE.to_owned(),
753            workflow_id: Some(generated::WorkflowId {
754                uuid: workflow_id().to_string(),
755            }),
756            run_id: None,
757            display_name: String::from("   "),
758        });
759        apply_metadata(request.metadata_mut())?;
760        let status = service
761            .rename(request)
762            .await
763            .err()
764            .ok_or_else(|| WireError::backend("expected a blank-name refusal"))?;
765        assert_eq!(status.code(), Code::InvalidArgument);
766        assert_eq!(
767            store.read_history(&workflow_id()).await?.len(),
768            before,
769            "a refused rename appends nothing"
770        );
771        Ok(())
772    }
773
774    /// #211 single writer, over the wire: renaming a durably-`Running` run that
775    /// no registry handle owns is refused as tonic `FailedPrecondition`
776    /// carrying the typed `InvalidState` detail — never served through a second
777    /// writer appending behind the run's live recorder. This is the wire
778    /// contract `docs/operations/API.md` states for the route.
779    #[tokio::test]
780    async fn in_process_tonic_rename_of_a_non_resident_running_run_is_failed_precondition()
781    -> Result<(), Box<dyn std::error::Error>> {
782        let (service, store) = rename_fixture_with_terminal(true, false).await?;
783        let before = store.read_history(&workflow_id()).await?.len();
784
785        let mut request = Request::new(generated::RenameRequest {
786            namespace: NAMESPACE.to_owned(),
787            workflow_id: Some(generated::WorkflowId {
788                uuid: workflow_id().to_string(),
789            }),
790            run_id: None,
791            display_name: String::from("Nightly settlement"),
792        });
793        apply_metadata(request.metadata_mut())?;
794        let status = service
795            .rename(request)
796            .await
797            .err()
798            .ok_or_else(|| WireError::backend("expected a residency refusal"))?;
799
800        assert_eq!(status.code(), Code::FailedPrecondition);
801        let detail = ProtoWireError::decode(status.details())?;
802        assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
803        assert!(
804            detail.message.contains("not resident"),
805            "the refusal must name the reason: {}",
806            detail.message
807        );
808        assert_eq!(
809            store.read_history(&workflow_id()).await?.len(),
810            before,
811            "a refused rename appends nothing"
812        );
813        Ok(())
814    }
815
816    /// Reopening a terminal-Completed workflow over the gRPC service surfaces
817    /// the engine's `InvalidState` precondition as tonic `FailedPrecondition`
818    /// carrying the typed `InvalidState` detail (AO-007 C35/C38).
819    #[tokio::test]
820    async fn in_process_tonic_reopen_completed_is_failed_precondition_invalid_state()
821    -> Result<(), Box<dyn std::error::Error>> {
822        let backing = Arc::new(InMemoryStore::default());
823        let store: Arc<dyn EventStore> = backing.clone();
824        let visibility_store: Arc<dyn VisibilityStore> = backing;
825        let engine = Arc::new(
826            EngineBuilder::new()
827                .store_arc(Arc::clone(&store))
828                .visibility_store_arc(Arc::clone(&visibility_store))
829                .scheduler_threads(1)
830                .build()
831                .await?,
832        );
833        // A terminal-Completed run whose history records its namespace, so the
834        // guard's durable-ownership verification passes and the request reaches
835        // the engine reopen op (which rejects Completed with InvalidState).
836        store
837            .append(
838                WriteToken::recorder(),
839                &workflow_id(),
840                &[
841                    started_event()?,
842                    Event::SearchAttributesUpdated {
843                        envelope: EventEnvelope {
844                            seq: 2,
845                            recorded_at: Utc::now(),
846                            workflow_id: workflow_id(),
847                        },
848                        workflow_id: workflow_id(),
849                        attributes: std::collections::HashMap::from([(
850                            crate::namespace::NAMESPACE_ATTRIBUTE.to_owned(),
851                            aion_core::SearchAttributeValue::String(NAMESPACE.to_owned()),
852                        )]),
853                    },
854                    Event::WorkflowCompleted {
855                        envelope: EventEnvelope {
856                            seq: 3,
857                            recorded_at: Utc::now(),
858                            workflow_id: workflow_id(),
859                        },
860                        result: payload()?,
861                    },
862                ],
863                0,
864            )
865            .await?;
866        let resolver = NamespaceResolver::from_config(
867            crate::config::NamespaceConfig {
868                mode: NamespaceMode::SharedEngine,
869            },
870            engine,
871        );
872        let state = server_state(resolver, runtime_config()).await?;
873        let service = WorkflowGrpcService::new(state);
874
875        let mut reopen = Request::new(generated::ReopenRequest {
876            namespace: NAMESPACE.to_owned(),
877            workflow_id: Some(generated::WorkflowId {
878                uuid: workflow_id().to_string(),
879            }),
880            run_id: None,
881        });
882        apply_metadata(reopen.metadata_mut())?;
883        let status = service
884            .reopen(reopen)
885            .await
886            .err()
887            .ok_or_else(|| WireError::backend("expected a reopen precondition error"))?;
888        assert_eq!(status.code(), Code::FailedPrecondition);
889        let detail = ProtoWireError::decode(status.details())?;
890        assert_eq!(detail.error_type.as_deref(), Some("InvalidState"));
891        assert_eq!(
892            detail.code,
893            aion_proto::ProtoWireErrorCode::InvalidState as i32
894        );
895        Ok(())
896    }
897
898    fn apply_metadata(
899        metadata: &mut tonic::metadata::MetadataMap,
900    ) -> Result<(), Box<dyn std::error::Error>> {
901        // Bearer credential accepted by the compiled authentication path: a
902        // JWT minted against the fixture JWKS under `feature = "auth"`, the
903        // development shared-secret token otherwise.
904        #[cfg(feature = "auth")]
905        let bearer = crate::auth::test_support::mint_token("alice", NAMESPACE)?;
906        #[cfg(not(feature = "auth"))]
907        let bearer = TOKEN.to_owned();
908        metadata.insert("authorization", format!("Bearer {bearer}").parse()?);
909        metadata.insert("x-aion-subject", "alice".parse()?);
910        metadata.insert("x-aion-namespaces", NAMESPACE.parse()?);
911        Ok(())
912    }
913
914    /// Test runtime settings with authentication enabled; under
915    /// `feature = "auth"` validation runs against the [`server_state`]-injected
916    /// JWKS cache, so the configured dev-secret `jwks_url` is never fetched.
917    fn runtime_config() -> RuntimeConfig {
918        RuntimeConfig {
919            listen: ListenConfig {
920                grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
921                http: SocketAddr::from(([127, 0, 0, 1], 8080)),
922            },
923            tls: None,
924            auth: AuthConfig {
925                enabled: true,
926                jwks_url: Some(TOKEN.to_owned()),
927                jwks_refresh_seconds: 300,
928            },
929            ops_console: OpsConsoleConfig {
930                source: OpsConsoleAssetSource::Embedded,
931            },
932            namespace: NamespaceConfig {
933                mode: NamespaceMode::SharedEngine,
934            },
935            worker: WorkerConfig {
936                heartbeat_window: std::time::Duration::from_secs(30),
937                ..WorkerConfig::default()
938            },
939            websocket: WebSocketConfig {
940                outbound_buffer_bound: 32,
941                event_broadcast_capacity: Some(64),
942                cluster_broadcast_capacity: Some(64),
943            },
944            workflow_packages: Vec::new(),
945            deploy: DeployConfig::default(),
946            authoring: AuthoringConfig::default(),
947            dev: crate::config::DevConfig::default(),
948            outbox: crate::config::OutboxConfig::default(),
949            observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
950            mcp: crate::config::ResolvedMcpConfig::default(),
951            scheduler_threads: 1,
952            jit_threshold: None,
953            query_timeout: Some(std::time::Duration::from_secs(10)),
954            default_namespace: "default".to_owned(),
955            auto_create: crate::config::AutoCreate::Open,
956            max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
957            drain_timeout: std::time::Duration::from_secs(30),
958            metrics: MetricsConfig { enabled: true },
959            owned_shards: Vec::new(),
960            cors_allowed_origins: Vec::new(),
961        }
962    }
963
964    fn started_event() -> Result<Event, aion_core::PayloadError> {
965        Ok(Event::WorkflowStarted {
966            envelope: EventEnvelope {
967                seq: 1,
968                recorded_at: Utc::now(),
969                workflow_id: workflow_id(),
970            },
971            workflow_type: "fixture".to_owned(),
972            input: payload()?,
973            run_id: aion_core::RunId::new(uuid::Uuid::from_u128(1)),
974            parent_run_id: None,
975            parent_workflow_id: None,
976            package_version: aion_core::PackageVersion::new("a".repeat(64)),
977        })
978    }
979
980    fn proto_payload() -> Result<aion_proto::ProtoPayload, aion_core::PayloadError> {
981        Ok(payload()?.into())
982    }
983
984    fn payload() -> Result<Payload, aion_core::PayloadError> {
985        Payload::from_json(&json!({ "fixture": "input" }))
986    }
987
988    fn workflow_id() -> WorkflowId {
989        WorkflowId::new(uuid::Uuid::from_u128(1))
990    }
991}