recycleactor commited on
Commit
6fcdbdf
·
verified ·
1 Parent(s): db60e8d

Upload 2 files

Browse files
Files changed (1) hide show
  1. src/main.rs +38 -11
src/main.rs CHANGED
@@ -30,7 +30,7 @@ type BoxErr = Box<dyn std::error::Error + Send + Sync + 'static>;
30
  mod api;
31
  use api::{APIClient, APIError, Format};
32
 
33
- #[derive(Clone)]
34
  struct ArlSession {
35
  check_form: String,
36
  license_token: String,
@@ -43,6 +43,7 @@ struct AppState {
43
  arl_sessions: Arc<RwLock<HashMap<String, ArlSession>>>,
44
  pair: Arc<RwLock<HashMap<String, PairSession>>>,
45
  pair_store_path: String,
 
46
  }
47
 
48
  fn license_token_from_ud(ud: &serde_json::Value) -> String {
@@ -68,6 +69,32 @@ async fn store_arl_session(state: &AppState, arl: &str, check_form: &str, licens
68
  license_token: license_token.to_string(),
69
  },
70
  );
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
71
  }
72
 
73
  fn json_error(status: StatusCode, message: impl Into<String>) -> axum::response::Response {
@@ -592,15 +619,12 @@ async fn stream(
592
  let mut url: Option<String> = None;
593
  let mut last_err: Option<String> = None;
594
 
595
- // Если ARL не передан — пробуем upstream (может быть отключён)
596
  if arl_key.is_none() {
597
- match upstream_media_url(q.id, &formats).await {
598
- Ok(u) => url = Some(u),
599
- Err(e) => last_err = Some(e),
600
- }
601
  }
602
 
603
- // До 3 попыток: 0 — renew+запрос, 1 — сброс кэша клиента, 2 — ещё раз renew
604
  for attempt in 0..3u8 {
605
  if url.is_some() { break; }
606
 
@@ -632,10 +656,11 @@ async fn stream(
632
 
633
  let url = match url {
634
  Some(u) => u,
635
- None => return (
636
- StatusCode::SERVICE_UNAVAILABLE,
637
- last_err.unwrap_or_else(|| "stream:no_url".to_string()),
638
- ).into_response(),
 
639
  };
640
 
641
  fn parse_range_header(v: &str) -> Option<(u64, Option<u64>)> {
@@ -2837,7 +2862,9 @@ async fn main() {
2837
  arl_sessions: Arc::new(RwLock::new(HashMap::new())),
2838
  pair: Arc::new(RwLock::new(HashMap::new())),
2839
  pair_store_path: std::env::var("PAIR_STORE_PATH").ok().unwrap_or_else(|| "pair_store.json".to_string()),
 
2840
  };
 
2841
  pair_store_load(&state).await;
2842
 
2843
  let cors = CorsLayer::new()
 
30
  mod api;
31
  use api::{APIClient, APIError, Format};
32
 
33
+ #[derive(Clone, Serialize, Deserialize)]
34
  struct ArlSession {
35
  check_form: String,
36
  license_token: String,
 
43
  arl_sessions: Arc<RwLock<HashMap<String, ArlSession>>>,
44
  pair: Arc<RwLock<HashMap<String, PairSession>>>,
45
  pair_store_path: String,
46
+ arl_store_path: String,
47
  }
48
 
49
  fn license_token_from_ud(ud: &serde_json::Value) -> String {
 
69
  license_token: license_token.to_string(),
70
  },
71
  );
72
+ arl_store_save(state).await;
73
+ }
74
+
75
+ async fn arl_store_save(state: &AppState) {
76
+ let path = state.arl_store_path.trim();
77
+ if path.is_empty() {
78
+ return;
79
+ }
80
+ let map = state.arl_sessions.read().await.clone();
81
+ let Ok(txt) = serde_json::to_string(&map) else { return };
82
+ let tmp = format!("{path}.tmp");
83
+ if tokio::fs::write(&tmp, txt).await.is_ok() {
84
+ let _ = tokio::fs::rename(&tmp, path).await;
85
+ }
86
+ }
87
+
88
+ async fn arl_store_load(state: &AppState) {
89
+ let path = state.arl_store_path.trim();
90
+ if path.is_empty() {
91
+ return;
92
+ }
93
+ if let Ok(txt) = tokio::fs::read_to_string(path).await {
94
+ if let Ok(map) = serde_json::from_str::<HashMap<String, ArlSession>>(&txt) {
95
+ *state.arl_sessions.write().await = map;
96
+ }
97
+ }
98
  }
99
 
100
  fn json_error(status: StatusCode, message: impl Into<String>) -> axum::response::Response {
 
619
  let mut url: Option<String> = None;
620
  let mut last_err: Option<String> = None;
621
 
 
622
  if arl_key.is_none() {
623
+ tracing::error!("stream {}: arl missing", q.id);
624
+ return (StatusCode::SERVICE_UNAVAILABLE, "arl required").into_response();
 
 
625
  }
626
 
627
+ // До 3 попыток
628
  for attempt in 0..3u8 {
629
  if url.is_some() { break; }
630
 
 
656
 
657
  let url = match url {
658
  Some(u) => u,
659
+ None => {
660
+ let err = last_err.unwrap_or_else(|| "stream:no_url".to_string());
661
+ tracing::error!("stream {} failed: {}", q.id, err);
662
+ return (StatusCode::SERVICE_UNAVAILABLE, err).into_response();
663
+ }
664
  };
665
 
666
  fn parse_range_header(v: &str) -> Option<(u64, Option<u64>)> {
 
2862
  arl_sessions: Arc::new(RwLock::new(HashMap::new())),
2863
  pair: Arc::new(RwLock::new(HashMap::new())),
2864
  pair_store_path: std::env::var("PAIR_STORE_PATH").ok().unwrap_or_else(|| "pair_store.json".to_string()),
2865
+ arl_store_path: std::env::var("ARL_STORE_PATH").ok().unwrap_or_else(|| "arl_store.json".to_string()),
2866
  };
2867
+ arl_store_load(&state).await;
2868
  pair_store_load(&state).await;
2869
 
2870
  let cors = CorsLayer::new()