Spaces:
Running
Running
Update storage_handler.py
Browse files- storage_handler.py +174 -54
storage_handler.py
CHANGED
|
@@ -1,7 +1,7 @@
|
|
| 1 |
import os
|
| 2 |
import base64
|
| 3 |
-
import
|
| 4 |
-
import
|
| 5 |
import logging
|
| 6 |
from typing import Optional, BinaryIO, Dict, Any, List
|
| 7 |
import time
|
|
@@ -16,17 +16,25 @@ logger = logging.getLogger("InfinityChat.Storage")
|
|
| 16 |
# Configuration
|
| 17 |
# ------------------------------------------------------------------------
|
| 18 |
HF_TOKEN = os.environ.get("HF_TOKEN")
|
| 19 |
-
if not HF_TOKEN:
|
| 20 |
-
raise ValueError("β HF_TOKEN environment variable is required")
|
| 21 |
-
|
| 22 |
BUCKET_NAME = os.environ.get("INFINITY_CHAT_BUCKET", "infinitychat-data")
|
| 23 |
FILE_ENCRYPTION_KEY_B64 = os.environ.get("FILE_ENCRYPTION_KEY", None)
|
| 24 |
|
|
|
|
|
|
|
|
|
|
|
|
|
| 25 |
if FILE_ENCRYPTION_KEY_B64 is None:
|
| 26 |
-
|
| 27 |
-
|
| 28 |
-
|
| 29 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 30 |
|
| 31 |
try:
|
| 32 |
FILE_ENCRYPTION_KEY = base64.urlsafe_b64decode(FILE_ENCRYPTION_KEY_B64.encode())
|
|
@@ -36,37 +44,64 @@ except Exception as e:
|
|
| 36 |
assert len(FILE_ENCRYPTION_KEY) == 32, "FILE_ENCRYPTION_KEY must decode to exactly 32 bytes for AES-256"
|
| 37 |
|
| 38 |
# ------------------------------------------------------------------------
|
| 39 |
-
#
|
| 40 |
# ------------------------------------------------------------------------
|
| 41 |
-
|
| 42 |
-
|
| 43 |
-
owner_info = _api.whoami()
|
| 44 |
-
OWNER = owner_info["name"]
|
| 45 |
-
except Exception as e:
|
| 46 |
-
logger.warning(f"Could not determine HF username: {e}")
|
| 47 |
-
OWNER = "infinitychat"
|
| 48 |
|
| 49 |
-
if "/" in BUCKET_NAME:
|
| 50 |
-
OWNER, BUCKET_ID = BUCKET_NAME.split("/", 1)
|
| 51 |
-
else:
|
| 52 |
-
BUCKET_ID = BUCKET_NAME
|
| 53 |
|
| 54 |
-
|
| 55 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 56 |
|
| 57 |
# ------------------------------------------------------------------------
|
| 58 |
-
#
|
| 59 |
# ------------------------------------------------------------------------
|
| 60 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 61 |
|
| 62 |
# ------------------------------------------------------------------------
|
| 63 |
# Bucket Initialization
|
| 64 |
# ------------------------------------------------------------------------
|
| 65 |
def ensure_bucket():
|
| 66 |
"""Create the private bucket if it doesn't exist."""
|
|
|
|
|
|
|
| 67 |
global OWNER, BUCKET_URI, BUCKET_ID
|
|
|
|
| 68 |
try:
|
| 69 |
-
|
| 70 |
bucket_id=f"{OWNER}/{BUCKET_ID}",
|
| 71 |
private=True,
|
| 72 |
exist_ok=True
|
|
@@ -80,7 +115,7 @@ def ensure_bucket():
|
|
| 80 |
logger.warning(f"β οΈ Cannot create bucket - permission issue: {e}")
|
| 81 |
else:
|
| 82 |
try:
|
| 83 |
-
|
| 84 |
bucket_id=BUCKET_ID,
|
| 85 |
private=True,
|
| 86 |
exist_ok=True
|
|
@@ -95,12 +130,17 @@ def ensure_bucket():
|
|
| 95 |
else:
|
| 96 |
logger.error(f"β Failed to create bucket: {e2}")
|
| 97 |
|
| 98 |
-
|
|
|
|
|
|
|
| 99 |
|
| 100 |
# ------------------------------------------------------------------------
|
| 101 |
# Heartbeat sync
|
| 102 |
# ------------------------------------------------------------------------
|
| 103 |
def _sync_bucket_interval():
|
|
|
|
|
|
|
|
|
|
| 104 |
def _sync():
|
| 105 |
while True:
|
| 106 |
time.sleep(60)
|
|
@@ -115,6 +155,7 @@ def _sync_bucket_interval():
|
|
| 115 |
t = threading.Thread(target=_sync, daemon=True)
|
| 116 |
t.start()
|
| 117 |
|
|
|
|
| 118 |
_sync_bucket_interval()
|
| 119 |
|
| 120 |
# ------------------------------------------------------------------------
|
|
@@ -126,12 +167,14 @@ def encrypt_bytes(data: bytes, aad: Optional[bytes] = None) -> bytes:
|
|
| 126 |
ciphertext = aesgcm.encrypt(nonce, data, aad or b"")
|
| 127 |
return nonce + ciphertext
|
| 128 |
|
|
|
|
| 129 |
def decrypt_bytes(encrypted_blob: bytes, aad: Optional[bytes] = None) -> bytes:
|
| 130 |
aesgcm = AESGCM(FILE_ENCRYPTION_KEY)
|
| 131 |
nonce = encrypted_blob[:12]
|
| 132 |
ciphertext = encrypted_blob[12:]
|
| 133 |
return aesgcm.decrypt(nonce, ciphertext, aad or b"")
|
| 134 |
|
|
|
|
| 135 |
def verify_file_integrity(encrypted_blob: bytes) -> bool:
|
| 136 |
return len(encrypted_blob) >= 29
|
| 137 |
|
|
@@ -139,8 +182,12 @@ def verify_file_integrity(encrypted_blob: bytes) -> bool:
|
|
| 139 |
# Path Utilities
|
| 140 |
# ------------------------------------------------------------------------
|
| 141 |
def _bucket_path(remote_path: str) -> str:
|
|
|
|
|
|
|
|
|
|
| 142 |
return f"{BUCKET_URI}/{remote_path.lstrip('/')}"
|
| 143 |
|
|
|
|
| 144 |
def _validate_path(remote_path: str) -> None:
|
| 145 |
if not remote_path:
|
| 146 |
raise ValueError("Path cannot be empty")
|
|
@@ -156,12 +203,14 @@ def _validate_path(remote_path: str) -> None:
|
|
| 156 |
# ------------------------------------------------------------------------
|
| 157 |
def store_file(remote_path: str, data: bytes, encrypt: bool = True) -> str:
|
| 158 |
_validate_path(remote_path)
|
|
|
|
|
|
|
| 159 |
try:
|
| 160 |
payload = encrypt_bytes(data, remote_path.encode('utf-8')) if encrypt else data
|
| 161 |
full_uri = _bucket_path(remote_path)
|
| 162 |
-
with
|
| 163 |
f.write(payload)
|
| 164 |
-
if not
|
| 165 |
raise IOError(f"Failed to verify file was stored: {full_uri}")
|
| 166 |
logger.debug(f"πΎ Stored: {remote_path} ({len(data)} bytes)")
|
| 167 |
return remote_path
|
|
@@ -169,13 +218,16 @@ def store_file(remote_path: str, data: bytes, encrypt: bool = True) -> str:
|
|
| 169 |
logger.error(f"β Failed to store {remote_path}: {e}")
|
| 170 |
raise IOError(f"Storage failed: {e}")
|
| 171 |
|
|
|
|
| 172 |
def retrieve_file(remote_path: str, decrypt: bool = True) -> bytes:
|
| 173 |
_validate_path(remote_path)
|
|
|
|
|
|
|
| 174 |
full_uri = _bucket_path(remote_path)
|
| 175 |
try:
|
| 176 |
-
if not
|
| 177 |
raise FileNotFoundError(f"File not found: {remote_path}")
|
| 178 |
-
with
|
| 179 |
payload = f.read()
|
| 180 |
if not payload:
|
| 181 |
raise IOError(f"Empty file: {remote_path}")
|
|
@@ -192,12 +244,16 @@ def retrieve_file(remote_path: str, decrypt: bool = True) -> bytes:
|
|
| 192 |
logger.error(f"β Failed to retrieve {remote_path}: {e}")
|
| 193 |
raise IOError(f"Retrieval failed: {e}")
|
| 194 |
|
|
|
|
| 195 |
def delete_file(remote_path: str) -> bool:
|
| 196 |
_validate_path(remote_path)
|
|
|
|
|
|
|
|
|
|
| 197 |
full_uri = _bucket_path(remote_path)
|
| 198 |
try:
|
| 199 |
-
if
|
| 200 |
-
|
| 201 |
logger.debug(f"ποΈ Deleted: {remote_path}")
|
| 202 |
return True
|
| 203 |
logger.warning(f"β οΈ Not found for deletion: {remote_path}")
|
|
@@ -206,37 +262,51 @@ def delete_file(remote_path: str) -> bool:
|
|
| 206 |
logger.error(f"β Failed to delete {remote_path}: {e}")
|
| 207 |
raise IOError(f"Deletion failed: {e}")
|
| 208 |
|
|
|
|
| 209 |
def file_exists(remote_path: str) -> bool:
|
| 210 |
-
|
|
|
|
| 211 |
try:
|
| 212 |
-
|
|
|
|
| 213 |
except Exception:
|
| 214 |
return False
|
| 215 |
|
|
|
|
| 216 |
def list_files(prefix: str = "", recursive: bool = True) -> List[str]:
|
| 217 |
-
|
|
|
|
| 218 |
try:
|
| 219 |
-
|
|
|
|
|
|
|
| 220 |
prefix_len = len(BUCKET_URI) + 1
|
| 221 |
-
return [item[prefix_len:] for item in items if not
|
| 222 |
except FileNotFoundError:
|
| 223 |
return []
|
| 224 |
except Exception as e:
|
| 225 |
logger.error(f"Failed to list files: {e}")
|
| 226 |
return []
|
| 227 |
|
|
|
|
| 228 |
def get_file_size(remote_path: str) -> Optional[int]:
|
| 229 |
-
|
|
|
|
| 230 |
try:
|
| 231 |
-
|
|
|
|
| 232 |
return info.get("size")
|
| 233 |
except Exception:
|
| 234 |
return None
|
| 235 |
|
|
|
|
| 236 |
def get_file_info(remote_path: str) -> Optional[Dict[str, Any]]:
|
| 237 |
-
|
|
|
|
| 238 |
try:
|
| 239 |
-
|
|
|
|
|
|
|
| 240 |
return {
|
| 241 |
"name": remote_path,
|
| 242 |
"size": info.get("size"),
|
|
@@ -247,25 +317,33 @@ def get_file_info(remote_path: str) -> Optional[Dict[str, Any]]:
|
|
| 247 |
except FileNotFoundError:
|
| 248 |
return None
|
| 249 |
except Exception as e:
|
| 250 |
-
logger.error(f"Failed to get file info: {e}")
|
| 251 |
return None
|
| 252 |
|
|
|
|
| 253 |
def store_file_stream(remote_path: str, data_stream: BinaryIO, encrypt: bool = True) -> str:
|
| 254 |
return store_file(remote_path, data_stream.read(), encrypt=encrypt)
|
| 255 |
|
|
|
|
| 256 |
def retrieve_file_stream(remote_path: str, decrypt: bool = True) -> BinaryIO:
|
| 257 |
import io
|
| 258 |
return io.BytesIO(retrieve_file(remote_path, decrypt=decrypt))
|
| 259 |
|
|
|
|
| 260 |
def get_storage_stats() -> Dict[str, Any]:
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 261 |
try:
|
| 262 |
-
|
|
|
|
| 263 |
total_size = 0
|
| 264 |
file_count = 0
|
| 265 |
-
for
|
| 266 |
-
|
| 267 |
-
|
| 268 |
-
total_size += size
|
| 269 |
file_count += 1
|
| 270 |
return {
|
| 271 |
"bucket": f"{OWNER}/{BUCKET_ID}",
|
|
@@ -283,7 +361,9 @@ def get_storage_stats() -> Dict[str, Any]:
|
|
| 283 |
"error": str(e)
|
| 284 |
}
|
| 285 |
|
|
|
|
| 286 |
def create_backup(backup_prefix: str = "backups") -> str:
|
|
|
|
| 287 |
timestamp = int(time.time())
|
| 288 |
backup_path = f"{backup_prefix}/backup_{timestamp}"
|
| 289 |
try:
|
|
@@ -308,11 +388,15 @@ _db_sync_lock = threading.Lock()
|
|
| 308 |
_last_db_sync = 0
|
| 309 |
DB_SYNC_INTERVAL = 30 # seconds
|
| 310 |
|
|
|
|
| 311 |
def download_database(local_path: str) -> bool:
|
| 312 |
"""
|
| 313 |
Download database from bucket to local path.
|
| 314 |
Returns True if database was found and downloaded.
|
| 315 |
"""
|
|
|
|
|
|
|
|
|
|
| 316 |
try:
|
| 317 |
if file_exists(DB_BUCKET_PATH):
|
| 318 |
logger.info("π₯ Downloading database from bucket...")
|
|
@@ -329,21 +413,44 @@ def download_database(local_path: str) -> bool:
|
|
| 329 |
logger.error(f"β Failed to download database: {e}")
|
| 330 |
return False
|
| 331 |
|
|
|
|
| 332 |
def upload_database(local_path: str) -> bool:
|
| 333 |
"""
|
| 334 |
-
Upload local database
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 335 |
Returns True on success.
|
| 336 |
"""
|
| 337 |
global _last_db_sync
|
|
|
|
|
|
|
| 338 |
with _db_sync_lock:
|
|
|
|
| 339 |
try:
|
| 340 |
if not os.path.exists(local_path):
|
| 341 |
logger.warning(f"β οΈ Database not found at {local_path}, skipping upload")
|
| 342 |
return False
|
| 343 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 344 |
data = f.read()
|
| 345 |
if not data:
|
| 346 |
-
logger.warning("β οΈ Database
|
| 347 |
return False
|
| 348 |
store_file(DB_BUCKET_PATH, data)
|
| 349 |
_last_db_sync = time.time()
|
|
@@ -352,11 +459,22 @@ def upload_database(local_path: str) -> bool:
|
|
| 352 |
except Exception as e:
|
| 353 |
logger.error(f"β Failed to upload database: {e}")
|
| 354 |
return False
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 355 |
|
| 356 |
def start_db_sync(local_path: str):
|
| 357 |
"""
|
| 358 |
Start background thread that periodically uploads the database to bucket.
|
| 359 |
"""
|
|
|
|
|
|
|
|
|
|
|
|
|
| 360 |
def _sync_loop():
|
| 361 |
while True:
|
| 362 |
time.sleep(DB_SYNC_INTERVAL)
|
|
@@ -369,9 +487,11 @@ def start_db_sync(local_path: str):
|
|
| 369 |
t.start()
|
| 370 |
logger.info(f"π Database auto-sync started (every {DB_SYNC_INTERVAL}s)")
|
| 371 |
|
|
|
|
| 372 |
def close():
|
| 373 |
try:
|
| 374 |
-
|
|
|
|
| 375 |
except Exception:
|
| 376 |
pass
|
| 377 |
|
|
@@ -384,5 +504,5 @@ __all__ = [
|
|
| 384 |
'store_file_stream', 'retrieve_file_stream',
|
| 385 |
'get_storage_stats', 'create_backup', 'close',
|
| 386 |
'download_database', 'upload_database', 'start_db_sync',
|
| 387 |
-
'DB_BUCKET_PATH'
|
| 388 |
]
|
|
|
|
| 1 |
import os
|
| 2 |
import base64
|
| 3 |
+
import sqlite3
|
| 4 |
+
import tempfile
|
| 5 |
import logging
|
| 6 |
from typing import Optional, BinaryIO, Dict, Any, List
|
| 7 |
import time
|
|
|
|
| 16 |
# Configuration
|
| 17 |
# ------------------------------------------------------------------------
|
| 18 |
HF_TOKEN = os.environ.get("HF_TOKEN")
|
|
|
|
|
|
|
|
|
|
| 19 |
BUCKET_NAME = os.environ.get("INFINITY_CHAT_BUCKET", "infinitychat-data")
|
| 20 |
FILE_ENCRYPTION_KEY_B64 = os.environ.get("FILE_ENCRYPTION_KEY", None)
|
| 21 |
|
| 22 |
+
# Offline mode: lets the app run without a Hugging Face token (dev/tests).
|
| 23 |
+
# File storage endpoints will return a clear error instead of crashing boot.
|
| 24 |
+
_OFFLINE = not HF_TOKEN or HF_TOKEN.strip() in ("", "none", "offline", "0")
|
| 25 |
+
|
| 26 |
if FILE_ENCRYPTION_KEY_B64 is None:
|
| 27 |
+
if _OFFLINE:
|
| 28 |
+
# Deterministic dev-only key so local databases stay readable between runs.
|
| 29 |
+
FILE_ENCRYPTION_KEY_B64 = base64.urlsafe_b64encode(b"0" * 32).decode()
|
| 30 |
+
logger.warning("β οΈ FILE_ENCRYPTION_KEY not set and running offline - using dev-only key")
|
| 31 |
+
else:
|
| 32 |
+
from cryptography.fernet import Fernet
|
| 33 |
+
FILE_ENCRYPTION_KEY_B64 = Fernet.generate_key().decode()
|
| 34 |
+
logger.warning(
|
| 35 |
+
"β οΈ FILE_ENCRYPTION_KEY not set - generated a random key for this run. "
|
| 36 |
+
"Files uploaded now will NOT be readable after restart unless you persist it!"
|
| 37 |
+
)
|
| 38 |
|
| 39 |
try:
|
| 40 |
FILE_ENCRYPTION_KEY = base64.urlsafe_b64decode(FILE_ENCRYPTION_KEY_B64.encode())
|
|
|
|
| 44 |
assert len(FILE_ENCRYPTION_KEY) == 32, "FILE_ENCRYPTION_KEY must decode to exactly 32 bytes for AES-256"
|
| 45 |
|
| 46 |
# ------------------------------------------------------------------------
|
| 47 |
+
# Storage availability
|
| 48 |
# ------------------------------------------------------------------------
|
| 49 |
+
class StorageUnavailableError(IOError):
|
| 50 |
+
"""Raised when a storage operation is attempted without a valid HF_TOKEN."""
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 51 |
|
|
|
|
|
|
|
|
|
|
|
|
|
| 52 |
|
| 53 |
+
def _require_available():
|
| 54 |
+
if _OFFLINE:
|
| 55 |
+
raise StorageUnavailableError(
|
| 56 |
+
"File storage is unavailable: HF_TOKEN is not configured."
|
| 57 |
+
)
|
| 58 |
+
|
| 59 |
+
|
| 60 |
+
def _get_api():
|
| 61 |
+
_require_available()
|
| 62 |
+
return HfApi(token=HF_TOKEN)
|
| 63 |
+
|
| 64 |
+
|
| 65 |
+
def _get_fs():
|
| 66 |
+
_require_available()
|
| 67 |
+
return HfFileSystem(token=HF_TOKEN)
|
| 68 |
+
|
| 69 |
|
| 70 |
# ------------------------------------------------------------------------
|
| 71 |
+
# Get bucket owner (HF username)
|
| 72 |
# ------------------------------------------------------------------------
|
| 73 |
+
OWNER = None
|
| 74 |
+
BUCKET_ID = BUCKET_NAME
|
| 75 |
+
BUCKET_URI = None
|
| 76 |
+
|
| 77 |
+
if not _OFFLINE:
|
| 78 |
+
try:
|
| 79 |
+
owner_info = _get_api().whoami()
|
| 80 |
+
OWNER = owner_info["name"]
|
| 81 |
+
except Exception as e:
|
| 82 |
+
logger.warning(f"Could not determine HF username: {e}")
|
| 83 |
+
OWNER = None
|
| 84 |
+
|
| 85 |
+
if OWNER:
|
| 86 |
+
if "/" in BUCKET_NAME:
|
| 87 |
+
OWNER, BUCKET_ID = BUCKET_NAME.split("/", 1)
|
| 88 |
+
else:
|
| 89 |
+
BUCKET_ID = BUCKET_NAME
|
| 90 |
+
BUCKET_URI = f"hf://buckets/{OWNER}/{BUCKET_ID}"
|
| 91 |
+
logger.info(f"π¦ Bucket URI: {BUCKET_URI}")
|
| 92 |
+
|
| 93 |
|
| 94 |
# ------------------------------------------------------------------------
|
| 95 |
# Bucket Initialization
|
| 96 |
# ------------------------------------------------------------------------
|
| 97 |
def ensure_bucket():
|
| 98 |
"""Create the private bucket if it doesn't exist."""
|
| 99 |
+
if _OFFLINE:
|
| 100 |
+
return
|
| 101 |
global OWNER, BUCKET_URI, BUCKET_ID
|
| 102 |
+
api = _get_api()
|
| 103 |
try:
|
| 104 |
+
api.create_bucket(
|
| 105 |
bucket_id=f"{OWNER}/{BUCKET_ID}",
|
| 106 |
private=True,
|
| 107 |
exist_ok=True
|
|
|
|
| 115 |
logger.warning(f"β οΈ Cannot create bucket - permission issue: {e}")
|
| 116 |
else:
|
| 117 |
try:
|
| 118 |
+
api.create_bucket(
|
| 119 |
bucket_id=BUCKET_ID,
|
| 120 |
private=True,
|
| 121 |
exist_ok=True
|
|
|
|
| 130 |
else:
|
| 131 |
logger.error(f"β Failed to create bucket: {e2}")
|
| 132 |
|
| 133 |
+
|
| 134 |
+
if not _OFFLINE:
|
| 135 |
+
ensure_bucket()
|
| 136 |
|
| 137 |
# ------------------------------------------------------------------------
|
| 138 |
# Heartbeat sync
|
| 139 |
# ------------------------------------------------------------------------
|
| 140 |
def _sync_bucket_interval():
|
| 141 |
+
if _OFFLINE:
|
| 142 |
+
return
|
| 143 |
+
|
| 144 |
def _sync():
|
| 145 |
while True:
|
| 146 |
time.sleep(60)
|
|
|
|
| 155 |
t = threading.Thread(target=_sync, daemon=True)
|
| 156 |
t.start()
|
| 157 |
|
| 158 |
+
|
| 159 |
_sync_bucket_interval()
|
| 160 |
|
| 161 |
# ------------------------------------------------------------------------
|
|
|
|
| 167 |
ciphertext = aesgcm.encrypt(nonce, data, aad or b"")
|
| 168 |
return nonce + ciphertext
|
| 169 |
|
| 170 |
+
|
| 171 |
def decrypt_bytes(encrypted_blob: bytes, aad: Optional[bytes] = None) -> bytes:
|
| 172 |
aesgcm = AESGCM(FILE_ENCRYPTION_KEY)
|
| 173 |
nonce = encrypted_blob[:12]
|
| 174 |
ciphertext = encrypted_blob[12:]
|
| 175 |
return aesgcm.decrypt(nonce, ciphertext, aad or b"")
|
| 176 |
|
| 177 |
+
|
| 178 |
def verify_file_integrity(encrypted_blob: bytes) -> bool:
|
| 179 |
return len(encrypted_blob) >= 29
|
| 180 |
|
|
|
|
| 182 |
# Path Utilities
|
| 183 |
# ------------------------------------------------------------------------
|
| 184 |
def _bucket_path(remote_path: str) -> str:
|
| 185 |
+
_require_available()
|
| 186 |
+
if not BUCKET_URI:
|
| 187 |
+
raise StorageUnavailableError("Bucket is not configured (offline mode).")
|
| 188 |
return f"{BUCKET_URI}/{remote_path.lstrip('/')}"
|
| 189 |
|
| 190 |
+
|
| 191 |
def _validate_path(remote_path: str) -> None:
|
| 192 |
if not remote_path:
|
| 193 |
raise ValueError("Path cannot be empty")
|
|
|
|
| 203 |
# ------------------------------------------------------------------------
|
| 204 |
def store_file(remote_path: str, data: bytes, encrypt: bool = True) -> str:
|
| 205 |
_validate_path(remote_path)
|
| 206 |
+
_require_available()
|
| 207 |
+
fs = _get_fs()
|
| 208 |
try:
|
| 209 |
payload = encrypt_bytes(data, remote_path.encode('utf-8')) if encrypt else data
|
| 210 |
full_uri = _bucket_path(remote_path)
|
| 211 |
+
with fs.open(full_uri, "wb") as f:
|
| 212 |
f.write(payload)
|
| 213 |
+
if not fs.exists(full_uri):
|
| 214 |
raise IOError(f"Failed to verify file was stored: {full_uri}")
|
| 215 |
logger.debug(f"πΎ Stored: {remote_path} ({len(data)} bytes)")
|
| 216 |
return remote_path
|
|
|
|
| 218 |
logger.error(f"β Failed to store {remote_path}: {e}")
|
| 219 |
raise IOError(f"Storage failed: {e}")
|
| 220 |
|
| 221 |
+
|
| 222 |
def retrieve_file(remote_path: str, decrypt: bool = True) -> bytes:
|
| 223 |
_validate_path(remote_path)
|
| 224 |
+
_require_available()
|
| 225 |
+
fs = _get_fs()
|
| 226 |
full_uri = _bucket_path(remote_path)
|
| 227 |
try:
|
| 228 |
+
if not fs.exists(full_uri):
|
| 229 |
raise FileNotFoundError(f"File not found: {remote_path}")
|
| 230 |
+
with fs.open(full_uri, "rb") as f:
|
| 231 |
payload = f.read()
|
| 232 |
if not payload:
|
| 233 |
raise IOError(f"Empty file: {remote_path}")
|
|
|
|
| 244 |
logger.error(f"β Failed to retrieve {remote_path}: {e}")
|
| 245 |
raise IOError(f"Retrieval failed: {e}")
|
| 246 |
|
| 247 |
+
|
| 248 |
def delete_file(remote_path: str) -> bool:
|
| 249 |
_validate_path(remote_path)
|
| 250 |
+
if _OFFLINE:
|
| 251 |
+
return False
|
| 252 |
+
fs = _get_fs()
|
| 253 |
full_uri = _bucket_path(remote_path)
|
| 254 |
try:
|
| 255 |
+
if fs.exists(full_uri):
|
| 256 |
+
fs.rm(full_uri)
|
| 257 |
logger.debug(f"ποΈ Deleted: {remote_path}")
|
| 258 |
return True
|
| 259 |
logger.warning(f"β οΈ Not found for deletion: {remote_path}")
|
|
|
|
| 262 |
logger.error(f"β Failed to delete {remote_path}: {e}")
|
| 263 |
raise IOError(f"Deletion failed: {e}")
|
| 264 |
|
| 265 |
+
|
| 266 |
def file_exists(remote_path: str) -> bool:
|
| 267 |
+
if _OFFLINE:
|
| 268 |
+
return False
|
| 269 |
try:
|
| 270 |
+
_validate_path(remote_path)
|
| 271 |
+
return _get_fs().exists(_bucket_path(remote_path))
|
| 272 |
except Exception:
|
| 273 |
return False
|
| 274 |
|
| 275 |
+
|
| 276 |
def list_files(prefix: str = "", recursive: bool = True) -> List[str]:
|
| 277 |
+
if _OFFLINE:
|
| 278 |
+
return []
|
| 279 |
try:
|
| 280 |
+
fs = _get_fs()
|
| 281 |
+
search_path = _bucket_path(prefix) if prefix else BUCKET_URI
|
| 282 |
+
items = fs.ls(search_path, detail=False, recursive=recursive)
|
| 283 |
prefix_len = len(BUCKET_URI) + 1
|
| 284 |
+
return [item[prefix_len:] for item in items if not fs.isdir(item)]
|
| 285 |
except FileNotFoundError:
|
| 286 |
return []
|
| 287 |
except Exception as e:
|
| 288 |
logger.error(f"Failed to list files: {e}")
|
| 289 |
return []
|
| 290 |
|
| 291 |
+
|
| 292 |
def get_file_size(remote_path: str) -> Optional[int]:
|
| 293 |
+
if _OFFLINE:
|
| 294 |
+
return None
|
| 295 |
try:
|
| 296 |
+
_validate_path(remote_path)
|
| 297 |
+
info = _get_fs().info(_bucket_path(remote_path))
|
| 298 |
return info.get("size")
|
| 299 |
except Exception:
|
| 300 |
return None
|
| 301 |
|
| 302 |
+
|
| 303 |
def get_file_info(remote_path: str) -> Optional[Dict[str, Any]]:
|
| 304 |
+
if _OFFLINE:
|
| 305 |
+
return None
|
| 306 |
try:
|
| 307 |
+
_validate_path(remote_path)
|
| 308 |
+
fs = _get_fs()
|
| 309 |
+
info = fs.info(_bucket_path(remote_path))
|
| 310 |
return {
|
| 311 |
"name": remote_path,
|
| 312 |
"size": info.get("size"),
|
|
|
|
| 317 |
except FileNotFoundError:
|
| 318 |
return None
|
| 319 |
except Exception as e:
|
| 320 |
+
logger.error(f"Failed to get file info: {remote_path}: {e}")
|
| 321 |
return None
|
| 322 |
|
| 323 |
+
|
| 324 |
def store_file_stream(remote_path: str, data_stream: BinaryIO, encrypt: bool = True) -> str:
|
| 325 |
return store_file(remote_path, data_stream.read(), encrypt=encrypt)
|
| 326 |
|
| 327 |
+
|
| 328 |
def retrieve_file_stream(remote_path: str, decrypt: bool = True) -> BinaryIO:
|
| 329 |
import io
|
| 330 |
return io.BytesIO(retrieve_file(remote_path, decrypt=decrypt))
|
| 331 |
|
| 332 |
+
|
| 333 |
def get_storage_stats() -> Dict[str, Any]:
|
| 334 |
+
if _OFFLINE:
|
| 335 |
+
return {
|
| 336 |
+
"bucket": None, "file_count": 0, "total_size": 0,
|
| 337 |
+
"total_size_mb": 0, "error": "Storage offline (HF_TOKEN not configured)"
|
| 338 |
+
}
|
| 339 |
try:
|
| 340 |
+
fs = _get_fs()
|
| 341 |
+
items = fs.ls(BUCKET_URI, detail=True, recursive=True)
|
| 342 |
total_size = 0
|
| 343 |
file_count = 0
|
| 344 |
+
for item in items:
|
| 345 |
+
if item.get("type") == "file":
|
| 346 |
+
total_size += item.get("size") or 0
|
|
|
|
| 347 |
file_count += 1
|
| 348 |
return {
|
| 349 |
"bucket": f"{OWNER}/{BUCKET_ID}",
|
|
|
|
| 361 |
"error": str(e)
|
| 362 |
}
|
| 363 |
|
| 364 |
+
|
| 365 |
def create_backup(backup_prefix: str = "backups") -> str:
|
| 366 |
+
_require_available()
|
| 367 |
timestamp = int(time.time())
|
| 368 |
backup_path = f"{backup_prefix}/backup_{timestamp}"
|
| 369 |
try:
|
|
|
|
| 388 |
_last_db_sync = 0
|
| 389 |
DB_SYNC_INTERVAL = 30 # seconds
|
| 390 |
|
| 391 |
+
|
| 392 |
def download_database(local_path: str) -> bool:
|
| 393 |
"""
|
| 394 |
Download database from bucket to local path.
|
| 395 |
Returns True if database was found and downloaded.
|
| 396 |
"""
|
| 397 |
+
if _OFFLINE:
|
| 398 |
+
logger.info("π Storage offline - skipping database download")
|
| 399 |
+
return False
|
| 400 |
try:
|
| 401 |
if file_exists(DB_BUCKET_PATH):
|
| 402 |
logger.info("π₯ Downloading database from bucket...")
|
|
|
|
| 413 |
logger.error(f"β Failed to download database: {e}")
|
| 414 |
return False
|
| 415 |
|
| 416 |
+
|
| 417 |
def upload_database(local_path: str) -> bool:
|
| 418 |
"""
|
| 419 |
+
Upload a consistent snapshot of the local database to the bucket.
|
| 420 |
+
|
| 421 |
+
The app runs SQLite in WAL mode, so the main .db file can lag behind
|
| 422 |
+
committed transactions (they live in the -wal file). Copying the raw file
|
| 423 |
+
could persist stale or torn data - instead we take a proper snapshot with
|
| 424 |
+
the sqlite backup API, which is safe even while writers are active.
|
| 425 |
Returns True on success.
|
| 426 |
"""
|
| 427 |
global _last_db_sync
|
| 428 |
+
if _OFFLINE:
|
| 429 |
+
return False
|
| 430 |
with _db_sync_lock:
|
| 431 |
+
tmp_path = None
|
| 432 |
try:
|
| 433 |
if not os.path.exists(local_path):
|
| 434 |
logger.warning(f"β οΈ Database not found at {local_path}, skipping upload")
|
| 435 |
return False
|
| 436 |
+
|
| 437 |
+
# snapshot to a temp file via the online backup API
|
| 438 |
+
fd, tmp_path = tempfile.mkstemp(suffix=".db",
|
| 439 |
+
dir=os.path.dirname(os.path.abspath(local_path)) or ".")
|
| 440 |
+
os.close(fd)
|
| 441 |
+
os.unlink(tmp_path) # backup() needs the destination to not exist
|
| 442 |
+
src = sqlite3.connect(f"file:{local_path}?mode=ro", uri=True)
|
| 443 |
+
dst = sqlite3.connect(tmp_path)
|
| 444 |
+
try:
|
| 445 |
+
src.backup(dst)
|
| 446 |
+
finally:
|
| 447 |
+
src.close()
|
| 448 |
+
dst.close()
|
| 449 |
+
|
| 450 |
+
with open(tmp_path, 'rb') as f:
|
| 451 |
data = f.read()
|
| 452 |
if not data:
|
| 453 |
+
logger.warning("β οΈ Database snapshot is empty, skipping upload")
|
| 454 |
return False
|
| 455 |
store_file(DB_BUCKET_PATH, data)
|
| 456 |
_last_db_sync = time.time()
|
|
|
|
| 459 |
except Exception as e:
|
| 460 |
logger.error(f"β Failed to upload database: {e}")
|
| 461 |
return False
|
| 462 |
+
finally:
|
| 463 |
+
if tmp_path and os.path.exists(tmp_path):
|
| 464 |
+
try:
|
| 465 |
+
os.unlink(tmp_path)
|
| 466 |
+
except Exception:
|
| 467 |
+
pass
|
| 468 |
+
|
| 469 |
|
| 470 |
def start_db_sync(local_path: str):
|
| 471 |
"""
|
| 472 |
Start background thread that periodically uploads the database to bucket.
|
| 473 |
"""
|
| 474 |
+
if _OFFLINE:
|
| 475 |
+
logger.info("π Database auto-sync skipped (storage offline)")
|
| 476 |
+
return
|
| 477 |
+
|
| 478 |
def _sync_loop():
|
| 479 |
while True:
|
| 480 |
time.sleep(DB_SYNC_INTERVAL)
|
|
|
|
| 487 |
t.start()
|
| 488 |
logger.info(f"π Database auto-sync started (every {DB_SYNC_INTERVAL}s)")
|
| 489 |
|
| 490 |
+
|
| 491 |
def close():
|
| 492 |
try:
|
| 493 |
+
if not _OFFLINE:
|
| 494 |
+
_get_fs().close()
|
| 495 |
except Exception:
|
| 496 |
pass
|
| 497 |
|
|
|
|
| 504 |
'store_file_stream', 'retrieve_file_stream',
|
| 505 |
'get_storage_stats', 'create_backup', 'close',
|
| 506 |
'download_database', 'upload_database', 'start_db_sync',
|
| 507 |
+
'DB_BUCKET_PATH', 'StorageUnavailableError'
|
| 508 |
]
|