Загрузка данных
"""
BACK. — Music streaming platform API
=====================================
Хранилище: PostgreSQL (backdb), строка подключения в BACK_DB_DSN.
Раньше был SQLite, но он упирается в одновременные записи
(избранное, плейлисты) задолго до того, как упрётся канал.
Пароли: PBKDF2-HMAC-SHA256, 200k итераций, соль на пользователя.
Токены: HMAC-SHA256 подписанные, со сроком жизни.
Поиск треков — гибридный:
* Spotify Web API даёт метаданные, обложки, альбомы (Client Credentials).
* Spotify с ноября 2024 отдаёт preview_url = null новым приложениям,
поэтому реальный аудиопоток (30 сек, m4a) подтягивается из iTunes Search API
и матчится к Spotify-трекам по нормализованным title+artist.
Итог: поиск и обложки — Spotify, звук — iTunes. Играет по-настоящему.
"""
import os
import re
import json
import time
import hmac
import base64
import hashlib
import logging
import secrets
import unicodedata
import asyncio
from contextlib import contextmanager, asynccontextmanager
from typing import Optional, List, Dict, Any
import httpx
import uvicorn
import psycopg2
import psycopg2.extras
import psycopg2.pool
import psycopg2.errors
from fastapi import FastAPI, HTTPException, Depends, Header, Request
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import RedirectResponse, HTMLResponse
from pydantic import BaseModel, Field
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s [%(name)s] %(message)s",
)
logger = logging.getLogger("back.api")
# ==============================================================================
# КОНФИГУРАЦИЯ
# ==============================================================================
DB_DSN = os.environ.get("BACK_DB_DSN", "postgresql://backapi@127.0.0.1:5432/backdb")
SECRET_PATH = os.environ.get("BACK_SECRET_PATH", "/etc/back-api/secret.key")
SPOTIFY_CLIENT_ID = os.environ.get("SPOTIFY_CLIENT_ID", "9240df9f9ded4dc9bd65cd041d451092")
SPOTIFY_CLIENT_SECRET = os.environ.get("SPOTIFY_CLIENT_SECRET", "f5eebbcf20c743d4a8f4813debc505c1")
SOUNDCLOUD_CLIENT_ID = os.environ.get("SOUNDCLOUD_CLIENT_ID", "XAlx9jKj5xLWgoYO1Dot48UYF4E6BmXX")
SOUNDCLOUD_CLIENT_SECRET = os.environ.get("SOUNDCLOUD_CLIENT_SECRET", "zBH3fGP4yCnfS1q8hkULyWcl5o0DpBzb")
SPOTIFY_AUTH_URL = "https://accounts.spotify.com/api/token"
SPOTIFY_SEARCH_URL = "https://api.spotify.com/v1/search"
ITUNES_SEARCH_URL = "https://itunes.apple.com/search"
SOUNDCLOUD_TOKEN_URL = "https://api.soundcloud.com/oauth2/token"
SOUNDCLOUD_TRACKS_URL = "https://api.soundcloud.com/tracks"
# Ссылки на поток SoundCloud подписаны и живут минуты, поэтому приложение
# получает не сам CDN-адрес, а наш /stream/... — он резолвится в момент запуска.
PUBLIC_BASE = os.environ.get("BACK_PUBLIC_BASE", "https://api.back-mus.fun")
TOKEN_TTL = 60 * 60 * 24 * 30 # 30 дней
PBKDF2_ROUNDS = 200_000
USERNAME_RE = re.compile(r"^[a-zA-Z0-9_.\-]{3,24}$")
def _load_or_create_key(path: str, generate) -> bytes:
"""
Читает ключ с диска, а если его нет — создаёт атомарно (O_EXCL).
Это важно при нескольких воркерах: без O_EXCL каждый сгенерировал бы
свой ключ, и подписи одного воркера не проходили бы проверку у другого.
"""
for _ in range(5):
try:
with open(path, "rb") as fh:
data = fh.read().strip()
if data:
return data
except FileNotFoundError:
pass
os.makedirs(os.path.dirname(path), exist_ok=True)
try:
fd = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
except FileExistsError:
continue # другой воркер успел первым — на следующем круге перечитаем
with os.fdopen(fd, "wb") as fh:
fh.write(generate())
logger.info("Создан ключ: %s", path)
with open(path, "rb") as fh:
return fh.read().strip()
def _load_secret() -> bytes:
env = os.environ.get("BACK_SECRET")
if env:
return env.encode()
return _load_or_create_key(SECRET_PATH, lambda: secrets.token_bytes(48))
SECRET = _load_secret()
# ==============================================================================
# БАЗА ДАННЫХ
# ==============================================================================
SCHEMA = """
CREATE TABLE IF NOT EXISTS users (
id SERIAL PRIMARY KEY,
username TEXT NOT NULL,
username_ci TEXT NOT NULL UNIQUE, -- нижний регистр, чтобы Vasya и vasya не сосуществовали
pw_salt TEXT NOT NULL,
pw_hash TEXT NOT NULL,
is_premium INTEGER NOT NULL DEFAULT 0,
created_at BIGINT NOT NULL,
last_login_at BIGINT
);
CREATE TABLE IF NOT EXISTS favorites (
id SERIAL PRIMARY KEY,
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
track_id TEXT NOT NULL,
payload TEXT NOT NULL, -- JSON трека целиком
created_at BIGINT NOT NULL,
UNIQUE(user_id, track_id)
);
CREATE TABLE IF NOT EXISTS playlists (
id SERIAL PRIMARY KEY,
user_id INTEGER NOT NULL REFERENCES users(id) ON DELETE CASCADE,
title TEXT NOT NULL,
created_at BIGINT NOT NULL
);
CREATE TABLE IF NOT EXISTS playlist_tracks (
id SERIAL PRIMARY KEY,
playlist_id INTEGER NOT NULL REFERENCES playlists(id) ON DELETE CASCADE,
track_id TEXT NOT NULL,
payload TEXT NOT NULL,
position INTEGER NOT NULL DEFAULT 0,
UNIQUE(playlist_id, track_id)
);
CREATE TABLE IF NOT EXISTS kv (
k TEXT PRIMARY KEY,
v TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_fav_user ON favorites(user_id);
CREATE INDEX IF NOT EXISTS idx_pl_user ON playlists(user_id);
CREATE INDEX IF NOT EXISTS idx_plt_pl ON playlist_tracks(playlist_id);
"""
_pool: Optional[psycopg2.pool.ThreadedConnectionPool] = None
def _get_pool() -> psycopg2.pool.ThreadedConnectionPool:
global _pool
if _pool is None:
# По воркеру свой пул; max_connections в постгресе = 60, воркеров 2
_pool = psycopg2.pool.ThreadedConnectionPool(1, 12, dsn=DB_DSN)
return _pool
class _Cursor:
"""
Обёртка, повторяющая поведение sqlite3: execute возвращает объект
с fetchone/fetchall, а строки ведут себя как словари.
Плейсхолдеры ? переводятся в %s, чтобы не переписывать все запросы.
"""
def __init__(self, cur):
self._cur = cur
def fetchone(self):
return self._cur.fetchone()
def fetchall(self):
return self._cur.fetchall()
@property
def rowcount(self):
return self._cur.rowcount
class _Conn:
def __init__(self, raw):
self._raw = raw
def execute(self, sql: str, params=()):
cur = self._raw.cursor(cursor_factory=psycopg2.extras.RealDictCursor)
cur.execute(sql.replace("?", "%s"), params)
return _Cursor(cur)
def executescript(self, sql: str):
with self._raw.cursor() as cur:
cur.execute(sql)
def init_db() -> None:
with db() as conn:
conn.executescript(SCHEMA)
logger.info("База данных готова: PostgreSQL")
@contextmanager
def db():
pool = _get_pool()
raw = pool.getconn()
try:
yield _Conn(raw)
raw.commit()
except Exception:
raw.rollback()
raise
finally:
pool.putconn(raw)
# ==============================================================================
# ПАРОЛИ И ТОКЕНЫ
# ==============================================================================
def hash_password(password: str, salt: Optional[bytes] = None) -> Dict[str, str]:
salt = salt or secrets.token_bytes(16)
digest = hashlib.pbkdf2_hmac("sha256", password.encode("utf-8"), salt, PBKDF2_ROUNDS)
return {"salt": salt.hex(), "hash": digest.hex()}
def verify_password(password: str, salt_hex: str, hash_hex: str) -> bool:
try:
salt = bytes.fromhex(salt_hex)
except ValueError:
return False
digest = hashlib.pbkdf2_hmac("sha256", password.encode("utf-8"), salt, PBKDF2_ROUNDS)
return hmac.compare_digest(digest.hex(), hash_hex)
def _b64e(raw: bytes) -> str:
return base64.urlsafe_b64encode(raw).decode().rstrip("=")
def _b64d(txt: str) -> bytes:
return base64.urlsafe_b64decode(txt + "=" * (-len(txt) % 4))
def issue_token(user_id: int, username: str) -> str:
payload = {"uid": user_id, "u": username, "exp": int(time.time()) + TOKEN_TTL}
body = _b64e(json.dumps(payload, separators=(",", ":")).encode())
sig = _b64e(hmac.new(SECRET, body.encode(), hashlib.sha256).digest())
return f"{body}.{sig}"
def decode_token(token: str) -> Optional[Dict[str, Any]]:
try:
body, sig = token.split(".", 1)
except ValueError:
return None
expected = _b64e(hmac.new(SECRET, body.encode(), hashlib.sha256).digest())
if not hmac.compare_digest(sig, expected):
return None
try:
payload = json.loads(_b64d(body))
except Exception:
return None
if payload.get("exp", 0) < time.time():
return None
return payload
# ==============================================================================
# FASTAPI
# ==============================================================================
@asynccontextmanager
async def lifespan(_: FastAPI):
init_db()
yield
app = FastAPI(title="BACK. Music Streaming API", version="2.0.0", lifespan=lifespan)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=False,
allow_methods=["*"],
allow_headers=["*"],
)
class UserAuth(BaseModel):
username: str = Field(..., min_length=1, max_length=64)
password: str = Field(..., min_length=1, max_length=256)
class PlaylistCreate(BaseModel):
title: str = Field(..., min_length=1, max_length=60)
class TrackPayload(BaseModel):
id: str
title: str
artist: str = ""
album: str = ""
artwork: str = ""
preview: str = ""
duration: int = 0
source: str = ""
def current_user(authorization: str = Header(default="")) -> Dict[str, Any]:
"""Достаёт пользователя по Bearer-токену. 401, если токен битый/просрочен."""
if not authorization.lower().startswith("bearer "):
raise HTTPException(status_code=401, detail="Требуется авторизация")
payload = decode_token(authorization[7:].strip())
if not payload:
raise HTTPException(status_code=401, detail="Сессия истекла, войдите заново")
with db() as conn:
row = conn.execute("SELECT * FROM users WHERE id = ?", (payload["uid"],)).fetchone()
if not row:
raise HTTPException(status_code=401, detail="Пользователь не найден")
return row
# ------------------------------------------------------------------ ЗАЩИТА ОТ БРУТА
_attempts: Dict[str, List[float]] = {}
def rate_limit(key: str, limit: int, window: int) -> None:
now = time.time()
hits = [t for t in _attempts.get(key, []) if now - t < window]
if len(hits) >= limit:
raise HTTPException(status_code=429, detail="Слишком много попыток, подождите минуту")
hits.append(now)
_attempts[key] = hits
def validate_credentials(username: str, password: str) -> str:
username = username.strip()
if not USERNAME_RE.match(username):
raise HTTPException(
status_code=400,
detail="Юзернейм: 3–24 символа, только латиница, цифры, _ . -",
)
if not re.search(r"[a-zA-Z]", username):
raise HTTPException(status_code=400, detail="Юзернейм должен содержать хотя бы одну букву")
if len(password) < 6:
raise HTTPException(status_code=400, detail="Пароль должен содержать минимум 6 символов")
if not re.search(r"\d", password):
raise HTTPException(status_code=400, detail="Пароль должен содержать хотя бы одну цифру")
return username
# ==============================================================================
# АВТОРИЗАЦИЯ
# ==============================================================================
@app.get("/health")
async def health() -> Dict[str, Any]:
with db() as conn:
users = conn.execute("SELECT COUNT(*) AS c FROM users").fetchone()["c"]
return {"status": "ok", "users": users, "version": app.version}
@app.get("/username-available")
async def username_available(username: str) -> Dict[str, Any]:
"""Живая проверка занятости ника — приложение дёргает её во время набора."""
uname = username.strip()
if not USERNAME_RE.match(uname):
return {"available": False, "reason": "invalid"}
with db() as conn:
row = conn.execute(
"SELECT 1 FROM users WHERE username_ci = ?", (uname.lower(),)
).fetchone()
return {"available": row is None}
@app.post("/register", status_code=201)
async def register(user: UserAuth, request: Request) -> Dict[str, Any]:
rate_limit(f"reg:{request.client.host}", limit=10, window=300)
username = validate_credentials(user.username, user.password)
creds = hash_password(user.password)
now = int(time.time())
try:
with db() as conn:
cur = conn.execute(
"INSERT INTO users (username, username_ci, pw_salt, pw_hash, created_at) "
"VALUES (?, ?, ?, ?, ?) RETURNING id",
(username, username.lower(), creds["salt"], creds["hash"], now),
)
user_id = cur.fetchone()["id"]
# Стартовый плейлист, чтобы медиатека не была пустой
conn.execute(
"INSERT INTO playlists (user_id, title, created_at) VALUES (?, ?, ?)",
(user_id, "Мой первый плейлист", now),
)
except psycopg2.errors.UniqueViolation:
raise HTTPException(status_code=409, detail="Этот юзернейм уже занят")
logger.info("Регистрация: %s (id=%s)", username, user_id)
return {
"status": "ok",
"message": "Регистрация успешна",
"token": issue_token(user_id, username),
"user": {"id": user_id, "username": username, "is_premium": False},
}
@app.post("/login")
async def login(user: UserAuth, request: Request) -> Dict[str, Any]:
rate_limit(f"login:{request.client.host}", limit=20, window=300)
uname = user.username.strip()
with db() as conn:
row = conn.execute(
"SELECT * FROM users WHERE username_ci = ?", (uname.lower(),)
).fetchone()
# Одинаковый ответ для «нет юзера» и «неверный пароль» — не подсказываем перебору
if not row or not verify_password(user.password, row["pw_salt"], row["pw_hash"]):
raise HTTPException(status_code=401, detail="Неверный юзернейм или пароль")
with db() as conn:
conn.execute(
"UPDATE users SET last_login_at = ? WHERE id = ?", (int(time.time()), row["id"])
)
logger.info("Вход: %s", row["username"])
return {
"status": "success",
"message": "Вход выполнен",
"token": issue_token(row["id"], row["username"]),
"expires_in": TOKEN_TTL,
"user": {
"id": row["id"],
"username": row["username"],
"is_premium": bool(row["is_premium"]),
},
}
@app.get("/me")
async def me(user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
"""Проверка сохранённой сессии при старте приложения."""
return {
"status": "ok",
"user": {
"id": user["id"],
"username": user["username"],
"is_premium": bool(user["is_premium"]),
},
}
# ==============================================================================
# ПОИСК: SPOTIFY (МЕТАДАННЫЕ) + ITUNES (АУДИО)
# ==============================================================================
_spotify_token: Dict[str, Any] = {"value": "", "expires": 0.0}
async def get_spotify_token(client: httpx.AsyncClient) -> str:
"""Client Credentials Flow с кэшем — не выбиваем новый токен на каждый поиск."""
if _spotify_token["value"] and _spotify_token["expires"] > time.time() + 30:
return _spotify_token["value"]
auth = base64.b64encode(
f"{SPOTIFY_CLIENT_ID}:{SPOTIFY_CLIENT_SECRET}".encode()
).decode()
resp = await client.post(
SPOTIFY_AUTH_URL,
headers={"Authorization": f"Basic {auth}",
"Content-Type": "application/x-www-form-urlencoded"},
data={"grant_type": "client_credentials"},
timeout=10.0,
)
resp.raise_for_status()
data = resp.json()
_spotify_token["value"] = data.get("access_token", "")
_spotify_token["expires"] = time.time() + int(data.get("expires_in", 3600))
return _spotify_token["value"]
def norm(text: str) -> str:
"""Нормализация для матчинга Spotify↔iTunes: убираем регистр, скобки, знаки."""
text = unicodedata.normalize("NFKD", text or "").lower()
text = re.sub(r"\((feat|ft|with)[^)]*\)", " ", text)
text = re.sub(r"\[[^\]]*\]", " ", text)
text = re.sub(r"\b(feat|ft|remaster(ed)?|radio edit|explicit)\b.*", " ", text)
text = re.sub(r"[^a-z0-9а-яё]+", " ", text)
return " ".join(text.split())
async def fetch_spotify(client: httpx.AsyncClient, query: str, limit: int) -> List[Dict[str, Any]]:
try:
token = await get_spotify_token(client)
if not token:
return []
resp = await client.get(
SPOTIFY_SEARCH_URL,
headers={"Authorization": f"Bearer {token}"},
params={"q": query, "type": "track", "limit": min(limit, 50), "market": "US"},
timeout=10.0,
)
resp.raise_for_status()
items = resp.json().get("tracks", {}).get("items", []) or []
except Exception as exc:
logger.warning("Spotify поиск не удался (%s): %s", query, exc)
return []
out = []
for it in items:
if not it:
continue
images = (it.get("album") or {}).get("images") or []
artists = ", ".join(a.get("name", "") for a in it.get("artists", []) if a.get("name"))
out.append({
"id": f"sp:{it.get('id')}",
"title": it.get("name", ""),
"artist": artists or "Unknown Artist",
"album": (it.get("album") or {}).get("name", ""),
"artwork": images[0]["url"] if images else "",
# Spotify отдаёт preview_url = null новым client-credentials приложениям,
# но если он вдруг есть — используем его как есть.
"preview": it.get("preview_url") or "",
"duration": int(it.get("duration_ms", 0) // 1000),
"source": "Spotify",
"spotify_url": (it.get("external_urls") or {}).get("spotify", ""),
})
return out
_soundcloud_token: Dict[str, Any] = {"value": "", "expires": 0.0}
async def get_soundcloud_token(client: httpx.AsyncClient) -> str:
"""
Официальный OAuth2 client_credentials SoundCloud.
Прошлая версия била в неофициальный api-v2 с сырым client_id — оттуда и брались 403.
"""
if _soundcloud_token["value"] and _soundcloud_token["expires"] > time.time() + 30:
return _soundcloud_token["value"]
resp = await client.post(
SOUNDCLOUD_TOKEN_URL,
headers={"Content-Type": "application/x-www-form-urlencoded"},
data={
"grant_type": "client_credentials",
"client_id": SOUNDCLOUD_CLIENT_ID,
"client_secret": SOUNDCLOUD_CLIENT_SECRET,
},
timeout=10.0,
)
resp.raise_for_status()
data = resp.json()
_soundcloud_token["value"] = data.get("access_token", "")
_soundcloud_token["expires"] = time.time() + int(data.get("expires_in", 3600))
return _soundcloud_token["value"]
async def fetch_soundcloud(client: httpx.AsyncClient, query: str, limit: int) -> List[Dict[str, Any]]:
try:
token = await get_soundcloud_token(client)
if not token:
return []
resp = await client.get(
SOUNDCLOUD_TRACKS_URL,
headers={"Authorization": f"OAuth {token}"},
params={"q": query, "limit": min(limit, 50), "access": "playable"},
timeout=12.0,
)
resp.raise_for_status()
items = resp.json() or []
except Exception as exc:
logger.warning("SoundCloud поиск не удался (%s): %s", query, exc)
return []
out = []
for it in items:
if not it or not it.get("id"):
continue
art = it.get("artwork_url") or (it.get("user") or {}).get("avatar_url") or ""
art = art.replace("-large.", "-t500x500.")
out.append({
"id": f"sc:{it['id']}",
"title": it.get("title", ""),
"artist": (it.get("user") or {}).get("username", "SoundCloud Artist"),
"album": "",
"artwork": art,
"preview": f"{PUBLIC_BASE}/stream/sc:{it['id']}",
"duration": int((it.get("duration") or 0) // 1000),
"source": "SoundCloud",
"permalink_url": it.get("permalink_url", ""),
})
return out
@app.get("/stream/{track_id}")
async def stream(track_id: str):
"""
Отдаёт 302 на подписанный CDN-адрес SoundCloud.
Токен приложения остаётся на сервере, а ссылка резолвится под каждое включение,
потому что подпись в ней живёт всего несколько минут.
"""
if not track_id.startswith("sc:"):
raise HTTPException(status_code=400, detail="Неизвестный источник потока")
sc_id = track_id[3:]
if not sc_id.isdigit():
raise HTTPException(status_code=400, detail="Некорректный идентификатор трека")
async with httpx.AsyncClient(follow_redirects=False) as client:
try:
token = await get_soundcloud_token(client)
resp = await client.get(
f"{SOUNDCLOUD_TRACKS_URL}/{sc_id}/stream",
headers={"Authorization": f"OAuth {token}"},
timeout=12.0,
)
except Exception as exc:
logger.warning("SoundCloud stream %s: %s", sc_id, exc)
raise HTTPException(status_code=502, detail="Источник недоступен")
location = resp.headers.get("location")
if resp.status_code in (301, 302, 303, 307, 308) and location:
return RedirectResponse(url=location, status_code=302)
raise HTTPException(status_code=404, detail="Поток для этого трека недоступен")
async def fetch_itunes(client: httpx.AsyncClient, query: str, limit: int) -> List[Dict[str, Any]]:
try:
resp = await client.get(
ITUNES_SEARCH_URL,
params={"term": query, "media": "music", "entity": "song", "limit": min(limit, 50)},
timeout=10.0,
)
resp.raise_for_status()
items = resp.json().get("results", []) or []
except Exception as exc:
logger.warning("iTunes поиск не удался (%s): %s", query, exc)
return []
out = []
for it in items:
preview = it.get("previewUrl") or ""
if not preview:
continue
art = (it.get("artworkUrl100") or "").replace("100x100bb", "600x600bb")
out.append({
"id": f"it:{it.get('trackId')}",
"title": it.get("trackName", ""),
"artist": it.get("artistName", ""),
"album": it.get("collectionName", ""),
"artwork": art,
"preview": preview,
"duration": int((it.get("trackTimeMillis") or 0) // 1000),
"source": "iTunes",
})
return out
@app.get("/search")
async def search(query: str = "", limit: int = 25, source: str = "spotify") -> Dict[str, Any]:
"""
source = spotify — каталог Spotify (метаданные и обложки), звук подставляется из iTunes
source = soundcloud — официальный SoundCloud API, звук через наш /stream
source = all — и то, и другое
Треки без проигрываемого потока в выдачу не попадают: в приложении
не должно быть кнопок, которые ничего не играют.
"""
query = (query or "").strip()
if not query:
raise HTTPException(status_code=400, detail="Параметр query обязателен")
source = (source or "spotify").lower()
if source not in ("spotify", "soundcloud", "all"):
source = "spotify"
if source == "soundcloud":
async with httpx.AsyncClient(follow_redirects=True) as client:
sc = await fetch_soundcloud(client, query, limit)
return {"status": "success", "query": query, "source": source,
"count": len(sc[:limit]), "results": sc[:limit]}
async with httpx.AsyncClient(follow_redirects=True) as client:
if source == "all":
spotify, itunes, soundcloud = await asyncio.gather(
fetch_spotify(client, query, limit),
fetch_itunes(client, query, 50),
fetch_soundcloud(client, query, limit),
)
else:
spotify, itunes = await asyncio.gather(
fetch_spotify(client, query, limit),
fetch_itunes(client, query, 50),
)
soundcloud = []
# Индекс iTunes по нормализованным ключам: (название+артист) и просто название
by_pair: Dict[str, Dict[str, Any]] = {}
by_title: Dict[str, Dict[str, Any]] = {}
for t in itunes:
pair = f"{norm(t['title'])}|{norm(t['artist'])}"
by_pair.setdefault(pair, t)
by_title.setdefault(norm(t["title"]), t)
results: List[Dict[str, Any]] = []
used_itunes: set = set()
for sp in spotify:
if sp["preview"]:
results.append(sp)
continue
pair = f"{norm(sp['title'])}|{norm(sp['artist'])}"
match = by_pair.get(pair)
if not match:
cand = by_title.get(norm(sp["title"]))
# Совпадение по названию принимаем, только если артист пересекается
if cand and norm(cand["artist"]).split() and (
set(norm(cand["artist"]).split()) & set(norm(sp["artist"]).split())
):
match = cand
if match:
used_itunes.add(match["id"])
merged = dict(sp)
merged["preview"] = match["preview"]
merged["duration"] = sp["duration"] or match["duration"]
merged["artwork"] = sp["artwork"] or match["artwork"]
merged["source"] = "Spotify"
results.append(merged)
# Добиваем выдачу iTunes-треками, которых не было в Spotify
for t in itunes:
if len(results) >= limit:
break
if t["id"] in used_itunes:
continue
if any(norm(r["title"]) == norm(t["title"]) and
norm(r["artist"]) == norm(t["artist"]) for r in results):
continue
results.append(t)
# В режиме "all" дописываем SoundCloud в хвост, не вытесняя каталог Spotify
for t in soundcloud:
if len(results) >= limit:
break
results.append(t)
return {"status": "success", "query": query, "source": source,
"count": len(results[:limit]), "results": results[:limit]}
# ==============================================================================
# ИЗБРАННОЕ
# ==============================================================================
@app.get("/favorites")
async def list_favorites(user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
with db() as conn:
rows = conn.execute(
"SELECT payload FROM favorites WHERE user_id = ? ORDER BY created_at DESC",
(user["id"],),
).fetchall()
return {"status": "ok", "results": [json.loads(r["payload"]) for r in rows]}
@app.post("/favorites")
async def add_favorite(track: TrackPayload,
user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
with db() as conn:
conn.execute(
"INSERT INTO favorites (user_id, track_id, payload, created_at) "
"VALUES (?, ?, ?, ?) ON CONFLICT (user_id, track_id) DO NOTHING",
(user["id"], track.id, json.dumps(track.model_dump(), ensure_ascii=False),
int(time.time())),
)
return {"status": "ok"}
@app.delete("/favorites/{track_id:path}")
async def remove_favorite(track_id: str,
user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
with db() as conn:
conn.execute(
"DELETE FROM favorites WHERE user_id = ? AND track_id = ?", (user["id"], track_id)
)
return {"status": "ok"}
# ==============================================================================
# ПЛЕЙЛИСТЫ
# ==============================================================================
@app.get("/playlists")
async def list_playlists(user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
with db() as conn:
pls = conn.execute(
"SELECT id, title, created_at FROM playlists WHERE user_id = ? ORDER BY created_at",
(user["id"],),
).fetchall()
out = []
for p in pls:
tracks = conn.execute(
"SELECT payload FROM playlist_tracks WHERE playlist_id = ? ORDER BY position, id",
(p["id"],),
).fetchall()
out.append({
"id": str(p["id"]),
"title": p["title"],
"count": len(tracks),
"tracks": [json.loads(t["payload"]) for t in tracks],
})
return {"status": "ok", "results": out}
@app.post("/playlists", status_code=201)
async def create_playlist(body: PlaylistCreate,
user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
title = body.title.strip()
if not title:
raise HTTPException(status_code=400, detail="Название плейлиста не может быть пустым")
with db() as conn:
cur = conn.execute(
"INSERT INTO playlists (user_id, title, created_at) VALUES (?, ?, ?) RETURNING id",
(user["id"], title, int(time.time())),
)
pid = cur.fetchone()["id"]
return {"status": "ok", "playlist": {"id": str(pid), "title": title, "count": 0, "tracks": []}}
@app.delete("/playlists/{playlist_id}")
async def delete_playlist(playlist_id: int,
user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
with db() as conn:
cur = conn.execute(
"DELETE FROM playlists WHERE id = ? AND user_id = ?", (playlist_id, user["id"])
)
if cur.rowcount == 0:
raise HTTPException(status_code=404, detail="Плейлист не найден")
return {"status": "ok"}
@app.post("/playlists/{playlist_id}/tracks")
async def add_to_playlist(playlist_id: int, track: TrackPayload,
user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
with db() as conn:
owner = conn.execute(
"SELECT 1 FROM playlists WHERE id = ? AND user_id = ?", (playlist_id, user["id"])
).fetchone()
if not owner:
raise HTTPException(status_code=404, detail="Плейлист не найден")
pos = conn.execute(
"SELECT COALESCE(MAX(position), 0) + 1 AS p FROM playlist_tracks WHERE playlist_id = ?",
(playlist_id,),
).fetchone()["p"]
conn.execute(
"INSERT INTO playlist_tracks (playlist_id, track_id, payload, position) "
"VALUES (?, ?, ?, ?) ON CONFLICT (playlist_id, track_id) DO NOTHING",
(playlist_id, track.id, json.dumps(track.model_dump(), ensure_ascii=False), pos),
)
return {"status": "ok"}
@app.delete("/playlists/{playlist_id}/tracks/{track_id:path}")
async def remove_from_playlist(playlist_id: int, track_id: str,
user: Dict[str, Any] = Depends(current_user)) -> Dict[str, Any]:
with db() as conn:
owner = conn.execute(
"SELECT 1 FROM playlists WHERE id = ? AND user_id = ?", (playlist_id, user["id"])
).fetchone()
if not owner:
raise HTTPException(status_code=404, detail="Плейлист не найден")
conn.execute(
"DELETE FROM playlist_tracks WHERE playlist_id = ? AND track_id = ?",
(playlist_id, track_id),
)
return {"status": "ok"}
# ==============================================================================
# ЭКСПЕРИМЕНТ: ПОЛЬЗОВАТЕЛЬСКИЙ OAUTH SOUNDCLOUD
# ==============================================================================
# client_credentials измеренно отдаёт 30 секунд (cf-preview-media/preview/0/30/).
# Здесь проверяется гипотеза: снимает ли ограничение токен, выданный под
# конкретного пользователя. Эндпоинты закрыты ключом — это диагностика, не фича.
# Общий адрес возврата: он прописан и в Spotify, и в SoundCloud,
# поэтому провайдер определяется по state, а не по пути.
SC_REDIRECT_URI = os.environ.get("SC_REDIRECT_URI", f"{PUBLIC_BASE}/callback")
SC_AUTHORIZE_URL = "https://secure.soundcloud.com/authorize"
SC_TOKEN_URLS = [
"https://secure.soundcloud.com/oauth/token",
"https://api.soundcloud.com/oauth2/token",
]
DIAG_KEY_PATH = os.environ.get("BACK_DIAG_KEY_PATH", "/etc/back-api/diag.key")
def _check_diag(key: str) -> None:
# Читаем ключ с диска на каждой проверке: диагностика вызывается редко,
# зато ключ можно поменять без рестарта и нет расхождений между воркерами.
expected = _load_or_create_key(
DIAG_KEY_PATH, lambda: secrets.token_urlsafe(18).encode()
).decode()
if not hmac.compare_digest(key or "", expected):
raise HTTPException(status_code=403, detail="Неверный ключ диагностики")
# state -> {provider, verifier, at}; живёт минуты, переживать рестарт не обязано
_oauth_pending: Dict[str, Any] = {}
def _kv_set(key: str, value: str) -> None:
with db() as conn:
conn.execute("INSERT INTO kv (k, v) VALUES (?, ?) "
"ON CONFLICT (k) DO UPDATE SET v = EXCLUDED.v", (key, value))
def _kv_get(key: str) -> Optional[str]:
with db() as conn:
row = conn.execute("SELECT v FROM kv WHERE k = ?", (key,)).fetchone()
return row["v"] if row else None
def _html(body: str) -> "HTMLResponse":
return HTMLResponse(
"<!doctype html><meta charset='utf-8'>"
"<meta name='viewport' content='width=device-width,initial-scale=1'>"
"<style>body{background:#07090f;color:#e8edf5;font:14px/1.6 -apple-system,"
"Segoe UI,Roboto,sans-serif;padding:28px;max-width:760px;margin:auto}"
"code,pre{background:#111827;padding:2px 6px;border-radius:6px;"
"font:12px/1.5 ui-monospace,Menlo,monospace;white-space:pre-wrap;"
"word-break:break-all;display:block;margin:10px 0;padding:12px}"
"b{color:#4d9fff}.ok{color:#00e070}.bad{color:#ff5a4d}"
"a{color:#4d9fff}</style>" + body
)
@app.get("/sc/connect")
async def sc_connect(key: str = ""):
"""Шаг 1: уводим на страницу авторизации SoundCloud."""
_check_diag(key)
verifier = secrets.token_urlsafe(64)[:96]
challenge = _b64e(hashlib.sha256(verifier.encode()).digest())
state = secrets.token_urlsafe(16)
_oauth_pending[state] = {"provider": "soundcloud", "verifier": verifier,
"at": time.time()}
# чистим протухшие state, чтобы словарь не рос
for st in [s for s, v in _oauth_pending.items() if time.time() - v["at"] > 900]:
_oauth_pending.pop(st, None)
params = {
"client_id": SOUNDCLOUD_CLIENT_ID,
"redirect_uri": SC_REDIRECT_URI,
"response_type": "code",
"code_challenge": challenge,
"code_challenge_method": "S256",
"state": state,
}
return RedirectResponse(url=str(httpx.URL(SC_AUTHORIZE_URL, params=params)),
status_code=302)
async def _sc_exchange(code: str, verifier: Optional[str]) -> Dict[str, Any]:
"""Меняем код на токен. Пробуем оба известных token-эндпоинта и с PKCE, и без."""
attempts = []
for token_url in SC_TOKEN_URLS:
for use_pkce in ([True, False] if verifier else [False]):
data = {
"grant_type": "authorization_code",
"client_id": SOUNDCLOUD_CLIENT_ID,
"client_secret": SOUNDCLOUD_CLIENT_SECRET,
"redirect_uri": SC_REDIRECT_URI,
"code": code,
}
if use_pkce:
data["code_verifier"] = verifier
async with httpx.AsyncClient(follow_redirects=True) as client:
try:
r = await client.post(
token_url,
headers={"Content-Type": "application/x-www-form-urlencoded",
"accept": "application/json; charset=utf-8"},
data=data,
timeout=15.0,
)
except Exception as exc:
attempts.append(f"{token_url} pkce={use_pkce}: {exc}")
continue
if r.status_code == 200:
body = r.json()
if body.get("access_token"):
body["_endpoint"] = token_url
body["_pkce"] = use_pkce
return body
attempts.append(f"{token_url} pkce={use_pkce}: HTTP {r.status_code} "
f"{r.text[:200]}")
raise HTTPException(status_code=400,
detail="Обмен кода не удался. " + " | ".join(attempts))
@app.get("/callback")
@app.get("/sc/callback")
async def oauth_callback(code: str = "", state: str = "", error: str = "",
error_description: str = ""):
"""
Шаг 2: провайдер вернул код. Адрес возврата общий для Spotify и SoundCloud,
поэтому кто именно вернулся — узнаём из state, который мы сами и выдавали.
"""
if error:
return _html("<h2 class='bad'>Провайдер вернул ошибку</h2>"
f"<code>{error}\n{error_description}</code>")
if not code:
return _html("<h2 class='bad'>Нет параметра code</h2>")
pending = _oauth_pending.pop(state, None)
if pending and pending.get("provider") != "soundcloud":
return _html("<h2 class='bad'>Этот провайдер пока не подключён</h2>"
f"<p>state указывает на <b>{pending.get('provider')}</b>.</p>")
if not pending:
# state не найден: сервис перезапускался или ссылку открыли повторно.
# Код всё равно пробуем обменять — без PKCE.
logger.info("OAuth callback: state %s неизвестен, пробуем без PKCE", state)
verifier = pending["verifier"] if pending else None
try:
tok = await _sc_exchange(code, verifier)
except HTTPException as exc:
return _html(f"<h2 class='bad'>Не удалось обменять код</h2><code>{exc.detail}</code>")
_kv_set("sc_user_token", tok["access_token"])
if tok.get("refresh_token"):
_kv_set("sc_refresh_token", tok["refresh_token"])
scope = tok.get("scope", "(пусто)")
return _html(
"<h2 class='ok'>Токен получен</h2>"
f"<p>Эндпоинт: <b>{tok.get('_endpoint')}</b>, PKCE: <b>{tok.get('_pkce')}</b></p>"
f"<p>Область доступа (scope): <b>{scope}</b></p>"
f"<p>Живёт секунд: <b>{tok.get('expires_in', '?')}</b></p>"
"<p>Теперь откройте <b>/sc/diag?key=…&q=любой запрос</b> — "
"он сравнит, что отдаётся с этим токеном и без него.</p>"
)
@app.get("/sc/diag")
async def sc_diag(key: str = "", q: str = "lofi"):
"""
Шаг 3: главный замер. Один и тот же трек запрашивается двумя токенами —
приложения и пользователя — и сравнивается, что вернул CDN.
"""
_check_diag(key)
user_token = _kv_get("sc_user_token")
if not user_token:
return _html("<h2 class='bad'>Пользовательского токена нет</h2>"
"<p>Сначала пройдите <b>/sc/connect?key=…</b></p>")
async with httpx.AsyncClient(follow_redirects=False) as client:
app_token = await get_soundcloud_token(client)
# ищем трек под пользовательским токеном
r = await client.get(
SOUNDCLOUD_TRACKS_URL,
headers={"Authorization": f"OAuth {user_token}"},
params={"q": q, "limit": 1, "access": "playable"},
timeout=15.0,
)
if r.status_code != 200:
return _html("<h2 class='bad'>Поиск под пользовательским токеном отклонён</h2>"
f"<code>HTTP {r.status_code}\n{r.text[:400]}</code>")
items = r.json() or []
if not items:
return _html(f"<h2 class='bad'>По запросу «{q}» ничего не нашлось</h2>")
track = items[0]
tid = track.get("id")
rows = []
for label, tok in (("токен приложения", app_token), ("токен пользователя", user_token)):
try:
sr = await client.get(
f"{SOUNDCLOUD_TRACKS_URL}/{tid}/stream",
headers={"Authorization": f"OAuth {tok}"},
timeout=15.0,
)
loc = sr.headers.get("location", "")
host = httpx.URL(loc).host if loc else "(нет редиректа)"
is_preview = "preview" in loc
size = ""
if loc:
async with httpx.AsyncClient(follow_redirects=True) as c2:
hr = await c2.get(loc, headers={"Range": "bytes=0-1"}, timeout=15.0)
size = hr.headers.get("content-range", "") or \
hr.headers.get("content-length", "")
rows.append((label, sr.status_code, host, is_preview, size, loc[:160]))
except Exception as exc:
rows.append((label, "ошибка", str(exc)[:120], None, "", ""))
dur = round((track.get("duration") or 0) / 1000)
html = [f"<h2>Трек: {track.get('title')}</h2>",
f"<p>Длительность по метаданным: <b>{dur} сек</b></p>"]
for label, status, host, is_preview, size, loc in rows:
verdict = ("<span class='bad'>ФРАГМЕНТ 30 сек</span>" if is_preview
else "<span class='ok'>ПОЛНЫЙ ТРЕК</span>" if is_preview is False
else "?")
html.append(f"<h3>{label}</h3><p>HTTP {status} · хост <b>{host}</b> · {verdict}</p>")
if size:
html.append(f"<p>Размер по заголовку: <code>{size}</code></p>")
if loc:
html.append(f"<code>{loc}</code>")
html.append("<p>Если у пользовательского токена хост не содержит "
"<b>preview</b> — ограничение снимается, и это меняет всё.</p>")
return _html("".join(html))
# ==============================================================================
# ВЕБХУК ЮKASSA
# ==============================================================================
class YooKassaNotification(BaseModel):
event: str
type: str = ""
object: dict = {}
@app.post("/webhook")
async def yookassa_webhook(notification: YooKassaNotification) -> Dict[str, Any]:
"""
ВНИМАНИЕ: подпись уведомления здесь не проверяется, поэтому премиум по нему
не выдаётся — только логирование. Перед приёмом реальных платежей нужно
сверять IP-диапазоны ЮKassa и запрашивать статус платежа через их API.
"""
if notification.event != "payment.succeeded":
return {"status": "ignored"}
obj = notification.object or {}
logger.info(
"[YooKassa] Платёж %s на %s %s (metadata=%s) — премиум НЕ выдан, нужна проверка подписи",
obj.get("id"),
(obj.get("amount") or {}).get("value"),
(obj.get("amount") or {}).get("currency", "RUB"),
obj.get("metadata"),
)
return {"status": "logged", "premium_granted": False}
# ==============================================================================
# СОВМЕСТИМОСТЬ СО СТАРЫМИ СБОРКАМИ APK
# ==============================================================================
@app.get("/spotify/search")
async def legacy_spotify_search(query: str = "", username: str = "") -> Dict[str, Any]:
return await search(query=query, limit=25)
@app.get("/soundcloud/search")
async def legacy_soundcloud_search(query: str = "", username: str = "") -> Dict[str, Any]:
return await search(query=query, limit=25, source="soundcloud")
@app.get("/current-track")
async def legacy_current_track(username: str = "") -> Dict[str, Any]:
return {"track": ""}
if __name__ == "__main__":
uvicorn.run(app, host="127.0.0.1", port=8000, log_level="info")