本次提交包含两批改动(7月2日遗留未入库 + 本次新功能),分述如下:
【补登记:7月2日已上线但未提交的功能】
- 分片上传:chunk_upload_controller/service/dao + UploadSession model/schema,
支持 4MiB 分片、乱序、断点续传、去重、幂等 complete;后台 reaper 清理过期会话。
- 反向隧道:tunnel_controller/service/dao + TunnelSession model/schema,
SSH remote forward 经 /api/userPort/{userName} 反代到 user 本地服务。
- 上传页:views/upload_html.py(拖拽/多文件/分片/断点续传 UI)。
- config.py:StorageConfig.chunk_session_dir/ttl、TunnelConfig;
requirements.txt 加 httpx;start.sh 清理 .work/ 残留;
schema.sql 加 upload_session/tunnel_session 表;sftp_server 承载隧道转发。
【本次新功能】
- 文件浏览页:GET /files(Basic Auth 同 docs)+ /api/admin/files(list/get/download/DELETE)。
硬删除(DB 行 + 磁盘文件),删除后列表不再显示。前端 static/file_browser.*。
- 共享白板:GET /whiteboard/{id}(公开,不存在则新建)+ WS /ws/whiteboard/{id}。
MySQL 持久化(whiteboard 表),Canvas 实时同步,心跳 3s/5 次失活移除,
清空/复制按钮,移动端兼容。WhiteboardHub 管理 {board_id: set[Connection]},
disconnect 幂等 + 空 set 清理防泄漏,broadcast 失败连接自动移除。
- 白板管理页:GET /whiteboard-admin(Basic Auth)+ /api/admin/whiteboards(list/DELETE)。
删除时 hub.close_board 踢出在线连接。
- 清理:合并 UploadService.get_out_with_disk_path(下载/删除复用,消除重复 DB 读),
移除无用 resolve_disk_path。
- config.py:WhiteboardConfig(heartbeat/threshold/board_id 长度/list_limit);
schema.sql 加 whiteboard 表;README 补新接口与心跳/内存说明。
- 验证:tests/manual_whiteboard_hub.py / _ws.py / _kick.py 全部通过。
64 lines
2.1 KiB
Python
64 lines
2.1 KiB
Python
"""TunnelSession 的 DAO。"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from datetime import datetime
|
|
|
|
from sqlalchemy import select, update
|
|
from sqlalchemy.orm import Session
|
|
|
|
from ..models.tunnel_session import TunnelSession
|
|
|
|
|
|
class TunnelSessionDAO:
|
|
def __init__(self, db: Session) -> None:
|
|
self.db = db
|
|
|
|
def create(self, session: TunnelSession) -> TunnelSession:
|
|
self.db.add(session)
|
|
self.db.commit()
|
|
self.db.refresh(session)
|
|
return session
|
|
|
|
def get_active_by_user(self, user_name: str) -> TunnelSession | None:
|
|
"""返回该 user 当前活跃的隧道会话(至多一条)。"""
|
|
stmt = (
|
|
select(TunnelSession)
|
|
.where(TunnelSession.user_name == user_name)
|
|
.where(TunnelSession.status == "active")
|
|
.order_by(TunnelSession.started_at.desc())
|
|
.limit(1)
|
|
)
|
|
return self.db.scalars(stmt).first()
|
|
|
|
def get_active_by_port(self, tunnel_port: int) -> TunnelSession | None:
|
|
stmt = (
|
|
select(TunnelSession)
|
|
.where(TunnelSession.tunnel_port == tunnel_port)
|
|
.where(TunnelSession.status == "active")
|
|
.limit(1)
|
|
)
|
|
return self.db.scalars(stmt).first()
|
|
|
|
def list_active(self) -> list[TunnelSession]:
|
|
stmt = select(TunnelSession).where(TunnelSession.status == "active")
|
|
return list(self.db.scalars(stmt).all())
|
|
|
|
def close(self, session: TunnelSession) -> None:
|
|
"""标记会话结束。"""
|
|
session.ended_at = datetime.now()
|
|
session.status = "closed"
|
|
self.db.commit()
|
|
|
|
def close_active_by_user(self, user_name: str) -> int:
|
|
"""关闭该 user 所有 active 会话(断开清理用),返回关闭条数。"""
|
|
stmt = (
|
|
update(TunnelSession)
|
|
.where(TunnelSession.user_name == user_name)
|
|
.where(TunnelSession.status == "active")
|
|
.values(status="closed", ended_at=datetime.now())
|
|
)
|
|
result = self.db.execute(stmt)
|
|
self.db.commit()
|
|
return result.rowcount or 0
|