commit 4be0ea979861c587146a42ed619b2551c4e131f7
parent 03dc98ef04beb28c526f317735d00ee5cba5cd8b
Author: triesap <tyson@radroots.org>
Date: Tue, 11 Aug 2026 02:24:02 +0000
runtime: harden websocket liveness and nostr signatures
Diffstat:
6 files changed, 217 insertions(+), 43 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -1353,7 +1353,6 @@ dependencies = [
name = "tangle_test_support"
version = "0.1.0"
dependencies = [
- "k256",
"pocket-types",
"serde_json",
"tangle_crypto",
diff --git a/crates/tangle_crypto/src/lib.rs b/crates/tangle_crypto/src/lib.rs
@@ -3,7 +3,7 @@
use core::fmt;
use std::sync::Arc;
-use k256::schnorr::signature::{Signer, Verifier};
+use k256::schnorr::signature::hazmat::{PrehashSigner, PrehashVerifier};
use k256::schnorr::{Signature, SigningKey, VerifyingKey};
use pocket_types::{
Kind as PocketKind, OwnedEvent as PocketOwnedEvent, Tags as PocketTags, Time as PocketTime,
@@ -60,7 +60,7 @@ pub fn verify_event_signature(event: &Event) -> Result<(), String> {
let signature = Signature::try_from(signature.as_slice())
.map_err(|_| "event signature is not a valid schnorr signature".to_owned())?;
verifying_key
- .verify(&event_id, &signature)
+ .verify_prehash(&event_id, &signature)
.map_err(|_| "event signature verification failed".to_owned())
}
@@ -89,7 +89,7 @@ pub fn verify_event_signature_bytes(
let signature = Signature::try_from(signature.as_slice())
.map_err(|_| "event signature is not a valid schnorr signature".to_owned())?;
verifying_key
- .verify(event_id, &signature)
+ .verify_prehash(event_id, &signature)
.map_err(|_| "event signature verification failed".to_owned())
}
@@ -125,7 +125,10 @@ impl RelaySigner {
let event_id = compute_event_id(&unsigned);
let event_id_bytes =
fixed_hex_bytes(event_id.as_str(), 32, "event id").expect("event id is valid hex");
- let signature: Signature = self.signing_key.sign(&event_id_bytes);
+ let signature: Signature = self
+ .signing_key
+ .sign_prehash(&event_id_bytes)
+ .expect("validated signing key signs a 32-byte event id");
let signature = SignatureHex::new(&lower_hex(signature.to_bytes().as_ref()))
.expect("schnorr signature emits valid hex");
Event::new(event_id, unsigned, signature)
@@ -255,7 +258,7 @@ mod tests {
RelaySigner, VerificationService, compute_event_id, compute_event_id_hex, event_id_matches,
fixed_hex_bytes, lower_hex, verify_event_id, verify_event_signature,
};
- use k256::schnorr::signature::Signer;
+ use k256::schnorr::signature::hazmat::PrehashSigner;
use k256::schnorr::{Signature, SigningKey};
use pocket_types::{Kind as PocketKind, OwnedTags as PocketOwnedTags, Time as PocketTime};
use std::time::Duration;
@@ -329,9 +332,20 @@ mod tests {
);
let event = signer.sign_unsigned_event(unsigned);
+ let pocket_tags = PocketOwnedTags::new(&[["t", "radroots"]]).expect("pocket tags");
+ let pocket = signer
+ .sign_pocket_event(
+ PocketKind::from_u16(1),
+ &pocket_tags,
+ PocketTime::from_u64(1_714_124_433),
+ b"relay generated",
+ )
+ .expect("pocket event");
assert_eq!(event.unsigned().pubkey(), signer.public_key());
assert_eq!(verify_event_signature(&event), Ok(()));
+ pocket.verify().expect("Pocket accepts protocol signature");
+ assert_eq!(event.id().as_str(), pocket.id().as_hex_string());
assert_eq!(
format!("{signer:?}"),
format!(
@@ -525,7 +539,9 @@ mod tests {
);
let event_id = compute_event_id(&unsigned);
let event_id_bytes = fixed_hex_bytes(event_id.as_str(), 32, "event id").expect("event id");
- let signature: Signature = signing_key.sign(&event_id_bytes);
+ let signature: Signature = signing_key
+ .sign_prehash(&event_id_bytes)
+ .expect("event signature");
let signature = SignatureHex::new(&lower_hex(signature.to_bytes().as_ref())).expect("sig");
Event::new(event_id, unsigned, signature)
}
diff --git a/crates/tangle_runtime/Cargo.toml b/crates/tangle_runtime/Cargo.toml
@@ -18,7 +18,7 @@ tangle_crypto = { path = "../tangle_crypto" }
tangle_groups = { path = "../tangle_groups" }
tangle_protocol = { path = "../tangle_protocol" }
tangle_store_pocket = { path = "../tangle_store_pocket" }
-tokio = { version = "1", features = ["net", "sync"] }
+tokio = { version = "1", features = ["net", "sync", "time"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] }
diff --git a/crates/tangle_runtime/src/session.rs b/crates/tangle_runtime/src/session.rs
@@ -20,13 +20,17 @@ use crate::{
use axum::extract::ws::{CloseFrame, Message, Utf8Bytes, WebSocket};
use std::{
collections::BTreeMap,
+ future::pending,
net::IpAddr,
sync::atomic::{AtomicU64, Ordering},
- time::{Instant, SystemTime, UNIX_EPOCH},
+ time::{Duration, Instant, SystemTime, UNIX_EPOCH},
};
use tangle_protocol::{RelayMessage, SubscriptionId, UnixTimestamp};
use tangle_store_pocket::PocketOwnedFilter;
-use tokio::sync::{mpsc, watch};
+use tokio::{
+ sync::{mpsc, watch},
+ time::{MissedTickBehavior, interval_at, timeout},
+};
#[derive(Debug)]
pub struct TangleWebSocketSession {
@@ -44,6 +48,7 @@ pub struct TangleWebSocketSession {
subscription_permits: BTreeMap<SubscriptionId, RelaySubscriptionPermit>,
events: TangleEventReceiver,
projection: RelayProjectionContext,
+ keepalive: Option<TangleWebSocketKeepaliveConfig>,
}
static NEXT_TANGLE_CONNECTION_ID: AtomicU64 = AtomicU64::new(1);
@@ -53,6 +58,55 @@ pub struct TangleWebSocketSessionOptions {
peer_ip: Option<IpAddr>,
resource_limiter: Option<RelayResourceLimiter>,
projection: RelayProjectionContext,
+ keepalive: Option<TangleWebSocketKeepaliveConfig>,
+}
+
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub struct TangleWebSocketKeepaliveConfig {
+ ping_interval: Duration,
+ peer_timeout: Duration,
+ write_timeout: Duration,
+}
+
+impl TangleWebSocketKeepaliveConfig {
+ pub fn new(
+ ping_interval: Duration,
+ peer_timeout: Duration,
+ write_timeout: Duration,
+ ) -> Result<Self, BaseRelayError> {
+ if ping_interval.is_zero() || peer_timeout.is_zero() || write_timeout.is_zero() {
+ return Err(BaseRelayError::invalid(
+ "websocket keepalive durations must be greater than zero",
+ ));
+ }
+ if ping_interval >= peer_timeout {
+ return Err(BaseRelayError::invalid(
+ "websocket ping interval must be shorter than peer timeout",
+ ));
+ }
+ if write_timeout > peer_timeout {
+ return Err(BaseRelayError::invalid(
+ "websocket write timeout must not exceed peer timeout",
+ ));
+ }
+ Ok(Self {
+ ping_interval,
+ peer_timeout,
+ write_timeout,
+ })
+ }
+
+ pub fn ping_interval(self) -> Duration {
+ self.ping_interval
+ }
+
+ pub fn peer_timeout(self) -> Duration {
+ self.peer_timeout
+ }
+
+ pub fn write_timeout(self) -> Duration {
+ self.write_timeout
+ }
}
impl TangleWebSocketSessionOptions {
@@ -74,6 +128,11 @@ impl TangleWebSocketSessionOptions {
self.projection = projection;
self
}
+
+ pub fn with_keepalive(mut self, keepalive: Option<TangleWebSocketKeepaliveConfig>) -> Self {
+ self.keepalive = keepalive;
+ self
+ }
}
impl TangleWebSocketSession {
@@ -151,6 +210,7 @@ impl TangleWebSocketSession {
subscription_permits: BTreeMap::new(),
events,
projection: options.projection,
+ keepalive: options.keepalive,
})
}
@@ -189,20 +249,41 @@ impl TangleWebSocketSession {
);
return;
}
+ let mut last_peer_activity = Instant::now();
+ let mut keepalive_interval = self.keepalive.map(|config| {
+ let mut interval = interval_at(
+ tokio::time::Instant::now() + config.ping_interval(),
+ config.ping_interval(),
+ );
+ interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
+ interval
+ });
loop {
if self.shutdown_requested() {
- let _ = socket.send(Message::Close(None)).await;
+ let _ = self
+ .send_socket_message(&mut socket, Message::Close(None))
+ .await;
break;
}
tokio::select! {
incoming = socket.recv() => {
match incoming {
Some(Ok(Message::Close(_))) | Some(Err(_)) | None => break,
+ Some(Ok(Message::Ping(payload))) => {
+ last_peer_activity = Instant::now();
+ if !self.send_socket_message(&mut socket, Message::Pong(payload)).await {
+ break;
+ }
+ }
+ Some(Ok(Message::Pong(_))) => {
+ last_peer_activity = Instant::now();
+ }
Some(Ok(message)) => {
+ last_peer_activity = Instant::now();
match self.handle_incoming_message(message).await {
TangleSessionControl::Continue => {}
TangleSessionControl::Close(message) => {
- let _ = socket.send(message).await;
+ let _ = self.send_socket_message(&mut socket, message).await;
break;
}
TangleSessionControl::Stop => break,
@@ -214,7 +295,7 @@ impl TangleWebSocketSession {
let Some(message) = outgoing else {
break;
};
- if socket.send(message).await.is_err() {
+ if !self.send_socket_message(&mut socket, message).await {
break;
}
}
@@ -222,7 +303,7 @@ impl TangleWebSocketSession {
match self.handle_event_receive_result(event).await {
TangleSessionControl::Continue => {}
TangleSessionControl::Close(message) => {
- let _ = socket.send(message).await;
+ let _ = self.send_socket_message(&mut socket, message).await;
break;
}
TangleSessionControl::Stop => break,
@@ -230,7 +311,24 @@ impl TangleWebSocketSession {
}
changed = self.shutdown.changed() => {
if changed.is_err() || self.shutdown_requested() {
- let _ = socket.send(Message::Close(None)).await;
+ let _ = self.send_socket_message(&mut socket, Message::Close(None)).await;
+ break;
+ }
+ }
+ _ = keepalive_tick(&mut keepalive_interval) => {
+ let Some(config) = self.keepalive else {
+ continue;
+ };
+ if last_peer_activity.elapsed() >= config.peer_timeout() {
+ let _ = self
+ .send_socket_message(&mut socket, keepalive_timeout_close_message())
+ .await;
+ break;
+ }
+ if !self
+ .send_socket_message(&mut socket, Message::Ping(Vec::new().into()))
+ .await
+ {
break;
}
}
@@ -510,6 +608,24 @@ impl TangleWebSocketSession {
TangleOutboundQueueError::Closed => TangleSessionControl::Stop,
}
}
+
+ async fn send_socket_message(&self, socket: &mut WebSocket, message: Message) -> bool {
+ match self.keepalive {
+ Some(config) => timeout(config.write_timeout(), socket.send(message))
+ .await
+ .is_ok_and(|result| result.is_ok()),
+ None => socket.send(message).await.is_ok(),
+ }
+ }
+}
+
+async fn keepalive_tick(interval: &mut Option<tokio::time::Interval>) {
+ match interval {
+ Some(interval) => {
+ interval.tick().await;
+ }
+ None => pending::<()>().await,
+ }
}
#[derive(Debug, Clone, PartialEq, Eq)]
@@ -540,6 +656,13 @@ fn outbound_encode_close_message() -> Message {
}))
}
+fn keepalive_timeout_close_message() -> Message {
+ Message::Close(Some(CloseFrame {
+ code: 1001,
+ reason: Utf8Bytes::from_static("websocket peer timeout"),
+ }))
+}
+
#[derive(Debug, Clone)]
pub struct TangleOutboundSender {
sender: mpsc::Sender<Message>,
@@ -660,8 +783,9 @@ fn protocol_client_message_to_runtime_for_session_test(
#[cfg(test)]
mod tests {
use super::{
- TangleOutboundQueueError, TangleSessionControl, TangleWebSocketSession,
- current_unix_timestamp, event_stream_lag_close_message, outbound_queue_full_close_message,
+ TangleOutboundQueueError, TangleSessionControl, TangleWebSocketKeepaliveConfig,
+ TangleWebSocketSession, current_unix_timestamp, event_stream_lag_close_message,
+ keepalive_timeout_close_message, outbound_queue_full_close_message,
};
use crate::{
config::{BaseRelayRuntimeConfig, parse_base_relay_runtime_config_json},
@@ -673,7 +797,10 @@ mod tests {
};
use axum::extract::ws::Message;
use serde_json::json;
- use std::path::{Path, PathBuf};
+ use std::{
+ path::{Path, PathBuf},
+ time::Duration,
+ };
use tangle_crypto::RelaySigner;
use tangle_groups::{KIND_GROUP_CREATE_GROUP, StoreOffset};
use tangle_protocol::{
@@ -703,6 +830,48 @@ mod tests {
}
#[test]
+ fn websocket_keepalive_config_is_bounded_and_explicit() {
+ let config = TangleWebSocketKeepaliveConfig::new(
+ Duration::from_secs(30),
+ Duration::from_secs(90),
+ Duration::from_secs(5),
+ )
+ .expect("keepalive");
+ assert_eq!(config.ping_interval(), Duration::from_secs(30));
+ assert_eq!(config.peer_timeout(), Duration::from_secs(90));
+ assert_eq!(config.write_timeout(), Duration::from_secs(5));
+ assert!(
+ TangleWebSocketKeepaliveConfig::new(
+ Duration::ZERO,
+ Duration::from_secs(90),
+ Duration::from_secs(5),
+ )
+ .is_err()
+ );
+ assert!(
+ TangleWebSocketKeepaliveConfig::new(
+ Duration::from_secs(90),
+ Duration::from_secs(90),
+ Duration::from_secs(5),
+ )
+ .is_err()
+ );
+ assert!(
+ TangleWebSocketKeepaliveConfig::new(
+ Duration::from_secs(30),
+ Duration::from_secs(90),
+ Duration::from_secs(91),
+ )
+ .is_err()
+ );
+ let Message::Close(Some(frame)) = keepalive_timeout_close_message() else {
+ panic!("keepalive close frame")
+ };
+ assert_eq!(frame.code, 1001);
+ assert_eq!(frame.reason, "websocket peer timeout");
+ }
+
+ #[test]
fn websocket_session_limit_config_rejects_zero_outbound_capacity() {
assert!(session_limits_result(0).is_err());
}
diff --git a/crates/tangle_test_support/Cargo.toml b/crates/tangle_test_support/Cargo.toml
@@ -8,7 +8,6 @@ license.workspace = true
description = "Test fixtures and utilities for Tangle"
[dependencies]
-k256 = { version = "0.13", features = ["schnorr"] }
pocket-types = { git = "https://github.com/triesap/pocket", rev = "d24cf774a31978cf83d2b7289a618d3959bcc645" }
serde_json = "1"
tangle_crypto = { path = "../tangle_crypto" }
diff --git a/crates/tangle_test_support/src/lib.rs b/crates/tangle_test_support/src/lib.rs
@@ -1,10 +1,8 @@
#![forbid(unsafe_code)]
use core::fmt;
-use k256::schnorr::signature::Signer;
-use k256::schnorr::{Signature, SigningKey};
use pocket_types::OwnedEvent as PocketOwnedEvent;
-use tangle_crypto::{RelaySigner, compute_event_id};
+use tangle_crypto::RelaySigner;
use tangle_groups::{
CanonicalRelayUrl, GroupGeneratedEventBuilder, GroupLimitsConfig, GroupOutboxPayload,
GroupPolicyConfig, GroupRuntimeConfig, GroupRuntimeSettingsConfig, KIND_GROUP_CREATE_GROUP,
@@ -13,8 +11,7 @@ use tangle_groups::{
RelaySecret,
};
use tangle_protocol::{
- Event, EventId, Kind, PublicKeyHex, SignatureHex, Tag, UnixTimestamp, UnsignedEvent,
- event_to_value,
+ Event, EventId, Kind, PublicKeyHex, Tag, UnixTimestamp, UnsignedEvent, event_to_value,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
@@ -28,19 +25,18 @@ pub enum FixtureKey {
impl FixtureKey {
pub fn public_key(self) -> PublicKeyHex {
- let signing_key = self.signing_key();
- PublicKeyHex::new(&lower_hex(signing_key.verifying_key().to_bytes().as_ref()))
- .expect("fixture public key is valid x-only lowercase hex")
+ self.signer().public_key().clone()
}
- fn signing_key(self) -> SigningKey {
- match self {
- Self::Relay => SigningKey::from_bytes(&[9_u8; 32]).expect("relay fixture key"),
- Self::Owner => SigningKey::from_bytes(&[10_u8; 32]).expect("owner fixture key"),
- Self::Admin => SigningKey::from_bytes(&[11_u8; 32]).expect("admin fixture key"),
- Self::Member => SigningKey::from_bytes(&[12_u8; 32]).expect("member fixture key"),
- Self::Outsider => SigningKey::from_bytes(&[13_u8; 32]).expect("outsider fixture key"),
- }
+ fn signer(self) -> RelaySigner {
+ let secret_byte = match self {
+ Self::Relay => 9_u8,
+ Self::Owner => 10_u8,
+ Self::Admin => 11_u8,
+ Self::Member => 12_u8,
+ Self::Outsider => 13_u8,
+ };
+ RelaySigner::from_secret_hex(&lower_hex(&[secret_byte; 32])).expect("fixture signing key")
}
}
@@ -327,16 +323,10 @@ pub fn fixture_event_json(event: &Event) -> serde_json::Value {
}
fn sign_unsigned_event(fixture_key: FixtureKey, unsigned: UnsignedEvent) -> Result<Event, String> {
- let signing_key = fixture_key.signing_key();
- let event_id = compute_event_id(&unsigned);
- let event_id_bytes =
- fixed_hex_bytes(event_id.as_str(), 32, "event id").expect("computed event id decodes");
- let signature: Signature = signing_key.sign(&event_id_bytes);
- let signature =
- SignatureHex::new(&lower_hex(signature.to_bytes().as_ref())).expect("signature hex");
- Ok(Event::new(event_id, unsigned, signature))
+ Ok(fixture_key.signer().sign_unsigned_event(unsigned))
}
+#[cfg(test)]
fn fixed_hex_bytes(value: &str, expected: usize, scalar: &str) -> Result<Vec<u8>, String> {
if value.len() != expected * 2 {
return Err(format!(
@@ -351,6 +341,7 @@ fn fixed_hex_bytes(value: &str, expected: usize, scalar: &str) -> Result<Vec<u8>
Ok(output)
}
+#[cfg(test)]
fn hex_value(value: u8, scalar: &str) -> Result<u8, String> {
match value {
b'0'..=b'9' => Ok(value - b'0'),