Skip to content

Commit 3de2615

Browse files
committed
feat: add video room feature with API, database model, and UI integration
1 parent d5a24c6 commit 3de2615

19 files changed

Lines changed: 692 additions & 40 deletions

File tree

api/routers.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
from domain.processors import api as processors
1212
from domain.share import api as share
1313
from domain.tasks import api as tasks
14+
from domain.video_rooms import api as video_rooms
1415
from domain.ai import api as ai
1516
from domain.agent import api as agent
1617
from domain.virtual_fs import api as virtual_fs
@@ -32,6 +33,7 @@ def include_routers(app: FastAPI):
3233
app.include_router(config.router)
3334
app.include_router(processors.router)
3435
app.include_router(tasks.router)
36+
app.include_router(video_rooms.router)
3537
app.include_router(share.router)
3638
app.include_router(share.public_router)
3739
app.include_router(backup.router)

domain/adapters/providers/quark.py

Lines changed: 57 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -360,25 +360,14 @@ async def stream_file(self, root: str, rel: str, range_header: str | None):
360360
if tr:
361361
url = tr
362362
dl_headers = self._download_headers()
363-
364-
# 预获取大小/是否支持范围
365-
total_size: Optional[int] = None
366-
async with httpx.AsyncClient(timeout=self._timeout, follow_redirects=True) as client:
367-
try:
368-
head_resp = await client.head(url, headers=dl_headers)
369-
if head_resp.status_code == 200:
370-
cl = head_resp.headers.get("Content-Length")
371-
if cl and cl.isdigit():
372-
total_size = int(cl)
373-
except Exception:
374-
pass
363+
file_size = int(it.get("size") or 0)
375364

376365
mime, _ = mimetypes.guess_type(rel)
377366
content_type = mime or "application/octet-stream"
378367

379368
# 解析 Range
380369
start = 0
381-
end: Optional[int] = None
370+
end: Optional[int] = file_size - 1 if file_size > 0 else None
382371
status_code = 200
383372
if range_header and range_header.startswith("bytes="):
384373
status_code = 206
@@ -388,35 +377,65 @@ async def stream_file(self, root: str, rel: str, range_header: str | None):
388377
start = int(s)
389378
if e.strip():
390379
end = int(e)
391-
392-
if total_size is not None and end is None and status_code == 206:
393-
end = total_size - 1
394-
if end is not None and total_size is not None and end >= total_size:
395-
end = total_size - 1
396-
if total_size is not None and start >= total_size:
380+
elif file_size > 0:
381+
end = file_size - 1
382+
if file_size > 0:
383+
if start >= file_size:
384+
raise HTTPException(416, detail="Requested Range Not Satisfiable")
385+
if end is None or end >= file_size:
386+
end = file_size - 1
387+
if start > end:
388+
raise HTTPException(416, detail="Requested Range Not Satisfiable")
389+
headers = dict(dl_headers)
390+
if status_code == 206:
391+
headers["Range"] = f"bytes={start}-" if end is None else f"bytes={start}-{end}"
392+
393+
client = httpx.AsyncClient(timeout=None, follow_redirects=True)
394+
req = client.build_request("GET", url, headers=headers)
395+
resp = await client.send(req, stream=True)
396+
if resp.status_code == 404:
397+
await resp.aclose()
398+
await client.aclose()
399+
raise FileNotFoundError(rel)
400+
if resp.status_code == 416:
401+
await resp.aclose()
402+
await client.aclose()
397403
raise HTTPException(416, detail="Requested Range Not Satisfiable")
404+
try:
405+
resp.raise_for_status()
406+
except Exception:
407+
await resp.aclose()
408+
await client.aclose()
409+
raise
398410

399-
resp_headers: Dict[str, str] = {"Accept-Ranges": "bytes", "Content-Type": content_type}
400-
if status_code == 206 and total_size is not None and end is not None:
401-
resp_headers["Content-Range"] = f"bytes {start}-{end}/{total_size}"
402-
resp_headers["Content-Length"] = str(end - start + 1)
403-
elif total_size is not None:
404-
resp_headers["Content-Length"] = str(total_size)
411+
resp_headers: Dict[str, str] = {
412+
"Accept-Ranges": resp.headers.get("Accept-Ranges", "bytes"),
413+
"Content-Type": resp.headers.get("Content-Type", content_type),
414+
}
415+
content_range = resp.headers.get("Content-Range")
416+
content_length = resp.headers.get("Content-Length")
417+
if content_range:
418+
resp_headers["Content-Range"] = content_range
419+
elif status_code == 206 and file_size > 0 and end is not None:
420+
resp_headers["Content-Range"] = f"bytes {start}-{end}/{file_size}"
421+
if content_length:
422+
resp_headers["Content-Length"] = content_length
423+
elif file_size > 0:
424+
if status_code == 206 and end is not None:
425+
resp_headers["Content-Length"] = str(end - start + 1)
426+
elif resp.status_code == 200:
427+
resp_headers["Content-Length"] = str(file_size)
405428

406429
async def iterator():
407-
headers = dict(dl_headers)
408-
if status_code == 206 and end is not None:
409-
headers["Range"] = f"bytes={start}-{end}"
410-
async with httpx.AsyncClient(timeout=None, follow_redirects=True) as client:
411-
async with client.stream("GET", url, headers=headers) as resp:
412-
if resp.status_code in (404, 416):
413-
await resp.aclose()
414-
raise HTTPException(resp.status_code, detail="Upstream not available")
415-
async for chunk in resp.aiter_bytes():
416-
if chunk:
417-
yield chunk
418-
419-
return StreamingResponse(iterator(), status_code=status_code, headers=resp_headers, media_type=content_type)
430+
try:
431+
async for chunk in resp.aiter_bytes():
432+
if chunk:
433+
yield chunk
434+
finally:
435+
await resp.aclose()
436+
await client.aclose()
437+
438+
return StreamingResponse(iterator(), status_code=resp.status_code, headers=resp_headers, media_type=content_type)
420439

421440
# -----------------
422441
# 上传(大文件分片)

domain/video_rooms/__init__.py

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,9 @@
1+
from .service import VideoRoomService
2+
from .types import VideoRoomCreate, VideoRoomInfo, VideoRoomState
3+
4+
__all__ = [
5+
"VideoRoomService",
6+
"VideoRoomCreate",
7+
"VideoRoomInfo",
8+
"VideoRoomState",
9+
]

domain/video_rooms/api.py

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,72 @@
1+
from typing import Annotated
2+
3+
from fastapi import APIRouter, Depends, Request, WebSocket, WebSocketDisconnect
4+
5+
from api.response import success
6+
from domain.audit import AuditAction, audit
7+
from domain.auth import User, get_current_active_user
8+
from domain.permission import require_path_permission
9+
from domain.permission.types import PathAction
10+
from models.database import UserAccount
11+
from .service import VideoRoomService
12+
from .types import VideoRoomCreate, VideoRoomInfo
13+
from .ws import video_room_ws_manager
14+
15+
router = APIRouter(prefix="/api/video-rooms", tags=["Video Rooms"])
16+
17+
18+
@router.post("", response_model=VideoRoomInfo)
19+
@audit(action=AuditAction.SHARE, description="创建视频房", body_fields=["name", "path"])
20+
@require_path_permission(PathAction.SHARE, "payload.path")
21+
async def create_video_room(
22+
request: Request,
23+
payload: VideoRoomCreate,
24+
current_user: Annotated[User, Depends(get_current_active_user)],
25+
):
26+
user_account = await UserAccount.get(id=current_user.id)
27+
room = await VideoRoomService.create_room(
28+
user=user_account,
29+
name=payload.name,
30+
path=payload.path,
31+
)
32+
return VideoRoomInfo.from_orm(room, VideoRoomService.get_effective_state(room))
33+
34+
35+
@router.get("/{token}", response_model=VideoRoomInfo)
36+
@audit(action=AuditAction.SHARE, description="获取视频房信息")
37+
async def get_video_room(request: Request, token: str):
38+
room = await VideoRoomService.get_room_by_token(token)
39+
return VideoRoomInfo.from_orm(room, VideoRoomService.get_effective_state(room))
40+
41+
42+
@router.get("/{token}/stream")
43+
@audit(action=AuditAction.DOWNLOAD, description="播放视频房文件")
44+
async def stream_video_room(token: str, request: Request):
45+
return await VideoRoomService.stream_room_file(token, request.headers.get("Range"))
46+
47+
48+
@router.websocket("/{token}/ws")
49+
async def video_room_ws(websocket: WebSocket, token: str):
50+
room = await VideoRoomService.get_room_by_token(token)
51+
await video_room_ws_manager.connect(token, websocket)
52+
try:
53+
state = VideoRoomService.get_effective_state(room)
54+
await websocket.send_json({"type": "state", "state": state.model_dump()})
55+
while True:
56+
data = await websocket.receive_json()
57+
if data.get("type") != "state":
58+
continue
59+
state = await VideoRoomService.update_state(
60+
room,
61+
current_time=float(data.get("current_time") or 0),
62+
paused=bool(data.get("paused")),
63+
)
64+
await video_room_ws_manager.broadcast(
65+
token,
66+
{"type": "state", "state": state.model_dump()},
67+
exclude=websocket,
68+
)
69+
except WebSocketDisconnect:
70+
pass
71+
finally:
72+
video_room_ws_manager.disconnect(token, websocket)

domain/video_rooms/service.py

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
import secrets
2+
from datetime import datetime, timezone
3+
from urllib.parse import quote
4+
5+
from fastapi import HTTPException
6+
from fastapi.responses import Response
7+
8+
from domain.virtual_fs import VirtualFSService
9+
from models.database import UserAccount, VideoRoom
10+
from .types import VideoRoomState
11+
12+
13+
VIDEO_EXTENSIONS = {".mp4", ".webm", ".ogg", ".m4v", ".mov", ".mkv", ".avi", ".flv"}
14+
15+
16+
class VideoRoomService:
17+
@classmethod
18+
def _is_video_path(cls, path: str) -> bool:
19+
lower = path.lower()
20+
return any(lower.endswith(ext) for ext in VIDEO_EXTENSIONS)
21+
22+
@classmethod
23+
async def create_room(cls, user: UserAccount, name: str, path: str) -> VideoRoom:
24+
if not path or path == "/" or ".." in path.split("/"):
25+
raise HTTPException(status_code=400, detail="无效的视频路径")
26+
if not cls._is_video_path(path):
27+
raise HTTPException(status_code=400, detail="仅支持视频文件创建视频房")
28+
29+
stat = await VirtualFSService.stat_file(path)
30+
if stat.get("is_dir"):
31+
raise HTTPException(status_code=400, detail="目录不能创建视频房")
32+
33+
token = secrets.token_urlsafe(16)
34+
return await VideoRoom.create(
35+
token=token,
36+
name=name or path.rsplit("/", 1)[-1],
37+
path=path,
38+
user=user,
39+
state_updated_at=datetime.now(timezone.utc),
40+
)
41+
42+
@classmethod
43+
async def get_room_by_token(cls, token: str) -> VideoRoom:
44+
room = await VideoRoom.get_or_none(token=token).prefetch_related("user")
45+
if not room:
46+
raise HTTPException(status_code=404, detail="视频房不存在")
47+
return room
48+
49+
@classmethod
50+
def get_effective_state(cls, room: VideoRoom) -> VideoRoomState:
51+
current_time = float(room.current_time or 0)
52+
updated_at = room.state_updated_at
53+
if not room.paused and updated_at:
54+
if updated_at.tzinfo is None:
55+
updated_at = updated_at.replace(tzinfo=timezone.utc)
56+
current_time += max(0, (datetime.now(timezone.utc) - updated_at).total_seconds())
57+
return VideoRoomState(
58+
current_time=max(0, current_time),
59+
paused=bool(room.paused),
60+
updated_at=updated_at.isoformat() if updated_at else None,
61+
)
62+
63+
@classmethod
64+
async def update_state(cls, room: VideoRoom, current_time: float, paused: bool) -> VideoRoomState:
65+
now = datetime.now(timezone.utc)
66+
room.current_time = max(0, float(current_time or 0))
67+
room.paused = bool(paused)
68+
room.state_updated_at = now
69+
await room.save(update_fields=["current_time", "paused", "state_updated_at"])
70+
return VideoRoomState(
71+
current_time=room.current_time,
72+
paused=room.paused,
73+
updated_at=now.isoformat(),
74+
)
75+
76+
@classmethod
77+
async def stream_room_file(cls, token: str, range_header: str | None) -> Response:
78+
room = await cls.get_room_by_token(token)
79+
response = await VirtualFSService.stream_file(room.path, range_header)
80+
filename = room.path.rsplit("/", 1)[-1]
81+
response.headers["Content-Disposition"] = f"inline; filename*=UTF-8''{quote(filename)}"
82+
return response

domain/video_rooms/types.py

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,39 @@
1+
from pydantic import BaseModel
2+
3+
from models.database import VideoRoom
4+
5+
6+
class VideoRoomCreate(BaseModel):
7+
name: str
8+
path: str
9+
10+
11+
class VideoRoomState(BaseModel):
12+
current_time: float
13+
paused: bool
14+
updated_at: str | None = None
15+
16+
17+
class VideoRoomInfo(BaseModel):
18+
id: int
19+
token: str
20+
name: str
21+
path: str
22+
created_at: str
23+
state: VideoRoomState
24+
25+
@classmethod
26+
def from_orm(cls, obj: VideoRoom, state: VideoRoomState | None = None):
27+
return cls(
28+
id=obj.id,
29+
token=obj.token,
30+
name=obj.name,
31+
path=obj.path,
32+
created_at=obj.created_at.isoformat(),
33+
state=state
34+
or VideoRoomState(
35+
current_time=obj.current_time,
36+
paused=obj.paused,
37+
updated_at=obj.state_updated_at.isoformat() if obj.state_updated_at else None,
38+
),
39+
)

domain/video_rooms/ws.py

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
from fastapi import WebSocket
2+
3+
4+
class VideoRoomWebSocketManager:
5+
def __init__(self):
6+
self.rooms: dict[str, set[WebSocket]] = {}
7+
8+
async def connect(self, token: str, websocket: WebSocket):
9+
await websocket.accept()
10+
self.rooms.setdefault(token, set()).add(websocket)
11+
12+
def disconnect(self, token: str, websocket: WebSocket):
13+
sockets = self.rooms.get(token)
14+
if not sockets:
15+
return
16+
sockets.discard(websocket)
17+
if not sockets:
18+
self.rooms.pop(token, None)
19+
20+
async def broadcast(self, token: str, message: dict, exclude: WebSocket | None = None):
21+
sockets = list(self.rooms.get(token, set()))
22+
for socket in sockets:
23+
if socket is exclude:
24+
continue
25+
try:
26+
await socket.send_json(message)
27+
except Exception:
28+
self.disconnect(token, socket)
29+
30+
31+
video_room_ws_manager = VideoRoomWebSocketManager()

models/database.py

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -234,6 +234,23 @@ class Meta:
234234
table = "share_links"
235235

236236

237+
class VideoRoom(Model):
238+
id = fields.IntField(pk=True)
239+
token = fields.CharField(max_length=100, unique=True, index=True)
240+
name = fields.CharField(max_length=255)
241+
path = fields.CharField(max_length=4096)
242+
user: fields.ForeignKeyRelation[UserAccount] = fields.ForeignKeyField(
243+
"models.UserAccount", related_name="video_rooms", on_delete=fields.CASCADE
244+
)
245+
current_time = fields.FloatField(default=0)
246+
paused = fields.BooleanField(default=True)
247+
state_updated_at = fields.DatetimeField(null=True)
248+
created_at = fields.DatetimeField(auto_now_add=True)
249+
250+
class Meta:
251+
table = "video_rooms"
252+
253+
237254
class RecentFile(Model):
238255
id = fields.IntField(pk=True)
239256
user: fields.ForeignKeyRelation[UserAccount] = fields.ForeignKeyField(

0 commit comments

Comments
 (0)