use std::convert::Infallible;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use axum::extract::{DefaultBodyLimit, Query, State};
use axum::http::{HeaderMap, StatusCode};
use axum::response::sse::{Event, KeepAlive, Sse};
use axum::routing::{get, post};
use axum::{Json, Router};
use futures_util::stream::{self, Stream, StreamExt};
use secrecy::{ExposeSecret, SecretString};
use serde::{Deserialize, Serialize};
use subtle::ConstantTimeEq;
use time::OffsetDateTime;
use tokio::sync::{broadcast, watch};
use tokio_stream::wrappers::errors::BroadcastStreamRecvError;
use tokio_stream::wrappers::BroadcastStream;
use crate::chat::turns::ChatTurn;
use crate::receipt::Receipt;
use crate::wave::journal::{MessageOp, PendingMessage};
use crate::wave::playhead::PlayheadView;
use crate::wave::registry::{process_alive, StoreObserver};
use crate::wave::runtime::{InboxItem, TurnBroadcast, TurnDeltaFrame, TurnFrame, WaveRuntime};
use crate::wave::state::LoopState;
use crate::wave::supervisor::SupervisorHandle;
use crate::wave::wire::{
AttachRequest, AttachResponse, ContextResponse, InboxFrame, PostDeltasRequest,
PostDeltasResponse, RESIDENT_TOKEN_FILE, RESIDENT_TOKEN_HEADER,
};
pub const ENDPOINT_FILE: &str = ".wave-endpoint";
#[derive(Debug, Clone)]
pub struct ResidentDoor {
token: SecretString,
seat: Arc<Mutex<Option<u32>>>,
}
impl ResidentDoor {
pub fn new(token: impl Into<String>) -> Self {
Self {
token: SecretString::new(token.into()),
seat: Arc::new(Mutex::new(None)),
}
}
pub fn seat_pid(&self) -> Option<u32> {
*self.seat.lock().expect("resident seat lock poisoned")
}
pub fn record_pid(&self, pid: u32) {
*self.seat.lock().expect("resident seat lock poisoned") = Some(pid);
}
pub fn clear_seat(&self) {
*self.seat.lock().expect("resident seat lock poisoned") = None;
}
fn authorize(&self, headers: &HeaderMap) -> Result<(), (StatusCode, String)> {
let presented = headers
.get(RESIDENT_TOKEN_HEADER)
.and_then(|value| value.to_str().ok())
.unwrap_or_default();
if token_matches(&self.token, presented) {
return Ok(());
}
Err((
StatusCode::UNAUTHORIZED,
format!("missing or wrong {RESIDENT_TOKEN_HEADER}"),
))
}
}
fn token_matches(expected: &SecretString, provided: &str) -> bool {
expected
.expose_secret()
.as_bytes()
.ct_eq(provided.as_bytes())
.into()
}
pub fn generate_resident_token() -> String {
format!(
"{}{}",
uuid::Uuid::new_v4().simple(),
uuid::Uuid::new_v4().simple()
)
}
#[derive(Debug, Clone)]
pub struct ShutdownDoor {
requested: watch::Sender<bool>,
}
impl ShutdownDoor {
pub fn new() -> Self {
let (requested, _) = watch::channel(false);
Self { requested }
}
fn request(&self) {
self.requested.send_replace(true);
}
pub async fn wait(&self) {
let mut receiver = self.requested.subscribe();
if *receiver.borrow() {
return;
}
let _ = receiver.changed().await;
}
}
impl Default for ShutdownDoor {
fn default() -> Self {
Self::new()
}
}
#[derive(Debug, Serialize)]
struct HealthBody {
status: String,
loop_state: Option<String>,
wave: String,
turns: usize,
paused: bool,
uptime_seconds: i64,
}
#[derive(Debug, Serialize)]
struct ConversationBody {
turns: Vec<ChatTurn>,
}
#[derive(Debug, Deserialize)]
struct ConversationQuery {
limit: Option<usize>,
}
#[derive(Debug, Deserialize)]
struct PostMessage {
op: MessageOp,
text: String,
from: Option<String>,
}
#[derive(Debug, Deserialize)]
struct EventsQuery {
inbox: Option<bool>,
limit: Option<usize>,
}
pub(crate) const HUMAN_THREAD_REPLAY_LIMIT: usize = 12;
#[derive(Debug, Serialize)]
struct MemoryBody {
content: String,
}
#[derive(Debug, Serialize)]
struct MemoryLogBody {
facts: Vec<String>,
}
#[derive(Debug, Clone, Copy, Deserialize)]
#[serde(rename_all = "snake_case")]
enum MemoryOp {
Update,
Add,
}
#[derive(Debug, Deserialize)]
struct PostMemory {
op: MemoryOp,
content: String,
summary: Option<String>,
receipts: Vec<Receipt>,
}
#[derive(Debug, Serialize)]
struct PostMemoryResponse {
summary: String,
}
#[derive(Debug, Serialize)]
struct PostMessageResponse {
turn: Option<ChatTurn>,
state: String,
}
#[derive(Clone)]
struct ServerState {
runtime: Arc<WaveRuntime>,
resident: ResidentDoor,
observer: Option<Arc<StoreObserver>>,
supervisor: Option<SupervisorHandle>,
shutdown: ShutdownDoor,
started_at: OffsetDateTime,
}
const MAX_BODY_BYTES: usize = 1_048_576;
pub fn router(
runtime: Arc<WaveRuntime>,
resident: ResidentDoor,
observer: Option<Arc<StoreObserver>>,
supervisor: Option<SupervisorHandle>,
shutdown: ShutdownDoor,
) -> Router {
let state = ServerState {
runtime,
resident,
observer,
supervisor,
shutdown,
started_at: OffsetDateTime::now_utc(),
};
Router::new()
.route("/health", get(health_handler))
.route("/stop", post(stop_handler))
.route("/conversation", get(conversation_handler))
.route("/playhead", get(playhead_handler))
.route("/events", get(events_handler))
.route("/messages", post(messages_handler))
.route("/observations", post(observations_handler))
.route("/memory", get(memory_handler).post(memory_write_handler))
.route("/memory/log", get(memory_log_handler))
.route("/resident/attach", post(resident_attach_handler))
.route("/resident/deltas", post(resident_deltas_handler))
.route("/resident/context", get(resident_context_handler))
.layer(DefaultBodyLimit::max(MAX_BODY_BYTES))
.with_state(state)
}
async fn observations_handler(
State(state): State<ServerState>,
) -> Result<StatusCode, (StatusCode, String)> {
let observer = state.observer.as_ref().ok_or_else(|| {
(
StatusCode::SERVICE_UNAVAILABLE,
"child observations require the shared Loopflow registry".to_string(),
)
})?;
observer.poll_once().await;
Ok(StatusCode::NO_CONTENT)
}
async fn stop_handler(State(state): State<ServerState>) -> StatusCode {
state.shutdown.request();
StatusCode::ACCEPTED
}
async fn playhead_handler(
State(state): State<ServerState>,
) -> Result<Json<PlayheadView>, (StatusCode, String)> {
state
.runtime
.ensure_playhead()
.map(Json)
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))
}
async fn health_handler(State(state): State<ServerState>) -> Json<HealthBody> {
let loop_state = state
.runtime
.resident_expected()
.then(|| state.runtime.loop_state().name().to_string());
Json(HealthBody {
status: "serving".to_string(),
loop_state,
wave: state.runtime.name().to_string(),
turns: state.runtime.thread_len(),
paused: state.runtime.paused(),
uptime_seconds: (OffsetDateTime::now_utc() - state.started_at).whole_seconds(),
})
}
async fn resident_attach_handler(
State(state): State<ServerState>,
headers: HeaderMap,
Json(body): Json<AttachRequest>,
) -> Result<Json<AttachResponse>, (StatusCode, String)> {
state.resident.authorize(&headers)?;
if let Some(seated) = state.resident.seat_pid() {
if seated != body.pid && process_alive(seated).await {
return Err((
StatusCode::CONFLICT,
format!(
"wave '{}' already has a live resident on the seat (pid {seated}); \
stop it before attaching, or use `lf wave <name> --force` to take over",
state.runtime.name()
),
));
}
}
state.resident.record_pid(body.pid);
state.runtime.set_resident_expected();
if let Some(supervisor) = &state.supervisor {
supervisor.on_attach(body.pid);
}
if matches!(state.runtime.loop_state(), LoopState::Failed { .. }) {
state
.runtime
.transition(LoopState::Idle, "resident attached");
}
state
.runtime
.ensure_playhead()
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
tracing::info!(pid = body.pid, "resident attached");
Ok(Json(AttachResponse {
wave: state.runtime.name().to_string(),
}))
}
async fn resident_deltas_handler(
State(state): State<ServerState>,
headers: HeaderMap,
Json(body): Json<PostDeltasRequest>,
) -> Result<Json<PostDeltasResponse>, (StatusCode, String)> {
state.resident.authorize(&headers)?;
let accepted = body.deltas.len() as u64;
for delta in body.deltas {
state.runtime.apply_resident_delta(delta);
}
Ok(Json(PostDeltasResponse { accepted }))
}
async fn resident_context_handler(
State(state): State<ServerState>,
headers: HeaderMap,
) -> Result<Json<ContextResponse>, (StatusCode, String)> {
state.resident.authorize(&headers)?;
if let Some(observer) = &state.observer {
observer.poll_once().await;
}
let playhead = state
.runtime
.ensure_playhead()
.map_err(|err| (StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
Ok(Json(ContextResponse {
playhead,
provider_session: state.runtime.latest_provider_session(),
}))
}
async fn conversation_handler(
State(state): State<ServerState>,
Query(query): Query<ConversationQuery>,
) -> Json<ConversationBody> {
Json(ConversationBody {
turns: state.runtime.thread_tail(query.limit),
})
}
async fn messages_handler(
State(state): State<ServerState>,
Json(body): Json<PostMessage>,
) -> Result<Json<PostMessageResponse>, (StatusCode, String)> {
if matches!(body.op, MessageOp::Say) {
return Err((
StatusCode::BAD_REQUEST,
"`say` is not a wire op: machine speech rides the bus (`lf radio pub`)".to_string(),
));
}
if body.from.is_some() {
return Err((
StatusCode::BAD_REQUEST,
"the thread is unattributed: bylines belong to the bus (`lf radio pub --from`)"
.to_string(),
));
}
if body.text.trim().is_empty() && !matches!(body.op, MessageOp::Interrupt) {
return Err((
StatusCode::BAD_REQUEST,
"text is required for every op but interrupt".to_string(),
));
}
let turn = state.runtime.deliver(body.op, body.text);
Ok(Json(PostMessageResponse {
turn,
state: state.runtime.loop_state().name().to_string(),
}))
}
async fn memory_handler(State(state): State<ServerState>) -> Json<MemoryBody> {
Json(MemoryBody {
content: state.runtime.memory().read(),
})
}
async fn memory_log_handler(State(state): State<ServerState>) -> Json<MemoryLogBody> {
Json(MemoryLogBody {
facts: state.runtime.memory_adds(),
})
}
async fn memory_write_handler(
State(state): State<ServerState>,
Json(body): Json<PostMemory>,
) -> Result<Json<PostMemoryResponse>, (StatusCode, String)> {
let summary = body
.summary
.filter(|s| !s.trim().is_empty())
.or_else(|| first_line(&body.content))
.unwrap_or_else(|| "memory cleared".to_string());
let result = match body.op {
MemoryOp::Update => state.runtime.update_memory(&body.content, &summary),
MemoryOp::Add => {
let fact = body.content.trim();
if fact.is_empty() {
return Err((
StatusCode::BAD_REQUEST,
"content is required for the add op".to_string(),
));
}
state.runtime.append_memory(fact, body.receipts.clone())
}
};
match result {
Ok(()) => Ok(Json(PostMemoryResponse { summary })),
Err(err) => Err((
StatusCode::INTERNAL_SERVER_ERROR,
format!("memory write failed: {err}"),
)),
}
}
fn first_line(content: &str) -> Option<String> {
content
.lines()
.map(str::trim)
.find(|line| !line.is_empty())
.map(str::to_string)
}
async fn events_handler(
State(state): State<ServerState>,
Query(query): Query<EventsQuery>,
) -> axum::response::Response {
let shutdown = state.shutdown.clone();
let include_inbox = query.inbox == Some(true);
let replay_limit = if include_inbox {
query.limit
} else {
Some(query.limit.unwrap_or(HUMAN_THREAD_REPLAY_LIMIT))
};
let sub = state.runtime.subscribe_with_snapshot(replay_limit);
let inbox_replay: Vec<Result<Event, Infallible>> = if include_inbox {
sub.pending
.iter()
.map(|message| {
let frame = if let Some(observation) = sub.tasks.get(&message.id) {
InboxFrame::Task {
observation: observation.clone(),
}
} else if let Some(observation) = sub.projects.get(&message.id) {
InboxFrame::Project {
observation: observation.clone(),
}
} else {
pending_inbox_frame(message)
};
Ok(inbox_event(&frame))
})
.collect()
} else {
Vec::new()
};
let replay = stream::iter(
std::iter::once(Ok(state_event(&sub.state)))
.chain(sub.playhead.into_iter().map(|p| Ok(playhead_event(&p))))
.chain(sub.turns.into_iter().map(|t| Ok(turn_event(&t))))
.chain(
sub.memory_adds
.into_iter()
.map(|fact| Ok(memory_add_event(&fact))),
)
.chain(inbox_replay),
);
let live_turns = turn_event_stream(sub.turn_rx);
let live_states = live_stream(sub.state_rx, |s| state_event(&s));
let live_playhead = live_stream(sub.playhead_rx, |p| playhead_event(&p));
let live_memory_adds = live_stream(sub.memory_add_rx, |fact| memory_add_event(&fact));
let live_memory = live_stream(sub.memory_rx, |summary| memory_event(&summary));
let mut live: BoxedEventStream = Box::pin(stream::select(
live_turns,
stream::select(
stream::select(live_states, live_playhead),
stream::select(live_memory, live_memory_adds),
),
));
if include_inbox {
let live_inbox = live_stream(sub.inbox_rx, |item| inbox_event(&inbox_item_frame(&item)));
live = Box::pin(stream::select(live, live_inbox));
}
let merged: BoxedEventStream = Box::pin(replay.chain(live).take_until(async move {
shutdown.wait().await;
}));
axum::response::IntoResponse::into_response(Sse::new(merged).keep_alive(KeepAlive::default()))
}
type BoxedEventStream =
std::pin::Pin<Box<dyn Stream<Item = Result<Event, Infallible>> + Send + 'static>>;
fn live_stream<T, F>(
rx: broadcast::Receiver<T>,
to_event: F,
) -> impl Stream<Item = Result<Event, Infallible>> + Send + 'static
where
T: Clone + Send + 'static,
F: Fn(T) -> Event + Send + 'static,
{
BroadcastStream::new(rx).filter_map(move |res| {
let out = res.ok().map(|value| Ok(to_event(value)));
async move { out }
})
}
fn turn_event_stream(
rx: broadcast::Receiver<TurnBroadcast>,
) -> impl Stream<Item = Result<Event, Infallible>> + Send + 'static {
stream::unfold(Some(BroadcastStream::new(rx)), |state| async move {
let mut inner = state?;
let recv = inner.next().await?;
match turn_wire_frame(recv) {
TurnWireFrame::Whole(frame) => Some((
Ok(Event::default().event("turn").data(frame.json.as_str())),
Some(inner),
)),
TurnWireFrame::Delta(frame) => Some((
Ok(Event::default()
.event("turn-delta")
.data(frame.json.as_str())),
Some(inner),
)),
TurnWireFrame::Resync => {
Some((Ok(Event::default().event("resync").data("reconnect")), None))
}
}
})
}
enum TurnWireFrame {
Whole(Arc<TurnFrame>),
Delta(Arc<TurnDeltaFrame>),
Resync,
}
fn turn_wire_frame(recv: Result<TurnBroadcast, BroadcastStreamRecvError>) -> TurnWireFrame {
match recv {
Ok(TurnBroadcast::Whole(frame)) => TurnWireFrame::Whole(frame),
Ok(TurnBroadcast::Delta(frame)) => TurnWireFrame::Delta(frame),
Err(BroadcastStreamRecvError::Lagged(_)) => TurnWireFrame::Resync,
}
}
fn turn_event(turn: &ChatTurn) -> Event {
Event::default()
.event("turn")
.data(serde_json::to_string(turn).expect("ChatTurn serializes to JSON"))
}
fn playhead_event(playhead: &PlayheadView) -> Event {
Event::default()
.event("playhead")
.data(serde_json::to_string(playhead).expect("PlayheadView serializes to JSON"))
}
fn state_event(state: &LoopState) -> Event {
Event::default().event("state").data(state.name())
}
fn memory_event(summary: &str) -> Event {
Event::default().event("memory").data(summary)
}
fn memory_add_event(fact: &str) -> Event {
Event::default().event("memory-add").data(fact)
}
fn inbox_event(frame: &InboxFrame) -> Event {
Event::default()
.event("inbox")
.data(serde_json::to_string(frame).expect("InboxFrame serializes to JSON"))
}
fn pending_inbox_frame(message: &PendingMessage) -> InboxFrame {
InboxFrame::Message {
id: message.id.0.clone(),
op: message.op,
text: message.text.clone(),
from: message.from.clone(),
}
}
fn inbox_item_frame(item: &InboxItem) -> InboxFrame {
match item {
InboxItem::Message(message) => pending_inbox_frame(message),
InboxItem::Task(observation) => InboxFrame::Task {
observation: observation.clone(),
},
InboxItem::Project(observation) => InboxFrame::Project {
observation: observation.clone(),
},
InboxItem::Interrupt => InboxFrame::Interrupt,
InboxItem::Skip => InboxFrame::Skip,
}
}
pub fn endpoint_path(repo_root: &Path, wave: &str) -> PathBuf {
repo_root.join("wave").join(wave).join(ENDPOINT_FILE)
}
const ENDPOINT_PROBE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
pub async fn live_endpoint(repo_root: &Path, wave: &str) -> Option<String> {
let addr = std::fs::read_to_string(endpoint_path(repo_root, wave)).ok()?;
let addr = addr.trim().to_string();
if addr.is_empty() {
return None;
}
let client = reqwest::Client::builder()
.timeout(ENDPOINT_PROBE_TIMEOUT)
.build()
.ok()?;
let body: serde_json::Value = client
.get(format!("http://{addr}/health"))
.send()
.await
.ok()?
.json()
.await
.ok()?;
(body.get("wave").and_then(serde_json::Value::as_str) == Some(wave)).then_some(addr)
}
pub fn write_endpoint(
repo_root: &Path,
wave: &str,
addr: std::net::SocketAddr,
) -> std::io::Result<()> {
let path = endpoint_path(repo_root, wave);
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(path, addr.to_string())
}
pub fn remove_endpoint(repo_root: &Path, wave: &str, own_addr: &str) {
let path = endpoint_path(repo_root, wave);
match std::fs::read_to_string(&path) {
Ok(contents) if contents.trim() == own_addr => {
let _ = std::fs::remove_file(path);
}
_ => {}
}
}
pub fn resident_token_path(repo_root: &Path, wave: &str) -> PathBuf {
repo_root.join("wave").join(wave).join(RESIDENT_TOKEN_FILE)
}
pub fn write_resident_token(repo_root: &Path, wave: &str, token: &str) -> std::io::Result<()> {
let path = resident_token_path(repo_root, wave);
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(&path, token)?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600))?;
}
Ok(())
}
pub fn read_resident_token(repo_root: &Path, wave: &str) -> Option<String> {
let token = std::fs::read_to_string(resident_token_path(repo_root, wave)).ok()?;
let token = token.trim().to_string();
(!token.is_empty()).then_some(token)
}
pub fn remove_resident_token(repo_root: &Path, wave: &str, own_token: &str) {
let path = resident_token_path(repo_root, wave);
match std::fs::read_to_string(&path) {
Ok(contents) if contents.trim() == own_token => {
let _ = std::fs::remove_file(path);
}
_ => {}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::chat::turns::ChatTurn;
fn whole_broadcast(id: &str) -> TurnBroadcast {
let turn = ChatTurn::user(id.to_string(), "hi".to_string());
let json = serde_json::to_string(&turn).expect("serialize");
TurnBroadcast::Whole(Arc::new(TurnFrame { turn, json }))
}
#[test]
fn a_lagged_turn_broadcast_maps_to_resync_not_a_dropped_frame() {
assert!(matches!(
turn_wire_frame(Err(BroadcastStreamRecvError::Lagged(3))),
TurnWireFrame::Resync
));
assert!(matches!(
turn_wire_frame(Ok(whole_broadcast("turn-1"))),
TurnWireFrame::Whole(_)
));
}
#[tokio::test]
async fn turn_stream_signals_one_resync_then_ends_on_lag() {
let (tx, rx) = broadcast::channel::<TurnBroadcast>(2);
for i in 0..8 {
let _ = tx.send(whole_broadcast(&format!("turn-{i}")));
}
drop(tx);
let mut stream = Box::pin(turn_event_stream(rx));
assert!(
stream.next().await.is_some(),
"a lagged turn stream signals resync, not silence"
);
assert!(
stream.next().await.is_none(),
"the substream ends after resync so the client reconnects"
);
}
#[test]
fn write_and_remove_endpoint_roundtrips() {
let tmp = tempfile::tempdir().expect("tempdir");
let addr: std::net::SocketAddr = "127.0.0.1:54321".parse().unwrap();
write_endpoint(tmp.path(), "ship", addr).expect("write endpoint");
let path = endpoint_path(tmp.path(), "ship");
assert_eq!(std::fs::read_to_string(&path).unwrap(), "127.0.0.1:54321");
remove_endpoint(tmp.path(), "ship", "127.0.0.1:54321");
assert!(!path.exists());
}
#[test]
fn remove_endpoint_leaves_a_foreign_pointer_untouched() {
let tmp = tempfile::tempdir().expect("tempdir");
let addr: std::net::SocketAddr = "127.0.0.1:50000".parse().unwrap();
write_endpoint(tmp.path(), "ship", addr).expect("write endpoint");
remove_endpoint(tmp.path(), "ship", "127.0.0.1:50001");
let path = endpoint_path(tmp.path(), "ship");
assert_eq!(
std::fs::read_to_string(&path).unwrap(),
"127.0.0.1:50000",
"foreign pointer survives our shutdown"
);
}
#[test]
fn resident_token_file_roundtrips_and_respects_ownership() {
let tmp = tempfile::tempdir().expect("tempdir");
assert!(read_resident_token(tmp.path(), "ship").is_none());
write_resident_token(tmp.path(), "ship", "tok-1").expect("write");
assert_eq!(
read_resident_token(tmp.path(), "ship").as_deref(),
Some("tok-1")
);
remove_resident_token(tmp.path(), "ship", "tok-other");
assert_eq!(
read_resident_token(tmp.path(), "ship").as_deref(),
Some("tok-1")
);
remove_resident_token(tmp.path(), "ship", "tok-1");
assert!(read_resident_token(tmp.path(), "ship").is_none());
}
#[tokio::test]
async fn stop_route_requests_listener_shutdown() {
let tmp = tempfile::tempdir().expect("tempdir");
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open");
let shutdown = ShutdownDoor::new();
let requested = shutdown.clone();
let app = router(runtime, ResidentDoor::new("resident"), None, None, shutdown);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
let response = reqwest::Client::new()
.post(format!("http://{addr}/stop"))
.send()
.await
.unwrap();
assert_eq!(response.status(), reqwest::StatusCode::ACCEPTED);
requested.wait().await;
server.abort();
}
#[tokio::test]
async fn observation_nudge_fails_loudly_without_the_shared_registry() {
let tmp = tempfile::tempdir().expect("tempdir");
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open");
let app = router(
runtime,
ResidentDoor::new("resident"),
None,
None,
ShutdownDoor::new(),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
let response = reqwest::Client::new()
.post(format!("http://{addr}/observations"))
.json(&serde_json::json!({}))
.send()
.await
.unwrap();
assert_eq!(response.status(), reqwest::StatusCode::SERVICE_UNAVAILABLE);
server.abort();
}
#[tokio::test]
async fn live_endpoint_is_none_for_a_stale_pointer() {
let tmp = tempfile::tempdir().expect("tempdir");
assert!(live_endpoint(tmp.path(), "ship").await.is_none(), "no file");
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let dead = listener.local_addr().unwrap();
drop(listener);
write_endpoint(tmp.path(), "ship", dead).expect("write endpoint");
assert!(
live_endpoint(tmp.path(), "ship").await.is_none(),
"dead address is stale"
);
}
}