P0-B: 数据损坏修复
- v1 上传大小校验: multipart body 长度与文件字节数不应相等,改为上限校验(上传功能恢复) - mirror_sync 调度循环/remove_sync_source 死锁修复(锁外触发 start/stop) - HTTP 断点续传损坏修复: 续传时先复制已有前缀再追加尾部,服务器忽略 Range 时从头重下 - sync_engine 同款续传 bug 同步修复;temp_path 提前初始化避免 except NameError
This commit is contained in:
@@ -1015,7 +1015,9 @@ class APIv1:
|
|||||||
total_written += len(chunk)
|
total_written += len(chunk)
|
||||||
file_size = total_written
|
file_size = total_written
|
||||||
|
|
||||||
if content_length > 0 and file_size != content_length:
|
# content_length 是整个 multipart body 的长度(含 boundary/字段头),
|
||||||
|
# 文件字节数不可能等于它;这里只做上限校验防止越界读取
|
||||||
|
if content_length > 0 and file_size > content_length:
|
||||||
raise IOError(f"文件大小不匹配。期望: {content_length}, 实际: {file_size}")
|
raise IOError(f"文件大小不匹配。期望: {content_length}, 实际: {file_size}")
|
||||||
|
|
||||||
if os.path.exists(full_path):
|
if os.path.exists(full_path):
|
||||||
|
|||||||
+49
-13
@@ -193,6 +193,7 @@ class MirrorSyncManager:
|
|||||||
while self.scheduler_running:
|
while self.scheduler_running:
|
||||||
try:
|
try:
|
||||||
now = datetime.now()
|
now = datetime.now()
|
||||||
|
due_sources = []
|
||||||
with self.sync_lock:
|
with self.sync_lock:
|
||||||
for name, status in self.sync_status.items():
|
for name, status in self.sync_status.items():
|
||||||
source = self.sync_sources.get(name, {})
|
source = self.sync_sources.get(name, {})
|
||||||
@@ -202,13 +203,21 @@ class MirrorSyncManager:
|
|||||||
|
|
||||||
next_sync = status.get('next_sync')
|
next_sync = status.get('next_sync')
|
||||||
if next_sync:
|
if next_sync:
|
||||||
next_time = datetime.fromisoformat(next_sync)
|
try:
|
||||||
|
next_time = datetime.fromisoformat(next_sync)
|
||||||
|
except ValueError:
|
||||||
|
continue
|
||||||
if now >= next_time:
|
if now >= next_time:
|
||||||
# 触发同步
|
due_sources.append(name)
|
||||||
print(f"[定时同步] 触发同步: {name}")
|
|
||||||
self.start_sync(name)
|
# 锁外触发同步,避免与 start_sync 内的 sync_lock 形成死锁
|
||||||
# 计算下次同步时间
|
for name in due_sources:
|
||||||
self._calculate_next_sync(name)
|
try:
|
||||||
|
print(f"[定时同步] 触发同步: {name}")
|
||||||
|
self.start_sync(name)
|
||||||
|
self._calculate_next_sync(name)
|
||||||
|
except Exception as e:
|
||||||
|
print(f"[定时同步] 触发 {name} 失败: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
print(f"定时同步调度错误: {e}")
|
print(f"定时同步调度错误: {e}")
|
||||||
time.sleep(60) # 每分钟检查一次
|
time.sleep(60) # 每分钟检查一次
|
||||||
@@ -226,8 +235,12 @@ class MirrorSyncManager:
|
|||||||
del self.sync_sources[name]
|
del self.sync_sources[name]
|
||||||
if name in self.sync_status:
|
if name in self.sync_status:
|
||||||
del self.sync_status[name]
|
del self.sync_status[name]
|
||||||
if name in self.sync_threads:
|
|
||||||
self.stop_sync(name)
|
# 锁外停止线程(stop_sync 内部会再取 sync_lock,避免不可重入锁死锁)
|
||||||
|
if name in self.sync_threads:
|
||||||
|
self.stop_sync(name)
|
||||||
|
|
||||||
|
with self.sync_lock:
|
||||||
self.save_sync_state()
|
self.save_sync_state()
|
||||||
self._save_sync_sources()
|
self._save_sync_sources()
|
||||||
|
|
||||||
@@ -1118,11 +1131,14 @@ class MirrorSyncManager:
|
|||||||
return source_mtime > target_mtime or os.path.getsize(source_path) != os.path.getsize(target_path)
|
return source_mtime > target_mtime or os.path.getsize(source_path) != os.path.getsize(target_path)
|
||||||
|
|
||||||
def _download_file_http(self, url, local_path, config):
|
def _download_file_http(self, url, local_path, config):
|
||||||
"""下载HTTP文件"""
|
"""下载HTTP文件(支持断点续传,续传数据追加到已有前缀,避免损坏)"""
|
||||||
import urllib.request
|
import urllib.request
|
||||||
import base64
|
import base64
|
||||||
|
import shutil
|
||||||
from email.utils import parsedate
|
from email.utils import parsedate
|
||||||
|
|
||||||
|
temp_path = local_path + '.tmp' # 提前初始化,避免 except 分支 NameError
|
||||||
|
|
||||||
try:
|
try:
|
||||||
req = urllib.request.Request(url)
|
req = urllib.request.Request(url)
|
||||||
headers = config.get('headers', {})
|
headers = config.get('headers', {})
|
||||||
@@ -1142,12 +1158,24 @@ class MirrorSyncManager:
|
|||||||
if start_byte > 0:
|
if start_byte > 0:
|
||||||
req.add_header('Range', f'bytes={start_byte}-')
|
req.add_header('Range', f'bytes={start_byte}-')
|
||||||
|
|
||||||
temp_path = local_path + '.tmp'
|
|
||||||
mode = 'ab' if start_byte > 0 else 'wb'
|
|
||||||
|
|
||||||
with urllib.request.urlopen(req, timeout=timeout) as response:
|
with urllib.request.urlopen(req, timeout=timeout) as response:
|
||||||
if response.getcode() not in [200, 206]:
|
status_code = response.getcode()
|
||||||
raise Exception(f"HTTP错误: {response.getcode()}")
|
if status_code not in [200, 206]:
|
||||||
|
raise Exception(f"HTTP错误: {status_code}")
|
||||||
|
|
||||||
|
if start_byte > 0 and status_code != 206:
|
||||||
|
# 服务器忽略了 Range,返回完整内容:从头开始写
|
||||||
|
if os.path.exists(temp_path):
|
||||||
|
os.remove(temp_path)
|
||||||
|
start_byte = 0
|
||||||
|
|
||||||
|
mode = 'ab' if start_byte > 0 else 'wb'
|
||||||
|
|
||||||
|
if start_byte > 0:
|
||||||
|
# 续传: 先把已有前缀复制进临时文件,再追加尾部数据,
|
||||||
|
# 避免用"仅尾部"覆盖完整文件导致数据损坏
|
||||||
|
shutil.copyfile(local_path, temp_path)
|
||||||
|
|
||||||
with open(temp_path, mode) as f:
|
with open(temp_path, mode) as f:
|
||||||
while True:
|
while True:
|
||||||
@@ -1159,6 +1187,14 @@ class MirrorSyncManager:
|
|||||||
f.write(chunk)
|
f.write(chunk)
|
||||||
|
|
||||||
if self.running:
|
if self.running:
|
||||||
|
# 校验续传结果: 续传后文件应不小于原有大小
|
||||||
|
if start_byte > 0:
|
||||||
|
try:
|
||||||
|
final_size = os.path.getsize(temp_path)
|
||||||
|
if final_size < start_byte:
|
||||||
|
raise Exception(f"续传后文件异常变小: {final_size} < {start_byte}")
|
||||||
|
except OSError:
|
||||||
|
raise
|
||||||
if os.path.exists(local_path):
|
if os.path.exists(local_path):
|
||||||
os.remove(local_path)
|
os.remove(local_path)
|
||||||
os.rename(temp_path, local_path)
|
os.rename(temp_path, local_path)
|
||||||
|
|||||||
+17
-4
@@ -689,12 +689,25 @@ class SyncEngine:
|
|||||||
req.add_header('Range', f'bytes={start_byte}-')
|
req.add_header('Range', f'bytes={start_byte}-')
|
||||||
|
|
||||||
timeout = config.get('timeout', 30)
|
timeout = config.get('timeout', 30)
|
||||||
temp_path = local_path + '.tmp'
|
temp_path = local_path + '.tmp' # 提前初始化,避免 except 分支 NameError
|
||||||
mode = 'ab' if start_byte > 0 else 'wb'
|
|
||||||
|
|
||||||
with urllib.request.urlopen(req, timeout=timeout) as response:
|
with urllib.request.urlopen(req, timeout=timeout) as response:
|
||||||
if response.getcode() not in [200, 206]:
|
status_code = response.getcode()
|
||||||
raise Exception(f"HTTP错误: {response.getcode()}")
|
if status_code not in [200, 206]:
|
||||||
|
raise Exception(f"HTTP错误: {status_code}")
|
||||||
|
|
||||||
|
if start_byte > 0 and status_code != 206:
|
||||||
|
# 服务器忽略了 Range,返回完整内容:从头开始写
|
||||||
|
if os.path.exists(temp_path):
|
||||||
|
os.remove(temp_path)
|
||||||
|
start_byte = 0
|
||||||
|
|
||||||
|
mode = 'ab' if start_byte > 0 else 'wb'
|
||||||
|
|
||||||
|
if start_byte > 0:
|
||||||
|
# 续传: 先把已有前缀复制进临时文件,再追加尾部数据,避免损坏
|
||||||
|
import shutil
|
||||||
|
shutil.copyfile(local_path, temp_path)
|
||||||
|
|
||||||
total_size = int(response.headers.get('Content-Length', 0)) + start_byte
|
total_size = int(response.headers.get('Content-Length', 0)) + start_byte
|
||||||
task.total_size = total_size
|
task.total_size = total_size
|
||||||
|
|||||||
Reference in New Issue
Block a user