archiver / bot /app.py
KHALED-ALIM
feat: smart disk cleanup system
09418c2
Raw
History Blame Contribute Delete
8.84 kB
"""Application entry point — FastAPI + Pyrogram in one process."""
from __future__ import annotations
import asyncio
import secrets
from contextlib import asynccontextmanager
from pathlib import Path
from typing import Annotated, AsyncGenerator
import httpx
from fastapi import Depends, FastAPI, HTTPException, Request
from fastapi.responses import HTMLResponse
from fastapi.security import APIKeyHeader
from fastapi.templating import Jinja2Templates
from fastapi.middleware.trustedhost import TrustedHostMiddleware
from loguru import logger
from pyrogram import Client
from starlette.responses import Response as StarletteResponse
from bot.client import create_client
from config import settings
import secrets
import mimetypes
from pathlib import Path
from fastapi.responses import FileResponse
from database.connection import db
from services.queue import get_queue
from services.streamtape_service import _get_accounts, StreamTapeError
import bot.handlers.message_handler # noqa: F401
import bot.handlers.callback_handler # noqa: F401
_client: Client | None = None
_templates_dir = Path(__file__).resolve().parent.parent / "templates"
templates = Jinja2Templates(directory=str(_templates_dir))
_api_key_header = APIKeyHeader(name="X-API-Key", auto_error=False)
async def verify_api_key(
x_api_key: Annotated[str | None, Depends(_api_key_header)],
) -> None:
"""Validate the X-API-Key header using constant-time comparison."""
if x_api_key is None:
raise HTTPException(status_code=401, detail="Missing API key")
if not secrets.compare_digest(x_api_key, settings.api_key):
raise HTTPException(status_code=403, detail="Invalid API key")
def get_client() -> Client:
global _client
if _client is None:
_client = create_client()
return _client
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncGenerator:
"""Start Pyrogram client on boot, stop on shutdown."""
logger.info("🚀 Starting Upload Bot...")
client = get_client()
await db.ensure_tables()
await client.start()
logger.info("✅ Pyrogram client started.")
try:
await client.set_bot_commands([
{"command": "start", "description": "🔄 بدء / إعادة تعيين البوت"},
])
except Exception as exc:
logger.warning("Failed to set bot commands: {}", exc)
# ── Background disk cleanup task ─────────────────────────────────
async def _periodic_cleanup():
"""ينظف ملفات الرفع القديمة كل 15 دقيقة."""
from utils.file_utils import smart_disk_cleanup
from config import settings
while True:
try:
await asyncio.sleep(900) # 15 دقيقة
result = await asyncio.to_thread(
smart_disk_cleanup,
settings.upload_dir,
min_free_mb=200,
max_age_minutes=120,
)
if result["deleted_files"] > 0 or result["deleted_dirs"] > 0:
logger.info("Periodic cleanup: {msg}", msg=result["message"])
except asyncio.CancelledError:
break
except Exception as exc:
logger.warning("Cleanup task error: {}", exc)
cleanup_task = asyncio.create_task(_periodic_cleanup())
yield
cleanup_task.cancel()
try:
await cleanup_task
except asyncio.CancelledError:
pass
logger.info("🛑 Shutting down...")
if client.is_connected:
await client.stop()
app = FastAPI(
title="Upload Bot",
description="Telegram upload bot with MTProto, FFmpeg, StreamTape, Turso",
version="1.0.0",
lifespan=lifespan,
)
@app.middleware("http")
async def add_security_headers(request: Request, call_next) -> StarletteResponse:
"""Attach security hardening headers to every response."""
response = await call_next(request)
response.headers["X-Content-Type-Options"] = "nosniff"
response.headers["X-Frame-Options"] = "DENY"
response.headers["X-XSS-Protection"] = "1; mode=block"
response.headers["Content-Security-Policy"] = (
"default-src 'self'; "
"script-src 'self' https://unpkg.com https://cdn.tailwindcss.com 'unsafe-inline'; "
"style-src 'self' 'unsafe-inline'; "
"img-src 'self' data:; "
"connect-src 'self'; "
"frame-ancestors 'none';"
)
response.headers["Referrer-Policy"] = "strict-origin-when-cross-origin"
response.headers["Permissions-Policy"] = (
"camera=(), microphone=(), geolocation=()"
)
return response
app.add_middleware(
TrustedHostMiddleware,
allowed_hosts=[
settings.host,
"localhost",
"127.0.0.1",
"::1",
"*.onrender.com",
"*.hf.space",
"khaled-alim-archiver.hf.space",
],
)
@app.get("/")
async def root():
return {"status": "alive", "bot": "running"}
@app.get("/health")
async def health():
client = get_client()
return {
"status": "ok",
"pyrogram_connected": client.is_connected,
}
@app.get("/status", dependencies=[Depends(verify_api_key)])
async def status():
client = get_client()
me = await client.get_me()
return {
"bot": me.username or me.first_name,
"connected": client.is_connected,
}
@app.get("/api/stats", dependencies=[Depends(verify_api_key)])
async def api_stats():
return await db.get_stats()
@app.get("/api/queue-status", dependencies=[Depends(verify_api_key)])
async def api_queue_status():
q = get_queue()
return q.status()
@app.get("/api/streamtape-play/{file_id}")
async def api_streamtape_play(file_id: str):
"""توليد رابط مباشر للتشغيل من معرف ملف StreamTape.
يستخدم dlticket + dl API لتوليد رابط tapecontent.net مؤقت.
الرابط صالح لدقائق معدودة — كرر الاستدعاء عند الحاجة.
"""
from services.streamtape_service import generate_direct_url, _get_accounts, extract_streamtape_file_id
# استخراج file_id من الرابط إذا كان المستخدم أرسل رابطاً كاملاً
parsed_id = extract_streamtape_file_id(file_id)
if not parsed_id:
return {"error": f"Invalid StreamTape file ID: {file_id}"}
file_id = parsed_id
accounts = _get_accounts()
results: list[dict] = []
errors: list[str] = []
for acc in accounts:
try:
url = await generate_direct_url(file_id, acc.username, acc.password)
results.append({
"account": acc.index,
"url": url,
"is_direct": "tapecontent.net" in url,
})
except Exception as exc:
errors.append(f"Account {acc.index}: {exc}")
return {
"file_id": file_id,
"results": results,
"errors": errors if errors else None,
}
@app.get("/api/diagnose-streamtape", dependencies=[Depends(verify_api_key)])
async def api_diagnose_streamtape():
"""Test StreamTape API connectivity from this HF Space."""
import httpx
from services.streamtape_service import _get_accounts
# اختبار حساب واحد فقط (الأول) — الاختبارات المتعددة تسبب بطئاً
results: list[dict] = []
try:
acc = _get_accounts()[0]
async with httpx.AsyncClient(timeout=5) as client:
resp = await client.get(
"https://api.streamtape.com/file/info",
params={"login": acc.username, "key": acc.password, "file": "test"},
)
data = resp.json()
results.append({
"account": acc.index,
"api_status": data.get("status"),
"reachable": data.get("status") == 200,
})
except Exception as exc:
results.append({"account": 1, "error": str(exc)[:100], "reachable": False})
# اختبار CDN
cdn_reachable = False
try:
async with httpx.AsyncClient(timeout=3) as c:
r = await c.get("https://tapecontent.net")
cdn_reachable = r.status_code < 500
except Exception:
cdn_reachable = False
return {
"streamtape_diagnostic": results,
"cdn_reachable": cdn_reachable,
"note": "API ✅ | CDN ❌ محجوب → البوت يحول لـ FTP تلقائياً"
}
@app.get("/dashboard", response_class=HTMLResponse)
async def dashboard(request: Request):
stats = await db.get_stats()
q = get_queue()
return templates.TemplateResponse(
"dashboard.html",
{"request": request, "stats": stats, "queue": q.status()},
)