recycleactor commited on
Commit
4352622
·
verified ·
1 Parent(s): df27c84

Upload main.rs

Browse files
Files changed (1) hide show
  1. src/main.rs +169 -38
src/main.rs CHANGED
@@ -40,6 +40,13 @@ struct ArlSession {
40
  }
41
 
42
  #[derive(Clone)]
 
 
 
 
 
 
 
43
  struct AppState {
44
  api: Arc<Mutex<APIClient>>,
45
  api_by_arl: Arc<RwLock<HashMap<String, Arc<Mutex<APIClient>>>>>,
@@ -47,6 +54,8 @@ struct AppState {
47
  pair: Arc<RwLock<HashMap<String, PairSession>>>,
48
  pair_store_path: String,
49
  arl_store_path: String,
 
 
50
  }
51
 
52
  fn license_token_from_ud(ud: &serde_json::Value) -> String {
@@ -1117,12 +1126,15 @@ async fn download(
1117
  .into_response()
1118
  }
1119
 
1120
- async fn stream(
 
 
 
1121
  State(state): State<AppState>,
1122
  Query(q): Query<StreamParams>,
1123
- headers: HeaderMap,
1124
  ) -> impl IntoResponse {
1125
- let formats = match q.format.as_deref().unwrap_or("AUTO") {
 
1126
  "FLAC" => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
1127
  "MP3_320" => vec![Format::MP3_320, Format::MP3_128],
1128
  "MP3_MISC" => vec![Format::MP3_MISC, Format::MP3_128],
@@ -1132,82 +1144,199 @@ async fn stream(
1132
 
1133
  let arl_key = q.arl.clone()
1134
  .map(|s| s.trim().to_string())
1135
- .filter(|s| !s.is_empty());
1136
-
1137
- let mut url: Option<String> = None;
1138
- let mut last_err: Option<String> = None;
1139
 
1140
  if arl_key.is_none() {
1141
- tracing::error!("stream {}: arl missing", q.id);
1142
- return (StatusCode::SERVICE_UNAVAILABLE, "arl required").into_response();
1143
  }
1144
 
 
 
 
 
 
 
 
 
 
 
 
 
1145
  let client = api_client_for_arl(&state, arl_key.clone()).await;
1146
  let target_id = {
1147
  let mut c = client.lock().await;
1148
  resolve_fallback_id(&mut c, q.id).await
1149
  };
1150
  let mut decrypt_id = target_id;
 
1151
 
1152
- // Try DZMEDIA_UPSTREAM first — gives highest available quality (FLAC/320)
1153
- // We already have correct decrypt_id from ARL above, so decryption will be correct
1154
- if url.is_none() {
1155
- if let Ok(text) = upstream_get_url_text(&formats, &vec![target_id]).await {
1156
- if let Ok(v) = serde_json::from_str::<serde_json::Value>(&text) {
1157
- let format_strings: Vec<String> = formats.iter().map(|f| format!("{:?}", f)).collect();
1158
- let format_strs: Vec<&str> = format_strings.iter().map(|s| s.as_str()).collect();
1159
- let (u, _, _) = extract_media_url_and_format(&v, &format_strs);
1160
- if !u.is_empty() {
1161
- url = Some(u);
1162
- }
1163
- }
1164
  }
1165
  }
1166
 
1167
- // Fallback: use ARL to get CDN URL via Deezer API (quality limited by account tier)
1168
  for attempt in 0..2u8 {
1169
  if url.is_some() { break; }
1170
-
1171
  if attempt >= 1 {
1172
- if let Some(ref arl) = arl_key {
1173
- state.api_by_arl.write().await.remove(arl);
1174
- } else {
1175
- *state.api.lock().await = APIClient::new();
1176
- }
1177
  }
1178
-
1179
  let client = api_client_for_arl(&state, arl_key.clone()).await;
1180
  let res = {
1181
  let mut c = client.lock().await;
1182
  if attempt >= 1 {
1183
- if let Err(e) = c.force_renew().await {
1184
- last_err = Some(format!("renew:{}", e));
1185
- continue;
1186
- }
1187
  }
1188
- match timeout(Duration::from_secs(12), media_url_for_track(&mut c, target_id, &formats)).await {
1189
  Ok(v) => v,
1190
  Err(_) => Err("timeout".to_string()),
1191
  }
1192
  };
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1193
 
1194
- match res {
1195
- Ok((u, _, id_used)) => { url = Some(u); decrypt_id = id_used; break; }
1196
- Err(e) => { last_err = Some(e); }
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1197
  }
1198
  }
1199
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1200
 
1201
  let url = match url {
1202
  Some(u) => u,
1203
  None => {
1204
  let err = last_err.unwrap_or_else(|| "stream:no_url".to_string());
1205
  tracing::error!("stream {} failed: {}", q.id, err);
1206
- // Return 404 so the browser audio element gracefully fails without triggering demuxer errors
1207
  return (StatusCode::NOT_FOUND, "Not Found").into_response();
1208
  }
1209
  };
1210
 
 
1211
  fn parse_range_header(v: &str) -> Option<(u64, Option<u64>)> {
1212
  let v = v.trim();
1213
  let v = v.strip_prefix("bytes=")?;
@@ -3644,6 +3773,7 @@ async fn main() {
3644
  pair: Arc::new(RwLock::new(HashMap::new())),
3645
  pair_store_path: std::env::var("PAIR_STORE_PATH").ok().unwrap_or_else(|| "pair_store.json".to_string()),
3646
  arl_store_path: std::env::var("ARL_STORE_PATH").ok().unwrap_or_else(|| "arl_store.json".to_string()),
 
3647
  };
3648
  arl_store_load(&state).await;
3649
  pair_store_load(&state).await;
@@ -3662,6 +3792,7 @@ async fn main() {
3662
  .route("/stream", get(stream))
3663
  .route("/download", get(download))
3664
  .route("/send_audio", get(send_audio))
 
3665
  .route("/stream_info", get(stream_info))
3666
  .route("/user_data", post(user_data))
3667
  .route("/playlists", post(playlists))
 
40
  }
41
 
42
  #[derive(Clone)]
43
+ struct CdnCacheEntry {
44
+ cdn_url: String,
45
+ decrypt_id: u64,
46
+ format: String,
47
+ expires_at: std::time::Instant,
48
+ }
49
+
50
  struct AppState {
51
  api: Arc<Mutex<APIClient>>,
52
  api_by_arl: Arc<RwLock<HashMap<String, Arc<Mutex<APIClient>>>>>,
 
54
  pair: Arc<RwLock<HashMap<String, PairSession>>>,
55
  pair_store_path: String,
56
  arl_store_path: String,
57
+ /// CDN URL кэш для предзагрузки. Ключ: track_id. TTL 20 мин.
58
+ cdn_cache: Arc<RwLock<HashMap<u64, CdnCacheEntry>>>,
59
  }
60
 
61
  fn license_token_from_ud(ud: &serde_json::Value) -> String {
 
1126
  .into_response()
1127
  }
1128
 
1129
+ /// /warm — прогревает CDN URL кэш для трека.
1130
+ /// Render вызывает этот эндпоинт ПЕРЕД отправкой sendAudio,
1131
+ /// чтобы к моменту Telegram-запроса кэш уже был горячим.
1132
+ async fn warm(
1133
  State(state): State<AppState>,
1134
  Query(q): Query<StreamParams>,
 
1135
  ) -> impl IntoResponse {
1136
+ let fmt_str = q.format.as_deref().unwrap_or("AUTO");
1137
+ let formats = match fmt_str {
1138
  "FLAC" => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
1139
  "MP3_320" => vec![Format::MP3_320, Format::MP3_128],
1140
  "MP3_MISC" => vec![Format::MP3_MISC, Format::MP3_128],
 
1144
 
1145
  let arl_key = q.arl.clone()
1146
  .map(|s| s.trim().to_string())
1147
+ .filter(|s| !s.is_empty())
1148
+ .or_else(|| std::env::var("DEEZER_ARL").ok().filter(|s| !s.is_empty()));
 
 
1149
 
1150
  if arl_key.is_none() {
1151
+ return json_error(StatusCode::SERVICE_UNAVAILABLE, "arl required");
 
1152
  }
1153
 
1154
+ // Проверяем — может уже в кэше есть свежая запись
1155
+ {
1156
+ let cache = state.cdn_cache.read().await;
1157
+ if let Some(e) = cache.get(&q.id) {
1158
+ if e.expires_at > std::time::Instant::now() {
1159
+ tracing::debug!("warm {}: already cached", q.id);
1160
+ return (StatusCode::OK, Json(json!({ "ok": true, "cached": true }))).into_response();
1161
+ }
1162
+ }
1163
+ }
1164
+
1165
+ // Resolve fallback_id
1166
  let client = api_client_for_arl(&state, arl_key.clone()).await;
1167
  let target_id = {
1168
  let mut c = client.lock().await;
1169
  resolve_fallback_id(&mut c, q.id).await
1170
  };
1171
  let mut decrypt_id = target_id;
1172
+ let mut url: Option<String> = None;
1173
 
1174
+ // Try upstream first
1175
+ if let Ok(text) = upstream_get_url_text(&formats, &vec![target_id]).await {
1176
+ if let Ok(v) = serde_json::from_str::<serde_json::Value>(&text) {
1177
+ let fstrs: Vec<String> = formats.iter().map(|f| format!("{:?}", f)).collect();
1178
+ let fstrs_ref: Vec<&str> = fstrs.iter().map(|s| s.as_str()).collect();
1179
+ let (u, _, _) = extract_media_url_and_format(&v, &fstrs_ref);
1180
+ if !u.is_empty() { url = Some(u); }
 
 
 
 
 
1181
  }
1182
  }
1183
 
1184
+ // Fallback: ARL
1185
  for attempt in 0..2u8 {
1186
  if url.is_some() { break; }
 
1187
  if attempt >= 1 {
1188
+ if let Some(ref arl) = arl_key { state.api_by_arl.write().await.remove(arl); }
 
 
 
 
1189
  }
 
1190
  let client = api_client_for_arl(&state, arl_key.clone()).await;
1191
  let res = {
1192
  let mut c = client.lock().await;
1193
  if attempt >= 1 {
1194
+ if let Err(_) = c.force_renew().await { continue; }
 
 
 
1195
  }
1196
+ match timeout(Duration::from_secs(15), media_url_for_track(&mut c, target_id, &formats)).await {
1197
  Ok(v) => v,
1198
  Err(_) => Err("timeout".to_string()),
1199
  }
1200
  };
1201
+ if let Ok((u, _, id_used)) = res { url = Some(u); decrypt_id = id_used; break; }
1202
+ }
1203
+
1204
+ match url {
1205
+ Some(u) => {
1206
+ let mut cache = state.cdn_cache.write().await;
1207
+ if cache.len() > 200 { cache.retain(|_, e| e.expires_at > std::time::Instant::now()); }
1208
+ cache.insert(q.id, CdnCacheEntry {
1209
+ cdn_url: u,
1210
+ decrypt_id,
1211
+ format: fmt_str.to_string(),
1212
+ expires_at: std::time::Instant::now() + Duration::from_secs(20 * 60),
1213
+ });
1214
+ tracing::info!("warm {}: cached OK", q.id);
1215
+ (StatusCode::OK, Json(json!({ "ok": true, "cached": false }))).into_response()
1216
+ }
1217
+ None => {
1218
+ tracing::error!("warm {}: failed to get CDN URL", q.id);
1219
+ json_error(StatusCode::NOT_FOUND, "track not found")
1220
+ }
1221
+ }
1222
+ }
1223
 
1224
+ async fn stream(
1225
+ State(state): State<AppState>,
1226
+ Query(q): Query<StreamParams>,
1227
+ headers: HeaderMap,
1228
+ ) -> impl IntoResponse {
1229
+ let formats = match q.format.as_deref().unwrap_or("AUTO") {
1230
+ "FLAC" => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
1231
+ "MP3_320" => vec![Format::MP3_320, Format::MP3_128],
1232
+ "MP3_MISC" => vec![Format::MP3_MISC, Format::MP3_128],
1233
+ "MP3_128" => vec![Format::MP3_128],
1234
+ _ => vec![Format::FLAC, Format::MP3_320, Format::MP3_128],
1235
+ };
1236
+
1237
+ let arl_key = q.arl.clone()
1238
+ .map(|s| s.trim().to_string())
1239
+ .filter(|s| !s.is_empty());
1240
+
1241
+ let mut url: Option<String> = None;
1242
+ let mut decrypt_id: u64 = q.id; // overwritten after resolve_fallback_id
1243
+ let mut last_err: Option<String> = None;
1244
+
1245
+ if arl_key.is_none() {
1246
+ tracing::error!("stream {}: arl missing", q.id);
1247
+ return (StatusCode::SERVICE_UNAVAILABLE, "arl required").into_response();
1248
+ }
1249
+
1250
+ // ── CDN кэш: если есть свежая запись — пропускаем все Deezer API вызовы ──
1251
+ {
1252
+ let cache = state.cdn_cache.read().await;
1253
+ if let Some(e) = cache.get(&q.id) {
1254
+ if e.expires_at > std::time::Instant::now() {
1255
+ tracing::debug!("stream {}: CDN cache hit ({})", q.id, e.format);
1256
+ url = Some(e.cdn_url.clone());
1257
+ decrypt_id = e.decrypt_id;
1258
+ }
1259
  }
1260
  }
1261
 
1262
+ if url.is_none() {
1263
+ // Resolve fallback_id for proper decryption key
1264
+ let client = api_client_for_arl(&state, arl_key.clone()).await;
1265
+ let target_id = {
1266
+ let mut c = client.lock().await;
1267
+ resolve_fallback_id(&mut c, q.id).await
1268
+ };
1269
+ decrypt_id = target_id;
1270
+
1271
+ // Try DZMEDIA_UPSTREAM first — gives highest available quality (FLAC/320)
1272
+ if url.is_none() {
1273
+ if let Ok(text) = upstream_get_url_text(&formats, &vec![target_id]).await {
1274
+ if let Ok(v) = serde_json::from_str::<serde_json::Value>(&text) {
1275
+ let format_strings: Vec<String> = formats.iter().map(|f| format!("{:?}", f)).collect();
1276
+ let format_strs: Vec<&str> = format_strings.iter().map(|s| s.as_str()).collect();
1277
+ let (u, _, _) = extract_media_url_and_format(&v, &format_strs);
1278
+ if !u.is_empty() {
1279
+ url = Some(u);
1280
+ }
1281
+ }
1282
+ }
1283
+ }
1284
+
1285
+ // Fallback: use ARL to get CDN URL via Deezer API
1286
+ for attempt in 0..2u8 {
1287
+ if url.is_some() { break; }
1288
+ if attempt >= 1 {
1289
+ if let Some(ref arl) = arl_key {
1290
+ state.api_by_arl.write().await.remove(arl);
1291
+ } else {
1292
+ *state.api.lock().await = APIClient::new();
1293
+ }
1294
+ }
1295
+ let client = api_client_for_arl(&state, arl_key.clone()).await;
1296
+ let res = {
1297
+ let mut c = client.lock().await;
1298
+ if attempt >= 1 {
1299
+ if let Err(e) = c.force_renew().await {
1300
+ last_err = Some(format!("renew:{}", e));
1301
+ continue;
1302
+ }
1303
+ }
1304
+ match timeout(Duration::from_secs(12), media_url_for_track(&mut c, target_id, &formats)).await {
1305
+ Ok(v) => v,
1306
+ Err(_) => Err("timeout".to_string()),
1307
+ }
1308
+ };
1309
+ match res {
1310
+ Ok((u, fmt, id_used)) => { url = Some(u); decrypt_id = id_used; let _ = fmt; break; }
1311
+ Err(e) => { last_err = Some(e); }
1312
+ }
1313
+ }
1314
+
1315
+ // ── Пишем в кэш если удалось получить URL ──
1316
+ if let Some(ref u) = url {
1317
+ let mut cache = state.cdn_cache.write().await;
1318
+ if cache.len() > 200 {
1319
+ cache.retain(|_, e| e.expires_at > std::time::Instant::now());
1320
+ }
1321
+ cache.insert(q.id, CdnCacheEntry {
1322
+ cdn_url: u.clone(),
1323
+ decrypt_id,
1324
+ format: q.format.clone().unwrap_or_else(|| "AUTO".to_string()),
1325
+ expires_at: std::time::Instant::now() + Duration::from_secs(20 * 60),
1326
+ });
1327
+ }
1328
+ }
1329
 
1330
  let url = match url {
1331
  Some(u) => u,
1332
  None => {
1333
  let err = last_err.unwrap_or_else(|| "stream:no_url".to_string());
1334
  tracing::error!("stream {} failed: {}", q.id, err);
 
1335
  return (StatusCode::NOT_FOUND, "Not Found").into_response();
1336
  }
1337
  };
1338
 
1339
+
1340
  fn parse_range_header(v: &str) -> Option<(u64, Option<u64>)> {
1341
  let v = v.trim();
1342
  let v = v.strip_prefix("bytes=")?;
 
3773
  pair: Arc::new(RwLock::new(HashMap::new())),
3774
  pair_store_path: std::env::var("PAIR_STORE_PATH").ok().unwrap_or_else(|| "pair_store.json".to_string()),
3775
  arl_store_path: std::env::var("ARL_STORE_PATH").ok().unwrap_or_else(|| "arl_store.json".to_string()),
3776
+ cdn_cache: Arc::new(RwLock::new(HashMap::new())),
3777
  };
3778
  arl_store_load(&state).await;
3779
  pair_store_load(&state).await;
 
3792
  .route("/stream", get(stream))
3793
  .route("/download", get(download))
3794
  .route("/send_audio", get(send_audio))
3795
+ .route("/warm", get(warm))
3796
  .route("/stream_info", get(stream_info))
3797
  .route("/user_data", post(user_data))
3798
  .route("/playlists", post(playlists))