Skip to main content

arete_server/
http_health.rs

1use crate::http::transactions::{self, TransactionState};
2use crate::{
3    config::{RuntimePlan, TransactionConfig},
4    health::HealthMonitor,
5    ProgramRuntimeCatalog, ProgramRuntimeDefinition,
6};
7use anyhow::Result;
8use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
9use base64::Engine as _;
10use dashmap::DashMap;
11use http_body_util::BodyExt;
12use http_body_util::Full;
13use hyper::body::Bytes;
14use hyper::header::{
15    HeaderValue, ACCESS_CONTROL_ALLOW_HEADERS, ACCESS_CONTROL_ALLOW_METHODS,
16    ACCESS_CONTROL_ALLOW_ORIGIN, ACCESS_CONTROL_EXPOSE_HEADERS, ACCESS_CONTROL_MAX_AGE,
17};
18use hyper::server::conn::http1;
19use hyper::service::service_fn;
20use hyper::{Method, Request, Response, StatusCode};
21use hyper_util::rt::TokioIo;
22use reqwest::Client;
23use serde::Deserialize;
24use serde_json::{json, Value};
25use std::convert::Infallible;
26use std::env;
27use std::net::SocketAddr;
28use std::sync::Arc;
29use std::time::{Duration, SystemTime, UNIX_EPOCH};
30use tokio::net::TcpListener;
31use tracing::{error, info};
32
33use crate::websocket::auth::{AuthDecision, AuthDeny, ConnectionAuthRequest, WebSocketAuthPlugin};
34use arete_auth::SCOPE_READ;
35
36/// Configuration for the HTTP health server
37#[derive(Clone, Debug)]
38pub struct HttpHealthConfig {
39    pub bind_address: SocketAddr,
40}
41
42impl Default for HttpHealthConfig {
43    fn default() -> Self {
44        Self {
45            bind_address: "[::]:8081".parse().expect("valid socket address"),
46        }
47    }
48}
49
50impl HttpHealthConfig {
51    pub fn new(bind_address: impl Into<SocketAddr>) -> Self {
52        Self {
53            bind_address: bind_address.into(),
54        }
55    }
56}
57
58#[derive(Clone)]
59struct HttpRequestState {
60    health_monitor: Arc<Option<HealthMonitor>>,
61    runtime_plan: RuntimePlan,
62    rpc_url: Arc<Option<String>>,
63    rpc_client: Client,
64    program_runtime_catalog: Arc<ProgramRuntimeCatalog>,
65    auth_plugin: Arc<Option<Arc<dyn WebSocketAuthPlugin>>>,
66    limit_state: Arc<HttpLimitState>,
67    transaction_state: Arc<Option<TransactionState>>,
68    solana_gateway_target_id: Arc<Option<String>>,
69    program_read_binding_target_id: Arc<Option<String>>,
70}
71
72/// HTTP server that exposes health endpoints
73pub struct HttpHealthServer {
74    bind_addr: SocketAddr,
75    health_monitor: Option<HealthMonitor>,
76    runtime_plan: RuntimePlan,
77    program_runtime_catalog: ProgramRuntimeCatalog,
78    auth_plugin: Option<Arc<dyn WebSocketAuthPlugin>>,
79    transaction_config: Option<TransactionConfig>,
80    solana_gateway_target_id: Option<String>,
81    program_read_binding_target_id: Option<String>,
82    #[cfg(feature = "otel")]
83    metrics: Option<Arc<crate::metrics::Metrics>>,
84}
85
86impl HttpHealthServer {
87    pub fn new(bind_addr: SocketAddr) -> Self {
88        Self {
89            bind_addr,
90            health_monitor: None,
91            runtime_plan: RuntimePlan::http(),
92            program_runtime_catalog: ProgramRuntimeCatalog::default(),
93            auth_plugin: None,
94            transaction_config: None,
95            solana_gateway_target_id: None,
96            program_read_binding_target_id: None,
97            #[cfg(feature = "otel")]
98            metrics: None,
99        }
100    }
101
102    pub fn with_health_monitor(mut self, monitor: HealthMonitor) -> Self {
103        self.health_monitor = Some(monitor);
104        self
105    }
106
107    pub fn with_runtime_plan(mut self, runtime_plan: RuntimePlan) -> Self {
108        self.runtime_plan = runtime_plan;
109        self
110    }
111
112    pub fn with_program_runtime_catalog(mut self, catalog: ProgramRuntimeCatalog) -> Self {
113        self.program_runtime_catalog = catalog;
114        self
115    }
116
117    pub fn with_auth_plugin(mut self, plugin: Arc<dyn WebSocketAuthPlugin>) -> Self {
118        self.auth_plugin = Some(plugin);
119        self
120    }
121
122    pub fn with_transaction_config(mut self, config: TransactionConfig) -> Self {
123        self.runtime_plan.transactions = config.enabled;
124        self.transaction_config = Some(config);
125        self
126    }
127
128    pub fn with_solana_gateway_target(mut self, target_id: impl Into<String>) -> Self {
129        self.solana_gateway_target_id = Some(target_id.into());
130        self
131    }
132
133    pub fn with_program_read_binding_target(mut self, target_id: impl Into<String>) -> Self {
134        self.program_read_binding_target_id = Some(target_id.into());
135        self
136    }
137
138    #[cfg(feature = "otel")]
139    pub fn with_metrics(mut self, metrics: Option<Arc<crate::metrics::Metrics>>) -> Self {
140        self.metrics = metrics;
141        self
142    }
143
144    pub async fn start(self) -> Result<()> {
145        info!("Starting HTTP health server on {}", self.bind_addr);
146
147        let listener = TcpListener::bind(&self.bind_addr).await?;
148        info!("HTTP health server listening on {}", self.bind_addr);
149
150        let transaction_state = self
151            .transaction_config
152            .filter(|config| config.enabled)
153            .map(TransactionState::new)
154            .transpose()?;
155        #[cfg(feature = "otel")]
156        let transaction_state = transaction_state.map(|state| state.with_metrics(self.metrics));
157        let request_state = HttpRequestState {
158            health_monitor: Arc::new(self.health_monitor),
159            runtime_plan: self.runtime_plan,
160            rpc_url: Arc::new(resolve_rpc_url()),
161            rpc_client: Client::builder().build()?,
162            program_runtime_catalog: Arc::new(self.program_runtime_catalog),
163            auth_plugin: Arc::new(self.auth_plugin),
164            limit_state: Arc::new(HttpLimitState::default()),
165            transaction_state: Arc::new(transaction_state),
166            solana_gateway_target_id: Arc::new(self.solana_gateway_target_id),
167            program_read_binding_target_id: Arc::new(self.program_read_binding_target_id),
168        };
169
170        loop {
171            match listener.accept().await {
172                Ok((stream, remote_addr)) => {
173                    let io = TokioIo::new(stream);
174                    let request_state = request_state.clone();
175
176                    tokio::spawn(async move {
177                        let service = service_fn(move |req| {
178                            let request_state = request_state.clone();
179                            async move { handle_request(remote_addr, req, request_state).await }
180                        });
181
182                        if let Err(e) = http1::Builder::new().serve_connection(io, service).await {
183                            error!("HTTP connection error: {}", e);
184                        }
185                    });
186                }
187                Err(e) => {
188                    error!("Failed to accept HTTP connection: {}", e);
189                }
190            }
191        }
192    }
193}
194
195async fn handle_request(
196    remote_addr: SocketAddr,
197    req: Request<hyper::body::Incoming>,
198    state: HttpRequestState,
199) -> Result<Response<Full<Bytes>>, Infallible> {
200    if req.method() == Method::OPTIONS {
201        return Ok(with_cors(
202            Response::builder()
203                .status(StatusCode::NO_CONTENT)
204                .body(Full::new(Bytes::new()))
205                .unwrap(),
206        ));
207    }
208
209    let response = handle_request_inner(remote_addr, req, state).await?;
210    Ok(with_cors(response))
211}
212
213fn with_cors(mut response: Response<Full<Bytes>>) -> Response<Full<Bytes>> {
214    let headers = response.headers_mut();
215    headers.insert(ACCESS_CONTROL_ALLOW_ORIGIN, HeaderValue::from_static("*"));
216    headers.insert(
217        ACCESS_CONTROL_ALLOW_METHODS,
218        HeaderValue::from_static("GET, POST, OPTIONS"),
219    );
220    headers.insert(
221        ACCESS_CONTROL_ALLOW_HEADERS,
222        HeaderValue::from_static("Authorization, Content-Type"),
223    );
224    headers.insert(
225        ACCESS_CONTROL_EXPOSE_HEADERS,
226        HeaderValue::from_static(
227            "Retry-After, X-Error-Code, X-Request-Id, X-Arete-Upstream-Attempted, X-Arete-Program-Release-Hash, X-Arete-Idl-Content-Hash, X-Arete-Account-Address, X-Arete-Account-Exists",
228        ),
229    );
230    headers.insert(ACCESS_CONTROL_MAX_AGE, HeaderValue::from_static("86400"));
231    response
232}
233
234async fn handle_request_inner(
235    remote_addr: SocketAddr,
236    req: Request<hyper::body::Incoming>,
237    state: HttpRequestState,
238) -> Result<Response<Full<Bytes>>, Infallible> {
239    let HttpRequestState {
240        health_monitor,
241        runtime_plan,
242        rpc_url,
243        rpc_client,
244        program_runtime_catalog,
245        auth_plugin,
246        limit_state,
247        transaction_state,
248        solana_gateway_target_id,
249        program_read_binding_target_id,
250    } = state;
251    let path = req.uri().path().to_string();
252
253    match path.as_str() {
254        "/health" | "/healthz" if runtime_plan.health => {
255            // Basic health check - server is running
256            Ok(Response::builder()
257                .status(StatusCode::OK)
258                .header("Content-Type", "text/plain")
259                .body(Full::new(Bytes::from("OK")))
260                .unwrap())
261        }
262        "/ready" | "/readiness" if runtime_plan.health => {
263            // Readiness check - check if stream is healthy
264            if let Some(monitor) = health_monitor.as_ref() {
265                if monitor.is_healthy().await {
266                    Ok(Response::builder()
267                        .status(StatusCode::OK)
268                        .header("Content-Type", "text/plain")
269                        .body(Full::new(Bytes::from("READY")))
270                        .unwrap())
271                } else {
272                    Ok(Response::builder()
273                        .status(StatusCode::SERVICE_UNAVAILABLE)
274                        .header("Content-Type", "text/plain")
275                        .body(Full::new(Bytes::from("NOT READY")))
276                        .unwrap())
277                }
278            } else {
279                // No health monitor configured, assume ready
280                Ok(Response::builder()
281                    .status(StatusCode::OK)
282                    .header("Content-Type", "text/plain")
283                    .body(Full::new(Bytes::from("READY")))
284                    .unwrap())
285            }
286        }
287        "/status" if runtime_plan.health => {
288            // Detailed status endpoint
289            if let Some(monitor) = health_monitor.as_ref() {
290                let status = monitor.status().await;
291                let error_count = monitor.error_count().await;
292                let is_healthy = monitor.is_healthy().await;
293
294                let status_json = serde_json::json!({
295                    "healthy": is_healthy,
296                    "status": format!("{:?}", status),
297                    "error_count": error_count
298                });
299
300                let status_code = if is_healthy {
301                    StatusCode::OK
302                } else {
303                    StatusCode::SERVICE_UNAVAILABLE
304                };
305
306                Ok(Response::builder()
307                    .status(status_code)
308                    .header("Content-Type", "application/json")
309                    .body(Full::new(Bytes::from(status_json.to_string())))
310                    .unwrap())
311            } else {
312                let status_json = serde_json::json!({
313                    "healthy": true,
314                    "status": "no_monitor",
315                    "error_count": 0
316                });
317
318                Ok(Response::builder()
319                    .status(StatusCode::OK)
320                    .header("Content-Type", "application/json")
321                    .body(Full::new(Bytes::from(status_json.to_string())))
322                    .unwrap())
323            }
324        }
325        _ if runtime_plan.transactions && path.starts_with("/transactions/") => {
326            let Some(transaction_state) = transaction_state.as_ref() else {
327                return Ok(error_response(StatusCode::NOT_FOUND, "Not Found"));
328            };
329            let client_addr = transaction_state.client_addr(remote_addr, req.headers());
330            let auth_context = match authorize_http_request(
331                client_addr,
332                &req,
333                auth_plugin.as_ref().as_ref(),
334                &limit_state,
335                None,
336                false,
337                solana_gateway_target_id.as_deref(),
338            )
339            .await
340            {
341                Ok(context) => context,
342                Err(response) => return Ok(transaction_auth_error(response, path.as_str())),
343            };
344            Ok(
345                transactions::handle(client_addr, req, auth_context, transaction_state.clone())
346                    .await,
347            )
348        }
349        _ if runtime_plan.chain_reads && path.starts_with("/chain/") => {
350            let auth_context = match authorize_http_request(
351                remote_addr,
352                &req,
353                auth_plugin.as_ref().as_ref(),
354                &limit_state,
355                Some(SCOPE_READ),
356                true,
357                solana_gateway_target_id.as_deref(),
358            )
359            .await
360            {
361                Ok(context) => context,
362                Err(response) => return Ok(response),
363            };
364            Ok(handle_chain_request(req, path.as_str(), rpc_url, rpc_client, auth_context).await)
365        }
366        _ if path.starts_with("/v1/releases/") => {
367            if !runtime_plan.program_reads {
368                return Ok(program_read_error_response(
369                    ProgramReadError::ProgramReadsDisabled,
370                ));
371            }
372            let auth_context = match authorize_http_request(
373                remote_addr,
374                &req,
375                auth_plugin.as_ref().as_ref(),
376                &limit_state,
377                Some(SCOPE_READ),
378                true,
379                None,
380            )
381            .await
382            {
383                Ok(context) => context,
384                Err(response) => return Ok(response),
385            };
386            Ok(handle_program_account_request(
387                req,
388                path.as_str(),
389                rpc_url,
390                rpc_client,
391                program_runtime_catalog,
392                auth_context,
393                program_read_binding_target_id,
394            )
395            .await)
396        }
397        _ => Ok(Response::builder()
398            .status(StatusCode::NOT_FOUND)
399            .header("Content-Type", "text/plain")
400            .body(Full::new(Bytes::from("Not Found")))
401            .unwrap()),
402    }
403}
404
405fn transaction_auth_error(response: Response<Full<Bytes>>, path: &str) -> Response<Full<Bytes>> {
406    let status = response.status();
407    let code = response
408        .headers()
409        .get("X-Error-Code")
410        .and_then(|value| value.to_str().ok())
411        .unwrap_or("authentication_failed")
412        .to_string();
413    let request_id = uuid::Uuid::new_v4().to_string();
414    let mut value = json!({
415        "code": code,
416        "message": "Transaction request authentication failed",
417        "retryable": status == StatusCode::TOO_MANY_REQUESTS || status == StatusCode::UNAUTHORIZED,
418        "requestId": request_id,
419    });
420    if path == "/transactions/v1/send" {
421        value["submissionState"] = json!("not_submitted");
422    }
423    Response::builder()
424        .status(status)
425        .header("Content-Type", "application/json")
426        .header("X-Error-Code", code)
427        .header("X-Request-Id", request_id)
428        .header("X-Arete-Upstream-Attempted", "false")
429        .body(Full::new(Bytes::from(value.to_string())))
430        .expect("valid transaction authentication response")
431}
432
433#[derive(Debug, Deserialize)]
434struct AddressesBody {
435    addresses: Vec<String>,
436}
437
438#[derive(Debug, Deserialize)]
439struct BalanceBody {
440    owner: String,
441    mint: String,
442    #[serde(default, rename = "tokenProgram")]
443    token_program: Option<String>,
444    #[serde(default, rename = "minContextSlot")]
445    min_context_slot: Option<String>,
446}
447
448#[derive(Debug, Deserialize)]
449struct NativeBalanceBody {
450    address: String,
451    #[serde(default, rename = "minContextSlot")]
452    min_context_slot: Option<String>,
453}
454
455#[derive(Default)]
456struct HttpLimitState {
457    per_subject_per_minute: DashMap<String, (u64, u32)>,
458}
459
460fn resolve_rpc_url() -> Option<String> {
461    ["ARETE_READ_RPC_URL", "SOLANA_RPC_URL", "RPC_URL"]
462        .iter()
463        .find_map(|key| env::var(key).ok())
464        .filter(|value| !value.is_empty())
465}
466
467fn json_response(status: StatusCode, value: Value) -> Response<Full<Bytes>> {
468    Response::builder()
469        .status(status)
470        .header("Content-Type", "application/json")
471        .body(Full::new(Bytes::from(value.to_string())))
472        .unwrap()
473}
474
475fn error_response(status: StatusCode, message: impl Into<String>) -> Response<Full<Bytes>> {
476    json_response(status, json!({ "error": message.into() }))
477}
478
479fn parse_min_context_slot(value: Option<&str>) -> std::result::Result<Option<u64>, &'static str> {
480    value
481        .map(|slot| {
482            slot.parse::<u64>()
483                .map_err(|_| "minContextSlot must be a decimal u64 string")
484        })
485        .transpose()
486}
487
488fn auth_deny_response(deny: &AuthDeny) -> Response<Full<Bytes>> {
489    let mut builder = Response::builder()
490        .status(deny.http_status)
491        .header("Content-Type", "application/json")
492        .header("X-Error-Code", deny.code.as_str());
493
494    if let Some(reset_at) = deny.reset_at {
495        if let Ok(duration) = reset_at.duration_since(SystemTime::now()) {
496            builder = builder.header("Retry-After", duration.as_secs().to_string());
497        }
498    }
499
500    builder
501        .body(Full::new(Bytes::from(
502            json!({
503                "error": deny.reason,
504                "message": deny.reason,
505                "code": deny.code.as_str(),
506                "retryable": deny.code.should_retry(),
507                "fatal": !deny.code.should_retry() && !deny.code.should_refresh_token(),
508            })
509            .to_string(),
510        )))
511        .unwrap()
512}
513
514async fn authorize_http_request(
515    remote_addr: SocketAddr,
516    req: &Request<hyper::body::Incoming>,
517    auth_plugin: Option<&Arc<dyn WebSocketAuthPlugin>>,
518    limit_state: &HttpLimitState,
519    required_scope: Option<&str>,
520    enforce_read_limits: bool,
521    solana_gateway_target_id: Option<&str>,
522) -> std::result::Result<Option<crate::websocket::auth::AuthContext>, Response<Full<Bytes>>> {
523    let Some(plugin) = auth_plugin else {
524        return Ok(None);
525    };
526
527    let mut auth_request = ConnectionAuthRequest::from_http_request(remote_addr, req);
528    // HTTP reads are bearer-only. Do not allow query-param session tokens here.
529    auth_request.query = None;
530    let decision = plugin.authorize(&auth_request).await;
531    let context = match decision {
532        AuthDecision::Allow(context) => context,
533        AuthDecision::Deny(deny) => return Err(auth_deny_response(&deny)),
534    };
535
536    if let Some(target_id) = solana_gateway_target_id {
537        if let Err(error) =
538            arete_auth::SolanaGatewayAuthorization::validate_target(&context, target_id)
539        {
540            return Err(Response::builder()
541                .status(StatusCode::FORBIDDEN)
542                .header("Content-Type", "application/json")
543                .header("X-Error-Code", "invalid_gateway_target")
544                .body(Full::new(Bytes::from(
545                    json!({
546                        "error": "invalid_gateway_target",
547                        "message": error.to_string(),
548                        "code": "invalid_gateway_target",
549                        "retryable": false,
550                        "fatal": true
551                    })
552                    .to_string(),
553                )))
554                .expect("valid gateway target error response"));
555        }
556    }
557
558    if let Some(required_scope) = required_scope {
559        if !context.has_scope(required_scope) {
560            return Err(json_response(
561                StatusCode::FORBIDDEN,
562                json!({
563                    "error": "insufficient_scope",
564                    "message": format!("Required scope: {required_scope}"),
565                    "code": "insufficient_scope",
566                    "retryable": false,
567                    "fatal": true
568                }),
569            ));
570        }
571    }
572
573    if enforce_read_limits {
574        enforce_http_limits(&context, limit_state).map_err(|deny| auth_deny_response(&deny))?;
575    }
576    Ok(Some(context))
577}
578
579fn enforce_http_limits(
580    context: &crate::websocket::auth::AuthContext,
581    limit_state: &HttpLimitState,
582) -> std::result::Result<(), Box<AuthDeny>> {
583    let Some(limit) = context.limits.max_http_requests_per_minute else {
584        return Ok(());
585    };
586
587    let now_bucket = SystemTime::now()
588        .duration_since(UNIX_EPOCH)
589        .unwrap_or(Duration::from_secs(0))
590        .as_secs()
591        / 60;
592    let key = format!("{}:{}", context.subject, context.metering_key);
593    let mut entry = limit_state
594        .per_subject_per_minute
595        .entry(key)
596        .or_insert((now_bucket, 0));
597    if entry.0 != now_bucket {
598        *entry = (now_bucket, 0);
599    }
600    if entry.1 >= limit {
601        return Err(Box::new(AuthDeny::rate_limited(
602            Duration::from_secs(60),
603            "http reads",
604        )));
605    }
606    entry.1 += 1;
607    Ok(())
608}
609
610async fn read_json_body<T: for<'de> Deserialize<'de>>(
611    req: Request<hyper::body::Incoming>,
612) -> std::result::Result<T, Response<Full<Bytes>>> {
613    let collected = req
614        .into_body()
615        .collect()
616        .await
617        .map_err(|err| error_response(StatusCode::BAD_REQUEST, err.to_string()))?;
618    serde_json::from_slice::<T>(&collected.to_bytes())
619        .map_err(|err| error_response(StatusCode::BAD_REQUEST, err.to_string()))
620}
621
622async fn handle_chain_request(
623    req: Request<hyper::body::Incoming>,
624    path: &str,
625    rpc_url: Arc<Option<String>>,
626    rpc_client: Client,
627    _auth_context: Option<crate::websocket::auth::AuthContext>,
628) -> Response<Full<Bytes>> {
629    let Some(rpc_url) = rpc_url.as_ref() else {
630        return error_response(
631            StatusCode::SERVICE_UNAVAILABLE,
632            "No RPC URL configured for chain reads",
633        );
634    };
635
636    match (req.method().as_str(), path) {
637        ("GET", path) if path.starts_with("/chain/exists/") => {
638            let address = path.trim_start_matches("/chain/exists/");
639            match rpc_get_account_info(&rpc_client, rpc_url, address).await {
640                Ok(value) => json_response(StatusCode::OK, json!({ "exists": !value.is_null() })),
641                Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
642            }
643        }
644        ("GET", path) if path.starts_with("/chain/lamports/") => {
645            let address = path.trim_start_matches("/chain/lamports/");
646            match rpc_call(
647                &rpc_client,
648                rpc_url,
649                "getBalance",
650                json!([address, { "commitment": "confirmed" }]),
651            )
652            .await
653            {
654                Ok(value) => json_response(
655                    StatusCode::OK,
656                    json!({ "lamports": value.pointer("/value").and_then(Value::as_u64).unwrap_or(0) }),
657                ),
658                Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
659            }
660        }
661        ("POST", "/chain/native-balance") => match read_json_body::<NativeBalanceBody>(req).await {
662            Ok(body) => {
663                let min_context_slot =
664                    match parse_min_context_slot(body.min_context_slot.as_deref()) {
665                        Ok(slot) => slot,
666                        Err(message) => return error_response(StatusCode::BAD_REQUEST, message),
667                    };
668                match rpc_get_native_balance(&rpc_client, rpc_url, &body.address, min_context_slot)
669                    .await
670                {
671                    Ok(balance) => json_response(StatusCode::OK, balance),
672                    Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
673                }
674            }
675            Err(response) => response,
676        },
677        ("GET", path) if path.starts_with("/chain/rent-exemption/") => {
678            let raw_space = path.trim_start_matches("/chain/rent-exemption/");
679            let Ok(space) = raw_space.parse::<u64>() else {
680                return error_response(
681                    StatusCode::BAD_REQUEST,
682                    "rent-exemption space must be an integer",
683                );
684            };
685            match rpc_call(
686                &rpc_client,
687                rpc_url,
688                "getMinimumBalanceForRentExemption",
689                json!([space, { "commitment": "confirmed" }]),
690            )
691            .await
692            {
693                Ok(value) => json_response(
694                    StatusCode::OK,
695                    json!({ "lamports": value.as_u64().unwrap_or(0) }),
696                ),
697                Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
698            }
699        }
700        ("GET", "/chain/clock") => {
701            let slot = rpc_call(
702                &rpc_client,
703                rpc_url,
704                "getSlot",
705                json!([{ "commitment": "confirmed" }]),
706            )
707            .await;
708            let epoch_info = rpc_call(
709                &rpc_client,
710                rpc_url,
711                "getEpochInfo",
712                json!([{ "commitment": "confirmed" }]),
713            )
714            .await;
715            match (slot, epoch_info) {
716                (Ok(slot_value), Ok(epoch_value)) => {
717                    let slot_num = slot_value.as_u64().unwrap_or(0);
718                    let unix_timestamp =
719                        rpc_call(&rpc_client, rpc_url, "getBlockTime", json!([slot_num]))
720                            .await
721                            .ok()
722                            .and_then(|value| value.as_i64())
723                            .unwrap_or_default();
724                    json_response(
725                        StatusCode::OK,
726                        json!({
727                            "slot": slot_num,
728                            "epoch": epoch_value.get("epoch").and_then(Value::as_u64),
729                            "leaderScheduleEpoch": epoch_value.get("leaderScheduleSlotOffset").and_then(Value::as_u64),
730                            "unixTimestamp": unix_timestamp,
731                        }),
732                    )
733                }
734                (Err(err), _) | (_, Err(err)) => {
735                    error_response(StatusCode::BAD_GATEWAY, err.to_string())
736                }
737            }
738        }
739        ("GET", path) if path.starts_with("/chain/accounts/") => {
740            let address = path.trim_start_matches("/chain/accounts/");
741            match rpc_get_account_info(&rpc_client, rpc_url, address).await {
742                Ok(value) if value.is_null() => json_response(StatusCode::OK, Value::Null),
743                Ok(value) => json_response(StatusCode::OK, raw_account_json(address, &value)),
744                Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
745            }
746        }
747        ("GET", path) if path.starts_with("/chain/mints/") => {
748            let address = path.trim_start_matches("/chain/mints/");
749            match rpc_get_parsed_account_info(&rpc_client, rpc_url, address).await {
750                Ok(Some(value)) => json_response(StatusCode::OK, mint_info_json(address, &value)),
751                Ok(None) => json_response(StatusCode::OK, Value::Null),
752                Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
753            }
754        }
755        ("GET", path) if path.starts_with("/chain/token-accounts/") => {
756            let address = path.trim_start_matches("/chain/token-accounts/");
757            match rpc_get_parsed_account_info(&rpc_client, rpc_url, address).await {
758                Ok(Some(value)) => {
759                    json_response(StatusCode::OK, token_account_json(address, &value))
760                }
761                Ok(None) => json_response(StatusCode::OK, Value::Null),
762                Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
763            }
764        }
765        ("POST", "/chain/balances") => match read_json_body::<BalanceBody>(req).await {
766            Ok(body) => {
767                let min_context_slot =
768                    match parse_min_context_slot(body.min_context_slot.as_deref()) {
769                        Ok(slot) => slot,
770                        Err(message) => return error_response(StatusCode::BAD_REQUEST, message),
771                    };
772                match rpc_get_token_balance(
773                    &rpc_client,
774                    rpc_url,
775                    &body.owner,
776                    &body.mint,
777                    body.token_program.as_deref(),
778                    min_context_slot,
779                )
780                .await
781                {
782                    Ok(balance) => json_response(StatusCode::OK, balance),
783                    Err(err) => error_response(StatusCode::BAD_GATEWAY, err.to_string()),
784                }
785            }
786            Err(response) => response,
787        },
788        _ => error_response(StatusCode::NOT_FOUND, "Not Found"),
789    }
790}
791
792async fn handle_program_account_request(
793    req: Request<hyper::body::Incoming>,
794    path: &str,
795    rpc_url: Arc<Option<String>>,
796    rpc_client: Client,
797    program_runtime_catalog: Arc<ProgramRuntimeCatalog>,
798    auth_context: Option<crate::websocket::auth::AuthContext>,
799    program_read_binding_target_id: Arc<Option<String>>,
800) -> Response<Full<Bytes>> {
801    let route = match ProgramReadRoute::parse(req.method(), path) {
802        Ok(route) => route,
803        Err(error) => return program_read_error_response(error),
804    };
805    let release_hash = match route.release_hash.parse() {
806        Ok(release_hash) => release_hash,
807        Err(_) => return program_read_error_response(ProgramReadError::InvalidReleaseHash),
808    };
809    let Some(definition) = program_runtime_catalog.get(&release_hash).cloned() else {
810        return program_read_error_response(ProgramReadError::ReleaseNotFound);
811    };
812    if let Err(error) = authorize_program_read(
813        auth_context.as_ref(),
814        program_read_binding_target_id.as_deref(),
815        &definition,
816    ) {
817        return with_release_metadata(program_read_error_response(error), &definition, None, None);
818    }
819    let Some(rpc_url) = rpc_url.as_ref() else {
820        return with_release_metadata(
821            program_read_error_response(ProgramReadError::RpcNotConfigured),
822            &definition,
823            None,
824            None,
825        );
826    };
827
828    match route.operation {
829        ProgramReadOperation::Fetch { address } => {
830            let account_value = match rpc_get_account_info(&rpc_client, rpc_url, &address).await {
831                Ok(value) => value,
832                Err(error) => {
833                    error!("Program account RPC read failed: {}", error);
834                    return with_release_metadata(
835                        program_read_error_response(ProgramReadError::RpcFailed),
836                        &definition,
837                        Some(&address),
838                        None,
839                    );
840                }
841            };
842            match decode_release_account(&definition, &route.account, &account_value).await {
843                AccountReadOutcome::Missing => with_release_metadata(
844                    json_response(StatusCode::OK, Value::Null),
845                    &definition,
846                    Some(&address),
847                    Some(false),
848                ),
849                AccountReadOutcome::Value(value) => with_release_metadata(
850                    json_response(StatusCode::OK, value),
851                    &definition,
852                    Some(&address),
853                    Some(true),
854                ),
855                AccountReadOutcome::Error(error) => with_release_metadata(
856                    program_read_error_response(error),
857                    &definition,
858                    Some(&address),
859                    Some(true),
860                ),
861            }
862        }
863        ProgramReadOperation::Exists { address } => {
864            let account_value = match rpc_get_account_info(&rpc_client, rpc_url, &address).await {
865                Ok(value) => value,
866                Err(error) => {
867                    error!("Program account existence RPC read failed: {}", error);
868                    return with_release_metadata(
869                        program_read_error_response(ProgramReadError::RpcFailed),
870                        &definition,
871                        Some(&address),
872                        None,
873                    );
874                }
875            };
876            let exists = !account_value.is_null();
877            if exists && !account_owner_matches(&definition, &account_value) {
878                return with_release_metadata(
879                    program_read_error_response(ProgramReadError::AccountOwnerMismatch),
880                    &definition,
881                    Some(&address),
882                    Some(true),
883                );
884            }
885            with_release_metadata(
886                json_response(StatusCode::OK, json!({ "exists": exists })),
887                &definition,
888                Some(&address),
889                Some(exists),
890            )
891        }
892        ProgramReadOperation::Batch => {
893            let body = match read_json_body::<AddressesBody>(req).await {
894                Ok(body) => body,
895                Err(_) => return program_read_error_response(ProgramReadError::InvalidRequest),
896            };
897            let configured_limit = auth_context
898                .as_ref()
899                .and_then(|ctx| ctx.limits.max_http_batch_addresses)
900                .map(|limit| limit as usize)
901                .unwrap_or(MAX_PROGRAM_BATCH_ADDRESSES)
902                .min(MAX_PROGRAM_BATCH_ADDRESSES);
903            if body.addresses.len() > configured_limit {
904                return with_release_metadata(
905                    program_read_error_response(ProgramReadError::BatchLimitExceeded),
906                    &definition,
907                    None,
908                    None,
909                );
910            }
911            if body.addresses.is_empty() {
912                return with_release_metadata(
913                    json_response(StatusCode::OK, json!({ "items": [] })),
914                    &definition,
915                    None,
916                    None,
917                );
918            }
919
920            let values =
921                match rpc_get_multiple_accounts(&rpc_client, rpc_url, &body.addresses).await {
922                    Ok(values) if values.len() == body.addresses.len() => values,
923                    Ok(_) => {
924                        return with_release_metadata(
925                            program_read_error_response(ProgramReadError::RpcResponseInvalid),
926                            &definition,
927                            None,
928                            None,
929                        )
930                    }
931                    Err(error) => {
932                        error!("Program account batch RPC read failed: {}", error);
933                        return with_release_metadata(
934                            program_read_error_response(ProgramReadError::RpcFailed),
935                            &definition,
936                            None,
937                            None,
938                        );
939                    }
940                };
941
942            let outcomes = futures_util::future::join_all(
943                values
944                    .iter()
945                    .map(|value| decode_release_account(&definition, &route.account, value)),
946            )
947            .await;
948            let mut items = Vec::with_capacity(outcomes.len());
949            for (address, outcome) in body.addresses.iter().zip(outcomes) {
950                let item = match outcome {
951                    AccountReadOutcome::Missing => {
952                        json!({ "address": address, "status": "missing" })
953                    }
954                    AccountReadOutcome::Value(value) => {
955                        json!({ "address": address, "status": "ok", "value": value })
956                    }
957                    AccountReadOutcome::Error(error) => json!({
958                        "address": address,
959                        "status": "error",
960                        "error": { "code": error.code() }
961                    }),
962                };
963                items.push(item);
964            }
965            with_release_metadata(
966                json_response(StatusCode::OK, json!({ "items": items })),
967                &definition,
968                None,
969                None,
970            )
971        }
972    }
973}
974
975const MAX_PROGRAM_BATCH_ADDRESSES: usize = 100;
976const MAX_PROGRAM_ACCOUNT_BYTES: usize = 10 * 1024 * 1024;
977const PROGRAM_DECODE_TIMEOUT: Duration = Duration::from_secs(2);
978
979fn authorize_program_read(
980    auth_context: Option<&crate::websocket::auth::AuthContext>,
981    expected_target_id: Option<&str>,
982    definition: &ProgramRuntimeDefinition,
983) -> std::result::Result<(), ProgramReadError> {
984    let Some(context) = auth_context else {
985        return Ok(());
986    };
987    let Some(expected_target_id) = expected_target_id.filter(|target_id| !target_id.is_empty())
988    else {
989        return Err(ProgramReadError::AuthorizationNotConfigured);
990    };
991
992    arete_auth::ProgramReadAuthorization::try_from_context(
993        context,
994        expected_target_id,
995        &definition.program_id,
996        &definition.program_release_hash.to_string(),
997    )
998    .map(|_| ())
999    .map_err(|_| ProgramReadError::Unauthorized)
1000}
1001
1002#[derive(Clone, Copy, Debug, PartialEq, Eq)]
1003enum ProgramReadError {
1004    NotFound,
1005    ProgramReadsDisabled,
1006    InvalidReleaseHash,
1007    ReleaseNotFound,
1008    InvalidRequest,
1009    Unauthorized,
1010    AuthorizationNotConfigured,
1011    BatchLimitExceeded,
1012    RpcNotConfigured,
1013    RpcFailed,
1014    RpcResponseInvalid,
1015    AccountOwnerMismatch,
1016    AccountDataInvalid,
1017    AccountDataTooLarge,
1018    AccountDecodeFailed,
1019    AccountDecodeTimeout,
1020}
1021
1022impl ProgramReadError {
1023    fn code(self) -> &'static str {
1024        match self {
1025            Self::NotFound => "NOT_FOUND",
1026            Self::ProgramReadsDisabled => "PROGRAM_READS_DISABLED",
1027            Self::InvalidReleaseHash => "INVALID_PROGRAM_RELEASE_HASH",
1028            Self::ReleaseNotFound => "PROGRAM_RELEASE_NOT_FOUND",
1029            Self::InvalidRequest => "INVALID_REQUEST",
1030            Self::Unauthorized => "PROGRAM_READ_UNAUTHORIZED",
1031            Self::AuthorizationNotConfigured => "PROGRAM_READ_AUTH_NOT_CONFIGURED",
1032            Self::BatchLimitExceeded => "BATCH_LIMIT_EXCEEDED",
1033            Self::RpcNotConfigured => "READ_RPC_NOT_CONFIGURED",
1034            Self::RpcFailed => "RPC_REQUEST_FAILED",
1035            Self::RpcResponseInvalid => "RPC_RESPONSE_INVALID",
1036            Self::AccountOwnerMismatch => "ACCOUNT_OWNER_MISMATCH",
1037            Self::AccountDataInvalid => "ACCOUNT_DATA_INVALID",
1038            Self::AccountDataTooLarge => "ACCOUNT_DATA_TOO_LARGE",
1039            Self::AccountDecodeFailed => "ACCOUNT_DECODE_FAILED",
1040            Self::AccountDecodeTimeout => "ACCOUNT_DECODE_TIMEOUT",
1041        }
1042    }
1043
1044    fn status(self) -> StatusCode {
1045        match self {
1046            Self::NotFound | Self::ReleaseNotFound => StatusCode::NOT_FOUND,
1047            Self::InvalidReleaseHash | Self::InvalidRequest => StatusCode::BAD_REQUEST,
1048            Self::Unauthorized => StatusCode::FORBIDDEN,
1049            Self::AuthorizationNotConfigured => StatusCode::SERVICE_UNAVAILABLE,
1050            Self::BatchLimitExceeded | Self::AccountDataTooLarge => StatusCode::PAYLOAD_TOO_LARGE,
1051            Self::ProgramReadsDisabled | Self::RpcNotConfigured => StatusCode::SERVICE_UNAVAILABLE,
1052            Self::RpcFailed | Self::RpcResponseInvalid | Self::AccountDataInvalid => {
1053                StatusCode::BAD_GATEWAY
1054            }
1055            Self::AccountOwnerMismatch | Self::AccountDecodeFailed => {
1056                StatusCode::UNPROCESSABLE_ENTITY
1057            }
1058            Self::AccountDecodeTimeout => StatusCode::GATEWAY_TIMEOUT,
1059        }
1060    }
1061}
1062
1063fn program_read_error_response(error: ProgramReadError) -> Response<Full<Bytes>> {
1064    Response::builder()
1065        .status(error.status())
1066        .header("Content-Type", "application/json")
1067        .header("X-Error-Code", error.code())
1068        .body(Full::new(Bytes::from(
1069            json!({ "error": { "code": error.code() } }).to_string(),
1070        )))
1071        .expect("valid program read error response")
1072}
1073
1074#[derive(Debug)]
1075struct ProgramReadRoute {
1076    release_hash: String,
1077    account: String,
1078    operation: ProgramReadOperation,
1079}
1080
1081#[derive(Debug)]
1082enum ProgramReadOperation {
1083    Fetch { address: String },
1084    Batch,
1085    Exists { address: String },
1086}
1087
1088impl ProgramReadRoute {
1089    fn parse(method: &Method, path: &str) -> std::result::Result<Self, ProgramReadError> {
1090        let segments: Vec<&str> = path.trim_start_matches('/').split('/').collect();
1091        if segments.len() < 5
1092            || segments[0] != "v1"
1093            || segments[1] != "releases"
1094            || segments[2].is_empty()
1095            || segments[3] != "accounts"
1096            || segments[4].is_empty()
1097        {
1098            return Err(ProgramReadError::NotFound);
1099        }
1100        let operation = match (method, segments.as_slice()) {
1101            (&Method::POST, [_, _, _, _, _]) => ProgramReadOperation::Batch,
1102            (&Method::GET, [_, _, _, _, _, address]) if !address.is_empty() => {
1103                ProgramReadOperation::Fetch {
1104                    address: (*address).to_string(),
1105                }
1106            }
1107            (&Method::GET, [_, _, _, _, _, address, "exists"]) if !address.is_empty() => {
1108                ProgramReadOperation::Exists {
1109                    address: (*address).to_string(),
1110                }
1111            }
1112            _ => return Err(ProgramReadError::NotFound),
1113        };
1114        Ok(Self {
1115            release_hash: segments[2].to_string(),
1116            account: segments[4].to_string(),
1117            operation,
1118        })
1119    }
1120}
1121
1122enum AccountReadOutcome {
1123    Missing,
1124    Value(Value),
1125    Error(ProgramReadError),
1126}
1127
1128fn account_owner_matches(definition: &ProgramRuntimeDefinition, value: &Value) -> bool {
1129    value.get("owner").and_then(Value::as_str) == Some(definition.program_id.as_str())
1130}
1131
1132async fn decode_release_account(
1133    definition: &ProgramRuntimeDefinition,
1134    account: &str,
1135    value: &Value,
1136) -> AccountReadOutcome {
1137    if value.is_null() {
1138        return AccountReadOutcome::Missing;
1139    }
1140    if !account_owner_matches(definition, value) {
1141        return AccountReadOutcome::Error(ProgramReadError::AccountOwnerMismatch);
1142    }
1143    let data = match decode_account_bytes(value) {
1144        Some(data) => data,
1145        None => return AccountReadOutcome::Error(ProgramReadError::AccountDataInvalid),
1146    };
1147    if data.len() > MAX_PROGRAM_ACCOUNT_BYTES {
1148        return AccountReadOutcome::Error(ProgramReadError::AccountDataTooLarge);
1149    }
1150
1151    let reader = definition.account_reader.clone();
1152    let account = account.to_string();
1153    let decode = tokio::task::spawn_blocking(move || reader(&account, &data));
1154    match tokio::time::timeout(PROGRAM_DECODE_TIMEOUT, decode).await {
1155        Ok(Ok(Ok(value))) => AccountReadOutcome::Value(value),
1156        Ok(Ok(Err(_))) | Ok(Err(_)) => {
1157            AccountReadOutcome::Error(ProgramReadError::AccountDecodeFailed)
1158        }
1159        Err(_) => AccountReadOutcome::Error(ProgramReadError::AccountDecodeTimeout),
1160    }
1161}
1162
1163fn with_release_metadata(
1164    mut response: Response<Full<Bytes>>,
1165    definition: &ProgramRuntimeDefinition,
1166    address: Option<&str>,
1167    exists: Option<bool>,
1168) -> Response<Full<Bytes>> {
1169    let headers = response.headers_mut();
1170    if let Ok(value) = HeaderValue::from_str(&definition.program_release_hash.to_string()) {
1171        headers.insert("X-Arete-Program-Release-Hash", value);
1172    }
1173    if let Ok(value) = HeaderValue::from_str(&definition.idl_content_hash.to_string()) {
1174        headers.insert("X-Arete-Idl-Content-Hash", value);
1175    }
1176    if let Some(address) = address {
1177        if let Ok(value) = HeaderValue::from_str(address) {
1178            headers.insert("X-Arete-Account-Address", value);
1179        }
1180    }
1181    if let Some(exists) = exists {
1182        headers.insert(
1183            "X-Arete-Account-Exists",
1184            HeaderValue::from_static(if exists { "true" } else { "false" }),
1185        );
1186    }
1187    response
1188}
1189
1190async fn rpc_call(
1191    client: &Client,
1192    rpc_url: &str,
1193    method: &str,
1194    params: Value,
1195) -> anyhow::Result<Value> {
1196    let response = client
1197        .post(rpc_url)
1198        .json(&json!({
1199            "jsonrpc": "2.0",
1200            "id": "arete-read",
1201            "method": method,
1202            "params": params,
1203        }))
1204        .send()
1205        .await?
1206        .error_for_status()?;
1207    let value = response.json::<Value>().await?;
1208    if let Some(error) = value.get("error") {
1209        return Err(anyhow::anyhow!(error.to_string()));
1210    }
1211    Ok(value.get("result").cloned().unwrap_or(Value::Null))
1212}
1213
1214fn rpc_read_config(encoding: Option<&str>, min_context_slot: Option<u64>) -> Value {
1215    let mut config = json!({ "commitment": "confirmed" });
1216    let object = config
1217        .as_object_mut()
1218        .expect("RPC read config is always an object");
1219    if let Some(encoding) = encoding {
1220        object.insert("encoding".to_string(), json!(encoding));
1221    }
1222    if let Some(min_context_slot) = min_context_slot {
1223        object.insert("minContextSlot".to_string(), json!(min_context_slot));
1224    }
1225    config
1226}
1227
1228async fn rpc_get_native_balance(
1229    client: &Client,
1230    rpc_url: &str,
1231    address: &str,
1232    min_context_slot: Option<u64>,
1233) -> anyhow::Result<Value> {
1234    let result = rpc_call(
1235        client,
1236        rpc_url,
1237        "getBalance",
1238        json!([address, rpc_read_config(None, min_context_slot)]),
1239    )
1240    .await?;
1241    contextual_native_balance_json(&result)
1242}
1243
1244fn contextual_native_balance_json(result: &Value) -> anyhow::Result<Value> {
1245    let lamports = result
1246        .get("value")
1247        .and_then(Value::as_u64)
1248        .ok_or_else(|| anyhow::anyhow!("getBalance response is missing a u64 value"))?;
1249    let context_slot = result
1250        .pointer("/context/slot")
1251        .and_then(Value::as_u64)
1252        .ok_or_else(|| anyhow::anyhow!("getBalance response is missing a u64 context slot"))?;
1253    Ok(json!({
1254        "lamports": lamports.to_string(),
1255        "contextSlot": context_slot.to_string(),
1256    }))
1257}
1258
1259async fn rpc_get_account_info(
1260    client: &Client,
1261    rpc_url: &str,
1262    address: &str,
1263) -> anyhow::Result<Value> {
1264    let result = rpc_call(
1265        client,
1266        rpc_url,
1267        "getAccountInfo",
1268        json!([address, { "encoding": "base64", "commitment": "confirmed" }]),
1269    )
1270    .await?;
1271    Ok(result.get("value").cloned().unwrap_or(Value::Null))
1272}
1273
1274async fn rpc_get_multiple_accounts(
1275    client: &Client,
1276    rpc_url: &str,
1277    addresses: &[String],
1278) -> anyhow::Result<Vec<Value>> {
1279    let result = rpc_call(
1280        client,
1281        rpc_url,
1282        "getMultipleAccounts",
1283        json!([addresses, { "encoding": "base64", "commitment": "confirmed" }]),
1284    )
1285    .await?;
1286    Ok(result
1287        .get("value")
1288        .and_then(Value::as_array)
1289        .cloned()
1290        .unwrap_or_default())
1291}
1292
1293async fn rpc_get_parsed_account_info(
1294    client: &Client,
1295    rpc_url: &str,
1296    address: &str,
1297) -> anyhow::Result<Option<Value>> {
1298    let result = rpc_call(
1299        client,
1300        rpc_url,
1301        "getAccountInfo",
1302        json!([address, { "encoding": "jsonParsed", "commitment": "confirmed" }]),
1303    )
1304    .await?;
1305    Ok(result
1306        .get("value")
1307        .cloned()
1308        .filter(|value| !value.is_null()))
1309}
1310
1311async fn rpc_get_token_balance(
1312    client: &Client,
1313    rpc_url: &str,
1314    owner: &str,
1315    mint: &str,
1316    token_program: Option<&str>,
1317    min_context_slot: Option<u64>,
1318) -> anyhow::Result<Value> {
1319    let filter = token_program
1320        .map(|program_id| json!({ "programId": program_id }))
1321        .unwrap_or_else(|| json!({ "mint": mint }));
1322    let result = rpc_call(
1323        client,
1324        rpc_url,
1325        "getTokenAccountsByOwner",
1326        json!([
1327            owner,
1328            filter,
1329            rpc_read_config(Some("jsonParsed"), min_context_slot)
1330        ]),
1331    )
1332    .await?;
1333
1334    contextual_token_balance_json(&result, owner, mint, token_program)
1335}
1336
1337fn contextual_token_balance_json(
1338    result: &Value,
1339    owner: &str,
1340    mint: &str,
1341    token_program: Option<&str>,
1342) -> anyhow::Result<Value> {
1343    let context_slot = result
1344        .pointer("/context/slot")
1345        .and_then(Value::as_u64)
1346        .ok_or_else(|| {
1347            anyhow::anyhow!("getTokenAccountsByOwner response is missing a u64 context slot")
1348        })?;
1349
1350    let account = result
1351        .get("value")
1352        .and_then(Value::as_array)
1353        .and_then(|items| {
1354            items.iter().find(|item| {
1355                item.pointer("/account/data/parsed/info/mint")
1356                    .and_then(Value::as_str)
1357                    == Some(mint)
1358            })
1359        })
1360        .cloned();
1361
1362    if let Some(account) = account {
1363        let pubkey = account.get("pubkey").and_then(Value::as_str);
1364        let info = account.pointer("/account/data/parsed/info");
1365        return Ok(json!({
1366            "exists": true,
1367            "address": pubkey,
1368            "owner": owner,
1369            "mint": mint,
1370            "tokenProgram": token_program,
1371            "amount": info.and_then(|value| value.pointer("/tokenAmount/amount")).and_then(Value::as_str).unwrap_or("0"),
1372            "decimals": info.and_then(|value| value.pointer("/tokenAmount/decimals")).and_then(Value::as_u64),
1373            "uiAmountString": info.and_then(|value| value.pointer("/tokenAmount/uiAmountString")).and_then(Value::as_str),
1374            "contextSlot": context_slot.to_string(),
1375        }));
1376    }
1377
1378    Ok(json!({
1379        "exists": false,
1380        "address": Value::Null,
1381        "owner": owner,
1382        "mint": mint,
1383        "tokenProgram": token_program,
1384        "amount": "0",
1385        "decimals": Value::Null,
1386        "uiAmountString": Value::Null,
1387        "contextSlot": context_slot.to_string(),
1388    }))
1389}
1390
1391fn decode_account_bytes(value: &Value) -> Option<Vec<u8>> {
1392    let data = value.get("data")?.as_array()?;
1393    let encoded = data.first()?.as_str()?;
1394    BASE64_STANDARD.decode(encoded).ok()
1395}
1396
1397fn raw_account_json(address: &str, value: &Value) -> Value {
1398    json!({
1399        "address": address,
1400        "ownerProgram": value.get("owner").and_then(Value::as_str).unwrap_or_default(),
1401        "lamports": value.get("lamports").and_then(Value::as_u64).unwrap_or(0),
1402        "executable": value.get("executable").and_then(Value::as_bool).unwrap_or(false),
1403        "data": value
1404            .pointer("/data/0")
1405            .and_then(Value::as_str)
1406            .unwrap_or_default(),
1407    })
1408}
1409
1410fn mint_info_json(address: &str, value: &Value) -> Value {
1411    let owner_program = value
1412        .get("owner")
1413        .and_then(Value::as_str)
1414        .unwrap_or_default();
1415    let info = value.pointer("/data/parsed/info");
1416    json!({
1417        "address": address,
1418        "ownerProgram": owner_program,
1419        "decimals": info.and_then(|v| v.get("decimals")).and_then(Value::as_u64),
1420        "supply": info.and_then(|v| v.get("supply")).and_then(Value::as_str),
1421        "mintAuthority": info.and_then(|v| v.get("mintAuthority")).and_then(Value::as_str),
1422        "freezeAuthority": info.and_then(|v| v.get("freezeAuthority")).and_then(Value::as_str),
1423    })
1424}
1425
1426fn token_account_json(address: &str, value: &Value) -> Value {
1427    let owner_program = value
1428        .get("owner")
1429        .and_then(Value::as_str)
1430        .unwrap_or_default();
1431    let info = value.pointer("/data/parsed/info");
1432    json!({
1433        "address": address,
1434        "ownerProgram": owner_program,
1435        "mint": info.and_then(|v| v.get("mint")).and_then(Value::as_str),
1436        "owner": info.and_then(|v| v.get("owner")).and_then(Value::as_str),
1437        "amount": info.and_then(|v| v.pointer("/tokenAmount/amount")).and_then(Value::as_str),
1438        "uiAmountString": info.and_then(|v| v.pointer("/tokenAmount/uiAmountString")).and_then(Value::as_str),
1439    })
1440}
1441
1442#[cfg(test)]
1443mod tests {
1444    use super::*;
1445
1446    #[test]
1447    fn cors_headers_allow_browser_sdk_reads() {
1448        let response = with_cors(
1449            Response::builder()
1450                .status(StatusCode::OK)
1451                .body(Full::new(Bytes::new()))
1452                .unwrap(),
1453        );
1454
1455        assert_eq!(response.headers()[ACCESS_CONTROL_ALLOW_ORIGIN], "*");
1456        assert_eq!(
1457            response.headers()[ACCESS_CONTROL_ALLOW_METHODS],
1458            "GET, POST, OPTIONS"
1459        );
1460        assert_eq!(
1461            response.headers()[ACCESS_CONTROL_ALLOW_HEADERS],
1462            "Authorization, Content-Type"
1463        );
1464    }
1465
1466    #[test]
1467    fn native_balance_serializes_u64_values_as_decimal_strings() {
1468        let value = contextual_native_balance_json(&json!({
1469            "context": { "slot": 9_007_199_254_740_995_u64 },
1470            "value": 9_007_199_254_740_993_u64,
1471        }))
1472        .unwrap();
1473
1474        assert_eq!(value["lamports"], "9007199254740993");
1475        assert_eq!(value["contextSlot"], "9007199254740995");
1476    }
1477
1478    #[test]
1479    fn rpc_read_config_propagates_min_context_slot() {
1480        let config = rpc_read_config(Some("jsonParsed"), Some(9_007_199_254_740_997));
1481
1482        assert_eq!(config["commitment"], "confirmed");
1483        assert_eq!(config["encoding"], "jsonParsed");
1484        assert_eq!(config["minContextSlot"], 9_007_199_254_740_997_u64);
1485    }
1486
1487    #[test]
1488    fn token_balance_preserves_raw_amount_and_stringifies_context_slot() {
1489        let value = contextual_token_balance_json(
1490            &json!({
1491                "context": { "slot": 9_007_199_254_740_995_u64 },
1492                "value": [{
1493                    "pubkey": "token-account",
1494                    "account": {
1495                        "data": {
1496                            "parsed": {
1497                                "info": {
1498                                    "mint": "mint",
1499                                    "tokenAmount": {
1500                                        "amount": "18446744073709551615",
1501                                        "decimals": 9,
1502                                        "uiAmountString": "18446744073.709551615"
1503                                    }
1504                                }
1505                            }
1506                        }
1507                    }
1508                }]
1509            }),
1510            "owner",
1511            "mint",
1512            None,
1513        )
1514        .unwrap();
1515
1516        assert_eq!(value["amount"], "18446744073709551615");
1517        assert_eq!(value["contextSlot"], "9007199254740995");
1518    }
1519
1520    #[test]
1521    fn balance_bodies_accept_decimal_string_min_context_slots() {
1522        let native: NativeBalanceBody = serde_json::from_value(json!({
1523            "address": "owner",
1524            "minContextSlot": "9007199254740997",
1525        }))
1526        .unwrap();
1527        let token: BalanceBody = serde_json::from_value(json!({
1528            "owner": "owner",
1529            "mint": "mint",
1530            "minContextSlot": "9007199254740999",
1531        }))
1532        .unwrap();
1533
1534        assert_eq!(
1535            parse_min_context_slot(native.min_context_slot.as_deref()).unwrap(),
1536            Some(9_007_199_254_740_997)
1537        );
1538        assert_eq!(
1539            parse_min_context_slot(token.min_context_slot.as_deref()).unwrap(),
1540            Some(9_007_199_254_740_999)
1541        );
1542    }
1543
1544    fn runtime_definition(reader: crate::ProgramAccountReaderFn) -> ProgramRuntimeDefinition {
1545        let program_spec_hash = crate::ProgramSpecHash::from_digest([1; 32]);
1546        let idl_content_hash = crate::IdlContentHash::from_digest([2; 32]);
1547        let normalized_idl_hash = crate::NormalizedIdlHash::from_digest([3; 32]);
1548        let program_release_hash = arete_hash::OssGeneratedProgramReleaseV1::new(
1549            "Program111",
1550            program_spec_hash,
1551            idl_content_hash,
1552            normalized_idl_hash,
1553        )
1554        .hash()
1555        .unwrap();
1556        ProgramRuntimeDefinition {
1557            program_id: "Program111".to_string(),
1558            program_spec_hash,
1559            idl_content_hash,
1560            normalized_idl_hash,
1561            program_release_hash,
1562            account_reader: reader,
1563        }
1564    }
1565
1566    fn program_read_context(
1567        target_id: &str,
1568        program_id: &str,
1569        release_hash: &str,
1570    ) -> crate::websocket::auth::AuthContext {
1571        crate::websocket::auth::AuthContext::from_claims(
1572            arete_auth::SessionClaims::program_read_builder(
1573                "issuer",
1574                "user:1",
1575                target_id,
1576                program_id,
1577                release_hash,
1578            )
1579            .build(),
1580        )
1581    }
1582
1583    #[test]
1584    fn program_read_auth_is_bound_to_target_program_and_release() {
1585        let definition = runtime_definition(Arc::new(|_, _| Ok(Value::Null)));
1586        let release_hash = definition.program_release_hash.to_string();
1587        let valid = program_read_context("binding-1", &definition.program_id, &release_hash);
1588
1589        assert_eq!(
1590            authorize_program_read(Some(&valid), Some("binding-1"), &definition),
1591            Ok(())
1592        );
1593
1594        let wrong_target = program_read_context("binding-2", &definition.program_id, &release_hash);
1595        let wrong_program = program_read_context("binding-1", "Program222", &release_hash);
1596        let wrong_release = program_read_context(
1597            "binding-1",
1598            &definition.program_id,
1599            "arete:h1:program-release:sha256:different",
1600        );
1601
1602        for context in [&wrong_target, &wrong_program, &wrong_release] {
1603            assert_eq!(
1604                authorize_program_read(Some(context), Some("binding-1"), &definition),
1605                Err(ProgramReadError::Unauthorized)
1606            );
1607        }
1608        assert_eq!(
1609            authorize_program_read(Some(&valid), None, &definition),
1610            Err(ProgramReadError::AuthorizationNotConfigured)
1611        );
1612        assert_eq!(authorize_program_read(None, None, &definition), Ok(()));
1613    }
1614
1615    #[test]
1616    fn release_routes_are_exact_and_legacy_program_routes_are_not_accepted() {
1617        let hash = crate::ProgramReleaseHash::from_digest([4; 32]);
1618        let fetch_path = format!("/v1/releases/{hash}/accounts/Vault/address");
1619        let exists_path = format!("{fetch_path}/exists");
1620        let batch_path = format!("/v1/releases/{hash}/accounts/Vault");
1621
1622        assert!(matches!(
1623            ProgramReadRoute::parse(&Method::GET, &fetch_path)
1624                .unwrap()
1625                .operation,
1626            ProgramReadOperation::Fetch { .. }
1627        ));
1628        assert!(matches!(
1629            ProgramReadRoute::parse(&Method::GET, &exists_path)
1630                .unwrap()
1631                .operation,
1632            ProgramReadOperation::Exists { .. }
1633        ));
1634        assert!(matches!(
1635            ProgramReadRoute::parse(&Method::POST, &batch_path)
1636                .unwrap()
1637                .operation,
1638            ProgramReadOperation::Batch
1639        ));
1640        assert_eq!(
1641            ProgramReadRoute::parse(&Method::GET, "/programs/demo/accounts/Vault/address")
1642                .unwrap_err(),
1643            ProgramReadError::NotFound
1644        );
1645    }
1646
1647    #[tokio::test]
1648    async fn owner_mismatch_is_rejected_before_decoder_execution() {
1649        use std::sync::atomic::{AtomicBool, Ordering};
1650
1651        let called = Arc::new(AtomicBool::new(false));
1652        let called_by_reader = called.clone();
1653        let definition = runtime_definition(Arc::new(move |_, _| {
1654            called_by_reader.store(true, Ordering::SeqCst);
1655            Ok(json!({ "decoded": true }))
1656        }));
1657        let value = json!({
1658            "owner": "DifferentProgram",
1659            "data": [BASE64_STANDARD.encode([1, 2, 3]), "base64"]
1660        });
1661
1662        assert!(matches!(
1663            decode_release_account(&definition, "Vault", &value).await,
1664            AccountReadOutcome::Error(ProgramReadError::AccountOwnerMismatch)
1665        ));
1666        assert!(!called.load(Ordering::SeqCst));
1667    }
1668
1669    #[tokio::test]
1670    async fn decode_failures_remain_typed_errors_and_metadata_is_public_only() {
1671        let definition = runtime_definition(Arc::new(|_, _| anyhow::bail!("private diagnostic")));
1672        let value = json!({
1673            "owner": "Program111",
1674            "data": [BASE64_STANDARD.encode([1, 2, 3]), "base64"]
1675        });
1676        assert!(matches!(
1677            decode_release_account(&definition, "Vault", &value).await,
1678            AccountReadOutcome::Error(ProgramReadError::AccountDecodeFailed)
1679        ));
1680
1681        let response = with_release_metadata(
1682            program_read_error_response(ProgramReadError::AccountDecodeFailed),
1683            &definition,
1684            Some("address"),
1685            Some(true),
1686        );
1687        assert_eq!(response.headers()["X-Error-Code"], "ACCOUNT_DECODE_FAILED");
1688        assert_eq!(
1689            response.headers()["X-Arete-Program-Release-Hash"],
1690            definition.program_release_hash.to_string()
1691        );
1692        assert_eq!(
1693            response.headers()["X-Arete-Idl-Content-Hash"],
1694            definition.idl_content_hash.to_string()
1695        );
1696        let body = response.into_body().collect().await.unwrap().to_bytes();
1697        let body = std::str::from_utf8(&body).unwrap();
1698        assert_eq!(body, r#"{"error":{"code":"ACCOUNT_DECODE_FAILED"}}"#);
1699        assert!(!body.contains("private diagnostic"));
1700        assert!(!body.contains("decoder"));
1701    }
1702}