1#![forbid(unsafe_code)]
24#![warn(clippy::pedantic)]
25#![deny(warnings)]
26#![warn(missing_docs)]
27#![allow(clippy::module_name_repetitions)]
28#![allow(clippy::missing_errors_doc)]
29
30use std::future::Future;
31use std::pin::Pin;
32use std::sync::Arc;
33
34use themql_core::{
35 Context, Error, Message, MessageHandler, Query, Resource, ResponseValue, Subject,
36};
37use thiserror::Error;
38
39use async_graphql::{Guard, Object, Subscription};
40
41use futures_util::stream::Stream;
42
43pub trait GraphqlResolverBridge: Send + Sync {
54 fn resolve_field<'a>(
60 &'a self,
61 field: &'a str,
62 args: &'a serde_json::Value,
63 ctx: &'a Context,
64 ) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, Error>> + Send + 'a>>;
65}
66
67#[derive(Clone)]
74pub struct GraphqlResolverBridgeImpl {
75 resolver: Arc<dyn themql_core::Resolver>,
76}
77
78impl GraphqlResolverBridgeImpl {
79 #[must_use]
81 pub fn new(resolver: Arc<dyn themql_core::Resolver>) -> Self {
82 Self { resolver }
83 }
84}
85
86impl GraphqlResolverBridge for GraphqlResolverBridgeImpl {
87 fn resolve_field<'a>(
88 &'a self,
89 field: &'a str,
90 args: &'a serde_json::Value,
91 ctx: &'a Context,
92 ) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, Error>> + Send + 'a>> {
93 Box::pin(async move {
94 let resource =
95 Resource::from_str(field).map_err(|e| Error::validation_error(e.to_string()))?;
96 let mut query = Query::new(resource);
97 query.context = ctx.clone();
98 if !args.is_null() {
99 query = query.with_arguments(args.clone());
100 }
101 let response = self.resolver.resolve(&query, ctx).await?;
102 Ok(response_value_to_json(response.result?))
103 })
104 }
105}
106
107fn response_value_to_json(value: ResponseValue) -> serde_json::Value {
109 match value {
110 ResponseValue::Json(v) => v,
111 ResponseValue::Bytes(bytes, tag) => {
112 let obj = serde_json::json!({
113 "format": tag.to_string(),
114 "bytes": bytes,
115 });
116 obj
117 }
118 ResponseValue::Unit => serde_json::Value::Null,
119 }
120}
121
122pub trait DispatchBridge: Send + Sync {
130 fn dispatch<'a>(
135 &'a self,
136 subject: &'a str,
137 payload: &'a serde_json::Value,
138 ) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, Error>> + Send + 'a>>;
139}
140
141#[derive(Clone)]
143pub struct DispatchBridgeImpl {
144 handler: Arc<dyn MessageHandler>,
145}
146
147impl DispatchBridgeImpl {
148 #[must_use]
150 pub fn new(handler: Arc<dyn MessageHandler>) -> Self {
151 Self { handler }
152 }
153}
154
155impl DispatchBridge for DispatchBridgeImpl {
156 fn dispatch<'a>(
157 &'a self,
158 subject: &'a str,
159 payload: &'a serde_json::Value,
160 ) -> Pin<Box<dyn Future<Output = Result<serde_json::Value, Error>> + Send + 'a>> {
161 Box::pin(async move {
162 use themql_core::{Operation, Payload};
163 let mut msg = Message::new(subject, Operation::Command)
164 .map_err(|e| Error::validation_error(e.to_string()))?;
165 msg.payload = Payload::Json(payload.clone());
166 let response = self.handler.handle(&msg).await?;
167 Ok(response_value_to_json(response.result?))
168 })
169 }
170}
171
172pub type JsonValueStream = Box<dyn Stream<Item = serde_json::Value> + Send + Unpin>;
178
179pub trait GraphqlSubscriptionSource: Send + Sync {
187 fn subscribe_stream<'a>(
193 &'a self,
194 subject: &'a Subject,
195 ) -> Pin<Box<dyn Future<Output = Result<JsonValueStream, Error>> + Send + 'a>>;
196}
197
198struct SseStreamAdapter {
205 rx: tokio::sync::mpsc::Receiver<serde_json::Value>,
206}
207
208impl SseStreamAdapter {
209 fn new(mut stream: Box<dyn themql_sse::SseStream>) -> Self {
210 let (tx, rx) = tokio::sync::mpsc::channel::<serde_json::Value>(64);
211 tokio::spawn(async move {
212 while let Ok(Some(event)) = stream.next_event().await {
213 let value = event_to_json(&event);
214 if tx.send(value).await.is_err() {
215 break;
216 }
217 }
218 });
219 Self { rx }
220 }
221}
222
223fn event_to_json(event: &themql_sse::SseEvent) -> serde_json::Value {
224 let data = serde_json::from_str(&event.data)
225 .unwrap_or_else(|_| serde_json::Value::String(event.data.clone()));
226 serde_json::json!({
227 "id": event.id,
228 "event": event.event,
229 "data": data,
230 "retry": event.retry,
231 })
232}
233
234impl Stream for SseStreamAdapter {
235 type Item = serde_json::Value;
236
237 fn poll_next(
238 mut self: Pin<&mut Self>,
239 cx: &mut std::task::Context<'_>,
240 ) -> std::task::Poll<Option<Self::Item>> {
241 self.rx.poll_recv(cx)
242 }
243}
244
245impl GraphqlSubscriptionSource for themql_sse::TokioSsePublisher {
246 fn subscribe_stream<'a>(
247 &'a self,
248 subject: &'a Subject,
249 ) -> Pin<Box<dyn Future<Output = Result<JsonValueStream, Error>> + Send + 'a>> {
250 Box::pin(async move {
251 use themql_sse::SsePublisher;
252 let stream = SsePublisher::add_subscriber(self, subject)
253 .await
254 .map_err(|e| Error::internal_error(e.to_string()))?;
255 let adapter = SseStreamAdapter::new(stream);
256 Ok(Box::new(adapter) as JsonValueStream)
257 })
258 }
259}
260
261#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
267pub enum AuthRole {
268 Admin,
270 Operator,
272 Observer,
274}
275
276impl AuthRole {
277 #[must_use]
280 pub fn satisfies(self, required: AuthRole) -> bool {
281 if self == AuthRole::Admin {
282 return true;
283 }
284 if matches!(
285 (self, required),
286 (AuthRole::Operator, AuthRole::Operator | AuthRole::Observer)
287 | (AuthRole::Observer, AuthRole::Observer)
288 ) {
289 return true;
290 }
291 false
292 }
293}
294
295impl std::str::FromStr for AuthRole {
296 type Err = ();
297
298 fn from_str(s: &str) -> Result<Self, Self::Err> {
302 match s.to_ascii_lowercase().as_str() {
303 "admin" => Ok(AuthRole::Admin),
304 "operator" => Ok(AuthRole::Operator),
305 "observer" => Ok(AuthRole::Observer),
306 _ => Err(()),
307 }
308 }
309}
310
311#[derive(Debug, Clone, Copy)]
314pub struct RoleGuard {
315 pub required: AuthRole,
317}
318
319impl RoleGuard {
320 #[must_use]
322 pub fn new(required: AuthRole) -> Self {
323 Self { required }
324 }
325}
326
327#[allow(clippy::unused_async_trait_impl)]
328impl Guard for RoleGuard {
329 async fn check(&self, ctx: &async_graphql::Context<'_>) -> async_graphql::Result<()> {
330 match ctx.data::<AuthRole>() {
331 Ok(role) => {
332 if role.satisfies(self.required) {
333 Ok(())
334 } else {
335 Err(async_graphql::Error::new(format!(
336 "insufficient role: requires {:?}, have {:?}",
337 self.required, role
338 )))
339 }
340 }
341 Err(_) => Err(async_graphql::Error::new(
342 "no auth role in context (auth not enabled?)",
343 )),
344 }
345 }
346}
347
348#[derive(Clone)]
355pub struct QueryRoot {
356 bridge: Arc<dyn GraphqlResolverBridge>,
357}
358
359impl QueryRoot {
360 #[must_use]
362 pub fn new(bridge: Arc<dyn GraphqlResolverBridge>) -> Self {
363 Self { bridge }
364 }
365}
366
367#[Object]
368impl QueryRoot {
369 #[graphql(guard = "RoleGuard::new(AuthRole::Observer)")]
373 async fn resource(
374 &self,
375 subject: String,
376 selection: Option<async_graphql::Json<serde_json::Value>>,
377 args: Option<async_graphql::Json<serde_json::Value>>,
378 ) -> Result<async_graphql::Json<serde_json::Value>, async_graphql::Error> {
379 let _ = selection;
380 let args_val = args.map_or(serde_json::Value::Null, |v| v.0);
381 let ctx = Context::new();
382 let value = self
383 .bridge
384 .resolve_field(&subject, &args_val, &ctx)
385 .await
386 .map_err(|e| async_graphql::Error::new(e.message))?;
387 Ok(async_graphql::Json(value))
388 }
389
390 #[graphql(guard = "RoleGuard::new(AuthRole::Observer)")]
395 async fn resources(
396 &self,
397 pattern: String,
398 ) -> Result<Vec<async_graphql::Json<serde_json::Value>>, async_graphql::Error> {
399 let ctx = Context::new();
400 let args_val = serde_json::json!({ "pattern": pattern });
401 let value = self
402 .bridge
403 .resolve_field("resources", &args_val, &ctx)
404 .await
405 .map_err(|e| async_graphql::Error::new(e.message))?;
406 match value {
407 serde_json::Value::Array(items) => {
408 Ok(items.into_iter().map(async_graphql::Json).collect())
409 }
410 other => Ok(vec![async_graphql::Json(other)]),
411 }
412 }
413}
414
415#[derive(Clone)]
418pub struct MutationRoot {
419 bridge: Arc<dyn DispatchBridge>,
420}
421
422impl MutationRoot {
423 #[must_use]
425 pub fn new(bridge: Arc<dyn DispatchBridge>) -> Self {
426 Self { bridge }
427 }
428}
429
430#[Object]
431impl MutationRoot {
432 #[graphql(guard = "RoleGuard::new(AuthRole::Operator)")]
435 async fn dispatch(
436 &self,
437 subject: String,
438 payload: async_graphql::Json<serde_json::Value>,
439 ) -> Result<async_graphql::Json<serde_json::Value>, async_graphql::Error> {
440 let value = self
441 .bridge
442 .dispatch(&subject, &payload.0)
443 .await
444 .map_err(|e| async_graphql::Error::new(e.message))?;
445 Ok(async_graphql::Json(value))
446 }
447}
448
449#[derive(Clone, Default)]
457pub struct SubscriptionRoot {
458 source: Option<Arc<dyn GraphqlSubscriptionSource>>,
459}
460
461impl SubscriptionRoot {
462 #[must_use]
465 pub fn with_source(source: Arc<dyn GraphqlSubscriptionSource>) -> Self {
466 Self {
467 source: Some(source),
468 }
469 }
470}
471
472#[Subscription]
473impl SubscriptionRoot {
474 #[graphql(guard = "RoleGuard::new(AuthRole::Observer)")]
483 async fn subscribe(
484 &self,
485 subject: String,
486 ) -> Result<impl Stream<Item = async_graphql::Json<serde_json::Value>>, async_graphql::Error>
487 {
488 use futures_util::StreamExt;
489 let source = self
490 .source
491 .as_ref()
492 .ok_or_else(|| async_graphql::Error::new("no subscription source configured"))?;
493 let subj =
494 Subject::from_str(&subject).map_err(|e| async_graphql::Error::new(e.to_string()))?;
495 let stream = source
496 .subscribe_stream(&subj)
497 .await
498 .map_err(|e| async_graphql::Error::new(e.message))?;
499 Ok(stream.map(async_graphql::Json))
500 }
501}
502
503pub trait GraphqlSchema: Send + Sync {
511 fn schema(&self) -> &async_graphql::Schema<QueryRoot, MutationRoot, SubscriptionRoot>;
513}
514
515pub struct GraphqlSchemaImpl {
518 schema: async_graphql::Schema<QueryRoot, MutationRoot, SubscriptionRoot>,
519}
520
521impl GraphqlSchemaImpl {
522 #[must_use]
527 pub fn new(query: QueryRoot, mutation: MutationRoot, subscription: SubscriptionRoot) -> Self {
528 let schema = async_graphql::Schema::build(query, mutation, subscription).finish();
529 Self { schema }
530 }
531
532 #[must_use]
541 pub fn with_role(
542 query: QueryRoot,
543 mutation: MutationRoot,
544 subscription: SubscriptionRoot,
545 role: AuthRole,
546 ) -> Self {
547 let schema = async_graphql::Schema::build(query, mutation, subscription)
548 .data(role)
549 .finish();
550 Self { schema }
551 }
552}
553
554impl GraphqlSchema for GraphqlSchemaImpl {
555 fn schema(&self) -> &async_graphql::Schema<QueryRoot, MutationRoot, SubscriptionRoot> {
556 &self.schema
557 }
558}
559
560use async_graphql_axum::{GraphQL, GraphQLSubscription};
565use axum::routing::get_service;
566
567#[must_use = "the router must be served to handle requests"]
575pub fn serve_graphql(
576 schema: async_graphql::Schema<QueryRoot, MutationRoot, SubscriptionRoot>,
577) -> axum::Router {
578 axum::Router::new().route_service(
579 "/graphql",
580 get_service(GraphQLSubscription::new(schema.clone())).post_service(GraphQL::new(schema)),
581 )
582}
583
584#[derive(Debug, Clone, PartialEq, Error)]
590pub enum GraphqlError {
591 #[error("graphql schema build failed: {0}")]
593 SchemaBuildFailed(String),
594 #[error("graphql resolution failed: {0}")]
596 ResolutionFailed(String),
597 #[error("graphql mapping error: {0}")]
599 MappingError(String),
600}
601
602impl From<GraphqlError> for Error {
603 fn from(e: GraphqlError) -> Self {
604 match e {
605 GraphqlError::SchemaBuildFailed(_) | GraphqlError::MappingError(_) => {
606 Error::internal_error(e.to_string())
607 }
608 GraphqlError::ResolutionFailed(_) => Error::resolver_error(e.to_string()),
609 }
610 }
611}
612
613#[cfg(test)]
618mod tests {
619 use super::*;
620 use themql_core::CorrelationId;
621 use themql_core::{ErrorCode, Payload, Response};
622
623 struct StubResolver;
626
627 #[allow(clippy::unused_async_trait_impl)]
628 impl themql_core::ResolverBoxed for StubResolver {
629 async fn resolve(&self, query: &Query, _ctx: &Context) -> Result<Response, Error> {
630 let subject = query.resource.subject.as_str();
631 Ok(Response::ok(
632 ResponseValue::Json(serde_json::json!({ "subject": subject })),
633 CorrelationId::new(),
634 ))
635 }
636 }
637
638 struct StubHandler;
639
640 impl MessageHandler for StubHandler {
641 fn handle<'a>(
642 &'a self,
643 msg: &'a Message,
644 ) -> Pin<Box<dyn Future<Output = Result<Response, Error>> + Send + 'a>> {
645 let subject = msg.subject.as_str();
646 let payload = msg.payload.clone();
647 Box::pin(async move {
648 let value = match payload {
649 Payload::Json(v) => v,
650 Payload::Unit => serde_json::Value::Null,
651 Payload::Bytes(b, tag) => serde_json::json!({
652 "format": tag.to_string(),
653 "len": b.len(),
654 }),
655 };
656 Ok(Response::ok(
657 ResponseValue::Json(serde_json::json!({
658 "subject": subject,
659 "echo": value,
660 })),
661 CorrelationId::new(),
662 ))
663 })
664 }
665 }
666
667 #[test]
670 fn marker_roots_are_constructible() {
671 let bridge: Arc<dyn GraphqlResolverBridge> =
672 Arc::new(GraphqlResolverBridgeImpl::new(Arc::new(StubResolver)));
673 let dispatch: Arc<dyn DispatchBridge> =
674 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
675 let _: QueryRoot = QueryRoot::new(bridge);
676 let _: MutationRoot = MutationRoot::new(dispatch);
677 let _: SubscriptionRoot = SubscriptionRoot::default();
678 }
679
680 #[test]
681 fn graphql_error_schema_build_maps_to_internal_error() {
682 let e: Error = GraphqlError::SchemaBuildFailed("bad sdl".to_owned()).into();
683 assert_eq!(e.code, ErrorCode::InternalError);
684 }
685
686 #[test]
687 fn graphql_error_mapping_maps_to_internal_error() {
688 let e: Error = GraphqlError::MappingError("no field".to_owned()).into();
689 assert_eq!(e.code, ErrorCode::InternalError);
690 }
691
692 #[test]
693 fn graphql_error_resolution_maps_to_resolver_error() {
694 let e: Error = GraphqlError::ResolutionFailed("field x".to_owned()).into();
695 assert_eq!(e.code, ErrorCode::ResolverError);
696 }
697
698 #[test]
701 fn resolver_bridge_impl_construction() {
702 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
703 let bridge = GraphqlResolverBridgeImpl::new(resolver);
704 let _: &dyn GraphqlResolverBridge = &bridge;
705 }
706
707 #[test]
708 fn dispatch_bridge_impl_construction() {
709 let handler: Arc<dyn MessageHandler> = Arc::new(StubHandler);
710 let bridge = DispatchBridgeImpl::new(handler);
711 let _: &dyn DispatchBridge = &bridge;
712 }
713
714 #[tokio::test]
715 async fn schema_builds_with_real_resolver_bridge() {
716 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
717 let bridge: Arc<dyn GraphqlResolverBridge> =
718 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
719 let dispatch: Arc<dyn DispatchBridge> =
720 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
721 let schema = GraphqlSchemaImpl::new(
722 QueryRoot::new(Arc::clone(&bridge)),
723 MutationRoot::new(Arc::clone(&dispatch)),
724 SubscriptionRoot::default(),
725 );
726 let _ = schema.schema();
727 }
728
729 #[tokio::test]
730 async fn query_resource_resolves_via_bridge() {
731 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
732 let bridge: Arc<dyn GraphqlResolverBridge> =
733 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
734 let dispatch: Arc<dyn DispatchBridge> =
735 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
736 let schema = GraphqlSchemaImpl::with_role(
737 QueryRoot::new(bridge),
738 MutationRoot::new(dispatch),
739 SubscriptionRoot::default(),
740 AuthRole::Observer,
741 );
742 let q = r#"{ resource(subject: "vehicle.sensors.imu") }"#;
743 let result = schema.schema().execute(q).await;
744 assert!(
745 result.errors.is_empty(),
746 "expected no errors, got {result:?}"
747 );
748 let data = result.data.into_json().expect("data to json");
749 assert_eq!(data["resource"]["subject"], "vehicle.sensors.imu");
750 }
751
752 #[tokio::test]
753 async fn mutation_dispatch_resolves_via_bridge() {
754 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
755 let bridge: Arc<dyn GraphqlResolverBridge> =
756 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
757 let dispatch: Arc<dyn DispatchBridge> =
758 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
759 let schema = GraphqlSchemaImpl::with_role(
760 QueryRoot::new(bridge),
761 MutationRoot::new(dispatch),
762 SubscriptionRoot::default(),
763 AuthRole::Operator,
764 );
765 let q = r#"mutation { dispatch(subject: "vehicle.command", payload: {x: 1}) }"#;
766 let result = schema.schema().execute(q).await;
767 assert!(
768 result.errors.is_empty(),
769 "expected no errors, got {result:?}"
770 );
771 let data = result.data.into_json().expect("data to json");
772 assert_eq!(data["dispatch"]["subject"], "vehicle.command");
773 assert_eq!(data["dispatch"]["echo"]["x"], 1);
774 }
775
776 #[tokio::test]
777 async fn serve_graphql_builds_router() {
778 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
779 let bridge: Arc<dyn GraphqlResolverBridge> =
780 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
781 let dispatch: Arc<dyn DispatchBridge> =
782 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
783 let schema = async_graphql::Schema::build(
784 QueryRoot::new(bridge),
785 MutationRoot::new(dispatch),
786 SubscriptionRoot::default(),
787 )
788 .finish();
789 let router: axum::Router = serve_graphql(schema);
790 let _ = router;
791 }
792
793 #[test]
794 fn response_value_to_json_json_variant() {
795 let v = response_value_to_json(ResponseValue::Json(serde_json::json!({"a": 1})));
796 assert_eq!(v["a"], 1);
797 }
798
799 #[test]
800 fn response_value_to_json_unit_variant() {
801 let v = response_value_to_json(ResponseValue::Unit);
802 assert!(v.is_null());
803 }
804
805 #[test]
806 fn resolver_bridge_rejects_invalid_subject() {
807 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
808 let bridge = GraphqlResolverBridgeImpl::new(resolver);
809 let ctx = Context::new();
810 let args = serde_json::Value::Null;
811 let rt = tokio::runtime::Builder::new_current_thread()
812 .enable_time()
813 .build()
814 .expect("rt build");
815 let err = rt
816 .block_on(bridge.resolve_field("not..valid", &args, &ctx))
817 .expect_err("error for invalid subject");
818 assert_eq!(err.code, ErrorCode::ValidationError);
819 }
820
821 #[tokio::test]
822 async fn subscription_with_source_streams_real_events() {
823 use futures_util::StreamExt;
824 use themql_core::Operation;
825 use themql_sse::{SsePublisher, TokioSsePublisher};
826
827 let publisher = Arc::new(TokioSsePublisher::new());
828 let source: Arc<dyn GraphqlSubscriptionSource> =
829 Arc::clone(&publisher) as Arc<dyn GraphqlSubscriptionSource>;
830 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
831 let bridge: Arc<dyn GraphqlResolverBridge> =
832 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
833 let dispatch: Arc<dyn DispatchBridge> =
834 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
835 let schema = GraphqlSchemaImpl::with_role(
836 QueryRoot::new(bridge),
837 MutationRoot::new(dispatch),
838 SubscriptionRoot::with_source(source),
839 AuthRole::Observer,
840 );
841
842 let publisher_clone = Arc::clone(&publisher);
843 tokio::spawn(async move {
844 tokio::time::sleep(std::time::Duration::from_millis(200)).await;
845 let subject = Subject::from_str("vehicle.sensors.imu").expect("subject");
846 let subject_str = subject.as_str();
847 let msg = Message::new(&subject_str, Operation::Event).expect("message");
848 publisher_clone
849 .broadcast(&subject, &msg)
850 .await
851 .expect("broadcast");
852 });
853
854 let q = r#"subscription { subscribe(subject: "vehicle.sensors.imu") }"#;
855 let mut stream = schema.schema().execute_stream(q);
856 let result = tokio::time::timeout(std::time::Duration::from_secs(5), stream.next()).await;
857
858 assert!(
859 result.is_ok(),
860 "subscription should produce at least one event"
861 );
862 let response = result.expect("timeout").expect("event");
863 assert!(
864 response.errors.is_empty(),
865 "no errors expected, got {response:?}"
866 );
867 }
868
869 #[tokio::test]
870 async fn subscription_without_source_returns_error() {
871 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
872 let bridge: Arc<dyn GraphqlResolverBridge> =
873 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
874 let dispatch: Arc<dyn DispatchBridge> =
875 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
876 let schema = GraphqlSchemaImpl::new(
877 QueryRoot::new(bridge),
878 MutationRoot::new(dispatch),
879 SubscriptionRoot::default(),
880 );
881 let q = r#"subscription { subscribe(subject: "vehicle.sensors.imu") }"#;
882 let result = schema.schema().execute(q).await;
883 assert!(
884 !result.errors.is_empty(),
885 "expected error for no source, got {result:?}"
886 );
887 }
888
889 #[test]
892 fn auth_role_satisfies_hierarchy() {
893 assert!(AuthRole::Admin.satisfies(AuthRole::Admin));
894 assert!(AuthRole::Admin.satisfies(AuthRole::Operator));
895 assert!(AuthRole::Admin.satisfies(AuthRole::Observer));
896 assert!(AuthRole::Operator.satisfies(AuthRole::Operator));
897 assert!(AuthRole::Operator.satisfies(AuthRole::Observer));
898 assert!(!AuthRole::Operator.satisfies(AuthRole::Admin));
899 assert!(AuthRole::Observer.satisfies(AuthRole::Observer));
900 assert!(!AuthRole::Observer.satisfies(AuthRole::Operator));
901 assert!(!AuthRole::Observer.satisfies(AuthRole::Admin));
902 }
903
904 #[test]
905 fn auth_role_from_str_maps_known_roles_case_insensitive() {
906 use std::str::FromStr;
907 assert_eq!(AuthRole::from_str("admin"), Ok(AuthRole::Admin));
908 assert_eq!(AuthRole::from_str("Operator"), Ok(AuthRole::Operator));
909 assert_eq!(AuthRole::from_str("OBSERVER"), Ok(AuthRole::Observer));
910 }
911
912 #[test]
913 fn auth_role_from_str_unknown_returns_err() {
914 use std::str::FromStr;
915 assert!(AuthRole::from_str("root").is_err());
916 assert!(AuthRole::from_str("").is_err());
917 assert!(AuthRole::from_str("superuser").is_err());
918 }
919
920 #[tokio::test]
921 async fn query_rejected_without_role() {
922 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
923 let bridge: Arc<dyn GraphqlResolverBridge> =
924 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
925 let dispatch: Arc<dyn DispatchBridge> =
926 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
927 let schema = GraphqlSchemaImpl::new(
928 QueryRoot::new(bridge),
929 MutationRoot::new(dispatch),
930 SubscriptionRoot::default(),
931 );
932 let q = r#"{ resource(subject: "vehicle.sensors.imu") }"#;
933 let result = schema.schema().execute(q).await;
934 assert!(
935 !result.errors.is_empty(),
936 "expected authz error, got {result:?}"
937 );
938 assert!(result.errors[0].message.contains("no auth role"));
939 }
940
941 #[tokio::test]
942 async fn mutation_rejected_for_observer() {
943 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
944 let bridge: Arc<dyn GraphqlResolverBridge> =
945 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
946 let dispatch: Arc<dyn DispatchBridge> =
947 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
948 let schema = GraphqlSchemaImpl::with_role(
949 QueryRoot::new(bridge),
950 MutationRoot::new(dispatch),
951 SubscriptionRoot::default(),
952 AuthRole::Observer,
953 );
954 let q = r#"mutation { dispatch(subject: "vehicle.command", payload: {x: 1}) }"#;
955 let result = schema.schema().execute(q).await;
956 assert!(
957 !result.errors.is_empty(),
958 "expected authz error for observer on mutation, got {result:?}"
959 );
960 assert!(result.errors[0].message.contains("insufficient role"));
961 }
962
963 #[tokio::test]
964 async fn mutation_allowed_for_admin() {
965 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
966 let bridge: Arc<dyn GraphqlResolverBridge> =
967 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
968 let dispatch: Arc<dyn DispatchBridge> =
969 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
970 let schema = GraphqlSchemaImpl::with_role(
971 QueryRoot::new(bridge),
972 MutationRoot::new(dispatch),
973 SubscriptionRoot::default(),
974 AuthRole::Admin,
975 );
976 let q = r#"mutation { dispatch(subject: "vehicle.command", payload: {x: 1}) }"#;
977 let result = schema.schema().execute(q).await;
978 assert!(
979 result.errors.is_empty(),
980 "expected no authz error for admin, got {result:?}"
981 );
982 }
983
984 #[tokio::test]
985 async fn query_allowed_for_operator() {
986 let resolver: Arc<dyn themql_core::Resolver> = Arc::new(StubResolver);
987 let bridge: Arc<dyn GraphqlResolverBridge> =
988 Arc::new(GraphqlResolverBridgeImpl::new(resolver));
989 let dispatch: Arc<dyn DispatchBridge> =
990 Arc::new(DispatchBridgeImpl::new(Arc::new(StubHandler)));
991 let schema = GraphqlSchemaImpl::with_role(
992 QueryRoot::new(bridge),
993 MutationRoot::new(dispatch),
994 SubscriptionRoot::default(),
995 AuthRole::Operator,
996 );
997 let q = r#"{ resource(subject: "vehicle.sensors.imu") }"#;
998 let result = schema.schema().execute(q).await;
999 assert!(
1000 result.errors.is_empty(),
1001 "expected no authz error for operator on query, got {result:?}"
1002 );
1003 }
1004}