pub mod chat;
pub(crate) mod discord;
pub mod journal;
pub(crate) mod memory;
pub mod playhead;
pub mod relocate;
pub(crate) mod registry;
pub mod resident;
pub mod runtime;
pub mod server;
pub mod state;
pub mod subscription;
pub(crate) mod supervisor;
mod types;
pub mod wire;
pub(crate) use types::PromotionWake;
pub use types::{Wave, WaveLocator, WaveLocatorError};
use std::future::Future;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use anyhow::{anyhow, Result};
use secrecy::SecretString;
use crate::engine::repo::find_repo_root;
use crate::engine::wave_config::{try_read_wave_config, WaveChatConfig};
use crate::engine::worktrees::main_repo_root;
use crate::ops::util::normalize_wave_name;
use crate::store::{open_existing_store, SharedStore};
use crate::wave::chat::ChatBacking;
use crate::wave::runtime::WaveRuntime;
pub(crate) const RESIDENT_SUBCOMMAND: &str = "__resident";
pub(crate) const WAVE_SERVER_ENDPOINT_ENV: &str = "LF_WAVE_SERVER_ENDPOINT";
#[derive(Debug)]
pub(crate) struct ListenerSignals<F> {
startup: Option<tokio::sync::oneshot::Sender<String>>,
shutdown: F,
}
impl<F> ListenerSignals<F> {
pub(crate) fn new(startup: Option<tokio::sync::oneshot::Sender<String>>, shutdown: F) -> Self {
Self { startup, shutdown }
}
}
pub fn run(name: &str, force: bool) -> Result<()> {
let repo_root = find_repo_root()?;
let main_repo = main_repo_root(&repo_root).unwrap_or_else(|_| repo_root.clone());
let wave = normalize_wave_name(name).ok_or_else(|| anyhow!("invalid wave name: '{name}'"))?;
let rt = tokio::runtime::Runtime::new()?;
rt.block_on(async {
let registry_config = resolve_registry(&main_repo, &wave).await;
run_listener(
main_repo,
wave,
registry_config,
force,
true,
None,
shutdown_signal(),
)
.await
})
}
pub fn stop(name: &str) -> Result<()> {
let repo_root = find_repo_root()?;
let main_repo = main_repo_root(&repo_root).unwrap_or_else(|_| repo_root.clone());
let wave = normalize_wave_name(name).ok_or_else(|| anyhow!("invalid wave name: '{name}'"))?;
let rt = tokio::runtime::Runtime::new()?;
let requested = rt.block_on(request_stop(&main_repo, &wave))?;
if requested {
println!("stopped wave {wave}");
} else {
println!("wave {wave} is already stopped");
}
Ok(())
}
pub(crate) async fn request_stop(repo_root: &Path, wave: &str) -> Result<bool> {
let Some(endpoint) = server::live_endpoint(repo_root, wave).await else {
return Ok(false);
};
let response = reqwest::Client::new()
.post(format!("http://{endpoint}/stop"))
.send()
.await
.map_err(|err| anyhow!("failed to stop wave '{wave}': {err}"))?;
if response.status() != reqwest::StatusCode::ACCEPTED {
return Err(anyhow!(
"wave '{wave}' refused stop with HTTP {}",
response.status()
));
}
for _ in 0..100 {
let recorded = std::fs::read_to_string(server::endpoint_path(repo_root, wave))
.ok()
.map(|value| value.trim().to_string());
if recorded.as_deref() != Some(endpoint.as_str()) {
return Ok(true);
}
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
if server::live_endpoint(repo_root, wave).await.is_none() {
return Ok(true);
}
Err(anyhow!(
"wave '{wave}' accepted stop but is still serving at http://{endpoint}"
))
}
async fn resolve_registry(main_repo: &Path, wave: &str) -> Option<registry::RegistryConfig> {
let Some(store) = open_existing_store().await else {
tracing::warn!(
wave,
"no local registry on this machine; running without child observations"
);
return None;
};
let store: SharedStore = Arc::new(store);
match registry::ensure_wave_row(&store, main_repo, wave).await {
Ok(row) => Some(registry::RegistryConfig { store, wave: row }),
Err(err) => {
tracing::warn!(wave, error = %err, "local registry unusable; running without child observations");
None
}
}
}
struct WaveRunGuard(Option<(SharedStore, crate::durable::RunLease)>);
impl WaveRunGuard {
fn as_ref(&self) -> Option<&(SharedStore, crate::durable::RunLease)> {
self.0.as_ref()
}
fn take(&mut self) -> Option<(SharedStore, crate::durable::RunLease)> {
self.0.take()
}
}
impl Drop for WaveRunGuard {
fn drop(&mut self) {
let Some((store, lease)) = self.0.take() else {
return;
};
let Ok(runtime) = tokio::runtime::Handle::try_current() else {
return;
};
runtime.spawn(async move {
if let Err(error) = store
.stop_run(
&lease,
crate::durable::StopCause::Requested,
crate::durable::ContainmentObservation::Absent,
)
.await
{
tracing::warn!(%error, "failed to roll back Wave Run authority");
}
});
}
}
fn resident_spawner(
lf_bin: PathBuf,
wave: String,
repo_root: PathBuf,
endpoint: String,
token: String,
resident_env: Vec<(String, String)>,
) -> supervisor::SpawnResident {
Box::new(move || {
let mut command =
resident_command(&lf_bin, &wave, &repo_root, &endpoint, &token, &resident_env);
#[cfg(unix)]
command.process_group(0);
command.spawn()
})
}
fn resident_command(
lf_bin: &Path,
wave: &str,
repo_root: &Path,
endpoint: &str,
token: &str,
resident_env: &[(String, String)],
) -> tokio::process::Command {
let mut command = tokio::process::Command::new(lf_bin);
command
.arg(RESIDENT_SUBCOMMAND)
.arg(wave)
.current_dir(repo_root)
.env(WAVE_SERVER_ENDPOINT_ENV, endpoint)
.env(wire::RESIDENT_TOKEN_ENV, token)
.env("PATH", crate::flowloop::wave::path_for_children())
.env_remove(crate::durable::RUN_CONTEXT_ENV)
.env_remove(crate::durable::RUN_LEASE_ENV)
.stdin(std::process::Stdio::null());
for (key, value) in resident_env {
command.env(key, value);
}
command.env_remove(discord::TOKEN_ENV);
command
}
pub(crate) async fn run_listener(
repo_root: PathBuf,
wave: String,
registry_config: Option<registry::RegistryConfig>,
force: bool,
spawn_resident: bool,
discord_token: Option<SecretString>,
shutdown: impl Future<Output = ()> + Send + 'static,
) -> Result<()> {
run_listener_with_startup(
repo_root,
wave,
registry_config,
force,
spawn_resident,
discord_token,
ListenerSignals::new(None, shutdown),
)
.await
}
pub(crate) async fn run_listener_with_startup<F>(
repo_root: PathBuf,
wave: String,
registry_config: Option<registry::RegistryConfig>,
force: bool,
spawn_resident: bool,
discord_token: Option<SecretString>,
signals: ListenerSignals<F>,
) -> Result<()>
where
F: Future<Output = ()> + Send + 'static,
{
let ListenerSignals { startup, shutdown } = signals;
let _locator_lock = match relocate::WaveLocatorLock::acquire(&repo_root, &wave) {
Ok(lock) => lock,
Err(lock_error) if force => {
if !request_stop(&repo_root, &wave).await? {
return Err(lock_error);
}
relocate::WaveLocatorLock::acquire(&repo_root, &wave)?
}
Err(lock_error) => return Err(lock_error),
};
if let Some(config) = registry_config.as_ref() {
let current = config
.store
.get_wave(config.wave.id())
.await?
.ok_or_else(|| {
anyhow!(
"Wave {} disappeared before listener start",
config.wave.id()
)
})?;
let locator = WaveLocator::discover(&repo_root, &wave)?;
if current.repo() != locator.repo().to_string() || current.name() != locator.slug() {
return Err(anyhow!(
"Wave {} moved to {}/{} before listener start",
current.id(),
current.repo(),
current.name()
));
}
}
if let Some(live) = server::live_endpoint(&repo_root, &wave).await {
if !force {
return Err(anyhow!(
"refusing to start: wave '{wave}' already has a live server at \
http://{live} (per wave/{wave}/{endpoint}). Stop that session (or let \
it finish), or rerun with --force to take over.",
endpoint = server::ENDPOINT_FILE,
));
}
tracing::warn!(
wave,
live,
"--force: taking over a live wave server endpoint"
);
}
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
let addr = listener.local_addr()?;
let resident_lf = spawn_resident
.then(crate::engine::process::resolve_current_home_lf_binary_checked)
.transpose()?;
let (chat_backing, discord_adapter) =
match try_read_wave_config(&repo_root, &wave)?.and_then(|config| config.chat) {
Some(WaveChatConfig::Discord {
home_id,
guild_id,
channel_id,
}) => {
let registry = registry_config.as_ref().ok_or_else(|| {
anyhow!("Discord chat requires the local registry to verify its owner Home")
})?;
let local_home = registry.store.local_home().await?;
let binding = journal::DiscordChatBinding {
guild_id,
channel_id,
};
let backing = ChatBacking::discord(&binding);
let adapter = discord::DiscordAdapter::preflight(
binding,
&home_id,
&local_home.id,
discord_token,
)
.await?;
(backing, Some(adapter))
}
Some(WaveChatConfig::Local) | None => (ChatBacking::Local, None),
};
let registered = registry_config
.as_ref()
.map(|config| (config.store.clone(), config.wave.id().clone()));
let mut wave_run = WaveRunGuard(if spawn_resident {
match registry_config.as_ref() {
Some(config) => {
let work = crate::durable::WorkRef::Wave(config.wave.id().clone());
let (_, lease) = config
.store
.reserve_run(&work, crate::durable::RunTrigger::User)
.await?;
Some((config.store.clone(), lease))
}
None => None,
}
} else {
None
});
let mut resident_env = registry_config
.as_ref()
.map(|config| {
vec![(
crate::engine::wave_context::WAVE_ID_ENV.to_string(),
config.wave.id().to_string(),
)]
})
.unwrap_or_default();
if let Some((_, lease)) = wave_run.as_ref() {
resident_env.extend([
(
crate::durable::RUN_CONTEXT_ENV.to_string(),
"agent".to_string(),
),
(
crate::durable::RUN_LEASE_ENV.to_string(),
lease.env_value().to_string(),
),
]);
}
let runtime = WaveRuntime::open_with_backing(wave.clone(), repo_root.clone(), chat_backing)?;
if let Some(adapter) = discord_adapter.as_ref() {
adapter.attach(&runtime)?;
}
runtime.journal_server_started(std::process::id(), &addr.to_string());
let observer = registered.map(|(store, wave_id)| {
Arc::new(registry::StoreObserver::new(
runtime.clone(),
store,
wave_id,
))
});
let observer = Arc::new(registry::ObserverSlot::new(runtime.clone(), observer));
let observer_task = tokio::spawn(Arc::clone(&observer).run(registry::POLL_CADENCE));
let discord_projection = discord_adapter
.as_ref()
.map(discord::DiscordAdapter::projection);
let discord_task = discord_adapter.map(|adapter| tokio::spawn(adapter.run(runtime.clone())));
let token = server::generate_resident_token();
let door = server::ResidentDoor::new(token.clone());
server::write_resident_token(&repo_root, &wave, &token)?;
let spawner = resident_lf.map(|lf_bin| {
resident_spawner(
lf_bin,
wave.clone(),
repo_root.clone(),
addr.to_string(),
token.clone(),
resident_env,
)
});
let supervisor = supervisor::Supervisor::new(
runtime.clone(),
door.clone(),
spawner,
supervisor::SupervisorConfig::default(),
);
let supervisor_handle = supervisor.handle();
let supervisor_task = tokio::spawn(supervisor.run());
server::write_endpoint(&repo_root, &wave, addr)?;
if let Some(startup) = startup {
let _ = startup.send(addr.to_string());
}
let own_addr = addr.to_string();
let cleanup_repo = repo_root.clone();
let cleanup_wave = wave.clone();
let cleanup_addr = own_addr.clone();
let cleanup_token = token.clone();
let cleanup_door = door.clone();
let cleanup_run = wave_run
.as_ref()
.map(|(store, lease)| (Arc::clone(store), lease.clone()));
crate::engine::agent::register_interrupt_cleanup(move || {
if let Some(pid) = cleanup_door.seat_pid() {
supervisor::terminate_resident_blocking(pid);
}
if let Some((store, lease)) = cleanup_run.as_ref() {
if let Err(error) = store.stop_run_on_interrupt(lease) {
tracing::warn!(%error, "failed to release Wave Run authority on interrupt");
}
}
server::remove_endpoint(&cleanup_repo, &cleanup_wave, &cleanup_addr);
server::remove_resident_token(&cleanup_repo, &cleanup_wave, &cleanup_token);
});
println!(
"lf wave · {wave} · listener on http://{addr}{} \
(Ctrl-C to stop, RUST_LOG=loopflow=debug for the firehose)",
if spawn_resident {
" · spawning resident"
} else {
" · no resident"
}
);
let shutdown_door = server::ShutdownDoor::new();
let shutdown_request = shutdown_door.clone();
let graceful_shutdown = async move {
tokio::select! {
_ = shutdown => {}
_ = shutdown_request.wait() => {}
}
};
let app = server::router_with_chat_projection(
runtime.clone(),
door.clone(),
observer,
Some(supervisor_handle),
shutdown_door,
discord_projection,
);
let result = axum::serve(listener, app)
.with_graceful_shutdown(graceful_shutdown)
.await;
if let Some(task) = discord_task {
task.abort();
}
supervisor_task.abort();
if let Some(pid) = door.seat_pid() {
supervisor::terminate_resident(pid).await;
}
observer_task.abort();
if let Some((store, lease)) = wave_run.take() {
if let Err(error) = store
.stop_run(
&lease,
crate::durable::StopCause::Requested,
crate::durable::ContainmentObservation::Absent,
)
.await
{
tracing::warn!(%error, "failed to release Wave Run authority");
}
}
server::remove_endpoint(&repo_root, &wave, &own_addr);
server::remove_resident_token(&repo_root, &wave, &token);
result.map_err(|err| anyhow!("wave server error: {err}"))
}
async fn shutdown_signal() {
let _ = tokio::signal::ctrl_c().await;
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use crate::chat::turns::{ChatRole, ChatTurn, TurnDelta};
use crate::chat::types::Lifecycle;
use crate::wave::chat::{ChatHistorySnapshot, PostMessageResponse, WaveChatMessage};
use crate::wave::journal::MessageOp;
use crate::wave::server::ResidentDoor;
use crate::wave::wire::{ResidentDelta, RESIDENT_TOKEN_HEADER};
#[test]
fn the_listener_spawns_its_resident_body_by_name() {
use crate::lf::{Cli, Commands};
use clap::Parser;
let cli = Cli::try_parse_from(["lf", RESIDENT_SUBCOMMAND, "goals"])
.expect("the spawner's subcommand must be one the CLI accepts");
assert!(matches!(
cli.command,
Some(Commands::Resident { name }) if name == "goals"
));
assert!(matches!(
Cli::try_parse_from(["lf", "loop", "goals"])
.expect("external fallback")
.command,
Some(Commands::External(parts)) if parts[0] == "loop"
));
}
#[test]
fn discord_chat_token_is_explicitly_scrubbed_from_the_resident() {
let command = resident_command(
Path::new("/paired/bin/lf"),
"goals",
Path::new("/tmp"),
"127.0.0.1:1234",
"resident-token",
&[(discord::TOKEN_ENV.to_string(), "must-not-pass".to_string())],
);
assert_eq!(command.as_std().get_program(), "/paired/bin/lf");
assert!(command
.as_std()
.get_envs()
.any(|(name, value)| { name == discord::TOKEN_ENV && value.is_none() }));
}
#[tokio::test]
async fn stopping_an_absent_wave_is_idempotent() {
let tmp = tempfile::tempdir().expect("tempdir");
assert!(!request_stop(tmp.path(), "ship").await.expect("stop"));
}
fn progress_turn(text: &str) -> ChatTurn {
ChatTurn {
id: String::new(),
role: ChatRole::Assistant,
text: text.to_string(),
status: Lifecycle::Completed,
items: Vec::new(),
created_at: "1970-01-01T00:00:00Z".to_string(),
body: None,
activity: None,
}
}
fn narrate(runtime: &WaveRuntime, text: &str) {
runtime.append_finalized_turn(progress_turn(text), Vec::new());
}
async fn boot() -> (String, Arc<WaveRuntime>, tempfile::TempDir) {
let tmp = tempfile::tempdir().expect("tempdir");
let dir = tmp.path().join("wave/ship");
std::fs::create_dir_all(&dir).expect("dir");
std::fs::write(dir.join("MEMORY.md"), "Goal: ship the reactive server.\n").expect("mem");
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open");
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let app = server::router_with_observer(
runtime.clone(),
ResidentDoor::new("test-token"),
Arc::new(registry::ObserverSlot::new(runtime.clone(), None)),
None,
server::ShutdownDoor::new(),
);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
(format!("http://{addr}"), runtime, tmp)
}
async fn wait_for<F: Fn() -> bool>(cond: F) {
for _ in 0..200 {
if cond() {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!("condition not met in time");
}
#[tokio::test]
async fn finalized_turn_appears_in_conversation() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "Implemented the reactive server.");
let body = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.text()
.await
.unwrap();
assert!(body.contains("Implemented the reactive server."));
assert!(body.contains("\"role\":\"assistant\""));
}
#[tokio::test]
async fn conversation_limit_tails_the_thread() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "first");
narrate(&runtime, "second");
narrate(&runtime, "third");
let body: ChatHistorySnapshot = reqwest::get(format!("{base}/conversation?limit=2"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns: Vec<_> = body.messages.iter().map(|message| &message.turn).collect();
assert_eq!(turns.len(), 2);
assert_eq!(turns[0].text, "second");
assert_eq!(turns[1].text, "third");
for url in [
format!("{base}/conversation?limit=99"),
format!("{base}/conversation"),
] {
let body: ChatHistorySnapshot = reqwest::get(url).await.unwrap().json().await.unwrap();
assert_eq!(body.messages.len(), 3);
}
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "in progress".into(),
});
let body: ChatHistorySnapshot = reqwest::get(format!("{base}/conversation?limit=1"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns: Vec<_> = body.messages.iter().map(|message| &message.turn).collect();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].status, Lifecycle::Running);
assert_eq!(turns[0].text, "in progress");
}
#[tokio::test]
async fn posted_message_appears_as_user_turn() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
let body: PostMessageResponse = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "message", "text": "how's it going?" }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
let posted = body.message.expect("posted message").turn;
assert_eq!(posted.role, ChatRole::User);
assert_eq!(posted.text, "how's it going?");
assert_eq!(body.state, "idle");
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].role, ChatRole::User);
}
#[tokio::test]
async fn the_thread_door_refuses_machine_speech() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
for body in [
serde_json::json!({
"op": "say",
"text": "run-7 landed: PR #12",
"from": "worker",
}),
serde_json::json!({ "op": "say", "text": "anon" }),
serde_json::json!({ "op": "message", "text": "hello", "from": "cli" }),
] {
let response = client
.post(format!("{base}/messages"))
.json(&body)
.send()
.await
.unwrap();
assert!(response.status().is_client_error());
}
assert!(runtime.thread_snapshot().is_empty(), "nothing journaled");
for body in [
serde_json::json!({ "op": "message", "text": "hello" }),
serde_json::json!({ "op": "steer", "text": "change course" }),
serde_json::json!({ "op": "interrupt", "text": "stop here" }),
] {
let response = client
.post(format!("{base}/messages"))
.json(&body)
.send()
.await
.unwrap();
assert!(response.status().is_success());
}
}
#[tokio::test]
async fn post_without_op_or_without_text_is_rejected() {
let (base, _runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
let no_op = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "text": "hello" }))
.send()
.await
.unwrap();
assert_eq!(no_op.status(), reqwest::StatusCode::UNPROCESSABLE_ENTITY);
for op in ["message", "steer"] {
let empty = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": op, "text": " " }))
.send()
.await
.unwrap();
assert_eq!(
empty.status(),
reqwest::StatusCode::BAD_REQUEST,
"empty text rejected for {op}"
);
}
}
#[tokio::test]
async fn bare_interrupt_while_idle_is_a_noop_success_with_state() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
let body: PostMessageResponse = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "interrupt", "text": "" }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert!(body.message.is_none(), "nothing said, nothing appended");
assert_eq!(body.state, "idle");
assert!(runtime.thread_snapshot().is_empty());
}
#[tokio::test]
async fn health_splits_listener_liveness_from_the_loop() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "first");
let body: serde_json::Value = reqwest::get(format!("{base}/health"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["status"], "serving", "status is listener liveness");
assert!(body["loop_state"].is_null(), "dormant listener has no loop");
assert_eq!(body["wave"], "ship");
assert_eq!(body["turns"], 1);
runtime.set_resident_expected();
let body: serde_json::Value = reqwest::get(format!("{base}/health"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["loop_state"], "idle", "loop is the resident's state");
runtime.transition(
crate::wave::state::LoopState::Failed {
reason: "vendor gone".into(),
},
"test",
);
let body: serde_json::Value = reqwest::get(format!("{base}/health"))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(body["status"], "serving");
assert_eq!(body["loop_state"], "failed");
}
#[tokio::test]
async fn resident_door_gates_attaches_and_applies_deltas() {
let (base, runtime, _tmp) = boot().await;
let client = reqwest::Client::new();
for request in [
client
.post(format!("{base}/resident/attach"))
.json(&serde_json::json!({ "pid": 1234 })),
client
.post(format!("{base}/resident/attach"))
.header(RESIDENT_TOKEN_HEADER, "wrong")
.json(&serde_json::json!({ "pid": 1234 })),
] {
let denied = request.send().await.unwrap();
assert_eq!(denied.status(), reqwest::StatusCode::UNAUTHORIZED);
}
assert!(!runtime.resident_expected());
runtime.set_resident_expected();
runtime.transition(
crate::wave::state::LoopState::Failed {
reason: "old resident died".into(),
},
"test",
);
let attach: serde_json::Value = client
.post(format!("{base}/resident/attach"))
.header(RESIDENT_TOKEN_HEADER, "test-token")
.json(&serde_json::json!({ "pid": std::process::id() }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(attach["wave"], "ship");
assert_eq!(runtime.loop_state().name(), "idle", "attach revives");
let deltas: serde_json::Value = client
.post(format!("{base}/resident/deltas"))
.header(RESIDENT_TOKEN_HEADER, "test-token")
.json(&serde_json::json!({ "deltas": [
{ "kind": "turn_opened", "answers": [] },
{ "kind": "turn_text", "text": "over the wire" },
{ "kind": "turn_usage", "input_tokens": 7, "output_tokens": 3, "cache_read_tokens": null },
{ "kind": "turn_finished", "status": "completed", "cost_usd": 0.01 },
] }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(deltas["accepted"], 4);
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].text, "over the wire");
assert_eq!(thread[0].status, Lifecycle::Completed);
let context: serde_json::Value = client
.get(format!("{base}/resident/context"))
.header(RESIDENT_TOKEN_HEADER, "test-token")
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert!(context.get("playhead").is_some());
assert!(context.get("provider_session").is_some());
assert!(context.get("in_flight").is_none());
}
#[tokio::test]
async fn sse_replays_on_connect_then_streams_live() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "replayed turn");
let host = base.strip_prefix("http://").unwrap().to_string();
let mut stream = tokio::net::TcpStream::connect(&host).await.unwrap();
stream
.write_all(b"GET /events HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n")
.await
.unwrap();
narrate(&runtime, "live turn");
let mut acc = String::new();
let mut buf = [0u8; 2048];
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
loop {
let read = tokio::time::timeout_at(deadline, stream.read(&mut buf)).await;
match read {
Ok(Ok(0)) | Err(_) => break,
Ok(Ok(n)) => {
acc.push_str(&String::from_utf8_lossy(&buf[..n]));
if acc.contains("replayed turn") && acc.contains("live turn") {
break;
}
}
Ok(Err(_)) => break,
}
}
assert!(
acc.contains("event: message"),
"human SSE frames are named `message`"
);
assert!(
acc.contains("replayed turn"),
"replays the thread on connect"
);
assert!(
acc.contains("live turn"),
"streams turns narrated after connect"
);
assert!(
!acc.contains("event: inbox"),
"the default stream carries no inbox frames"
);
}
#[tokio::test]
async fn events_inbox_scope_replays_pending_and_streams_ops() {
let (base, runtime, _tmp) = boot().await;
runtime
.deliver(MessageOp::Message, "queued before".into())
.expect("user turn");
let host = base.strip_prefix("http://").unwrap().to_string();
let mut stream = tokio::net::TcpStream::connect(&host).await.unwrap();
stream
.write_all(
b"GET /events?inbox=true HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n",
)
.await
.unwrap();
let mut acc = String::new();
let mut buf = [0u8; 4096];
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut interrupt_sent = false;
loop {
if acc.contains("queued before") && !interrupt_sent {
interrupt_sent = true;
runtime.deliver_interrupt();
}
let read = tokio::time::timeout_at(deadline, stream.read(&mut buf)).await;
match read {
Ok(Ok(0)) | Err(_) => break,
Ok(Ok(n)) => {
acc.push_str(&String::from_utf8_lossy(&buf[..n]));
if acc.contains("queued before") && acc.contains(r#""kind":"interrupt""#) {
break;
}
}
Ok(Err(_)) => break,
}
}
assert!(acc.contains("event: inbox"), "inbox frames are named");
assert!(
acc.contains("queued before") && acc.contains(r#""op":"message""#),
"the pending queue replays: {acc}"
);
assert!(
acc.contains(r#""kind":"interrupt""#),
"a live bare interrupt rides as a tagged control frame: {acc}"
);
}
struct SseClient {
stream: tokio::net::TcpStream,
raw: Vec<u8>,
}
impl SseClient {
async fn connect(base: &str) -> Self {
let host = base.strip_prefix("http://").unwrap();
let mut stream = tokio::net::TcpStream::connect(host).await.unwrap();
stream
.write_all(
b"GET /events HTTP/1.1\r\nHost: localhost\r\nConnection: keep-alive\r\n\r\n",
)
.await
.unwrap();
Self {
stream,
raw: Vec::new(),
}
}
async fn frames_until(&mut self, pred: impl Fn(&[ChatTurn]) -> bool) -> Vec<ChatTurn> {
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut buf = [0u8; 4096];
loop {
let frames = parse_message_frames(&dechunk(&self.raw));
if pred(&frames) {
return frames;
}
match tokio::time::timeout_at(deadline, self.stream.read(&mut buf)).await {
Ok(Ok(0)) | Err(_) => {
panic!("SSE ended before condition; {} frames so far", frames.len())
}
Ok(Ok(n)) => self.raw.extend_from_slice(&buf[..n]),
Ok(Err(err)) => panic!("SSE read error: {err}"),
}
}
}
async fn states_until(&mut self, pred: impl Fn(&[String]) -> bool) -> Vec<String> {
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut buf = [0u8; 4096];
loop {
let states = parse_state_frames(&dechunk(&self.raw));
if pred(&states) {
return states;
}
match tokio::time::timeout_at(deadline, self.stream.read(&mut buf)).await {
Ok(Ok(0)) | Err(_) => {
panic!("SSE ended before condition; states so far: {states:?}")
}
Ok(Ok(n)) => self.raw.extend_from_slice(&buf[..n]),
Ok(Err(err)) => panic!("SSE read error: {err}"),
}
}
}
}
fn dechunk(raw: &[u8]) -> String {
let text = String::from_utf8_lossy(raw);
let Some(head_end) = text.find("\r\n\r\n") else {
return String::new();
};
let mut body = &text[head_end + 4..];
let mut out = String::new();
while let Some(size_end) = body.find("\r\n") {
let Ok(size) = usize::from_str_radix(body[..size_end].trim(), 16) else {
break;
};
let start = size_end + 2;
if size == 0 || body.len() < start + size {
break;
}
out.push_str(&body[start..start + size]);
body = &body[(start + size + 2).min(body.len())..];
}
out
}
fn parse_message_frames(sse_body: &str) -> Vec<ChatTurn> {
let mut order: Vec<String> = Vec::new();
let mut by_id: std::collections::HashMap<String, ChatTurn> =
std::collections::HashMap::new();
let mut event = String::new();
for line in sse_body.lines() {
if let Some(name) = line.strip_prefix("event:") {
event = name.trim().to_string();
} else if let Some(data) = line.strip_prefix("data:") {
let data = data.trim();
match event.as_str() {
"message" => {
if let Ok(message) = serde_json::from_str::<WaveChatMessage>(data) {
let turn = message.turn;
if !by_id.contains_key(&turn.id) {
order.push(turn.id.clone());
}
by_id.insert(turn.id.clone(), turn);
}
}
"message-delta" => {
if let Ok(delta) = serde_json::from_str::<TurnDelta>(data) {
if let Some(turn) = by_id.get_mut(&delta.turn_id) {
turn.absorb_item(delta.item);
}
}
}
_ => {}
}
}
}
order
.into_iter()
.map(|id| by_id.remove(&id).expect("id tracked"))
.collect()
}
fn parse_state_frames(sse_body: &str) -> Vec<String> {
let mut states = Vec::new();
let mut in_state_event = false;
for line in sse_body.lines() {
if let Some(name) = line.strip_prefix("event:") {
in_state_event = name.trim() == "state";
} else if let Some(data) = line.strip_prefix("data:") {
if in_state_event {
states.push(data.trim().to_string());
}
}
}
states
}
#[tokio::test]
async fn sse_late_subscriber_watches_the_open_turn_grow_and_finalize() {
let (base, runtime, _tmp) = boot().await;
narrate(&runtime, "already finalized");
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "thinking".into(),
});
let mut client = SseClient::connect(&base).await;
let frames = client
.frames_until(|f| f.iter().any(|t| t.status == Lifecycle::Running))
.await;
assert!(
frames
.iter()
.any(|t| t.text == "already finalized" && t.status == Lifecycle::Completed),
"replay carries the finalized thread"
);
let open = frames
.iter()
.find(|t| t.status == Lifecycle::Running && t.text == "thinking")
.unwrap()
.clone();
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "more".into(),
});
client
.frames_until(|f| {
f.iter().any(|t| {
t.id == open.id && t.text == "thinkingmore" && t.status == Lifecycle::Running
})
})
.await;
runtime.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
reason: None,
});
let frames = client
.frames_until(|f| {
f.iter()
.any(|t| t.id == open.id && t.status == Lifecycle::Completed)
})
.await;
let last = frames.iter().rfind(|t| t.id == open.id).unwrap();
assert_eq!(last.status, Lifecycle::Completed, "terminal frame is last");
assert_eq!(last.text, "thinkingmore");
}
#[tokio::test]
async fn sse_state_events_track_the_loop_live() {
let (base, runtime, _tmp) = boot().await;
let mut client = SseClient::connect(&base).await;
let states = client.states_until(|s| !s.is_empty()).await;
assert_eq!(states, vec!["idle"]);
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
let states = client.states_until(|s| s.len() >= 2).await;
assert_eq!(states, vec!["idle", "turning"]);
runtime.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
reason: None,
});
let states = client.states_until(|s| s.len() >= 3).await;
assert_eq!(states, vec!["idle", "turning", "idle"]);
}
#[tokio::test]
async fn conversation_includes_the_open_running_turn() {
let (base, runtime, _tmp) = boot().await;
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "half a thought".into(),
});
let body: ChatHistorySnapshot = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns: Vec<_> = body.messages.iter().map(|message| &message.turn).collect();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].status, Lifecycle::Running);
assert_eq!(turns[0].text, "half a thought");
runtime.apply_resident_delta(ResidentDelta::TurnFinished {
status: Lifecycle::Completed,
cost_usd: None,
reason: None,
});
let body: ChatHistorySnapshot = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns: Vec<_> = body.messages.iter().map(|message| &message.turn).collect();
assert_eq!(turns.len(), 1);
assert_eq!(turns[0].status, Lifecycle::Completed);
}
#[tokio::test]
async fn restart_mid_turn_never_serves_a_stale_running_turn() {
let tmp = tempfile::tempdir().expect("tempdir");
{
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("open");
runtime.apply_resident_delta(ResidentDelta::TurnOpened {
answers: Vec::new(),
});
runtime.apply_resident_delta(ResidentDelta::TurnText {
text: "half a thought".into(),
});
}
let runtime = WaveRuntime::open("ship".into(), tmp.path().to_path_buf()).expect("reopen");
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let app = server::router_with_observer(
runtime.clone(),
ResidentDoor::new("test-token"),
Arc::new(registry::ObserverSlot::new(runtime.clone(), None)),
None,
server::ShutdownDoor::new(),
);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
let base = format!("http://{addr}");
let body: ChatHistorySnapshot = reqwest::get(format!("{base}/conversation"))
.await
.unwrap()
.json()
.await
.unwrap();
let turns: Vec<_> = body.messages.iter().map(|message| &message.turn).collect();
assert_eq!(turns.len(), 1);
assert_eq!(
turns[0].status,
Lifecycle::Failed,
"janitor closed the turn"
);
assert_eq!(turns[0].text, "half a thought");
let mut client = SseClient::connect(&base).await;
let frames = client
.frames_until(|f| f.iter().any(|t| t.status == Lifecycle::Failed))
.await;
assert!(
frames.iter().all(|t| t.status != Lifecycle::Running),
"no stale running turn in replay"
);
}
async fn boot_family() -> (String, Arc<WaveRuntime>, tempfile::TempDir) {
let tmp = tempfile::tempdir().expect("tempdir");
let origin = tmp.path().join("repo");
std::fs::create_dir_all(origin.join("wave/ship")).expect("wave dir");
let runtime = WaveRuntime::open("ship".into(), origin).expect("open runtime");
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let app = server::router_with_observer(
runtime.clone(),
ResidentDoor::new("test-token"),
Arc::new(registry::ObserverSlot::new(runtime.clone(), None)),
None,
server::ShutdownDoor::new(),
);
tokio::spawn(async move {
axum::serve(listener, app).await.ok();
});
(format!("http://{addr}"), runtime, tmp)
}
#[tokio::test]
async fn a_message_is_recorded_once_in_the_waves_journal() {
let (base, runtime, tmp) = boot_family().await;
let client = reqwest::Client::new();
let body: PostMessageResponse = client
.post(format!("{base}/messages"))
.json(&serde_json::json!({ "op": "message", "text": "landed the parser" }))
.send()
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(
body.message.expect("posted message").turn.text,
"landed the parser"
);
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
let wave_journal = journal::journal_path(runtime.repo_root(), "ship");
assert_eq!(
loopflow_test_support::journal_files_under(tmp.path()),
vec![wave_journal.clone()]
);
assert_eq!(
journal::read_events(&wave_journal)
.iter()
.filter(|e| matches!(e.kind, journal::EventKind::UserMessage { .. }))
.count(),
1,
"one copy, in the wave's journal",
);
}
#[tokio::test]
async fn same_named_waves_in_two_repositories_serve_independently() {
let tmp = tempfile::tempdir().expect("tempdir");
let repo_a = tmp.path().join("alpha");
let repo_b = tmp.path().join("beta");
for repo in [&repo_a, &repo_b] {
std::fs::create_dir_all(repo.join("wave/infrastructure")).unwrap();
std::fs::write(
repo.join("wave/infrastructure/GOAL.md"),
format!("# {}\n", repo.file_name().unwrap().to_string_lossy()),
)
.unwrap();
}
let store = Arc::new(
crate::store::open_store(&crate::store::StorageConfig::sqlite(
tmp.path().join("registry.db"),
))
.await
.unwrap(),
);
let wave_a = registry::ensure_wave_row(&store, &repo_a, "infrastructure")
.await
.unwrap();
let wave_b = registry::ensure_wave_row(&store, &repo_b, "infrastructure")
.await
.unwrap();
assert_ne!(wave_a.id(), wave_b.id());
let (stop_a_tx, stop_a_rx) = tokio::sync::oneshot::channel::<()>();
let (stop_b_tx, stop_b_rx) = tokio::sync::oneshot::channel::<()>();
let handle_a = tokio::spawn(run_listener(
repo_a.clone(),
"infrastructure".to_string(),
Some(registry::RegistryConfig {
store: store.clone(),
wave: wave_a.clone(),
}),
false,
false,
None,
async move {
let _ = stop_a_rx.await;
},
));
let handle_b = tokio::spawn(run_listener(
repo_b.clone(),
"infrastructure".to_string(),
Some(registry::RegistryConfig {
store,
wave: wave_b.clone(),
}),
false,
false,
None,
async move {
let _ = stop_b_rx.await;
},
));
let endpoint_a = server::endpoint_path(&repo_a, "infrastructure");
let endpoint_b = server::endpoint_path(&repo_b, "infrastructure");
wait_for(|| endpoint_a.exists() && endpoint_b.exists()).await;
let address_a = std::fs::read_to_string(&endpoint_a).unwrap();
let address_b = std::fs::read_to_string(&endpoint_b).unwrap();
assert_ne!(address_a, address_b);
for address in [&address_a, &address_b] {
let health: serde_json::Value =
reqwest::get(format!("http://{}/health", address.trim()))
.await
.unwrap()
.json()
.await
.unwrap();
assert_eq!(health["wave"], "infrastructure");
}
let response = reqwest::Client::new()
.post(format!("http://{}/messages", address_a.trim()))
.json(&serde_json::json!({ "op": "message", "text": "alpha only" }))
.send()
.await
.unwrap();
assert!(response.status().is_success());
assert_eq!(
journal::read_events(&journal::journal_path(&repo_a, "infrastructure"))
.iter()
.filter(|event| matches!(event.kind, journal::EventKind::UserMessage { .. }))
.count(),
1
);
assert_eq!(
journal::read_events(&journal::journal_path(&repo_b, "infrastructure"))
.iter()
.filter(|event| matches!(event.kind, journal::EventKind::UserMessage { .. }))
.count(),
0
);
stop_a_tx.send(()).unwrap();
stop_b_tx.send(()).unwrap();
handle_a.await.unwrap().unwrap();
handle_b.await.unwrap().unwrap();
}
#[tokio::test]
async fn serve_publishes_and_removes_discovery_pointer_and_token() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let repo2 = repo.clone();
let handle = tokio::spawn(async move {
run_listener(repo2, "ship".into(), None, false, false, None, async {
let _ = shutdown_rx.await;
})
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
let contents = std::fs::read_to_string(&endpoint).unwrap();
assert!(
contents.starts_with("127.0.0.1:"),
"pointer is just an address"
);
assert!(
server::read_resident_token(&repo, "ship").is_some(),
"the resident token publishes beside the pointer"
);
shutdown_tx.send(()).unwrap();
handle.await.unwrap().unwrap();
assert!(!endpoint.exists(), "pointer removed on shutdown");
assert!(
server::read_resident_token(&repo, "ship").is_none(),
"token removed on shutdown"
);
}
#[tokio::test]
async fn stop_command_path_gracefully_shuts_down_the_listener() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let repo2 = repo.clone();
let handle = tokio::spawn(async move {
run_listener(
repo2,
"ship".into(),
None,
false,
false,
None,
std::future::pending(),
)
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
assert!(request_stop(&repo, "ship").await.expect("request stop"));
handle.await.unwrap().unwrap();
assert!(!endpoint.exists(), "stop removes the endpoint");
assert!(
server::read_resident_token(&repo, "ship").is_none(),
"stop removes the resident token"
);
}
#[tokio::test]
async fn stop_closes_a_live_event_subscription() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let repo2 = repo.clone();
let mut handle = tokio::spawn(async move {
run_listener(
repo2,
"ship".into(),
None,
false,
false,
None,
std::future::pending(),
)
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
let addr = std::fs::read_to_string(&endpoint).unwrap();
let client = reqwest::Client::new();
let events = client
.get(format!("http://{}/events?inbox=true", addr.trim()))
.send()
.await
.expect("subscribe");
assert_eq!(events.status(), reqwest::StatusCode::OK);
assert!(request_stop(&repo, "ship").await.expect("request stop"));
match tokio::time::timeout(Duration::from_secs(2), &mut handle).await {
Ok(result) => result.unwrap().unwrap(),
Err(_) => {
handle.abort();
panic!("listener stayed alive behind its event subscription");
}
}
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] async fn resolve_registry_without_a_store_runs_unregistered() {
let _env = crate::journal::TestLedgerGuard::new();
let tmp = tempfile::tempdir().expect("tempdir");
let config = resolve_registry(tmp.path(), "ship").await;
assert!(config.is_none(), "missing store boots unregistered");
}
#[tokio::test]
async fn serve_overwrites_a_stale_endpoint_file() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let dead = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let dead_addr = dead.local_addr().unwrap();
drop(dead);
server::write_endpoint(&repo, "ship", dead_addr).expect("stale pointer");
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let repo2 = repo.clone();
let handle = tokio::spawn(async move {
run_listener(repo2, "ship".into(), None, false, false, None, async {
let _ = shutdown_rx.await;
})
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| {
std::fs::read_to_string(&endpoint)
.is_ok_and(|contents| contents != dead_addr.to_string())
})
.await;
shutdown_tx.send(()).unwrap();
handle.await.unwrap().unwrap();
}
#[tokio::test]
async fn serve_refuses_to_start_over_a_live_endpoint_file() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let repo2 = repo.clone();
let first = tokio::spawn(async move {
run_listener(repo2, "ship".into(), None, false, false, None, async {
let _ = shutdown_rx.await;
})
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
let first_addr = std::fs::read_to_string(&endpoint).unwrap();
let (_shutdown_tx2, shutdown_rx2) = tokio::sync::oneshot::channel::<()>();
let err = run_listener(
repo.clone(),
"ship".into(),
None,
false,
false,
None,
async {
let _ = shutdown_rx2.await;
},
)
.await
.expect_err("live endpoint refuses a second server");
assert!(
err.to_string().contains("--force"),
"error points at --force: {err}"
);
assert_eq!(
std::fs::read_to_string(&endpoint).unwrap(),
first_addr,
"refused boot never touched the pointer"
);
shutdown_tx.send(()).unwrap();
first.await.unwrap().unwrap();
assert!(!endpoint.exists(), "first server still owns its shutdown");
}
#[tokio::test]
async fn force_stops_the_current_listener_before_taking_its_locator_lock() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/ship")).unwrap();
let repo = tmp.path().to_path_buf();
let (_first_stop_tx, first_stop_rx) = tokio::sync::oneshot::channel::<()>();
let first_repo = repo.clone();
let first = tokio::spawn(async move {
run_listener(first_repo, "ship".into(), None, false, false, None, async {
let _ = first_stop_rx.await;
})
.await
});
let endpoint = server::endpoint_path(&repo, "ship");
wait_for(|| endpoint.exists()).await;
let first_address = std::fs::read_to_string(&endpoint).unwrap();
let (second_stop_tx, second_stop_rx) = tokio::sync::oneshot::channel::<()>();
let second_repo = repo.clone();
let second = tokio::spawn(async move {
run_listener(second_repo, "ship".into(), None, true, false, None, async {
let _ = second_stop_rx.await;
})
.await
});
wait_for(|| {
std::fs::read_to_string(&endpoint).is_ok_and(|address| address != first_address)
})
.await;
first.await.unwrap().unwrap();
second_stop_tx.send(()).unwrap();
second.await.unwrap().unwrap();
assert!(!endpoint.exists());
}
}