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#[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
72pub 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 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 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 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 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 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}