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