Skip to main content

hashtree_cli/server/
mod.rs

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, // 5 MB default
203                public_writes: true,               // Allow anyone with valid Nostr auth by default
204                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(), // No pubkeys allowed by default (use public_writes)
212                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    /// Set maximum upload size for Blossom uploads
254    pub fn with_max_upload_bytes(mut self, bytes: usize) -> Self {
255        self.state.max_upload_bytes = bytes;
256        self
257    }
258
259    /// Set whether to allow public writes (anyone with valid Nostr auth)
260    /// When false, only social graph members can write
261    pub fn with_public_writes(mut self, public: bool) -> Self {
262        self.state.public_writes = public;
263        self
264    }
265
266    /// Set whether mutable npub routes serve plaintext for unapproved pubkeys
267    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    /// Set allowed pubkeys for blossom write access (hex format)
319    pub fn with_allowed_pubkeys(mut self, pubkeys: HashSet<String>) -> Self {
320        self.state.allowed_pubkeys = pubkeys;
321        self
322    }
323
324    /// Set upstream Blossom servers for cascade fetching
325    pub fn with_upstream_blossom(mut self, servers: Vec<String>) -> Self {
326        self.state.upstream_blossom = servers;
327        self
328    }
329
330    /// Set write-behind Blossom servers for blobs accepted by this server.
331    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    /// Set social graph access control
355    pub fn with_social_graph(mut self, sg: Arc<socialgraph::SocialGraphAccessControl>) -> Self {
356        self.state.social_graph = Some(sg);
357        self
358    }
359
360    /// Configure social graph snapshot export (store handle + root)
361    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    /// Set Nostr relay state (shared for /ws and WebRTC)
374    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    /// Set active upstream Nostr relays for HTTP resolver operations.
385    pub fn with_nostr_relay_urls(mut self, relays: Vec<String>) -> Self {
386        self.state.nostr_relay_urls = relays;
387        self
388    }
389
390    /// Seed mutable root cache entries before the server starts.
391    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    /// Merge extra routes into the daemon router (e.g. Tauri embeds /nip07).
411    pub fn with_extra_routes(mut self, routes: Router<AppState>) -> Self {
412        self.extra_routes = Some(routes);
413        self
414    }
415
416    /// Apply a CORS layer to all routes (used by embedded clients like Tauri).
417    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        // `tokio::net::TcpListener::bind` inherits the platform's small
442        // default listen backlog (128 on Linux). A daemon restart can attract
443        // more reconnecting CDN/tunnel sockets than that before the first
444        // requests are accepted, leaving otherwise healthy clients stalled.
445        // Re-listening updates the queue limit without replacing the bound
446        // socket; the kernel still applies its configured `somaxconn` cap.
447        SockRef::from(&listener).listen(TCP_LISTEN_BACKLOG)?;
448        let local_addr = listener.local_addr()?;
449
450        // Public endpoints (no auth required)
451        // Note: /:id serves raw SHA256 blobs only. Logical tree/file assembly
452        // stays on explicitly configured mutable routes such as approved npub.
453        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            // /htree/nhash1...[/path] - content-addressed (immutable)
467            .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            // /htree/npub1.../tree[/path] - mutable (resolver-backed)
471            .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            // Nostr resolver endpoints - resolve npub/treename to content
478            .route("/n/:pubkey/:treename", get(handlers::resolve_and_serve))
479            // Direct npub route (clients should parse nhash and request by hex hash)
480            .route("/npub1:rest", get(handlers::serve_npub))
481            .route("/npub1:rest/*path", get(handlers::serve_npub))
482            // Blossom endpoints (BUD-01, BUD-02)
483            .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            // Hashtree API endpoints
535            .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            // Resolver API endpoints
549            .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        // Protected endpoints (require auth if enabled)
568        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        // Internal mutating endpoints require configured Basic auth. These
580        // routes stay closed even when optional API auth is disabled.
581        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)); // 10GB limit
613
614        if let Some(extra) = self.extra_routes {
615            app = app.merge(extra.with_state(state.clone()));
616        }
617
618        // This gate is deliberately outside every route-level auth/body
619        // layer, including caller-supplied routes. Audit mode rejects
620        // mutation ingress before any handler can alter blob or metadata
621        // state. Status metrics remain outermost so maintenance responses are
622        // still observable.
623        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        // Create and upload a test file
1204        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        // Verify we can get it
1211        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        // Send headers only: receiving a response proves authentication ran
1523        // before either buffering the declared body or enforcing its limit.
1524        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}