forked from hesabix/arc
1301 lines
57 KiB
Python
Executable file
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")
|
|
|
|
|