171 lines
5.4 KiB
Python
171 lines
5.4 KiB
Python
import os
|
|
import shutil
|
|
import time
|
|
from abc import ABC, abstractmethod
|
|
from .samba_manager import SambaManager
|
|
|
|
class Storage(ABC):
|
|
@abstractmethod
|
|
def exists(self, path): pass
|
|
@abstractmethod
|
|
def get_size(self, path): pass
|
|
@abstractmethod
|
|
def makedirs(self, path): pass
|
|
@abstractmethod
|
|
def rename(self, src, dst): pass
|
|
@abstractmethod
|
|
def delete(self, path): pass
|
|
@abstractmethod
|
|
def save(self, path, response, progress_callback=None): pass
|
|
@abstractmethod
|
|
def get_full_path(self, path): pass
|
|
|
|
class LocalStorage(Storage):
|
|
def exists(self, path):
|
|
return os.path.exists(path)
|
|
|
|
def get_size(self, path):
|
|
if os.path.exists(path):
|
|
return os.path.getsize(path)
|
|
return 0
|
|
|
|
def makedirs(self, path):
|
|
os.makedirs(path, exist_ok=True)
|
|
|
|
def rename(self, src, dst):
|
|
if os.path.exists(dst):
|
|
os.remove(dst)
|
|
os.rename(src, dst)
|
|
|
|
def delete(self, path):
|
|
if os.path.exists(path):
|
|
os.remove(path)
|
|
|
|
def get_full_path(self, path):
|
|
return os.path.abspath(path)
|
|
|
|
def save(self, path, response, progress_callback=None):
|
|
# Local save supports chunked writing which allows easy progress update
|
|
# We assume response is a requests response object with stream=True
|
|
|
|
mode = 'wb'
|
|
# Check if we are resuming? LocalStorage in legacy code handled this outside save.
|
|
# But here we are encapsulating the write loop.
|
|
# If headers Range is set, we might need 'ab'.
|
|
# But for simplicity, we assume 'wb' unless specific logic is added.
|
|
# Downloader logic passed 'ab' if resuming.
|
|
# We can inspect response.request.headers['Range']?
|
|
# Or just take a mode arg.
|
|
# For now, let's implement standard write. Resume support can be added if needed via mode arg.
|
|
|
|
downloaded = 0
|
|
total_size = int(response.headers.get('content-length', 0))
|
|
|
|
with open(path, 'wb') as f:
|
|
for chunk in response.iter_content(chunk_size=1048576):
|
|
if chunk:
|
|
f.write(chunk)
|
|
downloaded += len(chunk)
|
|
if progress_callback:
|
|
progress_callback(len(chunk))
|
|
|
|
class ProgressReader:
|
|
def __init__(self, raw_stream, callback):
|
|
self.raw_stream = raw_stream
|
|
self.callback = callback
|
|
|
|
def read(self, size=-1):
|
|
chunk = self.raw_stream.read(size)
|
|
if chunk and self.callback:
|
|
self.callback(len(chunk))
|
|
return chunk
|
|
|
|
class SambaStorage(Storage):
|
|
def __init__(self, samba_manager: SambaManager, base_path=""):
|
|
self.smb = samba_manager
|
|
self.base_path = base_path # e.g. "videos/coomerparty"
|
|
|
|
def _full_path(self, path):
|
|
# path comes from Downloader as "user/img/file.jpg"
|
|
# we append to base_path
|
|
if self.base_path:
|
|
return f"{self.base_path}/{path}".replace("\\", "/")
|
|
return path.replace("\\", "/")
|
|
|
|
def get_full_path(self, path):
|
|
# Return a string representing the full smb path
|
|
inner_path = self._full_path(path)
|
|
return f"smb://{self.smb.server_ip}/{self.smb.share_name}/{inner_path}"
|
|
|
|
def _exec(self, func, *args):
|
|
# Helper to execute with a fresh connection
|
|
conn = self.smb.clone()
|
|
try:
|
|
return func(conn, *args)
|
|
finally:
|
|
conn.close()
|
|
|
|
def exists(self, path):
|
|
def _op(smb):
|
|
try:
|
|
attr = smb.get_attributes(self._full_path(path))
|
|
if isinstance(attr, dict) and "error" in attr:
|
|
return False
|
|
return True
|
|
except:
|
|
return False
|
|
return self._exec(_op)
|
|
|
|
def get_size(self, path):
|
|
def _op(smb):
|
|
try:
|
|
attr = smb.get_attributes(self._full_path(path))
|
|
return attr.file_size
|
|
except:
|
|
return 0
|
|
return self._exec(_op)
|
|
|
|
def makedirs(self, path):
|
|
full = self._full_path(path)
|
|
parts = full.split("/")
|
|
def _op(smb):
|
|
current = ""
|
|
for part in parts:
|
|
if not part: continue
|
|
current = f"{current}/{part}" if current else part
|
|
try:
|
|
smb.create_directory(current)
|
|
except:
|
|
pass # Ignore if exists
|
|
self._exec(_op)
|
|
|
|
def rename(self, src, dst):
|
|
full_src = self._full_path(src)
|
|
full_dst = self._full_path(dst)
|
|
def _op(smb):
|
|
try:
|
|
# Check if dst exists, delete if so
|
|
try:
|
|
smb.get_attributes(full_dst)
|
|
smb.delete_file(full_dst)
|
|
except:
|
|
pass
|
|
|
|
smb.rename_file(full_src, full_dst)
|
|
except Exception as e:
|
|
raise e
|
|
self._exec(_op)
|
|
|
|
def delete(self, path):
|
|
def _op(smb):
|
|
smb.delete_file(self._full_path(path))
|
|
self._exec(_op)
|
|
|
|
def save(self, path, response, progress_callback=None):
|
|
full_path = self._full_path(path)
|
|
# Use ProgressReader to wrap response.raw
|
|
reader = ProgressReader(response.raw, progress_callback)
|
|
def _op(smb):
|
|
smb.upload_file(full_path, reader)
|
|
self._exec(_op)
|