Reconcile firm tasks fully with system defaults

This commit is contained in:
A R R R Associates
2026-09-19 16:02:07 +05:30
parent 5cee4f8de4
commit 53732a7270
3 changed files with 259 additions and 56 deletions
+234
View File
@@ -56,6 +56,7 @@ class FirmDefaultTaskSyncResult:
created: int = 0
updated: int = 0
duplicates_disabled: int = 0
retired: int = 0
unchanged: int = 0
custom_updates_available: int = 0
linked_legacy: int = 0
@@ -515,6 +516,239 @@ def sync_firm_tasks_from_system_defaults(
return result
def fully_align_firm_tasks_to_system_defaults(
db: Session,
*,
tenant_id: int,
service_catalogue_id: int,
updated_by_user_id: int,
sync_open_engagements: bool = True,
) -> FirmDefaultTaskSyncResult:
"""Fully align one firm's system-derived checklist to the current system master.
This is the bulk/explicit "Upgrade All to System" operation. It is intentionally
stronger than the normal background sync:
* every row already linked to a system task follows that system task, even when an
earlier firm customization exists;
* linked rows whose system source is retired are retired at firm level;
* legacy unlinked rows are linked only by an exact normalized task-name match;
* duplicate rows for the same system source are retired;
* active system-derived rows receive the exact current system sequence;
* genuine firm-only rows (no reliable system match) are preserved. If one occupies
a sequence now required by the system master it is moved to the next free sequence;
* engagement instances are synchronized once after the template reconciliation,
using the existing safe rollout rules so historical/started/reviewed work is not
rewritten.
No task/template row is hard-deleted.
"""
defaults = db.execute(
select(ServiceDefaultTaskTemplate)
.where(ServiceDefaultTaskTemplate.service_catalogue_id == service_catalogue_id)
.order_by(ServiceDefaultTaskTemplate.sequence_no.asc(), ServiceDefaultTaskTemplate.id.asc())
).scalars().all()
firm_rows = db.execute(
select(FirmServiceTaskTemplate)
.where(
FirmServiceTaskTemplate.tenant_id == tenant_id,
FirmServiceTaskTemplate.service_catalogue_id == service_catalogue_id,
)
.order_by(FirmServiceTaskTemplate.sequence_no.asc(), FirmServiceTaskTemplate.id.asc())
).scalars().all()
result = FirmDefaultTaskSyncResult()
reference_counts = _reference_counts(db, firm_task_ids=[int(r.id) for r in firm_rows if r.id])
defaults_by_id = {int(row.id): row for row in defaults}
defaults_by_name: dict[str, list[ServiceDefaultTaskTemplate]] = {}
for source in defaults:
key = _normalise_name(source.task_name)
if key:
defaults_by_name.setdefault(key, []).append(source)
# Capture original firm ordering before temporarily moving rows away from the
# unique (tenant, service, sequence) namespace.
original_sequence = {int(row.id): int(row.sequence_no or 0) for row in firm_rows if row.id is not None}
# Safely recover provenance for legacy rows. Exact name is deliberately the only
# automatic fallback here: anything ambiguous remains a genuine firm-only row.
for row in firm_rows:
if row.source_system_task_id is not None:
continue
matches = defaults_by_name.get(_normalise_name(row.task_name), [])
if len(matches) == 1:
row.source_system_task_id = matches[0].id
result.linked_legacy += 1
db.flush()
by_source: dict[int, list[FirmServiceTaskTemplate]] = {}
firm_only: list[FirmServiceTaskTemplate] = []
for row in firm_rows:
source_id = int(row.source_system_task_id) if row.source_system_task_id is not None else None
if source_id is not None and source_id in defaults_by_id:
by_source.setdefault(source_id, []).append(row)
else:
# A missing/deleted source becomes firm-only because the FK uses SET NULL;
# never guess and retire it without a reliable system identity.
if source_id is not None and source_id not in defaults_by_id:
row.source_system_task_id = None
firm_only.append(row)
# Free every sequence first. This allows exact system sequencing even when an old
# firm row currently occupies the target number. The unique constraint remains valid
# because each temporary sequence is unique.
highest_existing = max([int(r.sequence_no or 0) for r in firm_rows] + [0])
highest_default = max([int(r.sequence_no or 0) for r in defaults] + [0])
temporary_base = max(highest_existing, highest_default, 100000) + 1000000
for offset, row in enumerate(sorted(firm_rows, key=lambda r: int(r.id or 0)), start=1):
row.sequence_no = temporary_base + offset
db.flush()
canonical_by_source: dict[int, FirmServiceTaskTemplate] = {}
for source in defaults:
source_id = int(source.id)
latest_hash = system_task_hash(source)
candidates = by_source.get(source_id, [])
if candidates:
# Prefer the currently active row, then the row referenced by the most
# engagement history, then the oldest id. This avoids switching canonical
# identity unnecessarily while still preserving historical references.
target = sorted(
candidates,
key=lambda row: (
-int(bool(row.is_active)),
-reference_counts.get(int(row.id), 0),
int(row.id),
),
)[0]
canonical_by_source[source_id] = target
for duplicate in candidates:
if duplicate.id == target.id:
continue
if duplicate.is_active:
result.retired += 1
duplicate.is_active = False
duplicate.is_customized = False
duplicate.last_synced_system_hash = latest_hash
duplicate.last_reviewed_system_hash = latest_hash
duplicate.system_update_available = False
duplicate.system_update_detected_at_utc = None
duplicate.updated_by_user_id = updated_by_user_id
result.duplicates_disabled += 1
elif source.is_active:
# Create only active system defaults. Retired defaults with no firm row do
# not need a new historical firm record.
target = FirmServiceTaskTemplate(
tenant_id=tenant_id,
service_catalogue_id=service_catalogue_id,
task_name=source.task_name,
sequence_no=temporary_base + len(firm_rows) + len(canonical_by_source) + 1,
source_system_task_id=source.id,
is_customized=False,
last_synced_system_hash=latest_hash,
last_reviewed_system_hash=latest_hash,
system_update_available=False,
created_by_user_id=updated_by_user_id,
updated_by_user_id=updated_by_user_id,
)
db.add(target)
db.flush()
firm_rows.append(target)
by_source.setdefault(source_id, []).append(target)
canonical_by_source[source_id] = target
result.created += 1
else:
continue
target = canonical_by_source[source_id]
was_active = bool(target.is_active)
changed = _copy_default_columns(
db,
source,
target,
tenant_id=tenant_id,
user_id=updated_by_user_id,
allow_sequence_change=False,
)
target.source_system_task_id = source.id
target.is_customized = False
target.last_synced_system_hash = latest_hash
target.last_reviewed_system_hash = latest_hash
target.system_update_available = False
target.system_update_detected_at_utc = None
target.updated_by_user_id = updated_by_user_id
if was_active and not bool(source.is_active):
result.retired += 1
if changed:
result.updated += 1
else:
result.unchanged += 1
db.flush()
# Active system rows now own the exact master sequence numbers.
used_sequences: set[int] = set()
active_sources = [row for row in defaults if bool(row.is_active)]
for source in sorted(active_sources, key=lambda r: (int(r.sequence_no or 0), int(r.id))):
target = canonical_by_source.get(int(source.id))
if target is None:
continue
desired = int(source.sequence_no or 0)
target.sequence_no = desired
used_sequences.add(desired)
db.flush()
# Genuine firm-only active tasks survive the full upgrade. Keep their original
# sequence when it is still free; otherwise move them after the system checklist.
next_free = max(used_sequences or {0}) + 1
active_firm_only = sorted(
[row for row in firm_only if bool(row.is_active)],
key=lambda row: (original_sequence.get(int(row.id), 0), int(row.id)),
)
for row in active_firm_only:
preferred = original_sequence.get(int(row.id), 0)
if preferred > 0 and preferred not in used_sequences:
row.sequence_no = preferred
used_sequences.add(preferred)
else:
while next_free in used_sequences:
next_free += 1
row.sequence_no = next_free
used_sequences.add(next_free)
next_free += 1
# Retired/duplicate templates are retained for history but moved outside the live
# sequence range so they can never block exact system sequencing.
inactive_rows = [row for row in firm_rows if not bool(row.is_active)]
retired_base = max(used_sequences or {0}) + 100000
for offset, row in enumerate(sorted(inactive_rows, key=lambda r: int(r.id or 0)), start=1):
row.sequence_no = retired_base + offset
db.flush()
if sync_open_engagements:
from app.modules.services.execution import sync_open_engagement_tasks_for_service
engagement = sync_open_engagement_tasks_for_service(
db,
tenant_id=tenant_id,
catalogue_id=service_catalogue_id,
user_id=updated_by_user_id,
include_started_open_tasks=False,
safe_system_rollout=True,
)
result.engagement_created = engagement.get("created", 0)
result.engagement_updated_pending = engagement.get("updated_pending", 0)
result.engagement_deactivated_pending = engagement.get("deactivated_pending", 0)
result.engagement_preserved_history = engagement.get("preserved_history", 0)
return result
def sync_system_defaults_to_all_firms(
db: Session,
*,