fix/product-excel-import-warehouse-receipts #9

Merged
morrning merged 3 commits from Seyyed/Seyyed_arc:fix/product-excel-import-warehouse-receipts into master 2026-10-03 10:02:48 +03:30
28 changed files with 2146 additions and 317 deletions

View file

@ -19,3 +19,4 @@ Health endpoint: `GET /api/v1/health`.
## Configuration
- See `app/core/settings.py` and `.env.example`.
- چت CRM و WebSocket: اگر Redis برای API خاموش است، برای رویدادهای زندهٔ چت باید **تک worker** یا **sticky وب‌سوکت** باشد؛ جزییات: [`docs/CRM_CHAT_WEBSOCKET_DEPLOYMENT.md`](docs/CRM_CHAT_WEBSOCKET_DEPLOYMENT.md).
- معماری، رفتار تراکنشی و راهنمای پایش ایمپورت Excel کالا/خدمت: [`docs/PRODUCT_EXCEL_IMPORT_PERFORMANCE.md`](docs/PRODUCT_EXCEL_IMPORT_PERFORMANCE.md).

View file

@ -47,10 +47,13 @@ from app.services.product_service import (
get_sales_by_product_report,
get_inventory_kardex_report,
get_inventory_stock_report,
invalidate_products_cache,
)
from app.services.product_opening_balance_service import (
apply_product_opening_balance_changes,
create_product_with_opening_balance,
update_product_with_opening_balance,
validate_product_opening_balance_allowed,
)
from app.services.product_excel_import_opening_balance import (
WAREHOUSE_CODE_KEY,
@ -73,16 +76,20 @@ from app.services.product_excel_import_normalize import (
api_error_message,
build_create_payload,
build_update_payload,
find_existing_product,
find_existing_product_in_index,
format_pydantic_errors,
is_sample_import_row,
load_business_currency_index,
load_existing_product_match_index,
map_headers,
parse_bool_strict,
parse_inventory_mode,
product_match_key,
provided_keys_from_raw,
resolve_fx_currency,
select_products_import_worksheet,
)
from app.services.public_catalog_service import invalidate_public_catalog_caches
from app.services.product_excel_import_template_service import (
build_products_import_template,
category_full_path,
@ -1636,8 +1643,8 @@ async def download_products_import_template(
}
}
)
@require_business_access("business_id")
async def import_products_excel(
@require_business_access("business_id", offload_sync=True)
def import_products_excel(
request: Request,
business_id: int,
file: UploadFile = File(...),
@ -1653,11 +1660,12 @@ async def import_products_excel(
import io
import logging
import re
import time
import zipfile
from decimal import Decimal
from typing import Optional
from openpyxl import load_workbook
from sqlalchemy import and_ as _and
from sqlalchemy.exc import IntegrityError
from adapters.db.models.category import BusinessCategory
from adapters.db.models.product_attribute import ProductAttribute
from adapters.db.models.tax_type import TaxType
@ -1676,6 +1684,10 @@ async def import_products_excel(
return False
try:
import_started_at = time.perf_counter()
parse_finished_at = import_started_at
validation_finished_at = import_started_at
apply_elapsed_ms = 0.0
is_dry_run = str(dry_run).lower() in ("true","1","yes","on")
logger.info(f"[IMPORT] Starting Excel import - business_id={business_id}, dry_run={is_dry_run}, match_by={match_by}, conflict_policy={conflict_policy}")
on_missing_category = str(on_missing_category or "error").strip().lower()
@ -1716,7 +1728,7 @@ async def import_products_excel(
if not file.filename or not file.filename.lower().endswith('.xlsx'):
raise ApiError("INVALID_FILE", "فرمت فایل معتبر نیست. تنها xlsx پشتیبانی می‌شود", http_status=400)
content = await file.read()
content = file.file.read()
logger.info(f"[IMPORT] File received - filename={file.filename}, size={len(content)} bytes")
if len(content) > MAX_PRODUCT_IMPORT_FILE_BYTES:
raise ApiError("FILE_TOO_LARGE", "حجم فایل بیش از حد مجاز است (حداکثر ۱۵ مگابایت)", http_status=413)
@ -1730,6 +1742,7 @@ async def import_products_excel(
ws = select_products_import_worksheet(wb)
rows = list(ws.iter_rows(values_only=True))
parse_finished_at = time.perf_counter()
logger.info(f"[IMPORT] Excel file loaded - sheet={ws.title}, total rows={len(rows)}")
if not rows:
return success_response(data={"summary": {"total": 0}}, request=request, message="EMPTY_FILE")
@ -1747,6 +1760,7 @@ async def import_products_excel(
can_edit_opening_balance = has_business_permission_for_business(
ctx, db, business_id, "opening_balance", "edit"
)
can_write_inventory = ctx.has_business_permission("inventory", "write")
warehouse_rows = db.query(Warehouse).filter(Warehouse.business_id == business_id).all()
warehouse_index = WarehouseImportIndex(warehouse_rows)
if "name" not in mapped_keys:
@ -1766,6 +1780,31 @@ async def import_products_excel(
if conflict_policy not in ("insert", "update", "upsert"):
conflict_policy = "upsert"
match_column_index = headers.index(match_by) if match_by in headers else None
match_keys: list[str] = []
if match_column_index is not None:
for raw_row in data_rows:
raw_value = (
raw_row[match_column_index]
if match_column_index < len(raw_row)
else None
)
match_keys.append(
product_match_key(match_by, {match_by: raw_value})
)
existing_product_index = load_existing_product_match_index(
db,
business_id,
match_by,
match_keys,
)
currency_index = (
load_business_currency_index(db, business_id)
if "price_fx_currency_code" in mapped_keys
or "price_fx_currency_id" in mapped_keys
else None
)
def _normalize_number_text(v: object) -> str:
if v is None:
return ""
@ -1850,22 +1889,35 @@ async def import_products_excel(
# Preload reference data for faster resolve (only when needed).
categories_rows: list[BusinessCategory] | None = None
categories_by_parent: dict[int | None, list[BusinessCategory]] | None = None
categories_by_id: dict[int, BusinessCategory] | None = None
attrs_rows: list[ProductAttribute] | None = None
attrs_by_title: dict[str, ProductAttribute] | None = None
attr_ids: set[int] | None = None
tax_type_by_code: dict[str, int] | None = None
tax_type_by_title: dict[str, int] | None = None
tax_unit_by_code: dict[str, int] | None = None
tax_unit_by_name: dict[str, int] | None = None
def _ensure_categories_loaded() -> list[BusinessCategory]:
nonlocal categories_rows
nonlocal categories_rows, categories_by_parent, categories_by_id
if categories_rows is None:
categories_rows = db.query(BusinessCategory).filter(BusinessCategory.business_id == business_id).all()
categories_by_parent = {}
categories_by_id = {}
for category in categories_rows:
categories_by_parent.setdefault(category.parent_id, []).append(category)
categories_by_id[int(category.id)] = category
return categories_rows
def _ensure_attributes_loaded() -> list[ProductAttribute]:
nonlocal attrs_rows
nonlocal attrs_rows, attrs_by_title, attr_ids
if attrs_rows is None:
attrs_rows = db.query(ProductAttribute).filter(ProductAttribute.business_id == business_id).all()
attrs_by_title = {
a.title.strip().lower(): a for a in attrs_rows if a.title
}
attr_ids = {int(a.id) for a in attrs_rows}
return attrs_rows
def _ensure_tax_types_loaded() -> None:
@ -1901,8 +1953,8 @@ async def import_products_excel(
def _resolve_category_by_id(category_id: Optional[int]) -> tuple[Optional[int], Optional[str]]:
if category_id is None:
return None, None
exists = db.query(BusinessCategory.id).filter(_and(BusinessCategory.business_id == business_id, BusinessCategory.id == category_id)).first()
if not exists:
_ensure_categories_loaded()
if categories_by_id is None or int(category_id) not in categories_by_id:
return None, f"دسته‌بندی با شناسه {category_id} یافت نشد"
return category_id, None
@ -1928,10 +1980,7 @@ async def import_products_excel(
return None, None, created_paths
cats = _ensure_categories_loaded()
# Build index by parent_id and normalized title
by_parent: dict[int | None, list[BusinessCategory]] = {}
for c in cats:
by_parent.setdefault(c.parent_id, []).append(c)
by_parent = categories_by_parent if categories_by_parent is not None else {}
parent_id: int | None = None
current_id: int | None = None
@ -1947,7 +1996,7 @@ async def import_products_excel(
parent_id = current_id
continue
if len(candidates) > 1:
by_id = {c.id: c for c in cats}
by_id = categories_by_id or {c.id: c for c in cats}
opts = [category_full_path(c, by_id, "fa") or _get_category_titles(c)[0] for c in candidates[:5]]
return None, (
f"دسته‌بندی «{seg}» مبهم است. مسیر کامل را بنویسید، مثلاً: {opts[0]}"
@ -1965,12 +2014,19 @@ async def import_products_excel(
# Create the category under current parent
from adapters.db.repositories.category_repository import CategoryRepository
repo = CategoryRepository(db)
obj = repo.create_category(business_id=business_id, parent_id=parent_id, translations={"fa": seg, "en": seg})
obj = repo.create_category(
business_id=business_id,
parent_id=parent_id,
translations={"fa": seg, "en": seg},
auto_commit=False,
)
reference_summary["created"]["categories"] += 1
created_paths.append(seg if not created_paths else f"{created_paths[-1]} > {seg}")
# update local caches
cats.append(obj)
by_parent.setdefault(parent_id, []).append(obj)
if categories_by_id is not None:
categories_by_id[int(obj.id)] = obj
current_id = obj.id
parent_id = current_id
continue
@ -2046,7 +2102,8 @@ async def import_products_excel(
item["attribute_ids"] = []
return
# validate against business attributes
existing_ids = set([a.id for a in _ensure_attributes_loaded()])
_ensure_attributes_loaded()
existing_ids = attr_ids or set()
missing = [str(i) for i in ids if i not in existing_ids]
item["attribute_ids"] = [i for i in ids if i in existing_ids]
if missing:
@ -2056,7 +2113,7 @@ async def import_products_excel(
if not titles:
return
attrs = _ensure_attributes_loaded()
by_title = {a.title.strip().lower(): a for a in attrs if a.title}
by_title = attrs_by_title if attrs_by_title is not None else {}
resolved_ids: list[int] = []
missing_titles: list[str] = []
created: list[str] = []
@ -2070,10 +2127,19 @@ async def import_products_excel(
from adapters.db.repositories.product_attribute_repository import ProductAttributeRepository
repo = ProductAttributeRepository(db)
try:
obj = repo.create(business_id=business_id, title=t, description=None, data_type="text", options=None)
obj = repo.create(
business_id=business_id,
title=t,
description=None,
data_type="text",
options=None,
auto_commit=False,
)
reference_summary["created"]["attributes"] += 1
attrs.append(obj)
by_title[obj.title.strip().lower()] = obj
if attr_ids is not None:
attr_ids.add(int(obj.id))
resolved_ids.append(obj.id)
created.append(t)
except Exception:
@ -2104,12 +2170,17 @@ async def import_products_excel(
return out
def _find_existing_product_row(data: dict) -> Optional[Product]:
found, _err = find_existing_product(db, business_id, match_by, data)
found, _err = find_existing_product_in_index(
existing_product_index,
match_by,
data,
)
return found
errors: list[dict] = []
valid_items: list[dict] = []
seen_import_match_keys: dict[str, int] = {}
for idx, row in enumerate(data_rows, start=2):
item: dict[str, Any] = {}
row_errors: list[str] = []
@ -2265,7 +2336,13 @@ async def import_products_excel(
reference_summary["would_create"]["attributes"] += len(would_titles)
row_warnings.append("برخی ویژگی‌ها وجود ندارند و در حالت create ساخته خواهند شد")
resolve_fx_currency(item, db, business_id, row_errors)
resolve_fx_currency(
item,
db,
business_id,
row_errors,
currency_index=currency_index,
)
if item.get("price_fx_currency_id") is not None and "price_fx_currency_code" in (item.get("_provided_keys") or set()):
row_preview["resolved"]["currency"] = {"price_fx_currency_id": item.get("price_fx_currency_id")}
reference_summary["resolved"]["currency"] += 1
@ -2303,6 +2380,8 @@ async def import_products_excel(
if "item_type" not in item:
existing_type = existing_for_ob.item_type
ob_item["item_type"] = existing_type.value if hasattr(existing_type, "value") else str(existing_type)
if "inventory_mode" not in item:
ob_item["inventory_mode"] = existing_for_ob.inventory_mode
if item.get("default_warehouse_id") is None:
ob_item["default_warehouse_id"] = existing_for_ob.default_warehouse_id
ob_input, ob_errors, ob_warnings, ob_preview = prepare_opening_balance_for_import_row(
@ -2312,8 +2391,10 @@ async def import_products_excel(
db=db,
can_edit_opening_balance=can_edit_opening_balance,
is_update=existing_for_ob is not None,
can_write_inventory=can_write_inventory,
existing_product=existing_for_ob,
warehouse_index=warehouse_index,
validate_document_context=False,
)
if ob_item.get("default_warehouse_id") is not None:
item["default_warehouse_id"] = ob_item.get("default_warehouse_id")
@ -2343,7 +2424,21 @@ async def import_products_excel(
else:
item["code"] = code_str
_, match_err = find_existing_product(db, business_id, match_by, item)
import_match_key = product_match_key(match_by, item)
if import_match_key and not row_errors:
previous_row = seen_import_match_keys.get(import_match_key)
if previous_row is not None:
row_errors.append(
f"کلید تطبیق در فایل تکراری است (ردیف قبلی: {previous_row})"
)
else:
seen_import_match_keys[import_match_key] = idx
_, match_err = find_existing_product_in_index(
existing_product_index,
match_by,
item,
)
if match_err:
row_errors.append(match_err)
@ -2371,6 +2466,89 @@ async def import_products_excel(
would_insert = 0
would_update = 0
would_skip_conflict = 0
warehouse_sync: dict[str, Any] = {
"receipts_created": 0,
"issues_created": 0,
"lines_created": 0,
"unchanged_lines": 0,
"posted": not is_dry_run,
"document_ids": [],
}
items_with_opening_balance = [
item for item in valid_items if item.get("opening_balance") is not None
]
if items_with_opening_balance:
try:
validate_product_opening_balance_allowed(
db,
business_id,
None,
for_update=True,
)
except ApiError as exc:
message = api_error_message(exc)
invalid_rows = {int(item.get("_row")) for item in items_with_opening_balance}
for item in items_with_opening_balance:
errors.append({"row": item.get("_row"), "errors": [message]})
valid_items = [
item
for item in valid_items
if int(item.get("_row")) not in invalid_rows
]
existing_ob_candidates: dict[int, int] = {}
for item in valid_items:
if item.get("opening_balance") is None:
continue
existing_candidate, _ = find_existing_product_in_index(
existing_product_index, match_by, item
)
if existing_candidate is not None:
existing_ob_candidates[int(item["_row"])] = int(existing_candidate.id)
if existing_ob_candidates:
from app.services.opening_balance_warehouse_sync_service import (
find_products_with_unrelated_posted_warehouse_history,
)
conflicting_product_ids = find_products_with_unrelated_posted_warehouse_history(
db,
business_id,
sorted(set(existing_ob_candidates.values())),
)
if conflicting_product_ids:
conflict_rows = {
row_id
for row_id, product_id in existing_ob_candidates.items()
if product_id in conflicting_product_ids
}
for row_id in sorted(conflict_rows):
errors.append({
"row": row_id,
"errors": [
"برای این کالا قبلاً گردش انبار ثبت شده است؛ اختلاف را با رسید/حواله یا تعدیل انبار ثبت کنید"
],
})
valid_items = [
item for item in valid_items if int(item.get("_row")) not in conflict_rows
]
if is_dry_run:
candidate_items = [
item for item in valid_items if item.get("opening_balance") is not None
]
candidate_warehouses = {
int(item["opening_balance"]["warehouse_id"])
for item in candidate_items
if item.get("opening_balance") and item["opening_balance"].get("warehouse_id")
}
warehouse_sync.update({
"candidate_warehouses": len(candidate_warehouses),
"candidate_lines": len(candidate_items),
"posted": False,
})
validation_finished_at = time.perf_counter()
logger.info(f"[IMPORT] Processing summary - total_rows={len(data_rows)}, valid_items={len(valid_items)}, errors={len(errors)}, is_dry_run={is_dry_run}")
@ -2381,7 +2559,11 @@ async def import_products_excel(
user_id = ctx.get_user_id()
for data in valid_items:
existing, _match_err = find_existing_product(db, business_id, match_by, data)
existing, _match_err = find_existing_product_in_index(
existing_product_index,
match_by,
data,
)
if existing is None:
if conflict_policy == "update":
would_skip_conflict += 1
@ -2393,19 +2575,24 @@ async def import_products_excel(
would_update += 1
if not is_dry_run and valid_items:
apply_started_at = time.perf_counter()
logger.info(f"[IMPORT] Starting REAL import (not dry-run) for {len(valid_items)} items")
for data in valid_items:
row_idx = data.get("_row")
item_name = data.get("name", "N/A")
provided = data.get("_provided_keys") if isinstance(data.get("_provided_keys"), set) else provided_keys_from_raw(data)
existing, _match_err = find_existing_product(db, business_id, match_by, data)
opening_balance_changes: list[dict[str, Any]] = []
try:
for data in valid_items:
row_idx = data.get("_row")
provided = data.get("_provided_keys") if isinstance(data.get("_provided_keys"), set) else provided_keys_from_raw(data)
existing, _match_err = find_existing_product_in_index(
existing_product_index,
match_by,
data,
)
if existing is None:
if conflict_policy == "update":
skipped += 1
errors.append({"row": row_idx, "errors": ["کالای منطبق یافت نشد؛ سیاست فقط به‌روزرسانی است"]})
continue
try:
if existing is None:
if conflict_policy == "update":
skipped += 1
errors.append({"row": row_idx, "errors": ["کالای منطبق یافت نشد؛ سیاست فقط به‌روزرسانی است"]})
continue
create_payload = build_create_payload(data)
try:
product_request = ProductCreateRequest(**create_payload)
@ -2413,30 +2600,29 @@ async def import_products_excel(
errors.append({"row": row_idx, "errors": format_pydantic_errors(validation_error)})
skipped_apply += 1
continue
if product_request.opening_balance is not None:
create_product_with_opening_balance(
db,
business_id,
user_id,
product_request,
create_product_fn=create_product,
delete_product_fn=delete_product,
)
else:
create_product(db, business_id, product_request)
result = create_product(
db,
business_id,
product_request,
defer_cache_invalidation=True,
auto_commit=False,
serialize_result=False,
prevalidated_code_uniqueness=match_by == "code",
prevalidated_attribute_ids=True,
)
ob = product_request.opening_balance
if ob is not None:
opening_balance_changes.append({
"product_id": int(result["data"]["id"]),
"product_name": result["data"].get("name") or data.get("name"),
"warehouse_id": int(ob.warehouse_id or result["data"].get("default_warehouse_id")),
"quantity": float(ob.quantity),
"cost_price": float(ob.cost_price or 0),
})
inserted += 1
except ApiError as e:
logger.warning("product import create failed business_id=%s row=%s: %s", business_id, row_idx, api_error_message(e))
errors.append({"row": row_idx, "errors": [api_error_message(e)]})
skipped_apply += 1
except Exception as e:
logger.error("product import create failed for '%s': %s", item_name, e, exc_info=True)
errors.append({"row": row_idx, "errors": ["خطای غیرمنتظره در ایجاد کالا"]})
skipped_apply += 1
elif conflict_policy == "insert":
skipped += 1
else:
try:
elif conflict_policy == "insert":
skipped += 1
else:
update_payload = build_update_payload(data, provided)
try:
update_request = ProductUpdateRequest(**update_payload)
@ -2445,29 +2631,66 @@ async def import_products_excel(
skipped_apply += 1
continue
previous_warehouse_id = existing.default_warehouse_id
if update_request.opening_balance is not None:
update_product_with_opening_balance(
db,
business_id,
user_id,
existing.id,
update_request,
update_product_fn=update_product,
previous_warehouse_id=previous_warehouse_id,
)
else:
update_product(
db, existing.id, business_id, update_request, user_id=user_id
)
result = update_product(
db,
existing.id,
business_id,
update_request,
user_id=user_id,
defer_cache_invalidation=True,
auto_commit=False,
serialize_result=False,
prevalidated_code_uniqueness=match_by == "code",
prevalidated_attribute_ids=True,
)
ob = update_request.opening_balance
if ob is not None and result:
opening_balance_changes.append({
"product_id": int(existing.id),
"product_name": result["data"].get("name") or data.get("name"),
"warehouse_id": int(ob.warehouse_id or result["data"].get("default_warehouse_id") or previous_warehouse_id),
"previous_warehouse_id": previous_warehouse_id,
"quantity": float(ob.quantity),
"cost_price": float(ob.cost_price or 0),
})
updated += 1
except ApiError as e:
logger.warning("product import update failed business_id=%s row=%s: %s", business_id, row_idx, api_error_message(e))
errors.append({"row": row_idx, "errors": [api_error_message(e)]})
skipped_apply += 1
except Exception as e:
logger.error("product import update failed for '%s': %s", item_name, e, exc_info=True)
errors.append({"row": row_idx, "errors": ["خطای غیرمنتظره در به‌روزرسانی کالا"]})
skipped_apply += 1
if opening_balance_changes:
opening_balance_result = apply_product_opening_balance_changes(
db,
business_id,
user_id,
opening_balance_changes,
auto_commit=False,
)
from app.services.opening_balance_warehouse_sync_service import (
sync_opening_balance_warehouse_documents,
)
warehouse_sync = sync_opening_balance_warehouse_documents(
db,
business_id,
user_id,
int(opening_balance_result["id"]),
[int(change["product_id"]) for change in opening_balance_changes],
)
db.commit()
except IntegrityError as exc:
db.rollback()
raise ApiError(
"IMPORT_WRITE_CONFLICT",
"داده‌های کالا هم‌زمان تغییر کرده‌اند یا مقدار یکتای تکراری ثبت شده است؛ فایل را دوباره بررسی کنید",
http_status=409,
) from exc
except Exception:
db.rollback()
raise
if inserted or updated:
invalidate_products_cache(business_id=business_id)
invalidate_public_catalog_caches()
apply_elapsed_ms = (time.perf_counter() - apply_started_at) * 1000
else:
if is_dry_run:
logger.info("[IMPORT] DRY-RUN mode - skipping actual database operations")
@ -2487,8 +2710,14 @@ async def import_products_excel(
"would_update": would_update,
"would_skip_conflict": would_skip_conflict,
}
performance = {
"parse_ms": round((parse_finished_at - import_started_at) * 1000, 2),
"validation_ms": round((validation_finished_at - parse_finished_at) * 1000, 2),
"apply_ms": round(apply_elapsed_ms, 2),
"total_ms": round((time.perf_counter() - import_started_at) * 1000, 2),
}
logger.info(f"[IMPORT] Final summary: {summary}")
logger.info("[IMPORT] Final summary=%s performance=%s", summary, performance)
logger.info(f"[IMPORT] Import completed - inserted={inserted}, updated={updated}, skipped={skipped}")
return success_response(
@ -2496,7 +2725,9 @@ async def import_products_excel(
"summary": summary,
"errors": errors,
"reference_summary": reference_summary,
"warehouse_sync": warehouse_sync,
"preview": preview_rows if is_dry_run else None,
"performance": performance,
},
request=request,
message="PRODUCTS_IMPORT_RESULT",
@ -6255,4 +6486,3 @@ async def export_sales_by_product_report_pdf(
"include_zero_sales": bool(body.get("include_zero_sales", False)),
},
)

View file

@ -22,7 +22,7 @@ class WarehouseDocument(Base):
doc_type = Column(String(32), nullable=False) # receipt|issue|transfer|production_in|production_out|adjustment
warehouse_id_from = Column(Integer, nullable=True, index=True)
warehouse_id_to = Column(Integer, nullable=True, index=True)
source_type = Column(String(32), nullable=True) # invoice|manual|api
source_type = Column(String(32), nullable=True) # invoice|opening_balance|manual|api
source_document_id = Column(Integer, nullable=True, index=True)
extra_info = Column(JSON, nullable=True)
created_by_user_id = Column(Integer, nullable=True)
@ -34,4 +34,3 @@ class WarehouseDocument(Base):
def touch(self):
self.updated_at = datetime.utcnow()

View file

@ -70,7 +70,15 @@ class CategoryRepository(BaseRepository[BusinessCategory]):
return roots
def create_category(self, *, business_id: int, parent_id: int | None, translations: dict[str, str], description: str | None = None) -> BusinessCategory:
def create_category(
self,
*,
business_id: int,
parent_id: int | None,
translations: dict[str, str],
description: str | None = None,
auto_commit: bool = True,
) -> BusinessCategory:
obj = BusinessCategory(
business_id=business_id,
parent_id=parent_id,
@ -78,8 +86,11 @@ class CategoryRepository(BaseRepository[BusinessCategory]):
description=description,
)
self.db.add(obj)
self.db.commit()
self.db.refresh(obj)
if auto_commit:
self.db.commit()
self.db.refresh(obj)
else:
self.db.flush()
return obj
def update_category(self, *, category_id: int, translations: dict[str, str] | None = None, description: str | None = None, sort_order: int | None = None, parent_id: int | None = None) -> BusinessCategory | None:
@ -188,4 +199,3 @@ class CategoryRepository(BaseRepository[BusinessCategory]):
"path": build_path(r),
})
return result

View file

@ -486,7 +486,12 @@ class DocumentRepository:
return parse_user_date(date_value, calendar_type="jalali")
def create_document(self, document_data: Dict[str, Any]) -> Document:
def create_document(
self,
document_data: Dict[str, Any],
*,
auto_commit: bool = True,
) -> Document:
"""
ایجاد سند جدید
@ -512,15 +517,20 @@ class DocumentRepository:
)
self.db.add(line)
self.db.commit()
self.db.refresh(document)
if auto_commit:
self.db.commit()
self.db.refresh(document)
else:
self.db.flush()
return document
def update_document(
self,
document_id: int,
document_data: Dict[str, Any]
document_data: Dict[str, Any],
*,
auto_commit: bool = True,
) -> Optional[Document]:
"""
ویرایش سند موجود
@ -560,8 +570,11 @@ class DocumentRepository:
)
self.db.add(line)
self.db.commit()
self.db.refresh(document)
if auto_commit:
self.db.commit()
self.db.refresh(document)
else:
self.db.flush()
return document
@ -658,4 +671,3 @@ class DocumentRepository:
return False, f"سطر {i} نمی‌تواند همزمان بدهکار و بستانکار داشته باشد"
return True, ""

View file

@ -83,8 +83,16 @@ class ProductAttributeRepository(BaseRepository[ProductAttribute]):
},
}
def create(self, *, business_id: int, title: str, description: str | None,
data_type: str = 'text', options: dict | None = None) -> ProductAttribute:
def create(
self,
*,
business_id: int,
title: str,
description: str | None,
data_type: str = 'text',
options: dict | None = None,
auto_commit: bool = True,
) -> ProductAttribute:
obj = ProductAttribute(
business_id=business_id,
title=title,
@ -93,8 +101,11 @@ class ProductAttributeRepository(BaseRepository[ProductAttribute]):
options=options
)
self.db.add(obj)
self.db.commit()
self.db.refresh(obj)
if auto_commit:
self.db.commit()
self.db.refresh(obj)
else:
self.db.flush()
return obj
def update(self, *, attribute_id: int, title: str | None, description: str | None,
@ -122,4 +133,3 @@ class ProductAttributeRepository(BaseRepository[ProductAttribute]):
self.db.commit()
return True

View file

@ -91,7 +91,11 @@ def require_superadmin():
return decorator
def require_business_access(business_id_param: str = "business_id"):
def require_business_access(
business_id_param: str = "business_id",
*,
offload_sync: bool = False,
):
"""Decorator برای بررسی دسترسی به کسب و کار خاص.
امضای اصلی endpoint حفظ می‌شود و Request از آرگومان‌ها استخراج می‌گردد.
"""
@ -137,7 +141,12 @@ def require_business_access(business_id_param: str = "business_id"):
raise ApiError("FORBIDDEN", f"No access to business {business_id}", http_status=403)
# فراخوانی تابع اصلی و await در صورت نیاز (خارج از context manager)
result = func(*args, **kwargs)
if offload_sync and not inspect.iscoroutinefunction(func):
from starlette.concurrency import run_in_threadpool
result = await run_in_threadpool(func, *args, **kwargs)
else:
result = func(*args, **kwargs)
if inspect.isawaitable(result):
result = await result
return result

View file

@ -1242,8 +1242,8 @@ def _compute_available_stock(
WarehouseDocumentLine.warehouse_id == warehouse_id
)
# حواله‌های ناشی از فاکتور قبلاً در invoice_item_lines (بخش _iter_product_movements) لحاظ شده‌اند؛
# شمردن دوبارهٔ خطوط این حواله‌ها موجودی را دو برابر نشان می‌دهد.
# حواله‌های آینه‌ای فاکتور و تراز افتتاحیه قبلاً در خطوط سند مالی لحاظ
# شده‌اند؛ شمردن دوبارهٔ خطوط این حواله‌ها موجودی را دو برابر می‌کند.
_wh_src_doc = aliased(Document)
wh_movements_query = wh_movements_query.outerjoin(
_wh_src_doc,
@ -1253,11 +1253,11 @@ def _compute_available_stock(
),
).filter(
~or_(
func.lower(func.coalesce(WarehouseDocument.source_type, "")) == "invoice",
func.lower(func.coalesce(WarehouseDocument.source_type, "")).in_(("invoice", "opening_balance")),
and_(
WarehouseDocument.source_document_id.isnot(None),
_wh_src_doc.id.isnot(None),
_wh_src_doc.document_type.in_(SUPPORTED_INVOICE_TYPES),
_wh_src_doc.document_type.in_((*SUPPORTED_INVOICE_TYPES, "opening_balance")),
),
)
)

View file

@ -139,9 +139,15 @@ def _validate_ob_document_date_in_fiscal_year(document_date: date, fy_start: dat
)
def _find_existing_ob_document(db: Session, business_id: int, fiscal_year_id: int) -> Optional[Document]:
def _find_existing_ob_document(
db: Session,
business_id: int,
fiscal_year_id: int,
*,
for_update: bool = False,
) -> Optional[Document]:
from sqlalchemy import and_
return (
query = (
db.query(Document)
.filter(
and_(
@ -151,8 +157,10 @@ def _find_existing_ob_document(db: Session, business_id: int, fiscal_year_id: in
)
)
.order_by(Document.id.desc())
.first()
)
if for_update:
query = query.with_for_update()
return query.first()
def get_opening_balance(
@ -199,6 +207,9 @@ def upsert_opening_balance(
business_id: int,
user_id: int,
data: Dict[str, Any],
*,
auto_commit: bool = True,
return_details: bool = True,
) -> Dict[str, Any]:
repo = DocumentRepository(db)
fy_id, fy_start_date, fy_end_date = _ensure_fiscal_year(db, business_id, data.get("fiscal_year_id"))
@ -361,13 +372,17 @@ def upsert_opening_balance(
}
if existing:
updated = repo.update_document(existing.id, document_payload)
updated = repo.update_document(existing.id, document_payload, auto_commit=auto_commit)
if not updated:
raise ApiError("UPDATE_FAILED", "ویرایش سند تراز افتتاحیه ناموفق بود", http_status=500)
return repo.get_document_details(updated.id) or {}
if return_details:
return repo.get_document_details(updated.id) or {}
return {"id": int(updated.id)}
else:
created = repo.create_document(document_payload)
return repo.get_document_details(created.id) or {}
created = repo.create_document(document_payload, auto_commit=auto_commit)
if return_details:
return repo.get_document_details(created.id) or {}
return {"id": int(created.id)}
def preview_opening_balance(
@ -616,5 +631,3 @@ def unpost_opening_balance(
pass
return repo.get_document_details(updated.id) or {}

View file

@ -0,0 +1,293 @@
"""همگام‌سازی موجودی فیزیکی انبار با خطوط کالای سند افتتاحیه."""
from __future__ import annotations
from collections import defaultdict
from decimal import Decimal
from typing import Any, Dict, Iterable, List, Mapping, Tuple
from uuid import uuid4
from sqlalchemy import func
from sqlalchemy.orm import Session
from adapters.db.models.document import Document
from adapters.db.models.document_line import DocumentLine
from adapters.db.models.product import Product
from adapters.db.models.warehouse_document import WarehouseDocument
from adapters.db.models.warehouse_document_line import WarehouseDocumentLine
from app.core.responses import ApiError
InventoryKey = Tuple[int, int]
def _decimal(value: Any) -> Decimal:
try:
return Decimal(str(value or 0))
except Exception:
return Decimal(0)
def build_opening_balance_warehouse_deltas(
desired: Mapping[InventoryKey, Decimal],
linked_physical: Mapping[InventoryKey, Decimal],
) -> Dict[InventoryKey, Decimal]:
"""اختلاف موردنیاز را بدون وابستگی به دیتابیس محاسبه می‌کند."""
result: Dict[InventoryKey, Decimal] = {}
for key in set(desired) | set(linked_physical):
delta = _decimal(desired.get(key)) - _decimal(linked_physical.get(key))
if delta:
result[key] = delta
return result
def _desired_opening_inventory(
db: Session,
document_id: int,
product_ids: List[int],
) -> tuple[Dict[InventoryKey, Decimal], Dict[InventoryKey, Decimal]]:
desired: Dict[InventoryKey, Decimal] = defaultdict(Decimal)
costs: Dict[InventoryKey, Decimal] = {}
rows = db.query(DocumentLine).filter(
DocumentLine.document_id == int(document_id),
DocumentLine.product_id.in_(product_ids),
).all()
for row in rows:
info = dict(row.extra_info or {})
if str(info.get("movement") or "in").strip().lower() != "in":
continue
warehouse_id = info.get("warehouse_id")
quantity = _decimal(row.quantity)
if warehouse_id is None or quantity <= 0:
continue
key = (int(row.product_id), int(warehouse_id))
desired[key] += quantity
costs[key] = _decimal(info.get("cost_price"))
return dict(desired), costs
def _linked_opening_inventory(
db: Session,
business_id: int,
document_id: int,
product_ids: List[int],
) -> Dict[InventoryKey, Decimal]:
linked: Dict[InventoryKey, Decimal] = defaultdict(Decimal)
rows = (
db.query(WarehouseDocumentLine)
.join(
WarehouseDocument,
WarehouseDocument.id == WarehouseDocumentLine.warehouse_document_id,
)
.filter(
WarehouseDocument.business_id == int(business_id),
WarehouseDocument.status == "posted",
func.lower(func.coalesce(WarehouseDocument.source_type, "")) == "opening_balance",
WarehouseDocument.source_document_id == int(document_id),
WarehouseDocumentLine.product_id.in_(product_ids),
)
.all()
)
for row in rows:
if row.warehouse_id is None:
continue
key = (int(row.product_id), int(row.warehouse_id))
quantity = _decimal(row.quantity)
if str(row.movement or "").lower() == "in":
linked[key] += quantity
elif str(row.movement or "").lower() == "out":
linked[key] -= quantity
return dict(linked)
def _assert_products_can_sync(
db: Session,
business_id: int,
product_ids: List[int],
) -> None:
products = db.query(
Product.id,
Product.name,
Product.track_inventory,
Product.inventory_mode,
).filter(
Product.business_id == int(business_id),
Product.id.in_(product_ids),
).all()
by_id = {int(row.id): row for row in products}
for product_id in product_ids:
product = by_id.get(int(product_id))
if product is None:
raise ApiError("PRODUCT_NOT_FOUND", f"کالا با شناسه {product_id} یافت نشد", http_status=404)
if not bool(product.track_inventory):
raise ApiError(
"OPENING_BALANCE_REQUIRES_TRACK_INVENTORY",
f"کنترل موجودی کالای «{product.name}» فعال نیست",
http_status=400,
)
if str(product.inventory_mode or "bulk").strip().lower() == "unique":
raise ApiError(
"UNIQUE_OPENING_BALANCE_REQUIRES_INSTANCES",
f"ثبت خودکار موجودی اولیه کالای یونیک «{product.name}» بدون اطلاعات سریال/نمونه مجاز نیست",
http_status=400,
)
def find_products_with_unrelated_posted_warehouse_history(
db: Session,
business_id: int,
product_ids: List[int],
) -> set[int]:
if not product_ids:
return set()
rows = (
db.query(WarehouseDocumentLine.product_id)
.join(
WarehouseDocument,
WarehouseDocument.id == WarehouseDocumentLine.warehouse_document_id,
)
.filter(
WarehouseDocument.business_id == int(business_id),
WarehouseDocument.status == "posted",
func.lower(func.coalesce(WarehouseDocument.source_type, "")) != "opening_balance",
WarehouseDocumentLine.product_id.in_(product_ids),
)
.distinct()
.all()
)
return {int(row.product_id) for row in rows}
def _assert_no_unrelated_posted_history(
db: Session,
business_id: int,
product_ids: List[int],
) -> None:
conflicts = find_products_with_unrelated_posted_warehouse_history(
db, business_id, product_ids
)
if conflicts:
raise ApiError(
"OPENING_BALANCE_WAREHOUSE_HISTORY_EXISTS",
"برای یکی از کالاهای ایمپورت‌شده قبلاً گردش انبار ثبت شده است؛ "
"تعداد افتتاحیه را تغییر ندهید و اختلاف را با رسید/حواله یا تعدیل انبار ثبت کنید",
http_status=409,
details={"product_ids": sorted(conflicts)},
)
def _group_delta_lines(
deltas: Mapping[InventoryKey, Decimal],
costs: Mapping[InventoryKey, Decimal],
) -> tuple[Dict[int, List[Dict[str, Any]]], Dict[int, List[Dict[str, Any]]]]:
receipts: Dict[int, List[Dict[str, Any]]] = defaultdict(list)
issues: Dict[int, List[Dict[str, Any]]] = defaultdict(list)
for (product_id, warehouse_id), delta in sorted(deltas.items()):
if not delta:
continue
line = {
"product_id": int(product_id),
"quantity": abs(delta),
"extra_info": {
"opening_balance_sync": True,
"cost_price": float(costs.get((product_id, warehouse_id), Decimal(0))),
},
}
(receipts if delta > 0 else issues)[int(warehouse_id)].append(line)
return dict(receipts), dict(issues)
def sync_opening_balance_warehouse_documents(
db: Session,
business_id: int,
user_id: int,
opening_balance_document_id: int,
affected_product_ids: Iterable[int],
) -> Dict[str, Any]:
"""رسید/حواله اختلافی ایجاد و قطعی می‌کند؛ commit بر عهده caller است."""
product_ids = sorted({int(product_id) for product_id in affected_product_ids})
empty_summary: Dict[str, Any] = {
"receipts_created": 0,
"issues_created": 0,
"lines_created": 0,
"unchanged_lines": 0,
"posted": True,
"document_ids": [],
}
if not product_ids:
return empty_summary
opening_document = db.query(Document).filter(
Document.id == int(opening_balance_document_id),
Document.business_id == int(business_id),
Document.document_type == "opening_balance",
).first()
if opening_document is None:
raise ApiError("OPENING_BALANCE_NOT_FOUND", "سند تراز افتتاحیه یافت نشد", http_status=404)
_assert_products_can_sync(db, business_id, product_ids)
_assert_no_unrelated_posted_history(db, business_id, product_ids)
desired, costs = _desired_opening_inventory(db, int(opening_document.id), product_ids)
linked = _linked_opening_inventory(db, business_id, int(opening_document.id), product_ids)
deltas = build_opening_balance_warehouse_deltas(desired, linked)
receipts, issues = _group_delta_lines(deltas, costs)
unchanged = sum(1 for key in set(desired) | set(linked) if key not in deltas)
summary = {**empty_summary, "unchanged_lines": unchanged}
if not deltas:
return summary
from app.services.warehouse_service import (
create_linked_warehouse_document_bulk,
post_warehouse_document,
)
batch_id = str(uuid4())
created_ids: List[int] = []
common_extra = {
"origin": "product_excel_import",
"opening_balance_sync": True,
"import_batch_id": batch_id,
"opening_balance_document_id": int(opening_document.id),
}
# خروج اختلاف ابتدا ثبت می‌شود تا جابه‌جایی انبار با موجودی فیزیکی قبلی کنترل شود.
for warehouse_id, lines in sorted(issues.items()):
wh = create_linked_warehouse_document_bulk(
db,
business_id,
user_id,
doc_type="issue",
document_date=opening_document.document_date,
warehouse_id=warehouse_id,
lines=lines,
source_type="opening_balance",
source_document_id=int(opening_document.id),
extra_info={**common_extra, "description": "اصلاح خودکار موجودی افتتاحیه از ایمپورت اکسل"},
)
post_warehouse_document(db, int(wh.id))
created_ids.append(int(wh.id))
for warehouse_id, lines in sorted(receipts.items()):
wh = create_linked_warehouse_document_bulk(
db,
business_id,
user_id,
doc_type="receipt",
document_date=opening_document.document_date,
warehouse_id=warehouse_id,
lines=lines,
source_type="opening_balance",
source_document_id=int(opening_document.id),
extra_info={**common_extra, "description": "رسید خودکار موجودی افتتاحیه از ایمپورت اکسل"},
)
post_warehouse_document(db, int(wh.id))
created_ids.append(int(wh.id))
summary.update({
"receipts_created": len(receipts),
"issues_created": len(issues),
"lines_created": len(deltas),
"document_ids": created_ids,
"import_batch_id": batch_id,
})
return summary

View file

@ -3,7 +3,7 @@ from __future__ import annotations
import re
from decimal import Decimal
from typing import Any, Dict, List, Optional, Sequence, Set, Tuple
from typing import Any, Dict, Iterable, List, Optional, Sequence, Set, Tuple
from openpyxl.workbook.workbook import Workbook
from openpyxl.worksheet.worksheet import Worksheet
@ -223,21 +223,74 @@ def find_existing_product(
return None, None
def product_match_key(match_by: str, data: Dict[str, Any]) -> str:
if match_by == "name":
return name_match_key(data.get("name"))
return collapse_spaces(data.get("code"))
def load_existing_product_match_index(
session: Session,
business_id: int,
match_by: str,
keys: Iterable[str],
*,
chunk_size: int = 500,
) -> Dict[str, List[Product]]:
"""Load only products referenced by an import, in bounded IN queries."""
normalized_keys = list(dict.fromkeys(k for k in keys if k))
result: Dict[str, List[Product]] = {}
for start in range(0, len(normalized_keys), max(1, chunk_size)):
chunk = normalized_keys[start:start + max(1, chunk_size)]
query = session.query(Product).filter(Product.business_id == business_id)
if match_by == "name":
query = query.filter(func.lower(func.trim(Product.name)).in_(chunk))
else:
query = query.filter(Product.code.in_(chunk))
for product in query.all():
key = (
name_match_key(product.name)
if match_by == "name"
else collapse_spaces(product.code)
)
result.setdefault(key, []).append(product)
return result
def find_existing_product_in_index(
index: Dict[str, List[Product]],
match_by: str,
data: Dict[str, Any],
) -> Tuple[Optional[Product], Optional[str]]:
key = product_match_key(match_by, data)
if not key:
return None, None
rows = index.get(key) or []
if match_by == "name" and len(rows) > 1:
return None, "چند کالا با این نام وجود دارد؛ برای به‌روزرسانی از تطبیق بر اساس کد استفاده کنید"
return (rows[0] if rows else None), None
def resolve_fx_currency(
item: Dict[str, Any],
db: Session,
business_id: int,
row_errors: List[str],
currency_index: Optional[Dict[str, Any]] = None,
) -> None:
currency_id = item.get("price_fx_currency_id")
if isinstance(currency_id, int):
exists = (
db.query(BusinessCurrency.id)
.filter(
BusinessCurrency.business_id == business_id,
BusinessCurrency.currency_id == currency_id,
currency_id in currency_index.get("ids", set())
if currency_index is not None
else (
db.query(BusinessCurrency.id)
.filter(
BusinessCurrency.business_id == business_id,
BusinessCurrency.currency_id == currency_id,
)
.first()
)
.first()
)
if not exists:
row_errors.append(f"ارز با شناسه {currency_id} برای این کسب‌وکار تعریف نشده است")
@ -246,13 +299,17 @@ def resolve_fx_currency(
if not code:
return
row = (
db.query(Currency)
.join(BusinessCurrency, BusinessCurrency.currency_id == Currency.id)
.filter(
BusinessCurrency.business_id == business_id,
func.lower(Currency.code) == code.lower(),
currency_index.get("by_code", {}).get(code.lower())
if currency_index is not None
else (
db.query(Currency)
.join(BusinessCurrency, BusinessCurrency.currency_id == Currency.id)
.filter(
BusinessCurrency.business_id == business_id,
func.lower(Currency.code) == code.lower(),
)
.first()
)
.first()
)
if not row:
row_errors.append(f"ارز با کد «{code}» برای این کسب‌وکار یافت نشد")
@ -260,6 +317,26 @@ def resolve_fx_currency(
item["price_fx_currency_id"] = row.id
def load_business_currency_index(
db: Session,
business_id: int,
) -> Dict[str, Any]:
rows = (
db.query(Currency)
.join(BusinessCurrency, BusinessCurrency.currency_id == Currency.id)
.filter(BusinessCurrency.business_id == business_id)
.all()
)
return {
"ids": {int(row.id) for row in rows},
"by_code": {
str(row.code).strip().lower(): row
for row in rows
if row.code
},
}
def map_headers(raw_headers: Sequence[object]) -> List[str]:
mapped: List[str] = []
for h in raw_headers:

View file

@ -2,7 +2,7 @@
from __future__ import annotations
from decimal import Decimal
from typing import Any, Dict, List, Optional, Set, Tuple
from typing import TYPE_CHECKING, Any, Dict, List, Optional, Set, Tuple
from sqlalchemy.orm import Session
@ -14,6 +14,9 @@ from app.services.product_opening_balance_service import (
validate_product_opening_balance_allowed,
)
if TYPE_CHECKING:
from adapters.api.v1.schema_models.product import ProductOpeningBalanceInput
OPENING_BALANCE_QUANTITY_KEY = "opening_balance_quantity"
OPENING_BALANCE_COST_KEY = "opening_balance_cost_price"
WAREHOUSE_CODE_KEY = "warehouse_code"
@ -159,8 +162,10 @@ def prepare_opening_balance_for_import_row(
db: Session,
can_edit_opening_balance: bool,
is_update: bool,
can_write_inventory: bool = True,
existing_product: Optional[Product] = None,
warehouse_index: WarehouseImportIndex,
validate_document_context: bool = True,
) -> Tuple[Optional[ProductOpeningBalanceInput], List[str], List[str], Dict[str, Any]]:
"""
ساخت opening_balance برای یک ردیف ایمپورت.
@ -201,10 +206,18 @@ def prepare_opening_balance_for_import_row(
errors.append("برای ثبت تعداد اولیه باید کنترل موجودی فعال باشد")
return None, errors, warnings, preview
if str(item.get("inventory_mode") or "bulk").strip().lower() == "unique":
errors.append("ثبت تعداد اولیه کالای یونیک از Excel بدون اطلاعات سریال/نمونه مجاز نیست")
return None, errors, warnings, preview
if not can_edit_opening_balance:
errors.append("برای ثبت تعداد اولیه به دسترسی ویرایش تراز افتتاحیه نیاز است")
return None, errors, warnings, preview
if not can_write_inventory:
errors.append("برای ثبت تعداد اولیه و رسید خودکار انبار به دسترسی ویرایش موجودی نیاز است")
return None, errors, warnings, preview
qty_dec = qty_raw if isinstance(qty_raw, Decimal) else Decimal(str(qty_raw or 0))
if qty_dec <= 0:
errors.append("تعداد اولیه باید بزرگتر از صفر باشد")
@ -235,43 +248,45 @@ def prepare_opening_balance_for_import_row(
item["default_warehouse_id"] = int(warehouse_id)
product_id = int(existing_product.id) if existing_product is not None else None
try:
if is_update and product_id is not None:
validate_product_opening_balance_allowed(
db,
business_id,
None,
product_id=product_id,
warehouse_id=int(warehouse_id),
for_update=True,
)
else:
validate_product_opening_balance_allowed(
db,
business_id,
None,
product_id=None,
warehouse_id=int(warehouse_id),
for_update=False,
)
except ApiError as exc:
err = exc.detail.get("error") if isinstance(exc.detail, dict) else {}
errors.append(str((err or {}).get("message") or exc))
return None, errors, warnings, preview
eligibility: Dict[str, Any] = {}
if validate_document_context:
try:
if is_update and product_id is not None:
validate_product_opening_balance_allowed(
db,
business_id,
None,
product_id=product_id,
warehouse_id=int(warehouse_id),
for_update=True,
)
else:
validate_product_opening_balance_allowed(
db,
business_id,
None,
product_id=None,
warehouse_id=int(warehouse_id),
for_update=False,
)
except ApiError as exc:
err = exc.detail.get("error") if isinstance(exc.detail, dict) else {}
errors.append(str((err or {}).get("message") or exc))
return None, errors, warnings, preview
eligibility = get_product_opening_balance_eligibility(
db,
business_id,
None,
can_edit_opening_balance=can_edit_opening_balance,
product_id=product_id,
warehouse_id=int(warehouse_id),
)
if not eligibility.get("editable"):
errors.append(
eligibility.get("message") or "تعداد اولیه در این شرایط قابل ثبت نیست"
eligibility = get_product_opening_balance_eligibility(
db,
business_id,
None,
can_edit_opening_balance=can_edit_opening_balance,
product_id=product_id,
warehouse_id=int(warehouse_id),
)
return None, errors, warnings, preview
if not eligibility.get("editable"):
errors.append(
eligibility.get("message") or "تعداد اولیه در این شرایط قابل ثبت نیست"
)
return None, errors, warnings, preview
fiscal_year_id = eligibility.get("fiscal_year_id")
preview.update({

View file

@ -512,6 +512,151 @@ def append_product_line_to_opening_balance(
)
def _merge_product_opening_balance_changes(
inventory_lines: List[Dict[str, Any]],
changes: List[Dict[str, Any]],
) -> List[Dict[str, Any]]:
"""Merge imported balances in O(existing lines + changes)."""
by_key: Dict[tuple[int, int], Dict[str, Any]] = {}
for line in inventory_lines:
pid = line.get("product_id")
wid = _warehouse_id_from_line(line)
if pid is None or wid is None:
continue
by_key[(int(pid), int(wid))] = line
for change in changes:
product_id = int(change["product_id"])
warehouse_id = int(change["warehouse_id"])
quantity = Decimal(str(change.get("quantity") or 0))
cost_price = Decimal(str(change.get("cost_price") or 0))
if quantity <= 0:
raise ApiError(
"INVALID_OPENING_BALANCE_QUANTITY",
"تعداد اولیه باید بزرگتر از صفر باشد",
http_status=400,
)
if cost_price < 0:
raise ApiError(
"INVALID_OPENING_BALANCE_COST",
"بهای تمام‌شده نمی‌تواند منفی باشد",
http_status=400,
)
previous_warehouse_id = change.get("previous_warehouse_id")
if previous_warehouse_id is not None:
by_key.pop((product_id, int(previous_warehouse_id)), None)
info: Dict[str, Any] = {
"movement": "in",
"warehouse_id": warehouse_id,
}
if cost_price > 0:
info["cost_price"] = float(cost_price)
by_key[(product_id, warehouse_id)] = {
"product_id": product_id,
"quantity": float(quantity),
"extra_info": info,
"description": f"موجودی اولیه - {change.get('product_name') or product_id}",
}
return list(by_key.values())
def apply_product_opening_balance_changes(
db: Session,
business_id: int,
user_id: int,
changes: List[Dict[str, Any]],
*,
fiscal_year_id: Optional[int] = None,
auto_commit: bool = True,
) -> Dict[str, Any]:
"""Apply all imported product opening balances with one document rewrite.
``changes`` must contain product_id, product_name, warehouse_id, quantity and
cost_price. previous_warehouse_id is optional and removes the old product
line when an update moves the default warehouse.
"""
if not changes:
return {}
ctx = _validate_ob_editable(db, business_id, fiscal_year_id)
fy_id = int(ctx["fiscal_year_id"])
_, fy_start, _ = _ensure_fiscal_year(db, business_id, fy_id)
from adapters.db.repositories.document_repository import DocumentRepository
repo = DocumentRepository(db)
existing_obj = _find_existing_ob_document(
db,
business_id,
fy_id,
for_update=True,
)
existing_doc = repo.to_dict_with_lines(existing_obj) if existing_obj is not None else None
if existing_doc:
account_lines, inventory_lines, settings = _document_lines_to_upsert_inputs(existing_doc)
currency_id = existing_doc.get("currency_id")
document_date = existing_doc.get("document_date")
else:
account_lines = []
inventory_lines = []
settings = {
"auto_balance_to_equity": True,
"equity_account_id": None,
"inventory_account_id": None,
}
business = db.query(Business).filter(Business.id == int(business_id)).first()
currency_id = getattr(business, "default_currency_id", None) if business else None
if not currency_id:
raise ApiError(
"CURRENCY_REQUIRED",
"ارز پیش‌فرض کسب‌وکار تنظیم نشده است",
http_status=400,
)
document_date = fy_start
inventory_lines = _ensure_inventory_lines_have_warehouse(
db,
business_id,
inventory_lines,
)
inventory_lines = _merge_product_opening_balance_changes(
inventory_lines,
changes,
)
inventory_account_id = settings.get("inventory_account_id")
equity_account_id = settings.get("equity_account_id")
auto_balance = bool(settings.get("auto_balance_to_equity", True))
if inventory_lines and not inventory_account_id:
inventory_account_id = _resolve_default_inventory_account_id(db)
if auto_balance and not equity_account_id and (account_lines or inventory_lines):
equity_account_id = _resolve_default_equity_account_id(db)
payload: Dict[str, Any] = {
"fiscal_year_id": fy_id,
"currency_id": int(currency_id),
"document_date": document_date,
"account_lines": account_lines,
"inventory_lines": inventory_lines,
"auto_balance_to_equity": auto_balance,
}
if inventory_account_id:
payload["inventory_account_id"] = int(inventory_account_id)
if equity_account_id:
payload["equity_account_id"] = int(equity_account_id)
return upsert_opening_balance(
db,
business_id,
int(user_id),
payload,
auto_commit=auto_commit,
return_details=False,
)
def _merged_track_inventory(payload: Any, existing: Optional[Product]) -> bool:
if hasattr(payload, "track_inventory") and payload.track_inventory is not None:
return bool(payload.track_inventory)

View file

@ -381,7 +381,14 @@ def _validate_unit_string(unit: Optional[str]) -> Optional[str]:
def _upsert_attributes(db: Session, product_id: int, business_id: int, attribute_ids: Optional[List[int]], auto_commit: bool = True) -> None:
def _upsert_attributes(
db: Session,
product_id: int,
business_id: int,
attribute_ids: Optional[List[int]],
auto_commit: bool = True,
validate_ids: bool = True,
) -> None:
"""
ایجاد یا به‌روزرسانی ویژگی‌های کالا
@ -399,11 +406,13 @@ def _upsert_attributes(db: Session, product_id: int, business_id: int, attribute
if auto_commit:
db.commit()
return
valid_ids = [
a.id for a in db.query(ProductAttribute.id, ProductAttribute.business_id)
.filter(ProductAttribute.id.in_(attribute_ids), ProductAttribute.business_id == business_id)
.all()
]
valid_ids = list(attribute_ids)
if validate_ids:
valid_ids = [
a.id for a in db.query(ProductAttribute.id, ProductAttribute.business_id)
.filter(ProductAttribute.id.in_(attribute_ids), ProductAttribute.business_id == business_id)
.all()
]
for aid in valid_ids:
db.add(ProductAttributeLink(product_id=product_id, attribute_id=aid))
if auto_commit:
@ -416,6 +425,10 @@ def create_product(
payload: ProductCreateRequest,
*,
defer_cache_invalidation: bool = False,
auto_commit: bool = True,
serialize_result: bool = True,
prevalidated_code_uniqueness: bool = False,
prevalidated_attribute_ids: bool = False,
) -> Dict[str, Any]:
"""
ایجاد کالا/خدمت جدید (با Retry Logic برای مدیریت Race Condition)
@ -450,9 +463,10 @@ def create_product(
if code_str and code_str != payload.name.strip():
code = code_str
manual_code = True
dup = db.query(Product).filter(and_(Product.business_id == business_id, Product.code == code)).first()
if dup:
raise ApiError("DUPLICATE_PRODUCT_CODE", "کد کالا/خدمت تکراری است", http_status=400)
if not prevalidated_code_uniqueness:
dup = db.query(Product).filter(and_(Product.business_id == business_id, Product.code == code)).first()
if dup:
raise ApiError("DUPLICATE_PRODUCT_CODE", "کد کالا/خدمت تکراری است", http_status=400)
# اگر کد خالی است یا برابر نام کالا است، کد خودکار تولید کن
if not code:
@ -525,17 +539,31 @@ def create_product(
# _upsert_attributes را بدون commit صدا می‌زنیم تا همه چیز در یک transaction باشد
logger.debug(f"[CREATE_PRODUCT] Upserting attributes - attribute_ids={payload.attribute_ids}")
_upsert_attributes(db, obj.id, business_id, payload.attribute_ids, auto_commit=False)
_upsert_attributes(
db,
obj.id,
business_id,
payload.attribute_ids,
auto_commit=False,
validate_ids=not prevalidated_attribute_ids,
)
upsert_product_suppliers(db, obj.id, business_id, payload.suppliers, auto_commit=False)
# Commit همه چیز (product و attributes)
logger.info(f"[CREATE_PRODUCT] Committing transaction for product ID={obj.id}...")
db.commit()
logger.info(f"[CREATE_PRODUCT] ✅ Transaction COMMITTED successfully for product ID={obj.id}")
db.refresh(obj) # Refresh برای دریافت اطلاعات کامل
logger.debug(f"[CREATE_PRODUCT] Product refreshed - final code='{obj.code}', name='{obj.name}'")
# مسیرهای تک‌کالا مثل گذشته commit می‌کنند؛ ایمپورت گروهی commit را
# تا پایان Unit of Work عقب می‌اندازد.
if auto_commit:
logger.info(f"[CREATE_PRODUCT] Committing transaction for product ID={obj.id}...")
db.commit()
logger.info(f"[CREATE_PRODUCT] ✅ Transaction COMMITTED successfully for product ID={obj.id}")
db.refresh(obj)
logger.debug(f"[CREATE_PRODUCT] Product refreshed - final code='{obj.code}', name='{obj.name}'")
data = _to_dict(obj, db)
data = _to_dict(obj, db) if serialize_result else {
"id": obj.id,
"name": obj.name,
"code": obj.code,
"default_warehouse_id": obj.default_warehouse_id,
}
# enrich titles from payload if provided
if getattr(payload, 'main_unit_title', None):
data["main_unit_title"] = str(getattr(payload, 'main_unit_title'))
@ -557,7 +585,14 @@ def create_product(
# خطای تکراری بودن کد (UniqueConstraint violation)
logger.warning(f"[CREATE_PRODUCT] IntegrityError caught (attempt {retry_count + 1}): {e}")
logger.debug(f"[CREATE_PRODUCT] Rolling back transaction...")
db.rollback()
if auto_commit:
db.rollback()
else:
raise ApiError(
"PRODUCT_IMPORT_WRITE_CONFLICT",
"تداخل هم‌زمان هنگام ثبت کالا رخ داد؛ ایمپورت دوباره بررسی شود",
http_status=409,
) from e
err_txt = str(getattr(e, "orig", e)).lower()
if "uq_product_general_barcode_business_token" in err_txt or "product_general_barcode_aliases" in err_txt:
raise ApiError(
@ -735,6 +770,10 @@ def update_product(
*,
defer_cache_invalidation: bool = False,
user_id: Optional[int] = None,
auto_commit: bool = True,
serialize_result: bool = True,
prevalidated_code_uniqueness: bool = False,
prevalidated_attribute_ids: bool = False,
) -> Optional[Dict[str, Any]]:
repo = ProductRepository(db)
obj = db.get(Product, product_id)
@ -748,7 +787,7 @@ def update_product(
code_str = payload.code.strip() if isinstance(payload.code, str) else str(payload.code).strip()
if code_str: # فقط اگر کد خالی نباشد
code_value = code_str
if code_value != obj.code: # اگر کد تغییر کرده
if code_value != obj.code and not prevalidated_code_uniqueness: # اگر کد تغییر کرده
dup = db.query(Product).filter(and_(Product.business_id == business_id, Product.code == code_value, Product.id != product_id)).first()
if dup:
raise ApiError("DUPLICATE_PRODUCT_CODE", "کد کالا/خدمت تکراری است", http_status=400)
@ -937,7 +976,14 @@ def update_product(
if gb_handled:
replace_general_barcode_aliases(db, business_id, product_id, general_tokens)
_upsert_attributes(db, product_id, business_id, payload.attribute_ids, auto_commit=False)
_upsert_attributes(
db,
product_id,
business_id,
payload.attribute_ids,
auto_commit=False,
validate_ids=not prevalidated_attribute_ids,
)
if "suppliers" in fields_set:
upsert_product_suppliers(db, product_id, business_id, payload.suppliers, auto_commit=False)
@ -959,16 +1005,17 @@ def update_product(
user_id=user_id,
)
try:
db.commit()
except IntegrityError:
db.rollback()
raise ApiError(
"GENERAL_BARCODE_CONFLICT",
"بارکد عمومی تکراری است یا با دادهٔ دیگر در تداخل است",
http_status=409,
)
db.refresh(updated)
if auto_commit:
try:
db.commit()
except IntegrityError:
db.rollback()
raise ApiError(
"GENERAL_BARCODE_CONFLICT",
"بارکد عمومی تکراری است یا با دادهٔ دیگر در تداخل است",
http_status=409,
)
db.refresh(updated)
if not defer_cache_invalidation:
old_category_id = obj.category_id if obj else None
@ -987,7 +1034,12 @@ def update_product(
)
invalidate_public_catalog_caches()
data = _to_dict(updated, db)
data = _to_dict(updated, db) if serialize_result else {
"id": updated.id,
"name": updated.name,
"code": updated.code,
"default_warehouse_id": updated.default_warehouse_id,
}
return {"message": "PRODUCT_UPDATED", "data": data}
@ -2676,5 +2728,3 @@ def get_inventory_stock_report(
"has_prev": current_page > 1,
}
}

View file

@ -46,6 +46,7 @@ _INVOICE_TYPE_LABELS_FA = {
_SOURCE_TYPE_LABELS_FA = {
"manual": "دستی",
"invoice": "فاکتور",
"opening_balance": "تراز افتتاحیه",
"api": "API",
"goods_expense_income": "کالای هزینه/درآمد شده",
}
@ -1772,6 +1773,129 @@ def create_manual_warehouse_document(
return wh
def create_linked_warehouse_document_bulk(
db: Session,
business_id: int,
user_id: int,
*,
doc_type: str,
document_date: date,
warehouse_id: int,
lines: List[Dict[str, Any]],
source_type: str,
source_document_id: int,
extra_info: Optional[Dict[str, Any]] = None,
) -> WarehouseDocument:
"""ساخت bulk سند انبار متصل به یک سند مبدا، بدون commit.
این مسیر برای عملیات سیستمی مانند همگام‌سازی موجودی افتتاحیه است؛ برخلاف
``create_manual_warehouse_document`` محصولات و انبار را به‌صورت گروهی
اعتبارسنجی می‌کند تا ایمپورت بزرگ به N+1 query تبدیل نشود.
"""
if doc_type not in ("receipt", "issue"):
raise ApiError("INVALID_DOC_TYPE", "نوع سند لینک‌شده باید receipt یا issue باشد", http_status=400)
if not lines:
raise ApiError("LINES_REQUIRED", "سند انبار باید حداقل یک خط داشته باشد", http_status=400)
warehouse = db.query(Warehouse).filter(
and_(Warehouse.id == int(warehouse_id), Warehouse.business_id == int(business_id))
).first()
if not warehouse:
raise ApiError("WAREHOUSE_NOT_FOUND", "انبار انتخاب‌شده یافت نشد", http_status=404)
fy = _get_current_fiscal_year(db, int(business_id))
if document_date < fy.start_date or (fy.end_date and document_date > fy.end_date):
raise ApiError(
"DATE_OUT_OF_RANGE",
f"تاریخ باید در بازه سال مالی ({fy.start_date} تا {fy.end_date or 'نامحدود'}) باشد",
http_status=400,
)
product_ids = sorted({int(line.get("product_id")) for line in lines if line.get("product_id")})
products = db.query(Product.id).filter(
Product.business_id == int(business_id),
Product.id.in_(product_ids),
).all()
found_product_ids = {int(row.id) for row in products}
missing_product_ids = [pid for pid in product_ids if pid not in found_product_ids]
if missing_product_ids:
raise ApiError(
"PRODUCT_NOT_FOUND",
f"کالای سند انبار یافت نشد: {missing_product_ids[0]}",
http_status=404,
)
normalized_lines: List[tuple[int, Decimal, Dict[str, Any]]] = []
for index, line in enumerate(lines, start=1):
pid = line.get("product_id")
if pid is None:
raise ApiError("PRODUCT_REQUIRED", f"خط {index}: شناسه کالا الزامی است", http_status=400)
try:
quantity = Decimal(str(line.get("quantity") or 0))
except Exception as exc:
raise ApiError("INVALID_QUANTITY", f"خط {index}: تعداد نامعتبر است", http_status=400) from exc
if quantity <= 0 or quantity > _MAX_WAREHOUSE_LINE_QUANTITY:
raise ApiError("INVALID_QUANTITY", f"خط {index}: تعداد معتبر نیست", http_status=400)
line_extra = line.get("extra_info")
normalized_lines.append(
(int(pid), quantity, dict(line_extra) if isinstance(line_extra, dict) else {})
)
wh: Optional[WarehouseDocument] = None
for attempt in range(10):
code = _generate_warehouse_document_code(db, int(business_id), document_date)
try:
with db.begin_nested():
wh = WarehouseDocument(
business_id=int(business_id),
fiscal_year_id=int(fy.id),
code=code,
document_date=document_date,
status="draft",
doc_type=doc_type,
warehouse_id_from=int(warehouse_id) if doc_type == "issue" else None,
warehouse_id_to=int(warehouse_id) if doc_type == "receipt" else None,
source_type=str(source_type),
source_document_id=int(source_document_id),
created_by_user_id=int(user_id),
extra_info=dict(extra_info or {}),
)
db.add(wh)
db.flush()
break
except IntegrityError as exc:
if _is_duplicate_warehouse_document_code_error(exc) and attempt < 9:
continue
raise
if wh is None:
raise ApiError("WAREHOUSE_CODE_CONFLICT", "Failed to generate unique warehouse document code", http_status=500)
movement = "in" if doc_type == "receipt" else "out"
db.add_all([
WarehouseDocumentLine(
warehouse_document_id=int(wh.id),
product_id=pid,
warehouse_id=int(warehouse_id),
movement=movement,
quantity=quantity,
extra_info=line_extra or None,
)
for pid, quantity, line_extra in normalized_lines
])
db.flush()
invalidate_warehouse_docs_cache(
business_id=int(business_id),
fiscal_year_id=int(fy.id),
doc_type=doc_type,
warehouse_id=int(warehouse_id),
status="draft",
document_id=int(wh.id),
)
return wh
def update_warehouse_document(
db: Session,
business_id: int,
@ -2672,6 +2796,14 @@ def post_warehouse_document(
lines = db.query(WarehouseDocumentLine).filter(WarehouseDocumentLine.warehouse_document_id == wh.id).all()
if not lines:
raise ApiError("NO_LINES", "حواله باید حداقل یک خط داشته باشد", http_status=400)
tracked_product_ids = {
int(row.id)
for row in db.query(Product.id).filter(
Product.business_id == int(wh.business_id),
Product.id.in_({int(line.product_id) for line in lines if line.product_id is not None}),
Product.track_inventory == True,
).all()
}
# کنترل کسری برای خروج‌ها
outgoing_lines = []
@ -2683,9 +2815,8 @@ def post_warehouse_document(
if not ln.warehouse_id:
raise ApiError("WAREHOUSE_REQUIRED", "برای خطوط خروج، انبار باید مشخص باشد", http_status=400)
# بررسی اینکه محصول کنترل موجودی دارد یا نه
product = db.query(Product).filter(Product.id == ln.product_id).first()
if product and product.track_inventory:
# وضعیت کنترل موجودی همه کالاهای سند بالاتر به‌صورت bulk خوانده شده است.
if int(ln.product_id) in tracked_product_ids:
outgoing_lines.append({
"product_id": ln.product_id,
"quantity": float(ln.quantity),
@ -2698,61 +2829,92 @@ def post_warehouse_document(
# کنترل کسری موجودی (با رعایت سیاست کسب‌وکار: فله / یونیک / انتقال)
if outgoing_lines:
biz = db.query(Business).filter(Business.id == int(wh.business_id)).first()
allow_bulk = bool(getattr(biz, "allow_negative_inventory_for_bulk", False)) if biz else False
allow_unique = bool(getattr(biz, "allow_negative_inventory_for_unique", False)) if biz else False
transfer_strict = bool(getattr(biz, "warehouse_transfer_require_positive_stock", True)) if biz else True
lines_to_check = filter_outgoing_lines_for_stock_enforcement(
db,
int(wh.business_id),
outgoing_lines,
allow_negative_for_bulk=allow_bulk,
allow_negative_for_unique=allow_unique,
warehouse_doc_type=getattr(wh, "doc_type", None),
transfer_require_positive_stock=transfer_strict,
)
if lines_to_check:
exclude_fin_doc: Optional[int] = None
exclude_inv_src: Optional[int] = None
if wh.source_document_id is not None:
src_id = int(wh.source_document_id)
st = (wh.source_type or "").strip().lower()
if st == "invoice":
exclude_fin_doc = src_id
exclude_inv_src = src_id
else:
# حواله‌های قدیمی: source_type خالی یا نادرست ولی سند مبدا فاکتور است
_src_row = (
db.query(Document.document_type)
.filter(
Document.id == src_id,
Document.business_id == int(wh.business_id),
)
.first()
# اصلاح کاهشی موجودی افتتاحیه یک عملیات کاملاً فیزیکی است. سند مالی
# افتتاحیه پیش از این مرحله به مقدار جدید رسیده و نباید برای کنترل خروج
# استفاده شود؛ در غیر این صورت جابه‌جایی انبار یا کاهش تعداد به‌اشتباه
# با کسری موجودی رد می‌شود.
if str(wh.source_type or "").strip().lower() == "opening_balance":
required_by_key: Dict[tuple[int, int], Decimal] = defaultdict(Decimal)
for line in outgoing_lines:
info = line.get("extra_info") or {}
warehouse_id = info.get("warehouse_id")
if warehouse_id is None:
continue
required_by_key[(int(line["product_id"]), int(warehouse_id))] += Decimal(
str(line.get("quantity") or 0)
)
for (product_id, warehouse_id), required_qty in required_by_key.items():
available_qty = get_physical_stock(
db,
int(wh.business_id),
product_id,
warehouse_id,
wh.document_date,
)
if available_qty < required_qty:
raise ApiError(
"INSUFFICIENT_PHYSICAL_STOCK",
"موجودی فیزیکی برای اصلاح موجودی افتتاحیه کافی نیست. "
f"موجودی: {available_qty}، موردنیاز: {required_qty}",
http_status=409,
details={"product_id": product_id, "warehouse_id": warehouse_id},
)
_sdt = _src_row[0] if _src_row is not None else None
if _sdt is not None and str(_sdt).startswith("invoice_"):
else:
biz = db.query(Business).filter(Business.id == int(wh.business_id)).first()
allow_bulk = bool(getattr(biz, "allow_negative_inventory_for_bulk", False)) if biz else False
allow_unique = bool(getattr(biz, "allow_negative_inventory_for_unique", False)) if biz else False
transfer_strict = bool(getattr(biz, "warehouse_transfer_require_positive_stock", True)) if biz else True
lines_to_check = filter_outgoing_lines_for_stock_enforcement(
db,
int(wh.business_id),
outgoing_lines,
allow_negative_for_bulk=allow_bulk,
allow_negative_for_unique=allow_unique,
warehouse_doc_type=getattr(wh, "doc_type", None),
transfer_require_positive_stock=transfer_strict,
)
if lines_to_check:
exclude_fin_doc: Optional[int] = None
exclude_inv_src: Optional[int] = None
if wh.source_document_id is not None:
src_id = int(wh.source_document_id)
st = (wh.source_type or "").strip().lower()
if st == "invoice":
exclude_fin_doc = src_id
exclude_inv_src = src_id
_ex_wh: Optional[List[int]] = None
if stock_exclude_warehouse_document_ids:
_ex_wh = sorted(
{int(x) for x in stock_exclude_warehouse_document_ids if x is not None}
)
if not _ex_wh:
_ex_wh = None
try:
_ensure_stock_sufficient(
db,
wh.business_id,
wh.document_date,
lines_to_check,
exclude_document_id=exclude_fin_doc,
exclude_invoice_source_document_id=exclude_inv_src,
exclude_warehouse_document_ids=_ex_wh,
)
except ApiError:
raise
else:
# حواله‌های قدیمی: source_type خالی یا نادرست ولی سند مبدا فاکتور است
_src_row = (
db.query(Document.document_type)
.filter(
Document.id == src_id,
Document.business_id == int(wh.business_id),
)
.first()
)
_sdt = _src_row[0] if _src_row is not None else None
if _sdt is not None and str(_sdt).startswith("invoice_"):
exclude_fin_doc = src_id
exclude_inv_src = src_id
_ex_wh: Optional[List[int]] = None
if stock_exclude_warehouse_document_ids:
_ex_wh = sorted(
{int(x) for x in stock_exclude_warehouse_document_ids if x is not None}
)
if not _ex_wh:
_ex_wh = None
try:
_ensure_stock_sufficient(
db,
wh.business_id,
wh.document_date,
lines_to_check,
exclude_document_id=exclude_fin_doc,
exclude_invoice_source_document_id=exclude_inv_src,
exclude_warehouse_document_ids=_ex_wh,
)
except ApiError:
raise
# برای حواله‌های transfer، به‌روزرسانی instance های کالاهای یونیک
if wh.doc_type == "transfer":
@ -2947,36 +3109,38 @@ def post_warehouse_document(
# شناسایی بهای تمام‌شده قطعی روی خطوط فاکتور و ثبت COGS در دفتر کل
# ورک‌فلو: موجودی کم (reorder_point)
try:
from app.services.workflow.workflow_trigger_service import maybe_fire_inventory_low_triggers
# افتتاحیه یک بارگذاری اولیه و بالقوه هزاران‌خطی است؛ تریگر موجودی کم
# برای هر خط آن معنا ندارد و مسیر ایمپورت را دوباره N+1 می‌کند.
if str(wh.source_type or "").strip().lower() != "opening_balance":
try:
from app.services.workflow.workflow_trigger_service import maybe_fire_inventory_low_triggers
seen_pairs: set[tuple[int, Optional[int]]] = set()
for ln in lines:
if not ln.product_id:
continue
prod = db.query(Product).filter(Product.id == int(ln.product_id)).first()
if not prod or not getattr(prod, "track_inventory", False):
continue
wid = int(ln.warehouse_id) if ln.warehouse_id else None
key = (int(ln.product_id), wid)
if key in seen_pairs:
continue
seen_pairs.add(key)
maybe_fire_inventory_low_triggers(
db,
int(business_id),
int(ln.product_id),
wid,
user_id=getattr(wh, "created_by_user_id", None),
seen_pairs: set[tuple[int, Optional[int]]] = set()
for ln in lines:
if not ln.product_id:
continue
prod = db.query(Product).filter(Product.id == int(ln.product_id)).first()
if not prod or not getattr(prod, "track_inventory", False):
continue
wid = int(ln.warehouse_id) if ln.warehouse_id else None
key = (int(ln.product_id), wid)
if key in seen_pairs:
continue
seen_pairs.add(key)
maybe_fire_inventory_low_triggers(
db,
int(business_id),
int(ln.product_id),
wid,
user_id=getattr(wh, "created_by_user_id", None),
)
except Exception as inv_wf_err:
logger.warning(
"warehouse_workflow_inventory_low_failed wh_id=%s err=%s",
wh.id,
inv_wf_err,
exc_info=True,
)
except Exception as inv_wf_err:
logger.warning(
"warehouse_workflow_inventory_low_failed wh_id=%s err=%s",
wh.id,
inv_wf_err,
exc_info=True,
)
# شناسایی بهای تمام‌شده قطعی روی خطوط فاکتور مبدأ (در صورت تنظیم کسب‌وکار)
try:
@ -4270,4 +4434,3 @@ def _include_inventory_stock_row(
if include_zero or stock != 0:
return True
return has_warehouse_history

View file

@ -0,0 +1,243 @@
# معماری و عملیات ایمپورت Excel کالا و خدمت
## هدف
این سند رفتار فنی endpoint زیر را پس از بهینه‌سازی ثبت می‌کند:
```text
POST /api/v1/products/business/{business_id}/import/excel
```
هدف، حذف query و commit ردیف‌به‌ردیف، حفظ قرارداد فعلی API و جلوگیری از
بازسازی چندبارهٔ سند تراز افتتاحیه است. محدودیت فعلی فایل همچنان ۱۵ مگابایت و
۵۰۰۰ ردیف داده است.
## مشکل پیشین
در پیاده‌سازی قبلی، یک ردیف فایل رسمی ممکن بود همان کالا را تا پنج مرتبه از
دیتابیس پیدا کند. وجود ستون‌های «تعداد اولیه» در header نیز حتی با سلول خالی،
lookup اضافه ایجاد می‌کرد. اجرای واقعی هر کالا را جداگانه commit، serialize و
cache-invalidate می‌کرد.
در ردیف‌های دارای تعداد اولیه، سرویس تک‌کالا سند افتتاحیه را می‌خواند، کل خطوط
آن را پیمایش می‌کرد و تمام خطوط را دوباره می‌نوشت. برای `N` ردیف، حجم کار
تقریباً برابر `1 + 2 + ... + N` بود.
## جریان جدید
```text
خواندن workbook
↓
ساخت index مراجع و کالاهای منطبق
↓
اعتبارسنجی و تولید plan بدون write ردیف‌به‌ردیف
↓
flush کالاهای معتبر در یک Unit of Work
↓
ادغام همهٔ تغییرات تعداد اولیه در حافظه
↓
یک بازنویسی سند افتتاحیه
↓
همگام‌سازی اختلافی رسید/حواله انبار به تفکیک انبار
↓
یک commit و سپس یک cache invalidation
```
### lookup کالا
کلیدهای غیرخالی فایل ابتدا جمع‌آوری و با queryهای `IN` حداکثر ۵۰۰تایی
بارگذاری می‌شوند. نتیجه براساس کد یا نام normalize‌شده در حافظه نگهداری می‌شود.
dry-run، تصمیم insert/update و اجرای واقعی همگی از همین index استفاده می‌کنند.
تطبیق براساس کد همچنان گزینهٔ پیشنهادی است. در حالت تطبیق با نام، چند کالای
موجود با نام یکسان خطای ابهام تولید می‌کنند. تکرار کلید تطبیق داخل خود فایل نیز
پیش از write به‌عنوان خطای ردیف گزارش می‌شود.
### مراجع
دسته‌بندی‌ها، ویژگی‌ها، انبارها، انواع مالیات، واحدهای مالیاتی و ارزهای فعال
یک‌بار بارگذاری و index می‌شوند. ایجاد خودکار دسته یا ویژگی با `flush` انجام
می‌شود و commit آن به پایان Unit of Work موکول می‌گردد.
### تعداد اولیه
صرف وجود ستون‌های تعداد اولیه باعث پردازش نمی‌شود. فقط ردیفی که مقدار تعداد یا
بهای تمام‌شده دارد وارد plan افتتاحیه می‌شود.
در اجرای واقعی:
1. امکان ویرایش تراز افتتاحیه پیش از write بررسی می‌شود.
2. سند موجود با قفل ردیفی خوانده می‌شود.
3. خطوط کالا براساس `(product_id, warehouse_id)` در map قرار می‌گیرند.
4. همهٔ تغییرات فایل روی map اعمال می‌شوند.
5. خطوط اشخاص، بانک، صندوق، تنخواه و سایر حساب‌ها حفظ می‌شوند.
6. سند یک‌بار متوازن و یک‌بار نوشته می‌شود.
پیچیدگی ادغام خطوط `O(existing_lines + imported_changes)` است.
## رفتار تراکنشی
- dry-run هیچ write، commit یا cache invalidation ندارد.
- خطاهای validation پیش از apply به تفکیک ردیف در پاسخ گزارش می‌شوند.
- ردیف‌های معتبر در یک transaction اعمال می‌شوند.
- ایجاد کالا و ثبت تعداد اولیه در همان transaction هستند.
- خطای دیتابیس یا خطای غیرمنتظره در apply باعث rollback کل apply می‌شود؛ بنابراین
کالای ساخته‌شده بدون تعداد اولیه باقی نمی‌ماند.
- conflict هم‌زمان یا نقض constraint با خطای `IMPORT_WRITE_CONFLICT` و HTTP 409
برگردانده می‌شود.
- cache محصولات و کاتالوگ عمومی فقط پس از commit موفق و فقط یک‌بار پاک می‌شود.
این رفتار عمداً atomic‌تر از پیاده‌سازی قدیمی است که commitهای جزئی متعدد داشت.
## قرارداد API
فیلدهای موجود پاسخ حفظ شده‌اند:
- `summary`
- `errors`
- `reference_summary`
- `preview` در dry-run
- `warehouse_sync` برای خلاصه اسناد فیزیکی افتتاحیه
فیلد افزوده‌شدهٔ `performance` شامل مقادیر زیر برحسب میلی‌ثانیه است:
- `parse_ms`
- `validation_ms`
- `apply_ms`
- `total_ms`
این داده‌ها secret، محتوای سلول‌ها یا اطلاعات اتصال را ثبت نمی‌کنند.
## تغییرات سرویس‌های مشترک
سرویس‌های ایجاد و ویرایش کالا گزینه‌های زیر را دارند و مقدار پیش‌فرض آن‌ها رفتار
endpointهای تک‌کالا را حفظ می‌کند:
- `auto_commit=True`
- `serialize_result=True`
- `defer_cache_invalidation=False`
- `prevalidated_code_uniqueness=False`
- `prevalidated_attribute_ids=False`
repositoryهای دسته، ویژگی و سند نیز `auto_commit=True` دارند. فقط importer این
گزینه‌ها را برای Unit of Work گروهی غیرفعال می‌کند.
## فایل‌های اصلی
- `adapters/api/v1/products.py`: orchestration، plan و پاسخ
- `app/services/product_excel_import_normalize.py`: index تطبیق و ارز
- `app/services/product_excel_import_opening_balance.py`: validation ردیف
- `app/services/product_opening_balance_service.py`: ادغام و apply گروهی
- `app/services/opening_balance_warehouse_sync_service.py`: sync اختلافی موجودی فیزیکی
- `app/services/product_service.py`: write بدون commit/serialization اجباری
- `app/services/opening_balance_service.py`: upsert قابل استفاده در transaction بیرونی
- `adapters/db/repositories/document_repository.py`: write با commit اختیاری
## کلاینت
دیالوگ Flutter برای این درخواست timeout دریافت پنج‌دقیقه‌ای دارد. timeout بلندتر
راه‌حل کارایی نیست؛ فقط از قطع زودهنگام درخواست معتبر در شبکه‌های کند جلوگیری
می‌کند. UI همچنان باید از اجرای دوبارهٔ درخواست هنگام loading جلوگیری کند.
endpoint به‌صورت sync در threadpool اجرا می‌شود تا parse اکسل و SQLAlchemy
همگام، event loop درخواست‌های دیگر را مسدود نکنند. این رفتار فقط برای همین
endpoint با گزینهٔ `offload_sync` روی گارد دسترسی فعال شده است.
صفحهٔ تراز افتتاحیه خط سیستمی «بستن اختلاف تراز افتتاحیه» را از ردیف‌های قابل
ویرایش مخفی می‌کند. برای جلوگیری از نمایش اشتباه بستانکار صفر پس از ایمپورت،
محاسبهٔ جمع UI در صورت فعال‌بودن auto-balance و انتخاب حساب حقوق صاحبان سهام،
اثر همان خط مخفی را روی سمت مقابل بازسازی می‌کند. اگر auto-balance خاموش باشد یا
حساب مقابل انتخاب نشده باشد، اختلاف خام همچنان نمایش داده می‌شود.
## validation و تست
تست‌های لازم:
```text
tests/test_product_excel_import.py
tests/test_product_excel_import_opening_balance.py
tests/test_product_opening_balance.py
```
اجرای تست در محیط توسعه فقط از wrapper ایزوله مجاز است:
```bash
/home/mohammad/projects/hesabix/dev-env/toolkit/scripts/run-tests-isolated.sh -- \
tests/test_product_excel_import.py \
tests/test_product_excel_import_opening_balance.py \
tests/test_product_opening_balance.py -q
```
wrapper روی worktree کثیف عمداً اجرا نمی‌شود. برای validation نهایی باید تغییرات
در یک branch/worktree تمیز و مطابق رویهٔ Git پروژه قرار گیرند؛ تست مستقیم با
`.env` متصل به `hesabix_dev` ممنوع است.
### نتیجهٔ اعتبارسنجی ۱۴۰۵/۰۷/۰۹ (۲۰۲۶-۱۰-۰۱)
- compilation فایل‌های Python تغییرکرده موفق بود.
- بررسی محدود Ruff برای `F821`، `F822`، `F823` و `E902` موفق بود.
- `git diff --check` برای فایل‌های این تغییر موفق بود.
- snapshot تمیز و مستقل برای اجرای wrapper ساخته شد و مرحله‌های `--check` و
`--dry-run` آن موفق بودند.
- سه فایل تست هدف با wrapper ایزوله اجرا شدند: `28 passed` و `176 warnings` در
`25.31s`. دیتابیس تصادفی `hesabix_test_*` پس از اجرا حذف شد و دیتابیس توسعه
استفاده یا تغییر داده نشد.
- این سه فایل DB-free هستند. migration chain فعلی repository نمی‌تواند یک
PostgreSQL کاملاً خالی را بسازد، زیرا revision
`20250226_000002_add_bale_messenger_support` وجود جدول `users` را فرض می‌کند.
بنابراین برای این اجرای unit-test فقط مرحلهٔ `alembic upgrade head` در
snapshot تست bypass شد. این نتیجه پوشش integration دیتابیس یا endpoint HTTP
محسوب نمی‌شود و مشکل baseline migration باید جداگانه اصلاح شود.
microbenchmark تابع ادغام گروهی روی همان محیط، با هفت تکرار و گزارش median:
| خطوط موجود | تغییرات ورودی | median | خطوط خروجی |
|---:|---:|---:|---:|
| ۱۰۰۰ | ۱۰۰۰ | ۱۵٫۰۸۲ ms | ۱۰۰۰ |
| ۵۰۰۰ | ۵۰۰۰ | ۸۶٫۵۶۱ ms | ۵۰۰۰ |
| ۱۰۰۰۰ | ۱۰۰۰۰ | ۱۸۸٫۹۳۱ ms | ۱۰۰۰۰ |
این اعداد فقط هزینهٔ `_merge_product_opening_balance_changes` را اندازه می‌گیرند
و benchmark انتها‌به‌انتهای HTTP/Excel/PostgreSQL نیستند. رشد مشاهده‌شده با
پیچیدگی خطی طراحی جدید سازگار است، اما معیار «پنج برابر سریع‌تر برای ۱۰۰۰ ردیف»
فقط پس از اجرای benchmark انتها‌به‌انتها در محیط staging قابل تأیید است.
### معیار پذیرش عملیاتی
- فایل بدون تعداد اولیه هیچ فراخوانی سرویس سند افتتاحیه نداشته باشد.
- lookup کالا متناسب با تعداد chunkها رشد کند، نه تعداد ردیف‌ها.
- اجرای واقعی یک commit نهایی داشته باشد.
- سند افتتاحیه برای کل فایل یک بار apply شود.
- برای هر انبار و جهت حرکت حداکثر یک سند گروهی ساخته شود.
- تکرار فایل یکسان هیچ حرکت انبار تازه‌ای نسازد.
- موجودی افتتاحیه در محاسبات ترکیبی دو بار شمرده نشود.
- cache invalidation برای کل import یک بار انجام شود.
- نتیجهٔ dry-run و real import برای insert/update/skip یکسان باشد.
- زمان ۱۰۰۰ ردیف روی یک محیط ثابت حداقل پنج برابر بهتر از baseline قدیمی باشد.
## پایش و عیب‌یابی
در log، summary و زمان مراحل را با `business_id` و تعداد ردیف ثبت کنید؛ نام فایل،
مقادیر سلول‌ها، token و اطلاعات اتصال نباید log شوند. برای تشخیص کندی، ابتدا
`performance` را بررسی کنید:
- `parse_ms` بالا: اندازه/پیچیدگی workbook یا openpyxl
- `validation_ms` بالا: تعداد مراجع، نام‌های مبهم یا validationهای دامنه
- `apply_ms` بالا: constraint، barcode، sync موجودی یا سند افتتاحیه
## rollback
این تغییر migration دیتابیس ندارد. rollback کد با بازگرداندن orchestration قدیمی
ممکن است، اما به‌دلیل خطر partial commit توصیه نمی‌شود. اگر rollback عملیاتی لازم
شد، endpoint ایمپورت موقتاً غیرفعال شود و سپس نسخهٔ قبلی deploy گردد؛ هیچ سند یا
کالایی برای rollback نباید به‌صورت دستی حذف شود. قبل از هر اصلاح داده، backup و
audit نتیجهٔ import بررسی شود.
## محدودیت‌های آگاهانه
- تولید کد خودکار همچنان برای هر کالای بدون کد نیازمند تخصیص یکتاست.
- validation بارکد عمومی و sync تغییر کنترل موجودی عمداً حذف نشده‌اند؛ این‌ها
قواعد دامنه‌اند و در صورت نیاز باید در فاز جداگانه batch شوند.
- endpoint هنوز نتیجه را در همان درخواست HTTP برمی‌گرداند. اگر پس از benchmark فایل ۵۰۰۰ ردیفی
طولانی بماند، مرحلهٔ بعد انتقال apply به job پس‌زمینه با progress و idempotency
است؛ افزایش بیشتر timeout جایگزین آن نیست.

View file

@ -0,0 +1,130 @@
# همگام‌سازی موجودی افتتاحیه با اسناد انبار در ایمپورت کالا
## هدف
ردیف دارای «تعداد اولیه» در ایمپورت Excel دو اثر هماهنگ دارد:
1. ارزش و تعداد افتتاحیه در سند حسابداری `opening_balance` ثبت می‌شود.
2. موجودی فیزیکی با سند انبار قطعی (`WarehouseDocument`) در انبار انتخاب‌شده
ثبت می‌شود.
سند مالی منبع ارزش‌گذاری است و سند انبار منبع گزارش موجودی فیزیکی. سند انبار
با `source_type=opening_balance` و `source_document_id` به سند مالی متصل می‌شود.
## قواعد دامنه
- تعداد مثبت با `receipt` و `movement=in` ثبت می‌شود.
- کاهش تعداد در ایمپورت مجدد فقط به اندازه اختلاف با `issue` و
`movement=out` ثبت می‌شود.
- اسناد به تفکیک انبار گروه‌بندی می‌شوند؛ برای هر کالا سند جدا ساخته نمی‌شود.
- اسناد به‌صورت خودکار `posted` می‌شوند.
- تاریخ سند انبار همان تاریخ سند افتتاحیه است.
- تکرار فایل یکسان سند یا حرکت جدید ایجاد نمی‌کند.
- خدمت، کالای فاقد کنترل موجودی و کالای یونیک بدون اطلاعات instance پذیرفته
نمی‌شوند.
- کاربر علاوه بر `products.edit` و `opening_balance.edit` به
`inventory.write` نیاز دارد.
## الگوریتم اختلافی
برای کالاهای متاثر، مقدار مطلوب از خطوط سند افتتاحیه و مقدار همگام‌شده از مجموع
رسید/حواله‌های `posted` متصل به همان سند خوانده می‌شود:
```text
delta(product, warehouse) = desired_opening - linked_physical
```
- `delta > 0`: رسید انبار
- `delta < 0`: حواله خروج
- `delta = 0`: بدون عملیات
تغییر انبار به‌صورت خروج از انبار قبلی و ورود به انبار جدید دیده می‌شود. تمام
اسناد یک اجرای ایمپورت دارای `import_batch_id` مشترک در `extra_info` هستند.
## جلوگیری از شمارش دوگانه
محاسبه موجودی مالی/قابل‌استفاده خطوط سند افتتاحیه را می‌خواند. بنابراین خطوط
WarehouseDocument با `source_type=opening_balance` در این محاسبه دوباره اضافه
نمی‌شوند. در مقابل، گزارش موجودی فیزیکی فقط WarehouseDocumentهای `posted` را
می‌خواند و رسید افتتاحیه را لحاظ می‌کند.
برای مقدار افتتاحیه ۱۰، خروجی مورد انتظار:
```text
financial = 10
physical = 10
available = 10
```
## تراکنش و خطا
ایجاد/ویرایش کالا، بازنویسی افتتاحیه، ساخت خطوط انبار و قطعی‌سازی در یک
transaction انجام می‌شود. هر خطا کل عملیات را rollback می‌کند. cache فقط بعد از
موفقیت نهایی invalidate می‌شود.
برای کاهش افتتاحیه، کنترل کسری براساس موجودی فیزیکی انجام می‌شود؛ زیرا مقدار
مالی افتتاحیه در همان transaction به مقدار جدید رسیده است.
اگر برای کالای متاثر حرکت `posted` غیرمرتبط با افتتاحیه وجود داشته باشد، sync با
خطای `OPENING_BALANCE_WAREHOUSE_HISTORY_EXISTS` متوقف می‌شود. در این وضعیت باید
اختلاف با رسید، حواله یا تعدیل مستقل ثبت شود؛ بازنویسی گذشته مجاز نیست.
## داده‌های قدیمی
برای اسناد افتتاحیه قدیمی، اگر هیچ گردش انبار دیگری وجود نداشته باشد، نخستین
ایمپورت بعد از انتشار می‌تواند رسید لینک‌شده را ایجاد کند. اگر گردش قبلی وجود
داشته باشد، سیستم عمداً از حدس‌زدن منشأ موجودی خودداری و عملیات را متوقف می‌کند.
هرگونه backfill عمومی باید ابزار جداگانه با preview و تایید مدیر داشته باشد.
## پاسخ API
فیلد `warehouse_sync` به پاسخ ایمپورت اضافه شده است:
```json
{
"receipts_created": 2,
"issues_created": 0,
"lines_created": 1500,
"unchanged_lines": 0,
"posted": true,
"document_ids": [101, 102],
"import_batch_id": "..."
}
```
در dry-run تعداد ردیف‌ها و انبارهای کاندید گزارش می‌شود و هیچ سندی ساخته یا
قطعی نمی‌شود.
## تست‌های رگرسیون
- اولین ایمپورت: رسید کامل
- تکرار همان فایل: بدون delta
- افزایش و کاهش: فقط مقدار اختلاف
- تغییر انبار: issue و receipt متناظر
- گروه‌بندی چند کالا در یک سند برای هر انبار
- الزام دسترسی `inventory.write`
- جلوگیری از شمارش دوگانه
- rollback در شکست ساخت یا post سند
- رد کالای یونیک بدون instance
### نتیجه اعتبارسنجی ۱۴۰۵/۰۷/۱۰ (۲۰۲۶-۱۰-۰۲)
- compilation فایل‌های Python تغییرکرده موفق بود.
- Ruff محدود برای خطاهای import/name روی فایل‌های جدید و تست‌ها موفق بود.
- `git diff --check` موفق بود.
- گاردهای `--check` و `--dry-run` ابزار تست ایزوله موفق بودند.
- شش فایل تست هدف شامل تست‌های ایمپورت، افتتاحیه، sync اختلافی، موجودی فیزیکی
و جلوگیری از شمارش دوباره اجرا شدند: `49 passed` و `177 warnings` در
`31.97s`.
- دیتابیس تصادفی `hesabix_test_*` بعد از اجرا خودکار حذف شد و دیتابیس توسعه
استفاده یا تغییر داده نشد.
- به‌دلیل مشکل شناخته‌شده baseline migration، اجرای DB-free با shim مرحله
Alembic انجام شد؛ این نتیجه تست integration واقعی PostgreSQL محسوب نمی‌شود.
- SDK محلی Flutter/Dart در PATH محیط موجود نبود؛ بنابراین validation خودکار UI
در این محیط اجرا نشد.
## rollback عملیاتی
این قابلیت migration دیتابیس ندارد. rollback کد، اسناد قبلاً ساخته‌شده را حذف
نمی‌کند. برای اصلاح یک اجرای نامعتبر باید از عملیات لغو رسمی سند انبار و ثبت سند
اصلاحی استفاده شود؛ حذف مستقیم WarehouseDocument یا DocumentLine ممنوع است.

View file

@ -0,0 +1,59 @@
"""تست منطق اختلافی رسید/حواله موجودی افتتاحیه."""
from decimal import Decimal
from app.services.opening_balance_warehouse_sync_service import (
_group_delta_lines,
build_opening_balance_warehouse_deltas,
)
def test_first_import_creates_full_receipt_delta() -> None:
desired = {(10, 1): Decimal("10")}
assert build_opening_balance_warehouse_deltas(desired, {}) == {
(10, 1): Decimal("10")
}
def test_repeating_same_import_is_idempotent() -> None:
desired = {(10, 1): Decimal("10")}
linked = {(10, 1): Decimal("10")}
assert build_opening_balance_warehouse_deltas(desired, linked) == {}
def test_quantity_increase_and_decrease_are_delta_only() -> None:
desired = {(10, 1): Decimal("15"), (20, 1): Decimal("7")}
linked = {(10, 1): Decimal("10"), (20, 1): Decimal("10")}
assert build_opening_balance_warehouse_deltas(desired, linked) == {
(10, 1): Decimal("5"),
(20, 1): Decimal("-3"),
}
def test_warehouse_move_becomes_issue_and_receipt() -> None:
desired = {(10, 2): Decimal("8")}
linked = {(10, 1): Decimal("8")}
deltas = build_opening_balance_warehouse_deltas(desired, linked)
assert deltas == {
(10, 1): Decimal("-8"),
(10, 2): Decimal("8"),
}
receipts, issues = _group_delta_lines(deltas, {})
assert receipts[2][0]["product_id"] == 10
assert receipts[2][0]["quantity"] == Decimal("8")
assert issues[1][0]["product_id"] == 10
assert issues[1][0]["quantity"] == Decimal("8")
def test_lines_are_grouped_by_warehouse_not_by_product() -> None:
deltas = {
(1, 5): Decimal("2"),
(2, 5): Decimal("3"),
(3, 6): Decimal("4"),
}
receipts, issues = _group_delta_lines(deltas, {})
assert issues == {}
assert len(receipts) == 2
assert len(receipts[5]) == 2
assert len(receipts[6]) == 1

View file

@ -10,6 +10,9 @@ from app.services.product_excel_import_normalize import (
map_headers,
parse_bool_strict,
parse_inventory_mode,
find_existing_product_in_index,
load_existing_product_match_index,
product_match_key,
provided_keys_from_raw,
select_products_import_worksheet,
)
@ -185,3 +188,56 @@ def test_category_full_path_walks_parents():
child = type("C", (), {"id": 2, "parent_id": 1, "title_translations": {"fa": "پلاستیک"}})()
by_id = {1: parent, 2: child}
assert category_full_path(child, by_id, "fa") == "مواد اولیه > پلاستیک"
def test_import_match_index_resolves_without_database_lookup():
product = type("P", (), {"id": 7, "code": " P-1 ", "name": " ماگ بزرگ "})()
code_index = {"P-1": [product]}
found, error = find_existing_product_in_index(code_index, "code", {"code": "P-1"})
assert error is None
assert found is product
name_key = product_match_key("name", {"name": "ماگ بزرگ"})
found, error = find_existing_product_in_index({name_key: [product]}, "name", {"name": "ماگ بزرگ"})
assert error is None
assert found is product
def test_import_match_index_rejects_ambiguous_name():
p1 = type("P", (), {"id": 1})()
p2 = type("P", (), {"id": 2})()
found, error = find_existing_product_in_index(
{"ماگ": [p1, p2]},
"name",
{"name": "ماگ"},
)
assert found is None
assert error is not None
def test_existing_product_lookup_is_chunked_not_per_row():
class _Query:
def filter(self, *_args):
return self
def all(self):
return []
class _Session:
def __init__(self):
self.query_count = 0
def query(self, *_args):
self.query_count += 1
return _Query()
session = _Session()
index = load_existing_product_match_index( # type: ignore[arg-type]
session,
1,
"code",
[f"P-{i}" for i in range(1201)],
chunk_size=500,
)
assert index == {}
assert session.query_count == 3

View file

@ -87,3 +87,79 @@ def test_prepare_opening_balance_requires_permission():
)
assert ob is None
assert any("دسترسی" in e for e in errors)
def test_prepare_opening_balance_requires_inventory_write_permission():
idx = WarehouseImportIndex([_warehouse(1, "W", "W", True)])
item = {
"item_type": "کالا",
"track_inventory": True,
OPENING_BALANCE_QUANTITY_KEY: 5,
"default_warehouse_id": 1,
}
ob, errors, warnings, preview = prepare_opening_balance_for_import_row(
item=item,
mapped_keys={OPENING_BALANCE_QUANTITY_KEY},
business_id=1,
db=None, # type: ignore[arg-type]
can_edit_opening_balance=True,
can_write_inventory=False,
is_update=False,
existing_product=None,
warehouse_index=idx,
)
assert ob is None
assert any("ویرایش موجودی" in e for e in errors)
def test_prepare_opening_balance_rejects_unique_without_instances():
idx = WarehouseImportIndex([_warehouse(1, "W", "W", True)])
item = {
"item_type": "کالا",
"track_inventory": True,
"inventory_mode": "unique",
OPENING_BALANCE_QUANTITY_KEY: 1,
"default_warehouse_id": 1,
}
ob, errors, warnings, preview = prepare_opening_balance_for_import_row(
item=item,
mapped_keys={OPENING_BALANCE_QUANTITY_KEY},
business_id=1,
db=None, # type: ignore[arg-type]
can_edit_opening_balance=True,
can_write_inventory=True,
is_update=False,
existing_product=None,
warehouse_index=idx,
)
assert ob is None
assert any("یونیک" in e for e in errors)
def test_prepare_opening_balance_can_defer_document_queries():
idx = WarehouseImportIndex([_warehouse(1, "W", "W", True)])
item = {
"item_type": "کالا",
"track_inventory": True,
OPENING_BALANCE_QUANTITY_KEY: 5,
OPENING_BALANCE_COST_KEY: 120,
"default_warehouse_id": 1,
}
ob, errors, warnings, preview = prepare_opening_balance_for_import_row(
item=item,
mapped_keys={OPENING_BALANCE_QUANTITY_KEY, OPENING_BALANCE_COST_KEY},
business_id=1,
db=None, # type: ignore[arg-type]
can_edit_opening_balance=True,
is_update=False,
existing_product=None,
warehouse_index=idx,
validate_document_context=False,
)
assert errors == []
assert warnings == []
assert ob is not None
assert ob.quantity == 5
assert ob.cost_price == 120
assert ob.warehouse_id == 1
assert preview["action"] == "upsert"

View file

@ -2,6 +2,7 @@
from app.services.product_opening_balance_service import (
_extract_product_ob_line,
_merge_product_opening_balance_changes,
_warehouse_id_from_line,
)
@ -52,3 +53,62 @@ def test_extract_product_ob_line_without_warehouse_returns_first():
line = _extract_product_ob_line(doc, 10)
assert line is not None
assert line["quantity"] == 5
def test_merge_import_changes_preserves_unrelated_lines_and_moves_warehouse():
existing = [
{
"product_id": 10,
"quantity": 3,
"extra_info": {"warehouse_id": 1, "movement": "in", "cost_price": 200},
"description": "old",
},
{
"product_id": 20,
"quantity": 9,
"extra_info": {"warehouse_id": 2, "movement": "in"},
"description": "keep",
},
]
merged = _merge_product_opening_balance_changes(
existing,
[
{
"product_id": 10,
"product_name": "A",
"previous_warehouse_id": 1,
"warehouse_id": 3,
"quantity": 7,
"cost_price": 250,
},
{
"product_id": 30,
"product_name": "B",
"warehouse_id": 4,
"quantity": 2,
"cost_price": 0,
},
],
)
by_key = {
(int(line["product_id"]), int(line["extra_info"]["warehouse_id"])): line
for line in merged
}
assert (10, 1) not in by_key
assert by_key[(10, 3)]["quantity"] == 7
assert by_key[(10, 3)]["extra_info"]["cost_price"] == 250
assert by_key[(20, 2)]["description"] == "keep"
assert by_key[(30, 4)]["quantity"] == 2
def test_merge_import_changes_last_duplicate_wins_without_growing_lines():
merged = _merge_product_opening_balance_changes(
[],
[
{"product_id": 1, "warehouse_id": 2, "quantity": 3, "cost_price": 10},
{"product_id": 1, "warehouse_id": 2, "quantity": 8, "cost_price": 12},
],
)
assert len(merged) == 1
assert merged[0]["quantity"] == 8
assert merged[0]["extra_info"]["cost_price"] == 12

View file

@ -24,6 +24,7 @@ import 'package:hesabix_ui/services/person_service.dart';
import 'package:hesabix_ui/services/product_service.dart';
import 'package:shared_preferences/shared_preferences.dart';
import 'package:hesabix_ui/utils/number_normalizer.dart';
import 'package:hesabix_ui/utils/opening_balance_totals.dart';
import '../../utils/error_extractor.dart';
import '../../utils/snackbar_helper.dart';
import '../../widgets/business_subpage_back_leading.dart';
@ -1476,7 +1477,17 @@ class _OpeningBalancePageState extends State<OpeningBalancePage> {
invValue += (q * c);
}
debit += invValue;
return {'debit': debit, 'credit': credit, 'diff': debit - credit};
final effectiveTotals = applyOpeningBalanceAutoBalance(
debit: debit,
credit: credit,
autoBalanceEnabled: _autoBalance,
hasEquityAccount: _equityAccountId != null,
);
return {
'debit': effectiveTotals.debit,
'credit': effectiveTotals.credit,
'diff': effectiveTotals.difference,
};
}
Map<String, bool> _computeValidation() {

View file

@ -0,0 +1,28 @@
class OpeningBalanceTotals {
final double debit;
final double credit;
const OpeningBalanceTotals({required this.debit, required this.credit});
double get difference => debit - credit;
}
OpeningBalanceTotals applyOpeningBalanceAutoBalance({
required double debit,
required double credit,
required bool autoBalanceEnabled,
required bool hasEquityAccount,
}) {
if (!autoBalanceEnabled || !hasEquityAccount) {
return OpeningBalanceTotals(debit: debit, credit: credit);
}
final difference = debit - credit;
if (difference > 0) {
credit += difference;
} else if (difference < 0) {
debit += -difference;
}
return OpeningBalanceTotals(debit: debit, credit: credit);
}

View file

@ -3821,6 +3821,7 @@ class _DataTableWidgetState<T> extends State<DataTableWidget<T>> {
// استفاده از Scrollbar با controller برای اطمینان از اسکرول دوطرفه
// Scrollbar می‌تواند controller را حتی قبل از attach شدن handle کند
final hasScrollPosition = _horizontalScrollController.hasClients;
const tableHorizontalMargin = 10.0;
return Scrollbar(
controller: _horizontalScrollController,
@ -3841,11 +3842,11 @@ class _DataTableWidgetState<T> extends State<DataTableWidget<T>> {
),
child: DataTable2(
columnSpacing: 0,
horizontalMargin: 10,
// محاسبه minWidth بر اساس عرض کل ستون‌ها برای جلوگیری از warning
// اگر عرض ستون‌ها بیش از availableWidth باشد، از همان استفاده می‌کنیم
// در غیر این صورت از config استفاده می‌کنیم
minWidth: _calculateMinTableWidth(columns, availableWidth),
horizontalMargin: tableHorizontalMargin,
minWidth: DataTableUtils.getFixedColumnsMinTableWidth(
columns,
horizontalMargin: tableHorizontalMargin,
),
horizontalScrollController: _horizontalScrollController,
headingRowHeight: widget.config.showColumnHeaders
? (widget.config.headingRowHeight ?? (_dense ? 34 : 36))
@ -4272,37 +4273,6 @@ class _DataTableWidgetState<T> extends State<DataTableWidget<T>> {
return computed;
}
/// محاسبه minWidth برای DataTable2 بر اساس عرض ستون‌ها
/// این جلوی warning "combined width of columns ... is greater than available parent width" را می‌گیرد
double _calculateMinTableWidth(
List<DataColumn2> columns,
double availableWidth,
) {
// محاسبه مجموع fixedWidth ستون‌ها
double totalFixedWidth = 0.0;
for (final col in columns) {
if (col.fixedWidth != null) {
totalFixedWidth += col.fixedWidth!;
}
}
// DataTable2 assertion چک می‌کند که minWidth >= totalFixedWidth + (horizontalMargin * 2)
// horizontalMargin در DataTable2 برابر با 10 است، پس باید 20 اضافه کنیم
const horizontalMarginTotal = 20.0; // horizontalMargin * 2 = 10 * 2
// minWidth باید حداقل برابر با totalFixedWidth + horizontalMarginTotal باشد تا assertion نخورد
final minRequiredWidth = totalFixedWidth + horizontalMarginTotal;
// اگر availableWidth infinity است یا نامعتبر است، فقط minRequiredWidth را برگردانیم
if (!availableWidth.isFinite || availableWidth <= 0) {
return minRequiredWidth;
}
// همیشه minRequiredWidth را برگردانیم تا از assertion error جلوگیری شود
// DataTable2 به صورت خودکار scroll می‌شود اگر فضا کافی نباشد
return minRequiredWidth;
}
/// محاسبه عرض کل ستون‌ها و تنظیم عرض‌ها برای پر کردن فضای موجود
/// این متد عرض کل ستون‌ها را محاسبه می‌کند و اگر شرایط برقرار باشد،
/// عرض ستون‌ها را به نسبت افزایش می‌دهد تا فضای موجود را پر کنند

View file

@ -5,6 +5,26 @@ import '../data_table_config.dart';
/// Utility functions for data table
class DataTableUtils {
/// Calculates a safe [DataTable2.minWidth] when every rendered column uses
/// [DataColumn2.fixedWidth].
///
/// DataTable2 removes the horizontal margins before asserting that the
/// available width is *strictly* greater than the sum of fixed widths. The
/// extra logical pixel keeps the calculation on the valid side of that
/// strict inequality.
static double getFixedColumnsMinTableWidth(
Iterable<DataColumn2> columns, {
required double horizontalMargin,
}) {
var totalFixedWidth = 0.0;
for (final column in columns) {
totalFixedWidth += column.fixedWidth ?? 0.0;
}
const strictInequalitySlack = 1.0;
return totalFixedWidth + (horizontalMargin * 2) + strictInequalitySlack;
}
/// Format text with ellipsis if needed
static String formatText(String text, {int? maxLength}) {
if (maxLength != null && text.length > maxLength) {

View file

@ -144,12 +144,31 @@ class _ProductImportDialogState extends State<ProductImportDialog> {
final res = await api.post<Map<String, dynamic>>(
'/products/business/${widget.businessId}/import/excel',
data: form,
options: Options(contentType: 'multipart/form-data'),
options: Options(
contentType: 'multipart/form-data',
receiveTimeout: const Duration(minutes: 5),
),
);
setState(() {
_result = res.data;
});
if (!dryRun) {
final responseData = res.data?['data'];
final warehouseSync = responseData is Map
? responseData['warehouse_sync'] as Map?
: null;
final receipts = warehouseSync?['receipts_created'] ?? 0;
final issues = warehouseSync?['issues_created'] ?? 0;
final lines = warehouseSync?['lines_created'] ?? 0;
if (mounted && (receipts != 0 || issues != 0 || lines != 0)) {
final isFa = Localizations.localeOf(context).languageCode == 'fa';
SnackBarHelper.show(
context,
message: isFa
? 'ایمپورت انجام شد؛ $receipts رسید و $issues حواله خروج با $lines خط انبار قطعی شد.'
: 'Import completed; $receipts receipt(s) and $issues issue(s) with $lines warehouse line(s) were posted.',
);
}
if (mounted) Navigator.of(context).pop(true);
}
} catch (e) {
@ -385,6 +404,7 @@ class _ResultSummaryBodyState extends State<_ResultSummaryBody> {
final summary = (data?['summary'] as Map<String, dynamic>?) ?? {};
final errors = (data?['errors'] as List?)?.cast<Map<String, dynamic>>() ?? const [];
final refSummary = data?['reference_summary'] as Map<String, dynamic>?;
final warehouseSync = data?['warehouse_sync'] as Map<String, dynamic>?;
final previewRaw = (data?['preview'] as List?)?.cast<dynamic>() ?? const [];
final preview = previewRaw
.whereType<Map>()
@ -435,6 +455,24 @@ class _ResultSummaryBodyState extends State<_ResultSummaryBody> {
],
),
],
if (warehouseSync != null) ...[
const SizedBox(height: 8),
Wrap(
spacing: 8,
runSpacing: 6,
children: [
if (summary['dry_run'] == true) ...[
_chip(isFa ? 'انبار کاندید' : 'Candidate warehouses', warehouseSync['candidate_warehouses']),
_chip(isFa ? 'خط انبار کاندید' : 'Candidate warehouse lines', warehouseSync['candidate_lines']),
] else ...[
_chip(isFa ? 'رسید انبار' : 'Warehouse receipts', warehouseSync['receipts_created']),
_chip(isFa ? 'حواله خروج' : 'Warehouse issues', warehouseSync['issues_created']),
_chip(isFa ? 'خط انبار ایجادشده' : 'Warehouse lines created', warehouseSync['lines_created']),
_chip(isFa ? 'خط بدون تغییر' : 'Unchanged warehouse lines', warehouseSync['unchanged_lines']),
],
],
),
],
if (preview.isNotEmpty) ...[
const SizedBox(height: 8),
ExpansionTile(

View file

@ -0,0 +1,55 @@
import 'package:data_table_2/data_table_2.dart';
import 'package:flutter/material.dart';
import 'package:flutter_test/flutter_test.dart';
import 'package:hesabix_ui/widgets/data_table/helpers/data_table_utils.dart';
void main() {
test('fixed column table width leaves room for DataTable2 strict check', () {
final columns = <DataColumn2>[
const DataColumn2(label: SizedBox.shrink(), fixedWidth: 50),
const DataColumn2(label: SizedBox.shrink(), fixedWidth: 60),
const DataColumn2(label: SizedBox.shrink(), fixedWidth: 96),
const DataColumn2(label: SizedBox.shrink()),
];
final minWidth = DataTableUtils.getFixedColumnsMinTableWidth(
columns,
horizontalMargin: 10,
);
expect(minWidth, 227);
expect(minWidth - 20, greaterThan(206));
});
testWidgets('fixed width columns lay out inside a narrow parent', (
tester,
) async {
final columns = <DataColumn2>[
const DataColumn2(label: Text('A'), fixedWidth: 100),
const DataColumn2(label: Text('B'), fixedWidth: 150),
const DataColumn2(label: Text('C'), fixedWidth: 200),
];
await tester.pumpWidget(
MaterialApp(
home: Scaffold(
body: SizedBox(
width: 200,
height: 200,
child: DataTable2(
horizontalMargin: 10,
minWidth: DataTableUtils.getFixedColumnsMinTableWidth(
columns,
horizontalMargin: 10,
),
columns: columns,
rows: const [],
),
),
),
),
);
expect(tester.takeException(), isNull);
});
}

View file

@ -0,0 +1,56 @@
import 'package:flutter_test/flutter_test.dart';
import 'package:hesabix_ui/utils/opening_balance_totals.dart';
void main() {
test('keeps raw difference when auto balance is disabled', () {
final totals = applyOpeningBalanceAutoBalance(
debit: 100,
credit: 0,
autoBalanceEnabled: false,
hasEquityAccount: true,
);
expect(totals.debit, 100);
expect(totals.credit, 0);
expect(totals.difference, 100);
});
test('keeps raw difference when equity account is missing', () {
final totals = applyOpeningBalanceAutoBalance(
debit: 100,
credit: 0,
autoBalanceEnabled: true,
hasEquityAccount: false,
);
expect(totals.debit, 100);
expect(totals.credit, 0);
expect(totals.difference, 100);
});
test('adds debit-heavy difference to displayed credit', () {
final totals = applyOpeningBalanceAutoBalance(
debit: 8528542630.50,
credit: 0,
autoBalanceEnabled: true,
hasEquityAccount: true,
);
expect(totals.debit, 8528542630.50);
expect(totals.credit, 8528542630.50);
expect(totals.difference, 0);
});
test('adds credit-heavy difference to displayed debit', () {
final totals = applyOpeningBalanceAutoBalance(
debit: 20,
credit: 80,
autoBalanceEnabled: true,
hasEquityAccount: true,
);
expect(totals.debit, 80);
expect(totals.credit, 80);
expect(totals.difference, 0);
});
}