v0.2.2 performance boost smb
This commit is contained in:
+171
-69
@@ -128,7 +128,7 @@ class FlaskServer:
|
||||
}
|
||||
return meta
|
||||
|
||||
def _start_backup_run(self, backup_id: int) -> tuple[bool, str]:
|
||||
def _start_backup_run(self, backup_id: int, trigger: str = "manual") -> tuple[bool, str]:
|
||||
with self._run_lock:
|
||||
existing = self._active_runs.get(backup_id)
|
||||
if existing and bool(existing.get("is_running", False)):
|
||||
@@ -160,12 +160,16 @@ class FlaskServer:
|
||||
name=f"backup-run-{backup_id}",
|
||||
)
|
||||
|
||||
trigger_normalized = (trigger or "manual").strip().lower()
|
||||
is_scheduled = trigger_normalized == "scheduled"
|
||||
start_label = "Geplanter Lauf" if is_scheduled else "Manueller Lauf"
|
||||
|
||||
self._set_run_state(
|
||||
backup_id,
|
||||
is_running=True,
|
||||
progress_percent=2,
|
||||
progress_text="Startet...",
|
||||
status_message="Manueller Lauf wird gestartet...",
|
||||
status_message=f"{start_label} wird gestartet...",
|
||||
stop_event=stop_event,
|
||||
thread=thread,
|
||||
)
|
||||
@@ -175,9 +179,10 @@ class FlaskServer:
|
||||
category="backup",
|
||||
event="run_started",
|
||||
backup_id=backup_id,
|
||||
message=f"Manueller Backup-Lauf wurde gestartet: {job.name}",
|
||||
message=f"{start_label} wurde gestartet: {job.name}",
|
||||
details={"trigger": trigger_normalized},
|
||||
)
|
||||
return True, f"Manueller Lauf für '{job.name}' wurde gestartet."
|
||||
return True, f"{start_label} für '{job.name}' wurde gestartet."
|
||||
|
||||
def _stop_backup_run(self, backup_id: int) -> tuple[bool, str]:
|
||||
with self._run_lock:
|
||||
@@ -382,84 +387,181 @@ class FlaskServer:
|
||||
destination_location: str,
|
||||
stop_event: threading.Event,
|
||||
) -> int:
|
||||
# Fast path: SMB server-side copy when source and destination are on the same host.
|
||||
if source_protocol == "SMB" and destination_protocol == "SMB":
|
||||
src_parts = source_location.split("|", 2)
|
||||
dst_parts = destination_location.split("|", 2)
|
||||
if len(src_parts) == 3 and len(dst_parts) == 3:
|
||||
src_host, src_share, src_subpath = src_parts
|
||||
dst_host, dst_share, dst_subpath = dst_parts
|
||||
if src_host.strip().lower() == dst_host.strip().lower():
|
||||
if stop_event.is_set():
|
||||
raise RuntimeError("Lauf wurde abgebrochen.")
|
||||
src_unc = self._smb_unc_path(src_host, src_share, src_subpath)
|
||||
dst_unc = self._smb_unc_path(dst_host, dst_share, dst_subpath)
|
||||
try:
|
||||
smbclient.copyfile(src_unc, dst_unc)
|
||||
try:
|
||||
return int(smbclient.stat(dst_unc).st_size or 0)
|
||||
except Exception: # noqa: BLE001
|
||||
return int(smbclient.stat(src_unc).st_size or 0)
|
||||
except Exception: # noqa: BLE001
|
||||
# Fallback to streamed copy for incompatible servers or share policies.
|
||||
pass
|
||||
|
||||
total_bytes = 0
|
||||
source_handle = None
|
||||
destination_handle = None
|
||||
cleanup_callbacks: list = []
|
||||
|
||||
if source_protocol == "LOCAL":
|
||||
source_handle = open(source_location, "rb")
|
||||
elif source_protocol in {"FTP", "SFTP"}:
|
||||
# FTP/SFTP: download file to an in-memory buffer first
|
||||
import io
|
||||
from ftplib import FTP as _FTP
|
||||
buf = io.BytesIO()
|
||||
ftp_dl = _FTP()
|
||||
# source_location is the FTP path; connection info comes from the job context
|
||||
# We store it as "host:port|username|password|path" for FTP sources
|
||||
parts = source_location.split("|", 3)
|
||||
if len(parts) == 4:
|
||||
ftp_host_port, ftp_user, ftp_pass, ftp_path = parts
|
||||
_h, _p = (ftp_host_port.rsplit(":", 1) + [None])[:2]
|
||||
ftp_dl.connect(_h, int(_p) if _p else 21, timeout=60)
|
||||
ftp_dl.login(ftp_user, ftp_pass)
|
||||
ftp_dl.retrbinary(f"RETR {ftp_path}", buf.write)
|
||||
ftp_dl.quit()
|
||||
buf.seek(0)
|
||||
source_handle = buf
|
||||
else:
|
||||
raise RuntimeError(f"Ungültiger FTP-Source-Pfad: {source_location}")
|
||||
else:
|
||||
source_handle = smbclient.open_file(
|
||||
self._smb_unc_path(*source_location.split("|", 2)),
|
||||
mode="rb",
|
||||
)
|
||||
class _StopAwareReader:
|
||||
def __init__(self, wrapped, stopper: threading.Event):
|
||||
self._wrapped = wrapped
|
||||
self._stopper = stopper
|
||||
|
||||
if destination_protocol == "LOCAL":
|
||||
destination_handle = open(destination_location, "wb")
|
||||
elif destination_protocol in {"FTP", "SFTP"}:
|
||||
import io as _io
|
||||
from ftplib import FTP as _FTP
|
||||
parts = destination_location.split("|", 3)
|
||||
def read(self, size: int = -1):
|
||||
if self._stopper.is_set():
|
||||
raise RuntimeError("Lauf wurde abgebrochen.")
|
||||
return self._wrapped.read(size)
|
||||
|
||||
def _parse_endpoint(endpoint: str, default_port: int) -> tuple[str, int, str, str, str]:
|
||||
parts = endpoint.split("|", 3)
|
||||
if len(parts) != 4:
|
||||
raise RuntimeError(f"Ungültiger FTP-Ziel-Pfad: {destination_location}")
|
||||
ftp_host_port, ftp_user, ftp_pass, ftp_path = parts
|
||||
_h, _p = (ftp_host_port.rsplit(":", 1) + [None])[:2]
|
||||
ftp_dest = _FTP()
|
||||
try:
|
||||
ftp_dest.connect(_h, int(_p) if _p else 21, timeout=60)
|
||||
ftp_dest.login(ftp_user, ftp_pass)
|
||||
# Ensure parent directories exist
|
||||
parent_dirs = [d for d in (ftp_path.rsplit("/", 1)[0] if "/" in ftp_path else "").split("/") if d]
|
||||
cur = ""
|
||||
for _d in parent_dirs:
|
||||
cur = f"{cur}/{_d}"
|
||||
try: ftp_dest.mkd(cur)
|
||||
except Exception: pass
|
||||
src_buf = _io.BytesIO(source_handle.read())
|
||||
source_handle.close()
|
||||
ftp_dest.storbinary(f"STOR {ftp_path}", src_buf)
|
||||
return len(src_buf.getvalue())
|
||||
except Exception as _exc:
|
||||
raise RuntimeError(f"FTP-Upload fehlgeschlagen: {_exc}") from _exc
|
||||
finally:
|
||||
try: ftp_dest.quit()
|
||||
except Exception: pass
|
||||
else:
|
||||
destination_handle = smbclient.open_file(
|
||||
self._smb_unc_path(*destination_location.split("|", 2)),
|
||||
mode="wb",
|
||||
)
|
||||
raise RuntimeError(f"Ungültiger Verbindungs-String: {endpoint}")
|
||||
host_port, user, password, remote_path = parts
|
||||
host, port_text = (host_port.rsplit(":", 1) + [None])[:2]
|
||||
port = int(port_text) if port_text else default_port
|
||||
return host, port, user, password, remote_path
|
||||
|
||||
try:
|
||||
if source_protocol == "LOCAL":
|
||||
source_handle = open(source_location, "rb")
|
||||
elif source_protocol == "FTP":
|
||||
from ftplib import FTP as _FTP
|
||||
|
||||
host, port, user, password, remote_path = _parse_endpoint(source_location, 21)
|
||||
ftp_dl = _FTP()
|
||||
spool = tempfile.SpooledTemporaryFile(max_size=32 * 1024 * 1024, mode="w+b")
|
||||
cleanup_callbacks.append(spool.close)
|
||||
try:
|
||||
ftp_dl.connect(host, port, timeout=60)
|
||||
ftp_dl.login(user, password)
|
||||
ftp_dl.retrbinary(f"RETR {remote_path}", spool.write, blocksize=4 * 1024 * 1024)
|
||||
finally:
|
||||
try:
|
||||
ftp_dl.quit()
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
spool.seek(0)
|
||||
source_handle = spool
|
||||
elif source_protocol == "SFTP":
|
||||
import paramiko
|
||||
|
||||
host, port, user, password, remote_path = _parse_endpoint(source_location, 22)
|
||||
transport = paramiko.Transport((host, port))
|
||||
transport.connect(username=user, password=password)
|
||||
sftp_client = paramiko.SFTPClient.from_transport(transport)
|
||||
cleanup_callbacks.append(sftp_client.close)
|
||||
cleanup_callbacks.append(transport.close)
|
||||
source_handle = sftp_client.open(remote_path, mode="rb")
|
||||
else:
|
||||
source_handle = smbclient.open_file(
|
||||
self._smb_unc_path(*source_location.split("|", 2)),
|
||||
mode="rb",
|
||||
)
|
||||
|
||||
if destination_protocol == "LOCAL":
|
||||
destination_handle = open(destination_location, "wb")
|
||||
elif destination_protocol == "FTP":
|
||||
from ftplib import FTP as _FTP
|
||||
|
||||
host, port, user, password, remote_path = _parse_endpoint(destination_location, 21)
|
||||
ftp_dest = _FTP()
|
||||
try:
|
||||
ftp_dest.connect(host, port, timeout=60)
|
||||
ftp_dest.login(user, password)
|
||||
parent_dirs = [d for d in (remote_path.rsplit("/", 1)[0] if "/" in remote_path else "").split("/") if d]
|
||||
current = ""
|
||||
for directory in parent_dirs:
|
||||
current = f"{current}/{directory}"
|
||||
try:
|
||||
ftp_dest.mkd(current)
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
|
||||
stream = _StopAwareReader(source_handle, stop_event)
|
||||
ftp_dest.storbinary(f"STOR {remote_path}", stream, blocksize=4 * 1024 * 1024)
|
||||
try:
|
||||
total_bytes = int(ftp_dest.size(remote_path) or 0)
|
||||
except Exception: # noqa: BLE001
|
||||
total_bytes = 0
|
||||
if total_bytes <= 0:
|
||||
source_handle.seek(0)
|
||||
while True:
|
||||
chunk = source_handle.read(4 * 1024 * 1024)
|
||||
if not chunk:
|
||||
break
|
||||
total_bytes += len(chunk)
|
||||
return total_bytes
|
||||
except Exception as exc:
|
||||
raise RuntimeError(f"FTP-Upload fehlgeschlagen: {exc}") from exc
|
||||
finally:
|
||||
try:
|
||||
ftp_dest.quit()
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
elif destination_protocol == "SFTP":
|
||||
import paramiko
|
||||
|
||||
host, port, user, password, remote_path = _parse_endpoint(destination_location, 22)
|
||||
transport = paramiko.Transport((host, port))
|
||||
transport.connect(username=user, password=password)
|
||||
sftp_client = paramiko.SFTPClient.from_transport(transport)
|
||||
cleanup_callbacks.append(sftp_client.close)
|
||||
cleanup_callbacks.append(transport.close)
|
||||
|
||||
parent_dir = posixpath.dirname(remote_path)
|
||||
if parent_dir:
|
||||
parts = [p for p in parent_dir.split("/") if p]
|
||||
current = ""
|
||||
for part in parts:
|
||||
current = f"{current}/{part}"
|
||||
try:
|
||||
sftp_client.mkdir(current)
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
|
||||
destination_handle = sftp_client.open(remote_path, mode="wb")
|
||||
else:
|
||||
destination_handle = smbclient.open_file(
|
||||
self._smb_unc_path(*destination_location.split("|", 2)),
|
||||
mode="wb",
|
||||
)
|
||||
|
||||
while True:
|
||||
if stop_event.is_set():
|
||||
raise RuntimeError("Lauf wurde abgebrochen.")
|
||||
chunk = source_handle.read(1024 * 1024)
|
||||
chunk = source_handle.read(4 * 1024 * 1024)
|
||||
if not chunk:
|
||||
break
|
||||
destination_handle.write(chunk)
|
||||
total_bytes += len(chunk)
|
||||
finally:
|
||||
source_handle.close()
|
||||
destination_handle.close()
|
||||
if source_handle is not None:
|
||||
try:
|
||||
source_handle.close()
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
if destination_handle is not None:
|
||||
try:
|
||||
destination_handle.close()
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
for cb in reversed(cleanup_callbacks):
|
||||
try:
|
||||
cb()
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
|
||||
return total_bytes
|
||||
|
||||
@@ -2507,7 +2609,7 @@ class FlaskServer:
|
||||
if fired.get(b.backup_id) == sched_key:
|
||||
continue
|
||||
fired[b.backup_id] = sched_key
|
||||
ok, start_message = self._start_backup_run(b.backup_id)
|
||||
ok, start_message = self._start_backup_run(b.backup_id, trigger="scheduled")
|
||||
self._log(
|
||||
level="INFO" if ok else "WARNING",
|
||||
category="scheduler",
|
||||
|
||||
Reference in New Issue
Block a user