diff --git a/constellation/src/server/mod.rs b/constellation/src/server/mod.rs index 9390056..5d37f07 100644 --- a/constellation/src/server/mod.rs +++ b/constellation/src/server/mod.rs @@ -14,7 +14,7 @@ use serde_with::serde_as; use std::collections::{HashMap, HashSet}; use std::time::{Duration, UNIX_EPOCH}; use tokio::net::{TcpListener, ToSocketAddrs}; -use tokio::task::block_in_place; +use tokio::task::spawn_blocking; use tokio_util::sync::CancellationToken; use crate::storage::{LinkReader, StorageStats}; @@ -34,6 +34,11 @@ fn get_default_cursor_limit() -> u64 { const INDEX_BEGAN_AT_TS: u64 = 1738083600; // TODO: not this +fn to500(e: tokio::task::JoinError) -> http::StatusCode { + eprintln!("handler join error: {e}"); + http::StatusCode::INTERNAL_SERVER_ERROR +} + pub async fn serve(store: S, addr: A, stay_alive: CancellationToken) -> anyhow::Result<()> where S: LinkReader, @@ -45,14 +50,22 @@ where "/", get({ let store = store.clone(); - move |accept| async { block_in_place(|| hello(accept, store)) } + move |accept| async { + spawn_blocking(|| hello(accept, store)) + .await + .map_err(to500)? + } }), ) .route( "/links/count", get({ let store = store.clone(); - move |accept, query| async { block_in_place(|| count_links(accept, query, store)) } + move |accept, query| async { + spawn_blocking(|| count_links(accept, query, store)) + .await + .map_err(to500)? + } }), ) .route( @@ -60,7 +73,9 @@ where get({ let store = store.clone(); move |accept, query| async { - block_in_place(|| count_distinct_dids(accept, query, store)) + spawn_blocking(|| count_distinct_dids(accept, query, store)) + .await + .map_err(to500)? } }), ) @@ -69,7 +84,9 @@ where get({ let store = store.clone(); move |accept, query| async { - block_in_place(|| get_backlinks(accept, query, store)) + spawn_blocking(|| get_backlinks(accept, query, store)) + .await + .map_err(to500)? } }), ) @@ -77,7 +94,11 @@ where "/links", get({ let store = store.clone(); - move |accept, query| async { block_in_place(|| get_links(accept, query, store)) } + move |accept, query| async { + spawn_blocking(|| get_links(accept, query, store)) + .await + .map_err(to500)? + } }), ) .route( @@ -85,7 +106,9 @@ where get({ let store = store.clone(); move |accept, query| async { - block_in_place(|| get_distinct_dids(accept, query, store)) + spawn_blocking(|| get_distinct_dids(accept, query, store)) + .await + .map_err(to500)? } }), ) @@ -95,7 +118,9 @@ where get({ let store = store.clone(); move |accept, query| async { - block_in_place(|| count_all_links(accept, query, store)) + spawn_blocking(|| count_all_links(accept, query, store)) + .await + .map_err(to500)? } }), ) @@ -104,7 +129,9 @@ where get({ let store = store.clone(); move |accept, query| async { - block_in_place(|| explore_links(accept, query, store)) + spawn_blocking(|| explore_links(accept, query, store)) + .await + .map_err(to500)? } }), )