Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 62 additions & 8 deletions crates/client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,12 @@ pub enum ConnectionState {
Connected {
path: Option<tcode_protocol::PathInfo>,
},
Syncing,
/// The link is up and the baseline is being loaded over it. The path is
/// the same as `Connected` carries and a path change is a new `Syncing`
/// value; only entering `Syncing` from another state restarts a sync.
Syncing {
path: Option<tcode_protocol::PathInfo>,
},
Reconnecting {
attempt: u32,
reason: Option<ConnectionFailure>,
Expand Down Expand Up @@ -427,7 +432,7 @@ impl HostLink {
fn flush_outbox(&self) {
if !matches!(
self.connection_state(),
ConnectionState::Connected { .. } | ConnectionState::Syncing
ConnectionState::Connected { .. } | ConnectionState::Syncing { .. }
) {
return;
}
Expand Down Expand Up @@ -456,7 +461,7 @@ impl HostLink {
fn check_deadlines(&self, now: web_time::Instant) {
if !matches!(
self.connection_state(),
ConnectionState::Connected { .. } | ConnectionState::Syncing
ConnectionState::Connected { .. } | ConnectionState::Syncing { .. }
) {
return;
}
Expand Down Expand Up @@ -548,7 +553,7 @@ impl HostLink {
if state_guard.as_ref().is_some_and(|state| {
!matches!(
**state,
ConnectionState::Connected { .. } | ConnectionState::Syncing
ConnectionState::Connected { .. } | ConnectionState::Syncing { .. }
)
}) {
return Err(error("disconnected", "Read requires a connection"));
Expand Down Expand Up @@ -806,14 +811,16 @@ impl HostLink {
let _ = self.send_subscription(subscription);
}
}
if state == ConnectionState::Syncing {
if matches!(state, ConnectionState::Syncing { .. })
&& !matches!(previous, ConnectionState::Syncing { .. })
{
for subscription in self.subscriptions() {
let _ = self.send_subscription(subscription);
}
}
if matches!(
state,
ConnectionState::Syncing | ConnectionState::Connected { .. }
ConnectionState::Syncing { .. } | ConnectionState::Connected { .. }
) {
self.flush_outbox();
}
Expand Down Expand Up @@ -1037,7 +1044,7 @@ mod tests {
let (incoming, from_host) = async_channel::unbounded();
let recreated = HostLink::new(to_host, from_host);
recreated.restore_outbox(storage.clone()).unwrap();
recreated.set_connection_state(ConnectionState::Syncing);
recreated.set_connection_state(ConnectionState::Syncing { path: None });
let mut pump = std::pin::pin!(recreated.pump_with_timer(std::future::pending::<()>));
// Every restored write is on the wire before the first Ack arrives.
let requests: Vec<_> = entries
Expand Down Expand Up @@ -1190,7 +1197,7 @@ mod tests {
.iter()
.all(|write| write.sent.is_none())
);
link.set_connection_state(ConnectionState::Syncing);
link.set_connection_state(ConnectionState::Syncing { path: None });
let resent: Vec<_> = (0..2).map(|_| request(&outgoing)).collect();
assert!(outgoing.try_recv().is_err());
for (original, again) in first[1..].iter().zip(&resent) {
Expand Down Expand Up @@ -1285,6 +1292,53 @@ mod tests {
assert!(link.inner.pending.lock().unwrap().is_empty());
}

/// A path change while the baseline loads is a new `Syncing` value for
/// the screen, not a lost link: nothing is resubscribed, so the reply
/// already in flight stays current.
#[test]
fn a_path_change_while_syncing_does_not_restart_the_sync() {
let (to_host, outgoing) = async_channel::unbounded();
let (_incoming, from_host) = async_channel::unbounded();
let link = HostLink::new(to_host, from_host);
let subscription = Subscription {
topic: Topic::SessionEvents {
session_id: "one".into(),
},
after: None,
};
link.subscribe(subscription.clone()).unwrap();
outgoing.try_recv().unwrap();
link.set_connection_state(ConnectionState::Syncing { path: None });
let request = tcode_protocol::decode_client_line(&outgoing.try_recv().unwrap()).unwrap();
assert!(outgoing.try_recv().is_err());
let lan = tcode_protocol::PathInfo {
direct: true,
relay: None,
lan: true,
probing_direct: false,
};
link.set_connection_state(ConnectionState::Syncing {
path: Some(lan.clone()),
});
assert!(outgoing.try_recv().is_err(), "the path alone resubscribed");
assert_eq!(
link.connection_state(),
ConnectionState::Syncing { path: Some(lan) }
);
let reply = EventEnvelope {
request_id: Some(request.id),
topic: subscription.topic,
event: ServerEvent::SessionSnapshot {
total: 0,
total_turns: 0,
truncated: false,
from: 0,
records: vec![],
},
};
assert!(link.subscription_reply_is_current(&reply));
}

#[test]
fn retired_subscription_rejects_queued_reply_and_releases_generation() {
let (to_host, outgoing) = async_channel::unbounded();
Expand Down
5 changes: 5 additions & 0 deletions crates/traverse-server/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,11 @@ relay. `tls.mode = "off"` also turns address discovery off.
public `https://` URL), or for development. QUIC address discovery is
unavailable in this mode.

The manifest URL carries the `tls.bind` port unless it is 443. With `manual`
or `reloading` behind a proxy that owns 443 and forwards TLS to `tls.bind`
(so the instance still holds the certificate, and QUIC address discovery
works), set `tls.public_port = 443`.

## Manifest

`GET /relays.json` describes the instance to clients:
Expand Down
8 changes: 8 additions & 0 deletions crates/traverse-server/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,9 @@ cert = ""
key = ""
# ACME directory URL; unset means Let's Encrypt production.
# acme_directory = "https://acme-staging-v02.api.letsencrypt.org/directory"
# Port in the advertised https://<hostname> URL; unset means the tls.bind
# port. Set 443 when a reverse proxy on 443 forwards TLS to tls.bind.
# public_port = 443

[relay]
# QUIC address discovery lets clients learn their public address. It needs
Expand Down Expand Up @@ -146,6 +149,10 @@ pub struct TlsConfig {
/// ACME directory other than Let's Encrypt production (staging, pebble).
#[serde(default)]
pub acme_directory: Option<String>,
/// Port advertised in the manifest URLs when it differs from `bind`, for
/// a TLS listener behind a proxy that forwards a public port to it.
#[serde(default)]
pub public_port: Option<u16>,
}

#[derive(Debug, Clone, PartialEq, Deserialize)]
Expand Down Expand Up @@ -275,6 +282,7 @@ impl Default for TlsConfig {
cert: String::new(),
key: String::new(),
acme_directory: None,
public_port: None,
}
}
}
Expand Down
7 changes: 6 additions & 1 deletion crates/traverse-server/src/manifest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ pub fn public_base(config: &Config, http_addr: SocketAddr) -> Url {
match config.hostname.as_deref().filter(|host| !host.is_empty()) {
Some(hostname) => {
let port = if config.tls_enabled() {
config.tls.bind.port()
config.tls.public_port.unwrap_or(config.tls.bind.port())
} else {
443
};
Expand Down Expand Up @@ -122,5 +122,10 @@ mod tests {
compose(&config, &base, "x".into()).relays[0].quic_port,
Some(7842)
);

// A proxy on 443 forwarding to tls.bind: the manifest names the public port.
config.tls.public_port = Some(443);
let base = public_base(&config, "127.0.0.1:8080".parse().unwrap());
assert_eq!(base.as_str(), "https://h.example/");
}
}
116 changes: 84 additions & 32 deletions crates/traverse/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -580,55 +580,90 @@ struct Established {
reader: LineReader,
}

/// The state channel. `Connected` carries the selected path, so a path change
/// is published as a new `Connected` value; before the first host line the
/// path is only remembered, since repeating `Syncing` would restart a sync.
/// The state channel. `Syncing` and `Connected` carry the selected path, so
/// a path change is published as a new value of whichever the link is in;
/// a repeated `Syncing` with only its path changed does not restart a sync.
struct StateSender {
tx: Sender<ConnectionState>,
path: Mutex<Option<PathInfo>>,
connected: std::sync::atomic::AtomicBool,
link: Mutex<Link>,
}

#[derive(Clone)]
struct Link {
phase: Phase,
path: Option<PathInfo>,
}

#[derive(Clone, Copy, PartialEq, Eq)]
enum Phase {
Down,
Syncing,
Connected,
}

impl StateSender {
fn new(tx: Sender<ConnectionState>) -> Self {
Self {
tx,
path: Mutex::new(None),
connected: std::sync::atomic::AtomicBool::new(false),
link: Mutex::new(Link {
phase: Phase::Down,
path: None,
}),
}
}

/// Any state but `Connected`: the path belongs to the connection that
/// `Reconnecting` or `Offline`: the path belongs to the connection that
/// was up, and the next one reports its own.
async fn send(
&self,
state: ConnectionState,
) -> Result<(), async_channel::SendError<ConnectionState>> {
debug_assert!(!state.is_connected(), "connected() carries the path");
self.connected
.store(false, std::sync::atomic::Ordering::Relaxed);
*self.path.lock().unwrap() = None;
debug_assert!(
!matches!(
state,
ConnectionState::Connected { .. } | ConnectionState::Syncing { .. }
),
"syncing() and connected() carry the path"
);
*self.link.lock().unwrap() = Link {
phase: Phase::Down,
path: None,
};
self.tx.send(state).await
}

async fn syncing(&self, path: PathInfo) {
*self.link.lock().unwrap() = Link {
phase: Phase::Syncing,
path: Some(path),
};
self.publish().await;
}

async fn connected(&self) {
self.connected
.store(true, std::sync::atomic::Ordering::Relaxed);
let path = self.path.lock().unwrap().clone();
let _ = self.tx.send(ConnectionState::Connected { path }).await;
self.link.lock().unwrap().phase = Phase::Connected;
self.publish().await;
}

async fn set_path(&self, path: PathInfo) {
{
let mut current = self.path.lock().unwrap();
if current.as_ref() == Some(&path) {
let mut link = self.link.lock().unwrap();
if link.phase == Phase::Down || link.path.as_ref() == Some(&path) {
return;
}
*current = Some(path);
}
if self.connected.load(std::sync::atomic::Ordering::Relaxed) {
self.connected().await;
link.path = Some(path);
}
self.publish().await;
}

async fn publish(&self) {
let link = self.link.lock().unwrap().clone();
let state = match link.phase {
Phase::Down => return,
Phase::Syncing => ConnectionState::Syncing { path: link.path },
Phase::Connected => ConnectionState::Connected { path: link.path },
};
let _ = self.tx.send(state).await;
}
}

Expand Down Expand Up @@ -712,7 +747,9 @@ async fn connection_loop(
persist_addresses(&device, &host);
// Tunnels are available by the time Syncing is observable.
tunnels.set(Some(established.connection.clone()));
let _ = state.send(ConnectionState::Syncing).await;
state
.syncing(crate::host::path_info(&established.connection))
.await;
let paths = tokio::spawn(watch_paths(
device.clone(),
established.connection.clone(),
Expand Down Expand Up @@ -1129,32 +1166,45 @@ mod tests {
}

/// The phone's direct path dying is a state change, not a side channel:
/// the link publishes a fresh `Connected` naming the relay, and only
/// while it is up — before the first host line the path waits for it.
/// the link publishes a fresh value of the state it is in naming the
/// relay — `Syncing` before the first host line, `Connected` after it —
/// and nothing while it is down.
#[tokio::test]
async fn a_path_change_is_published_as_a_new_connected_state() {
async fn a_path_change_is_published_as_a_new_value_of_the_current_state() {
let (tx, rx) = async_channel::unbounded();
let state = StateSender::new(tx);
state.send(ConnectionState::Syncing).await.unwrap();
state.set_path(direct()).await;
assert_eq!(rx.try_recv(), Ok(ConnectionState::Syncing));
assert!(rx.try_recv().is_err(), "syncing is not repeated");
state.connected().await;
state.syncing(direct()).await;
assert_eq!(
rx.try_recv(),
Ok(ConnectionState::Connected {
Ok(ConnectionState::Syncing {
path: Some(direct())
})
);
state.set_path(direct()).await;
assert!(rx.try_recv().is_err(), "an unchanged path is not repeated");
state.set_path(relay()).await;
assert_eq!(
rx.try_recv(),
Ok(ConnectionState::Syncing {
path: Some(relay())
})
);
state.connected().await;
assert_eq!(
rx.try_recv(),
Ok(ConnectionState::Connected {
path: Some(relay())
})
);
state.set_path(relay()).await;
assert!(rx.try_recv().is_err(), "an unchanged path is not repeated");
state.set_path(direct()).await;
assert_eq!(
rx.try_recv(),
Ok(ConnectionState::Connected {
path: Some(direct())
})
);
state
.send(ConnectionState::Reconnecting {
attempt: 1,
Expand All @@ -1163,6 +1213,8 @@ mod tests {
.await
.unwrap();
rx.try_recv().unwrap();
state.set_path(direct()).await;
assert!(rx.try_recv().is_err(), "a down link has no path");
state.connected().await;
assert_eq!(
rx.try_recv(),
Expand Down
4 changes: 2 additions & 2 deletions crates/traverse/tests/lan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -140,7 +140,7 @@ fn stale_addresses_without_a_browse_do_not_reach_the_machine() {
assert!(
!seen.iter().any(|state| matches!(
state,
ConnectionState::Syncing | ConnectionState::Connected { .. }
ConnectionState::Syncing { .. } | ConnectionState::Connected { .. }
)),
"the device connected without a working address: {seen:?}"
);
Expand Down Expand Up @@ -178,7 +178,7 @@ fn a_saved_address_further_down_the_list_reaches_the_machine() {
let client = tcode_traverse::connect(&saved, &device);
let (_, took) = wait_state(
&client,
|state| *state == ConnectionState::Syncing,
|state| matches!(state, ConnectionState::Syncing { .. }),
Duration::from_secs(15),
);
eprintln!("connected through a saved address in {took:?}");
Expand Down
Loading