服务端 deb 化: debian/ 全套 (control/rules/postinst/changelog), 双包 all源码 + amd64编译(Nuitka onefile sz-server 22MB), postinst 自动装依赖 (aria2+p7zip+venv fastapi/uvicorn)
This commit is contained in:
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Executable
BIN
Binary file not shown.
Binary file not shown.
Executable
BIN
Binary file not shown.
+8
@@ -0,0 +1,8 @@
|
||||
7z-encrypt-server (1.3.0) stable; urgency=medium
|
||||
|
||||
* 服务端 deb 化: 源码版 (all, 手机/容器) + 编译版 (amd64, Nuitka onefile)
|
||||
* postinst 自动安装依赖: aria2 + p7zip-full + venv (fastapi/uvicorn/cryptography/tqdm)
|
||||
* 目录管理: dirs 表 + CRUD 端点 + init dir_id + move/clone
|
||||
* 账户严格隔离 + pull-status 实时进度
|
||||
|
||||
-- edgevoid <edgevoid@users.noreply.gitee.com> Mon, 10 Aug 2026 21:50:00 +0800
|
||||
+8
@@ -0,0 +1,8 @@
|
||||
7z-encrypt-server (1.3.0) stable; urgency=medium
|
||||
|
||||
* 服务端 deb 化: 源码版 (all, 手机/容器) + 编译版 (amd64, Nuitka onefile)
|
||||
* postinst 自动安装依赖: aria2 + p7zip-full + venv (fastapi/uvicorn/cryptography/tqdm)
|
||||
* 目录管理: dirs 表 + CRUD 端点 + init dir_id + move/clone
|
||||
* 账户严格隔离 + pull-status 实时进度
|
||||
|
||||
-- edgevoid <edgevoid@users.noreply.gitee.com> Mon, 10 Aug 2026 21:50:00 +0800
|
||||
+3
@@ -0,0 +1,3 @@
|
||||
shlibs:Depends=libc6 (>= 2.34)
|
||||
misc:Depends=
|
||||
misc:Pre-Depends=
|
||||
+13
@@ -0,0 +1,13 @@
|
||||
Package: 7z-encrypt-server-bin
|
||||
Source: 7z-encrypt-server
|
||||
Version: 1.3.0
|
||||
Architecture: amd64
|
||||
Maintainer: edgevoid <edgevoid@users.noreply.gitee.com>
|
||||
Installed-Size: 21948
|
||||
Depends: aria2, p7zip-full
|
||||
Section: utils
|
||||
Priority: optional
|
||||
Description: 端到端加密文件传输服务端 (x86_64 原生编译版)
|
||||
Nuitka 编译的原生二进制 (Python->C->gcc, onefile 自包含, 免 Python 依赖)。
|
||||
功能与源码版一致: sz-server。
|
||||
装完即用: 系统依赖 (aria2/p7zip) postinst 自动安装。
|
||||
+3
@@ -0,0 +1,3 @@
|
||||
907e46fc1ee535f5997471beafe88b7c usr/bin/sz-server
|
||||
ec86c28685a1356540729bd74b9fbad7 usr/share/doc/7z-encrypt-server-bin/changelog.gz
|
||||
5f9ffdb75dd3b928d6decf94d22a0605 usr/share/doc/7z-encrypt-server-bin/copyright
|
||||
BIN
Binary file not shown.
BIN
Binary file not shown.
@@ -0,0 +1,7 @@
|
||||
Format: https://www.debian.org/doc/packaging-manuals/copyright-format/1.0/
|
||||
Upstream-Name: 7z-encrypt-server
|
||||
Source: https://gitee.com/edgevoid/7z-encrypt-server
|
||||
|
||||
Files: *
|
||||
Copyright: 2026 edgevoid
|
||||
License: MIT
|
||||
Vendored
+2
@@ -0,0 +1,2 @@
|
||||
misc:Depends=
|
||||
misc:Pre-Depends=
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
Package: 7z-encrypt-server
|
||||
Version: 1.3.0
|
||||
Architecture: all
|
||||
Maintainer: edgevoid <edgevoid@users.noreply.gitee.com>
|
||||
Installed-Size: 59
|
||||
Depends: python3 (>= 3.10), python3-venv, aria2, p7zip-full
|
||||
Section: utils
|
||||
Priority: optional
|
||||
Description: 端到端加密文件传输服务端 (源码版, 手机/容器用)
|
||||
零知识 + 零合并的服务端: 只存卷/传卷, 不合并不解密, 不持有密钥。
|
||||
账户严格隔离 (一个账户一个空间), 目录管理 (dirs/move/clone),
|
||||
aria2 反向拉加速 + pull-status 实时进度。
|
||||
.
|
||||
包含: sz-server (入口) + 配套模块
|
||||
端口 8000 (SZ_PORT), 存储/数据目录可配 (SZ_DATA_DIR)。
|
||||
装完即用: 依赖 (fastapi/uvicorn/cryptography/tqdm) postinst 自动安装。
|
||||
+11
@@ -0,0 +1,11 @@
|
||||
7a46467b852172dd11fde61a9437c57c usr/bin/sz-server
|
||||
f8b20e094d5084238c23a8396abc4a4c usr/share/7z-encrypt-server/api.py
|
||||
249597b1c15cfe6bc74a55cc7cc3f840 usr/share/7z-encrypt-server/auth.py
|
||||
87ff0814c04fa600efa7cb8170609f93 usr/share/7z-encrypt-server/db.py
|
||||
de33280101fdc7841d3def37256f1a1e usr/share/7z-encrypt-server/main.py
|
||||
7f1accde37a9cd37590b517309df5242 usr/share/7z-encrypt-server/receiver.py
|
||||
c13cce2e211dad6da61e76f64ecc965f usr/share/7z-encrypt-server/settings.py
|
||||
eca8daf83ca027fe7d62abd1acfefae7 usr/share/7z-encrypt-server/storage.py
|
||||
beff1bc0f315b31e241bf6abb0bf5258 usr/share/7z-encrypt-server/task_manager.py
|
||||
ec86c28685a1356540729bd74b9fbad7 usr/share/doc/7z-encrypt-server/changelog.gz
|
||||
5f9ffdb75dd3b928d6decf94d22a0605 usr/share/doc/7z-encrypt-server/copyright
|
||||
+18
@@ -0,0 +1,18 @@
|
||||
#!/bin/sh
|
||||
set -e
|
||||
|
||||
case "$1" in
|
||||
configure)
|
||||
echo "正在安装系统依赖 (aria2 + p7zip-full) ..."
|
||||
apt-get install -y aria2 p7zip-full >/dev/null 2>&1 || true
|
||||
P=/usr/share/7z-encrypt-server/.venv
|
||||
if [ ! -x "$P/bin/python" ]; then
|
||||
echo "正在创建运行环境 (venv + fastapi/uvicorn/cryptography/tqdm) ..."
|
||||
python3 -m venv "$P"
|
||||
"$P/bin/pip" install --quiet --upgrade pip || true
|
||||
"$P/bin/pip" install --quiet fastapi "uvicorn[standard]" cryptography tqdm
|
||||
fi
|
||||
;;
|
||||
esac
|
||||
|
||||
exit 0
|
||||
+2
@@ -0,0 +1,2 @@
|
||||
#!/bin/sh
|
||||
exec /usr/share/7z-encrypt-server/.venv/bin/python /usr/share/7z-encrypt-server/main.py "$@"
|
||||
@@ -0,0 +1,452 @@
|
||||
"""
|
||||
API 层 (api)
|
||||
模块: 服务端 / API 层
|
||||
输入: HTTP 请求 (带认证)
|
||||
输出: 响应 / 路由到各模块
|
||||
|
||||
零知识 + 零合并设计: 服务端只存卷/传卷, 不合并不解密, 不持有密钥。
|
||||
端点:
|
||||
POST /api/transfer/init {init json} -> {transfer_id}
|
||||
GET /api/transfer/{id}/chunks -> {received: [n...]}
|
||||
PUT /api/transfer/{id}/chunk/{n} 卷密文 -> {ok} / 409
|
||||
POST /api/transfer/{id}/complete -> {status, file_id} / 409 (秒回, 不合并)
|
||||
GET /api/files -> {files: [...]}
|
||||
GET /api/files/{id}/chunk/{n} -> 单卷密文 (头 X-Enc-Params / X-Chunk-Count)
|
||||
全部端点需 Authorization: Bearer <token> (与客户端预共享)
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import base64
|
||||
import json
|
||||
import logging
|
||||
import shutil
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
|
||||
from fastapi import Depends, FastAPI, Header, HTTPException, Request
|
||||
from fastapi.responses import FileResponse, JSONResponse
|
||||
|
||||
from auth import hash_password, new_token, verify_password
|
||||
|
||||
# 业务日志 (随 uvicorn 输出到 stderr -> /tmp/sz-server.log)
|
||||
log = logging.getLogger("sz")
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s [业务] %(message)s")
|
||||
|
||||
from db import ServerDB
|
||||
from receiver import Receiver, ReceiverError
|
||||
from settings import QUOTA_BYTES, STORAGE_ROOT, TMP_ROOT, TOKEN
|
||||
from storage import Storage
|
||||
from task_manager import TaskManager
|
||||
|
||||
# ---------- 依赖装配 (单例) ----------
|
||||
db = ServerDB()
|
||||
tasks = TaskManager(db)
|
||||
receiver = Receiver(tasks, TMP_ROOT)
|
||||
storage = Storage(tasks, STORAGE_ROOT)
|
||||
|
||||
app = FastAPI(title="7z-encrypt 服务端 (零知识 + 零合并)")
|
||||
|
||||
|
||||
# ---------- 认证 ----------
|
||||
|
||||
def verify_token(authorization: str | None = Header(default=None)) -> str:
|
||||
"""Bearer token -> username; 失败 401.
|
||||
SZ_TOKEN (系统 token) = admin, 兼容旧配置; 注册用户走 users 表 token"""
|
||||
if authorization == f"Bearer {TOKEN}":
|
||||
return "admin"
|
||||
if authorization and authorization.startswith("Bearer "):
|
||||
user = db.get_username_by_token(authorization[7:])
|
||||
if user is not None:
|
||||
return user
|
||||
raise HTTPException(status_code=401, detail="未授权: token 无效或缺失")
|
||||
|
||||
|
||||
def _check_owner(rec: dict, username: str) -> None:
|
||||
"""文件归属校验: 每个账户只操作自己的空间 (admin 只认 admin/无主旧文件)"""
|
||||
if username == "admin":
|
||||
if rec.get("user_id") not in ("admin", None):
|
||||
raise HTTPException(status_code=403, detail="无权访问他人文件")
|
||||
return
|
||||
if rec.get("user_id") != username:
|
||||
raise HTTPException(status_code=403, detail="无权访问他人文件")
|
||||
|
||||
|
||||
# ---------- 账号: 注册 / 登录 / 注销 ----------
|
||||
|
||||
@app.post("/api/auth/register")
|
||||
async def register(request: Request):
|
||||
"""注册 {username, password} -> 自动登录 {username, token}"""
|
||||
try:
|
||||
body = await request.json()
|
||||
username = str(body["username"]).strip()
|
||||
password = str(body["password"])
|
||||
except (KeyError, TypeError, ValueError):
|
||||
return JSONResponse(status_code=400, content={"error": "缺少 username/password"})
|
||||
if not 3 <= len(username) <= 32 or not username.replace("_", "").isalnum():
|
||||
return JSONResponse(status_code=400, content={"error": "用户名需 3-32 位字母数字或下划线"})
|
||||
if not 6 <= len(password) <= 128:
|
||||
return JSONResponse(status_code=400, content={"error": "密码需 6-128 位"})
|
||||
if not db.create_user(username, hash_password(password)):
|
||||
return JSONResponse(status_code=409, content={"error": "用户名已存在"})
|
||||
token = new_token()
|
||||
db.set_token(username, token)
|
||||
log.info("register %s 用户=%s", request.client.host if request.client else "?", username)
|
||||
return {"username": username, "token": token}
|
||||
|
||||
|
||||
@app.post("/api/auth/login")
|
||||
async def login(request: Request):
|
||||
"""登录 {username, password} -> {username, token}"""
|
||||
try:
|
||||
body = await request.json()
|
||||
username = str(body["username"]).strip()
|
||||
password = str(body["password"])
|
||||
except (KeyError, TypeError, ValueError):
|
||||
return JSONResponse(status_code=400, content={"error": "缺少 username/password"})
|
||||
user = db.get_user(username)
|
||||
if user is None or not verify_password(password, user["password_hash"]):
|
||||
return JSONResponse(status_code=401, content={"error": "用户名或密码错误"})
|
||||
token = new_token()
|
||||
db.set_token(username, token)
|
||||
log.info("login %s 用户=%s", request.client.host if request.client else "?", username)
|
||||
return {"username": username, "token": token}
|
||||
|
||||
|
||||
@app.get("/api/auth/me", dependencies=[Depends(verify_token)])
|
||||
async def me(username: str = Depends(verify_token)):
|
||||
"""当前登录用户 (token 有效性验证)"""
|
||||
return {"username": username}
|
||||
|
||||
|
||||
@app.post("/api/auth/logout", dependencies=[Depends(verify_token)])
|
||||
async def logout(username: str = Depends(verify_token)):
|
||||
"""注销: 使当前 token 失效"""
|
||||
if username != "admin":
|
||||
db.clear_token_by_username(username)
|
||||
return {"ok": True}
|
||||
|
||||
def _enc_headers(rec: dict) -> dict[str, str]:
|
||||
"""下载响应头: 解密参数 + 卷数"""
|
||||
ep = json.dumps(rec["enc_params"], ensure_ascii=False)
|
||||
return {
|
||||
"X-Enc-Params": base64.b64encode(ep.encode()).decode(),
|
||||
"X-Chunk-Count": str(rec["chunk_count"]),
|
||||
}
|
||||
|
||||
|
||||
# ---------- 文件: ls / 逐卷下载 ----------
|
||||
|
||||
@app.get("/api/files")
|
||||
async def list_files(request: Request, username: str = Depends(verify_token)):
|
||||
"""文件列表 (ls)"""
|
||||
files = db.list_files(username)
|
||||
log.info("ls %s %s 个文件", request.client.host if request.client else "?", len(files))
|
||||
return {"files": files}
|
||||
|
||||
|
||||
@app.get("/api/quota")
|
||||
async def quota(request: Request, username: str = Depends(verify_token)):
|
||||
"""用户空间配额: 已用 (files size 总和) / 配额 / 剩余"""
|
||||
used = db.sum_files_size(username)
|
||||
remain = max(QUOTA_BYTES - used, 0)
|
||||
log.info("quota %s 已用 %.1fMB/%.1fMB",
|
||||
request.client.host if request.client else "?", used / 1048576, QUOTA_BYTES / 1048576)
|
||||
return {
|
||||
"used_bytes": used,
|
||||
"quota_bytes": QUOTA_BYTES,
|
||||
"remain_bytes": remain,
|
||||
"percent": round(used * 100 / QUOTA_BYTES, 1) if QUOTA_BYTES else 0,
|
||||
}
|
||||
|
||||
|
||||
# ---------- 目录 ----------
|
||||
|
||||
|
||||
@app.post("/api/dirs")
|
||||
async def create_dir(request: Request, username: str = Depends(verify_token)):
|
||||
"""创建目录 (按用户隔离, 重名 409)"""
|
||||
body = await request.json()
|
||||
name = str(body.get("name", "")).strip()
|
||||
if not name or "/" in name or "\\" in name:
|
||||
return JSONResponse(status_code=400, content={"error": "目录名非法 (不能含 / 或 \\\\)"})
|
||||
if len(name) > 64:
|
||||
return JSONResponse(status_code=400, content={"error": "目录名过长 (≤64)"})
|
||||
dir_id = db.create_dir(name, username)
|
||||
if dir_id is None:
|
||||
return JSONResponse(status_code=409, content={"error": f"目录已存在: {name}"})
|
||||
log.info("mkdir %s %s (dir_id=%s)", request.client.host if request.client else "?", name, dir_id)
|
||||
return {"ok": True, "dir_id": dir_id, "name": name}
|
||||
|
||||
|
||||
@app.get("/api/dirs")
|
||||
async def list_dirs(request: Request, username: str = Depends(verify_token)):
|
||||
"""目录列表"""
|
||||
dirs = db.list_dirs(username)
|
||||
return {"dirs": dirs}
|
||||
|
||||
|
||||
@app.delete("/api/dirs/{dir_id}")
|
||||
async def delete_dir(dir_id: int, request: Request, username: str = Depends(verify_token)):
|
||||
"""删除目录 (归属校验, 不存在/跨账户 404)"""
|
||||
rec = db.get_dir(dir_id)
|
||||
if rec is None or rec["user_id"] != username:
|
||||
return JSONResponse(status_code=404, content={"error": "目录不存在"})
|
||||
db.delete_dir(dir_id, username)
|
||||
log.info("rmdir %s dir_id=%s (%s)", request.client.host if request.client else "?", dir_id, rec["name"])
|
||||
return {"ok": True, "dir_id": dir_id}
|
||||
|
||||
|
||||
@app.delete("/api/files/{file_id}")
|
||||
async def delete_file(file_id: str, request: Request, username: str = Depends(verify_token)):
|
||||
"""删除文件 (卷目录 + 入库记录)"""
|
||||
rec = db.get_file(file_id)
|
||||
if rec is None:
|
||||
return JSONResponse(status_code=404, content={"error": "文件不存在"})
|
||||
_check_owner(rec, username)
|
||||
try:
|
||||
p = Path(rec["path"])
|
||||
if p.is_dir():
|
||||
shutil.rmtree(p, ignore_errors=True)
|
||||
elif p.is_file():
|
||||
p.unlink(missing_ok=True)
|
||||
db.delete_file(file_id)
|
||||
log.info("delete %s %s", request.client.host if request.client else "?", file_id)
|
||||
except OSError as e:
|
||||
return JSONResponse(status_code=409, content={"error": f"删除失败: {e}"})
|
||||
return {"ok": True, "file_id": file_id}
|
||||
|
||||
|
||||
def _resolve_dir(dir_id: int | None, username: str) -> int | None:
|
||||
"""目录归属校验: 0/None=根目录, 否则必须属于该用户 (跨账户 403)"""
|
||||
if dir_id in (None, 0):
|
||||
return None
|
||||
d = db.get_dir(dir_id)
|
||||
if d is None or d["user_id"] != username:
|
||||
raise HTTPException(status_code=403, detail="目录不存在或无权使用")
|
||||
return dir_id
|
||||
|
||||
|
||||
@app.put("/api/files/{file_id}/move")
|
||||
async def move_file(file_id: str, request: Request, username: str = Depends(verify_token)):
|
||||
"""移动文件到目录 (body {dir_id}; 0/缺省=根目录)"""
|
||||
rec = db.get_file(file_id)
|
||||
if rec is None:
|
||||
return JSONResponse(status_code=404, content={"error": "文件不存在"})
|
||||
_check_owner(rec, username)
|
||||
body = await request.json()
|
||||
try:
|
||||
dir_id = _resolve_dir(body.get("dir_id"), username)
|
||||
except HTTPException as e:
|
||||
return JSONResponse(status_code=e.status_code, content={"error": e.detail})
|
||||
db.move_file(file_id, dir_id)
|
||||
log.info("move %s %s -> dir=%s", request.client.host if request.client else "?", file_id, dir_id)
|
||||
return {"ok": True, "file_id": file_id, "dir_id": dir_id}
|
||||
|
||||
|
||||
@app.post("/api/files/{file_id}/clone")
|
||||
async def clone_file(file_id: str, request: Request, username: str = Depends(verify_token)):
|
||||
"""克隆文件 (物理复制卷目录 + 新记录), body {dir_id} 可选"""
|
||||
rec = db.get_file(file_id)
|
||||
if rec is None:
|
||||
return JSONResponse(status_code=404, content={"error": "文件不存在"})
|
||||
_check_owner(rec, username)
|
||||
body = await request.json()
|
||||
try:
|
||||
dir_id = _resolve_dir(body.get("dir_id"), username)
|
||||
except HTTPException as e:
|
||||
return JSONResponse(status_code=e.status_code, content={"error": e.detail})
|
||||
src_dir = Path(rec["path"])
|
||||
new_dir = src_dir.with_name(src_dir.name + "_clone")
|
||||
try:
|
||||
shutil.copytree(src_dir, new_dir)
|
||||
except OSError as e:
|
||||
return JSONResponse(status_code=409, content={"error": f"克隆失败: {e}"})
|
||||
new_id = db.clone_file(file_id, dir_id, str(new_dir))
|
||||
if not new_id:
|
||||
shutil.rmtree(new_dir, ignore_errors=True)
|
||||
return JSONResponse(status_code=404, content={"error": "源文件记录丢失"})
|
||||
log.info("clone %s %s -> %s", request.client.host if request.client else "?", file_id, new_id)
|
||||
return {"ok": True, "file_id": new_id}
|
||||
|
||||
|
||||
@app.get("/api/files/{file_id}/chunk/{idx}")
|
||||
async def download_chunk(file_id: str, idx: int, request: Request, username: str = Depends(verify_token)):
|
||||
"""下载单卷密文 (零合并: 服务端不拼接, 客户端逐卷拉取本地合并)"""
|
||||
rec = db.get_file(file_id)
|
||||
if rec is None:
|
||||
return JSONResponse(status_code=404, content={"error": "文件不存在"})
|
||||
_check_owner(rec, username)
|
||||
if not 1 <= idx <= rec["chunk_count"]:
|
||||
return JSONResponse(status_code=404, content={"error": f"卷号越界: {idx}"})
|
||||
chunk_path = Path(rec["path"]) / f"chunk_{idx:04d}"
|
||||
if not chunk_path.exists():
|
||||
return JSONResponse(status_code=404, content={"error": f"卷文件缺失: {chunk_path}"})
|
||||
log.info("download %s %s 卷%d/%d", request.client.host if request.client else "?",
|
||||
file_id, idx, rec["chunk_count"])
|
||||
return FileResponse(
|
||||
chunk_path,
|
||||
media_type="application/octet-stream",
|
||||
headers=_enc_headers(rec),
|
||||
)
|
||||
|
||||
|
||||
# ---------- 传输 ----------
|
||||
|
||||
@app.post("/api/transfer/init")
|
||||
async def init_transfer(request: Request, username: str = Depends(verify_token)):
|
||||
# 防超大 json: init 元数据不该超过 1MB (chunk spec 每条约 120B, 1MB 可容纳 ~8000 卷)
|
||||
length = request.headers.get("Content-Length")
|
||||
if length and int(length) > 1 << 20:
|
||||
return JSONResponse(status_code=413, content={"error": "init json 过大"})
|
||||
try:
|
||||
init_json = await request.json()
|
||||
dir_id = init_json.pop("dir_id", None)
|
||||
if dir_id is not None:
|
||||
d = db.get_dir(int(dir_id))
|
||||
if d is None or d["user_id"] != username:
|
||||
return JSONResponse(status_code=403, content={"error": "目录不存在或无权使用"})
|
||||
transfer_id = tasks.create(init_json, username, dir_id)
|
||||
log.info("init %s %s卷 任务=%s", request.client.host if request.client else "?",
|
||||
init_json.get("chunk_count"), transfer_id)
|
||||
except (KeyError, TypeError, ValueError) as e:
|
||||
return JSONResponse(status_code=400, content={"error": f"init json 非法: {e}"})
|
||||
return {"transfer_id": transfer_id}
|
||||
|
||||
|
||||
@app.get("/api/transfer/{transfer_id}/chunks")
|
||||
async def get_chunks(transfer_id: str, username: str = Depends(verify_token)):
|
||||
if tasks.get(transfer_id) is None:
|
||||
return JSONResponse(status_code=404, content={"error": "任务不存在"})
|
||||
return {"received": sorted(tasks.received(transfer_id))}
|
||||
|
||||
|
||||
@app.put("/api/transfer/{transfer_id}/chunk/{idx}")
|
||||
async def put_chunk(transfer_id: str, idx: int, request: Request, username: str = Depends(verify_token)):
|
||||
if tasks.get(transfer_id) is None:
|
||||
return JSONResponse(status_code=404, content={"error": "任务不存在"})
|
||||
data = await request.body()
|
||||
try:
|
||||
receiver.receive(transfer_id, idx, data)
|
||||
log.info("chunk %s %s 卷%d (%dB)", request.client.host if request.client else "?",
|
||||
transfer_id, idx, len(data))
|
||||
except ReceiverError as e:
|
||||
return JSONResponse(status_code=409, content={"error": str(e)})
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@app.post("/api/transfer/{transfer_id}/pull")
|
||||
async def pull_chunks(transfer_id: str, request: Request, username: str = Depends(verify_token)):
|
||||
"""aria2 反向拉卷: 客户端起临时 HTTP 服务, 服务端 aria2c 并发拉取 (打满带宽)
|
||||
|
||||
body: {base_url, token} —— 客户端临时服务地址 + 一次性拉取令牌
|
||||
卷文件拉取后逐卷走 receiver.receive 校验 (大小+SHA-256) 落盘。
|
||||
aria2 子进程走 asyncio.to_thread: 同步阻塞会卡死事件循环,
|
||||
导致 pull-status/其他请求全部排队超时 (进度条 0%)。
|
||||
"""
|
||||
task = tasks.get(transfer_id)
|
||||
if task is None:
|
||||
return JSONResponse(status_code=404, content={"error": "任务不存在"})
|
||||
try:
|
||||
body = await request.json()
|
||||
base_url = str(body["base_url"]).rstrip("/")
|
||||
pull_token = str(body["token"])
|
||||
except (KeyError, TypeError, ValueError):
|
||||
return JSONResponse(status_code=400, content={"error": "缺少 base_url/token"})
|
||||
aria2 = shutil.which("aria2c")
|
||||
if aria2 is None:
|
||||
return JSONResponse(status_code=409, content={"error": "服务端未安装 aria2 (apt install aria2)"})
|
||||
chunk_count = task.get("chunk_count", 0)
|
||||
chunk_dir = TMP_ROOT / transfer_id
|
||||
chunk_dir.mkdir(parents=True, exist_ok=True)
|
||||
urls_file = chunk_dir / "urls.txt"
|
||||
with open(urls_file, "w", encoding="utf-8") as f:
|
||||
for idx in range(1, chunk_count + 1):
|
||||
f.write(f"{base_url}/pull/{transfer_id}/{idx}?token={pull_token}\n")
|
||||
try:
|
||||
r = await asyncio.to_thread(
|
||||
subprocess.run,
|
||||
[aria2, "-i", str(urls_file), "-d", str(chunk_dir),
|
||||
"--max-concurrent-downloads=8", "--split=4", "--max-connection-per-server=4",
|
||||
"--min-split-size=1M", "--file-allocation=none", "--allow-overwrite=true",
|
||||
"--auto-file-renaming=false", "--console-log-level=notice",
|
||||
"--summary-interval=1", "--quiet=false", "--no-conf=true"],
|
||||
timeout=3600,
|
||||
)
|
||||
except OSError as e:
|
||||
return JSONResponse(status_code=500, content={"error": f"aria2 启动失败: {e}"})
|
||||
if r.returncode != 0:
|
||||
return JSONResponse(status_code=409, content={"error": f"aria2 拉取失败 (exit {r.returncode}), 详见服务端日志"})
|
||||
# 逐卷校验落盘 (大小 + SHA-256); aria2 落盘文件名 = URL 末段 (卷序号)
|
||||
pulled = 0
|
||||
for p in sorted(chunk_dir.iterdir()):
|
||||
if p.name == "urls.txt" or p.suffix in (".aria2", ".tmp"):
|
||||
continue
|
||||
try:
|
||||
idx = int(p.name)
|
||||
except ValueError:
|
||||
continue
|
||||
try:
|
||||
receiver.receive(transfer_id, idx, p.read_bytes())
|
||||
p.unlink(missing_ok=True)
|
||||
pulled += 1
|
||||
except ReceiverError:
|
||||
pass # 校验失败的卷忽略, 客户端会补传
|
||||
urls_file.unlink(missing_ok=True)
|
||||
log.info("pull %s %s %d/%d 卷 (aria2)", request.client.host if request.client else "?",
|
||||
transfer_id, pulled, chunk_count)
|
||||
return {"ok": True, "pulled": pulled, "total": chunk_count}
|
||||
|
||||
|
||||
@app.get("/api/transfer/{transfer_id}/pull-status")
|
||||
async def pull_status(transfer_id: str, request: Request, username: str = Depends(verify_token)):
|
||||
"""aria2 拉取进度: 已落盘卷文件数 (aria2 每拉完一卷落盘, received 标记在全部拉完后才打)"""
|
||||
task = tasks.get(transfer_id)
|
||||
if task is None:
|
||||
return JSONResponse(status_code=404, content={"error": "任务不存在"})
|
||||
chunk_dir = TMP_ROOT / transfer_id
|
||||
n = 0
|
||||
if chunk_dir.is_dir():
|
||||
for p in chunk_dir.iterdir():
|
||||
if p.name != "urls.txt" and p.suffix not in (".aria2", ".tmp") and p.name.isdigit():
|
||||
n += 1
|
||||
return {"pulled": n, "total": task.get("chunk_count", 0)}
|
||||
|
||||
|
||||
@app.post("/api/transfer/{transfer_id}/complete")
|
||||
async def complete(transfer_id: str, request: Request, username: str = Depends(verify_token)):
|
||||
if tasks.get(transfer_id) is None:
|
||||
return JSONResponse(status_code=404, content={"error": "任务不存在"})
|
||||
# 幂等: 已完成的任务直接返回已有 file_id
|
||||
done_rec = db.get_file_by_transfer(transfer_id)
|
||||
if done_rec is not None:
|
||||
return {"status": "done", "file_id": done_rec["file_id"]}
|
||||
try:
|
||||
transfer = tasks.get(transfer_id)
|
||||
assert transfer is not None
|
||||
# 零合并: 只校验卷齐, 卷目录直接 move 入存储区, 秒回
|
||||
got = tasks.received(transfer_id)
|
||||
if len(got) != transfer["chunk_count"]:
|
||||
return JSONResponse(status_code=409, content={
|
||||
"status": "incomplete",
|
||||
"received": sorted(got),
|
||||
"error": f"缺卷: {transfer['chunk_count'] - len(got)} 卷未上传",
|
||||
})
|
||||
tasks.set_status(transfer_id, "storing")
|
||||
final_dir = storage.store_chunks(transfer_id, TMP_ROOT / transfer_id, username)
|
||||
file_id = db.insert_file(
|
||||
transfer_id,
|
||||
transfer["file_name"],
|
||||
str(final_dir),
|
||||
transfer["file_size"],
|
||||
transfer["total_sha256"],
|
||||
username,
|
||||
transfer.get("dir_id"),
|
||||
)
|
||||
tasks.set_status(transfer_id, "done")
|
||||
log.info("complete %s 任务=%s %s卷 -> file=%s",
|
||||
request.client.host if request.client else "?", transfer_id,
|
||||
transfer["chunk_count"], file_id)
|
||||
return {"status": "done", "file_id": file_id}
|
||||
except OSError as e:
|
||||
tasks.set_status(transfer_id, "failed")
|
||||
return JSONResponse(status_code=409, content={"error": str(e), "status": "failed"})
|
||||
@@ -0,0 +1,41 @@
|
||||
"""
|
||||
认证工具 (auth)
|
||||
模块: 服务端 / 认证
|
||||
输入: 明文密码
|
||||
输出: pbkdf2 哈希 / 随机 token
|
||||
|
||||
标准库实现 (pbkdf2-hmac-sha256, 20 万轮迭代), 不引第三方
|
||||
密码不落库: 只存 salt$hash
|
||||
"""
|
||||
|
||||
import hashlib
|
||||
import hmac
|
||||
import secrets
|
||||
|
||||
_ITERATIONS = 200_000
|
||||
|
||||
|
||||
def hash_password(password: str) -> str:
|
||||
"""密码 -> 'salt$hash' (pbkdf2-hmac-sha256, 随机盐)"""
|
||||
salt = secrets.token_hex(16)
|
||||
dk = hashlib.pbkdf2_hmac(
|
||||
"sha256", password.encode("utf-8"), bytes.fromhex(salt), _ITERATIONS
|
||||
)
|
||||
return f"{salt}${dk.hex()}"
|
||||
|
||||
|
||||
def verify_password(password: str, stored: str) -> bool:
|
||||
"""校验明文密码与存储哈希 (常数时间比较)"""
|
||||
try:
|
||||
salt, expected = stored.split("$", 1)
|
||||
dk = hashlib.pbkdf2_hmac(
|
||||
"sha256", password.encode("utf-8"), bytes.fromhex(salt), _ITERATIONS
|
||||
)
|
||||
return hmac.compare_digest(dk.hex(), expected)
|
||||
except ValueError:
|
||||
return False
|
||||
|
||||
|
||||
def new_token() -> str:
|
||||
"""随机 token (32 字节 hex, 会话凭证)"""
|
||||
return secrets.token_hex(32)
|
||||
@@ -0,0 +1,399 @@
|
||||
"""
|
||||
入库层 (db)
|
||||
模块: 服务端 / 入库层
|
||||
输入: 文件路径 + 元数据
|
||||
输出: file_id
|
||||
|
||||
SQLite 实现 (老板拍板: PRoot 容器装 PG/MySQL 服务进程都会 fsync 卡死,
|
||||
轻量文件库最稳)。表结构与 PostgreSQL 版一致, 未来可平滑切回 PG。
|
||||
|
||||
表:
|
||||
transfers 任务主表 (状态机 uploading->assembling->decrypting->done/failed)
|
||||
transfer_chunks 每卷信息 + 接收时间 (null = 未收, 断点续传依据)
|
||||
files 明文文件记录
|
||||
"""
|
||||
|
||||
import json
|
||||
import sqlite3
|
||||
import uuid
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from settings import DB_PATH
|
||||
|
||||
|
||||
class ServerDB:
|
||||
"""服务端数据库访问 (sqlite3, 每次操作独立连接, WAL 并发安全)"""
|
||||
|
||||
def __init__(self, db_path: str | Path = DB_PATH) -> None:
|
||||
self.db_path = Path(db_path)
|
||||
self._init_schema()
|
||||
|
||||
def _conn(self) -> sqlite3.Connection:
|
||||
conn = sqlite3.connect(self.db_path)
|
||||
conn.row_factory = sqlite3.Row
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA busy_timeout=5000")
|
||||
return conn
|
||||
|
||||
def _init_schema(self) -> None:
|
||||
self.db_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS transfers (
|
||||
transfer_id TEXT PRIMARY KEY,
|
||||
file_name TEXT NOT NULL,
|
||||
file_size INTEGER NOT NULL,
|
||||
chunk_count INTEGER NOT NULL,
|
||||
total_sha256 TEXT NOT NULL,
|
||||
enc_params TEXT NOT NULL,
|
||||
user_id TEXT,
|
||||
status TEXT NOT NULL,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
|
||||
)
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS transfer_chunks (
|
||||
transfer_id TEXT NOT NULL REFERENCES transfers(transfer_id),
|
||||
idx INTEGER NOT NULL,
|
||||
size INTEGER NOT NULL,
|
||||
sha256 TEXT NOT NULL,
|
||||
received_at TEXT,
|
||||
PRIMARY KEY (transfer_id, idx)
|
||||
)
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS files (
|
||||
file_id TEXT PRIMARY KEY,
|
||||
transfer_id TEXT NOT NULL REFERENCES transfers(transfer_id),
|
||||
file_name TEXT NOT NULL,
|
||||
path TEXT NOT NULL,
|
||||
size INTEGER NOT NULL,
|
||||
sha256 TEXT NOT NULL,
|
||||
user_id TEXT,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
||||
)
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS users (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
username TEXT UNIQUE NOT NULL,
|
||||
password_hash TEXT NOT NULL,
|
||||
token TEXT,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
||||
)
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE TABLE IF NOT EXISTS dirs (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
name TEXT NOT NULL,
|
||||
user_id TEXT,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
UNIQUE (name, user_id)
|
||||
)
|
||||
"""
|
||||
)
|
||||
# 旧库升级: files/transfers 补 user_id 列 (老数据 user_id=NULL -> admin 可见)
|
||||
try:
|
||||
conn.execute("ALTER TABLE transfers ADD COLUMN user_id TEXT")
|
||||
except sqlite3.OperationalError:
|
||||
pass # 列已存在
|
||||
try:
|
||||
conn.execute("ALTER TABLE files ADD COLUMN user_id TEXT")
|
||||
except sqlite3.OperationalError:
|
||||
pass # 列已存在
|
||||
try:
|
||||
conn.execute("ALTER TABLE transfers ADD COLUMN dir_id INTEGER")
|
||||
except sqlite3.OperationalError:
|
||||
pass # 列已存在
|
||||
try:
|
||||
conn.execute("ALTER TABLE files ADD COLUMN dir_id INTEGER")
|
||||
except sqlite3.OperationalError:
|
||||
pass # 列已存在
|
||||
|
||||
# ---------- users ----------
|
||||
|
||||
def create_user(self, username: str, password_hash: str) -> bool:
|
||||
"""注册用户, 返回是否成功 (False = 用户名已存在)"""
|
||||
try:
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"INSERT INTO users (username, password_hash) VALUES (?, ?)",
|
||||
(username, password_hash),
|
||||
)
|
||||
return True
|
||||
except sqlite3.IntegrityError:
|
||||
return False
|
||||
|
||||
def get_user(self, username: str) -> dict[str, Any] | None:
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT id, username, password_hash, token FROM users WHERE username = ?",
|
||||
(username,),
|
||||
).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
def set_token(self, username: str, token: str) -> None:
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"UPDATE users SET token = ? WHERE username = ?", (token, username)
|
||||
)
|
||||
|
||||
def get_username_by_token(self, token: str) -> str | None:
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT username FROM users WHERE token = ?", (token,)
|
||||
).fetchone()
|
||||
return str(row["username"]) if row else None
|
||||
|
||||
def clear_token(self, token: str) -> None:
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"UPDATE users SET token = NULL WHERE token = ?", (token,)
|
||||
)
|
||||
|
||||
def clear_token_by_username(self, username: str) -> None:
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"UPDATE users SET token = NULL WHERE username = ?", (username,)
|
||||
)
|
||||
|
||||
# ---------- dirs ----------
|
||||
|
||||
def create_dir(self, name: str, user_id: str) -> int | None:
|
||||
"""创建目录, 返回 dir id; None = 重名"""
|
||||
try:
|
||||
with self._conn() as conn:
|
||||
cur = conn.execute(
|
||||
"INSERT INTO dirs (name, user_id) VALUES (?, ?)", (name, user_id)
|
||||
)
|
||||
rid = cur.lastrowid
|
||||
return int(rid) if rid is not None else None
|
||||
except sqlite3.IntegrityError:
|
||||
return None
|
||||
|
||||
def get_dir(self, dir_id: int) -> dict[str, Any] | None:
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT id, name, user_id, created_at FROM dirs WHERE id = ?", (dir_id,)
|
||||
).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
def list_dirs(self, user_id: str) -> list[dict[str, Any]]:
|
||||
with self._conn() as conn:
|
||||
rows = conn.execute(
|
||||
"SELECT id, name, created_at FROM dirs WHERE user_id = ? ORDER BY id",
|
||||
(user_id,),
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
def delete_dir(self, dir_id: int, user_id: str) -> bool:
|
||||
"""删除目录 (归属校验), 返回是否删除"""
|
||||
with self._conn() as conn:
|
||||
cur = conn.execute(
|
||||
"DELETE FROM dirs WHERE id = ? AND user_id = ?", (dir_id, user_id)
|
||||
)
|
||||
return cur.rowcount > 0
|
||||
|
||||
# ---------- transfers ----------
|
||||
|
||||
def create_transfer(self, transfer_id: str, init_json: dict[str, Any], user_id: str,
|
||||
dir_id: int | None = None) -> None:
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"INSERT INTO transfers (transfer_id, file_name, file_size, chunk_count, total_sha256, enc_params, user_id, dir_id, status) "
|
||||
"VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
(
|
||||
transfer_id,
|
||||
init_json["file_name"],
|
||||
init_json["file_size"],
|
||||
init_json["chunk_count"],
|
||||
init_json["total_sha256"],
|
||||
json.dumps(init_json["enc"], ensure_ascii=False),
|
||||
user_id,
|
||||
dir_id,
|
||||
"uploading",
|
||||
),
|
||||
)
|
||||
conn.executemany(
|
||||
"INSERT INTO transfer_chunks (transfer_id, idx, size, sha256) VALUES (?, ?, ?, ?)",
|
||||
[
|
||||
(transfer_id, ch["index"], ch["size"], ch["sha256"])
|
||||
for ch in init_json["chunks"]
|
||||
],
|
||||
)
|
||||
|
||||
def get_transfer(self, transfer_id: str) -> dict[str, Any] | None:
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT * FROM transfers WHERE transfer_id = ?", (transfer_id,)
|
||||
).fetchone()
|
||||
if not row:
|
||||
return None
|
||||
d = dict(row)
|
||||
d["enc_params"] = json.loads(d["enc_params"])
|
||||
return d
|
||||
|
||||
def set_status(self, transfer_id: str, status: str) -> None:
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"UPDATE transfers SET status = ?, updated_at = datetime('now') WHERE transfer_id = ?",
|
||||
(status, transfer_id),
|
||||
)
|
||||
|
||||
# ---------- chunks ----------
|
||||
|
||||
def get_chunk_spec(self, transfer_id: str, idx: int) -> dict[str, Any] | None:
|
||||
"""卷规格 (size/sha256), 接收校验用"""
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT size, sha256 FROM transfer_chunks WHERE transfer_id = ? AND idx = ?",
|
||||
(transfer_id, idx),
|
||||
).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
def is_chunk_received(self, transfer_id: str, idx: int) -> bool:
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT 1 FROM transfer_chunks WHERE transfer_id = ? AND idx = ? AND received_at IS NOT NULL",
|
||||
(transfer_id, idx),
|
||||
).fetchone()
|
||||
return row is not None
|
||||
|
||||
def mark_chunk_received(self, transfer_id: str, idx: int) -> None:
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"UPDATE transfer_chunks SET received_at = datetime('now') "
|
||||
"WHERE transfer_id = ? AND idx = ?",
|
||||
(transfer_id, idx),
|
||||
)
|
||||
|
||||
def get_received_chunks(self, transfer_id: str) -> set[int]:
|
||||
with self._conn() as conn:
|
||||
rows = conn.execute(
|
||||
"SELECT idx FROM transfer_chunks WHERE transfer_id = ? AND received_at IS NOT NULL",
|
||||
(transfer_id,),
|
||||
).fetchall()
|
||||
return {r["idx"] for r in rows}
|
||||
|
||||
# ---------- files ----------
|
||||
|
||||
def insert_file(
|
||||
self,
|
||||
transfer_id: str,
|
||||
file_name: str,
|
||||
path: str,
|
||||
size: int,
|
||||
sha256: str,
|
||||
user_id: str,
|
||||
dir_id: int | None = None,
|
||||
) -> str:
|
||||
file_id = uuid.uuid4().hex[:12]
|
||||
with self._conn() as conn:
|
||||
conn.execute(
|
||||
"INSERT INTO files (file_id, transfer_id, file_name, path, size, sha256, user_id, dir_id) "
|
||||
"VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
(file_id, transfer_id, file_name, path, size, sha256, user_id, dir_id),
|
||||
)
|
||||
return file_id
|
||||
|
||||
def get_file_by_transfer(self, transfer_id: str) -> dict[str, Any] | None:
|
||||
"""按任务查已入库文件 (complete 幂等: 已完成直接返回)"""
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT file_id, file_name, path, size FROM files WHERE transfer_id = ?",
|
||||
(transfer_id,),
|
||||
).fetchone()
|
||||
return dict(row) if row else None
|
||||
|
||||
def list_files(self, user_id: str) -> list[dict[str, Any]]:
|
||||
"""文件记录 (ls 端点): 每个账户只看自己的空间 (admin 含无主旧文件), 带目录名"""
|
||||
with self._conn() as conn:
|
||||
if user_id == "admin":
|
||||
rows = conn.execute(
|
||||
"SELECT f.file_id, f.file_name, f.size, f.sha256, f.created_at, f.dir_id, d.name AS dir_name "
|
||||
"FROM files f LEFT JOIN dirs d ON d.id = f.dir_id "
|
||||
"WHERE f.user_id = ? OR f.user_id IS NULL ORDER BY f.created_at DESC",
|
||||
("admin",),
|
||||
).fetchall()
|
||||
else:
|
||||
rows = conn.execute(
|
||||
"SELECT f.file_id, f.file_name, f.size, f.sha256, f.created_at, f.dir_id, d.name AS dir_name "
|
||||
"FROM files f LEFT JOIN dirs d ON d.id = f.dir_id "
|
||||
"WHERE f.user_id = ? ORDER BY f.created_at DESC",
|
||||
(user_id,),
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
def move_file(self, file_id: str, dir_id: int | None) -> bool:
|
||||
"""移动文件到目录 (dir_id=None=根目录), 返回是否更新"""
|
||||
with self._conn() as conn:
|
||||
cur = conn.execute(
|
||||
"UPDATE files SET dir_id = ? WHERE file_id = ?", (dir_id, file_id)
|
||||
)
|
||||
return cur.rowcount > 0
|
||||
|
||||
def clone_file(self, file_id: str, dir_id: int | None, new_path: str) -> str:
|
||||
"""克隆文件记录 (卷目录已物理复制到 new_path), 返回新 file_id"""
|
||||
new_id = uuid.uuid4().hex[:12]
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT transfer_id, file_name, size, sha256, user_id FROM files WHERE file_id = ?",
|
||||
(file_id,),
|
||||
).fetchone()
|
||||
if row is None:
|
||||
return ""
|
||||
conn.execute(
|
||||
"INSERT INTO files (file_id, transfer_id, file_name, path, size, sha256, user_id, dir_id) "
|
||||
"VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
|
||||
(new_id, row["transfer_id"], row["file_name"], new_path,
|
||||
row["size"], row["sha256"], row["user_id"], dir_id),
|
||||
)
|
||||
return new_id
|
||||
|
||||
def delete_file(self, file_id: str) -> bool:
|
||||
"""删除文件记录, 返回是否删除 (del 端点)"""
|
||||
with self._conn() as conn:
|
||||
cur = conn.execute("DELETE FROM files WHERE file_id = ?", (file_id,))
|
||||
return cur.rowcount > 0
|
||||
|
||||
def sum_files_size(self, user_id: str) -> int:
|
||||
"""文件 size 总和 (配额统计用): 每个账户只算自己的空间"""
|
||||
with self._conn() as conn:
|
||||
if user_id == "admin":
|
||||
row = conn.execute(
|
||||
"SELECT COALESCE(SUM(size), 0) FROM files WHERE user_id = ? OR user_id IS NULL",
|
||||
("admin",),
|
||||
).fetchone()
|
||||
else:
|
||||
row = conn.execute(
|
||||
"SELECT COALESCE(SUM(size), 0) FROM files WHERE user_id = ?",
|
||||
(user_id,),
|
||||
).fetchone()
|
||||
return int(row[0])
|
||||
|
||||
def get_file(self, file_id: str) -> dict[str, Any] | None:
|
||||
"""单文件记录 (下载端点, 含解密参数 enc_params + 卷数)"""
|
||||
with self._conn() as conn:
|
||||
row = conn.execute(
|
||||
"SELECT f.file_id, f.file_name, f.path, f.size, f.sha256, f.user_id, "
|
||||
"t.enc_params, t.chunk_count "
|
||||
"FROM files f JOIN transfers t ON f.transfer_id = t.transfer_id "
|
||||
"WHERE f.file_id = ?",
|
||||
(file_id,),
|
||||
).fetchone()
|
||||
if not row:
|
||||
return None
|
||||
d = dict(row)
|
||||
d["enc_params"] = json.loads(d["enc_params"])
|
||||
return d
|
||||
@@ -0,0 +1,14 @@
|
||||
"""服务端入口: uvicorn 启动
|
||||
|
||||
用法 (server/ 目录下):
|
||||
.venv/bin/python main.py
|
||||
或:
|
||||
DATABASE_URL=... SZ_PORT=8000 .venv/bin/python main.py
|
||||
"""
|
||||
|
||||
import uvicorn
|
||||
|
||||
from settings import HOST, PORT
|
||||
|
||||
if __name__ == "__main__":
|
||||
uvicorn.run("api:app", host=HOST, port=PORT, log_level="info")
|
||||
@@ -0,0 +1,49 @@
|
||||
"""
|
||||
分卷接收模块 (receiver)
|
||||
模块: 服务端 / 分卷接收
|
||||
输入: 卷数据 (transfer_id, index)
|
||||
输出: 落盘卷 + ok/409
|
||||
|
||||
逐卷校验 SHA-256 (对照 init 登记的规格); 重复卷直接 ok (幂等, 断点续传重发不炸)。
|
||||
"""
|
||||
|
||||
import hashlib
|
||||
import shutil
|
||||
from pathlib import Path
|
||||
|
||||
from task_manager import TaskManager
|
||||
|
||||
|
||||
class ReceiverError(RuntimeError):
|
||||
"""卷校验失败 (哈希不一致/规格缺失)"""
|
||||
|
||||
|
||||
class Receiver:
|
||||
def __init__(self, tasks: TaskManager, tmp_root: str | Path) -> None:
|
||||
self.tasks = tasks
|
||||
self.tmp_root = Path(tmp_root)
|
||||
|
||||
def receive(self, transfer_id: str, idx: int, data: bytes) -> None:
|
||||
"""接收并校验一卷。失败抛 ReceiverError (api 转 409)。"""
|
||||
if self.tasks.is_chunk_received(transfer_id, idx):
|
||||
return # 幂等: 重复卷直接 ok
|
||||
spec = self.tasks.chunk_spec(transfer_id, idx)
|
||||
if spec is None:
|
||||
raise ReceiverError(f"任务 {transfer_id} 卷 {idx} 规格不存在")
|
||||
|
||||
# 校验大小 + SHA-256
|
||||
if len(data) != spec["size"]:
|
||||
raise ReceiverError(f"卷 {idx} 大小不符: 期望 {spec['size']}, 实际 {len(data)}")
|
||||
sha256 = hashlib.sha256(data).hexdigest()
|
||||
if sha256 != spec["sha256"]:
|
||||
raise ReceiverError(f"卷 {idx} SHA-256 不一致 (数据损坏或串卷)")
|
||||
|
||||
# 落盘临时目录
|
||||
chunk_dir = self.tmp_root / transfer_id
|
||||
chunk_dir.mkdir(parents=True, exist_ok=True)
|
||||
chunk_path = chunk_dir / f"chunk_{idx:04d}"
|
||||
tmp_path = chunk_path.with_suffix(".tmp")
|
||||
tmp_path.write_bytes(data)
|
||||
shutil.move(str(tmp_path), str(chunk_path))
|
||||
|
||||
self.tasks.mark_chunk_received(transfer_id, idx)
|
||||
@@ -0,0 +1,25 @@
|
||||
"""服务端配置 (环境变量驱动, 带本地默认值)"""
|
||||
import os
|
||||
from pathlib import Path
|
||||
|
||||
# 应用绑定上下文 (必须与客户端一致)
|
||||
CONTEXT = os.environ.get("SZ_CONTEXT", "7z-encrypt:v1").encode()
|
||||
|
||||
# SQLite 数据库文件 (轻量部署: PRoot 容器装 PG/MySQL 会 fsync 卡死)
|
||||
DB_PATH = Path(os.environ.get("SZ_DB_PATH", "data/app.db"))
|
||||
|
||||
# 认证 token (预共享: 客户端 config.json 同值; 生产必须设置)
|
||||
TOKEN = os.environ.get("SZ_TOKEN", "sz-dev-token-change-me")
|
||||
|
||||
# 存储根目录 (密文落盘, 零知识: 服务端不碰明文)
|
||||
STORAGE_ROOT = Path(os.environ.get("SZ_STORAGE_ROOT", "data/storage"))
|
||||
|
||||
# 临时目录 (分卷/合并中间产物)
|
||||
TMP_ROOT = Path(os.environ.get("SZ_TMP_ROOT", "data/tmp"))
|
||||
|
||||
# 用户空间配额 (已用 = files 表 size 总和, 超配额拒绝新上传? 暂只读展示)
|
||||
QUOTA_BYTES = int(os.environ.get("SZ_QUOTA", str(10 * 1024**3))) # 默认 10GB
|
||||
|
||||
# API 服务
|
||||
HOST = os.environ.get("SZ_HOST", "0.0.0.0")
|
||||
PORT = int(os.environ.get("SZ_PORT", "8000"))
|
||||
@@ -0,0 +1,36 @@
|
||||
"""
|
||||
存储层模块 (storage)
|
||||
模块: 服务端 / 存储层
|
||||
输入: 已校验的密文卷目录 + 元数据
|
||||
输出: 最终存储目录
|
||||
|
||||
零知识 + 零合并: 只存卷, 不做任何拼接/合并, 服务端无密钥。
|
||||
卷目录整体 move 到日期/transfer_id 目录 (transfer_id 唯一, 无同名冲突)。
|
||||
"""
|
||||
|
||||
import shutil
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
|
||||
from task_manager import TaskManager
|
||||
|
||||
|
||||
class Storage:
|
||||
def __init__(self, tasks: TaskManager, storage_root: str | Path) -> None:
|
||||
self.tasks = tasks
|
||||
self.storage_root = Path(storage_root)
|
||||
|
||||
def store_chunks(self, transfer_id: str, chunk_dir: Path, user_id: str) -> Path:
|
||||
"""卷目录 -> 用户目录/日期/transfer_id (整体 move, 返回目标目录路径)
|
||||
|
||||
每用户独立目录 (user_id 已过注册校验: 3-32 位字母数字/下划线, 无路径风险)。
|
||||
注意: 目标目录不能预创建, 否则 shutil.move 会嵌套成 dst/dst/。
|
||||
"""
|
||||
user_dir = self.storage_root / user_id
|
||||
day_dir = user_dir / datetime.now().strftime("%Y-%m-%d")
|
||||
final_dir = day_dir / transfer_id
|
||||
if final_dir.exists():
|
||||
raise OSError(f"存储目录已存在: {final_dir}")
|
||||
day_dir.mkdir(parents=True, exist_ok=True)
|
||||
shutil.move(str(chunk_dir), str(final_dir))
|
||||
return final_dir
|
||||
@@ -0,0 +1,67 @@
|
||||
"""
|
||||
任务管理模块 (task_manager)
|
||||
模块: 服务端 / 任务管理
|
||||
输入: init 请求 / 卷到达事件
|
||||
输出: transfer 状态机
|
||||
|
||||
状态机: uploading -> storing -> done / failed
|
||||
(零合并: 无 assembling/decrypting 阶段, 卷齐即入库)
|
||||
transfer_id 生命周期贯穿接收->入库; 孤儿任务由清理逻辑回收 (见 server/data/tmp)。
|
||||
"""
|
||||
|
||||
import uuid
|
||||
from typing import Any
|
||||
|
||||
from db import ServerDB
|
||||
|
||||
STATUS_UPLOADING = "uploading"
|
||||
STATUS_STORING = "storing"
|
||||
STATUS_DONE = "done"
|
||||
STATUS_FAILED = "failed"
|
||||
|
||||
|
||||
class TaskManager:
|
||||
"""transfer 生命周期管理 (基于 ServerDB)"""
|
||||
|
||||
def __init__(self, db: ServerDB) -> None:
|
||||
self.db = db
|
||||
|
||||
def create(self, init_json: dict[str, Any], user_id: str,
|
||||
dir_id: int | None = None) -> str:
|
||||
"""POST /init: 建任务, 返回 transfer_id (输入清洗: 文件名/卷数上限)
|
||||
|
||||
文件名只做可打印字符过滤 + 长度截断。不删 / \\ 等路径分隔符:
|
||||
加密文件名是 URL-safe base64 (可能含 - _ =), 路径安全由下载端 _safe_name 兜底
|
||||
"""
|
||||
# 文件名字段清洗: 长度上限 + 去路径分隔符/控制字符
|
||||
raw = str(init_json.get("file_name", ""))
|
||||
clean = "".join(c for c in raw if c.isprintable())[:255]
|
||||
init_json["file_name"] = clean or "unnamed"
|
||||
# 卷数上限: 防超大清单放大 (1MB init / 120B 每卷 spec ≈ 8000)
|
||||
chunks = init_json.get("chunks") or []
|
||||
if len(chunks) > 10000:
|
||||
raise ValueError(f"卷数超上限: {len(chunks)} > 10000")
|
||||
init_json["chunk_count"] = len(chunks)
|
||||
transfer_id = uuid.uuid4().hex[:12]
|
||||
self.db.create_transfer(transfer_id, init_json, user_id, dir_id)
|
||||
return transfer_id
|
||||
|
||||
def get(self, transfer_id: str) -> dict[str, Any] | None:
|
||||
return self.db.get_transfer(transfer_id)
|
||||
|
||||
def chunk_spec(self, transfer_id: str, idx: int) -> dict[str, Any] | None:
|
||||
return self.db.get_chunk_spec(transfer_id, idx)
|
||||
|
||||
def is_chunk_received(self, transfer_id: str, idx: int) -> bool:
|
||||
return self.db.is_chunk_received(transfer_id, idx)
|
||||
|
||||
def mark_chunk_received(self, transfer_id: str, idx: int) -> None:
|
||||
"""事件: 卷到达 (校验通过后)"""
|
||||
self.db.mark_chunk_received(transfer_id, idx)
|
||||
|
||||
def received(self, transfer_id: str) -> set[int]:
|
||||
"""已收卷集合 (GET /chunks 依据)"""
|
||||
return self.db.get_received_chunks(transfer_id)
|
||||
|
||||
def set_status(self, transfer_id: str, status: str) -> None:
|
||||
self.db.set_status(transfer_id, status)
|
||||
Binary file not shown.
@@ -0,0 +1,7 @@
|
||||
Format: https://www.debian.org/doc/packaging-manuals/copyright-format/1.0/
|
||||
Upstream-Name: 7z-encrypt-server
|
||||
Source: https://gitee.com/edgevoid/7z-encrypt-server
|
||||
|
||||
Files: *
|
||||
Copyright: 2026 edgevoid
|
||||
License: MIT
|
||||
Vendored
+8
@@ -0,0 +1,8 @@
|
||||
7z-encrypt-server (1.3.0) stable; urgency=medium
|
||||
|
||||
* 服务端 deb 化: 源码版 (all, 手机/容器) + 编译版 (amd64, Nuitka onefile)
|
||||
* postinst 自动安装依赖: aria2 + p7zip-full + venv (fastapi/uvicorn/cryptography/tqdm)
|
||||
* 目录管理: dirs 表 + CRUD 端点 + init dir_id + move/clone
|
||||
* 账户严格隔离 + pull-status 实时进度
|
||||
|
||||
-- edgevoid <edgevoid@users.noreply.gitee.com> Mon, 10 Aug 2026 21:50:00 +0800
|
||||
Vendored
+26
@@ -0,0 +1,26 @@
|
||||
Source: 7z-encrypt-server
|
||||
Section: utils
|
||||
Priority: optional
|
||||
Maintainer: edgevoid <edgevoid@users.noreply.gitee.com>
|
||||
Build-Depends: debhelper-compat (= 13)
|
||||
Standards-Version: 4.6.2
|
||||
|
||||
Package: 7z-encrypt-server
|
||||
Architecture: all
|
||||
Depends: python3 (>= 3.10), python3-venv, aria2, p7zip-full
|
||||
Description: 端到端加密文件传输服务端 (源码版, 手机/容器用)
|
||||
零知识 + 零合并的服务端: 只存卷/传卷, 不合并不解密, 不持有密钥。
|
||||
账户严格隔离 (一个账户一个空间), 目录管理 (dirs/move/clone),
|
||||
aria2 反向拉加速 + pull-status 实时进度。
|
||||
.
|
||||
包含: sz-server (入口) + 配套模块
|
||||
端口 8000 (SZ_PORT), 存储/数据目录可配 (SZ_DATA_DIR)。
|
||||
装完即用: 依赖 (fastapi/uvicorn/cryptography/tqdm) postinst 自动安装。
|
||||
|
||||
Package: 7z-encrypt-server-bin
|
||||
Architecture: amd64
|
||||
Depends: aria2, p7zip-full
|
||||
Description: 端到端加密文件传输服务端 (x86_64 原生编译版)
|
||||
Nuitka 编译的原生二进制 (Python->C->gcc, onefile 自包含, 免 Python 依赖)。
|
||||
功能与源码版一致: sz-server。
|
||||
装完即用: 系统依赖 (aria2/p7zip) postinst 自动安装。
|
||||
Vendored
+7
@@ -0,0 +1,7 @@
|
||||
Format: https://www.debian.org/doc/packaging-manuals/copyright-format/1.0/
|
||||
Upstream-Name: 7z-encrypt-server
|
||||
Source: https://gitee.com/edgevoid/7z-encrypt-server
|
||||
|
||||
Files: *
|
||||
Copyright: 2026 edgevoid
|
||||
License: MIT
|
||||
Vendored
+2
@@ -0,0 +1,2 @@
|
||||
7z-encrypt-server
|
||||
7z-encrypt-server-bin
|
||||
Vendored
+3
@@ -0,0 +1,3 @@
|
||||
7z-encrypt-server-bin_1.3.0_amd64.deb utils optional architecture=amd64
|
||||
7z-encrypt-server_1.3.0_all.deb utils optional architecture=all
|
||||
7z-encrypt-server_1.3.0_amd64.buildinfo utils optional
|
||||
Vendored
+18
@@ -0,0 +1,18 @@
|
||||
#!/bin/sh
|
||||
set -e
|
||||
|
||||
case "$1" in
|
||||
configure)
|
||||
echo "正在安装系统依赖 (aria2 + p7zip-full) ..."
|
||||
apt-get install -y aria2 p7zip-full >/dev/null 2>&1 || true
|
||||
P=/usr/share/7z-encrypt-server/.venv
|
||||
if [ ! -x "$P/bin/python" ]; then
|
||||
echo "正在创建运行环境 (venv + fastapi/uvicorn/cryptography/tqdm) ..."
|
||||
python3 -m venv "$P"
|
||||
"$P/bin/pip" install --quiet --upgrade pip || true
|
||||
"$P/bin/pip" install --quiet fastapi "uvicorn[standard]" cryptography tqdm
|
||||
fi
|
||||
;;
|
||||
esac
|
||||
|
||||
exit 0
|
||||
+21
@@ -0,0 +1,21 @@
|
||||
#!/usr/bin/make -f
|
||||
%:
|
||||
dh $@
|
||||
|
||||
override_dh_auto_build:
|
||||
|
||||
override_dh_auto_test:
|
||||
|
||||
override_dh_auto_install:
|
||||
# --- 源码版 (all): 服务端源码 + wrapper ---
|
||||
mkdir -p debian/7z-encrypt-server/usr/share/7z-encrypt-server
|
||||
cp api.py auth.py db.py main.py receiver.py settings.py storage.py task_manager.py \
|
||||
debian/7z-encrypt-server/usr/share/7z-encrypt-server/
|
||||
mkdir -p debian/7z-encrypt-server/usr/bin
|
||||
printf '#!/bin/sh\nexec /usr/share/7z-encrypt-server/.venv/bin/python /usr/share/7z-encrypt-server/main.py "$$@"\n' \
|
||||
> debian/7z-encrypt-server/usr/bin/sz-server
|
||||
chmod 755 debian/7z-encrypt-server/usr/bin/sz-server
|
||||
# --- 编译版 (amd64): Nuitka 二进制 ---
|
||||
mkdir -p debian/7z-encrypt-server-bin/usr/bin
|
||||
cp bin/sz-server debian/7z-encrypt-server-bin/usr/bin/
|
||||
chmod 755 debian/7z-encrypt-server-bin/usr/bin/sz-server
|
||||
Vendored
+1
@@ -0,0 +1 @@
|
||||
3.0 (quilt)
|
||||
Reference in New Issue
Block a user