Watch
1
0
Fork
You've already forked Seyyed_arc
0
forked from hesabix/arc
Seyyed_arc/hesabixAPI/adapters/api/v1/business_backups.py

1301 lines
57 KiB
Python
Executable file

from __future__ import annotations
import io
import json
import os
import tempfile
import zipfile
from datetime import datetime, date
from decimal import Decimal
from typing import Any, Dict, List, Optional
from psycopg2.extras import Json
from fastapi import APIRouter, Depends, Request, UploadFile, File, Body, BackgroundTasks
from fastapi.responses import StreamingResponse
from sqlalchemy.orm import Session
from sqlalchemy import inspect, text
from pydantic import BaseModel, Field
from adapters.db.session import get_db
from app.core.auth_dependency import get_current_user, AuthContext
from app.core.permissions import require_business_access_dep, require_business_backup_restore_dep
from app.services.business_backup_access import assert_backup_file_belongs_to_business
from app.core.responses import success_response, ApiError
from app.services.file_storage_service import FileStorageService
from app.services.business_service import create_business
from adapters.api.v1.schemas import BusinessCreateRequest, BusinessType, BusinessField
import logging
from app.services.job_manager import JobManager
from adapters.api.v1.business_ftp_backup import assert_can_use_ftp_on_backup, upload_saved_backup_to_ftp
from app.services.business_ftp_service import load_decrypted_params
from app.services.business_backup_financial_policy import (
build_backup_metadata,
compute_backup_checksum,
finalize_financial_state_after_restore,
guard_new_business_import,
is_backup_excluded_table,
register_backup_import,
validate_backup_owner,
)
from app.services.business_backup_schema_compat import (
analyze_schema_diff,
backup_columns_for_table,
build_table_restore_plan,
build_table_schemas_snapshot,
iter_jsonl_rows,
load_table_column_meta,
log_restore_plan_warnings,
validate_row_for_insert,
)
logger = logging.getLogger(__name__)
def _merge_backup_ftp_metadata(db: Session, file_id: str, ftp_info: dict[str, Any]) -> None:
from adapters.db.models.file_storage import FileStorage
from sqlalchemy.orm.attributes import flag_modified
fs = db.query(FileStorage).filter(FileStorage.id == file_id).first()
if not fs:
return
dev = dict(fs.developer_data or {})
dev["ftp_uploaded"] = True
dev["ftp_remote_path"] = ftp_info.get("remote_path")
dev["ftp_remote_filename"] = ftp_info.get("remote_filename")
dev["ftp_transport"] = ftp_info.get("transport")
fs.developer_data = dev
flag_modified(fs, "developer_data")
db.add(fs)
db.commit()
router = APIRouter(prefix="/businesses/{business_id}/backups", tags=["پشتیبان‌گیری"])
class RestoreBackupRequest(BaseModel):
backup_id: str | None = Field(default=None, description="شناسه بکاپ برای بازیابی")
mode: str = Field(default="new_business", description="حالت بازیابی: replace یا new_business")
class CreateBackupBody(BaseModel):
upload_to_ftp: bool = Field(default=False, description="پس از ذخیره داخلی، ارسال کپی به سرور FTP")
def _normalize_business_type(value: Any) -> BusinessType:
"""
تبدیل business_type از بکاپ به BusinessType enum.
پشتیبانی از enum name (مثل 'SHOP') و مقدار فارسی (مثل 'مغازه').
"""
if value is None:
return BusinessType.COMPANY
value_str = str(value).strip()
# اگر مقدار فارسی است، مستقیماً استفاده می‌کنیم
persian_values = {
"شرکت": BusinessType.COMPANY,
"مغازه": BusinessType.SHOP,
"فروشگاه": BusinessType.STORE,
"اتحادیه": BusinessType.UNION,
"باشگاه": BusinessType.CLUB,
"موسسه": BusinessType.INSTITUTE,
"شخصی": BusinessType.INDIVIDUAL,
}
if value_str in persian_values:
return persian_values[value_str]
# اگر enum name است (مثل 'SHOP', 'COMPANY')
enum_name_to_type = {
"COMPANY": BusinessType.COMPANY,
"SHOP": BusinessType.SHOP,
"STORE": BusinessType.STORE,
"UNION": BusinessType.UNION,
"CLUB": BusinessType.CLUB,
"INSTITUTE": BusinessType.INSTITUTE,
"INDIVIDUAL": BusinessType.INDIVIDUAL,
}
if value_str.upper() in enum_name_to_type:
return enum_name_to_type[value_str.upper()]
# تلاش برای پیدا کردن در enum
try:
return BusinessType(value_str)
except ValueError:
# اگر پیدا نشد، پیش‌فرض شرکت
return BusinessType.COMPANY
def _normalize_business_field(value: Any) -> BusinessField:
"""
تبدیل business_field از بکاپ به BusinessField enum.
پشتیبانی از enum name (مثل 'MANUFACTURING', 'COMMERCIAL', 'TRADING') و مقدار فارسی (مثل 'تولیدی').
"""
if value is None:
return BusinessField.OTHER
value_str = str(value).strip()
# اگر مقدار فارسی است، مستقیماً استفاده می‌کنیم
persian_values = {
"تولیدی": BusinessField.MANUFACTURING,
"بازرگانی": BusinessField.COMMERCIAL, # در schemas.py این COMMERCIAL است نه TRADING
"خدماتی": BusinessField.SERVICE,
"سایر": BusinessField.OTHER,
}
if value_str in persian_values:
return persian_values[value_str]
# اگر enum name است (مثل 'MANUFACTURING', 'COMMERCIAL', 'TRADING')
# توجه: در schemas.py این COMMERCIAL است، اما در db/models ممکن است TRADING باشد
enum_name_to_field = {
"MANUFACTURING": BusinessField.MANUFACTURING,
"COMMERCIAL": BusinessField.COMMERCIAL, # نام صحیح در schemas.py
"TRADING": BusinessField.COMMERCIAL, # تبدیل TRADING به COMMERCIAL (برای سازگاری)
"SERVICE": BusinessField.SERVICE,
"OTHER": BusinessField.OTHER,
}
if value_str.upper() in enum_name_to_field:
return enum_name_to_field[value_str.upper()]
# تلاش برای پیدا کردن در enum
try:
return BusinessField(value_str)
except ValueError:
# اگر پیدا نشد، پیش‌فرض سایر
return BusinessField.OTHER
def _json_default(o: Any):
if isinstance(o, (datetime, date)):
return o.isoformat()
if isinstance(o, Decimal):
return str(o)
return str(o)
def _read_first_business_row_from_zip(zf: zipfile.ZipFile) -> Dict[str, Any] | None:
"""اولین ردیف businesses.jsonl برای اعتبارسنجی مالک بکاپ‌های قدیمی."""
try:
with zf.open("tables/businesses.jsonl", "r") as f:
for raw in f:
if not raw:
continue
return json.loads(raw.decode("utf-8"))
except KeyError:
return None
return None
def _try_set_session_replication_role_replica(conn) -> bool:
"""
تلاش برای SET session_replication_role=replica (غیرفعال موقت بررسی FK در PG).
اگر کاربر دیتابیس مجوز ندارد، خطا تراکنش را abort می‌کند؛ با SAVEPOINT + ROLLBACK TO SAVEPOINT
تراکنش را سالم نگه می‌داریم تا بازیابی ادامه یابد (با FK فعال).
"""
conn.execute(text("SAVEPOINT sp_hesabix_replica_role"))
try:
conn.execute(text("SET session_replication_role = 'replica'"))
conn.execute(text("RELEASE SAVEPOINT sp_hesabix_replica_role"))
return True
except Exception as e:
conn.execute(text("ROLLBACK TO SAVEPOINT sp_hesabix_replica_role"))
logger.warning(
"SET session_replication_role=replica denied or failed; continuing with FK checks enabled: %s",
e,
)
return False
def _reset_session_replication_role(conn, had_replica: bool) -> None:
"""بازگرداندن نقش سشن به origin فقط اگر قبلاً replica ست شده بود."""
if not had_replica:
return
try:
conn.execute(text("SET session_replication_role = 'origin'"))
except Exception as e:
logger.warning("SET session_replication_role=origin failed: %s", e)
def _sort_tables_for_insert_by_fks(engine, table_names: List[str]) -> List[str]:
"""
ترتیب INSERT طوری که جداول ارجاع‌شده توسط FK قبل از جدول وابسته درج شوند.
وقتی session_replication_role=replica مجاز نیست، این ترتیب برای عبور از FK ضروری است.
"""
inspector = inspect(engine)
names = list(dict.fromkeys(table_names))
tset = set(names)
deps: Dict[str, set[str]] = {t: set() for t in names}
for t in names:
try:
for fk in inspector.get_foreign_keys(t):
ref = fk["referred_table"]
if ref in tset and ref != t:
deps[t].add(ref)
except Exception:
continue
done: set[str] = set()
out: List[str] = []
pending = set(names)
while pending:
layer = sorted(x for x in pending if deps[x] <= done)
if not layer:
logger.warning(
"Could not fully order tables by FKs (cycle or metadata); appending rest alphabetically: %s",
sorted(pending)[:30],
)
out.extend(sorted(pending))
break
for t in layer:
out.append(t)
done.add(t)
pending.remove(t)
return out
def _build_insert_params_for_restore(
schema_inspector,
table: str,
insert_col_list: List[str],
rec: Dict[str, Any],
json_cols_cache: Dict[str, set[str]],
) -> Dict[str, Any]:
"""
مقادیر ستون‌های JSON/SQLAlchemy.JSON را برای raw SQL + psycopg2 به Json() می‌پیچد
تا خطای «can't adapt type 'dict'» در bulk insert رخ ندهد.
"""
if table not in json_cols_cache:
json_col_names: set[str] = set()
try:
for col in schema_inspector.get_columns(table):
typ = col.get("type")
tn = type(typ).__name__ if typ is not None else ""
ts = str(typ) if typ is not None else ""
if "JSON" in tn or "JSON" in ts.upper():
json_col_names.add(col["name"])
except Exception:
pass
json_cols_cache[table] = json_col_names
jc = json_cols_cache[table]
params: Dict[str, Any] = {}
for c in insert_col_list:
v = rec.get(c)
if v is None:
params[c] = None
elif c in jc and isinstance(v, (dict, list)):
params[c] = Json(v)
else:
params[c] = v
return params
def _insert_table_from_backup_zip(
conn,
zf: zipfile.ZipFile,
*,
table: str,
metadata: Dict[str, Any],
tables_info: Dict[str, Dict[str, Any]],
schema_inspector,
new_business_id: int,
mode: str,
json_cols_cache: Dict[str, set[str]],
column_meta_cache: Dict[str, Dict[str, Any]],
) -> None:
"""درج ردیف‌های یک جدول از jsonl با سازگاری مبتنی بر اسکیمای DB."""
col_list = tables_info[table]["columns"]
pk_omit = (
_pk_columns_to_omit_for_new_business(schema_inspector, table)
if mode == "new_business"
else []
)
preview_rows = iter_jsonl_rows(zf, table)
backup_cols = backup_columns_for_table(table, metadata, preview_rows)
if table not in column_meta_cache:
column_meta_cache[table] = load_table_column_meta(schema_inspector, table, conn)
table_meta = column_meta_cache[table]
restore_plan = build_table_restore_plan(
table, col_list, backup_cols, pk_omit, table_meta
)
log_restore_plan_warnings(restore_plan)
insert_col_list = restore_plan.insert_columns
if not insert_col_list:
return
placeholders = ", ".join([f":{c}" for c in insert_col_list])
columns_sql = ", ".join([f'"{c}"' for c in insert_col_list])
if mode == "new_business":
insert_sql = text(
f'INSERT INTO "{table}" ({columns_sql}) VALUES ({placeholders}) ON CONFLICT DO NOTHING'
)
else:
insert_sql = text(f'INSERT INTO "{table}" ({columns_sql}) VALUES ({placeholders})')
validation_logged = False
batch: List[Dict[str, Any]] = []
with zf.open(f"tables/{table}.jsonl", "r") as f:
for raw in f:
if not raw:
continue
rec = json.loads(raw.decode("utf-8"))
if "business_id" in rec:
rec["business_id"] = new_business_id
for _pkc in pk_omit:
if _pkc in rec:
del rec[_pkc]
rec = restore_plan.enrich_row(rec)
if not validation_logged:
row_issues = validate_row_for_insert(table, rec, restore_plan)
for issue in row_issues[:5]:
logger.warning("backup_restore validation: %s", issue)
validation_logged = True
params = _build_insert_params_for_restore(
schema_inspector,
table,
insert_col_list,
rec,
json_cols_cache,
)
batch.append(params)
if len(batch) >= 500:
conn.execute(insert_sql, batch)
batch.clear()
if batch:
conn.execute(insert_sql, batch)
def _pk_columns_to_omit_for_new_business(inspector, table: str) -> List[str]:
"""
در حالت new_business فقط ستون‌های PK خودکار (serial/identity) را از INSERT حذف می‌کنیم تا
دیتابیس مقدار جدید بدهد. PK از نوع رشته/UUID (مثل file_storage.id) باید از بکاپ حفظ شود.
برای PK مرکب فعلاً چیزی حذف نمی‌شود (نیاز به نگاشت شناسه‌ها).
"""
try:
pk_cols = inspector.get_pk_constraint(table).get("constrained_columns") or []
except Exception:
return []
if len(pk_cols) != 1:
return []
pk_name = pk_cols[0]
for c in inspector.get_columns(table):
if c["name"] != pk_name:
continue
typ = c.get("type")
tn = type(typ).__name__ if typ is not None else ""
ts = str(typ).upper()
if "VARCHAR" in ts or "CHAR" in ts or "TEXT" in ts or "CITEXT" in ts:
return []
if "UUID" in tn.upper():
return []
if c.get("autoincrement"):
return [pk_name]
# Integer بدون autoincrement: احتمال PK غیرخودکار — مقدار بکاپ را نگه دار
return []
return []
def _discover_scoped_tables(db: Session) -> Dict[str, Dict[str, Any]]:
"""
جداول مرتبط با دامنه کسب‌وکار را به‌صورت پویا کشف می‌کند.
راهکار پایه: هر جدول دارای ستون business_id + خود جدول businesses.
این پیاده‌سازی قابل توسعه است تا مسیرهای غیرمستقیم FK را نیز پوشش دهد.
"""
engine = db.get_bind()
inspector = inspect(engine)
tables_info: Dict[str, Dict[str, Any]] = {}
for table_name in inspector.get_table_names():
try:
cols = inspector.get_columns(table_name)
except Exception:
continue
col_names = {c["name"] for c in cols}
if "business_id" in col_names or table_name == "businesses":
pk_cols = inspector.get_pk_constraint(table_name).get("constrained_columns") or []
tables_info[table_name] = {
"columns": [c["name"] for c in cols],
"pk": pk_cols,
}
return tables_info
def _dump_business_data(db: Session, business_id: int) -> Dict[str, Any]:
"""
داده‌های tenant-scoped را به‌صورت پویا با شرط business_id استخراج می‌کند.
جداول مالی/اعتباری (کیف پول، اشتراک AI، …) در بکاپ قرار نمی‌گیرند.
خروجی: { metadata, tables: {table: [rows...] } }
"""
from adapters.db.models.business import Business
tables = _discover_scoped_tables(db)
data_out: Dict[str, List[Dict[str, Any]]] = {}
owner_id: int | None = None
biz = db.get(Business, int(business_id))
if biz is not None:
owner_id = getattr(biz, "owner_id", None)
for table_name, meta in tables.items():
if is_backup_excluded_table(table_name):
continue
if table_name == "businesses":
stmt = text(f"SELECT * FROM {table_name} WHERE id = :bid")
rows = [dict(r._mapping) for r in db.execute(stmt, {"bid": business_id}).all()]
else:
stmt = text(f"SELECT * FROM {table_name} WHERE business_id = :bid")
try:
rows = [dict(r._mapping) for r in db.execute(stmt, {"bid": business_id}).all()]
except Exception:
rows = []
data_out[table_name] = rows
table_schemas = build_table_schemas_snapshot(tables, data_out)
metadata = build_backup_metadata(
business_id=int(business_id),
table_names=data_out.keys(),
owner_id=owner_id,
table_schemas=table_schemas,
)
return {"metadata": metadata, "tables": data_out}
@router.get("", dependencies=[Depends(require_business_access_dep)])
async def list_backups(
request: Request,
business_id: int,
db: Session = Depends(get_db),
ctx: AuthContext = Depends(get_current_user),
):
"""
فهرست بکاپ‌های ثبت‌شده برای این کسب‌وکار.
توجه: از developer_data برای فیلتر بر اساس business_id استفاده می‌کنیم.
"""
from adapters.db.models.file_storage import FileStorage
q = (
db.query(FileStorage)
.filter(
FileStorage.module_context == "business_backup",
FileStorage.deleted_at.is_(None),
)
.order_by(FileStorage.created_at.desc())
)
items = []
# فیلتر JSON (سازگار با MySQL JSON)
for f in q.all():
dev = f.developer_data or {}
bid_val = dev.get("business_id")
same_business = False
if isinstance(bid_val, int):
same_business = (bid_val == business_id)
else:
try:
same_business = (int(str(bid_val)) == int(business_id))
except Exception:
same_business = False
if same_business:
items.append(
{
"id": str(f.id),
"filename": f.original_name,
"size": f.file_size,
"created_at": f.created_at.isoformat() if f.created_at else None,
"checksum": f.checksum,
"mime_type": f.mime_type,
"ftp_uploaded": bool(dev.get("ftp_uploaded")),
"ftp_remote_filename": dev.get("ftp_remote_filename"),
"ftp_transport": dev.get("ftp_transport"),
}
)
return success_response({"items": items}, request=request)
def _perform_backup(db: Session, ctx: AuthContext, business_id: int, job_id: str | None = None) -> Dict[str, Any]:
"""اجرای هم‌زمان بکاپ (برای استفاده در background task نیز)"""
jm = JobManager.instance()
try:
if job_id:
jm.start(job_id, "Starting backup")
jm.update(job_id, 10, "Collecting data")
snapshot = _dump_business_data(db, business_id)
if job_id:
jm.update(job_id, 50, "Packaging archive")
buf = io.BytesIO()
with zipfile.ZipFile(buf, mode="w", compression=zipfile.ZIP_DEFLATED) as zf:
zf.writestr("metadata.json", json.dumps(snapshot["metadata"], ensure_ascii=False, indent=2, default=_json_default))
for table_name, rows in snapshot["tables"].items():
out = io.StringIO()
for row in rows:
out.write(json.dumps(row, ensure_ascii=False, default=_json_default))
out.write("\n")
zf.writestr(f"tables/{table_name}.jsonl", out.getvalue().encode("utf-8"))
if job_id:
jm.update(job_id, 70, "Saving file")
buf.seek(0)
filename = f"business_{business_id}_backup_{datetime.utcnow().strftime('%Y%m%dT%H%M%SZ')}.hbx"
# اینجا فقط خروجی ساخت ZIP را برمی‌گردانیم.
result = {
"zip_bytes": buf.getvalue(),
"metadata": snapshot["metadata"],
"filename": filename,
}
if job_id:
jm.update(job_id, 90, "Finalizing")
return result
except Exception as e:
if job_id:
# استخراج error message به صورت مناسب
error_msg = None
error_code = None
error_details = None
# بررسی ApiError یا HTTPException
from fastapi import HTTPException
if isinstance(e, (ApiError, HTTPException)):
if hasattr(e, 'detail') and isinstance(e.detail, dict):
error_info = e.detail.get("error", {})
if isinstance(error_info, dict):
error_code = error_info.get("code")
error_msg = error_info.get("message")
error_details = error_info.get("details")
elif isinstance(e.detail, str):
error_msg = e.detail
# اگر error message پیدا نشد، از str(e) استفاده کن
if not error_msg:
error_msg = str(e) if e else "Unknown error"
if not error_msg or error_msg.strip() == "":
error_msg = f"Backup failed: {type(e).__name__}"
# ساخت error message مناسب برای frontend
if error_code:
final_error = f"{error_code}: {error_msg}"
if error_details:
final_error += f" | Details: {error_details}"
else:
final_error = error_msg
jm.fail(job_id, final_error, "Backup failed")
raise
def run_workflow_business_backup(
db: Session,
business_id: int,
user_id: Optional[int],
upload_to_ftp: bool = False,
) -> Dict[str, Any]:
"""
بکاپ کامل کسب‌وکار برای اجرای ورک‌فلو (هم‌زمان، بدون JobManager).
در صورت نبود user_id از مالک کسب‌وکار استفاده می‌شود.
"""
from adapters.db.models.business import Business
from adapters.db.models.user import User
from app.core.auth_dependency import AuthContext
from adapters.db.session import get_db_session
from app.services.async_isolated import run_coroutine_isolated
uid: Optional[int] = user_id
if uid is None:
b = db.get(Business, int(business_id))
if b and getattr(b, "owner_id", None):
uid = int(b.owner_id)
if not uid:
return {"success": False, "error": "NO_USER_FOR_BACKUP", "message": "user_id یا مالک کسب‌وکار یافت نشد"}
user = db.get(User, int(uid))
if not user:
return {"success": False, "error": "USER_NOT_FOUND", "message": "کاربر یافت نشد"}
ctx = AuthContext(user, api_key_id=0, db=db, business_id=int(business_id))
if upload_to_ftp:
try:
assert_can_use_ftp_on_backup(db, ctx, business_id)
except ApiError as e:
return {
"success": False,
"error": "FTP_FORBIDDEN",
"message": getattr(e, "detail", None) or str(e),
}
if not load_decrypted_params(db, business_id):
return {
"success": False,
"error": "FTP_NOT_CONFIGURED",
"message": "تنظیمات FTP برای این کسب‌وکار کامل نیست",
}
try:
data = _perform_backup(db, ctx, business_id, job_id=None)
except Exception as e:
logger.exception("run_workflow_business_backup _perform_backup failed")
return {"success": False, "error": "BACKUP_FAILED", "message": str(e)}
# آپلود async است؛ ورک‌فلو از داخل درخواست async اجرا می‌شود و در همان ترد حلقهٔ asyncio فعال است.
# anyio.run / asyncio.run در آن حالت خطا می‌دهند؛ session هم باید در همان ترد آپلود ساخته شود.
bid = int(business_id)
schema_ver = data["metadata"]["schema_version"]
fname = data["filename"]
zip_bytes = data["zip_bytes"]
async def _upload_async():
with get_db_session() as thread_db:
user = thread_db.get(User, int(uid))
if not user:
raise RuntimeError("USER_NOT_FOUND")
tctx = AuthContext(user, api_key_id=0, db=thread_db, business_id=bid)
faux = UploadFile(filename=fname, file=io.BytesIO(zip_bytes))
storage = FileStorageService(thread_db)
return await storage.upload_file(
faux,
user_id=tctx.get_user_id(),
module_context="business_backup",
developer_data={
"business_id": bid,
"schema_version": schema_ver,
"source": "workflow_action",
},
is_temporary=False,
expires_in_days=3650,
business_id=bid,
check_storage_limit=True,
)
try:
saved = run_coroutine_isolated(lambda: _upload_async())
except Exception as e:
logger.exception("run_workflow_business_backup upload failed")
return {"success": False, "error": "UPLOAD_FAILED", "message": str(e)}
out: Dict[str, Any] = {
"success": True,
"file_id": (saved or {}).get("file_id"),
"filename": data["filename"],
"metadata": data["metadata"],
}
if upload_to_ftp:
tmp_path: str | None = None
try:
with tempfile.NamedTemporaryFile(suffix=".hbx", delete=False) as tf:
tf.write(data["zip_bytes"])
tmp_path = tf.name
ftp_info = upload_saved_backup_to_ftp(
db, business_id, data["filename"], data=None, file_path=tmp_path
)
finally:
if tmp_path and os.path.isfile(tmp_path):
try:
os.unlink(tmp_path)
except OSError:
pass
if not ftp_info:
return {
"success": False,
"error": "FTP_UPLOAD_FAILED",
"message": "آپلود FTP انجام نشد",
"file_id": out.get("file_id"),
}
out["ftp"] = ftp_info
fid = (saved or {}).get("file_id")
if fid:
_merge_backup_ftp_metadata(db, str(fid), ftp_info)
return out
@router.post("", dependencies=[Depends(require_business_access_dep)])
async def create_backup(
request: Request,
business_id: int,
async_mode: bool = True,
background: BackgroundTasks = None,
db: Session = Depends(get_db),
ctx: AuthContext = Depends(get_current_user),
body: CreateBackupBody | None = Body(None),
):
"""
ایجاد بکاپ پویا. به‌صورت پیش‌فرض غیرهمزمان اجرا می‌شود و job_id برمی‌گرداند.
"""
upload_to_ftp = bool(body.upload_to_ftp) if body else False
if upload_to_ftp:
assert_can_use_ftp_on_backup(db, ctx, business_id)
if not load_decrypted_params(db, business_id):
raise ApiError("FTP_NOT_CONFIGURED", "ابتدا تنظیمات FTP را در بخش تنظیمات ذخیره کنید", http_status=400)
if async_mode:
jm = JobManager.instance()
job_id = jm.create("Backup queued")
logger.info(f"[BACKUP] Created backup job {job_id} for business_id: {business_id}")
do_ftp = upload_to_ftp
def task():
# استفاده از get_db_session برای background task
from adapters.db.session import get_db_session
with get_db_session() as db:
# اجرای بکاپ و ذخیره فایل
try:
data = _perform_backup(db, ctx, business_id, job_id=job_id)
jm.update(job_id, 92, "Uploading file")
faux_upload = UploadFile(filename=data["filename"], file=io.BytesIO(data["zip_bytes"]))
storage = FileStorageService(db)
import anyio
async def _upload():
saved = await storage.upload_file(
faux_upload,
user_id=ctx.get_user_id(),
module_context="business_backup",
developer_data={
"business_id": business_id,
"schema_version": data["metadata"]["schema_version"],
},
is_temporary=False,
expires_in_days=3650,
business_id=business_id,
check_storage_limit=True,
)
return saved
saved = anyio.run(_upload)
result_payload: Dict[str, Any] = {"file": saved, "metadata": data["metadata"]}
if do_ftp:
jm.update(job_id, 94, "Uploading to FTP")
tmp_path: str | None = None
try:
with tempfile.NamedTemporaryFile(suffix=".hbx", delete=False) as tf:
tf.write(data["zip_bytes"])
tmp_path = tf.name
ftp_info = upload_saved_backup_to_ftp(
db, business_id, data["filename"], data=None, file_path=tmp_path
)
finally:
if tmp_path and os.path.isfile(tmp_path):
try:
os.unlink(tmp_path)
except OSError:
pass
if not ftp_info:
raise ApiError("FTP_NOT_CONFIGURED", "تنظیمات FTP در دسترس نبود", http_status=400)
result_payload["ftp"] = ftp_info
fid = (saved or {}).get("file_id")
if fid:
_merge_backup_ftp_metadata(db, str(fid), ftp_info)
jm.succeed(job_id, result_payload, "Backup completed")
except Exception as e:
# استخراج error message به صورت مناسب
error_msg = None
error_code = None
error_details = None
# بررسی ApiError یا HTTPException
from fastapi import HTTPException
if isinstance(e, (ApiError, HTTPException)):
if hasattr(e, 'detail') and isinstance(e.detail, dict):
error_info = e.detail.get("error", {})
if isinstance(error_info, dict):
error_code = error_info.get("code")
error_msg = error_info.get("message")
error_details = error_info.get("details")
elif isinstance(e.detail, str):
error_msg = e.detail
# اگر error message پیدا نشد، از str(e) استفاده کن
if not error_msg:
error_msg = str(e) if e else "Unknown error"
if not error_msg or error_msg.strip() == "":
error_msg = f"Backup failed: {type(e).__name__}"
# ساخت error message مناسب برای frontend
if error_code:
final_error = f"{error_code}: {error_msg}"
if error_details:
final_error += f" | Details: {error_details}"
else:
final_error = error_msg
jm.fail(job_id, final_error, "Backup failed")
background.add_task(task)
return success_response({"job_id": job_id}, request=request, message="Backup started")
else:
# مسیر هم‌زمان قدیمی
data = _perform_backup(db, ctx, business_id)
faux_upload = UploadFile(filename=data["filename"], file=io.BytesIO(data["zip_bytes"]))
storage = FileStorageService(db)
saved = await storage.upload_file(
faux_upload,
user_id=ctx.get_user_id(),
module_context="business_backup",
developer_data={
"business_id": business_id,
"schema_version": data["metadata"]["schema_version"],
},
is_temporary=False,
expires_in_days=3650,
business_id=business_id,
check_storage_limit=True,
)
out: Dict[str, Any] = {"file": saved, "metadata": data["metadata"]}
if upload_to_ftp:
tmp_path: str | None = None
try:
with tempfile.NamedTemporaryFile(suffix=".hbx", delete=False) as tf:
tf.write(data["zip_bytes"])
tmp_path = tf.name
ftp_info = upload_saved_backup_to_ftp(
db, business_id, data["filename"], data=None, file_path=tmp_path
)
finally:
if tmp_path and os.path.isfile(tmp_path):
try:
os.unlink(tmp_path)
except OSError:
pass
if not ftp_info:
raise ApiError("FTP_NOT_CONFIGURED", "تنظیمات FTP در دسترس نبود", http_status=400)
out["ftp"] = ftp_info
fid = (saved or {}).get("file_id")
if fid:
_merge_backup_ftp_metadata(db, str(fid), ftp_info)
return success_response(out, request=request, message="Backup created")
@router.get("/{backup_id}/download", dependencies=[Depends(require_business_access_dep)])
async def download_backup(
request: Request,
business_id: int,
backup_id: str,
db: Session = Depends(get_db),
ctx: AuthContext = Depends(get_current_user),
):
"""
دانلود فایل بکاپ.
"""
from uuid import UUID
assert_backup_file_belongs_to_business(db, backup_id, business_id)
storage = FileStorageService(db)
try:
file_data = await storage.download_file(UUID(backup_id))
except Exception as e:
raise ApiError("FILE_NOT_FOUND", f"Backup not found: {e}", http_status=404)
return StreamingResponse(
io.BytesIO(file_data["content"]),
media_type=file_data["mime_type"] or "application/zip",
headers={"Content-Disposition": f'attachment; filename="{file_data["filename"]}"'},
)
@router.delete("/{backup_id}", dependencies=[Depends(require_business_access_dep)])
async def delete_backup(
request: Request,
business_id: int,
backup_id: str,
db: Session = Depends(get_db),
ctx: AuthContext = Depends(get_current_user),
):
"""
حذف نرم فایل بکاپ.
"""
from uuid import UUID
import os
from adapters.db.models.file_storage import FileStorage
# ابتدا فایل را بدون فیلتر deleted پیدا می‌کنیم تا ایدمپوتنت باشد
file = db.query(FileStorage).filter(FileStorage.id == str(UUID(backup_id))).first()
if not file:
raise ApiError("FILE_NOT_FOUND", "Backup not found", http_status=404)
# اعتبارسنجی اینکه این فایل بکاپ همین کسب‌وکار است
dev = file.developer_data or {}
bid_val = dev.get("business_id")
same_business = False
if isinstance(bid_val, int):
same_business = (bid_val == business_id)
else:
try:
same_business = (int(str(bid_val)) == int(business_id))
except Exception:
same_business = False
if file.module_context != "business_backup" or not same_business:
# برای جلوگیری از حذف فایل‌های سایر کسب‌وکارها
raise ApiError("FILE_NOT_FOUND", "Backup not found", http_status=404)
# لاگ برای بررسی وضعیت فایل روی دیسک قبل از حذف
try:
logger = logging.getLogger(__name__)
logger.info(
"backup_delete_request",
extra={
"backup_id": backup_id,
"business_id": business_id,
"path": file.file_path,
"storage_type": file.storage_type,
"exists_before": os.path.exists(file.file_path),
"deleted_at": str(file.deleted_at) if file.deleted_at else None,
},
)
except Exception:
pass
# تلاش برای حذف فیزیکی حتی اگر قبلاً soft-delete شده باشد
storage = FileStorageService(db)
try:
await storage._delete_file_from_storage(file.file_path, file.storage_type) # type: ignore[attr-defined]
except Exception:
pass
# اگر قبلاً soft-delete نشده، انجام ده؛ در غیر این صورت نتیجه موفق ایدمپوتنت برگردان
if file.deleted_at is None:
ok = await storage.delete_file(UUID(backup_id))
if not ok:
return success_response({"deleted": True}, request=request)
return success_response({"deleted": True}, request=request)
@router.post("/restore", dependencies=[Depends(require_business_backup_restore_dep)])
async def restore_backup(
request: Request,
business_id: int,
file: UploadFile | None = File(default=None),
async_mode: bool = True,
background: BackgroundTasks = None,
db: Session = Depends(get_db),
ctx: AuthContext = Depends(get_current_user),
):
"""
بازیابی داده‌ها از بکاپ.
- mode == "replace": جایگزینی داده‌های مرتبط با business_id جاری
- mode == "new_business": ایجاد کسب‌وکار جدید از بکاپ (ترتیب INSERT بر اساس FK)
می‌تواند به صورت JSON (با backup_id) یا form-data (با file) فراخوانی شود.
"""
logger = logging.getLogger(__name__)
# پردازش درخواست: می‌تواند JSON یا form-data باشد
backup_id: str | None = None
mode: str = "new_business"
# بررسی نوع محتوا
content_type = request.headers.get("content-type", "").lower()
if "application/json" in content_type:
# درخواست JSON - پارس کردن از body
try:
body_data = await request.json()
body = RestoreBackupRequest(**body_data)
backup_id = body.backup_id
mode = body.mode
except Exception as e:
raise ApiError("INVALID_INPUT", f"Invalid JSON body: {str(e)}")
elif "multipart/form-data" in content_type:
# درخواست form-data
form_data = await request.form()
backup_id_str = form_data.get("backup_id")
if backup_id_str:
backup_id = str(backup_id_str)
mode_str = form_data.get("mode")
if mode_str:
mode = str(mode_str)
else:
# تلاش برای پارس کردن به عنوان JSON
try:
body_data = await request.json()
body = RestoreBackupRequest(**body_data)
backup_id = body.backup_id
mode = body.mode
except Exception:
# اگر JSON نبود، از form-data استفاده می‌کنیم
try:
form_data = await request.form()
backup_id_str = form_data.get("backup_id")
if backup_id_str:
backup_id = str(backup_id_str)
mode_str = form_data.get("mode")
if mode_str:
mode = str(mode_str)
except Exception:
pass
if not (backup_id or file):
raise ApiError("INVALID_INPUT", "backup_id or file is required")
if mode not in ("replace", "new_business"):
raise ApiError("INVALID_INPUT", "mode must be one of: replace, new_business")
if backup_id:
assert_backup_file_belongs_to_business(db, backup_id, business_id)
# محدودسازی نرخ ساده (هر کاربر هر 60 ثانیه یکبار)
_rate_limiter_key = ("restore", ctx.get_user_id())
if not hasattr(restore_backup, "_last_calls"):
restore_backup._last_calls = {} # type: ignore[attr-defined]
last_calls = restore_backup._last_calls # type: ignore[attr-defined]
from time import time
now = time()
if _rate_limiter_key in last_calls and (now - last_calls[_rate_limiter_key]) < 60:
raise ApiError("RATE_LIMIT", "Too many restore attempts. Please wait.", http_status=429)
last_calls[_rate_limiter_key] = now
if async_mode:
jm = JobManager.instance()
job_id = jm.create("Restore queued")
logger.info(f"[RESTORE] Created restore job {job_id} for business_id: {business_id}, mode: {mode}")
# خواندن محتوای فایل قبل از شروع background task
# تا از بسته شدن فایل جلوگیری کنیم
file_bytes: bytes | None = None
if file:
try:
from app.services.system_settings_service import get_max_file_size_mb
from app.services.file_storage_service import read_upload_bounded
max_bytes = get_max_file_size_mb(db) * 1024 * 1024
file_bytes = await read_upload_bounded(file, max_bytes)
except Exception as e:
logger.error(f"Error reading file: {e}")
raise ApiError("FILE_READ_ERROR", f"Failed to read uploaded file: {str(e)}")
async def _load_zip() -> bytes:
if backup_id:
from uuid import UUID
storage = FileStorageService(db)
file_data = await storage.download_file(UUID(backup_id))
return file_data["content"]
elif file_bytes:
return file_bytes
else:
raise ApiError("NO_FILE_DATA", "No file data available")
def task():
# استفاده از get_db_session برای background task
from adapters.db.session import get_db_session
with get_db_session() as db:
try:
jm.start(job_id, "Starting restore")
jm.update(job_id, 10, "Loading backup")
import anyio
zip_bytes = anyio.run(_load_zip)
buf = io.BytesIO(zip_bytes)
with zipfile.ZipFile(buf, mode="r") as zf:
try:
metadata = json.loads(zf.read("metadata.json").decode("utf-8"))
except KeyError:
raise ApiError("INVALID_BACKUP", "metadata.json not found in backup")
snapshot_business_id = metadata.get("business_id")
backup_business_row = _read_first_business_row_from_zip(zf)
validate_backup_owner(
metadata,
ctx.get_user_id(),
db=db,
backup_business_row=backup_business_row,
)
backup_checksum = compute_backup_checksum(zip_bytes)
guard_new_business_import(
db,
user_id=ctx.get_user_id(),
backup_checksum=backup_checksum,
import_mode=mode,
)
if mode == "replace" and int(business_id) != int(snapshot_business_id):
raise ApiError("BUSINESS_MISMATCH", "Backup belongs to a different business")
# برای mode new_business، باید یک کسب‌وکار جدید ایجاد کنیم
new_business_id = business_id
if mode == "new_business":
jm.update(job_id, 20, "Creating new business")
# خواندن اطلاعات کسب‌وکار از بکاپ
try:
business_rows = []
with zf.open("tables/businesses.jsonl", "r") as f:
for raw in f:
if not raw:
continue
business_rows.append(json.loads(raw.decode("utf-8")))
if not business_rows:
raise ApiError("INVALID_BACKUP", "No business data found in backup")
# استفاده از اولین ردیف کسب‌وکار (معمولاً فقط یکی وجود دارد)
original_business = business_rows[0]
# ایجاد کسب‌وکار جدید با اطلاعات بکاپ
# تبدیل به BusinessCreateRequest
business_create_data = BusinessCreateRequest(
name=f"{original_business.get('name', 'کسب‌وکار جدید')} (بازیابی شده)",
business_type=_normalize_business_type(original_business.get('business_type')),
business_field=_normalize_business_field(original_business.get('business_field')),
address=original_business.get('address'),
phone=original_business.get('phone'),
mobile=original_business.get('mobile'),
national_id=original_business.get('national_id'),
registration_number=original_business.get('registration_number'),
economic_id=original_business.get('economic_id'),
country=original_business.get('country'),
province=original_business.get('province'),
city=original_business.get('city'),
postal_code=original_business.get('postal_code'),
default_currency_id=original_business.get('default_currency_id'),
default_credit_limit=float(original_business.get('default_credit_limit')) if original_business.get('default_credit_limit') else None,
check_credit_enabled_by_default=bool(original_business.get('check_credit_enabled_by_default', False)),
)
# ایجاد کسب‌وکار جدید (بدون commit؛ کل بازیابی یک تراکنش واحد است)
new_business = create_business(
db,
business_create_data,
ctx.get_user_id(),
defer_commit=True,
)
new_business_id = new_business['id']
jm.update(job_id, 30, f"New business created (ID: {new_business_id})")
except KeyError:
raise ApiError("INVALID_BACKUP", "businesses.jsonl not found in backup")
except Exception as e:
raise ApiError("BUSINESS_CREATION_FAILED", f"Failed to create new business: {str(e)}")
tables_info = _discover_scoped_tables(db)
target_tables = [
t
for t in tables_info.keys()
if t != "businesses" and not is_backup_excluded_table(t)
]
target_tables = _sort_tables_for_insert_by_fks(db.get_bind(), target_tables)
conn = db.connection()
replica_role_ok = _try_set_session_replication_role_replica(conn)
schema_inspector = inspect(db.get_bind())
json_cols_cache: Dict[str, set[str]] = {}
try:
# فقط برای mode replace باید داده‌های قبلی را پاک کنیم
if mode == "replace":
jm.update(job_id, 40, "Cleaning current data")
for table in target_tables:
cols = tables_info[table]["columns"]
if "business_id" in cols:
conn.execute(text(f"DELETE FROM {table} WHERE business_id = :bid"), {"bid": new_business_id})
else:
jm.update(job_id, 40, "Preparing to restore data")
# Update business row if present (فقط برای mode replace)
if mode == "replace":
jm.update(job_id, 55, "Updating business info")
if "businesses" in tables_info:
try:
rows = []
with zf.open("tables/businesses.jsonl", "r") as f:
for raw in f:
if not raw:
continue
rows.append(json.loads(raw.decode("utf-8")))
row = next((r for r in rows if int(r.get("id")) == int(snapshot_business_id)), None)
if row:
cols = [c for c in tables_info["businesses"]["columns"] if c not in ("id", "created_at", "updated_at", "owner_id")]
assignments = ", ".join([f"{c} = :{c}" for c in cols])
params = {c: row.get(c) for c in cols}
params["id"] = new_business_id
conn.execute(text(f"UPDATE businesses SET {assignments} WHERE id = :id"), params)
except KeyError:
pass
else:
jm.update(job_id, 55, "Preparing business data")
# Insert data for other tables
jm.update(job_id, 70, "Restoring data")
column_meta_cache: Dict[str, Dict[str, Any]] = {}
for _t in target_tables:
if _t not in column_meta_cache:
column_meta_cache[_t] = load_table_column_meta(
schema_inspector, _t, conn
)
analyze_schema_diff(
metadata, tables_info, target_tables, column_meta_cache
)
for table in target_tables:
try:
zf.getinfo(f"tables/{table}.jsonl")
except KeyError:
continue
try:
_insert_table_from_backup_zip(
conn,
zf,
table=table,
metadata=metadata,
tables_info=tables_info,
schema_inspector=schema_inspector,
new_business_id=new_business_id,
mode=mode,
json_cols_cache=json_cols_cache,
column_meta_cache=column_meta_cache,
)
except Exception as e:
logger.error("Error restoring table %s: %s", table, e)
raise
_reset_session_replication_role(conn, replica_role_ok)
jm.update(job_id, 85, "Resetting financial state")
finalize_financial_state_after_restore(db, new_business_id)
if mode == "new_business":
register_backup_import(
db,
user_id=ctx.get_user_id(),
backup_checksum=backup_checksum,
import_mode=mode,
source_business_id=int(snapshot_business_id) if snapshot_business_id else None,
target_business_id=new_business_id,
)
from app.services.business_service import ensure_business_default_document_policies
ensure_business_default_document_policies(
db,
new_business_id,
user_id=ctx.get_user_id(),
commit=True,
)
except Exception as e:
db.rollback()
_reset_session_replication_role(conn, replica_role_ok)
raise
result_data = {"restored": True, "mode": mode, "business_id": new_business_id}
if mode == "new_business":
result_data["new_business_id"] = new_business_id
jm.succeed(job_id, result_data, "Restore completed")
except Exception as e:
# استخراج error message به صورت مناسب
error_msg = None
error_code = None
error_details = None
# بررسی ApiError یا HTTPException
from fastapi import HTTPException
if isinstance(e, (ApiError, HTTPException)):
if hasattr(e, 'detail') and isinstance(e.detail, dict):
error_info = e.detail.get("error", {})
if isinstance(error_info, dict):
error_code = error_info.get("code")
error_msg = error_info.get("message")
error_details = error_info.get("details")
elif isinstance(e.detail, str):
error_msg = e.detail
# اگر error message پیدا نشد، از str(e) استفاده کن
if not error_msg:
error_msg = str(e) if e else "Unknown error"
if not error_msg or error_msg.strip() == "":
error_msg = f"Restore failed: {type(e).__name__}"
# ساخت error message مناسب برای frontend
if error_code:
final_error = f"{error_code}: {error_msg}"
if error_details:
final_error += f" | Details: {error_details}"
else:
final_error = error_msg
jm.fail(job_id, final_error, "Restore failed")
raise
background.add_task(task)
return success_response({"job_id": job_id}, request=request, message="Restore started")
else:
# مسیر هم‌زمان قبلی (همان منطق، بدون job)
# برای اختصار، از مسیر async استفاده کنید؛ این شاخه برای سازگاری باقی می‌ماند
raise ApiError("SYNC_DISABLED", "Use async_mode=true for restore")