Add full sync retirement for missing system default tasks

This commit is contained in:
A R R R Associates
2026-09-19 14:13:16 +05:30
parent 2d6043962b
commit 7bd7f52278
6 changed files with 195 additions and 33 deletions
+173 -18
View File
@@ -79,6 +79,7 @@ SERVICE_MASTER_COLUMNS = [
]
DEFAULT_TASK_COLUMNS = [
"system_task_id",
"service_code",
"sequence_no",
"task_name",
@@ -431,8 +432,8 @@ def build_template(template_type: str) -> bytes:
elif template_type == "system_default_tasks":
ws.title = "system_default_tasks"
ws.append(DEFAULT_TASK_COLUMNS)
ws.append(["GST-MONTHLY", 1, "Collect data", "Staff", "TRUE", "FALSE", "TRUE", "Collect sales/purchase data"])
ws.append(["GST-MONTHLY", 2, "Review and file", "Manager", "TRUE", "TRUE", "TRUE", "Review and file return"])
ws.append(["", "GST-MONTHLY", 1, "Collect data", "Staff", "TRUE", "FALSE", "", "", "FALSE", "NONE", "FALSE", "FALSE", "NONE", "FALSE", "FALSE", "FALSE", "FALSE", "FALSE", "FALSE", "FALSE", "", "TRUE", "Collect sales/purchase data"])
ws.append(["", "GST-MONTHLY", 2, "Review and file", "Manager", "TRUE", "TRUE", "", "", "FALSE", "NONE", "FALSE", "FALSE", "NONE", "FALSE", "FALSE", "FALSE", "FALSE", "FALSE", "FALSE", "FALSE", "", "TRUE", "Review and file return"])
elif template_type == "firm_task_templates":
ws.title = "firm_task_templates"
@@ -1059,14 +1060,51 @@ def _resolve_import_task_category(db: Session, *, catalogue_id: int, tenant_id:
return ensure_task_category(db, catalogue_id=catalogue_id, tenant_id=tenant_id, name=name or None, user_id=user_id)
def import_system_default_tasks(db: Session, *, current_user, file_bytes: bytes, update_existing: bool = True, expected_service_code: str | None = None) -> dict:
def import_system_default_tasks(
db: Session,
*,
current_user,
file_bytes: bytes,
update_existing: bool = True,
expected_service_code: str | None = None,
full_sync: bool = False,
) -> dict:
"""Import centrally maintained default tasks.
Normal mode preserves the historical behaviour: existing rows are matched by
service + sequence (unless a stable system_task_id is supplied).
Full Synchronization is intentionally explicit. For each service represented in
the workbook it treats the workbook as the complete desired active task list:
* existing rows are matched by stable system_task_id first;
* older workbooks without IDs fall back to exact task-name matching;
* unmatched workbook rows are created;
* existing tasks omitted from the workbook are soft-retired (is_active=False);
* existing rows are temporarily resequenced so a 120 -> 78 renumber does not
collide with the database unique constraint;
* history is preserved because system task rows are never physically deleted.
Retired system defaults are then rolled out through the existing inheritance
engine: inherited firm tasks follow automatically, customized firm tasks receive
the existing Firm Admin update decision, and only safe pending/unstarted tasks in
open engagements are changed.
"""
ws, headers = _load_sheet(file_bytes, "system_default_tasks")
missing = _validate_headers(headers, ["service_code", "sequence_no", "task_name"])
if missing:
return {"created": 0, "updated": 0, "skipped": 0, "errors": [{"row": 1, "message": f"Missing columns: {', '.join(missing)}"}]}
created = updated = skipped = 0
return {
"created": 0, "updated": 0, "retired": 0, "skipped": 0,
"errors": [{"row": 1, "message": f"Missing columns: {', '.join(missing)}"}],
}
created = updated = retired = skipped = 0
errors: list[dict] = []
touched_catalogue_ids: set[int] = set()
parsed_rows: list[dict[str, Any]] = []
seen_sequences: set[tuple[int, int]] = set()
seen_ids: set[int] = set()
# Phase 1: validate/parse the workbook without changing task rows.
for row_no, row in enumerate(ws.iter_rows(min_row=2, values_only=True), start=2):
if not any(v not in (None, "") for v in row):
continue
@@ -1076,35 +1114,136 @@ def import_system_default_tasks(db: Session, *, current_user, file_bytes: bytes,
raise ValueError(f"This import accepts only service code {normalize_code(expected_service_code)}.")
sequence_no = _int(_cell(row, headers, "sequence_no"), None)
task_name = _clean(_cell(row, headers, "task_name"))
system_task_id = _int(_cell(row, headers, "system_task_id"), None)
catalogue = find_service(db, service_code=service_code)
if not catalogue:
raise ValueError("Service code not found in service catalogue.")
if not sequence_no or sequence_no <= 0 or not task_name:
raise ValueError("sequence_no must be a positive number and task_name is required.")
seq_key = (int(catalogue.id), int(sequence_no))
if seq_key in seen_sequences:
raise ValueError("Duplicate sequence_no for this service in the workbook.")
seen_sequences.add(seq_key)
if system_task_id:
if system_task_id in seen_ids:
raise ValueError("The same system_task_id appears more than once in the workbook.")
seen_ids.add(system_task_id)
existing_by_id = db.get(ServiceDefaultTaskTemplate, system_task_id)
if existing_by_id is None:
raise ValueError(f"system_task_id {system_task_id} does not exist. Leave the ID blank for a new task.")
if int(existing_by_id.service_catalogue_id) != int(catalogue.id):
raise ValueError(f"system_task_id {system_task_id} belongs to a different service.")
parsed_rows.append({
"row_no": row_no,
"row": row,
"catalogue": catalogue,
"service_code": service_code,
"sequence_no": int(sequence_no),
"task_name": task_name,
"system_task_id": system_task_id,
})
touched_catalogue_ids.add(int(catalogue.id))
if not sequence_no or not task_name:
raise ValueError("sequence_no and task_name are required.")
task = db.execute(
select(ServiceDefaultTaskTemplate).where(
ServiceDefaultTaskTemplate.service_catalogue_id == catalogue.id,
ServiceDefaultTaskTemplate.sequence_no == sequence_no,
)
).scalar_one_or_none()
except Exception as exc:
errors.append({"row": row_no, "message": str(exc)})
if errors:
db.rollback()
return {"created": 0, "updated": 0, "retired": 0, "skipped": 0, "errors": errors}
# Cache existing rows before applying changes. Full sync needs this original set
# to know which rows were genuinely omitted from the workbook.
existing_by_catalogue: dict[int, list[ServiceDefaultTaskTemplate]] = {}
for catalogue_id in sorted(touched_catalogue_ids):
existing_by_catalogue[catalogue_id] = list(db.execute(
select(ServiceDefaultTaskTemplate)
.where(ServiceDefaultTaskTemplate.service_catalogue_id == catalogue_id)
.order_by(ServiceDefaultTaskTemplate.sequence_no.asc(), ServiceDefaultTaskTemplate.id.asc())
).scalars().all())
# Full-sync re-numbering can move task 90 to sequence 40 while old task 40 still
# exists. Move the original rows to a high temporary range first to avoid the
# service+sequence unique constraint during the transaction.
if full_sync:
for catalogue_id, existing_rows in existing_by_catalogue.items():
max_seq = max([int(t.sequence_no or 0) for t in existing_rows] or [0])
temporary_base = max(max_seq + 100000, 100000)
for offset, task in enumerate(existing_rows, start=1):
task.sequence_no = temporary_base + offset
db.flush()
matched_existing_ids: dict[int, set[int]] = {cid: set() for cid in touched_catalogue_ids}
# Phase 2: upsert workbook rows.
for item in parsed_rows:
row_no = item["row_no"]
row = item["row"]
catalogue = item["catalogue"]
catalogue_id = int(catalogue.id)
sequence_no = item["sequence_no"]
task_name = item["task_name"]
system_task_id = item["system_task_id"]
try:
task = None
if system_task_id:
task = db.get(ServiceDefaultTaskTemplate, int(system_task_id))
elif full_sync:
# Backward-compatible path for pre-ID exports (including the current
# reduced Tax Audit workbook): exact task name is safer than sequence
# because the whole list may have been resequenced.
normalised = task_name.strip().casefold()
name_matches = [
t for t in existing_by_catalogue.get(catalogue_id, [])
if (t.task_name or "").strip().casefold() == normalised
and int(t.id) not in matched_existing_ids[catalogue_id]
]
if len(name_matches) == 1:
task = name_matches[0]
else:
task = db.execute(
select(ServiceDefaultTaskTemplate).where(
ServiceDefaultTaskTemplate.service_catalogue_id == catalogue_id,
ServiceDefaultTaskTemplate.sequence_no == sequence_no,
)
).scalar_one_or_none()
if task and not update_existing:
matched_existing_ids[catalogue_id].add(int(task.id))
skipped += 1
continue
if task:
matched_existing_ids[catalogue_id].add(int(task.id))
updated += 1
else:
task = ServiceDefaultTaskTemplate(service_catalogue_id=catalogue.id, sequence_no=sequence_no, task_name=task_name)
task = ServiceDefaultTaskTemplate(
service_catalogue_id=catalogue_id,
sequence_no=sequence_no,
task_name=task_name,
)
db.add(task)
db.flush()
created += 1
task.sequence_no = sequence_no
task.task_name = task_name
task.description = _clean(_cell(row, headers, "description")) or None
task.default_role_name = _clean(_cell(row, headers, "default_role_name")) or None
if "eligible_role_names" in headers:
task.eligible_role_names = _clean(_cell(row, headers, "eligible_role_names")) or None
task.is_mandatory = _bool(_cell(row, headers, "is_mandatory"), True)
task.requires_review = _bool(_cell(row, headers, "requires_review"), False)
if "normal_review_role" in headers:
normal_review_role = (_clean(_cell(row, headers, "normal_review_role")) or "").lower()
if normal_review_role and normal_review_role not in {"manager", "partner", "manager_or_partner"}:
raise ValueError("normal_review_role must be manager, partner or manager_or_partner.")
task.normal_review_role = normal_review_role or None
category = _resolve_import_task_category(
db,
catalogue_id=catalogue.id,
catalogue_id=catalogue_id,
tenant_id=None,
category_code=_clean(_cell(row, headers, "task_category_code")),
category_name=_clean(_cell(row, headers, "task_category")),
@@ -1130,6 +1269,20 @@ def import_system_default_tasks(db: Session, *, current_user, file_bytes: bytes,
task.is_active = _bool(_cell(row, headers, "is_active"), True)
except Exception as exc:
errors.append({"row": row_no, "message": str(exc)})
# Full Synchronization: anything that existed before this import but was not
# represented in the workbook is retired, never hard-deleted. Keep its temporary
# high sequence so active workbook rows own the concise 1..N sequence safely.
if full_sync and not errors:
for catalogue_id, existing_rows in existing_by_catalogue.items():
matched = matched_existing_ids.get(catalogue_id, set())
for task in existing_rows:
if int(task.id) in matched:
continue
if task.is_active:
task.is_active = False
retired += 1
rollout_summary = {
"firms_processed": 0,
"firms_changed": 0,
@@ -1137,11 +1290,9 @@ def import_system_default_tasks(db: Session, *, current_user, file_bytes: bytes,
"engagement_created": 0,
"engagement_updated_pending": 0,
"engagement_deactivated_pending": 0,
"engagement_preserved_history": 0,
}
if not errors:
# Apply the same transaction to every enabled firm and safely refresh only
# unstarted tasks in open/unlocked engagements. Customized firm tasks are
# never overwritten; they receive an update-available decision instead.
from app.modules.services.default_task_sync import sync_system_defaults_to_all_firms
for catalogue_id in sorted(touched_catalogue_ids):
rollout = sync_system_defaults_to_all_firms(
@@ -1156,13 +1307,17 @@ def import_system_default_tasks(db: Session, *, current_user, file_bytes: bytes,
rollout_summary["engagement_created"] += rollout.engagement_created
rollout_summary["engagement_updated_pending"] += rollout.engagement_updated_pending
rollout_summary["engagement_deactivated_pending"] += rollout.engagement_deactivated_pending
rollout_summary["engagement_preserved_history"] += rollout.engagement_preserved_history
db.commit()
else:
db.rollback()
return {
"created": created if not errors else 0,
"updated": updated if not errors else 0,
"retired": retired if not errors else 0,
"skipped": skipped,
"full_sync": bool(full_sync),
"errors": errors,
**(rollout_summary if not errors else {}),
}