1use std::{
8 convert::Infallible,
9 error::Error,
10 panic::{catch_unwind, AssertUnwindSafe},
11 sync::Arc,
12 task::{Context, Poll},
13};
14
15use arc_swap::ArcSwap;
16use async_trait::async_trait;
17use axum::{body::Body as AxumBody, routing::future::RouteFuture, Router};
18use http::{Request, Response};
19use rcgen::{generate_simple_self_signed, CertifiedKey};
20use rustls::{pki_types::PrivateKeyDer, ServerConfig};
21use rustls_pki_types::PrivatePkcs8KeyDer;
22use saluki_api::{APIHandler, DynamicRoute, EndpointType};
23use saluki_common::{collections::FastIndexMap, sync::shutdown::ShutdownHandle};
24use saluki_core::runtime::{
25 state::{DataspaceRegistry, DataspaceUpdate, Identifier, IdentifierFilter, Subscription},
26 AutoShutdown, InitializationError, Supervisable, Supervisor, SupervisorFuture,
27};
28use saluki_error::{generic_error, GenericError};
29use saluki_io::net::{
30 server::{grpc::unmatched_route, http::HttpServer},
31 ListenAddress,
32};
33use saluki_tls::ensure_server_config_fips_compliant;
34use tokio::{pin, select};
35use tonic::{body::Body as GrpcBody, server::NamedService, service::RoutesBuilder};
36use tower::Service;
37use tracing::{debug, info, warn};
38
39pub struct APIBuilder {
75 endpoint_type: EndpointType,
76 listen_address: ListenAddress,
77 tls_config: Option<ServerConfig>,
78 http_router: Router,
79 grpc_router: RoutesBuilder,
80}
81
82impl APIBuilder {
83 pub fn new(endpoint_type: EndpointType, listen_address: ListenAddress) -> Self {
85 Self {
86 endpoint_type,
87 listen_address,
88 tls_config: None,
89 http_router: Router::new(),
90 grpc_router: RoutesBuilder::default(),
91 }
92 }
93
94 pub fn with_handler<H>(mut self, handler: H) -> Self
99 where
100 H: APIHandler,
101 {
102 let handler_router = handler.generate_routes();
103 let handler_state = handler.generate_initial_state();
104 self.http_router = self.http_router.merge(handler_router.with_state(handler_state));
105 self
106 }
107
108 pub fn with_optional_handler<H>(self, handler: Option<H>) -> Self
113 where
114 H: APIHandler,
115 {
116 if let Some(handler) = handler {
117 self.with_handler(handler)
118 } else {
119 self
120 }
121 }
122
123 pub fn with_grpc_service<S>(mut self, svc: S) -> Self
125 where
126 S: Service<Request<GrpcBody>, Response = Response<GrpcBody>, Error = Infallible>
127 + NamedService
128 + Clone
129 + Send
130 + Sync
131 + 'static,
132 S::Future: Send + 'static,
133 S::Error: Into<Box<dyn Error + Send + Sync>> + Send,
134 {
135 self.grpc_router.add_service(svc);
136 self
137 }
138
139 pub fn with_tls_config(mut self, config: ServerConfig) -> Self {
141 self.tls_config = Some(config);
142 self
143 }
144
145 pub fn with_self_signed_tls(self) -> Self {
147 self.try_with_self_signed_tls()
148 .expect("self-signed server TLS configuration should build and pass FIPS validation")
149 }
150
151 pub fn try_with_self_signed_tls(self) -> Result<Self, GenericError> {
158 let CertifiedKey { cert, signing_key } = generate_simple_self_signed(["localhost".to_owned()])?;
159 let cert_chain = vec![cert.der().clone()];
160 let key = PrivateKeyDer::Pkcs8(PrivatePkcs8KeyDer::from(signing_key.serialize_der()));
161
162 let mut config = ServerConfig::builder()
163 .with_no_client_auth()
164 .with_single_cert(cert_chain, key)?;
165
166 ensure_server_config_fips_compliant(&mut config)?;
167
168 Ok(self.with_tls_config(config))
169 }
170}
171
172impl APIBuilder {
173 pub fn into_supervisor(self) -> Supervisor {
178 let base = self
185 .http_router
186 .reset_fallback()
187 .merge(self.grpc_router.routes().into_axum_router().reset_fallback());
188
189 let (inner, outer) = create_dynamic_router(apply_fallback(base.clone()));
191
192 let mut http_server = HttpServer::from_listen_address(self.listen_address.clone())
193 .with_routes(outer)
194 .with_server_id(format!("{}-api", self.endpoint_type.name()));
195 if let Some(tls_config) = self.tls_config {
196 http_server = http_server.with_tls_config(tls_config);
197 }
198
199 let endpoint_type = self.endpoint_type;
200 let name = match endpoint_type {
201 EndpointType::Unprivileged => "unprivileged-api",
202 EndpointType::Privileged => "privileged-api",
203 };
204
205 let mut supervisor = Supervisor::new(name)
206 .expect("API supervisor name is a non-empty constant")
207 .with_auto_shutdown(AutoShutdown::AnySignificant);
208
209 supervisor.add_worker(http_server.into_supervisor());
210 supervisor.add_worker(RouteUpdaterWorker {
211 inner,
212 base,
213 endpoint_type,
214 listen_address: self.listen_address,
215 });
216
217 supervisor
218 }
219}
220
221struct RouteUpdaterWorker {
223 inner: Arc<ArcSwap<Router>>,
224 base: Router,
225 endpoint_type: EndpointType,
226 listen_address: ListenAddress,
227}
228
229#[async_trait]
230impl Supervisable for RouteUpdaterWorker {
231 fn name(&self) -> &str {
232 "route_updater"
233 }
234
235 async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
236 let dataspace = DataspaceRegistry::try_current().ok_or_else(|| generic_error!("Dataspace not available."))?;
237
238 let inner = Arc::clone(&self.inner);
239 let base = self.base.clone();
240 let endpoint_type = self.endpoint_type;
241 let listen_address = self.listen_address.clone();
242
243 Ok(Box::pin(async move {
244 info!("Serving {} API on {}.", endpoint_type.name(), listen_address);
245
246 let route_assertions = dataspace.subscribe::<DynamicRoute>(IdentifierFilter::All);
248
249 run_event_loop(process_shutdown, inner, base, route_assertions, endpoint_type).await
250 }))
251 }
252}
253
254#[derive(Clone)]
261struct DynamicRouterService {
262 inner_router: Arc<ArcSwap<Router>>,
263}
264
265impl DynamicRouterService {
266 fn from_inner(inner_router: &Arc<ArcSwap<Router>>) -> Self {
267 Self {
268 inner_router: Arc::clone(inner_router),
269 }
270 }
271}
272
273impl Service<http::Request<AxumBody>> for DynamicRouterService {
274 type Response = Response<AxumBody>;
275 type Error = Infallible;
276 type Future = RouteFuture<Infallible>;
277
278 fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
279 Poll::Ready(Ok(()))
280 }
281
282 fn call(&mut self, request: http::Request<AxumBody>) -> Self::Future {
283 let mut router = Arc::unwrap_or_clone(self.inner_router.load_full());
284 router.call(request)
285 }
286}
287
288async fn run_event_loop(
290 process_shutdown: ShutdownHandle, inner: Arc<ArcSwap<Router>>, base: Router,
291 mut route_assertions: Subscription<DynamicRoute>, endpoint_type: EndpointType,
292) -> Result<(), GenericError> {
293 let mut handlers = FastIndexMap::default();
296
297 pin!(process_shutdown);
298
299 loop {
300 select! {
301 _ = &mut process_shutdown => break,
302 maybe_update = route_assertions.recv() => match maybe_update {
303 None => {
304 debug!("Dataspace subscription ended unexpectedly.");
305 break
306 },
307 Some(update) => {
308 match update {
309 DataspaceUpdate::Asserted(id, route) => {
310 if route.endpoint_type() != endpoint_type {
311 continue;
312 }
313
314 debug!(?id, "Registering dynamic {} handler.", route.endpoint_protocol().name());
315
316 handlers.insert(id, route.into_router());
317 }
318 DataspaceUpdate::Retracted(id) => {
319 if handlers.swap_remove(&id).is_none() {
320 continue;
321 }
322
323 debug!(?id, "Withdrawing dynamic handler.");
324 }
325 DataspaceUpdate::Message(..) => continue,
327 }
328
329 rebuild_router(&inner, &base, &handlers);
330 }
331 }
332 }
333 }
334
335 Ok(())
336}
337
338fn create_dynamic_router(initial: Router) -> (Arc<ArcSwap<Router>>, Router) {
341 let inner = Arc::new(ArcSwap::from_pointee(initial));
342 let outer = Router::new().fallback_service(DynamicRouterService::from_inner(&inner));
343 (inner, outer)
344}
345
346fn try_merge_router(base: &Router, id: &Identifier, other: &Router) -> Result<Router, String> {
363 let candidate = base.clone();
364 match catch_unwind(AssertUnwindSafe(|| candidate.merge(other.clone()))) {
365 Ok(merged) => Ok(merged),
366 Err(payload) => {
367 let reason = payload
368 .downcast_ref::<String>()
369 .map(|s| s.as_str())
370 .or_else(|| payload.downcast_ref::<&str>().copied())
371 .unwrap_or("unknown");
372 Err(format!("failed to merge dynamic handler {id:?}: {reason}"))
373 }
374 }
375}
376
377fn rebuild_router(inner_router: &Arc<ArcSwap<Router>>, base: &Router, handlers: &FastIndexMap<Identifier, Router>) {
380 let mut merged = base.clone();
381 let mut skipped = 0usize;
382
383 for (id, router) in handlers.iter() {
384 let resetable = router.clone().reset_fallback();
385 match try_merge_router(&merged, id, &resetable) {
386 Ok(new_merged) => merged = new_merged,
387 Err(reason) => {
388 warn!(%reason, "Skipping dynamic handler due to overlapping route.");
389 skipped += 1;
390 }
391 }
392 }
393
394 inner_router.store(Arc::new(apply_fallback(merged)));
395 debug!(handler_count = handlers.len(), skipped, "Rebuilt inner router.");
396}
397
398fn apply_fallback(router: Router) -> Router {
403 router.fallback(unmatched_route)
404}
405
406#[cfg(test)]
407mod tests {
408 use std::{net::SocketAddr, time::Duration};
409
410 use async_trait::async_trait;
411 use axum::Router;
412 use http_body_util::{BodyExt as _, Empty};
413 use hyper::{body::Bytes, StatusCode};
414 use hyper_util::{client::legacy::Client, rt::TokioExecutor};
415 use saluki_api::{APIHandler, DynamicRoute, EndpointType};
416 use saluki_core::runtime::{
417 state::{DataspaceRegistry, DataspaceUpdate, Identifier, IdentifierFilter},
418 InitializationError, Supervisable, Supervisor, SupervisorFuture,
419 };
420 use saluki_io::net::BoundListenAddress;
421 use tokio::{
422 pin, select,
423 sync::{mpsc, oneshot},
424 task::JoinHandle,
425 time::{sleep, timeout, Instant},
426 };
427
428 use super::*;
429
430 struct SimpleHandler {
431 path: &'static str,
432 body: &'static str,
433 }
434
435 impl APIHandler for SimpleHandler {
436 type State = ();
437
438 fn generate_initial_state(&self) -> Self::State {}
439
440 fn generate_routes(&self) -> Router<Self::State> {
441 let body = self.body;
442 Router::new().route(self.path, axum::routing::get(move || async move { body }))
443 }
444 }
445
446 enum RouteCommand {
447 Assert { id: Identifier, route: DynamicRoute },
448 Retract { id: Identifier },
449 }
450
451 struct RouteAsserter {
452 commands_rx: std::sync::Mutex<Option<mpsc::Receiver<RouteCommand>>>,
453 addr_tx: std::sync::Mutex<Option<oneshot::Sender<SocketAddr>>>,
454 endpoint_type: EndpointType,
455 }
456
457 #[async_trait]
458 impl Supervisable for RouteAsserter {
459 fn name(&self) -> &str {
460 "route-asserter"
461 }
462
463 async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
464 let mut commands_rx =
465 self.commands_rx
466 .lock()
467 .unwrap()
468 .take()
469 .ok_or_else(|| InitializationError::Failed {
470 source: generic_error!("RouteAsserter can only be initialized once"),
471 })?;
472 let addr_tx = self.addr_tx.lock().unwrap().take();
473 let endpoint_type = self.endpoint_type;
474
475 Ok(Box::pin(async move {
476 let dataspace =
477 DataspaceRegistry::try_current().ok_or_else(|| generic_error!("Dataspace not available."))?;
478
479 let bound_addr_name = match endpoint_type {
481 EndpointType::Unprivileged => "http-server-unprivileged-api",
482 EndpointType::Privileged => "http-server-privileged-api",
483 };
484 let mut addr_sub = dataspace
485 .subscribe::<BoundListenAddress>(IdentifierFilter::exact(Identifier::named(bound_addr_name)));
486
487 let addr = match addr_sub.recv().await {
488 Some(DataspaceUpdate::Asserted(_, BoundListenAddress::Tcp(mut addr))) => {
489 if addr.ip().is_unspecified() {
491 addr.set_ip(std::net::Ipv4Addr::LOCALHOST.into());
492 }
493 addr
494 }
495 other => return Err(generic_error!("unexpected bound address update: {:?}", other)),
496 };
497
498 if let Some(tx) = addr_tx {
499 let _ = tx.send(addr);
500 }
501
502 pin!(process_shutdown);
504
505 loop {
506 select! {
507 _ = &mut process_shutdown => break,
508 cmd = commands_rx.recv() => {
509 let Some(cmd) = cmd else { break };
510 match cmd {
511 RouteCommand::Assert { id, route } => {
512 dataspace.assert(route, id);
513 }
514 RouteCommand::Retract { id } => {
515 dataspace.retract::<DynamicRoute>(id);
516 }
517 }
518 }
519 }
520 }
521
522 Ok(())
523 }))
524 }
525 }
526
527 struct TestHarness {
528 addr: SocketAddr,
529 commands: mpsc::Sender<RouteCommand>,
530 _shutdown: oneshot::Sender<()>,
531 _handle: JoinHandle<()>,
532 }
533
534 impl TestHarness {
535 async fn assert_route(&self, id: impl Into<Identifier>, route: DynamicRoute) {
536 self.commands
537 .send(RouteCommand::Assert { id: id.into(), route })
538 .await
539 .unwrap();
540 }
541
542 async fn retract_route(&self, id: impl Into<Identifier>) {
543 self.commands
544 .send(RouteCommand::Retract { id: id.into() })
545 .await
546 .unwrap();
547 }
548 }
549
550 async fn setup_test_harness(endpoint_type: EndpointType) -> TestHarness {
551 setup_test_harness_with(endpoint_type, |b| b).await
552 }
553
554 async fn setup_test_harness_with<F>(endpoint_type: EndpointType, configure: F) -> TestHarness
555 where
556 F: FnOnce(APIBuilder) -> APIBuilder,
557 {
558 let (commands_tx, commands_rx) = mpsc::channel(16);
559 let (addr_tx, addr_rx) = oneshot::channel();
560
561 let api_builder = configure(APIBuilder::new(endpoint_type, ListenAddress::tcp_any(0)));
562 let route_asserter = RouteAsserter {
563 commands_rx: std::sync::Mutex::new(Some(commands_rx)),
564 addr_tx: std::sync::Mutex::new(Some(addr_tx)),
565 endpoint_type,
566 };
567
568 let mut sup = Supervisor::new("test-dynamic-api").unwrap();
569 sup.add_worker(api_builder.into_supervisor());
570 sup.add_worker(route_asserter);
571
572 let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
573 let handle = tokio::spawn(async move {
574 let _ = sup.run_with_shutdown(shutdown_rx).await;
575 });
576
577 let addr = timeout(Duration::from_secs(5), addr_rx)
578 .await
579 .expect("timed out waiting for bound address")
580 .expect("addr channel closed");
581
582 TestHarness {
583 addr,
584 commands: commands_tx,
585 _shutdown: shutdown_tx,
586 _handle: handle,
587 }
588 }
589
590 async fn http_get(addr: SocketAddr, path: &str) -> (StatusCode, String) {
591 let client: Client<_, Empty<Bytes>> = Client::builder(TokioExecutor::new()).build_http();
592 let uri = format!("http://{}{}", addr, path);
593 let resp = client.get(uri.parse().unwrap()).await.unwrap();
594 let status = resp.status();
595 let body = resp.into_body().collect().await.unwrap().to_bytes();
596 let body_str = String::from_utf8_lossy(&body).into_owned();
597 (status, body_str)
598 }
599
600 async fn grpc_post(addr: SocketAddr, path: &str) -> (StatusCode, http::HeaderMap) {
601 let client: Client<_, Empty<Bytes>> = Client::builder(TokioExecutor::new()).build_http();
602 let uri: hyper::Uri = format!("http://{}{}", addr, path).parse().unwrap();
603 let req = hyper::Request::builder()
604 .uri(uri)
605 .method(hyper::Method::POST)
606 .header(hyper::header::CONTENT_TYPE, "application/grpc")
607 .body(Empty::<Bytes>::new())
608 .unwrap();
609 let resp = client.request(req).await.unwrap();
610 (resp.status(), resp.headers().clone())
611 }
612
613 async fn assert_status_eventually(addr: SocketAddr, path: &str, expected: StatusCode) -> String {
614 let deadline = Instant::now() + Duration::from_secs(2);
615 loop {
616 let (status, body) = http_get(addr, path).await;
617 if status == expected {
618 return body;
619 }
620 if Instant::now() > deadline {
621 panic!("expected {} for {} but got {}", expected, path, status);
622 }
623 sleep(Duration::from_millis(50)).await;
624 }
625 }
626
627 #[tokio::test]
630 async fn serves_asserted_http_route() {
631 let harness = setup_test_harness(EndpointType::Unprivileged).await;
632
633 let route = DynamicRoute::http(
634 EndpointType::Unprivileged,
635 SimpleHandler {
636 path: "/health",
637 body: "ok",
638 },
639 );
640 harness.assert_route("health", route).await;
641
642 let body = assert_status_eventually(harness.addr, "/health", StatusCode::OK).await;
643 assert_eq!(body, "ok");
644 }
645
646 #[tokio::test]
647 async fn returns_404_for_unknown_route() {
648 let harness = setup_test_harness(EndpointType::Unprivileged).await;
649 let (status, _) = http_get(harness.addr, "/nonexistent").await;
650 assert_eq!(status, StatusCode::NOT_FOUND);
651 }
652
653 #[tokio::test]
654 async fn route_retraction_removes_route() {
655 let harness = setup_test_harness(EndpointType::Unprivileged).await;
656
657 let route = DynamicRoute::http(
658 EndpointType::Unprivileged,
659 SimpleHandler {
660 path: "/temp",
661 body: "temporary",
662 },
663 );
664 harness.assert_route("temp", route).await;
665 assert_status_eventually(harness.addr, "/temp", StatusCode::OK).await;
666
667 harness.retract_route("temp").await;
668 assert_status_eventually(harness.addr, "/temp", StatusCode::NOT_FOUND).await;
669 }
670
671 #[tokio::test]
672 async fn multiple_routes_independent_lifecycle() {
673 let harness = setup_test_harness(EndpointType::Unprivileged).await;
674
675 let route_a = DynamicRoute::http(
676 EndpointType::Unprivileged,
677 SimpleHandler {
678 path: "/a",
679 body: "alpha",
680 },
681 );
682 let route_b = DynamicRoute::http(
683 EndpointType::Unprivileged,
684 SimpleHandler {
685 path: "/b",
686 body: "bravo",
687 },
688 );
689 harness.assert_route("a", route_a).await;
690 harness.assert_route("b", route_b).await;
691
692 assert_status_eventually(harness.addr, "/a", StatusCode::OK).await;
693 assert_status_eventually(harness.addr, "/b", StatusCode::OK).await;
694
695 harness.retract_route("a").await;
697 assert_status_eventually(harness.addr, "/a", StatusCode::NOT_FOUND).await;
698
699 let body = assert_status_eventually(harness.addr, "/b", StatusCode::OK).await;
700 assert_eq!(body, "bravo");
701 }
702
703 #[tokio::test]
704 async fn ignores_routes_for_different_endpoint_type() {
705 let harness = setup_test_harness(EndpointType::Unprivileged).await;
706
707 let wrong_route = DynamicRoute::http(
709 EndpointType::Privileged,
710 SimpleHandler {
711 path: "/secret",
712 body: "secret",
713 },
714 );
715 harness.assert_route("secret", wrong_route).await;
716
717 let (status, _) = http_get(harness.addr, "/secret").await;
718 assert_eq!(status, StatusCode::NOT_FOUND);
719
720 let right_route = DynamicRoute::http(
722 EndpointType::Unprivileged,
723 SimpleHandler {
724 path: "/secret",
725 body: "not secret",
726 },
727 );
728 harness.assert_route("secret-unpriv", right_route).await;
729
730 let body = assert_status_eventually(harness.addr, "/secret", StatusCode::OK).await;
731 assert_eq!(body, "not secret");
732 }
733
734 #[tokio::test]
735 async fn overlapping_routes_do_not_crash_server() {
736 let harness = setup_test_harness(EndpointType::Unprivileged).await;
737
738 let route_1 = DynamicRoute::http(
740 EndpointType::Unprivileged,
741 SimpleHandler {
742 path: "/health",
743 body: "health-1",
744 },
745 );
746 harness.assert_route("health-1", route_1).await;
747 let body = assert_status_eventually(harness.addr, "/health", StatusCode::OK).await;
748 assert_eq!(body, "health-1");
749
750 let route_2 = DynamicRoute::http(
753 EndpointType::Unprivileged,
754 SimpleHandler {
755 path: "/health",
756 body: "health-2",
757 },
758 );
759 harness.assert_route("health-2", route_2).await;
760
761 sleep(Duration::from_millis(200)).await;
763
764 let (status, body) = http_get(harness.addr, "/health").await;
766 assert_eq!(status, StatusCode::OK);
767 assert_eq!(body, "health-1");
768
769 let route_info = DynamicRoute::http(
771 EndpointType::Unprivileged,
772 SimpleHandler {
773 path: "/info",
774 body: "info",
775 },
776 );
777 harness.assert_route("info", route_info).await;
778 let body = assert_status_eventually(harness.addr, "/info", StatusCode::OK).await;
779 assert_eq!(body, "info");
780
781 harness.retract_route("health-1").await;
784 let body = assert_status_eventually(harness.addr, "/health", StatusCode::OK).await;
785 assert_eq!(body, "health-2");
786 }
787
788 #[tokio::test]
789 async fn overlapping_route_retraction_then_reassertion() {
790 let harness = setup_test_harness(EndpointType::Unprivileged).await;
791
792 let route_a = DynamicRoute::http(
794 EndpointType::Unprivileged,
795 SimpleHandler {
796 path: "/overlap",
797 body: "a",
798 },
799 );
800 let route_b = DynamicRoute::http(
801 EndpointType::Unprivileged,
802 SimpleHandler {
803 path: "/overlap",
804 body: "b",
805 },
806 );
807 harness.assert_route("ov-a", route_a).await;
808 harness.assert_route("ov-b", route_b).await;
809
810 let body = assert_status_eventually(harness.addr, "/overlap", StatusCode::OK).await;
812 assert_eq!(body, "a");
813
814 harness.retract_route("ov-a").await;
816 harness.retract_route("ov-b").await;
817 assert_status_eventually(harness.addr, "/overlap", StatusCode::NOT_FOUND).await;
818
819 let route_c = DynamicRoute::http(
821 EndpointType::Unprivileged,
822 SimpleHandler {
823 path: "/overlap",
824 body: "c",
825 },
826 );
827 harness.assert_route("ov-c", route_c).await;
828 let body = assert_status_eventually(harness.addr, "/overlap", StatusCode::OK).await;
829 assert_eq!(body, "c");
830 }
831
832 #[tokio::test]
833 async fn static_handler_served_without_dynamic_routes() {
834 let harness = setup_test_harness_with(EndpointType::Unprivileged, |b| {
835 b.with_handler(SimpleHandler {
836 path: "/static",
837 body: "static",
838 })
839 })
840 .await;
841
842 let body = assert_status_eventually(harness.addr, "/static", StatusCode::OK).await;
843 assert_eq!(body, "static");
844 }
845
846 #[tokio::test]
847 async fn static_and_dynamic_routes_coexist() {
848 let harness = setup_test_harness_with(EndpointType::Unprivileged, |b| {
849 b.with_handler(SimpleHandler {
850 path: "/static",
851 body: "static",
852 })
853 })
854 .await;
855
856 let body = assert_status_eventually(harness.addr, "/static", StatusCode::OK).await;
858 assert_eq!(body, "static");
859
860 let dynamic_route = DynamicRoute::http(
862 EndpointType::Unprivileged,
863 SimpleHandler {
864 path: "/dynamic",
865 body: "dynamic",
866 },
867 );
868 harness.assert_route("dyn", dynamic_route).await;
869
870 let body = assert_status_eventually(harness.addr, "/dynamic", StatusCode::OK).await;
871 assert_eq!(body, "dynamic");
872
873 let (status, body) = http_get(harness.addr, "/static").await;
874 assert_eq!(status, StatusCode::OK);
875 assert_eq!(body, "static");
876
877 harness.retract_route("dyn").await;
879 assert_status_eventually(harness.addr, "/dynamic", StatusCode::NOT_FOUND).await;
880
881 let (status, body) = http_get(harness.addr, "/static").await;
882 assert_eq!(status, StatusCode::OK);
883 assert_eq!(body, "static");
884 }
885
886 #[tokio::test]
887 async fn unknown_grpc_method_returns_unimplemented() {
888 let harness = setup_test_harness(EndpointType::Unprivileged).await;
889 let (status, headers) = grpc_post(harness.addr, "/some.Service/Method").await;
890
891 assert_eq!(status, StatusCode::OK);
893 let grpc_status = headers.get("grpc-status").and_then(|v| v.to_str().ok());
894 assert_eq!(grpc_status, Some("12"));
895 }
896
897 #[tokio::test]
898 async fn static_route_wins_overlap_with_dynamic() {
899 let harness = setup_test_harness_with(EndpointType::Unprivileged, |b| {
900 b.with_handler(SimpleHandler {
901 path: "/overlap",
902 body: "static",
903 })
904 })
905 .await;
906
907 let body = assert_status_eventually(harness.addr, "/overlap", StatusCode::OK).await;
909 assert_eq!(body, "static");
910
911 let dynamic_route = DynamicRoute::http(
913 EndpointType::Unprivileged,
914 SimpleHandler {
915 path: "/overlap",
916 body: "dynamic",
917 },
918 );
919 harness.assert_route("dyn-overlap", dynamic_route).await;
920
921 sleep(Duration::from_millis(200)).await;
922
923 let (status, body) = http_get(harness.addr, "/overlap").await;
924 assert_eq!(status, StatusCode::OK);
925 assert_eq!(body, "static");
926 }
927}