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)