1use axum::{
2 body::Body,
3 extract::ws::{Message, WebSocketUpgrade},
4 extract::{FromRequestParts, Path, Query, Request, State},
5 http::{header, request::Parts, StatusCode},
6 middleware::{self, Next},
7 response::{IntoResponse, Json, Response},
8 routing::{delete, get, post, put},
9 Router,
10};
11use dashmap::DashMap;
12use meow_common::TunnelMode;
13use meow_config::{
14 proxy_provider::ProxyProvider,
15 raw::{RawConfig, RawProxyGroup, RawSubscription},
16 rule_provider::RuleProvider,
17 NamedListener,
18};
19use meow_tunnel::Tunnel;
20use parking_lot::RwLock;
21use serde::{Deserialize, Serialize};
22use std::collections::{BTreeMap, HashMap};
23use std::sync::Arc;
24use std::time::Duration;
25use tokio::sync::{broadcast, Mutex};
26use tower_http::cors::CorsLayer;
27use tracing::{debug, info};
28
29#[cfg(feature = "listener-tun")]
30use meow_listener::TunListener;
31
32use crate::log_stream::{parse_log_level, LogMessage};
33use crate::ui;
34
35struct MaybeWebSocket(Option<WebSocketUpgrade>);
36
37impl<S> FromRequestParts<S> for MaybeWebSocket
38where
39 S: Send + Sync,
40{
41 type Rejection = std::convert::Infallible;
42
43 async fn from_request_parts(parts: &mut Parts, state: &S) -> Result<Self, Self::Rejection> {
44 let is_websocket = parts
45 .headers
46 .get(header::UPGRADE)
47 .and_then(|v| v.to_str().ok())
48 .is_some_and(|v| v.eq_ignore_ascii_case("websocket"));
49 if !is_websocket {
50 return Ok(Self(None));
51 }
52 Ok(Self(
53 WebSocketUpgrade::from_request_parts(parts, state)
54 .await
55 .ok(),
56 ))
57 }
58}
59
60pub struct AppState {
61 pub tunnel: Tunnel,
62 pub secret: Option<String>,
64 pub config_path: String,
65 pub raw_config: Arc<RwLock<RawConfig>>,
66 pub log_tx: broadcast::Sender<LogMessage>,
68 pub proxy_providers: Arc<DashMap<String, Arc<ProxyProvider>>>,
70 pub rule_providers: Arc<RwLock<HashMap<String, Arc<RuleProvider>>>>,
71 pub listeners: Vec<NamedListener>,
73 pub external_ui: Option<std::path::PathBuf>,
76 pub config_mutation_lock: tokio::sync::Mutex<()>,
83}
84
85static CONFIG_MUTATION: Mutex<()> = Mutex::const_new(());
88
89impl AppState {
90 fn auth_required(&self) -> bool {
91 self.secret.as_deref().is_some_and(|s| !s.is_empty())
92 }
93}
94
95async fn require_auth_ws(
99 State(state): State<Arc<AppState>>,
100 Query(query): Query<HashMap<String, String>>,
101 req: Request,
102 next: Next,
103) -> Response {
104 if !state.auth_required() {
105 return next.run(req).await;
106 }
107 let expected = state.secret.as_deref().unwrap_or("");
108
109 let bearer = req
110 .headers()
111 .get(header::AUTHORIZATION)
112 .and_then(|v| v.to_str().ok())
113 .and_then(|v| v.strip_prefix("Bearer "));
114
115 let is_websocket = req
116 .headers()
117 .get(header::UPGRADE)
118 .and_then(|v| v.to_str().ok())
119 .is_some_and(|v| v.eq_ignore_ascii_case("websocket"));
120 let token_param = if is_websocket {
121 query.get("token").map(std::string::String::as_str)
122 } else {
123 None
124 };
125 let provided = bearer.or(token_param);
126
127 let ok = match provided {
128 Some(t) if t.len() == expected.len() => {
129 use subtle::ConstantTimeEq;
130 t.as_bytes().ct_eq(expected.as_bytes()).into()
131 }
132 _ => false,
133 };
134 if ok {
135 next.run(req).await
136 } else {
137 (
138 StatusCode::UNAUTHORIZED,
139 Json(serde_json::json!({"message": "Unauthorized"})),
140 )
141 .into_response()
142 }
143}
144
145pub fn create_router(state: Arc<AppState>) -> Router {
146 let api = Router::new()
149 .route("/", get(hello))
150 .route("/version", get(version))
151 .route("/proxies", get(get_proxies))
152 .route(
153 "/proxies/{name}",
154 get(get_proxy).put(update_proxy).delete(unfix_proxy),
155 )
156 .route("/proxies/{name}/delay", get(get_proxy_delay))
157 .route("/group", get(get_groups))
158 .route("/group/{name}", get(get_group))
159 .route("/group/{name}/delay", get(get_group_delay))
160 .route(
161 "/rules",
162 get(get_rules).post(replace_rules).put(update_rule_at_index),
163 )
164 .route("/rules/{index}", delete(delete_rule))
165 .route("/rules/reorder", post(reorder_rules))
166 .route("/connections", get(get_connections))
167 .route("/connections/{id}", delete(close_connection))
168 .route("/connections", delete(close_all_connections))
169 .route(
170 "/configs",
171 get(get_configs).patch(update_configs).put(put_configs),
172 )
173 .route("/metrics", get(get_metrics))
174 .route("/traffic", get(get_traffic))
175 .route("/logs", get(get_logs))
176 .route("/memory", get(get_memory))
177 .route("/dns/results", get(get_dns_results))
178 .route("/dns/query", get(dns_query_get).post(dns_query))
179 .route("/cache/dns/flush", post(flush_dns_cache))
180 .route("/cache/fakeip/flush", post(flush_fakeip_cache))
181 .route("/api/config/save", post(save_config))
183 .route(
185 "/api/subscriptions",
186 get(get_subscriptions).post(add_subscription),
187 )
188 .route("/api/subscriptions/{name}", delete(delete_subscription))
189 .route(
190 "/api/subscriptions/{name}/refresh",
191 post(refresh_subscription),
192 )
193 .route(
195 "/api/proxy-groups",
196 get(get_proxy_groups).post(create_proxy_group),
197 )
198 .route(
199 "/api/proxy-groups/{name}",
200 put(update_proxy_group).delete(delete_proxy_group),
201 )
202 .route(
203 "/api/proxy-groups/{name}/select",
204 put(select_proxy_in_group),
205 )
206 .route("/providers/proxies", get(get_providers))
208 .route(
209 "/providers/proxies/{name}",
210 get(get_provider).put(refresh_provider),
211 )
212 .route(
213 "/providers/proxies/{name}/healthcheck",
214 get(provider_healthcheck),
215 )
216 .route(
217 "/providers/proxies/{provider_name}/{proxy_name}",
218 get(get_provider_proxy),
219 )
220 .route(
221 "/providers/proxies/{provider_name}/{proxy_name}/healthcheck",
222 get(provider_proxy_healthcheck),
223 )
224 .route("/providers/rules", get(get_rule_providers))
226 .route(
227 "/providers/rules/{name}",
228 get(get_rule_provider).put(refresh_rule_provider),
229 )
230 .route("/listeners", get(get_listeners))
232 .route_layer(middleware::from_fn_with_state(
233 Arc::clone(&state),
234 require_auth_ws,
235 ));
236
237 let router = api;
244 let router = if let Some(dir) = state.external_ui.clone() {
245 router.nest_service("/ui", tower_http::services::ServeDir::new(dir))
250 } else {
251 router
252 .route("/ui", get(ui::serve_ui))
253 .route("/ui/{*rest}", get(ui::serve_ui))
254 };
255
256 router.layer(CorsLayer::permissive()).with_state(state)
257}
258
259#[derive(Serialize)]
262struct HelloResponse {
263 hello: &'static str,
264}
265
266async fn hello() -> Json<HelloResponse> {
267 Json(HelloResponse { hello: "meow" })
268}
269
270#[derive(Serialize)]
271struct VersionResponse {
272 version: String,
273 meta: bool,
274}
275
276async fn version() -> Json<VersionResponse> {
277 Json(VersionResponse {
278 version: format!("v{}", env!("CARGO_PKG_VERSION")),
279 meta: true,
280 })
281}
282
283#[derive(Serialize)]
284struct ProxyInfo {
285 name: String,
286 #[serde(rename = "type")]
287 proxy_type: String,
288 alive: bool,
289 history: Vec<meow_common::DelayHistory>,
290 udp: bool,
291 #[serde(skip_serializing_if = "Option::is_none")]
293 all: Option<Vec<String>>,
294 #[serde(skip_serializing_if = "Option::is_none")]
296 now: Option<String>,
297 #[serde(skip_serializing_if = "Option::is_none")]
299 fixed: Option<String>,
300 #[serde(rename = "testUrl", skip_serializing_if = "Option::is_none")]
301 test_url: Option<String>,
302 #[serde(rename = "expectedStatus", skip_serializing_if = "Option::is_none")]
303 expected_status: Option<String>,
304 #[serde(skip_serializing_if = "Option::is_none")]
306 delay: Option<u16>,
307}
308
309impl ProxyInfo {
310 fn from_proxy(proxy: &Arc<dyn meow_common::Proxy>) -> Self {
311 let members = proxy.members();
312 let current = proxy.current();
313 debug!(
314 name = proxy.name(),
315 proxy_type = %proxy.adapter_type(),
316 member_count = members.as_ref().map(std::vec::Vec::len),
317 current = ?current,
318 "building ProxyInfo",
319 );
320 let delay = Some(proxy.last_delay()).filter(|&d| d > 0);
321 Self {
322 name: proxy.name().to_string(),
323 proxy_type: proxy.adapter_type().to_string(),
324 alive: proxy.alive(),
325 history: proxy.delay_history(),
326 udp: proxy.support_udp(),
327 all: members,
328 now: current,
329 fixed: proxy
330 .selection()
331 .and_then(meow_common::ProxySelection::fixed),
332 test_url: proxy.test_url().map(str::to_string),
333 expected_status: proxy.expected_status().map(str::to_string),
334 delay,
335 }
336 }
337}
338
339#[derive(Serialize)]
340struct ProxiesResponse {
341 proxies: std::collections::HashMap<String, ProxyInfo>,
342}
343
344async fn get_proxies(State(state): State<Arc<AppState>>) -> Json<ProxiesResponse> {
345 let route = state.tunnel.route_snapshot();
346 let mut result = std::collections::HashMap::new();
347 for (name, proxy) in &route.proxies {
348 result.insert(name.to_string(), ProxyInfo::from_proxy(proxy));
349 }
350 Json(ProxiesResponse { proxies: result })
351}
352
353async fn get_proxy(
354 State(state): State<Arc<AppState>>,
355 Path(name): Path<String>,
356) -> Result<Json<ProxyInfo>, StatusCode> {
357 let route = state.tunnel.route_snapshot();
358 let proxy = route
359 .proxies
360 .get(name.as_str())
361 .ok_or(StatusCode::NOT_FOUND)?;
362 Ok(Json(ProxyInfo::from_proxy(proxy)))
363}
364
365#[derive(Deserialize)]
366struct UpdateProxyRequest {
367 name: String,
368}
369
370async fn update_proxy(
371 State(state): State<Arc<AppState>>,
372 Path(group_name): Path<String>,
373 Json(body): Json<UpdateProxyRequest>,
374) -> Response {
375 let route = state.tunnel.route_snapshot();
376 let Some(proxy) = route.proxies.get(group_name.as_str()).cloned() else {
377 return msg_err(StatusCode::NOT_FOUND, "Resource not found");
378 };
379 let Some(selection) = proxy.selection() else {
380 return msg_err(StatusCode::BAD_REQUEST, "Must be a Selector");
381 };
382 match selection.set(&body.name).await {
383 Ok(()) => {
384 info!("Proxy group '{}' switched to '{}'", group_name, body.name);
385 StatusCode::NO_CONTENT.into_response()
386 }
387 Err(e) => (
388 StatusCode::BAD_REQUEST,
389 Json(serde_json::json!({"message": format!("Selector update error: {e}")})),
390 )
391 .into_response(),
392 }
393}
394
395async fn unfix_proxy(
396 State(state): State<Arc<AppState>>,
397 Path(group_name): Path<String>,
398) -> Response {
399 let route = state.tunnel.route_snapshot();
400 let Some(proxy) = route.proxies.get(group_name.as_str()).cloned() else {
401 return msg_err(StatusCode::NOT_FOUND, "Resource not found");
402 };
403 let Some(selection) = proxy.selection() else {
404 return msg_err(StatusCode::BAD_REQUEST, "Body invalid");
405 };
406 if !selection.can_unfix() {
407 return msg_err(StatusCode::BAD_REQUEST, "Body invalid");
408 }
409 selection.force_set(None);
410 StatusCode::NO_CONTENT.into_response()
411}
412
413async fn get_groups(State(state): State<Arc<AppState>>) -> Json<ProxiesResponse> {
414 let route = state.tunnel.route_snapshot();
415 let proxies = route
416 .proxies
417 .iter()
418 .filter(|(_, proxy)| proxy.members().is_some())
419 .map(|(name, proxy)| (name.to_string(), ProxyInfo::from_proxy(proxy)))
420 .collect();
421 Json(ProxiesResponse { proxies })
422}
423
424async fn get_group(State(state): State<Arc<AppState>>, Path(name): Path<String>) -> Response {
425 let route = state.tunnel.route_snapshot();
426 match route.proxies.get(name.as_str()) {
427 Some(proxy) if proxy.members().is_some() => {
428 Json(ProxyInfo::from_proxy(proxy)).into_response()
429 }
430 _ => msg_err(StatusCode::NOT_FOUND, "Resource not found"),
431 }
432}
433
434#[derive(Serialize)]
435struct RuleInfo<'a> {
436 index: usize,
437 #[serde(rename = "type")]
438 rule_type: &'static str,
439 payload: &'a str,
440 proxy: &'a str,
441 size: i64,
442}
443
444#[derive(Serialize)]
445struct RulesResponse<'a> {
446 rules: Vec<RuleInfo<'a>>,
447}
448
449async fn get_rules(State(state): State<Arc<AppState>>) -> Response {
450 let route = state.tunnel.route_snapshot();
453 let result: Vec<RuleInfo> = route
454 .rules
455 .iter()
456 .enumerate()
457 .map(|(index, r)| RuleInfo {
458 index,
459 rule_type: r.rule_type().as_str(),
460 payload: r.payload(),
461 proxy: r.adapter(),
462 size: -1,
463 })
464 .collect();
465 Json(RulesResponse { rules: result }).into_response()
466}
467
468#[derive(Serialize)]
469#[serde(rename_all = "camelCase")]
470struct ConnectionsResponse<'a> {
471 upload_total: i64,
472 download_total: i64,
473 memory: u64,
474 connections: meow_tunnel::statistics::ActiveConnectionsView<'a>,
479}
480
481#[derive(Deserialize)]
482struct ConnectionsParams {
483 interval: Option<String>,
484}
485
486const MIN_CONNECTIONS_INTERVAL_MS: u64 = 100;
493
494fn parse_connections_interval(raw: Option<&str>) -> Option<u64> {
499 match raw {
500 Some(raw) => match raw.parse::<u64>() {
501 Ok(0) | Err(_) => None,
502 Ok(value) => Some(value.max(MIN_CONNECTIONS_INTERVAL_MS)),
503 },
504 None => Some(1000),
505 }
506}
507
508async fn connections_json(state: &AppState) -> String {
509 let stats = state.tunnel.statistics();
510 let (up, down) = stats.snapshot();
511 let memory = read_rss_bytes().await;
512 #[allow(
513 clippy::unnecessary_cast,
514 reason = "no-op on 64-bit; widens i32 on targets without 64-bit atomics"
515 )]
516 let upload = up as i64;
517 #[allow(
518 clippy::unnecessary_cast,
519 reason = "no-op on 64-bit; widens i32 on targets without 64-bit atomics"
520 )]
521 let download = down as i64;
522 serde_json::to_string(&ConnectionsResponse {
523 upload_total: upload,
524 download_total: download,
525 memory,
526 connections: stats.active_connections_view(),
527 })
528 .unwrap_or_else(|_| {
529 "{\"uploadTotal\":0,\"downloadTotal\":0,\"memory\":0,\"connections\":[]}".into()
530 })
531}
532
533async fn get_connections(
534 State(state): State<Arc<AppState>>,
535 Query(params): Query<ConnectionsParams>,
536 MaybeWebSocket(ws): MaybeWebSocket,
537) -> Response {
538 let Some(interval_ms) = parse_connections_interval(params.interval.as_deref()) else {
539 return msg_err(StatusCode::BAD_REQUEST, "Body invalid");
540 };
541
542 if let Some(ws) = ws {
543 return ws.on_upgrade(move |mut socket| async move {
544 if socket
545 .send(Message::Text(connections_json(&state).await.into()))
546 .await
547 .is_err()
548 {
549 return;
550 }
551 let mut ticker = tokio::time::interval(Duration::from_millis(interval_ms));
552 ticker.tick().await;
553 loop {
554 ticker.tick().await;
555 if socket
556 .send(Message::Text(connections_json(&state).await.into()))
557 .await
558 .is_err()
559 {
560 break;
561 }
562 }
563 });
564 }
565
566 let body = connections_json(&state).await;
567 ([(header::CONTENT_TYPE, "application/json")], body).into_response()
568}
569
570async fn close_connection(
571 State(state): State<Arc<AppState>>,
572 Path(id): Path<String>,
573) -> StatusCode {
574 match uuid::Uuid::parse_str(&id) {
575 Ok(uuid) => {
576 state.tunnel.statistics().close_connection(uuid);
577 StatusCode::NO_CONTENT
578 }
579 Err(_) => StatusCode::BAD_REQUEST,
580 }
581}
582
583#[derive(Serialize)]
584struct ConfigResponse {
585 mode: String,
586 #[serde(rename = "log-level")]
587 log_level: String,
588 #[serde(rename = "mixed-port", skip_serializing_if = "Option::is_none")]
589 mixed_port: Option<u16>,
590 #[serde(rename = "socks-port", skip_serializing_if = "Option::is_none")]
591 socks_port: Option<u16>,
592 #[serde(rename = "port", skip_serializing_if = "Option::is_none")]
593 http_port: Option<u16>,
594 #[serde(rename = "redir-port")]
595 redir_port: u16,
596 #[serde(rename = "tproxy-port")]
597 tproxy_port: u16,
598 #[serde(
599 rename = "external-controller",
600 skip_serializing_if = "Option::is_none"
601 )]
602 external_controller: Option<String>,
603 #[serde(rename = "allow-lan")]
604 allow_lan: bool,
605 #[serde(rename = "bind-address")]
606 bind_address: String,
607 #[serde(rename = "ipv6")]
608 ipv6: bool,
609 #[serde(rename = "tun-enable")]
613 tun_enable: bool,
614}
615
616async fn get_configs(State(state): State<Arc<AppState>>) -> Json<ConfigResponse> {
617 let raw = state.raw_config.read();
618 Json(ConfigResponse {
619 mode: state.tunnel.mode().to_string(),
620 log_level: raw.log_level.clone().unwrap_or_else(|| "info".to_string()),
621 mixed_port: raw.mixed_port,
622 socks_port: raw.socks_port,
623 http_port: raw.port,
624 redir_port: 0,
625 tproxy_port: raw.tproxy_port.unwrap_or(0),
626 external_controller: raw.external_controller.clone(),
627 allow_lan: raw.allow_lan.unwrap_or(false),
628 bind_address: raw
629 .bind_address
630 .clone()
631 .unwrap_or_else(|| "0.0.0.0".to_string()),
632 ipv6: raw.ipv6.unwrap_or(false),
633 tun_enable: state.tunnel.has_tun(),
634 })
635}
636
637#[derive(Deserialize)]
638struct UpdateConfigRequest {
639 mode: Option<String>,
640 #[serde(rename = "log-level")]
641 log_level: Option<String>,
642}
643
644async fn update_configs(
645 State(state): State<Arc<AppState>>,
646 Json(body): Json<UpdateConfigRequest>,
647) -> Response {
648 let mode = body.mode.map(|s| s.parse::<TunnelMode>());
650 if let Some(Err(_)) = mode {
651 return msg_err(StatusCode::BAD_REQUEST, "Body invalid");
652 }
653 if let Some(ref level) = body.log_level {
654 if !matches!(
655 level.to_ascii_lowercase().as_str(),
656 "debug" | "info" | "warning" | "warn" | "error" | "silent"
657 ) {
658 return msg_err(StatusCode::BAD_REQUEST, "Body invalid");
659 }
660 }
661
662 let mut raw = state.raw_config.write();
664 if let Some(Ok(parsed_mode)) = mode {
665 state.tunnel.set_mode(parsed_mode);
666 raw.mode = Some(parsed_mode.to_string());
667 info!("Mode changed to {}", parsed_mode);
668 }
669 if let Some(level) = body.log_level {
670 if let Err(e) = crate::log_stream::reload_log_level(&level) {
671 return (
672 StatusCode::INTERNAL_SERVER_ERROR,
673 Json(serde_json::json!({"message": e})),
674 )
675 .into_response();
676 }
677 raw.log_level = Some(level);
678 }
679 StatusCode::NO_CONTENT.into_response()
680}
681
682#[derive(Serialize)]
683struct TrafficResponse {
684 up: i64,
685 down: i64,
686 #[serde(rename = "upTotal")]
687 up_total: i64,
688 #[serde(rename = "downTotal")]
689 down_total: i64,
690}
691
692fn traffic_json(state: &AppState) -> String {
693 let (up, down, up_total, down_total) = state.tunnel.statistics().traffic_snapshot();
694 #[allow(
695 clippy::useless_conversion,
696 reason = "identity on 64-bit; widens i32 on targets without 64-bit atomics"
697 )]
698 serde_json::to_string(&TrafficResponse {
699 up: up.into(),
700 down: down.into(),
701 up_total: up_total.into(),
702 down_total: down_total.into(),
703 })
704 .unwrap_or_default()
705}
706
707async fn get_traffic(
708 State(state): State<Arc<AppState>>,
709 MaybeWebSocket(ws): MaybeWebSocket,
710) -> Response {
711 if let Some(ws) = ws {
712 return ws.on_upgrade(move |mut socket| async move {
713 let mut ticker = tokio::time::interval(Duration::from_secs(1));
714 ticker.tick().await;
715 loop {
716 ticker.tick().await;
717 let frame = traffic_json(&state);
718 if socket.send(Message::Text(frame.into())).await.is_err() {
719 break;
720 }
721 }
722 });
723 }
724
725 let stream = futures::stream::unfold(state, |state| async move {
726 tokio::time::sleep(Duration::from_secs(1)).await;
727 let line = format!("{}\n", traffic_json(&state));
728 Some((Ok::<String, std::convert::Infallible>(line), state))
729 });
730 Response::builder()
731 .header(header::CONTENT_TYPE, "application/json")
732 .body(Body::from_stream(stream))
733 .expect("valid traffic stream response")
734}
735
736#[derive(Deserialize)]
737struct DnsQueryRequest {
738 name: String,
739 #[serde(rename = "type")]
740 qtype: Option<String>,
741}
742
743#[derive(Deserialize)]
744struct DnsResultsQuery {
745 search: Option<String>,
746 limit: Option<usize>,
747}
748
749#[derive(Serialize)]
750struct DnsResultEntry {
751 name: String,
752 ips: Vec<String>,
753 #[serde(skip_serializing_if = "Option::is_none")]
754 from_server: Option<String>,
755 ttl: u64,
756}
757
758async fn get_dns_results(
759 State(state): State<Arc<AppState>>,
760 Query(params): Query<DnsResultsQuery>,
761) -> Json<Vec<DnsResultEntry>> {
762 let limit = params.limit.unwrap_or(256).min(1024);
763 let results = state
764 .tunnel
765 .resolver()
766 .dns_results(params.search.as_deref(), limit)
767 .into_iter()
768 .map(|entry| DnsResultEntry {
769 name: entry.name,
770 ips: entry.ips.into_iter().map(|ip| ip.to_string()).collect(),
771 from_server: entry.source,
772 ttl: entry.ttl.as_secs(),
773 })
774 .collect();
775 Json(results)
776}
777
778async fn dns_query(
779 State(state): State<Arc<AppState>>,
780 Json(body): Json<DnsQueryRequest>,
781) -> Json<serde_json::Value> {
782 let resolver = state.tunnel.resolver();
783 let result = resolver.resolve_ip(&body.name).await;
784 let _ = body.qtype;
785 Json(serde_json::json!({ "name": body.name, "answer": result.map(|ip| ip.to_string()) }))
786}
787
788async fn dns_query_get(
791 State(state): State<Arc<AppState>>,
792 Query(params): Query<DnsQueryRequest>,
793) -> Response {
794 let enabled = state
795 .raw_config
796 .read()
797 .dns
798 .as_ref()
799 .is_some_and(|dns| dns.enable.unwrap_or(false));
800 if !enabled {
801 return (
802 StatusCode::INTERNAL_SERVER_ERROR,
803 Json(serde_json::json!({"message": "DNS section is disabled"})),
804 )
805 .into_response();
806 }
807
808 use hickory_proto::rr::RecordType;
809 let qtype_text = params.qtype.as_deref().unwrap_or("A").to_ascii_uppercase();
810 let Ok(record_type) = qtype_text.parse::<RecordType>() else {
811 return (
812 StatusCode::BAD_REQUEST,
813 Json(serde_json::json!({"message": "invalid query type"})),
814 )
815 .into_response();
816 };
817
818 let resolver = state.tunnel.resolver();
819 let fqdn = if params.name.ends_with('.') {
820 params.name.clone()
821 } else {
822 format!("{}.", params.name)
823 };
824 let question = serde_json::json!({
825 "Name": fqdn,
826 "Qtype": u16::from(record_type),
827 "Qclass": 1,
828 });
829
830 let mut response = serde_json::Map::new();
831 response.insert("Status".into(), 0.into());
832 response.insert("Question".into(), serde_json::Value::Array(vec![question]));
833 response.insert("TC".into(), false.into());
834 response.insert("RD".into(), true.into());
835 response.insert("RA".into(), true.into());
836 response.insert("AD".into(), false.into());
837 response.insert("CD".into(), false.into());
838
839 if matches!(record_type, RecordType::A | RecordType::AAAA) {
840 let ips = resolver.resolve_ips(¶ms.name).await.unwrap_or_default();
841 let answers: Vec<_> = ips
842 .into_iter()
843 .filter(|ip| {
844 matches!(record_type, RecordType::A) && ip.is_ipv4()
845 || matches!(record_type, RecordType::AAAA) && ip.is_ipv6()
846 })
847 .map(|ip| {
848 serde_json::json!({
849 "name": fqdn,
850 "type": u16::from(record_type),
851 "TTL": 60,
852 "data": ip.to_string(),
853 })
854 })
855 .collect();
856 if !answers.is_empty() {
857 response.insert("Answer".into(), serde_json::Value::Array(answers));
858 }
859 } else if let Some(message) = resolver.forward_generic(¶ms.name, record_type).await {
860 let metadata = &message.metadata;
861 response.insert("Status".into(), u16::from(metadata.response_code).into());
862 response.insert("TC".into(), metadata.truncation.into());
863 response.insert("RD".into(), metadata.recursion_desired.into());
864 response.insert("RA".into(), metadata.recursion_available.into());
865 response.insert("AD".into(), metadata.authentic_data.into());
866 response.insert("CD".into(), metadata.checking_disabled.into());
867 insert_dns_records(&mut response, "Answer", &message.answers);
868 insert_dns_records(&mut response, "Authority", &message.authorities);
869 insert_dns_records(&mut response, "Additional", &message.additionals);
870 } else {
871 return (
872 StatusCode::INTERNAL_SERVER_ERROR,
873 Json(serde_json::json!({"message": "DNS query failed"})),
874 )
875 .into_response();
876 }
877
878 Json(serde_json::Value::Object(response)).into_response()
879}
880
881fn insert_dns_records(
882 target: &mut serde_json::Map<String, serde_json::Value>,
883 key: &str,
884 records: &[hickory_proto::rr::Record],
885) {
886 if records.is_empty() {
887 return;
888 }
889 target.insert(
890 key.to_string(),
891 serde_json::Value::Array(
892 records
893 .iter()
894 .map(|record| {
895 serde_json::json!({
896 "name": record.name.to_string(),
897 "type": u16::from(record.record_type()),
898 "TTL": record.ttl,
899 "data": record.data.to_string(),
900 })
901 })
902 .collect(),
903 ),
904 );
905}
906
907async fn flush_dns_cache(State(state): State<Arc<AppState>>) -> StatusCode {
908 state.tunnel.resolver().clear_cache();
909 StatusCode::NO_CONTENT
910}
911
912async fn flush_fakeip_cache(
916 State(state): State<Arc<AppState>>,
917) -> Result<StatusCode, (StatusCode, Json<serde_json::Value>)> {
918 match state.tunnel.resolver().flush_fake_ip() {
919 Ok(()) => Ok(StatusCode::NO_CONTENT),
920 Err(e) => Err((
921 StatusCode::BAD_REQUEST,
922 Json(serde_json::json!({ "message": e.to_string() })),
923 )),
924 }
925}
926
927async fn close_all_connections(State(state): State<Arc<AppState>>) -> StatusCode {
928 state.tunnel.statistics().close_all_connections();
929 StatusCode::NO_CONTENT
930}
931
932async fn save_config(
935 State(state): State<Arc<AppState>>,
936) -> Result<Json<serde_json::Value>, (StatusCode, String)> {
937 let raw = state.raw_config.read().clone();
938 meow_config::save_raw_config_async(&state.config_path, &raw)
939 .await
940 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
941 Ok(Json(serde_json::json!({"message": "config saved"})))
942}
943
944async fn apply_raw_to_tunnel(
951 mut raw: RawConfig,
952 state: &AppState,
953) -> Result<(), (StatusCode, String)> {
954 let expected_groups: Vec<String> = raw
955 .proxy_groups
956 .as_deref()
957 .unwrap_or_default()
958 .iter()
959 .map(|group| group.name.clone())
960 .collect();
961 if let Some(ps) = raw.proxies.as_mut() {
962 meow_config::ech_dns::preresolve_ech(ps).await;
963 }
964 let providers = state
965 .proxy_providers
966 .iter()
967 .map(|entry| (entry.key().clone(), Arc::clone(entry.value())))
968 .collect();
969 let (proxies, rules) =
970 rebuild_from_raw_with_resolver_async(raw, Arc::clone(state.tunnel.resolver()), providers)
971 .await
972 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e))?;
973 if let Some(missing) = expected_groups
974 .iter()
975 .find(|name| !proxies.contains_key(name.as_str()))
976 {
977 return Err((
978 StatusCode::BAD_REQUEST,
979 format!("proxy group '{missing}' failed validation"),
980 ));
981 }
982 state.tunnel.update_proxies(proxies);
983 state.tunnel.update_rules(rules);
984 Ok(())
985}
986
987async fn commit_raw_candidate(
988 state: &AppState,
989 candidate: RawConfig,
990) -> Result<(), (StatusCode, String)> {
991 apply_raw_to_tunnel(candidate.clone(), state).await?;
992 swap_config_and_reconcile_tun(state, candidate).await;
993 Ok(())
994}
995
996async fn rebuild_from_raw_with_resolver_async(
997 raw: RawConfig,
998 resolver: Arc<meow_dns::Resolver>,
999 providers: HashMap<String, Arc<ProxyProvider>>,
1000) -> Result<meow_config::RebuildResult, String> {
1001 tokio::task::spawn_blocking(move || {
1002 meow_config::rebuild_from_raw_runtime(&raw, Some(resolver), &providers)
1003 })
1004 .await
1005 .map_err(|e| format!("config rebuild task failed: {e}"))?
1006 .map_err(|e| e.to_string())
1007}
1008
1009#[derive(Serialize)]
1013struct SubscriptionInfo {
1014 name: String,
1015 url: String,
1016 interval: Option<u64>,
1017 last_updated: Option<i64>,
1018 proxy_count: usize,
1019 group_count: usize,
1020 rule_count: usize,
1021}
1022
1023async fn get_subscriptions(State(state): State<Arc<AppState>>) -> Json<Vec<SubscriptionInfo>> {
1024 let raw = state.raw_config.read();
1025 let subs = raw.subscriptions.as_deref().unwrap_or(&[]);
1026 let result: Vec<SubscriptionInfo> = subs
1027 .iter()
1028 .map(|s| SubscriptionInfo {
1029 name: s.name.clone(),
1030 url: s.url.clone(),
1031 interval: s.interval,
1032 last_updated: s.last_updated,
1033 proxy_count: raw.proxies.as_ref().map_or(0, std::vec::Vec::len),
1034 group_count: raw.proxy_groups.as_ref().map_or(0, std::vec::Vec::len),
1035 rule_count: raw.rules.as_ref().map_or(0, std::vec::Vec::len),
1036 })
1037 .collect();
1038 Json(result)
1039}
1040
1041#[derive(Deserialize)]
1042struct AddSubscriptionRequest {
1043 name: String,
1044 url: String,
1045 interval: Option<u64>,
1046}
1047
1048async fn add_subscription(
1049 State(state): State<Arc<AppState>>,
1050 Json(body): Json<AddSubscriptionRequest>,
1051) -> Result<Json<serde_json::Value>, (StatusCode, String)> {
1052 let fetched = meow_config::subscription::fetch_subscription(&body.url)
1053 .await
1054 .map_err(|e| (StatusCode::BAD_REQUEST, format!("fetch failed: {e}")))?;
1055
1056 let now = std::time::SystemTime::now()
1057 .duration_since(std::time::UNIX_EPOCH)
1058 .unwrap_or_default()
1059 .as_secs() as i64;
1060
1061 let pc = fetched.proxies.len();
1062 let gc = fetched.proxy_groups.len();
1063 let rc = fetched.rules.len();
1064
1065 let _mutation = CONFIG_MUTATION.lock().await;
1066 let snapshot = {
1067 let mut raw = state.raw_config.read().clone();
1068
1069 if let Some(ref subs) = raw.subscriptions {
1070 if subs.iter().any(|s| s.name == body.name) {
1071 return Err((
1072 StatusCode::CONFLICT,
1073 "subscription name already exists".into(),
1074 ));
1075 }
1076 }
1077
1078 let sub = RawSubscription {
1079 name: body.name.clone(),
1080 url: body.url.clone(),
1081 interval: body.interval,
1082 last_updated: Some(now),
1083 };
1084 raw.subscriptions.get_or_insert_with(Vec::new).push(sub);
1085
1086 raw.proxies = Some(fetched.proxies);
1088 raw.proxy_groups = Some(fetched.proxy_groups);
1089 raw.rules = Some(fetched.rules);
1090
1091 raw
1092 };
1093 commit_raw_candidate(&state, snapshot.clone()).await?;
1094
1095 meow_config::save_raw_config_async(&state.config_path, &snapshot)
1097 .await
1098 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
1099
1100 Ok(Json(serde_json::json!({
1101 "message": "subscription added",
1102 "proxy_count": pc, "group_count": gc, "rule_count": rc
1103 })))
1104}
1105
1106async fn delete_subscription(
1107 State(state): State<Arc<AppState>>,
1108 Path(name): Path<String>,
1109) -> Result<StatusCode, (StatusCode, String)> {
1110 let _mutation = CONFIG_MUTATION.lock().await;
1111 let snapshot = {
1112 let mut raw = state.raw_config.read().clone();
1113
1114 if let Some(ref mut subs) = raw.subscriptions {
1115 let before = subs.len();
1116 subs.retain(|s| s.name != name);
1117 if subs.len() == before {
1118 return Err((StatusCode::NOT_FOUND, "subscription not found".into()));
1119 }
1120 } else {
1121 return Err((StatusCode::NOT_FOUND, "no subscriptions".into()));
1122 }
1123
1124 raw.proxies = Some(Vec::new());
1126 raw.proxy_groups = Some(Vec::new());
1127 raw.rules = Some(Vec::new());
1128
1129 raw
1130 };
1131 commit_raw_candidate(&state, snapshot.clone()).await?;
1132 meow_config::save_raw_config_async(&state.config_path, &snapshot)
1133 .await
1134 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
1135 Ok(StatusCode::NO_CONTENT)
1136}
1137
1138async fn refresh_subscription(
1139 State(state): State<Arc<AppState>>,
1140 Path(name): Path<String>,
1141) -> Result<Json<serde_json::Value>, (StatusCode, String)> {
1142 let url = {
1143 let raw = state.raw_config.read();
1144 raw.subscriptions
1145 .as_ref()
1146 .and_then(|subs| subs.iter().find(|s| s.name == name))
1147 .map(|s| s.url.clone())
1148 .ok_or_else(|| (StatusCode::NOT_FOUND, "subscription not found".into()))?
1149 };
1150
1151 let fetched = meow_config::subscription::fetch_subscription(&url)
1152 .await
1153 .map_err(|e| (StatusCode::BAD_REQUEST, format!("fetch failed: {e}")))?;
1154
1155 let now = std::time::SystemTime::now()
1156 .duration_since(std::time::UNIX_EPOCH)
1157 .unwrap_or_default()
1158 .as_secs() as i64;
1159
1160 let pc = fetched.proxies.len();
1161 let gc = fetched.proxy_groups.len();
1162 let rc = fetched.rules.len();
1163
1164 let _mutation = CONFIG_MUTATION.lock().await;
1165 let snapshot = {
1166 let mut raw = state.raw_config.read().clone();
1167
1168 if let Some(ref mut subs) = raw.subscriptions {
1169 if let Some(sub) = subs.iter_mut().find(|s| s.name == name) {
1170 sub.last_updated = Some(now);
1171 }
1172 }
1173
1174 raw.proxies = Some(fetched.proxies);
1175 raw.proxy_groups = Some(fetched.proxy_groups);
1176 raw.rules = Some(fetched.rules);
1177
1178 raw
1179 };
1180 commit_raw_candidate(&state, snapshot.clone()).await?;
1181
1182 meow_config::save_raw_config_async(&state.config_path, &snapshot)
1184 .await
1185 .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
1186
1187 Ok(Json(serde_json::json!({
1188 "message": "subscription refreshed",
1189 "proxy_count": pc, "group_count": gc, "rule_count": rc
1190 })))
1191}
1192
1193#[derive(Serialize)]
1196struct ProxyGroupInfo {
1197 name: String,
1198 #[serde(rename = "type")]
1199 group_type: String,
1200 proxies: Vec<String>,
1201 now: Option<String>,
1202 url: Option<String>,
1203 interval: Option<u64>,
1204 tolerance: Option<u16>,
1205}
1206
1207async fn get_proxy_groups(State(state): State<Arc<AppState>>) -> Json<Vec<ProxyGroupInfo>> {
1208 let raw = state.raw_config.read();
1209 let groups = raw.proxy_groups.as_deref().unwrap_or(&[]);
1210 let route = state.tunnel.route_snapshot();
1211 let tunnel_proxies = &route.proxies;
1212
1213 let result: Vec<ProxyGroupInfo> = groups
1214 .iter()
1215 .map(|g| {
1216 let runtime = tunnel_proxies.get(g.name.as_str());
1217 let now = runtime.and_then(|p| p.current());
1218 let proxies = runtime
1219 .and_then(|p| p.members())
1220 .unwrap_or_else(|| g.proxies.clone().unwrap_or_default());
1221 ProxyGroupInfo {
1222 name: g.name.clone(),
1223 group_type: g.group_type.clone(),
1224 proxies,
1225 now,
1226 url: g.url.clone(),
1227 interval: g.interval,
1228 tolerance: g.tolerance,
1229 }
1230 })
1231 .collect();
1232 Json(result)
1233}
1234
1235#[derive(Deserialize)]
1236struct CreateProxyGroupRequest {
1237 name: String,
1238 #[serde(rename = "type")]
1239 group_type: String,
1240 proxies: Vec<String>,
1241 url: Option<String>,
1242 interval: Option<u64>,
1243 tolerance: Option<u16>,
1244}
1245
1246async fn create_proxy_group(
1247 State(state): State<Arc<AppState>>,
1248 Json(body): Json<CreateProxyGroupRequest>,
1249) -> Result<Json<serde_json::Value>, (StatusCode, String)> {
1250 let group_name = body.name.clone();
1251 let _mutation = CONFIG_MUTATION.lock().await;
1252 let snapshot = {
1253 let mut raw = state.raw_config.read().clone();
1254 if let Some(ref groups) = raw.proxy_groups {
1255 if groups.iter().any(|g| g.name == body.name) {
1256 return Err((StatusCode::CONFLICT, "group name already exists".into()));
1257 }
1258 }
1259 let group = RawProxyGroup {
1260 name: body.name,
1261 group_type: body.group_type,
1262 proxies: Some(body.proxies),
1263 url: body.url,
1264 interval: body.interval,
1265 tolerance: body.tolerance,
1266 ..Default::default()
1267 };
1268 raw.proxy_groups.get_or_insert_with(Vec::new).push(group);
1269 raw
1270 };
1271 commit_raw_candidate(&state, snapshot).await?;
1272 Ok(Json(
1273 serde_json::json!({"message": "group created", "name": group_name}),
1274 ))
1275}
1276
1277async fn update_proxy_group(
1278 State(state): State<Arc<AppState>>,
1279 Path(name): Path<String>,
1280 Json(body): Json<CreateProxyGroupRequest>,
1281) -> Result<StatusCode, (StatusCode, String)> {
1282 let _mutation = CONFIG_MUTATION.lock().await;
1283 let snapshot = {
1284 let mut raw = state.raw_config.read().clone();
1285 let group = raw
1286 .proxy_groups
1287 .as_mut()
1288 .and_then(|groups| groups.iter_mut().find(|g| g.name == name))
1289 .ok_or_else(|| (StatusCode::NOT_FOUND, "group not found".into()))?;
1290 group.group_type = body.group_type;
1291 group.proxies = Some(body.proxies);
1292 group.url = body.url;
1293 group.interval = body.interval;
1294 group.tolerance = body.tolerance;
1295 raw
1296 };
1297 commit_raw_candidate(&state, snapshot).await?;
1298 Ok(StatusCode::NO_CONTENT)
1299}
1300
1301async fn delete_proxy_group(
1302 State(state): State<Arc<AppState>>,
1303 Path(name): Path<String>,
1304) -> Result<StatusCode, (StatusCode, String)> {
1305 let _mutation = CONFIG_MUTATION.lock().await;
1306 let snapshot = {
1307 let mut raw = state.raw_config.read().clone();
1308 if let Some(ref mut groups) = raw.proxy_groups {
1309 let before = groups.len();
1310 groups.retain(|g| g.name != name);
1311 if groups.len() == before {
1312 return Err((StatusCode::NOT_FOUND, "group not found".into()));
1313 }
1314 } else {
1315 return Err((StatusCode::NOT_FOUND, "no groups".into()));
1316 }
1317 if let Some(ref mut rules) = raw.rules {
1318 rules.retain(|r| {
1319 let parts: Vec<&str> = r.split(',').collect();
1320 parts.last().is_none_or(|target| target.trim() != name)
1321 });
1322 }
1323 raw
1324 };
1325 commit_raw_candidate(&state, snapshot).await?;
1326 Ok(StatusCode::NO_CONTENT)
1327}
1328
1329#[derive(Deserialize)]
1330struct SelectProxyRequest {
1331 name: String,
1332}
1333
1334async fn select_proxy_in_group(
1335 State(state): State<Arc<AppState>>,
1336 Path(group_name): Path<String>,
1337 Json(body): Json<SelectProxyRequest>,
1338) -> StatusCode {
1339 let route = state.tunnel.route_snapshot();
1340 let Some(proxy) = route.proxies.get(group_name.as_str()).cloned() else {
1341 return StatusCode::NOT_FOUND;
1342 };
1343 let Some(selection) = proxy.selection() else {
1344 return StatusCode::BAD_REQUEST;
1345 };
1346 match selection.set(&body.name).await {
1347 Ok(()) => {
1348 info!("Proxy group '{}' switched to '{}'", group_name, body.name);
1349 StatusCode::NO_CONTENT
1350 }
1351 Err(_) => StatusCode::BAD_REQUEST,
1352 }
1353}
1354
1355#[derive(Deserialize)]
1358struct ReplaceRulesRequest {
1359 rules: Vec<String>,
1360}
1361
1362async fn replace_rules(
1363 State(state): State<Arc<AppState>>,
1364 Json(body): Json<ReplaceRulesRequest>,
1365) -> Result<StatusCode, (StatusCode, String)> {
1366 let _mutation = CONFIG_MUTATION.lock().await;
1367 let snapshot = {
1368 let mut raw = state.raw_config.read().clone();
1369 raw.rules = Some(body.rules);
1370 raw
1371 };
1372 commit_raw_candidate(&state, snapshot).await?;
1373 Ok(StatusCode::NO_CONTENT)
1374}
1375
1376#[derive(Deserialize)]
1377struct UpdateRuleRequest {
1378 index: usize,
1379 rule: String,
1380}
1381
1382async fn update_rule_at_index(
1383 State(state): State<Arc<AppState>>,
1384 Json(body): Json<UpdateRuleRequest>,
1385) -> Result<StatusCode, (StatusCode, String)> {
1386 let _mutation = CONFIG_MUTATION.lock().await;
1387 let snapshot = {
1388 let mut raw = state.raw_config.read().clone();
1389 let rules = raw.rules.get_or_insert_with(Vec::new);
1390 if body.index >= rules.len() {
1391 return Err((StatusCode::BAD_REQUEST, "index out of range".into()));
1392 }
1393 rules[body.index] = body.rule;
1394 raw
1395 };
1396 commit_raw_candidate(&state, snapshot).await?;
1397 Ok(StatusCode::NO_CONTENT)
1398}
1399
1400async fn delete_rule(
1401 State(state): State<Arc<AppState>>,
1402 Path(index): Path<usize>,
1403) -> Result<StatusCode, (StatusCode, String)> {
1404 let _mutation = CONFIG_MUTATION.lock().await;
1405 let snapshot = {
1406 let mut raw = state.raw_config.read().clone();
1407 let rules = raw.rules.get_or_insert_with(Vec::new);
1408 if index >= rules.len() {
1409 return Err((StatusCode::BAD_REQUEST, "index out of range".into()));
1410 }
1411 rules.remove(index);
1412 raw
1413 };
1414 commit_raw_candidate(&state, snapshot).await?;
1415 Ok(StatusCode::NO_CONTENT)
1416}
1417
1418#[derive(Deserialize)]
1419struct ReorderRulesRequest {
1420 from: usize,
1421 to: usize,
1422}
1423
1424async fn reorder_rules(
1425 State(state): State<Arc<AppState>>,
1426 Json(body): Json<ReorderRulesRequest>,
1427) -> Result<StatusCode, (StatusCode, String)> {
1428 let _mutation = CONFIG_MUTATION.lock().await;
1429 let snapshot = {
1430 let mut raw = state.raw_config.read().clone();
1431 let rules = raw.rules.get_or_insert_with(Vec::new);
1432 if body.from >= rules.len() || body.to >= rules.len() {
1433 return Err((StatusCode::BAD_REQUEST, "index out of range".into()));
1434 }
1435 let rule = rules.remove(body.from);
1436 rules.insert(body.to, rule);
1437 raw
1438 };
1439 commit_raw_candidate(&state, snapshot).await?;
1440 Ok(StatusCode::NO_CONTENT)
1441}
1442
1443#[derive(Deserialize)]
1451struct DelayParams {
1452 url: Option<String>,
1453 timeout: Option<String>,
1454 expected: Option<String>,
1455}
1456
1457#[derive(Serialize)]
1458struct DelayResp {
1459 delay: u16,
1460}
1461
1462fn msg_err(status: StatusCode, message: &'static str) -> Response {
1464 (status, Json(serde_json::json!({ "message": message }))).into_response()
1465}
1466
1467fn parse_delay_params(params: &DelayParams) -> Result<Duration, Box<Response>> {
1471 let url = params.url.as_deref().unwrap_or("").trim();
1474 if url.is_empty() {
1475 return Err(Box::new(msg_err(StatusCode::BAD_REQUEST, "Body invalid")));
1476 }
1477
1478 let timeout_str = params
1481 .timeout
1482 .as_deref()
1483 .ok_or_else(|| Box::new(msg_err(StatusCode::BAD_REQUEST, "Body invalid")))?;
1484 let timeout_ms: u16 = timeout_str
1485 .trim()
1486 .parse()
1487 .map_err(|_| Box::new(msg_err(StatusCode::BAD_REQUEST, "Body invalid")))?;
1488 if timeout_ms == 0 {
1489 return Err(Box::new(msg_err(StatusCode::BAD_REQUEST, "Body invalid")));
1490 }
1491 Ok(Duration::from_millis(timeout_ms as u64))
1492}
1493
1494async fn probe_and_record(
1498 proxy: &Arc<dyn meow_common::Proxy>,
1499 url: &str,
1500 expected: Option<&str>,
1501 timeout: Duration,
1502) -> Result<u16, meow_proxy::health::UrlTestError> {
1503 meow_proxy::health::probe_and_record(proxy, url, expected, timeout).await
1504}
1505
1506async fn get_proxy_delay(
1507 State(state): State<Arc<AppState>>,
1508 Path(name): Path<String>,
1509 Query(params): Query<DelayParams>,
1510) -> Response {
1511 let timeout = match parse_delay_params(¶ms) {
1512 Ok(t) => t,
1513 Err(resp) => return *resp,
1514 };
1515 let url = params.url.as_deref().unwrap_or("").to_string();
1516 let expected = params.expected.clone();
1517
1518 let route = state.tunnel.route_snapshot();
1519 let Some(proxy) = route.proxies.get(name.as_str()).cloned() else {
1521 return msg_err(StatusCode::NOT_FOUND, "resource not found");
1522 };
1523 drop(route);
1524
1525 match probe_and_record(&proxy, &url, expected.as_deref(), timeout).await {
1526 Ok(delay) => Json(DelayResp { delay }).into_response(),
1527 Err(meow_proxy::health::UrlTestError::Timeout) => {
1529 msg_err(StatusCode::GATEWAY_TIMEOUT, "Timeout")
1530 }
1531 Err(meow_proxy::health::UrlTestError::Transport(_)) => msg_err(
1533 StatusCode::SERVICE_UNAVAILABLE,
1534 "An error occurred in the delay test",
1535 ),
1536 }
1537}
1538
1539async fn get_group_delay(
1540 State(state): State<Arc<AppState>>,
1541 Path(name): Path<String>,
1542 Query(params): Query<DelayParams>,
1543) -> Response {
1544 let route = state.tunnel.route_snapshot();
1545 let Some(group) = route.proxies.get(name.as_str()).cloned() else {
1546 return msg_err(StatusCode::NOT_FOUND, "resource not found");
1547 };
1548 let Some(member_names) = group.members() else {
1550 return msg_err(StatusCode::NOT_FOUND, "resource not found");
1551 };
1552
1553 let timeout = match parse_delay_params(¶ms) {
1554 Ok(t) => t,
1555 Err(resp) => return *resp,
1556 };
1557
1558 if let Some(selection) = group.selection().filter(|s| s.can_unfix()) {
1562 selection.force_set(None);
1563 }
1564
1565 let url = params.url.as_deref().unwrap_or("").to_string();
1566 let expected = params.expected.clone();
1567
1568 let members: Vec<(String, Arc<dyn meow_common::Proxy>)> = member_names
1571 .into_iter()
1572 .filter_map(|n| route.proxies.get(n.as_str()).cloned().map(|p| (n, p)))
1573 .collect();
1574 drop(route);
1575
1576 let collected = tokio::time::timeout(
1579 timeout,
1580 meow_proxy::health::probe_many_bounded_detailed(
1581 members,
1582 &url,
1583 expected.as_deref(),
1584 timeout,
1585 meow_proxy::health::GROUP_DELAY_CONCURRENCY,
1586 ),
1587 )
1588 .await;
1589
1590 let Ok(pairs) = collected else {
1591 return msg_err(StatusCode::GATEWAY_TIMEOUT, "Timeout");
1594 };
1595
1596 let mut result: BTreeMap<String, u16> = BTreeMap::new();
1597 for pair in pairs {
1598 if matches!(pair.error, Some(meow_proxy::health::UrlTestError::Timeout)) {
1599 return msg_err(StatusCode::GATEWAY_TIMEOUT, "Timeout");
1600 }
1601 result.insert(pair.name, pair.delay);
1602 }
1603 Json(result).into_response()
1604}
1605
1606#[cfg(feature = "listener-tun")]
1616async fn spawn_tun_from_raw(
1617 tunnel: &Tunnel,
1618 raw: &RawConfig,
1619) -> Option<tokio::task::JoinHandle<()>> {
1620 let tun_cfg = match meow_config::parse_tun_config(raw.tun.as_ref()) {
1621 Ok(c) => c,
1622 Err(e) => {
1623 tracing::error!("tun config parse error: {e}");
1624 return None;
1625 }
1626 };
1627 if !tun_cfg.enable {
1628 return None;
1629 }
1630
1631 let (ready_tx, ready_rx) = tokio::sync::oneshot::channel();
1632 let listener = TunListener::new(
1633 tunnel.clone(),
1634 crate::tun_config_to_listener_config(&tun_cfg),
1635 "tun".to_string(),
1636 )
1637 .with_readiness_signal(ready_tx);
1638
1639 let handle = tokio::spawn(async move {
1640 if let Err(e) = listener.run().await {
1641 tracing::error!("TUN listener error: {e}");
1642 }
1643 });
1644
1645 match tokio::time::timeout(crate::TUN_STARTUP_TIMEOUT, ready_rx).await {
1649 Ok(Ok(())) => {}
1650 Ok(Err(_)) => {
1651 tracing::error!(
1652 "TUN listener failed to start (device creation failed — \
1653 check permissions / admin / CAP_NET_ADMIN)"
1654 );
1655 handle.abort();
1656 return None;
1657 }
1658 Err(_) => {
1659 tracing::error!(
1660 "TUN listener startup timed out after {} s",
1661 crate::TUN_STARTUP_TIMEOUT.as_secs()
1662 );
1663 handle.abort();
1664 return None;
1665 }
1666 }
1667
1668 Some(handle)
1669}
1670
1671#[cfg(not(feature = "listener-tun"))]
1672async fn spawn_tun_from_raw(
1673 _tunnel: &Tunnel,
1674 raw: &RawConfig,
1675) -> Option<tokio::task::JoinHandle<()>> {
1676 if raw.tun.as_ref().is_some_and(|t| t.enable) {
1677 tracing::warn!("tun.enable is set but this build lacks the 'listener-tun' feature");
1678 }
1679 None
1680}
1681
1682async fn swap_config_and_reconcile_tun(state: &AppState, candidate: RawConfig) {
1689 let _guard = state.config_mutation_lock.lock().await;
1690
1691 let new_enable = candidate.tun.as_ref().is_some_and(|t| t.enable);
1692 let (old_enable, snapshot) = {
1696 let mut guard = state.raw_config.write();
1697 let old = guard.tun.as_ref().is_some_and(|t| t.enable);
1698 let snapshot = (new_enable && !old).then(|| candidate.clone());
1699 *guard = candidate;
1700 (old, snapshot)
1701 };
1702
1703 if old_enable == new_enable {
1704 return;
1705 }
1706 if let Some(snapshot) = snapshot {
1707 if let Some(handle) = spawn_tun_from_raw(&state.tunnel, &snapshot).await {
1708 state.tunnel.set_tun_handle(handle).await;
1709 info!("TUN listener started via config reload");
1710 }
1711 } else {
1712 state.tunnel.stop_tun().await;
1713 info!("TUN listener stopped via config reload");
1714 }
1715}
1716
1717#[derive(Deserialize)]
1718struct PutConfigsBody {
1719 path: Option<String>,
1720 payload: Option<String>,
1721}
1722
1723async fn put_configs(
1724 State(state): State<Arc<AppState>>,
1725 Query(params): Query<HashMap<String, String>>,
1726 Json(body): Json<PutConfigsBody>,
1727) -> Response {
1728 let force = params.get("force").is_some_and(|v| v == "true");
1729
1730 let yaml =
1731 match (body.path, body.payload) {
1732 (Some(p), _) => match tokio::fs::read_to_string(&p).await {
1733 Ok(s) => s,
1734 Err(e) => {
1735 return (
1736 StatusCode::BAD_REQUEST,
1737 Json(serde_json::json!({"message": e.to_string()})),
1738 )
1739 .into_response()
1740 }
1741 },
1742 (_, Some(b64)) => {
1743 use base64::engine::general_purpose::STANDARD;
1744 use base64::Engine as _;
1745 let Ok(bytes) = STANDARD.decode(&b64) else {
1746 return (
1747 StatusCode::BAD_REQUEST,
1748 Json(serde_json::json!({"message": "payload is not valid base64"})),
1749 )
1750 .into_response();
1751 };
1752 match String::from_utf8(bytes) {
1753 Ok(s) => s,
1754 Err(_) => {
1755 return (
1756 StatusCode::BAD_REQUEST,
1757 Json(serde_json::json!({"message": "payload is not valid UTF-8"})),
1758 )
1759 .into_response()
1760 }
1761 }
1762 }
1763 _ => return (
1764 StatusCode::BAD_REQUEST,
1765 Json(
1766 serde_json::json!({"message": "request body must contain 'path' or 'payload'"}),
1767 ),
1768 )
1769 .into_response(),
1770 };
1771
1772 let mut raw_config: RawConfig = match serde_yaml::from_str(&yaml) {
1774 Ok(c) => c,
1775 Err(e) => {
1776 return (
1777 StatusCode::BAD_REQUEST,
1778 Json(serde_json::json!({"message": format!("config parse error: {e}")})),
1779 )
1780 .into_response()
1781 }
1782 };
1783
1784 if let Some(ps) = raw_config.proxies.as_mut() {
1786 meow_config::ech_dns::preresolve_ech(ps).await;
1787 }
1788
1789 let _mutation = CONFIG_MUTATION.lock().await;
1790
1791 let resolver = Arc::clone(state.tunnel.resolver());
1793 let providers = state
1794 .proxy_providers
1795 .iter()
1796 .map(|entry| (entry.key().clone(), Arc::clone(entry.value())))
1797 .collect();
1798 let (proxies, rules) =
1799 match rebuild_from_raw_with_resolver_async(raw_config.clone(), resolver, providers).await {
1800 Ok(r) => r,
1801 Err(e) => {
1802 if force {
1803 tracing::error!("config reload forced despite validation error: {e}");
1804 (Default::default(), Vec::new())
1805 } else {
1806 return (
1807 StatusCode::BAD_REQUEST,
1808 Json(
1809 serde_json::json!({"message": format!("config validation error: {e}")}),
1810 ),
1811 )
1812 .into_response();
1813 }
1814 }
1815 };
1816
1817 let stats = state.tunnel.statistics();
1819 let dropped = stats.active_connection_count();
1820 stats.close_all_connections();
1821 if dropped > 0 {
1822 tracing::warn!(
1823 connections_dropped = dropped,
1824 "connections force-closed after reload drain timeout"
1825 );
1826 }
1827
1828 state.tunnel.update_proxies(proxies);
1829 state.tunnel.update_rules(rules);
1830 if let Some(mode_str) = &raw_config.mode {
1831 if let Ok(mode) = mode_str.parse::<TunnelMode>() {
1832 state.tunnel.set_mode(mode);
1833 }
1834 }
1835
1836 swap_config_and_reconcile_tun(&state, raw_config).await;
1837
1838 StatusCode::NO_CONTENT.into_response()
1839}
1840
1841async fn get_metrics(State(_state): State<Arc<AppState>>) -> Response {
1845 #[cfg(not(target_has_atomic = "64"))]
1850 {
1851 return (
1852 StatusCode::NOT_IMPLEMENTED,
1853 "metrics require 64-bit atomic support",
1854 )
1855 .into_response();
1856 }
1857
1858 #[cfg(target_has_atomic = "64")]
1859 {
1860 use prometheus_client::encoding::text::encode;
1861 use prometheus_client::metrics::counter::Counter;
1862 use prometheus_client::metrics::family::Family;
1863 use prometheus_client::metrics::gauge::Gauge;
1864 use prometheus_client::registry::Registry;
1865 use std::sync::atomic::{AtomicI64, AtomicU64};
1866
1867 let mut registry = Registry::default();
1868 let stats = _state.tunnel.statistics();
1869 let (upload_total, download_total) = stats.snapshot();
1870
1871 let traffic = Family::<Vec<(String, String)>, Counter<u64, AtomicU64>>::default();
1873 traffic
1874 .get_or_create(&vec![("direction".to_string(), "upload".to_string())])
1875 .inc_by(upload_total.max(0) as u64);
1876 traffic
1877 .get_or_create(&vec![("direction".to_string(), "download".to_string())])
1878 .inc_by(download_total.max(0) as u64);
1879 registry.register(
1880 "meow_traffic_bytes",
1881 "Cumulative bytes transferred since process start",
1882 traffic,
1883 );
1884
1885 let connections_active = Gauge::<i64, AtomicI64>::default();
1887 connections_active.set(stats.active_connection_count() as i64);
1888 registry.register(
1889 "meow_connections_active",
1890 "Number of currently open connections",
1891 connections_active,
1892 );
1893
1894 let proxy_alive = Family::<Vec<(String, String)>, Gauge<i64, AtomicI64>>::default();
1896 let proxy_delay = Family::<Vec<(String, String)>, Gauge<i64, AtomicI64>>::default();
1897 let route = _state.tunnel.route_snapshot();
1898 for (name, proxy) in &route.proxies {
1899 let labels = vec![
1900 ("proxy_name".to_string(), name.to_string()),
1901 ("adapter_type".to_string(), proxy.adapter_type().to_string()),
1902 ];
1903 proxy_alive
1904 .get_or_create(&labels)
1905 .set(if proxy.alive() { 1 } else { 0 });
1906 if !proxy.delay_history().is_empty() {
1909 proxy_delay
1910 .get_or_create(&labels)
1911 .set(proxy.last_delay() as i64);
1912 }
1913 }
1914 registry.register(
1915 "meow_proxy_alive",
1916 "Proxy alive status (1=alive, 0=dead)",
1917 proxy_alive,
1918 );
1919 registry.register(
1920 "meow_proxy_delay_ms",
1921 "Last measured proxy round-trip delay in milliseconds",
1922 proxy_delay,
1923 );
1924
1925 let rules_matched = Family::<Vec<(String, String)>, Counter<u64, AtomicU64>>::default();
1927 for ((rule_type, action), count) in stats.rule_match.snapshot() {
1928 rules_matched
1929 .get_or_create(&vec![
1930 ("rule_type".to_string(), rule_type.to_string()),
1931 ("action".to_string(), action.to_string()),
1932 ])
1933 .inc_by(count);
1934 }
1935 registry.register(
1936 "meow_rules_matched",
1937 "Cumulative rule matches by type and action",
1938 rules_matched,
1939 );
1940
1941 let memory_rss = Gauge::<i64, AtomicI64>::default();
1943 memory_rss.set(read_rss_bytes().await as i64);
1944 registry.register(
1945 "meow_memory_rss_bytes",
1946 "Current process RSS in bytes",
1947 memory_rss,
1948 );
1949
1950 let info = Family::<Vec<(String, String)>, Gauge<i64, AtomicI64>>::default();
1952 info.get_or_create(&vec![
1953 ("version".to_string(), env!("CARGO_PKG_VERSION").to_string()),
1954 ("mode".to_string(), _state.tunnel.mode().to_string()),
1955 ])
1956 .set(1);
1957 registry.register("meow_info", "meow-rs runtime info", info);
1958
1959 let mut body = String::new();
1960 encode(&mut body, ®istry).expect("prometheus text encoding is infallible");
1961 (
1962 StatusCode::OK,
1963 [(
1964 header::CONTENT_TYPE,
1965 "text/plain; version=0.0.4; charset=utf-8",
1966 )],
1967 body,
1968 )
1969 .into_response()
1970 }
1971}
1972
1973#[derive(Deserialize)]
1976struct LogsParams {
1977 level: Option<String>,
1978 format: Option<String>,
1979}
1980
1981fn parse_requested_log_level(
1982 value: Option<&str>,
1983) -> Result<crate::log_stream::LogLevel, Box<Response>> {
1984 let value = value.unwrap_or("info");
1985 match value.to_ascii_lowercase().as_str() {
1986 "debug" | "info" | "warning" | "warn" | "error" | "silent" => Ok(parse_log_level(value)),
1987 _ => Err(Box::new(msg_err(StatusCode::BAD_REQUEST, "Body invalid"))),
1988 }
1989}
1990
1991fn log_json(msg: &LogMessage, structured: bool) -> String {
1992 if !structured {
1993 return serde_json::json!({"type": msg.level.as_str(), "payload": msg.payload}).to_string();
1994 }
1995 let level = if msg.level.as_str() == "warning" {
1996 "warn"
1997 } else {
1998 msg.level.as_str()
1999 };
2000 let t = msg.time.time();
2001 serde_json::json!({
2002 "time": format!("{:02}:{:02}:{:02}", t.hour(), t.minute(), t.second()),
2003 "level": level,
2004 "message": msg.payload,
2005 "fields": [],
2006 })
2007 .to_string()
2008}
2009
2010async fn get_logs(
2012 State(state): State<Arc<AppState>>,
2013 Query(params): Query<LogsParams>,
2014 MaybeWebSocket(ws): MaybeWebSocket,
2015) -> Response {
2016 let level = match parse_requested_log_level(params.level.as_deref()) {
2017 Ok(level) => level,
2018 Err(response) => return *response,
2019 };
2020 let structured = params.format.as_deref() == Some("structured");
2021 let mut rx = state.log_tx.subscribe();
2022 if let Some(ws) = ws {
2023 return ws.on_upgrade(move |mut socket| async move {
2024 loop {
2025 match rx.recv().await {
2026 Ok(msg) if msg.level >= level => {
2027 if socket
2028 .send(Message::Text(log_json(&msg, structured).into()))
2029 .await
2030 .is_err()
2031 {
2032 break;
2033 }
2034 }
2035 Ok(_) | Err(broadcast::error::RecvError::Lagged(_)) => {}
2036 Err(broadcast::error::RecvError::Closed) => break,
2037 }
2038 }
2039 });
2040 }
2041
2042 let stream = futures::stream::unfold(rx, move |mut rx| async move {
2043 loop {
2044 match rx.recv().await {
2045 Ok(msg) if msg.level >= level => {
2046 return Some((
2047 Ok::<String, std::convert::Infallible>(format!(
2048 "{}\n",
2049 log_json(&msg, structured)
2050 )),
2051 rx,
2052 ));
2053 }
2054 Ok(_) | Err(broadcast::error::RecvError::Lagged(_)) => {}
2055 Err(broadcast::error::RecvError::Closed) => return None,
2056 }
2057 }
2058 });
2059 Response::builder()
2060 .header(header::CONTENT_TYPE, "application/json")
2061 .body(Body::from_stream(stream))
2062 .expect("valid log stream response")
2063}
2064
2065static MEMORY_FEED: std::sync::Mutex<Option<broadcast::Sender<Arc<str>>>> =
2076 std::sync::Mutex::new(None);
2077
2078fn subscribe_memory_feed() -> broadcast::Receiver<Arc<str>> {
2079 let mut guard = MEMORY_FEED.lock().expect("memory feed lock poisoned");
2080 if let Some(tx) = guard.as_ref() {
2081 return tx.subscribe();
2083 }
2084 let (tx, rx) = broadcast::channel(2);
2085 *guard = Some(tx.clone());
2086 tokio::spawn(async move {
2087 let mut interval = tokio::time::interval(Duration::from_secs(1));
2088 loop {
2089 interval.tick().await;
2090 if tx.receiver_count() == 0 {
2091 let mut guard = MEMORY_FEED.lock().expect("memory feed lock poisoned");
2095 if tx.receiver_count() == 0 {
2096 *guard = None;
2097 break;
2098 }
2099 }
2100 let inuse = read_rss_bytes().await;
2101 let oslimit = read_os_memory_limit().await;
2102 let msg: Arc<str> = Arc::from(format!("{{\"inuse\":{inuse},\"oslimit\":{oslimit}}}"));
2103 let _ = tx.send(msg);
2104 }
2105 });
2106 rx
2107}
2108
2109async fn get_memory(
2110 State(_state): State<Arc<AppState>>,
2111 MaybeWebSocket(ws): MaybeWebSocket,
2112) -> Response {
2113 let first: Arc<str> = Arc::from("{\"inuse\":0,\"oslimit\":0}");
2114 if let Some(ws) = ws {
2115 return ws.on_upgrade(move |mut socket| async move {
2116 if socket
2117 .send(Message::Text(first.as_ref().into()))
2118 .await
2119 .is_err()
2120 {
2121 return;
2122 }
2123 let mut feed = subscribe_memory_feed();
2124 loop {
2125 let msg = match feed.recv().await {
2126 Ok(msg) => msg,
2127 Err(broadcast::error::RecvError::Lagged(_)) => continue,
2128 Err(broadcast::error::RecvError::Closed) => break,
2129 };
2130 if socket
2131 .send(Message::Text(msg.as_ref().into()))
2132 .await
2133 .is_err()
2134 {
2135 break;
2136 }
2137 }
2138 });
2139 }
2140
2141 let feed = subscribe_memory_feed();
2142 let stream = futures::stream::unfold((Some(first), feed), |(first, mut feed)| async move {
2143 if let Some(first) = first {
2144 return Some((
2145 Ok::<String, std::convert::Infallible>(format!("{first}\n")),
2146 (None, feed),
2147 ));
2148 }
2149 loop {
2150 match feed.recv().await {
2151 Ok(msg) => {
2152 return Some((
2153 Ok::<String, std::convert::Infallible>(format!("{msg}\n")),
2154 (None, feed),
2155 ));
2156 }
2157 Err(broadcast::error::RecvError::Lagged(_)) => continue,
2158 Err(broadcast::error::RecvError::Closed) => return None,
2159 }
2160 }
2161 });
2162 Response::builder()
2163 .header(header::CONTENT_TYPE, "application/json")
2164 .body(Body::from_stream(stream))
2165 .expect("valid memory stream response")
2166}
2167
2168async fn read_rss_bytes() -> u64 {
2169 tokio::task::spawn_blocking(|| {
2170 use sysinfo::{Pid, ProcessesToUpdate, System};
2171 let pid = Pid::from_u32(std::process::id());
2172 let mut sys = System::new();
2173 sys.refresh_processes(ProcessesToUpdate::Some(&[pid]), false);
2174 sys.process(pid).map_or(0, sysinfo::Process::memory)
2175 })
2176 .await
2177 .unwrap_or(0)
2178}
2179
2180async fn read_os_memory_limit() -> u64 {
2181 #[cfg(target_os = "linux")]
2182 {
2183 read_os_memory_limit_linux().await
2184 }
2185 #[cfg(not(target_os = "linux"))]
2186 {
2187 0
2188 }
2189}
2190
2191#[cfg(target_os = "linux")]
2192async fn read_os_memory_limit_linux() -> u64 {
2193 if let Ok(s) = tokio::fs::read_to_string("/sys/fs/cgroup/memory.max").await {
2195 if let Ok(n) = s.trim().parse::<u64>() {
2196 return n;
2197 }
2198 }
2199 unsafe {
2201 let mut rl = libc::rlimit {
2202 rlim_cur: 0,
2203 rlim_max: 0,
2204 };
2205 if libc::getrlimit(libc::RLIMIT_AS, &mut rl) == 0 && rl.rlim_cur != libc::RLIM_INFINITY {
2206 #[cfg(target_pointer_width = "32")]
2207 {
2208 return rl.rlim_cur as u64;
2209 }
2210 #[cfg(not(target_pointer_width = "32"))]
2211 {
2212 return rl.rlim_cur;
2213 }
2214 }
2215 }
2216 0
2217}
2218
2219#[derive(Serialize)]
2222#[serde(rename_all = "camelCase")]
2223struct ProviderInfo {
2224 name: String,
2225 #[serde(rename = "type")]
2226 provider_type: String,
2227 vehicle_type: String,
2228 proxies: Vec<ProxyInfo>,
2229 #[serde(rename = "testUrl")]
2230 test_url: String,
2231 #[serde(rename = "expectedStatus")]
2232 expected_status: String,
2233 #[serde(rename = "updatedAt", skip_serializing_if = "Option::is_none")]
2234 updated_at: Option<String>,
2235}
2236
2237fn unix_rfc3339(seconds: u64) -> Option<String> {
2238 use time::format_description::well_known::Rfc3339;
2239 (seconds > 0)
2240 .then(|| time::OffsetDateTime::from_unix_timestamp(seconds as i64).ok())
2241 .flatten()
2242 .and_then(|time| time.format(&Rfc3339).ok())
2243}
2244
2245fn provider_to_info(name: &str, provider: &ProxyProvider) -> ProviderInfo {
2246 let proxies = provider
2247 .proxies()
2248 .iter()
2249 .map(ProxyInfo::from_proxy)
2250 .collect();
2251 ProviderInfo {
2252 name: name.to_string(),
2253 provider_type: "Proxy".to_string(),
2254 vehicle_type: provider.vehicle_type.to_string(),
2255 proxies,
2256 test_url: provider
2257 .health_check
2258 .as_ref()
2259 .map_or_else(String::new, |hc| hc.url.clone()),
2260 expected_status: provider
2261 .health_check
2262 .as_ref()
2263 .map_or_else(String::new, |hc| hc.expected_status.clone()),
2264 updated_at: unix_rfc3339(provider.updated_at_secs()),
2265 }
2266}
2267
2268async fn get_providers(State(state): State<Arc<AppState>>) -> Json<serde_json::Value> {
2269 let mut map = serde_json::Map::new();
2270 for entry in state.proxy_providers.iter() {
2271 let info = provider_to_info(entry.key(), entry.value());
2272 map.insert(
2273 entry.key().clone(),
2274 serde_json::to_value(info).unwrap_or_default(),
2275 );
2276 }
2277 Json(serde_json::json!({ "providers": map }))
2278}
2279
2280async fn get_provider(State(state): State<Arc<AppState>>, Path(name): Path<String>) -> Response {
2281 match state.proxy_providers.get(&name) {
2282 Some(entry) => Json(provider_to_info(&name, entry.value())).into_response(),
2283 None => msg_err(StatusCode::NOT_FOUND, "resource not found"),
2284 }
2285}
2286
2287async fn refresh_provider(
2288 State(state): State<Arc<AppState>>,
2289 Path(name): Path<String>,
2290) -> Response {
2291 let provider = match state.proxy_providers.get(&name) {
2292 Some(entry) => Arc::clone(entry.value()),
2293 None => return msg_err(StatusCode::NOT_FOUND, "resource not found"),
2294 };
2295 match provider.refresh().await {
2296 Ok(()) => StatusCode::NO_CONTENT.into_response(),
2297 Err(e) => (
2298 StatusCode::SERVICE_UNAVAILABLE,
2299 Json(serde_json::json!({"message": e})),
2300 )
2301 .into_response(),
2302 }
2303}
2304
2305async fn provider_healthcheck(
2308 State(state): State<Arc<AppState>>,
2309 Path(name): Path<String>,
2310) -> Response {
2311 let provider = match state.proxy_providers.get(&name) {
2312 Some(entry) => Arc::clone(entry.value()),
2313 None => return msg_err(StatusCode::NOT_FOUND, "resource not found"),
2314 };
2315
2316 let Some(health) = provider.health_check.as_ref() else {
2317 return StatusCode::NO_CONTENT.into_response();
2318 };
2319 let timeout = Duration::from_millis(health.timeout.max(1));
2320 let url = health.url.clone();
2321 let expected = (!health.expected_status.is_empty()).then(|| health.expected_status.clone());
2322
2323 let members = provider
2324 .proxies()
2325 .into_iter()
2326 .map(|proxy| (proxy.name().to_string(), proxy))
2327 .collect();
2328
2329 let _ = meow_proxy::health::probe_many_bounded(
2330 members,
2331 &url,
2332 expected.as_deref(),
2333 timeout,
2334 meow_proxy::health::PROVIDER_HEALTHCHECK_CONCURRENCY,
2335 )
2336 .await;
2337
2338 StatusCode::NO_CONTENT.into_response()
2339}
2340
2341async fn get_provider_proxy(
2342 State(state): State<Arc<AppState>>,
2343 Path((provider_name, proxy_name)): Path<(String, String)>,
2344) -> Response {
2345 let Some(provider) = state.proxy_providers.get(&provider_name) else {
2346 return msg_err(StatusCode::NOT_FOUND, "Resource not found");
2347 };
2348 match provider
2349 .proxies()
2350 .into_iter()
2351 .find(|p| p.name() == proxy_name)
2352 {
2353 Some(proxy) => Json(ProxyInfo::from_proxy(&proxy)).into_response(),
2354 None => msg_err(StatusCode::NOT_FOUND, "Resource not found"),
2355 }
2356}
2357
2358async fn provider_proxy_healthcheck(
2359 State(state): State<Arc<AppState>>,
2360 Path((provider_name, proxy_name)): Path<(String, String)>,
2361 Query(params): Query<DelayParams>,
2362) -> Response {
2363 let timeout = match parse_delay_params(¶ms) {
2364 Ok(timeout) => timeout,
2365 Err(response) => return *response,
2366 };
2367 let Some(provider) = state.proxy_providers.get(&provider_name) else {
2368 return msg_err(StatusCode::NOT_FOUND, "Resource not found");
2369 };
2370 let Some(proxy) = provider
2371 .proxies()
2372 .into_iter()
2373 .find(|p| p.name() == proxy_name)
2374 else {
2375 return msg_err(StatusCode::NOT_FOUND, "Resource not found");
2376 };
2377 match probe_and_record(
2378 &proxy,
2379 params.url.as_deref().unwrap_or(""),
2380 params.expected.as_deref(),
2381 timeout,
2382 )
2383 .await
2384 {
2385 Ok(delay) => Json(DelayResp { delay }).into_response(),
2386 Err(meow_proxy::health::UrlTestError::Timeout) => {
2387 msg_err(StatusCode::GATEWAY_TIMEOUT, "Timeout")
2388 }
2389 Err(meow_proxy::health::UrlTestError::Transport(_)) => msg_err(
2390 StatusCode::SERVICE_UNAVAILABLE,
2391 "An error occurred in the delay test",
2392 ),
2393 }
2394}
2395
2396#[derive(Serialize)]
2399struct RuleProviderInfo {
2400 name: String,
2401 #[serde(rename = "type")]
2402 provider_type: String,
2403 behavior: String,
2404 format: String,
2405 #[serde(rename = "ruleCount")]
2406 rule_count: usize,
2407 #[serde(rename = "updatedAt")]
2408 updated_at: String,
2409 #[serde(rename = "vehicleType")]
2410 vehicle_type: String,
2411}
2412
2413impl RuleProviderInfo {
2414 fn from_provider(p: &Arc<RuleProvider>, format: Option<&str>) -> Self {
2415 let vehicle_type = match p.provider_type {
2416 meow_config::rule_provider::ProviderType::Http => "HTTP",
2417 meow_config::rule_provider::ProviderType::File => "File",
2418 meow_config::rule_provider::ProviderType::Inline => "Inline",
2419 };
2420 Self {
2421 name: p.name.clone(),
2422 provider_type: "Rule".to_string(),
2423 behavior: p.behavior.to_string(),
2424 format: format.unwrap_or("yaml").to_string(),
2425 rule_count: p.rule_count(),
2426 updated_at: unix_rfc3339(p.updated_at_secs()).unwrap_or_default(),
2427 vehicle_type: vehicle_type.to_string(),
2428 }
2429 }
2430}
2431
2432#[derive(Serialize)]
2433struct RuleProvidersResponse {
2434 providers: HashMap<String, RuleProviderInfo>,
2435}
2436
2437async fn get_rule_providers(State(state): State<Arc<AppState>>) -> Json<RuleProvidersResponse> {
2438 let providers = state.rule_providers.read();
2439 let raw = state.raw_config.read();
2440 let map: HashMap<String, RuleProviderInfo> = providers
2441 .iter()
2442 .map(|(name, p): (&String, &Arc<RuleProvider>)| {
2443 let format = raw
2444 .rule_providers
2445 .as_ref()
2446 .and_then(|all| all.get(name))
2447 .and_then(|provider| provider.format.as_deref());
2448 (name.clone(), RuleProviderInfo::from_provider(p, format))
2449 })
2450 .collect();
2451 Json(RuleProvidersResponse { providers: map })
2452}
2453
2454async fn get_rule_provider(
2455 State(state): State<Arc<AppState>>,
2456 Path(name): Path<String>,
2457) -> Result<Json<RuleProviderInfo>, StatusCode> {
2458 let providers = state.rule_providers.read();
2459 let p = providers.get(&name).ok_or(StatusCode::NOT_FOUND)?;
2460 let raw = state.raw_config.read();
2461 let format = raw
2462 .rule_providers
2463 .as_ref()
2464 .and_then(|all| all.get(&name))
2465 .and_then(|provider| provider.format.as_deref());
2466 Ok(Json(RuleProviderInfo::from_provider(p, format)))
2467}
2468
2469async fn refresh_rule_provider(
2470 State(state): State<Arc<AppState>>,
2471 Path(name): Path<String>,
2472) -> StatusCode {
2473 let provider = {
2474 let providers = state.rule_providers.read();
2475 providers.get(&name).cloned()
2476 };
2477 let Some(p) = provider else {
2478 return StatusCode::NOT_FOUND;
2479 };
2480 let ctx = meow_rules::ParserContext::empty();
2481 match p.refresh(&ctx).await {
2482 Ok(()) => StatusCode::NO_CONTENT,
2483 Err(e) => {
2484 tracing::warn!(provider = %name, "rule-provider refresh failed: {:#}", e);
2485 StatusCode::SERVICE_UNAVAILABLE
2486 }
2487 }
2488}
2489
2490async fn get_listeners(State(state): State<Arc<AppState>>) -> Json<serde_json::Value> {
2493 let items: Vec<serde_json::Value> = state
2494 .listeners
2495 .iter()
2496 .map(|l| {
2497 serde_json::json!({
2498 "name": l.name,
2499 "type": l.listener_type.to_string(),
2500 "port": l.port,
2501 "listen": l.listen,
2502 })
2503 })
2504 .collect();
2505 Json(serde_json::json!(items))
2506}
2507
2508#[cfg(test)]
2509mod tests {
2510 use super::*;
2511
2512 #[test]
2513 fn connections_interval_rejects_zero_and_garbage() {
2514 assert_eq!(parse_connections_interval(Some("0")), None);
2517 assert_eq!(parse_connections_interval(Some("abc")), None);
2518 assert_eq!(parse_connections_interval(Some("-1")), None);
2519 assert_eq!(parse_connections_interval(Some("")), None);
2520 }
2521
2522 #[test]
2523 fn connections_interval_clamps_to_floor() {
2524 assert_eq!(parse_connections_interval(Some("1")), Some(100));
2525 assert_eq!(parse_connections_interval(Some("99")), Some(100));
2526 assert_eq!(parse_connections_interval(Some("100")), Some(100));
2527 }
2528
2529 #[test]
2530 fn connections_interval_passes_through_above_floor() {
2531 assert_eq!(parse_connections_interval(Some("101")), Some(101));
2532 assert_eq!(parse_connections_interval(Some("5000")), Some(5000));
2533 assert_eq!(parse_connections_interval(None), Some(1000));
2534 }
2535}