recycleactor commited on
Commit
ceb1161
·
verified ·
1 Parent(s): 6326736

Upload 2 files

Browse files
Files changed (2) hide show
  1. src/api.rs +27 -3
  2. src/main.rs +86 -112
src/api.rs CHANGED
@@ -174,14 +174,38 @@ impl APIClient {
174
  self.no_renew_api_call(method, params).await
175
  }
176
 
 
 
 
 
 
 
177
  async fn renew(&mut self) -> Result<(), APIError> {
178
  let user_data: serde_json::Value = self.user_data().await?;
179
 
180
- self.check_form = user_data["checkForm"].as_str().unwrap().to_string();
181
- self.license_token = user_data["USER"]["OPTIONS"]["license_token"].as_str().unwrap().to_string();
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
182
 
 
 
183
  self.renew_instant = Some(Instant::now());
184
-
185
  Ok(())
186
  }
187
 
 
174
  self.no_renew_api_call(method, params).await
175
  }
176
 
177
+ /// Принудительное обновление токенов — вызывается при ошибках стрима
178
+ pub async fn force_renew(&mut self) -> Result<(), APIError> {
179
+ self.renew_instant = None; // сбрасываем таймер
180
+ self.renew().await
181
+ }
182
+
183
  async fn renew(&mut self) -> Result<(), APIError> {
184
  let user_data: serde_json::Value = self.user_data().await?;
185
 
186
+ let check_form = user_data["checkForm"].as_str().unwrap_or("").to_string();
187
+ let license_token = user_data
188
+ .pointer("/USER/OPTIONS/license_token")
189
+ .and_then(|v| v.as_str())
190
+ .unwrap_or("")
191
+ .to_string();
192
+
193
+ if check_form.is_empty() {
194
+ return Err(APIError::DeezerError {
195
+ code: "renew".to_string(),
196
+ message: "checkForm missing in getUserData response".to_string(),
197
+ });
198
+ }
199
+ if license_token.is_empty() {
200
+ return Err(APIError::DeezerError {
201
+ code: "renew".to_string(),
202
+ message: "license_token missing in getUserData response".to_string(),
203
+ });
204
+ }
205
 
206
+ self.check_form = check_form;
207
+ self.license_token = license_token;
208
  self.renew_instant = Some(Instant::now());
 
209
  Ok(())
210
  }
211
 
src/main.rs CHANGED
@@ -176,10 +176,9 @@ async fn media_url_for_track(
176
  fn upstream_dzmedia_base() -> Option<String> {
177
  let env = std::env::var("DZMEDIA_UPSTREAM").ok().unwrap_or_default();
178
  let base = env.trim();
179
- if base.is_empty() {
180
- return Some("https://lufts-dzmedia.fly.dev".to_string());
181
- }
182
- if base == "-" {
183
  return None;
184
  }
185
  Some(base.trim_end_matches('/').to_string())
@@ -389,63 +388,58 @@ struct StreamParams {
389
 
390
  async fn stream_info(State(state): State<AppState>, Query(q): Query<StreamParams>) -> impl IntoResponse {
391
  let formats = match q.format.as_deref().unwrap_or("AUTO") {
392
- "FLAC" => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
393
- "MP3_320" => vec![Format::MP3_320, Format::MP3_128],
394
  "MP3_MISC" => vec![Format::MP3_MISC, Format::MP3_128],
395
- "MP3_128" => vec![Format::MP3_128],
396
- "AUTO" | _ => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
397
  };
398
 
399
  let requested = q.format.clone().unwrap_or_else(|| "AUTO".to_string());
400
- let mut url: String = String::new();
401
  let mut used: String = String::new();
402
  let mut last_err: Option<String> = None;
403
 
 
404
  if let Ok(text) = upstream_get_url_text(&formats, &vec![q.id]).await {
405
  if let Ok(v) = serde_json::from_str::<serde_json::Value>(&text) {
406
- url = extract_media_url(&v);
407
  used = extract_media_format(&v);
408
  }
409
  }
410
 
 
411
  if url.is_empty() {
412
  if let Some(arl) = q.arl.clone().filter(|s| !s.trim().is_empty()) {
413
- let client = api_client_for_arl(&state, Some(arl)).await;
414
- let res = async {
415
- let mut client = client.lock().await;
416
- let resp: Result<DeezerTrackList, APIError> = client
417
- .api_call(
418
- "song.getListData",
419
- &json!({"sng_ids":[q.id],"array_default":["SNG_ID","TRACK_TOKEN","DURATION"]}),
420
- )
421
- .await;
422
- let track_list = resp.map_err(|e| e.to_string())?;
423
- if track_list.data.is_empty() {
424
- return Err("No valid ID".to_string());
425
- }
426
- let track_token = track_list.data[0].TRACK_TOKEN.as_str();
427
- let media_resp = client
428
- .get_media(&formats, vec![track_token])
429
- .await
430
- .map_err(|e| e.to_string())?;
431
- let media_json: serde_json::Value = media_resp
432
- .json()
433
- .await
434
- .map_err(|_| "Bad media response".to_string())?;
435
- let url = extract_media_url(&media_json);
436
- let used = extract_media_format(&media_json);
437
- if url.is_empty() {
438
- return Err("No url in media response".to_string());
439
  }
440
- Ok::<(String, String), String>((url, used))
441
- }
442
- .await;
443
- match res {
444
- Ok((u, f)) => {
445
- url = u;
446
- used = f;
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
447
  }
448
- Err(e) => last_err = Some(e),
449
  }
450
  }
451
  }
@@ -475,37 +469,31 @@ async fn stream_info(State(state): State<AppState>, Query(q): Query<StreamParams
475
  if let Some(ct) = r.headers().get(reqwest::header::CONTENT_TYPE).and_then(|v| v.to_str().ok()) {
476
  mime = ct.to_string();
477
  }
478
- if let Some(cr) = r
479
- .headers()
480
- .get(reqwest::header::CONTENT_RANGE)
481
  .and_then(|v| v.to_str().ok())
482
  .and_then(total_from_content_range)
483
  {
484
  total_bytes = cr;
485
- } else if let Some(cl) = r.headers().get(reqwest::header::CONTENT_LENGTH).and_then(|v| v.to_str().ok()) {
 
 
486
  total_bytes = cl.parse::<u64>().unwrap_or(0);
487
  }
488
  }
489
 
490
  let bitrate_kbps = if duration > 0 && total_bytes > 0 {
491
  ((total_bytes as f64) * 8.0 / (duration as f64) / 1000.0).round() as u64
492
- } else {
493
- 0
494
- };
495
 
496
- (
497
- StatusCode::OK,
498
- Json(json!({
499
- "id": q.id,
500
- "requested": requested,
501
- "used": used,
502
- "duration": duration,
503
- "bytes": total_bytes,
504
- "bitrate_kbps": bitrate_kbps,
505
- "mime": mime
506
- })),
507
- )
508
- .into_response()
509
  }
510
 
511
  async fn stream(
@@ -514,80 +502,66 @@ async fn stream(
514
  headers: HeaderMap,
515
  ) -> impl IntoResponse {
516
  let formats = match q.format.as_deref().unwrap_or("AUTO") {
517
- "FLAC" => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
518
- "MP3_320" => vec![Format::MP3_320, Format::MP3_128],
519
  "MP3_MISC" => vec![Format::MP3_MISC, Format::MP3_128],
520
- "MP3_128" => vec![Format::MP3_128],
521
- "AUTO" | _ => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
522
  };
523
 
524
- let arl_key = q
525
- .arl
526
- .clone()
527
  .map(|s| s.trim().to_string())
528
  .filter(|s| !s.is_empty());
529
 
530
  let mut url: Option<String> = None;
531
  let mut last_err: Option<String> = None;
532
 
 
533
  if arl_key.is_none() {
534
  match upstream_media_url(q.id, &formats).await {
535
- Ok(u) => url = Some(u),
536
  Err(e) => last_err = Some(e),
537
  }
538
  }
539
 
540
- for attempt in 0..2 {
541
- if url.is_some() {
542
- break;
 
 
 
 
 
 
 
 
543
  }
 
544
  let client = api_client_for_arl(&state, arl_key.clone()).await;
545
  let res = {
546
- let mut client = client.lock().await;
547
- media_url_for_track(&mut client, q.id, &formats).await
 
 
 
 
 
 
 
548
  };
549
 
550
  match res {
551
- Ok(u) => {
552
- url = Some(u);
553
- break;
554
- }
555
- Err(e) => {
556
- last_err = Some(e);
557
- if attempt == 0 {
558
- if let Some(arl) = arl_key.clone() {
559
- let mut map = state.api_by_arl.write().await;
560
- map.remove(&arl);
561
- } else {
562
- let mut global = state.api.lock().await;
563
- *global = APIClient::new();
564
- }
565
- }
566
- }
567
  }
568
  }
569
 
570
  let url = match url {
571
  Some(u) => u,
572
- None => {
573
- if arl_key.is_some() {
574
- if let Ok(u) = upstream_media_url(q.id, &formats).await {
575
- u
576
- } else {
577
- return (
578
- StatusCode::SERVICE_UNAVAILABLE,
579
- last_err.unwrap_or_else(|| "stream:error".to_string()),
580
- )
581
- .into_response();
582
- }
583
- } else {
584
- return (
585
- StatusCode::SERVICE_UNAVAILABLE,
586
- last_err.unwrap_or_else(|| "stream:error".to_string()),
587
- )
588
- .into_response();
589
- }
590
- }
591
  };
592
 
593
  fn parse_range_header(v: &str) -> Option<(u64, Option<u64>)> {
 
176
  fn upstream_dzmedia_base() -> Option<String> {
177
  let env = std::env::var("DZMEDIA_UPSTREAM").ok().unwrap_or_default();
178
  let base = env.trim();
179
+ // Если переменная не задана — upstream отключён (используем только ARL)
180
+ // Чтобы включить upstream — задайте DZMEDIA_UPSTREAM=https://your-instance.fly.dev
181
+ if base.is_empty() || base == "-" {
 
182
  return None;
183
  }
184
  Some(base.trim_end_matches('/').to_string())
 
388
 
389
  async fn stream_info(State(state): State<AppState>, Query(q): Query<StreamParams>) -> impl IntoResponse {
390
  let formats = match q.format.as_deref().unwrap_or("AUTO") {
391
+ "FLAC" => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
392
+ "MP3_320" => vec![Format::MP3_320, Format::MP3_128],
393
  "MP3_MISC" => vec![Format::MP3_MISC, Format::MP3_128],
394
+ "MP3_128" => vec![Format::MP3_128],
395
+ _ => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
396
  };
397
 
398
  let requested = q.format.clone().unwrap_or_else(|| "AUTO".to_string());
399
+ let mut url: String = String::new();
400
  let mut used: String = String::new();
401
  let mut last_err: Option<String> = None;
402
 
403
+ // Upstream (опционально — если DZMEDIA_UPSTREAM задан)
404
  if let Ok(text) = upstream_get_url_text(&formats, &vec![q.id]).await {
405
  if let Ok(v) = serde_json::from_str::<serde_json::Value>(&text) {
406
+ url = extract_media_url(&v);
407
  used = extract_media_format(&v);
408
  }
409
  }
410
 
411
+ // ARL путь — до 2 попыток
412
  if url.is_empty() {
413
  if let Some(arl) = q.arl.clone().filter(|s| !s.trim().is_empty()) {
414
+ 'arl: for attempt in 0..2u8 {
415
+ if attempt == 1 {
416
+ state.api_by_arl.write().await.remove(&arl);
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
417
  }
418
+ let client = api_client_for_arl(&state, Some(arl.clone())).await;
419
+ let res = async {
420
+ let mut c = client.lock().await;
421
+ if attempt == 1 { let _ = c.force_renew().await; }
422
+ let resp: Result<DeezerTrackList, APIError> = c
423
+ .api_call(
424
+ "song.getListData",
425
+ &json!({"sng_ids":[q.id],"array_default":["SNG_ID","TRACK_TOKEN","DURATION"]}),
426
+ )
427
+ .await;
428
+ let track_list = resp.map_err(|e| e.to_string())?;
429
+ if track_list.data.is_empty() { return Err("No valid ID".to_string()); }
430
+ let token = track_list.data[0].TRACK_TOKEN.as_str();
431
+ let media_resp = c.get_media(&formats, vec![token]).await.map_err(|e| e.to_string())?;
432
+ let media_json: serde_json::Value = media_resp.json().await
433
+ .map_err(|_| "Bad media response".to_string())?;
434
+ let u = extract_media_url(&media_json);
435
+ let f = extract_media_format(&media_json);
436
+ if u.is_empty() { return Err("No url in media response".to_string()); }
437
+ Ok::<(String, String), String>((u, f))
438
+ }.await;
439
+ match res {
440
+ Ok((u, f)) => { url = u; used = f; break 'arl; }
441
+ Err(e) => { last_err = Some(e); }
442
  }
 
443
  }
444
  }
445
  }
 
469
  if let Some(ct) = r.headers().get(reqwest::header::CONTENT_TYPE).and_then(|v| v.to_str().ok()) {
470
  mime = ct.to_string();
471
  }
472
+ if let Some(cr) = r.headers().get(reqwest::header::CONTENT_RANGE)
 
 
473
  .and_then(|v| v.to_str().ok())
474
  .and_then(total_from_content_range)
475
  {
476
  total_bytes = cr;
477
+ } else if let Some(cl) = r.headers().get(reqwest::header::CONTENT_LENGTH)
478
+ .and_then(|v| v.to_str().ok())
479
+ {
480
  total_bytes = cl.parse::<u64>().unwrap_or(0);
481
  }
482
  }
483
 
484
  let bitrate_kbps = if duration > 0 && total_bytes > 0 {
485
  ((total_bytes as f64) * 8.0 / (duration as f64) / 1000.0).round() as u64
486
+ } else { 0 };
 
 
487
 
488
+ (StatusCode::OK, Json(json!({
489
+ "id": q.id,
490
+ "requested": requested,
491
+ "used": used,
492
+ "duration": duration,
493
+ "bytes": total_bytes,
494
+ "bitrate_kbps": bitrate_kbps,
495
+ "mime": mime
496
+ }))).into_response()
 
 
 
 
497
  }
498
 
499
  async fn stream(
 
502
  headers: HeaderMap,
503
  ) -> impl IntoResponse {
504
  let formats = match q.format.as_deref().unwrap_or("AUTO") {
505
+ "FLAC" => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
506
+ "MP3_320" => vec![Format::MP3_320, Format::MP3_128],
507
  "MP3_MISC" => vec![Format::MP3_MISC, Format::MP3_128],
508
+ "MP3_128" => vec![Format::MP3_128],
509
+ _ => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
510
  };
511
 
512
+ let arl_key = q.arl.clone()
 
 
513
  .map(|s| s.trim().to_string())
514
  .filter(|s| !s.is_empty());
515
 
516
  let mut url: Option<String> = None;
517
  let mut last_err: Option<String> = None;
518
 
519
+ // Если ARL не передан — пробуем upstream (может быть отключён)
520
  if arl_key.is_none() {
521
  match upstream_media_url(q.id, &formats).await {
522
+ Ok(u) => url = Some(u),
523
  Err(e) => last_err = Some(e),
524
  }
525
  }
526
 
527
+ // До 3 попыток: 0 — обычная, 1 — после сброса кэша клиента, 2 — принудительный renew внутри
528
+ for attempt in 0..3u8 {
529
+ if url.is_some() { break; }
530
+
531
+ // На попытке 1 — выбрасываем кэшированный клиент чтобы пересоздать его
532
+ if attempt == 1 {
533
+ if let Some(ref arl) = arl_key {
534
+ state.api_by_arl.write().await.remove(arl);
535
+ } else {
536
+ *state.api.lock().await = APIClient::new();
537
+ }
538
  }
539
+
540
  let client = api_client_for_arl(&state, arl_key.clone()).await;
541
  let res = {
542
+ let mut c = client.lock().await;
543
+ // На попытке 2 — принудительно обновляем токены перед запросом
544
+ if attempt == 2 {
545
+ if let Err(e) = c.force_renew().await {
546
+ last_err = Some(format!("renew:{}", e));
547
+ break;
548
+ }
549
+ }
550
+ media_url_for_track(&mut c, q.id, &formats).await
551
  };
552
 
553
  match res {
554
+ Ok(u) => { url = Some(u); break; }
555
+ Err(e) => { last_err = Some(e); }
 
 
 
 
 
 
 
 
 
 
 
 
 
 
556
  }
557
  }
558
 
559
  let url = match url {
560
  Some(u) => u,
561
+ None => return (
562
+ StatusCode::SERVICE_UNAVAILABLE,
563
+ last_err.unwrap_or_else(|| "stream:no_url".to_string()),
564
+ ).into_response(),
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
565
  };
566
 
567
  fn parse_range_header(v: &str) -> Option<(u64, Option<u64>)> {