1mod auth;
2mod blob_read;
3pub mod blossom;
4mod handlers;
5mod ingest_filter;
6mod mime;
7mod nostr_query;
8mod peer_status;
9mod request_paths;
10mod status_metrics;
11mod ui;
12pub mod ws_relay;
13
14use crate::nostr_relay::NostrRelay;
15use crate::socialgraph;
16use crate::storage::HashtreeStore;
17use anyhow::Result;
18use axum::{
19 body::Body,
20 extract::{DefaultBodyLimit, State},
21 http::{header, HeaderValue, Method, Request, StatusCode},
22 middleware,
23 middleware::Next,
24 response::{IntoResponse, Json, Response},
25 routing::{get, post, put},
26 Router,
27};
28use futures::{future::poll_fn, pin_mut, FutureExt};
29use hashtree_core::Cid;
30use hyper::body::Incoming;
31use hyper_util::{
32 rt::{TokioExecutor, TokioIo, TokioTimer},
33 server::conn::auto::Builder as HyperBuilder,
34 service::TowerToHyperService,
35};
36use socket2::{SockRef, TcpKeepalive};
37use std::collections::{HashMap, HashSet};
38use std::convert::Infallible;
39use std::future;
40use std::io;
41use std::net::SocketAddr;
42use std::sync::{Arc, OnceLock, RwLock};
43use std::time::{Duration, Instant};
44use tokio::sync::watch;
45use tower::{Service, ServiceExt as _};
46use tower_http::cors::CorsLayer;
47use tracing::{debug, error, trace};
48
49pub use auth::{
50 new_lookup_cache, new_upstream_http_client, AppState, AuthCredentials, CachedTreeRootEntry,
51};
52
53static VIRTUAL_TREE_HOSTS: OnceLock<RwLock<HashMap<String, String>>> = OnceLock::new();
54const DEFAULT_OPTIMISTIC_UPLOAD_QUEUE_BYTES: usize = 256 * 1024 * 1024;
55const DEFAULT_BLOSSOM_UPLOAD_REPLICA_QUEUE_BYTES: usize = 256 * 1024 * 1024;
56const INTERNAL_JSON_BODY_LIMIT_BYTES: usize = 64 * 1024;
57const POOL_AUDIT_READ_ONLY_HTTP_ERROR: &str = "PoolStore is in audit-serving read-only mode";
58const POOL_AUDIT_READ_ONLY_REASON: &str = "pool-audit-read-only";
59const POOL_AUDIT_READ_ONLY_REASON_HEADER: &str = "x-hashtree-maintenance-reason";
60
61#[cfg(not(test))]
62const HTTP1_HEADER_READ_TIMEOUT: Duration = Duration::from_secs(30);
63#[cfg(test)]
64const HTTP1_HEADER_READ_TIMEOUT: Duration = Duration::from_millis(200);
65const HTTP2_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(30);
66const HTTP2_KEEPALIVE_TIMEOUT: Duration = Duration::from_secs(10);
67const TCP_LISTEN_BACKLOG: i32 = 1_024;
68const TCP_KEEPALIVE_TIME: Duration = Duration::from_secs(60);
69const TCP_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(15);
70
71pub fn bounded_upload_queue_bytes(bytes: u64) -> usize {
72 usize::try_from(bytes)
73 .unwrap_or(usize::MAX)
74 .clamp(1, tokio::sync::Semaphore::MAX_PERMITS)
75}
76
77fn virtual_tree_hosts() -> &'static RwLock<HashMap<String, String>> {
78 VIRTUAL_TREE_HOSTS.get_or_init(|| RwLock::new(HashMap::new()))
79}
80
81fn pool_audit_request_is_allowed(method: &Method, path: &str) -> bool {
82 if path == "/ws" || path == "/ws/" {
83 return false;
84 }
85 matches!(*method, Method::GET | Method::HEAD | Method::OPTIONS)
86 || (*method == Method::POST && matches!(path, "/blob/batch" | "/upload/check"))
87}
88
89async fn pool_audit_read_only_middleware(
90 State(state): State<AppState>,
91 request: Request<Body>,
92 next: Next,
93) -> Response {
94 if !state.store.is_pool_audit_read_only()
95 || pool_audit_request_is_allowed(request.method(), request.uri().path())
96 {
97 return next.run(request).await;
98 }
99
100 let mut response = (
101 StatusCode::SERVICE_UNAVAILABLE,
102 Json(serde_json::json!({
103 "error": POOL_AUDIT_READ_ONLY_HTTP_ERROR,
104 })),
105 )
106 .into_response();
107 response.headers_mut().insert(
108 header::HeaderName::from_static(POOL_AUDIT_READ_ONLY_REASON_HEADER),
109 HeaderValue::from_static(POOL_AUDIT_READ_ONLY_REASON),
110 );
111 response
112}
113
114fn normalize_virtual_tree_host(host: &str) -> Option<String> {
115 let trimmed = host.trim().trim_end_matches('.').to_ascii_lowercase();
116 if trimmed.is_empty() {
117 return None;
118 }
119
120 if let Some(stripped) = trimmed
121 .strip_prefix('[')
122 .and_then(|value| value.split_once(']'))
123 {
124 let host_only = stripped.0.trim();
125 if host_only.is_empty() {
126 return None;
127 }
128 return Some(host_only.to_string());
129 }
130
131 if let Some((host_only, port)) = trimmed.rsplit_once(':') {
132 if !host_only.is_empty() && !port.is_empty() && port.chars().all(|ch| ch.is_ascii_digit()) {
133 return Some(host_only.to_string());
134 }
135 }
136
137 Some(trimmed)
138}
139
140pub fn register_virtual_tree_host(host: &str, internal_root: &str) {
141 let Some(normalized_host) = normalize_virtual_tree_host(host) else {
142 return;
143 };
144
145 let normalized_root = internal_root.trim().trim_end_matches('/');
146 if normalized_root.is_empty() {
147 return;
148 }
149
150 if let Ok(mut hosts) = virtual_tree_hosts().write() {
151 hosts.insert(normalized_host, normalized_root.to_string());
152 }
153}
154
155pub fn resolve_virtual_tree_host(host: &str) -> Option<String> {
156 let normalized_host = normalize_virtual_tree_host(host)?;
157 let configured = virtual_tree_hosts()
158 .read()
159 .ok()
160 .and_then(|hosts| hosts.get(&normalized_host).cloned());
161 configured.or_else(|| resolve_iris_localhost_tree_root(&normalized_host))
162}
163
164fn resolve_iris_localhost_tree_root(host: &str) -> Option<String> {
165 let labels: Vec<&str> = host.strip_suffix(".iris.localhost")?.split('.').collect();
166 match labels.as_slice() {
167 [nhash] if nhash.starts_with("nhash1") => Some(format!("/htree/{nhash}")),
168 [site, npub] if !site.is_empty() && npub.starts_with("npub1") => {
169 Some(format!("/htree/{npub}/{site}"))
170 }
171 _ => None,
172 }
173}
174
175#[cfg(test)]
176pub fn clear_virtual_tree_hosts_for_test() {
177 if let Ok(mut hosts) = virtual_tree_hosts().write() {
178 hosts.clear();
179 }
180}
181
182pub struct HashtreeServer {
183 state: AppState,
184 addr: String,
185 extra_routes: Option<Router<AppState>>,
186 cors: Option<CorsLayer>,
187}
188
189impl HashtreeServer {
190 pub fn new(store: Arc<HashtreeStore>, addr: String) -> Self {
191 Self {
192 state: AppState {
193 store,
194 auth: None,
195 daemon_started_at: current_unix_secs(),
196 peer_mode: crate::config::ServerMode::Normal,
197 hash_get_enabled: true,
198 fips_endpoint: None,
199 fips_blob_resolver: None,
200 fetch_from_fips_peers: true,
201 ws_relay: Arc::new(auth::WsRelayState::new()),
202 max_upload_bytes: 5 * 1024 * 1024, public_writes: true, public_plaintext_reads: false,
205 require_random_untrusted_ingest: true,
206 optimistic_blossom_uploads: false,
207 optimistic_upload_queue_bytes: DEFAULT_OPTIMISTIC_UPLOAD_QUEUE_BYTES,
208 optimistic_upload_queue: Arc::new(tokio::sync::Semaphore::new(
209 DEFAULT_OPTIMISTIC_UPLOAD_QUEUE_BYTES,
210 )),
211 allowed_pubkeys: HashSet::new(), upstream_blossom: Vec::new(),
213 upstream_http_client: new_upstream_http_client(),
214 upstream_blossom_miss_cache: Arc::new(std::sync::Mutex::new(new_lookup_cache())),
215 upstream_blossom_fetch_metrics: Arc::new(
216 auth::UpstreamBlossomFetchMetrics::default(),
217 ),
218 blossom_upload_replicas: Vec::new(),
219 blossom_upload_replica_queue_bytes: DEFAULT_BLOSSOM_UPLOAD_REPLICA_QUEUE_BYTES,
220 blossom_upload_replica_queue: Arc::new(tokio::sync::Semaphore::new(
221 DEFAULT_BLOSSOM_UPLOAD_REPLICA_QUEUE_BYTES,
222 )),
223 blossom_upload_replica_keys: None,
224 blossom_upload_replica_scheduler: Arc::new(
225 blossom::BlossomUploadReplicaScheduler::new(),
226 ),
227 social_graph: None,
228 social_graph_store: None,
229 social_graph_root: None,
230 socialgraph_snapshot_public: false,
231 nostr_relay: None,
232 nostr_provider: None,
233 nostr_relay_urls: Vec::new(),
234 tree_root_cache: Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())),
235 inflight_blob_fetches: Arc::new(tokio::sync::Mutex::new(
236 std::collections::HashMap::new(),
237 )),
238 inflight_blob_reads: Arc::new(tokio::sync::Mutex::new(
239 std::collections::HashMap::new(),
240 )),
241 blob_cache: Arc::new(crate::blob_cache::BlobCache::from_env()),
242 directory_listing_cache: Arc::new(std::sync::Mutex::new(new_lookup_cache())),
243 resolved_path_cache: Arc::new(std::sync::Mutex::new(new_lookup_cache())),
244 thumbnail_path_cache: Arc::new(std::sync::Mutex::new(new_lookup_cache())),
245 cid_size_cache: Arc::new(std::sync::Mutex::new(new_lookup_cache())),
246 },
247 addr,
248 extra_routes: None,
249 cors: None,
250 }
251 }
252
253 pub fn with_max_upload_bytes(mut self, bytes: usize) -> Self {
255 self.state.max_upload_bytes = bytes;
256 self
257 }
258
259 pub fn with_public_writes(mut self, public: bool) -> Self {
262 self.state.public_writes = public;
263 self
264 }
265
266 pub fn with_public_plaintext_reads(mut self, public: bool) -> Self {
268 self.state.public_plaintext_reads = public;
269 self
270 }
271
272 pub fn with_require_random_untrusted_ingest(mut self, require: bool) -> Self {
273 self.state.require_random_untrusted_ingest = require;
274 self
275 }
276
277 pub fn with_optimistic_blossom_uploads(mut self, enabled: bool) -> Self {
278 self.state.optimistic_blossom_uploads = enabled;
279 self
280 }
281
282 pub fn with_server_mode(mut self, mode: crate::config::ServerMode) -> Self {
283 self.state.peer_mode = mode;
284 self
285 }
286
287 pub fn with_hash_get_enabled(mut self, enabled: bool) -> Self {
288 self.state.hash_get_enabled = enabled;
289 self
290 }
291
292 pub fn with_fetch_from_fips_peers(mut self, enabled: bool) -> Self {
293 self.state.fetch_from_fips_peers = enabled;
294 self
295 }
296
297 pub fn with_fips_endpoint(
298 mut self,
299 endpoint: Arc<hashtree_fips_transport::FipsEndpoint>,
300 ) -> Self {
301 self.state.fips_endpoint = Some(endpoint);
302 self
303 }
304
305 pub fn with_fips_blob_resolver(
306 mut self,
307 resolver: Arc<crate::fips_transport::DaemonBlobResolver>,
308 ) -> Self {
309 self.state.fips_blob_resolver = Some(resolver);
310 self
311 }
312
313 pub fn with_auth(mut self, username: String, password: String) -> Self {
314 self.state.auth = Some(AuthCredentials { username, password });
315 self
316 }
317
318 pub fn with_allowed_pubkeys(mut self, pubkeys: HashSet<String>) -> Self {
320 self.state.allowed_pubkeys = pubkeys;
321 self
322 }
323
324 pub fn with_upstream_blossom(mut self, servers: Vec<String>) -> Self {
326 self.state.upstream_blossom = servers;
327 self
328 }
329
330 pub fn with_blossom_upload_replicas(
332 mut self,
333 servers: Vec<String>,
334 queue_bytes: usize,
335 keys: nostr::Keys,
336 ) -> Self {
337 let queue_bytes = queue_bytes.clamp(1, tokio::sync::Semaphore::MAX_PERMITS);
338 let mut servers: Vec<String> = servers
339 .into_iter()
340 .map(|server| server.trim().trim_end_matches('/').to_string())
341 .filter(|server| !server.is_empty())
342 .collect();
343 servers.sort();
344 servers.dedup();
345 let replica_keys = (!servers.is_empty()).then(|| Arc::new(keys));
346 self.state.blossom_upload_replicas = servers;
347 self.state.blossom_upload_replica_queue_bytes = queue_bytes;
348 self.state.blossom_upload_replica_queue =
349 Arc::new(tokio::sync::Semaphore::new(queue_bytes));
350 self.state.blossom_upload_replica_keys = replica_keys;
351 self
352 }
353
354 pub fn with_social_graph(mut self, sg: Arc<socialgraph::SocialGraphAccessControl>) -> Self {
356 self.state.social_graph = Some(sg);
357 self
358 }
359
360 pub fn with_socialgraph_snapshot(
362 mut self,
363 store: Arc<dyn socialgraph::SocialGraphBackend>,
364 root: [u8; 32],
365 public: bool,
366 ) -> Self {
367 self.state.social_graph_store = Some(store);
368 self.state.social_graph_root = Some(root);
369 self.state.socialgraph_snapshot_public = public;
370 self
371 }
372
373 pub fn with_nostr_relay(mut self, relay: Arc<NostrRelay>) -> Self {
375 self.state.nostr_relay = Some(relay);
376 self
377 }
378
379 pub fn with_nostr_provider(mut self, provider: Arc<dyn nostr_pubsub::PubsubProvider>) -> Self {
380 self.state.nostr_provider = Some(provider);
381 self
382 }
383
384 pub fn with_nostr_relay_urls(mut self, relays: Vec<String>) -> Self {
386 self.state.nostr_relay_urls = relays;
387 self
388 }
389
390 pub fn with_cached_tree_roots(self, roots: Vec<(String, Cid)>) -> Self {
392 if let Ok(mut cache) = self.state.tree_root_cache.lock() {
393 let now = Instant::now();
394 for (key, cid) in roots {
395 cache.insert(
396 key,
397 CachedTreeRootEntry {
398 cid,
399 source: "embedded-bootstrap",
400 root_event: None,
401 event: None,
402 cached_at: now,
403 },
404 );
405 }
406 }
407 self
408 }
409
410 pub fn with_extra_routes(mut self, routes: Router<AppState>) -> Self {
412 self.extra_routes = Some(routes);
413 self
414 }
415
416 pub fn with_cors(mut self, cors: CorsLayer) -> Self {
418 self.cors = Some(cors);
419 self
420 }
421
422 pub async fn run(self) -> Result<()> {
423 let listener = tokio::net::TcpListener::bind(&self.addr).await?;
424 let _ = self.run_with_listener(listener).await?;
425 Ok(())
426 }
427
428 pub async fn run_with_listener(self, listener: tokio::net::TcpListener) -> Result<u16> {
429 self.run_with_listener_until(listener, future::pending::<()>())
430 .await
431 }
432
433 pub async fn run_with_listener_until<F>(
434 self,
435 listener: tokio::net::TcpListener,
436 shutdown: F,
437 ) -> Result<u16>
438 where
439 F: std::future::Future<Output = ()> + Send + 'static,
440 {
441 SockRef::from(&listener).listen(TCP_LISTEN_BACKLOG)?;
448 let local_addr = listener.local_addr()?;
449
450 let state = self.state.clone();
454 let public_routes = Router::new()
455 .route("/", get(handlers::serve_root_or_virtual_host))
456 .route("/ws", get(ws_relay::ws_data))
457 .route("/ws/", get(ws_relay::ws_data))
458 .route(
459 "/__iris/store/:hash",
460 get(handlers::iris_store_get).head(handlers::iris_store_head),
461 )
462 .route(
463 "/htree/test",
464 get(handlers::htree_test).head(handlers::htree_test),
465 )
466 .route("/htree/nhash1:nhash", get(handlers::htree_nhash))
468 .route("/htree/nhash1:nhash/", get(handlers::htree_nhash))
469 .route("/htree/nhash1:nhash/*path", get(handlers::htree_nhash_path))
470 .route("/htree/npub1:npub/:treename", get(handlers::htree_npub))
472 .route("/htree/npub1:npub/:treename/", get(handlers::htree_npub))
473 .route(
474 "/htree/npub1:npub/:treename/*path",
475 get(handlers::htree_npub_path),
476 )
477 .route("/n/:pubkey/:treename", get(handlers::resolve_and_serve))
479 .route("/npub1:rest", get(handlers::serve_npub))
481 .route("/npub1:rest/*path", get(handlers::serve_npub))
482 .route(
484 "/:id",
485 get(handlers::serve_content_or_blob)
486 .head(blossom::head_blob)
487 .delete(blossom::delete_blob)
488 .options(blossom::cors_preflight),
489 )
490 .route(
491 "/upload",
492 put(blossom::upload_blob)
493 .layer::<_, std::convert::Infallible>(middleware::from_fn(
494 blossom::require_upload_auth_middleware,
495 ))
496 .layer(DefaultBodyLimit::max(blossom::MAX_SINGLE_UPLOAD_BODY_BYTES))
497 .head(blossom::head_upload)
498 .options(blossom::cors_preflight),
499 )
500 .route(
501 "/upload/batch",
502 post(blossom::upload_blob_batch)
503 .layer::<_, std::convert::Infallible>(middleware::from_fn(
504 blossom::require_upload_auth_middleware,
505 ))
506 .options(blossom::cors_preflight)
507 .layer(DefaultBodyLimit::max(
508 blossom::MAX_BATCH_UPLOAD_JSON_BODY_BYTES,
509 )),
510 )
511 .route(
512 "/upload/batch-binary",
513 post(blossom::upload_blob_batch_binary)
514 .layer::<_, std::convert::Infallible>(middleware::from_fn(
515 blossom::require_upload_auth_middleware,
516 ))
517 .options(blossom::cors_preflight)
518 .layer(DefaultBodyLimit::max(
519 blossom::MAX_BATCH_UPLOAD_BINARY_BODY_BYTES,
520 )),
521 )
522 .route(
523 "/upload/check",
524 post(blossom::upload_check).options(blossom::cors_preflight),
525 )
526 .route(
527 "/blob/batch",
528 post(handlers::download_blob_batch).options(blossom::cors_preflight),
529 )
530 .route(
531 "/list/:pubkey",
532 get(blossom::list_blobs).options(blossom::cors_preflight),
533 )
534 .route("/health", get(handlers::health_check))
536 .route("/api/pins", get(handlers::list_pins))
537 .route("/api/stats", get(handlers::storage_stats))
538 .route("/api/status", get(handlers::daemon_status))
539 .route("/api/socialgraph", get(handlers::socialgraph_stats))
540 .route(
541 "/api/socialgraph/snapshot",
542 get(handlers::socialgraph_snapshot),
543 )
544 .route(
545 "/api/socialgraph/distance/:pubkey",
546 get(handlers::follow_distance),
547 )
548 .route(
550 "/api/resolve/:pubkey/:treename",
551 get(handlers::resolve_to_hash),
552 )
553 .route(
554 "/api/nostr/resolve/:pubkey/:treename",
555 get(handlers::resolve_to_hash),
556 )
557 .route("/api/nostr/profile/:pubkey", get(handlers::nostr_profile))
558 .route(
559 "/api/nostr/events",
560 post(handlers::publish_nostr_event)
561 .layer(DefaultBodyLimit::max(INTERNAL_JSON_BODY_LIMIT_BYTES)),
562 )
563 .route("/api/trees/:pubkey", get(handlers::list_trees))
564 .fallback(get(handlers::serve_virtual_host_fallback))
565 .with_state(state.clone());
566
567 let protected_routes = Router::new()
569 .route("/upload", post(handlers::upload_file))
570 .route("/api/pin/:cid", post(handlers::pin_cid))
571 .route("/api/unpin/:cid", post(handlers::unpin_cid))
572 .route("/api/gc", post(handlers::garbage_collect))
573 .layer(middleware::from_fn_with_state(
574 state.clone(),
575 auth::auth_middleware,
576 ))
577 .with_state(state.clone());
578
579 let internal_routes = Router::new()
582 .route(
583 "/__iris/store/:hash",
584 put(handlers::iris_store_put)
585 .delete(handlers::iris_store_delete)
586 .layer(DefaultBodyLimit::max(blossom::MAX_BATCH_UPLOAD_BYTES)),
587 )
588 .route(
589 "/api/pin-tree",
590 post(handlers::pin_tree)
591 .layer(DefaultBodyLimit::max(INTERNAL_JSON_BODY_LIMIT_BYTES)),
592 )
593 .route(
594 "/api/cache-tree-root",
595 post(handlers::cache_tree_root)
596 .layer(DefaultBodyLimit::max(INTERNAL_JSON_BODY_LIMIT_BYTES)),
597 )
598 .route(
599 "/api/clear-tree-root-cache",
600 post(handlers::clear_tree_root_cache)
601 .layer(DefaultBodyLimit::max(INTERNAL_JSON_BODY_LIMIT_BYTES)),
602 )
603 .layer(middleware::from_fn_with_state(
604 state.clone(),
605 auth::require_auth_middleware,
606 ))
607 .with_state(state.clone());
608
609 let mut app = public_routes
610 .merge(protected_routes)
611 .merge(internal_routes)
612 .layer(DefaultBodyLimit::max(10 * 1024 * 1024 * 1024)); if let Some(extra) = self.extra_routes {
615 app = app.merge(extra.with_state(state.clone()));
616 }
617
618 app = app
624 .layer(middleware::from_fn_with_state(
625 state,
626 pool_audit_read_only_middleware,
627 ))
628 .layer(middleware::from_fn(status_metrics::record_http_status));
629
630 if let Some(cors) = self.cors {
631 app = app.layer(cors);
632 }
633
634 let make_service = app.into_make_service_with_connect_info::<std::net::SocketAddr>();
635 serve_with_connection_limits(listener, make_service, shutdown).await?;
636
637 Ok(local_addr.port())
638 }
639
640 pub fn addr(&self) -> &str {
641 &self.addr
642 }
643}
644
645async fn serve_with_connection_limits<M, S, F>(
646 listener: tokio::net::TcpListener,
647 mut make_service: M,
648 shutdown: F,
649) -> io::Result<()>
650where
651 M: Service<SocketAddr, Error = Infallible, Response = S> + Send + 'static,
652 M::Future: Send,
653 S: Service<Request<Body>, Response = Response, Error = Infallible> + Clone + Send + 'static,
654 S::Future: Send,
655 F: std::future::Future<Output = ()> + Send + 'static,
656{
657 let (signal_tx, signal_rx) = watch::channel(());
658 let signal_tx = Arc::new(signal_tx);
659 tokio::spawn(async move {
660 shutdown.await;
661 trace!("received graceful shutdown signal; stopping daemon listener");
662 drop(signal_rx);
663 });
664
665 let (close_tx, close_rx) = watch::channel(());
666
667 loop {
668 let (tcp_stream, remote_addr) = tokio::select! {
669 accepted = accept_tcp(&listener) => {
670 match accepted? {
671 Some(connection) => connection,
672 None => continue,
673 }
674 }
675 _ = signal_tx.closed() => {
676 trace!("shutdown signal received; no longer accepting daemon connections");
677 break;
678 }
679 };
680
681 configure_tcp_stream(&tcp_stream);
682 let tcp_stream = TokioIo::new(tcp_stream);
683
684 poll_fn(|cx| make_service.poll_ready(cx))
685 .await
686 .unwrap_or_else(|err| match err {});
687
688 let tower_service = make_service
689 .call(remote_addr)
690 .await
691 .unwrap_or_else(|err| match err {})
692 .map_request(|req: Request<Incoming>| req.map(Body::new));
693 let hyper_service = TowerToHyperService::new(tower_service);
694
695 let signal_tx = Arc::clone(&signal_tx);
696 let close_rx = close_rx.clone();
697
698 tokio::spawn(async move {
699 let mut builder = HyperBuilder::new(TokioExecutor::new());
700 builder
701 .http1()
702 .timer(TokioTimer::new())
703 .header_read_timeout(HTTP1_HEADER_READ_TIMEOUT);
704 builder
705 .http2()
706 .timer(TokioTimer::new())
707 .keep_alive_interval(Some(HTTP2_KEEPALIVE_INTERVAL))
708 .keep_alive_timeout(HTTP2_KEEPALIVE_TIMEOUT);
709
710 let conn = builder.serve_connection_with_upgrades(tcp_stream, hyper_service);
711 pin_mut!(conn);
712
713 let signal_closed = signal_tx.closed().fuse();
714 pin_mut!(signal_closed);
715
716 loop {
717 tokio::select! {
718 result = conn.as_mut() => {
719 if let Err(err) = result {
720 trace!("daemon connection closed with error: {err:#}");
721 }
722 break;
723 }
724 _ = &mut signal_closed => {
725 trace!("shutdown signal received by connection task");
726 conn.as_mut().graceful_shutdown();
727 }
728 }
729 }
730
731 drop(close_rx);
732 });
733 }
734
735 drop(close_rx);
736 drop(listener);
737 close_tx.closed().await;
738
739 Ok(())
740}
741
742fn configure_tcp_stream(tcp_stream: &tokio::net::TcpStream) {
743 if let Err(err) = tcp_stream.set_nodelay(true) {
744 debug!("failed to set TCP_NODELAY on daemon connection: {err:#}");
745 }
746
747 let socket = SockRef::from(tcp_stream);
748 if let Err(err) = socket.set_tcp_keepalive(
749 &TcpKeepalive::new()
750 .with_time(TCP_KEEPALIVE_TIME)
751 .with_interval(TCP_KEEPALIVE_INTERVAL),
752 ) {
753 debug!("failed to set TCP keepalive on daemon connection: {err:#}");
754 }
755}
756
757async fn accept_tcp(
758 listener: &tokio::net::TcpListener,
759) -> io::Result<Option<(tokio::net::TcpStream, SocketAddr)>> {
760 match listener.accept().await {
761 Ok(connection) => Ok(Some(connection)),
762 Err(err) => {
763 if is_connection_error(&err) {
764 return Ok(None);
765 }
766 if is_resource_exhaustion_error(&err) {
767 error!(
768 "daemon accept failed due to file descriptor exhaustion; exiting for supervisor restart: {err}"
769 );
770 return Err(err);
771 }
772 error!("daemon accept error: {err}");
773 tokio::time::sleep(Duration::from_secs(1)).await;
774 Ok(None)
775 }
776 }
777}
778
779fn is_resource_exhaustion_error(err: &io::Error) -> bool {
780 matches!(
781 err.raw_os_error(),
782 Some(code) if code == libc::EMFILE || code == libc::ENFILE
783 )
784}
785
786fn is_connection_error(err: &io::Error) -> bool {
787 matches!(
788 err.kind(),
789 io::ErrorKind::ConnectionRefused
790 | io::ErrorKind::ConnectionAborted
791 | io::ErrorKind::ConnectionReset
792 )
793}
794
795fn current_unix_secs() -> u64 {
796 std::time::SystemTime::now()
797 .duration_since(std::time::UNIX_EPOCH)
798 .unwrap_or(std::time::Duration::ZERO)
799 .as_secs()
800}
801
802#[cfg(test)]
803mod tests {
804 use super::*;
805 use crate::nostr_relay::{NostrRelay, NostrRelayConfig};
806 use crate::storage::HashtreeStore;
807 use async_trait::async_trait;
808 use hashtree_config::StorageBackend;
809 use hashtree_core::types::Hash;
810 use hashtree_core::{
811 from_hex, nhash_encode, nhash_encode_full, sha256, to_hex, DirEntry, HashTree,
812 HashTreeConfig, LinkType, NHashData,
813 };
814 use nostr::{nips::nip19::ToBech32, EventBuilder, Keys, Kind, Timestamp};
815 use nostr_pubsub::{
816 EventBus, EventSource, PublishReport, PubsubProvider, PubsubProviderMode, QueryEvent,
817 QueryOptions, QueryReport, VerifiedEvent,
818 };
819 use serde_json::json;
820 use std::path::{Path, PathBuf};
821 use std::process::Command;
822 use std::sync::atomic::{AtomicUsize, Ordering};
823 use tempfile::TempDir;
824 use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
825 use walkdir::WalkDir;
826
827 const AUDIT_SERVER_CHILD_DATA_ENV: &str = "HASHTREE_AUDIT_SERVER_CHILD_DATA";
828 const AUDIT_SERVER_CHILD_HASH_ENV: &str = "HASHTREE_AUDIT_SERVER_CHILD_HASH";
829
830 fn pool_data_file_snapshot(root: &Path) -> Vec<(PathBuf, Hash, std::time::SystemTime)> {
831 let mut snapshot = WalkDir::new(root)
832 .into_iter()
833 .filter_map(Result::ok)
834 .filter(|entry| {
835 if !entry.file_type().is_file() || entry.file_name() != "data.mdb" {
836 return false;
837 }
838 matches!(
839 entry
840 .path()
841 .parent()
842 .and_then(Path::file_name)
843 .and_then(|name| name.to_str()),
844 Some(hashtree_lmdb::SHARED_BLOB_POOL_DIR_NAME) | Some("blobs")
845 )
846 })
847 .map(|entry| {
848 let path = entry.into_path();
849 let bytes = std::fs::read(&path).expect("read Pool LMDB data file");
850 let modified = std::fs::metadata(&path)
851 .expect("read Pool LMDB metadata")
852 .modified()
853 .expect("read Pool LMDB mtime");
854 (path, sha256(&bytes), modified)
855 })
856 .collect::<Vec<_>>();
857 snapshot.sort_unstable_by(|left, right| left.0.cmp(&right.0));
858 snapshot
859 }
860
861 async fn assert_pool_audit_maintenance_response(response: reqwest::Response) -> Result<()> {
862 assert_eq!(response.status(), reqwest::StatusCode::SERVICE_UNAVAILABLE);
863 assert_eq!(
864 response
865 .headers()
866 .get(POOL_AUDIT_READ_ONLY_REASON_HEADER)
867 .and_then(|value| value.to_str().ok()),
868 Some(POOL_AUDIT_READ_ONLY_REASON)
869 );
870 let body: serde_json::Value = response.json().await?;
871 assert_eq!(
872 body,
873 json!({
874 "error": POOL_AUDIT_READ_ONLY_HTTP_ERROR,
875 })
876 );
877 Ok(())
878 }
879
880 struct StaticProvider {
881 event: Option<VerifiedEvent>,
882 queries: AtomicUsize,
883 publishes: AtomicUsize,
884 mode: PubsubProviderMode,
885 }
886
887 #[test]
888 fn pool_audit_request_gate_only_allows_read_ingress() {
889 for method in [Method::GET, Method::HEAD, Method::OPTIONS] {
890 assert!(pool_audit_request_is_allowed(&method, "/health"));
891 }
892 assert!(pool_audit_request_is_allowed(&Method::POST, "/blob/batch"));
893 assert!(pool_audit_request_is_allowed(
894 &Method::POST,
895 "/upload/check"
896 ));
897 for (method, path) in [
898 (Method::GET, "/ws"),
899 (Method::GET, "/ws/"),
900 (Method::PUT, "/upload"),
901 (Method::POST, "/upload/batch"),
902 (Method::POST, "/upload/batch-binary"),
903 (Method::POST, "/api/pin/hash"),
904 (Method::POST, "/api/nostr/events"),
905 (Method::DELETE, "/hash"),
906 ] {
907 assert!(
908 !pool_audit_request_is_allowed(&method, path),
909 "{method} {path} must be blocked"
910 );
911 }
912 }
913
914 #[tokio::test]
915 #[ignore = "subprocess entry point for Pool audit-serving HTTP verification"]
916 async fn pool_audit_read_only_server_subprocess() -> Result<()> {
917 let Some(data_dir) = std::env::var_os(AUDIT_SERVER_CHILD_DATA_ENV) else {
918 return Ok(());
919 };
920 let hash_hex = std::env::var(AUDIT_SERVER_CHILD_HASH_ENV)?;
921 let store = Arc::new(HashtreeStore::new_with_backend(
922 data_dir,
923 StorageBackend::Lmdb,
924 16 * 1024 * 1024 * 1024,
925 )?);
926 assert!(store.is_pool_audit_read_only());
927
928 let owner = Keys::generate();
929 let tree_name = "pool-audit-external";
930 let root_event = EventBuilder::new(Kind::Custom(30064), "")
931 .tags(vec![
932 nostr::Tag::identifier(tree_name),
933 nostr::Tag::custom(nostr::TagKind::custom("l"), vec!["hashtree".to_string()]),
934 nostr::Tag::custom(nostr::TagKind::custom("hash"), vec![hash_hex.clone()]),
935 ])
936 .sign_with_keys(&owner)?;
937 let provider = Arc::new(StaticProvider {
938 event: Some(VerifiedEvent::try_from(root_event.clone())?),
939 queries: AtomicUsize::new(0),
940 publishes: AtomicUsize::new(0),
941 mode: PubsubProviderMode::LocalOnly,
942 });
943 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
944 let port = listener.local_addr()?.port();
945 let server = HashtreeServer::new(Arc::clone(&store), "127.0.0.1:0".to_string())
946 .with_nostr_provider(provider.clone());
947 let handle =
948 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
949 let client = reqwest::Client::new();
950 let base = format!("http://127.0.0.1:{port}");
951
952 let response = client.get(format!("{base}/{hash_hex}")).send().await?;
953 assert_eq!(response.status(), reqwest::StatusCode::OK);
954 assert_eq!(response.bytes().await?.as_ref(), b"audit-serving bytes");
955
956 let status: serde_json::Value = client
957 .get(format!("{base}/api/status"))
958 .send()
959 .await?
960 .error_for_status()?
961 .json()
962 .await?;
963 assert_eq!(status["pool_audit_read_only"], true);
964 assert_eq!(status["capabilities"]["writes"], false);
965
966 let npub = owner.public_key().to_bech32()?;
967 let resolved: serde_json::Value = client
968 .get(format!(
969 "{base}/api/nostr/resolve/{npub}/{tree_name}?refresh=1"
970 ))
971 .send()
972 .await?
973 .error_for_status()?
974 .json()
975 .await?;
976 assert_eq!(resolved["hash"], hash_hex);
977 assert_eq!(resolved["source"], "fips-pubsub");
978 assert_eq!(resolved["event_id"], root_event.id.to_hex());
979 assert_eq!(provider.queries.load(Ordering::Relaxed), 1);
980 assert_eq!(provider.publishes.load(Ordering::Relaxed), 0);
981
982 assert_pool_audit_maintenance_response(
983 client
984 .post(format!("{base}/api/nostr/events"))
985 .json(&root_event)
986 .send()
987 .await?,
988 )
989 .await?;
990 assert_eq!(provider.publishes.load(Ordering::Relaxed), 0);
991
992 assert_pool_audit_maintenance_response(
993 client
994 .put(format!("{base}/upload"))
995 .body(b"blocked upload".to_vec())
996 .send()
997 .await?,
998 )
999 .await?;
1000 assert_pool_audit_maintenance_response(
1001 client
1002 .post(format!("{base}/upload/batch"))
1003 .json(&json!({"blobs": []}))
1004 .send()
1005 .await?,
1006 )
1007 .await?;
1008 assert_pool_audit_maintenance_response(
1009 client.delete(format!("{base}/{hash_hex}")).send().await?,
1010 )
1011 .await?;
1012 assert_pool_audit_maintenance_response(client.get(format!("{base}/ws")).send().await?)
1013 .await?;
1014
1015 let response = client.get(format!("{base}/{hash_hex}")).send().await?;
1016 assert_eq!(response.status(), reqwest::StatusCode::OK);
1017 assert_eq!(response.bytes().await?.as_ref(), b"audit-serving bytes");
1018
1019 handle.abort();
1020 Ok(())
1021 }
1022
1023 #[test]
1024 fn pool_audit_read_only_server_resolves_external_nostr_without_mutation() -> Result<()> {
1025 let temp = TempDir::new()?;
1026 let data_dir = temp.path().join("data");
1027 let store = HashtreeStore::new_with_backend(
1028 &data_dir,
1029 StorageBackend::Lmdb,
1030 16 * 1024 * 1024 * 1024,
1031 )?;
1032 let hash_hex = store.put_blob(b"audit-serving bytes")?;
1033 store.force_sync()?;
1034 drop(store);
1035
1036 let before = pool_data_file_snapshot(&data_dir);
1037 assert!(
1038 before.len() >= 2,
1039 "expected Pool catalog and at least one member data file"
1040 );
1041
1042 let output = Command::new(std::env::current_exe()?)
1043 .arg("--ignored")
1044 .arg("--exact")
1045 .arg("server::tests::pool_audit_read_only_server_subprocess")
1046 .arg("--nocapture")
1047 .env(AUDIT_SERVER_CHILD_DATA_ENV, &data_dir)
1048 .env(AUDIT_SERVER_CHILD_HASH_ENV, &hash_hex)
1049 .env(hashtree_lmdb::POOL_AUDIT_READ_ONLY_ENV, "1")
1050 .env_remove("HTREE_LMDB_HOT_BLOB_DIR")
1051 .env_remove("HTREE_LMDB_HOT_BLOB_LEGACY_DIR")
1052 .env_remove("HTREE_LMDB_HOT_EXTERNAL_BLOB_DIR")
1053 .env_remove("HTREE_LMDB_LEGACY_EXTERNAL_BLOB_DIR")
1054 .env("RUST_TEST_THREADS", "1")
1055 .output()?;
1056 assert!(
1057 output.status.success(),
1058 "audit-serving subprocess failed: stdout={} stderr={}",
1059 String::from_utf8_lossy(&output.stdout),
1060 String::from_utf8_lossy(&output.stderr)
1061 );
1062
1063 let after = pool_data_file_snapshot(&data_dir);
1064 assert_eq!(
1065 after, before,
1066 "audit-serving process mutated Pool data files"
1067 );
1068 Ok(())
1069 }
1070
1071 #[async_trait]
1072 impl EventBus for StaticProvider {
1073 async fn publish(
1074 &self,
1075 _event: VerifiedEvent,
1076 _source: EventSource,
1077 ) -> nostr_pubsub::Result<PublishReport> {
1078 self.publishes.fetch_add(1, Ordering::Relaxed);
1079 Ok(PublishReport {
1080 accepted: true,
1081 priority: 0,
1082 reason: None,
1083 })
1084 }
1085
1086 async fn query(
1087 &self,
1088 _filters: Vec<nostr_pubsub::Filter>,
1089 _options: QueryOptions,
1090 ) -> nostr_pubsub::Result<QueryReport> {
1091 self.queries.fetch_add(1, Ordering::Relaxed);
1092 Ok(QueryReport {
1093 events: self
1094 .event
1095 .clone()
1096 .map(|event| QueryEvent {
1097 event,
1098 source: EventSource::fips_endpoint("browser-router"),
1099 priority: 0,
1100 })
1101 .into_iter()
1102 .collect(),
1103 })
1104 }
1105 }
1106
1107 impl PubsubProvider for StaticProvider {
1108 fn mode(&self) -> PubsubProviderMode {
1109 self.mode
1110 }
1111 }
1112
1113 #[test]
1114 fn resource_exhaustion_errors_are_fatal_accept_errors() {
1115 assert!(is_resource_exhaustion_error(&io::Error::from_raw_os_error(
1116 libc::EMFILE
1117 )));
1118 assert!(is_resource_exhaustion_error(&io::Error::from_raw_os_error(
1119 libc::ENFILE
1120 )));
1121 assert!(!is_resource_exhaustion_error(
1122 &io::Error::from_raw_os_error(libc::ECONNRESET)
1123 ));
1124 }
1125
1126 #[test]
1127 fn upload_queue_semaphores_fit_all_tokio_targets() -> Result<()> {
1128 let temp_dir = TempDir::new()?;
1129 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1130 let server = HashtreeServer::new(store, "127.0.0.1:0".to_string());
1131
1132 assert_eq!(
1133 server.state.optimistic_upload_queue.available_permits(),
1134 256 * 1024 * 1024
1135 );
1136 assert_eq!(
1137 server
1138 .state
1139 .blossom_upload_replica_queue
1140 .available_permits(),
1141 256 * 1024 * 1024
1142 );
1143 assert_eq!(
1144 bounded_upload_queue_bytes(u64::MAX),
1145 tokio::sync::Semaphore::MAX_PERMITS
1146 );
1147 Ok(())
1148 }
1149
1150 #[test]
1151 fn server_builder_seeds_initial_tree_roots() -> Result<()> {
1152 let temp_dir = TempDir::new()?;
1153 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1154 let hash = from_hex("1111111111111111111111111111111111111111111111111111111111111111")?;
1155 let key = from_hex("2222222222222222222222222222222222222222222222222222222222222222")?;
1156 let cid = Cid {
1157 hash,
1158 key: Some(key),
1159 };
1160
1161 let server = HashtreeServer::new(store, "127.0.0.1:0".to_string())
1162 .with_cached_tree_roots(vec![("npub1example/sites".to_string(), cid.clone())]);
1163 let cached = server
1164 .state
1165 .tree_root_cache
1166 .lock()
1167 .unwrap()
1168 .get("npub1example/sites")
1169 .cloned()
1170 .expect("seeded root");
1171
1172 assert_eq!(cached.cid, cid);
1173 assert_eq!(cached.source, "embedded-bootstrap");
1174 assert!(cached.root_event.is_none());
1175 Ok(())
1176 }
1177
1178 #[test]
1179 fn iris_localhost_hosts_map_to_existing_tree_routes() {
1180 assert_eq!(
1181 resolve_virtual_tree_host("NHASH1EXAMPLE.iris.localhost:8080"),
1182 Some("/htree/nhash1example".to_string())
1183 );
1184 assert_eq!(
1185 resolve_virtual_tree_host("audio.NPUB1EXAMPLE.iris.localhost:8080"),
1186 Some("/htree/npub1example/audio".to_string())
1187 );
1188 assert_eq!(
1189 resolve_virtual_tree_host("audio.extra.npub1example.iris.localhost"),
1190 None
1191 );
1192 assert_eq!(
1193 resolve_virtual_tree_host("nhash1example.htree.localhost"),
1194 None
1195 );
1196 }
1197
1198 #[tokio::test]
1199 async fn test_server_serve_file() -> Result<()> {
1200 let temp_dir = TempDir::new()?;
1201 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1202
1203 let test_file = temp_dir.path().join("test.txt");
1205 std::fs::write(&test_file, b"Hello, Hashtree!")?;
1206
1207 let cid = store.upload_file(&test_file)?;
1208 let hash = from_hex(&cid)?;
1209
1210 let content = store.get_file(&hash)?;
1212 assert!(content.is_some());
1213 assert_eq!(content.unwrap(), b"Hello, Hashtree!");
1214
1215 Ok(())
1216 }
1217
1218 #[tokio::test]
1219 async fn test_server_list_pins() -> Result<()> {
1220 let temp_dir = TempDir::new()?;
1221 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1222
1223 let test_file = temp_dir.path().join("test.txt");
1224 std::fs::write(&test_file, b"Test")?;
1225
1226 let cid = store.upload_file(&test_file)?;
1227 let hash = from_hex(&cid)?;
1228
1229 let pins = store.list_pins_raw()?;
1230 assert_eq!(pins.len(), 1);
1231 assert_eq!(pins[0], hash);
1232
1233 Ok(())
1234 }
1235
1236 async fn spawn_test_server(
1237 store: Arc<HashtreeStore>,
1238 ) -> Result<(u16, tokio::task::JoinHandle<Result<()>>)> {
1239 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
1240 let port = listener.local_addr()?.port();
1241 let server = HashtreeServer::new(store, "127.0.0.1:0".to_string());
1242 let handle =
1243 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
1244 Ok((port, handle))
1245 }
1246
1247 async fn spawn_test_server_with_auth(
1248 store: Arc<HashtreeStore>,
1249 ) -> Result<(u16, tokio::task::JoinHandle<Result<()>>)> {
1250 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
1251 let port = listener.local_addr()?.port();
1252 let server = HashtreeServer::new(store, "127.0.0.1:0".to_string())
1253 .with_auth("test-user".to_string(), "test-password".to_string());
1254 let handle =
1255 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
1256 Ok((port, handle))
1257 }
1258
1259 async fn encrypted_test_directory(store: &HashtreeStore) -> Result<(Cid, Cid, Vec<u8>)> {
1260 let tree = HashTree::new(
1261 HashTreeConfig::new(store.store_arc())
1262 .with_chunk_size(4)
1263 .with_max_links(2),
1264 );
1265 let content = b"encrypted descendant content spanning chunks".to_vec();
1266 let (file_cid, file_size) = tree.put(&content).await?;
1267 let nested_cid = tree
1268 .put_directory(vec![DirEntry::from_cid("post.json", &file_cid)
1269 .with_size(file_size)
1270 .with_link_type(LinkType::File)])
1271 .await?;
1272 let root_cid = tree
1273 .put_directory(vec![
1274 DirEntry::from_cid("events", &nested_cid).with_link_type(LinkType::Dir)
1275 ])
1276 .await?;
1277 Ok((root_cid, file_cid, content))
1278 }
1279
1280 async fn spawn_test_server_with_nostr_relay(
1281 store: Arc<HashtreeStore>,
1282 relay: Arc<NostrRelay>,
1283 ) -> Result<(u16, tokio::task::JoinHandle<Result<()>>)> {
1284 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
1285 let port = listener.local_addr()?.port();
1286 let server = HashtreeServer::new(store, "127.0.0.1:0".to_string()).with_nostr_relay(relay);
1287 let handle =
1288 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
1289 Ok((port, handle))
1290 }
1291
1292 #[tokio::test]
1293 async fn unauthenticated_native_store_mutation_is_rejected() -> Result<()> {
1294 let temp_dir = TempDir::new()?;
1295 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1296 let body = b"unauthorized native store write";
1297 let hash = sha256(body);
1298 let hash_hex = to_hex(&hash);
1299 let (port, handle) = spawn_test_server(Arc::clone(&store)).await?;
1300
1301 let response = reqwest::Client::new()
1302 .put(format!("http://127.0.0.1:{port}/__iris/store/{hash_hex}"))
1303 .body(body.to_vec())
1304 .send()
1305 .await?;
1306
1307 assert_eq!(response.status(), reqwest::StatusCode::FORBIDDEN);
1308 assert!(store.get_blob(&hash)?.is_none());
1309 handle.abort();
1310 Ok(())
1311 }
1312
1313 #[tokio::test]
1314 async fn unauthenticated_cache_tree_root_is_rejected() -> Result<()> {
1315 let temp_dir = TempDir::new()?;
1316 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1317 let (port, handle) = spawn_test_server(store).await?;
1318
1319 let response = reqwest::Client::new()
1320 .post(format!("http://127.0.0.1:{port}/api/cache-tree-root"))
1321 .json(&json!({
1322 "npub": "npub1example",
1323 "treeName": "video",
1324 "hash": "988db3f24dc222715f1c1e1fa5876690d3147122243d72d85fd44283867cd61a",
1325 "visibility": "public"
1326 }))
1327 .send()
1328 .await?;
1329
1330 assert_eq!(response.status(), reqwest::StatusCode::FORBIDDEN);
1331 handle.abort();
1332 Ok(())
1333 }
1334
1335 #[tokio::test]
1336 async fn pin_tree_requires_configured_valid_auth() -> Result<()> {
1337 let temp_dir = TempDir::new()?;
1338 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1339 let nhash = nhash_encode(&sha256(b"missing"))?;
1340 let client = reqwest::Client::new();
1341
1342 let (port, handle) = spawn_test_server(Arc::clone(&store)).await?;
1343 let response = client
1344 .post(format!("http://127.0.0.1:{port}/api/pin-tree"))
1345 .json(&json!({"nhash": nhash}))
1346 .send()
1347 .await?;
1348 assert_eq!(response.status(), reqwest::StatusCode::FORBIDDEN);
1349 handle.abort();
1350
1351 let (port, handle) = spawn_test_server_with_auth(store).await?;
1352 let response = client
1353 .post(format!("http://127.0.0.1:{port}/api/pin-tree"))
1354 .json(&json!({"nhash": nhash}))
1355 .send()
1356 .await?;
1357 assert_eq!(response.status(), reqwest::StatusCode::UNAUTHORIZED);
1358 handle.abort();
1359 Ok(())
1360 }
1361
1362 #[tokio::test]
1363 async fn pin_tree_rejects_malformed_noncanonical_and_missing_roots() -> Result<()> {
1364 let temp_dir = TempDir::new()?;
1365 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1366 let missing_hash = sha256(b"not stored");
1367 let canonical = nhash_encode(&missing_hash)?;
1368 let (port, handle) = spawn_test_server_with_auth(Arc::clone(&store)).await?;
1369 let client = reqwest::Client::new();
1370
1371 for nhash in ["not-an-nhash".to_string(), format!("hashtree:{canonical}")] {
1372 let response = client
1373 .post(format!("http://127.0.0.1:{port}/api/pin-tree"))
1374 .basic_auth("test-user", Some("test-password"))
1375 .json(&json!({"nhash": nhash}))
1376 .send()
1377 .await?;
1378 assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST);
1379 }
1380
1381 let response = client
1382 .post(format!("http://127.0.0.1:{port}/api/pin-tree"))
1383 .basic_auth("test-user", Some("test-password"))
1384 .json(&json!({"nhash": canonical}))
1385 .send()
1386 .await?;
1387 assert_eq!(response.status(), reqwest::StatusCode::NOT_FOUND);
1388 assert!(!store.is_pinned(&missing_hash)?);
1389 assert!(store.get_tree_meta(&missing_hash)?.is_none());
1390 handle.abort();
1391 Ok(())
1392 }
1393
1394 #[tokio::test]
1395 async fn pin_tree_indexes_encrypted_descendants_and_is_idempotent() -> Result<()> {
1396 let temp_dir = TempDir::new()?;
1397 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1398 let (root_cid, file_cid, expected_content) = encrypted_test_directory(&store).await?;
1399 let nhash = nhash_encode_full(&NHashData {
1400 hash: root_cid.hash,
1401 decrypt_key: root_cid.key,
1402 })?;
1403 let (port, handle) = spawn_test_server_with_auth(Arc::clone(&store)).await?;
1404 let client = reqwest::Client::new();
1405
1406 for expected_already_pinned in [false, true] {
1407 let response = client
1408 .post(format!("http://127.0.0.1:{port}/api/pin-tree"))
1409 .basic_auth("test-user", Some("test-password"))
1410 .json(&json!({"nhash": nhash}))
1411 .send()
1412 .await?;
1413 assert_eq!(response.status(), reqwest::StatusCode::OK);
1414 let body: serde_json::Value = response.json().await?;
1415 assert_eq!(body["already_pinned"], expected_already_pinned);
1416 assert!(body["indexed_hashes"].as_u64().unwrap_or_default() > 2);
1417 }
1418
1419 assert!(store.is_pinned(&root_cid.hash)?);
1420 assert_eq!(store.list_pins_raw()?, vec![root_cid.hash]);
1421 assert_eq!(store.list_indexed_trees()?.len(), 1);
1422 let tree = HashTree::new(HashTreeConfig::new(store.store_arc()));
1423 assert_eq!(
1424 tree.get(&file_cid, None).await?,
1425 Some(expected_content),
1426 "encrypted descendant remains readable"
1427 );
1428 handle.abort();
1429 Ok(())
1430 }
1431
1432 #[tokio::test]
1433 async fn pin_tree_missing_descendant_leaves_no_pin_or_index() -> Result<()> {
1434 let temp_dir = TempDir::new()?;
1435 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1436 let (root_cid, file_cid, _) = encrypted_test_directory(&store).await?;
1437 assert!(store.router().delete_sync(&file_cid.hash)?);
1438 let nhash = nhash_encode_full(&NHashData {
1439 hash: root_cid.hash,
1440 decrypt_key: root_cid.key,
1441 })?;
1442 let (port, handle) = spawn_test_server_with_auth(Arc::clone(&store)).await?;
1443
1444 let response = reqwest::Client::new()
1445 .post(format!("http://127.0.0.1:{port}/api/pin-tree"))
1446 .basic_auth("test-user", Some("test-password"))
1447 .json(&json!({"nhash": nhash}))
1448 .send()
1449 .await?;
1450
1451 assert_eq!(response.status(), reqwest::StatusCode::UNPROCESSABLE_ENTITY);
1452 assert!(!store.is_pinned(&root_cid.hash)?);
1453 assert!(store.get_tree_meta(&root_cid.hash)?.is_none());
1454 handle.abort();
1455 Ok(())
1456 }
1457
1458 #[tokio::test]
1459 async fn pin_tree_resource_failure_leaves_no_pin_or_index() -> Result<()> {
1460 let temp_dir = TempDir::new()?;
1461 let store = Arc::new(HashtreeStore::with_options(
1462 temp_dir.path().join("db"),
1463 None,
1464 8,
1465 )?);
1466 let tree = HashTree::new(HashTreeConfig::new(store.store_arc()));
1467 let (cid, _) = tree.put(b"larger than the pin tree byte budget").await?;
1468 let nhash = nhash_encode_full(&NHashData {
1469 hash: cid.hash,
1470 decrypt_key: cid.key,
1471 })?;
1472 let (port, handle) = spawn_test_server_with_auth(Arc::clone(&store)).await?;
1473
1474 let response = reqwest::Client::new()
1475 .post(format!("http://127.0.0.1:{port}/api/pin-tree"))
1476 .basic_auth("test-user", Some("test-password"))
1477 .json(&json!({"nhash": nhash}))
1478 .send()
1479 .await?;
1480
1481 assert_eq!(response.status(), reqwest::StatusCode::PAYLOAD_TOO_LARGE);
1482 assert!(!store.is_pinned(&cid.hash)?);
1483 assert!(store.get_tree_meta(&cid.hash)?.is_none());
1484 handle.abort();
1485 Ok(())
1486 }
1487
1488 #[tokio::test]
1489 async fn unauthenticated_upload_batch_is_rejected_before_json_extraction() -> Result<()> {
1490 let temp_dir = TempDir::new()?;
1491 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1492 let (port, handle) = spawn_test_server(store).await?;
1493
1494 let response = reqwest::Client::new()
1495 .post(format!("http://127.0.0.1:{port}/upload/batch"))
1496 .header(reqwest::header::CONTENT_TYPE, "application/json")
1497 .body("not valid json")
1498 .send()
1499 .await?;
1500
1501 assert_eq!(response.status(), reqwest::StatusCode::UNAUTHORIZED);
1502 handle.abort();
1503 Ok(())
1504 }
1505
1506 #[tokio::test]
1507 async fn unauthenticated_single_upload_is_rejected_before_body_limit() -> Result<()> {
1508 let temp_dir = TempDir::new()?;
1509 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1510 let (port, handle) = spawn_test_server(store).await?;
1511 let declared_bytes = blossom::MAX_SINGLE_UPLOAD_BODY_BYTES + 1;
1512 let mut stream = tokio::net::TcpStream::connect(("127.0.0.1", port)).await?;
1513 stream
1514 .write_all(
1515 format!(
1516 "PUT /upload HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nContent-Length: {declared_bytes}\r\nConnection: close\r\n\r\n"
1517 )
1518 .as_bytes(),
1519 )
1520 .await?;
1521
1522 let mut status_line = String::new();
1525 BufReader::new(stream).read_line(&mut status_line).await?;
1526 assert_eq!(status_line.trim_end(), "HTTP/1.1 401 Unauthorized");
1527 handle.abort();
1528 Ok(())
1529 }
1530
1531 #[tokio::test]
1532 async fn upload_options_preflight_does_not_require_auth() -> Result<()> {
1533 let temp_dir = TempDir::new()?;
1534 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1535 let (port, handle) = spawn_test_server(store).await?;
1536
1537 let response = reqwest::Client::new()
1538 .request(
1539 reqwest::Method::OPTIONS,
1540 format!("http://127.0.0.1:{port}/upload"),
1541 )
1542 .send()
1543 .await?;
1544
1545 assert_eq!(response.status(), reqwest::StatusCode::NO_CONTENT);
1546 handle.abort();
1547 Ok(())
1548 }
1549
1550 #[tokio::test]
1551 async fn virtual_tree_hosts_serve_root_assets_and_spa_fallbacks() -> Result<()> {
1552 clear_virtual_tree_hosts_for_test();
1553
1554 let temp_dir = TempDir::new()?;
1555 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1556 let tree = HashTree::new(HashTreeConfig::new(store.store_arc()).public());
1557
1558 let (index_cid, _) = tree
1559 .put(b"<!doctype html><title>Virtual host ok</title>")
1560 .await?;
1561 let (favicon_cid, _) = tree.put(b"ico").await?;
1562 let (main_js_cid, _) = tree.put(b"console.log('ok');").await?;
1563 let assets_dir = tree
1564 .put_directory(vec![
1565 DirEntry::from_cid("main.js", &main_js_cid).with_link_type(LinkType::File)
1566 ])
1567 .await?;
1568 let root_cid = tree
1569 .put_directory(vec![
1570 DirEntry::from_cid("index.html", &index_cid).with_link_type(LinkType::File),
1571 DirEntry::from_cid("favicon.ico", &favicon_cid).with_link_type(LinkType::File),
1572 DirEntry::from_cid("assets", &assets_dir).with_link_type(LinkType::Dir),
1573 ])
1574 .await?;
1575 let nhash = nhash_encode(&root_cid.hash)?;
1576 let host = "tree-test.htree.localhost";
1577 register_virtual_tree_host(host, &format!("/htree/{nhash}"));
1578
1579 let (port, handle) = spawn_test_server(store).await?;
1580 let base_url = format!("http://127.0.0.1:{port}");
1581 let host_header = format!("{host}:{port}");
1582 let client = reqwest::Client::new();
1583
1584 let root_response = client
1585 .get(format!("{base_url}/"))
1586 .header("Host", &host_header)
1587 .header("Accept", "text/html")
1588 .send()
1589 .await?;
1590 assert_eq!(root_response.status(), reqwest::StatusCode::OK);
1591 assert_eq!(
1592 root_response.bytes().await?.as_ref(),
1593 b"<!doctype html><title>Virtual host ok</title>"
1594 );
1595
1596 let favicon_response = client
1597 .get(format!("{base_url}/favicon.ico"))
1598 .header("Host", &host_header)
1599 .send()
1600 .await?;
1601 assert_eq!(favicon_response.status(), reqwest::StatusCode::OK);
1602 assert_eq!(favicon_response.bytes().await?.as_ref(), b"ico");
1603
1604 let js_response = client
1605 .get(format!("{base_url}/assets/main.js"))
1606 .header("Host", &host_header)
1607 .send()
1608 .await?;
1609 assert_eq!(js_response.status(), reqwest::StatusCode::OK);
1610 assert_eq!(js_response.bytes().await?.as_ref(), b"console.log('ok');");
1611
1612 let profile_response = client
1613 .get(format!("{base_url}/users/npub1example"))
1614 .header("Host", &host_header)
1615 .header("Accept", "text/html")
1616 .send()
1617 .await?;
1618 assert_eq!(profile_response.status(), reqwest::StatusCode::OK);
1619 assert_eq!(
1620 profile_response.bytes().await?.as_ref(),
1621 b"<!doctype html><title>Virtual host ok</title>"
1622 );
1623
1624 handle.abort();
1625 clear_virtual_tree_hosts_for_test();
1626
1627 Ok(())
1628 }
1629
1630 #[tokio::test]
1631 async fn iris_localhost_hosts_serve_immutable_and_named_sites() -> Result<()> {
1632 let temp_dir = TempDir::new()?;
1633 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1634 let tree = HashTree::new(HashTreeConfig::new(store.store_arc()).public());
1635
1636 let (index_cid, _) = tree
1637 .put(b"<!doctype html><title>Iris localhost ok</title>")
1638 .await?;
1639 let (main_js_cid, _) = tree.put(b"console.log('iris localhost');").await?;
1640 let assets_dir = tree
1641 .put_directory(vec![
1642 DirEntry::from_cid("main.js", &main_js_cid).with_link_type(LinkType::File)
1643 ])
1644 .await?;
1645 let root_cid = tree
1646 .put_directory(vec![
1647 DirEntry::from_cid("index.html", &index_cid).with_link_type(LinkType::File),
1648 DirEntry::from_cid("assets", &assets_dir).with_link_type(LinkType::Dir),
1649 ])
1650 .await?;
1651
1652 let nhash = nhash_encode(&root_cid.hash)?;
1653 let npub = Keys::generate().public_key().to_bech32()?;
1654 let site = "audio";
1655 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
1656 let port = listener.local_addr()?.port();
1657 let server = HashtreeServer::new(Arc::clone(&store), "127.0.0.1:0".to_string())
1658 .with_public_plaintext_reads(true)
1659 .with_cached_tree_roots(vec![(format!("{npub}/{site}"), root_cid)]);
1660 let handle =
1661 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
1662 let client = reqwest::Client::new();
1663 let base_url = format!("http://127.0.0.1:{port}");
1664
1665 for host in [
1666 format!("{nhash}.iris.localhost:{port}"),
1667 format!("{site}.{npub}.iris.localhost:{port}"),
1668 ] {
1669 let root_response = client
1670 .get(format!("{base_url}/"))
1671 .header("Host", &host)
1672 .header("Accept", "text/html")
1673 .send()
1674 .await?;
1675 assert_eq!(root_response.status(), reqwest::StatusCode::OK, "{host}");
1676 assert_eq!(
1677 root_response.bytes().await?.as_ref(),
1678 b"<!doctype html><title>Iris localhost ok</title>"
1679 );
1680
1681 let asset_response = client
1682 .get(format!("{base_url}/assets/main.js"))
1683 .header("Host", &host)
1684 .send()
1685 .await?;
1686 assert_eq!(asset_response.status(), reqwest::StatusCode::OK, "{host}");
1687 assert_eq!(
1688 asset_response.bytes().await?.as_ref(),
1689 b"console.log('iris localhost');"
1690 );
1691 }
1692
1693 handle.abort();
1694 Ok(())
1695 }
1696
1697 #[tokio::test]
1698 async fn named_iris_site_resolves_from_fips_event_provider_without_relays() -> Result<()> {
1699 let temp_dir = TempDir::new()?;
1700 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1701 let tree = HashTree::new(HashTreeConfig::new(store.store_arc()).public());
1702 let (index_cid, _) = tree.put(b"fips event provider site").await?;
1703 let root_cid = tree
1704 .put_directory(vec![
1705 DirEntry::from_cid("index.html", &index_cid).with_link_type(LinkType::File)
1706 ])
1707 .await?;
1708 let owner = Keys::generate();
1709 let site = "radio";
1710 let root_event = EventBuilder::new(Kind::Custom(30064), "")
1711 .tags(vec![
1712 nostr::Tag::identifier(site),
1713 nostr::Tag::custom(nostr::TagKind::custom("l"), vec!["hashtree".to_string()]),
1714 nostr::Tag::custom(nostr::TagKind::custom("hash"), vec![to_hex(&root_cid.hash)]),
1715 ])
1716 .sign_with_keys(&owner)?;
1717 let provider = Arc::new(StaticProvider {
1718 event: Some(VerifiedEvent::try_from(root_event)?),
1719 queries: AtomicUsize::new(0),
1720 publishes: AtomicUsize::new(0),
1721 mode: PubsubProviderMode::LocalOnly,
1722 });
1723
1724 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
1725 let port = listener.local_addr()?.port();
1726 let server = HashtreeServer::new(Arc::clone(&store), "127.0.0.1:0".to_string())
1727 .with_public_plaintext_reads(true)
1728 .with_nostr_provider(provider.clone());
1729 assert!(server.state.nostr_relay_urls.is_empty());
1730 assert_eq!(
1731 server.state.nostr_provider.as_ref().unwrap().mode(),
1732 PubsubProviderMode::LocalOnly
1733 );
1734 let handle =
1735 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
1736
1737 let npub = owner.public_key().to_bech32()?;
1738 let response = reqwest::Client::new()
1739 .get(format!("http://127.0.0.1:{port}/"))
1740 .header("Host", format!("{site}.{npub}.iris.localhost"))
1741 .header("Accept", "text/html")
1742 .send()
1743 .await?;
1744
1745 assert_eq!(response.status(), reqwest::StatusCode::OK);
1746 assert_eq!(
1747 response.bytes().await?.as_ref(),
1748 b"fips event provider site"
1749 );
1750 assert_eq!(provider.queries.load(Ordering::Relaxed), 1);
1751 handle.abort();
1752 Ok(())
1753 }
1754
1755 #[tokio::test]
1756 async fn loopback_event_publication_uses_configured_provider() -> Result<()> {
1757 let temp_dir = TempDir::new()?;
1758 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1759 let keys = Keys::generate();
1760 let event = EventBuilder::new(Kind::TextNote, "provider test").sign_with_keys(&keys)?;
1761 let provider = Arc::new(StaticProvider {
1762 event: Some(VerifiedEvent::try_from(event.clone())?),
1763 queries: AtomicUsize::new(0),
1764 publishes: AtomicUsize::new(0),
1765 mode: PubsubProviderMode::LocalOnly,
1766 });
1767 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
1768 let port = listener.local_addr()?.port();
1769 let server = HashtreeServer::new(store, "127.0.0.1:0".to_string())
1770 .with_nostr_provider(provider.clone());
1771 assert!(server.state.nostr_relay_urls.is_empty());
1772 let handle =
1773 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
1774
1775 let response = reqwest::Client::new()
1776 .post(format!("http://127.0.0.1:{port}/api/nostr/events"))
1777 .json(&event)
1778 .send()
1779 .await?;
1780
1781 assert_eq!(response.status(), reqwest::StatusCode::ACCEPTED);
1782 assert_eq!(provider.publishes.load(Ordering::Relaxed), 1);
1783 handle.abort();
1784 Ok(())
1785 }
1786
1787 #[tokio::test]
1788 async fn loopback_root_publication_is_immediately_resolvable_without_provider_replay(
1789 ) -> Result<()> {
1790 let temp_dir = TempDir::new()?;
1791 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1792 let keys = Keys::generate();
1793 let tree_name = "webvm-e2e";
1794 let root_hash = "11".repeat(32);
1795 let encryption_key = "22".repeat(32);
1796 let event = EventBuilder::new(Kind::Custom(30064), &root_hash)
1797 .tags(vec![
1798 nostr::Tag::identifier(tree_name),
1799 nostr::Tag::custom(nostr::TagKind::custom("l"), vec!["hashtree".to_string()]),
1800 nostr::Tag::custom(nostr::TagKind::custom("l"), vec!["git".to_string()]),
1801 nostr::Tag::custom(nostr::TagKind::custom("hash"), vec![root_hash.clone()]),
1802 nostr::Tag::custom(nostr::TagKind::custom("key"), vec![encryption_key.clone()]),
1803 ])
1804 .sign_with_keys(&keys)?;
1805 let provider = Arc::new(StaticProvider {
1806 event: None,
1807 queries: AtomicUsize::new(0),
1808 publishes: AtomicUsize::new(0),
1809 mode: PubsubProviderMode::LocalOnly,
1810 });
1811 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
1812 let port = listener.local_addr()?.port();
1813 let server = HashtreeServer::new(store, "127.0.0.1:0".to_string())
1814 .with_nostr_provider(provider.clone());
1815 assert!(server.state.nostr_relay_urls.is_empty());
1816 let handle =
1817 tokio::spawn(async move { server.run_with_listener(listener).await.map(|_| ()) });
1818 let client = reqwest::Client::new();
1819
1820 let publish = client
1821 .post(format!("http://127.0.0.1:{port}/api/nostr/events"))
1822 .json(&event)
1823 .send()
1824 .await?;
1825 assert_eq!(publish.status(), reqwest::StatusCode::ACCEPTED);
1826
1827 let npub = keys.public_key().to_bech32()?;
1828 let resolved = client
1829 .get(format!(
1830 "http://127.0.0.1:{port}/api/nostr/resolve/{npub}/{tree_name}?refresh=1"
1831 ))
1832 .send()
1833 .await?
1834 .json::<serde_json::Value>()
1835 .await?;
1836
1837 assert_eq!(resolved["hash"], root_hash);
1838 assert_eq!(resolved["key_tag"], encryption_key);
1839 assert_eq!(resolved["source"], "local-relay");
1840 assert_eq!(provider.publishes.load(Ordering::Relaxed), 1);
1841 assert_eq!(provider.queries.load(Ordering::Relaxed), 1);
1842 handle.abort();
1843 Ok(())
1844 }
1845
1846 #[tokio::test]
1847 async fn nostr_profile_route_returns_latest_metadata_event() -> Result<()> {
1848 let temp_dir = TempDir::new()?;
1849 let store = Arc::new(HashtreeStore::new(temp_dir.path().join("db"))?);
1850 let graph_store = {
1851 let _guard = crate::socialgraph::test_lock().await;
1852 crate::socialgraph::open_test_social_graph_store_with_mapsize(
1853 &temp_dir.path().join("relay-db"),
1854 Some(128 * 1024 * 1024),
1855 )?
1856 };
1857 let backend: Arc<dyn crate::socialgraph::SocialGraphBackend> = graph_store;
1858 let relay = Arc::new(NostrRelay::new(
1859 backend,
1860 temp_dir.path().to_path_buf(),
1861 HashSet::new(),
1862 None,
1863 NostrRelayConfig {
1864 spambox_db_max_bytes: 0,
1865 ..Default::default()
1866 },
1867 )?);
1868
1869 let author = Keys::generate();
1870 let older = EventBuilder::new(
1871 Kind::Metadata,
1872 json!({ "name": "older", "about": "before" }).to_string(),
1873 )
1874 .custom_created_at(Timestamp::from_secs(10))
1875 .sign_with_keys(&author)?;
1876 let newer = EventBuilder::new(
1877 Kind::Metadata,
1878 json!({ "name": "newer", "about": "after" }).to_string(),
1879 )
1880 .custom_created_at(Timestamp::from_secs(20))
1881 .sign_with_keys(&author)?;
1882
1883 relay.ingest_trusted_event(older).await?;
1884 relay.ingest_trusted_event(newer.clone()).await?;
1885
1886 let (port, handle) = spawn_test_server_with_nostr_relay(store, relay).await?;
1887 let response = reqwest::get(format!(
1888 "http://127.0.0.1:{port}/api/nostr/profile/{}",
1889 author.public_key().to_hex()
1890 ))
1891 .await?;
1892
1893 assert_eq!(response.status(), reqwest::StatusCode::OK);
1894 let payload: serde_json::Value = response.json().await?;
1895 assert_eq!(payload["profile"]["name"].as_str(), Some("newer"),);
1896 assert_eq!(payload["profile"]["about"].as_str(), Some("after"));
1897 assert_eq!(payload["created_at"].as_u64(), Some(20));
1898 let expected_event_id = newer.id.to_hex();
1899 assert_eq!(
1900 payload["event_id"].as_str(),
1901 Some(expected_event_id.as_str())
1902 );
1903
1904 handle.abort();
1905 Ok(())
1906 }
1907}