import datetime
import time
import json
import os
import re
from json import JSONDecodeError

import django
import sys
from django.utils import timezone
from confluent_kafka import Consumer
from django.db.models import Prefetch, Q
from rest_framework import status
from django.db import transaction, InterfaceError, OperationalError

# if want access django, must put this 2
os.environ.setdefault("DJANGO_SETTINGS_MODULE", "backend.settings")
sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
django.setup()

from api.models import KafkaLog
from django.conf import settings
from api.utils import log_critical_error, get_exception_detail

from api.translation import __

from api.models import (
    Setting,
    TtAcademicTerm,
    TtActivity,
    TtActivityTemplate,
    TtActivityType,
    TtDepartment,
    TtLocation,
    TtModule,
    TtModuleGroup,
    TtPos,
    TtStaff,
    TtStudent,
    TtStudentSet,
    TtWeek,
    TtWeekPattern,
)
from api.models.relations import (
    TtAcademicTermWeek,
    TtActivityLocation,
    TtActivityStaff,
    TtActivityTemplateWeek,
    TtActivityWeek,
    TtModuleWeek,
    TtPosModuleGroup,
    TtPosModuleGroupModule,
    TtStaffSharedWithDepartment,
    TtLocationSharedWithDepartment,
    TtStudentSetActivity,
    TtStudentSetModule,
    TtStudentAcademicTerm,
    TtStudentPos,
    TtStudentPosModule,
    TtStudentSetStudent,
)

host = settings.KAFKA_CONFIG["HOST"]
port = str(settings.KAFKA_CONFIG["PORT"])
worker_1 = settings.KAFKA_CONFIG["WORKER_1"]
conf = {
    'bootstrap.servers': host + ":" + port,
    'group.id': worker_1, # 1 group id only can poll 1 time of same msg data, so if this group_id already poll request A 1 time, next time restart if using earliest also wont get again request A
    'auto.offset.reset': 'latest', #'earliest': consume from beginning, 'latest': consume only new messages, 'none': error if no offset, mostly in manual poll scenario will used
    'enable.auto.commit': False # make auto commit to false, if got problem of connection/timeout, wont commit it and after sleep will try again, so wont have problem of skip record
}

# websocket_host = settings.WEBSOCKET_CONFIG["HOST"]
# websocket_port = str(settings.WEBSOCKET_CONFIG["PORT"])
# websocket_schedule_endpoint = str(settings.WEBSOCKET_CONFIG["SCHEDULE_RESPONSE_ENDPOINT"])
# websocket_preschedule_endpoint = str(settings.WEBSOCKET_CONFIG["PRESCHEDULE_RESPONSE_ENDPOINT"])
# websocket_preschedule_details_endpoint = str(settings.WEBSOCKET_CONFIG["PRESCHEDULE_DETAILS_RESPONSE_ENDPOINT"])

# websocket_schedule_url = websocket_host + ":" + websocket_port + "/" + websocket_schedule_endpoint
# websocket_preschedule_url = websocket_host + ":" + websocket_port + "/" + websocket_preschedule_endpoint
# websocket_preschedule_details_url = websocket_host + ":" + websocket_port + "/" + websocket_preschedule_details_endpoint

consumer = Consumer(conf)
topic = settings.KAFKA_CONFIG["MICROSERVICES_TT_TOPIC"]
timeout = settings.KAFKA_CONFIG["TIMEOUT"]
# topic = "test-topic"
consumer.subscribe([topic])

def append_log(response_data, header_data, data):
    request_id = response_data.get("request_id")

    kafka_log = KafkaLog.insert_response(
        request_id,
        topic,
        response_data,
        None,
        header_data
    )

    try:
        extra_data = json.loads(kafka_log.extra_data) if kafka_log.extra_data else []
        if not isinstance(extra_data, list):
            extra_data = []
    except Exception:
        extra_data = []

    extra_data.append(data)

    KafkaLog.objects.filter(id=kafka_log.id).update(
        extra_data=json.dumps(extra_data, default=str)
    )

    return kafka_log

def normalize_payload(response_data, key):
    payload = response_data.get(key)
    
    if not payload:
        return []
    if isinstance(payload, dict):
        return [payload]
    if isinstance(payload, list):
        return payload
    raise Exception(f"{key} must be dict or list")

def get_existing_ids(Model, rows, id_key="id"):
    incoming_ids = [row.get(id_key) for row in rows if row.get(id_key) is not None]
    if not incoming_ids:
        return set()
    return set(
        Model.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
    )

def bulk_create_skip_duplicates(Model, data_rows, batch_size=500):
    if not data_rows:
        return []

    existing_ids = get_existing_ids(Model, data_rows, "id")
    to_create = [Model(**row) for row in data_rows if row.get("id") not in existing_ids]

    if not to_create:
        return []

    Model.objects.bulk_create(to_create, batch_size=batch_size,ignore_conflicts=True)
    return to_create

def build_update_data(row, field_map, single_item=False):
    data = {}

    for payload_key, model_field in field_map.items():
        if payload_key not in row:
            continue

        if payload_key == "code" and not single_item:
            continue

        data[model_field] = row.get(payload_key)

    return data

def replace_relation_rows(RelModel, parent_field, child_field, relation_map):
    parent_ids = list(relation_map.keys())

    if not parent_ids:
        return 0

    RelModel.objects.filter(
        **{f"{parent_field}_id__in": parent_ids}
    ).delete()

    new_rows = []
    for parent_id, child_ids in relation_map.items():
        for child_id in child_ids or []:
            if child_id:
                new_rows.append(
                    RelModel(
                        **{
                            f"{parent_field}_id": parent_id,
                            f"{child_field}_id": child_id,
                        }
                    )
                )

    if new_rows:
        RelModel.objects.bulk_create(
            new_rows,
            batch_size=500,
            ignore_conflicts=True,
        )

    return len(new_rows)

def get_valid_week_pattern_ids(rows, key="week_pattern"):
    week_pattern_ids = {
        row.get(key)
        for row in rows
        if row.get(key)
    }

    if not week_pattern_ids:
        return set()

    return set(
        TtWeekPattern.objects
        .filter(id__in=week_pattern_ids)
        .values_list("id", flat=True)
    )

def get_valid_ids(Model, ids):
    ids = list(dict.fromkeys(ids or []))

    if not ids:
        return []

    return list(
        Model.objects.filter(id__in=ids)
        .values_list("id", flat=True)
    )

def should_use_activity_week(row, valid_week_pattern_ids, key="week_pattern"):
    week_pattern_id = row.get(key)

    # week_pattern priority > direct activity week
    if week_pattern_id and week_pattern_id in valid_week_pattern_ids:
        return False

    return True

def bulk_delete_by_ids(Model, ids):
    if not ids:
        return 0, 0

    existing_ids = set(
        Model.objects.filter(id__in=ids).values_list("id", flat=True)
    )

    if not existing_ids:
        return 0, len(ids)

    deleted_count, _ = Model.objects.filter(id__in=existing_ids).delete()
    skipped_count = len(ids) - len(existing_ids)

    return deleted_count, skipped_count

# ============================
# CREATE FUNCTIONS
# ============================

def department_create(response_data):
    try:
        payload = normalize_payload(response_data, "department")

        cleaned_rows = []
        for row in payload:
            if not row.get("id"):
                continue

            cleaned_rows.append({
                "id": row.get("id"),
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "status": row.get("status", 1),
                "department_id": row.get("department"),
            })

        with transaction.atomic():
            created = bulk_create_skip_duplicates(TtDepartment, cleaned_rows)

        return {
            "step": "department_create",
            "created": len(created),
            "skipped": len(cleaned_rows) - len(created),
        }

    except Exception:
        raise

def location_create(response_data):
    try:
        payload = normalize_payload(response_data, "location")

        cleaned_rows = []
        shared_map = {}

        for row in payload:
            location_id = row.get("id")
            if not location_id:
                continue

            cleaned_rows.append({
                "id": location_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "capacity": row.get("capacity"),
                "area": row.get("area"),
                "maximum_period": row.get("maximum_period"),
                "contract_period": row.get("contract_period"),
                "shared_with_all_department": row.get("shared_with_all_department", False),
                "is_online": row.get("is_online", False),
                "status": row.get("status", 1),
                "department_id": row.get("department"),
            })

            shared_map[location_id] = row.get("shared_with_department", []) or []

        with transaction.atomic():
            created_locations = bulk_create_skip_duplicates(TtLocation, cleaned_rows)
            valid_ids = {row["id"] for row in cleaned_rows}

            relation_candidates = []
            for location_id, dept_ids in shared_map.items():
                if location_id not in valid_ids:
                    continue

                for dept_id in dept_ids:
                    relation_candidates.append({
                        "location_id": location_id,
                        "department_id": dept_id,
                    })

            existing_relations = set(
                TtLocationSharedWithDepartment.objects.filter(
                    location_id__in=[r["location_id"] for r in relation_candidates]
                ).values_list("location_id", "department_id")
            )

            to_create_rel = [
                TtLocationSharedWithDepartment(**r)
                for r in relation_candidates
                if (r["location_id"], r["department_id"]) not in existing_relations
            ]

            if to_create_rel:
                TtLocationSharedWithDepartment.objects.bulk_create(
                    to_create_rel,
                    batch_size=500
                )

        return {
            "step": "location_create",
            "created": len(created_locations),
            "skipped": len(cleaned_rows) - len(created_locations),
            "shared_department_rel_created": len(to_create_rel) if relation_candidates else 0,
        }

    except Exception:
        raise

def staff_create(response_data):
    try:
        payload = normalize_payload(response_data, "staff")

        cleaned_rows = []
        shared_map = {}

        for row in payload:
            staff_id = row.get("id")
            if not staff_id:
                continue

            cleaned_rows.append({
                "id": staff_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "email": row.get("email"),
                "desc": row.get("desc"),
                "department_id": row.get("department"),
                "is_part_time": row.get("is_part_time", False),
                "maximum_period": row.get("maximum_period"),
                "contract_period": row.get("contract_period"),
                "shared_with_all_department": row.get("shared_with_all_department", False),
                "status": row.get("status", 1),
            })

            shared_map[staff_id] = row.get("shared_with_department", []) or []

        with transaction.atomic():
            created_staff = bulk_create_skip_duplicates(TtStaff, cleaned_rows)
            created_ids = {obj.id for obj in created_staff}

            relation_candidates = []
            for staff_id, dept_ids in shared_map.items():
                if staff_id not in created_ids:
                    continue
                for dept_id in dept_ids:
                    relation_candidates.append({
                        "staff_id": staff_id,
                        "department_id": dept_id,
                    })

            existing_relations = set(
                TtStaffSharedWithDepartment.objects.filter(
                    staff_id__in=[r["staff_id"] for r in relation_candidates]
                ).values_list("staff_id", "department_id")
            )

            to_create_rel = [
                TtStaffSharedWithDepartment(**r)
                for r in relation_candidates
                if (r["staff_id"], r["department_id"]) not in existing_relations
            ]

            if to_create_rel:
                TtStaffSharedWithDepartment.objects.bulk_create(to_create_rel, batch_size=500)

        return {
            "step": "staff_create",
            "created": len(created_staff),
            "skipped": len(cleaned_rows) - len(created_staff),
            "shared_department_rel_created": len(to_create_rel) if relation_candidates else 0,
        }

    except Exception:
        raise

def week_pattern_create(response_data):
    try:
        payload = normalize_payload(response_data, "week_pattern")

        cleaned_rows = []
        week_map = {}

        for row in payload:
            week_pattern_id = row.get("id")
            if not week_pattern_id:
                continue

            cleaned_rows.append({
                "id": week_pattern_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "department_id": row.get("department"),
                "academic_term_id": row.get("academic_term"),
                "status": row.get("status", 1),
            })

            week_map[week_pattern_id] = row.get("week", []) or []

        with transaction.atomic():
            created_patterns = bulk_create_skip_duplicates(TtWeekPattern, cleaned_rows)
            created_ids = {obj.id for obj in created_patterns}

            week_rel_created = 0

            if created_ids:
                week_through = TtWeekPattern.week.through

                week_relation_candidates = []
                for week_pattern_id, week_ids in week_map.items():
                    if week_pattern_id not in created_ids:
                        continue

                    for week_id in week_ids:
                        if week_id:
                            week_relation_candidates.append(
                                week_through(
                                    week_pattern_id=week_pattern_id,
                                    week_id=week_id,
                                )
                            )

                if week_relation_candidates:
                    week_through.objects.bulk_create(
                        week_relation_candidates,
                        batch_size=500,
                        ignore_conflicts=True,
                    )
                    week_rel_created = len(week_relation_candidates)

        return {
            "step": "week_pattern_create",
            "created": len(created_patterns),
            "skipped": len(cleaned_rows) - len(created_patterns),
            "week_rel": week_rel_created,
        }

    except Exception:
        raise

def week_create(response_data):
    try:
        payload = normalize_payload(response_data, "week")

        cleaned_rows = []
        for row in payload:
            if not row.get("id"):
                continue

            cleaned_rows.append({
                "id": row.get("id"),
                "week": row.get("week"),
                "start_date": row.get("start_date"),
            })

        with transaction.atomic():
            created = bulk_create_skip_duplicates(TtWeek, cleaned_rows)

        return {
            "step": "week_create",
            "created": len(created),
            "skipped": len(cleaned_rows) - len(created),
        }

    except Exception:
        raise

def academic_term_create(response_data):
    try:
        week_result = week_create(response_data)

        payload = normalize_payload(response_data, "academic_term")

        cleaned_rows = []
        week_map = {}

        for row in payload:
            academic_term_id = row.get("id")
            if not academic_term_id:
                continue

            cleaned_rows.append({
                "id": academic_term_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "start_date": row.get("start_date"),
                "end_date": row.get("end_date"),
                "start_time": row.get("start_time"),
                "end_time": row.get("end_time"),
                "status": row.get("status", 1),
            })

            week_map[academic_term_id] = row.get("week", []) or []

        with transaction.atomic():
            created_terms = bulk_create_skip_duplicates(TtAcademicTerm, cleaned_rows)
            created_ids = {obj.id for obj in created_terms}

            relation_candidates = []
            for academic_term_id, week_ids in week_map.items():
                if academic_term_id not in created_ids:
                    continue
                for week_id in week_ids:
                    relation_candidates.append({
                        "academic_term_id": academic_term_id,
                        "week_id": week_id,
                    })

            existing_relations = set(
                TtAcademicTermWeek.objects.filter(
                    academic_term_id__in=[r["academic_term_id"] for r in relation_candidates]
                ).values_list("academic_term_id", "week_id")
            )

            to_create_rel = [
                TtAcademicTermWeek(**r)
                for r in relation_candidates
                if (r["academic_term_id"], r["week_id"]) not in existing_relations
            ]

            if to_create_rel:
                TtAcademicTermWeek.objects.bulk_create(
                    to_create_rel,
                    batch_size=500
                )

        return {
            "step": "academic_term_create",
            "week_result": week_result,
            "created": len(created_terms),
            "skipped": len(cleaned_rows) - len(created_terms),
            "week_rel_created": len(to_create_rel) if relation_candidates else 0,
        }

    except Exception:
        raise

def module_group_create(response_data):
    try:
        payload = normalize_payload(response_data, "module_group")

        cleaned_rows = []

        for row in payload:
            module_group_id = row.get("id")
            if not module_group_id:
                continue

            cleaned_rows.append({
                "id": module_group_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "status": row.get("status", 1),
            })

        with transaction.atomic():
            created = bulk_create_skip_duplicates(TtModuleGroup, cleaned_rows)

        return {
            "step": "module_group_create",
            "created": len(created),
            "skipped": len(cleaned_rows) - len(created),
        }

    except Exception:
        raise

def pos_create(response_data):
    try:
        payload = normalize_payload(response_data, "pos")

        pos_rows = []
        pos_module_group_rows = []

        for row in payload:
            pos_id = row.get("id")
            if not pos_id:
                continue

            pos_rows.append({
                "id": pos_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "department_id": row.get("department"),
                "planned_size": row.get("planned_size"),
                "max_credit": row.get("max_credit"),
                "academic_term_id": row.get("academic_term"),
                "status": row.get("status", 1),
            })

            module_group_ids = row.get("module_group") or []

            for module_group_id in module_group_ids:
                if module_group_id:
                    pos_module_group_rows.append({
                        "pos_id": pos_id,
                        "module_group_id": module_group_id,
                    })

        with transaction.atomic():
            created = bulk_create_skip_duplicates(TtPos, pos_rows)

            pos_module_group_created = []
            if pos_module_group_rows:
                pos_module_group_created = bulk_create_skip_duplicates(
                    TtPosModuleGroup,
                    pos_module_group_rows
                )

        return {
            "step": "pos_create",
            "created": len(created),
            "skipped": len(pos_rows) - len(created),
            "pos_module_group_created": len(pos_module_group_created),
            "pos_module_group_skipped": len(pos_module_group_rows) - len(pos_module_group_created),
        }

    except Exception:
        raise

def activity_create(response_data):
    try:
        payload = normalize_payload(response_data, "activity")

        cleaned_rows = []
        rel_staff = []
        rel_location = []
        rel_student_set = []
        rel_week = {}

        valid_week_pattern_ids = get_valid_week_pattern_ids(payload)

        for row in payload:
            activity_id = row.get("id")
            if not activity_id:
                continue

            cleaned_rows.append({
                "id": activity_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "activity_template_id": row.get("activity_template"),
                "activity_type_id": row.get("activity_type"),
                "department_id": row.get("department"),
                "academic_term_id": row.get("academic_term"),
                "module_id": row.get("module"),
                "duration": row.get("duration"),
                "slot_required": row.get("slot_required"),
                "planned_size": row.get("planned_size"),
                "scheduled_start_time": row.get("scheduled_start_time"),
                "scheduled_day": row.get("scheduled_day"),
                "scheduled_start_slot": row.get("scheduled_start_slot"),
                "week_pattern_id": row.get("week_pattern"),
                "is_jta": row.get("is_jta", 0),
                "jta_parent_id": row.get("jta_parent"),
                "is_variant": row.get("is_variant", 0),
                "variant_parent_id": row.get("variant_parent"),
                "scheduled": row.get("scheduled", 0),
                "status": row.get("status", 1),
            })

            rel_staff.append((activity_id, row.get("staff", []) or []))
            rel_location.append((activity_id, row.get("location", []) or []))
            rel_student_set.append((activity_id, row.get("student_set", []) or []))

            if should_use_activity_week(row, valid_week_pattern_ids):
                rel_week[activity_id] = row.get("week", []) or []
            else:
                rel_week[activity_id] = []

        with transaction.atomic():
            created_activities = bulk_create_skip_duplicates(TtActivity, cleaned_rows)
            created_ids = {obj.id for obj in created_activities}

            staff_rows = []
            for activity_id, staff_ids in rel_staff:
                if activity_id not in created_ids:
                    continue
                for staff_id in staff_ids:
                    staff_rows.append({
                        "activity_id": activity_id,
                        "staff_id": staff_id,
                    })

            location_rows = []
            for activity_id, location_ids in rel_location:
                if activity_id not in created_ids:
                    continue
                for location_id in location_ids:
                    location_rows.append({
                        "activity_id": activity_id,
                        "location_id": location_id,
                    })

            student_rows = []
            for activity_id, student_set_ids in rel_student_set:
                if activity_id not in created_ids:
                    continue
                for student_set_id in student_set_ids:
                    student_rows.append({
                        "activity_id": activity_id,
                        "student_set_id": student_set_id,
                    })

            week_rows = []
            for activity_id, week_ids in rel_week.items():
                if activity_id not in created_ids:
                    continue
                for week_id in week_ids:
                    week_rows.append({
                        "activity_id": activity_id,
                        "week_id": week_id,
                    })

            if staff_rows:
                TtActivityStaff.objects.bulk_create(
                    [TtActivityStaff(**r) for r in staff_rows],
                    batch_size=500,
                    ignore_conflicts=True,
                )

            if location_rows:
                TtActivityLocation.objects.bulk_create(
                    [TtActivityLocation(**r) for r in location_rows],
                    batch_size=500,
                    ignore_conflicts=True,
                )

            if student_rows:
                TtStudentSetActivity.objects.bulk_create(
                    [TtStudentSetActivity(**r) for r in student_rows],
                    batch_size=500,
                    ignore_conflicts=True,
                )

            if week_rows:
                TtActivityWeek.objects.bulk_create(
                    [TtActivityWeek(**r) for r in week_rows],
                    batch_size=500,
                    ignore_conflicts=True,
                )

        return {
            "step": "activity_create",
            "created": len(created_activities),
            "skipped": len(cleaned_rows) - len(created_activities),
            "staff_rel": len(staff_rows),
            "location_rel": len(location_rows),
            "student_set_rel": len(student_rows),
            "week_rel": len(week_rows),
        }

    except Exception:
        raise

def build_activity_defaults(row, is_variant=False, variant_parent_id=None):
    return {
        "code": row.get("code"),
        "name": row.get("name"),
        "desc": row.get("desc"),
        "activity_template_id": row.get("activity_template"),
        "activity_type_id": row.get("activity_type"),
        "department_id": row.get("department"),
        "academic_term_id": row.get("academic_term"),
        "module_id": row.get("module"),
        "duration": row.get("duration"),
        "slot_required": row.get("slot_required"),
        "planned_size": row.get("planned_size"),
        "scheduled_start_time": row.get("scheduled_start_time"),
        "scheduled_day": row.get("scheduled_day"),
        "scheduled_start_slot": row.get("scheduled_start_slot"),
        "week_pattern_id": row.get("week_pattern") if "week_pattern" in row else row.get("week_pattern_id"),
        "is_jta": row.get("is_jta", 0),
        "jta_parent_id": row.get("jta_parent"),
        "is_variant": 1 if is_variant else row.get("is_variant", 0),
        "variant_parent_id": variant_parent_id if variant_parent_id is not None else row.get("variant_parent"),
        "scheduled": row.get("scheduled", 0),
        "status": row.get("status", 1),
    }

def variant_create(response_data):
    try:
        parent_row = response_data.get("variant_parent") or {}
        child_row = response_data.get("variant_child") or {}

        week_pattern_ids = set()

        if parent_row.get("week_pattern_id"):
            week_pattern_ids.add(parent_row.get("week_pattern_id"))

        if child_row.get("week_pattern"):
            week_pattern_ids.add(child_row.get("week_pattern"))

        valid_week_pattern_ids = set(
            TtWeekPattern.objects
            .filter(id__in=week_pattern_ids)
            .values_list("id", flat=True)
        ) if week_pattern_ids else set()

        with transaction.atomic():
            parent_id = parent_row.get("id")
            if parent_id:
                parent_obj = TtActivity.objects.filter(id=parent_id).first()
                if not parent_obj:
                    raise Exception(f"variant_parent activity id={parent_id} does not exist")

                parent_update_data = build_update_data(
                    parent_row,
                    {
                        "code": "code",
                        "name": "name",
                        "week_pattern_id": "week_pattern_id",
                        "is_variant": "is_variant",
                    },
                    single_item=True,
                )

                if parent_update_data:
                    TtActivity.objects.filter(id=parent_id).update(**parent_update_data)

                parent_obj.refresh_from_db()

                parent_week_pattern_id = (
                    parent_row.get("week_pattern_id")
                    if "week_pattern_id" in parent_row
                    else parent_obj.week_pattern_id
                )

                if parent_week_pattern_id and parent_week_pattern_id in valid_week_pattern_ids:
                    parent_obj.week.clear()
                else:
                    delete_week_ids = parent_row.get("delete_week_ids") or []
                    if delete_week_ids:
                        parent_obj.week.remove(*delete_week_ids)

                    add_week_ids = parent_row.get("week_ids") or []
                    if add_week_ids:
                        parent_obj.week.add(*add_week_ids)

            child_id = child_row.get("id")
            child_obj = None
            created = False

            if child_id:
                child_defaults = build_activity_defaults(
                    child_row,
                    is_variant=True,
                    variant_parent_id=child_row.get("variant_parent") or parent_id
                )

                child_defaults = {
                    k: v for k, v in child_defaults.items()
                    if v is not None
                }

                child_obj, created = TtActivity.objects.update_or_create(
                    id=child_id,
                    defaults=child_defaults
                )

                child_week_pattern_id = child_row.get("week_pattern") or child_obj.week_pattern_id

                if hasattr(child_obj, "week"):
                    if child_week_pattern_id and child_week_pattern_id in valid_week_pattern_ids:
                        child_obj.week.clear()
                    else:
                        child_obj.week.set(child_row.get("week", []) or child_row.get("week_ids", []))

                if hasattr(child_obj, "student_set"):
                    child_obj.student_set.set(child_row.get("student_set", []))

                staff_ids = get_valid_ids(
                    TtStaff,
                    child_row.get("staff", []) or []
                )

                if hasattr(child_obj, "staff"):
                    child_obj.staff.set(staff_ids)

                location_ids = get_valid_ids(
                    TtLocation,
                    child_row.get("location", []) or []
                )

                if hasattr(child_obj, "location"):
                    child_obj.location.set(location_ids)

        return {
            "step": "variant_create",
            "parent_updated": 1 if parent_id else 0,
            "child_action": "created" if created else "updated" if child_obj else "skipped",
            "child_count": 1 if child_obj else 0,
        }

    except Exception:
        raise

def jta_create(response_data):
    try:
        parent_row = response_data.get("create_activity") or {}
        child_rows = response_data.get("update_activity") or []

        week_pattern_ids = set()

        if parent_row.get("week_pattern"):
            week_pattern_ids.add(parent_row.get("week_pattern"))

        valid_week_pattern_ids = set(
            TtWeekPattern.objects
            .filter(id__in=week_pattern_ids)
            .values_list("id", flat=True)
        ) if week_pattern_ids else set()

        with transaction.atomic():
            parent_id = parent_row.get("id")
            if not parent_id:
                raise Exception("create_activity.id is required")

            parent_defaults = build_activity_defaults(parent_row, is_variant=False)
            parent_defaults["is_jta"] = 1

            parent_defaults = {
                k: v for k, v in parent_defaults.items()
                if v is not None
            }

            parent_obj, created = TtActivity.objects.update_or_create(
                id=parent_id,
                defaults=parent_defaults
            )

            parent_week_pattern_id = parent_row.get("week_pattern") or parent_obj.week_pattern_id

            if hasattr(parent_obj, "week"):
                if parent_week_pattern_id and parent_week_pattern_id in valid_week_pattern_ids:
                    parent_obj.week.clear()
                else:
                    parent_obj.week.set(parent_row.get("week", []) or parent_row.get("week_ids", []))

            if hasattr(parent_obj, "student_set"):
                parent_obj.student_set.set(parent_row.get("student_set", []))

            staff_ids = get_valid_ids(
                TtStaff,
                parent_row.get("staff", []) or []
            )

            if hasattr(parent_obj, "staff"):
                parent_obj.staff.set(staff_ids)

            location_ids = get_valid_ids(
                TtLocation,
                parent_row.get("location", []) or []
            )

            if hasattr(parent_obj, "location"):
                parent_obj.location.set(location_ids)

            child_ids = [row.get("id") for row in child_rows if row.get("id")]
            existing_child_ids = set(
                TtActivity.objects.filter(id__in=child_ids).values_list("id", flat=True)
            )

            updated_count = 0
            skipped_count = 0

            for row in child_rows:
                child_id = row.get("id")
                if not child_id:
                    skipped_count += 1
                    continue

                if child_id not in existing_child_ids:
                    skipped_count += 1
                    continue

                update_data = {
                    "is_jta": row.get("is_jta", 1),
                    "jta_parent_id": row.get("jta_parent_id") or parent_id,
                }

                TtActivity.objects.filter(id=child_id).update(**update_data)
                updated_count += 1

        return {
            "step": "jta_create",
            "parent_action": "created" if created else "updated",
            "parent_id": parent_id,
            "children_updated": updated_count,
            "children_skipped": skipped_count,
        }

    except Exception:
        raise

def activity_type_create(response_data):
    try:
        payload = normalize_payload(response_data, "activity_type")

        cleaned_rows = []
        for row in payload:
            activity_type_id = row.get("id")
            if not activity_type_id:
                continue

            cleaned_rows.append({
                "id": activity_type_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "department_id": row.get("department"),
                "booking": row.get("booking", False),
                "online": row.get("online", False),
                "contact": row.get("contact", False),
                "cover": row.get("cover", False),
                "color": row.get("color"),
                "status": row.get("status", 1),
            })

        with transaction.atomic():
            created = bulk_create_skip_duplicates(TtActivityType, cleaned_rows)

        return {
            "step": "activity_type_create",
            "created": len(created),
            "skipped": len(cleaned_rows) - len(created),
        }

    except Exception:
        raise

def module_create(response_data):
    try:
        payload = normalize_payload(response_data, "module")

        cleaned_rows = []
        week_map = {}

        for row in payload:
            module_id = row.get("id")
            if not module_id:
                continue

            cleaned_rows.append({
                "id": module_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "department_id": row.get("department"),
                "academic_term_id": row.get("academic_term"),
                "planned_size": row.get("planned_size"),
                "credit_provided": row.get("credit_provided"),
                "week_pattern_id": row.get("week_pattern"),
                "status": row.get("status", 1),
            })

            week_map[module_id] = row.get("week", [])

        with transaction.atomic():
            created_modules = bulk_create_skip_duplicates(TtModule, cleaned_rows)
            created_ids = {obj.id for obj in created_modules}

            relation_candidates = []
            for module_id, week_ids in week_map.items():
                if module_id not in created_ids:
                    continue
                for week_id in week_ids:
                    relation_candidates.append({
                        "module_id": module_id,
                        "week_id": week_id,
                    })

            existing_relations = set(
                TtModuleWeek.objects.filter(
                    module_id__in=[r["module_id"] for r in relation_candidates]
                ).values_list("module_id", "week_id")
            )

            to_create_rel = [
                TtModuleWeek(**r)
                for r in relation_candidates
                if (r["module_id"], r["week_id"]) not in existing_relations
            ]

            if to_create_rel:
                TtModuleWeek.objects.bulk_create(to_create_rel, batch_size=500)

        return {
            "step": "module_create",
            "created": len(created_modules),
            "skipped": len(cleaned_rows) - len(created_modules),
            "week_rel_created": len(to_create_rel) if relation_candidates else 0,
        }

    except Exception:
        raise

def activity_template_create(response_data):
    try:
        payload = normalize_payload(response_data, "activity_template")

        cleaned_rows = []
        week_map = {}

        for row in payload:
            activity_template_id = row.get("id")
            if not activity_template_id:
                continue

            cleaned_rows.append({
                "id": activity_template_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "activity_type_id": row.get("activity_type"),
                "academic_term_id": row.get("academic_term"),
                "department_id": row.get("department"),
                "duration": row.get("duration"),
                "slot_required": row.get("slot_required"),
                "module_id": row.get("module"),
                "planned_size": row.get("planned_size"),
                "week_pattern_id": row.get("week_pattern"),
                "status": row.get("status", 1),
            })

            week_map[activity_template_id] = row.get("week", [])

        with transaction.atomic():
            created_templates = bulk_create_skip_duplicates(TtActivityTemplate, cleaned_rows)
            created_ids = {obj.id for obj in created_templates}

            relation_candidates = []
            for activity_template_id, week_ids in week_map.items():
                if activity_template_id not in created_ids:
                    continue
                for week_id in week_ids:
                    relation_candidates.append({
                        "activity_template_id": activity_template_id,
                        "week_id": week_id,
                    })

            existing_relations = set(
                TtActivityTemplateWeek.objects.filter(
                    activity_template_id__in=[r["activity_template_id"] for r in relation_candidates]
                ).values_list("activity_template_id", "week_id")
            )

            to_create_rel = [
                TtActivityTemplateWeek(**r)
                for r in relation_candidates
                if (r["activity_template_id"], r["week_id"]) not in existing_relations
            ]

            if to_create_rel:
                TtActivityTemplateWeek.objects.bulk_create(to_create_rel, batch_size=500)

        return {
            "step": "activity_template_create",
            "created": len(created_templates),
            "skipped": len(cleaned_rows) - len(created_templates),
            "week_rel_created": len(to_create_rel) if relation_candidates else 0,
        }

    except Exception:
        raise

def student_create(response_data):
    try:
        payload = normalize_payload(response_data, "student")

        cleaned_rows = []

        academic_term_map = {}
        pos_map = {}
        pos_module_map = {}
        student_set_map = {}

        for row in payload:
            student_id = row.get("id")
            if not student_id:
                continue

            cleaned_rows.append({
                "id": student_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "email": row.get("email"),
                "department_id": row.get("department"),
                "status": row.get("status", 1),
            })

            academic_term_map[student_id] = (
                row.get("academic_term", []) or []
            )

            pos_map[student_id] = (
                row.get("pos", []) or []
            )

            pos_module_map[student_id] = (
                row.get("pos_module", []) or []
            )

            student_set_map[student_id] = (
                row.get("student_set", []) or []
            )

        with transaction.atomic():
            created_students = bulk_create_skip_duplicates(TtStudent, cleaned_rows)
            relation_student_ids = set(
                TtStudent.objects.filter(
                    id__in=[row["id"] for row in cleaned_rows]
                ).values_list("id", flat=True)
            )

            pos_module_keys = {
                (
                    pos_module.get("pos_id"),
                    pos_module.get("module_group_id"),
                    pos_module.get("module_id"),
                )
                for student_id in relation_student_ids
                for pos_module in pos_module_map.get(student_id, [])
                if isinstance(pos_module, dict)
                and pos_module.get("pos_id")
                and pos_module.get("module_group_id")
                and pos_module.get("module_id")
            }

            pos_module_lookup = {}
            if pos_module_keys:
                pos_ids = {key[0] for key in pos_module_keys}
                module_group_ids = {key[1] for key in pos_module_keys}
                module_ids = {key[2] for key in pos_module_keys}

                pos_module_lookup = {
                    (pos_id, module_group_id, module_id): relation_id
                    for pos_id, module_group_id, module_id, relation_id in (
                        TtPosModuleGroupModule.objects.filter(
                            pos_module_group__pos_id__in=pos_ids,
                            pos_module_group__module_group_id__in=module_group_ids,
                            module_id__in=module_ids,
                        ).values_list(
                            "pos_module_group__pos_id",
                            "pos_module_group__module_group_id",
                            "module_id",
                            "id",
                        )
                    )
                }

            raw_pos_module_ids = {
                pos_module
                for student_id in relation_student_ids
                for pos_module in pos_module_map.get(student_id, [])
                if not isinstance(pos_module, dict) and pos_module
            }
            valid_raw_pos_module_ids = set(
                TtPosModuleGroupModule.objects.filter(
                    id__in=raw_pos_module_ids
                ).values_list("id", flat=True)
            )

            academic_term_candidates = []
            student_set_candidates = []
            pos_candidates = []
            pos_module_candidates = []
            pos_module_rel_skipped = 0

            for student_id in relation_student_ids:
                for academic_term_id in academic_term_map.get(
                    student_id, []
                ):
                    if academic_term_id:
                        academic_term_candidates.append({
                            "student_id": student_id,
                            "academic_term_id": academic_term_id,
                        })

                for student_set_id in student_set_map.get(
                    student_id, []
                ):
                    if student_set_id:
                        student_set_candidates.append({
                            "student_id": student_id,
                            "student_set_id": student_set_id,
                        })

                for pos_id in pos_map.get(student_id, []):
                    pos_id = (
                        pos_id.get("pos_id")
                        if isinstance(pos_id, dict)
                        else pos_id
                    )

                    if pos_id:
                        pos_candidates.append({
                            "student_id": student_id,
                            "pos_id": pos_id,
                        })

                for pos_module in pos_module_map.get(student_id, []):
                    if isinstance(pos_module, dict):
                        pos_module_group_module_id = pos_module_lookup.get(
                            (
                                pos_module.get("pos_id"),
                                pos_module.get("module_group_id"),
                                pos_module.get("module_id"),
                            )
                        )
                    else:
                        pos_module_group_module_id = (
                            pos_module
                            if pos_module in valid_raw_pos_module_ids
                            else None
                        )

                    if pos_module_group_module_id:
                        pos_module_candidates.append({
                            "student_id": student_id,
                            "pos_module_group_module_id": pos_module_group_module_id,
                        })
                    else:
                        pos_module_rel_skipped += 1

            existing_academic_term_rel = set(
                TtStudentAcademicTerm.objects.filter(
                    student_id__in=[
                        candidate["student_id"]
                        for candidate in academic_term_candidates
                    ]
                ).values_list("student_id", "academic_term_id")
            )

            to_create_academic_term_rel = [
                TtStudentAcademicTerm(**candidate)
                for candidate in academic_term_candidates
                if (
                    candidate["student_id"],
                    candidate["academic_term_id"],
                ) not in existing_academic_term_rel
            ]

            if to_create_academic_term_rel:
                TtStudentAcademicTerm.objects.bulk_create(
                    to_create_academic_term_rel,
                    batch_size=500,
                )

            existing_student_set_rel = set(
                TtStudentSetStudent.objects.filter(
                    student_id__in=[
                        candidate["student_id"]
                        for candidate in student_set_candidates
                    ]
                ).values_list("student_id", "student_set_id")
            )


            to_create_student_set_rel = [
                TtStudentSetStudent(**candidate)
                for candidate in student_set_candidates
                if (
                       candidate["student_id"],
                       candidate["student_set_id"],
                   ) not in existing_student_set_rel
            ]

            if to_create_student_set_rel:
                TtStudentSetStudent.objects.bulk_create(
                    to_create_student_set_rel,
                    batch_size=500,
                )

            existing_pos_rel = set(
                TtStudentPos.objects.filter(
                    student_id__in=[
                        candidate["student_id"]
                        for candidate in pos_candidates
                    ]
                ).values_list("student_id", "pos_id")
            )

            to_create_pos_rel = [
                TtStudentPos(**candidate)
                for candidate in pos_candidates
                if (
                    candidate["student_id"],
                    candidate["pos_id"],
                ) not in existing_pos_rel
            ]

            if to_create_pos_rel:
                TtStudentPos.objects.bulk_create(
                    to_create_pos_rel,
                    batch_size=500,
                )

            existing_pos_module_rel = set(
                TtStudentPosModule.objects.filter(
                    student_id__in=[
                        candidate["student_id"]
                        for candidate in pos_module_candidates
                    ]
                ).values_list(
                    "student_id",
                    "pos_module_group_module_id",
                )
            )

            to_create_pos_module_rel = [
                TtStudentPosModule(**candidate)
                for candidate in pos_module_candidates
                if (
                    candidate["student_id"],
                    candidate["pos_module_group_module_id"],
                ) not in existing_pos_module_rel
            ]

            if to_create_pos_module_rel:
                TtStudentPosModule.objects.bulk_create(
                    to_create_pos_module_rel,
                    batch_size=500,
                )

        return {
            "step": "student_create",
            "created": len(created_students),
            "skipped": len(cleaned_rows) - len(created_students),
            "academic_term_rel": len(to_create_academic_term_rel),
            "student_set_rel": len(to_create_student_set_rel),
            "pos_rel": len(to_create_pos_rel),
            "pos_module_rel": len(to_create_pos_module_rel),
            "pos_module_rel_skipped": pos_module_rel_skipped,
        }

    except Exception:
        raise

def student_set_create(response_data):
    try:
        payload = normalize_payload(response_data, "student_set")
        student_allocate_payload = normalize_payload(response_data, "student_allocate")

        cleaned_rows = []
        module_map = {}
        activity_map = {}

        for row in payload:
            student_set_id = row.get("id")
            if not student_set_id:
                continue

            cleaned_rows.append({
                "id": student_set_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "department_id": row.get("department"),
                "academic_term_id": row.get("academic_term"),
                "pos_id": row.get("pos"),
                "planned_size": row.get("planned_size"),
                "status": row.get("status", 1),
            })

            module_map[student_set_id] = row.get("module", [])
            activity_map[student_set_id] = row.get("activity", [])

        with transaction.atomic():
            created_sets = bulk_create_skip_duplicates(TtStudentSet, cleaned_rows)
            created_ids = {obj.id for obj in created_sets}

            module_candidates = []
            for student_set_id, module_ids in module_map.items():
                if student_set_id not in created_ids:
                    continue
                for module_id in module_ids:
                    module_candidates.append({
                        "student_set_id": student_set_id,
                        "module_id": module_id,
                    })

            activity_candidates = []
            for student_set_id, activity_ids in activity_map.items():
                if student_set_id not in created_ids:
                    continue
                for activity_id in activity_ids:
                    activity_candidates.append({
                        "student_set_id": student_set_id,
                        "activity_id": activity_id,
                    })

            existing_module_rel = set(
                TtStudentSetModule.objects.filter(
                    student_set_id__in=[r["student_set_id"] for r in module_candidates]
                ).values_list("student_set_id", "module_id")
            )

            to_create_module_rel = [
                TtStudentSetModule(**r)
                for r in module_candidates
                if (r["student_set_id"], r["module_id"]) not in existing_module_rel
            ]

            if to_create_module_rel:
                TtStudentSetModule.objects.bulk_create(
                    to_create_module_rel,
                    batch_size=500
                )

            existing_activity_rel = set(
                TtStudentSetActivity.objects.filter(
                    student_set_id__in=[r["student_set_id"] for r in activity_candidates]
                ).values_list("student_set_id", "activity_id")
            )

            to_create_activity_rel = [
                TtStudentSetActivity(**r)
                for r in activity_candidates
                if (r["student_set_id"], r["activity_id"]) not in existing_activity_rel
            ]

            if to_create_activity_rel:
                TtStudentSetActivity.objects.bulk_create(
                    to_create_activity_rel,
                    batch_size=500
                )

        # after success create student set, check got allocate_student param or not, if got, allocate student into student set
        to_create_student_rel = []
        if student_allocate_payload:
            student_candidates = []
            for row in student_allocate_payload:
                for student_id in row["student_ids"]:
                    student_candidates.append({
                        "student_set_id": row["student_set_id"],
                        "student_id": student_id,
                    })

            existing_student_rel = set(
                TtStudentSetStudent.objects.filter(
                    student_set_id__in=[row["student_set_id"] for row in student_candidates]
                ).values_list("student_set_id", "student_id")
            )

            to_create_student_rel = [
                TtStudentSetStudent(**row)
                for row in student_candidates
                if (row["student_set_id"], row["student_id"]) not in existing_student_rel
            ]

            if to_create_student_rel:
                TtStudentSetStudent.objects.bulk_create(
                    to_create_student_rel,
                    batch_size=500
                )

        return {
            "step": "student_set_create",
            "created": len(created_sets),
            "skipped": len(cleaned_rows) - len(created_sets),
            "module_rel": len(to_create_module_rel) if module_candidates else 0,
            "activity_rel": len(to_create_activity_rel) if activity_candidates else 0,
            "student_rel": len(to_create_student_rel) if to_create_student_rel else 0,
        }

    except Exception:
        raise

def booking_create(response_data):
    try:
        payload = normalize_payload(response_data, "activity")

        cleaned_rows = []
        week_map = {}

        for row in payload:
            activity_id = row.get("id")
            if not activity_id:
                continue

            cleaned_rows.append({
                "id": activity_id,
                "code": row.get("code"),
                "name": row.get("name"),
                "desc": row.get("desc"),
                "activity_type_id": row.get("activity_type"),
                "department_id": row.get("department"),
                "duration": row.get("duration"),
                "slot_required": row.get("slot_required"),
                "planned_size": row.get("planned_size"),
                "is_jta": row.get("is_jta", 0),
                "is_variant": row.get("is_variant", 0),
                "scheduled": row.get("scheduled", 0),
                "is_booking": row.get("is_booking", 1),
                "status": row.get("status", 1),
            })

            week_map[activity_id] = row.get("week", []) or []

        with transaction.atomic():
            created_activities = bulk_create_skip_duplicates(TtActivity, cleaned_rows)
            created_ids = {obj.id for obj in created_activities}

            week_rows = []
            for activity_id, week_ids in week_map.items():
                if activity_id not in created_ids:
                    continue

                for week_id in week_ids:
                    if week_id:
                        week_rows.append(
                            TtActivityWeek(
                                activity_id=activity_id,
                                week_id=week_id,
                            )
                        )

            if week_rows:
                TtActivityWeek.objects.bulk_create(
                    week_rows,
                    batch_size=500,
                    ignore_conflicts=True,
                )

        return {
            "step": "booking_create",
            "created": len(created_activities),
            "skipped": len(cleaned_rows) - len(created_activities),
            "week_rel": len(week_rows),
        }

    except Exception:
        raise

# ============================
# UPDATE FUNCTIONS
# ============================

def department_update(response_data):
    try:
        payload = normalize_payload(response_data, "department")

        if not payload:
            return {
                "step": "department_update",
                "updated": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtDepartment.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                department_id = row.get("id")
                if not department_id or department_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "department_id": "department_id",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtDepartment.objects.filter(id=department_id).update(**update_data)

                updated += 1

        return {
            "step": "department_update",
            "updated": updated,
            "skipped": skipped,
        }

    except Exception:
        raise

def location_update(response_data):
    try:
        payload = normalize_payload(response_data, "location")

        if not payload:
            return {
                "step": "location_update",
                "updated": 0,
                "shared_department_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtLocation.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        shared_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                location_id = row.get("id")
                if not location_id or location_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "department_id": "department_id",
                        "capacity": "capacity",
                        "area": "area",
                        "maximum_period": "maximum_period",
                        "contract_period": "contract_period",
                        "shared_with_all_department": "shared_with_all_department",
                        "is_online": "is_online",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtLocation.objects.filter(id=location_id).update(**update_data)

                if "shared_with_department" in row:
                    shared_map[location_id] = row.get("shared_with_department") or []

                updated += 1

            shared_department_rel = replace_relation_rows(
                TtLocationSharedWithDepartment,
                "location",
                "department",
                shared_map,
            ) if shared_map else 0

        return {
            "step": "location_update",
            "updated": updated,
            "shared_department_rel": shared_department_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def staff_update(response_data):
    try:
        payload = normalize_payload(response_data, "staff")

        if not payload:
            return {
                "step": "staff_update",
                "updated": 0,
                "shared_department_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtStaff.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        shared_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                staff_id = row.get("id")
                if not staff_id or staff_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "email": "email",
                        "desc": "desc",
                        "department_id": "department_id",
                        "is_part_time": "is_part_time",
                        "maximum_period": "maximum_period",
                        "contract_period": "contract_period",
                        "shared_with_all_department": "shared_with_all_department",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtStaff.objects.filter(id=staff_id).update(**update_data)

                if "shared_with_department" in row:
                    shared_map[staff_id] = row.get("shared_with_department") or []

                updated += 1

            shared_department_rel = replace_relation_rows(
                TtStaffSharedWithDepartment,
                "staff",
                "department",
                shared_map,
            ) if shared_map else 0

        return {
            "step": "staff_update",
            "updated": updated,
            "shared_department_rel": shared_department_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def week_pattern_update(response_data):
    try:
        payload = normalize_payload(response_data, "week_pattern")

        if not payload:
            return {
                "step": "week_pattern_update",
                "updated": 0,
                "week_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtWeekPattern.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        week_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                week_pattern_id = row.get("id")
                if not week_pattern_id or week_pattern_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "department_id": "department_id",
                        "academic_term_id": "academic_term_id",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtWeekPattern.objects.filter(id=week_pattern_id).update(**update_data)

                if "week" in row:
                    week_map[week_pattern_id] = row.get("week") or []

                updated += 1

            week_rel = 0
            if week_map:
                week_through = TtWeekPattern.week.through

                week_through.objects.filter(
                    week_pattern_id__in=list(week_map.keys())
                ).delete()

                new_rows = []
                for week_pattern_id, week_ids in week_map.items():
                    for week_id in week_ids:
                        if week_id:
                            new_rows.append(
                                week_through(
                                    week_pattern_id=week_pattern_id,
                                    week_id=week_id,
                                )
                            )

                if new_rows:
                    week_through.objects.bulk_create(
                        new_rows,
                        batch_size=500,
                        ignore_conflicts=True,
                    )

                week_rel = len(new_rows)

        return {
            "step": "week_pattern_update",
            "updated": updated,
            "week_rel": week_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def academic_term_update(response_data):
    try:
        payload = normalize_payload(response_data, "academic_term")

        if not payload:
            return {
                "step": "academic_term_update",
                "updated": 0,
                "week_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtAcademicTerm.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        week_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                academic_term_id = row.get("id")
                if not academic_term_id or academic_term_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "start_date": "start_date",
                        "end_date": "end_date",
                        "start_time": "start_time",
                        "end_time": "end_time",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtAcademicTerm.objects.filter(id=academic_term_id).update(**update_data)

                if "week" in row:
                    week_map[academic_term_id] = row.get("week") or []

                updated += 1

            week_rel = replace_relation_rows(
                TtAcademicTermWeek,
                "academic_term",
                "week",
                week_map,
            ) if week_map else 0

        return {
            "step": "academic_term_update",
            "updated": updated,
            "week_rel": week_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def pos_update(response_data):
    try:
        payload = normalize_payload(response_data, "pos")

        if not payload:
            return {
                "step": "pos_update",
                "updated": 0,
                "module_group_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtPos.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        module_group_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                pos_id = row.get("id")
                if not pos_id or pos_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "department_id": "department_id",
                        "planned_size": "planned_size",
                        "max_credit": "max_credit",
                        "academic_term_id": "academic_term_id",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtPos.objects.filter(id=pos_id).update(**update_data)

                if "module_group" in row:
                    module_group_map[pos_id] = row.get("module_group") or []

                updated += 1

            module_group_rel = 0
            if module_group_map:
                TtPosModuleGroup.objects.filter(
                    pos_id__in=list(module_group_map.keys())
                ).delete()

                new_rows = []
                for pos_id, module_group_ids in module_group_map.items():
                    for module_group_id in module_group_ids:
                        if module_group_id:
                            new_rows.append(
                                TtPosModuleGroup(
                                    pos_id=pos_id,
                                    module_group_id=module_group_id,
                                )
                            )

                if new_rows:
                    TtPosModuleGroup.objects.bulk_create(
                        new_rows,
                        batch_size=500,
                        ignore_conflicts=True,
                    )

                module_group_rel = len(new_rows)

        return {
            "step": "pos_update",
            "updated": updated,
            "module_group_rel": module_group_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def module_update(response_data):
    try:
        payload = normalize_payload(response_data, "module")

        if not payload:
            return {
                "step": "module_update",
                "updated": 0,
                "week_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtModule.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        week_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                module_id = row.get("id")
                if not module_id or module_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "department_id": "department_id",
                        "academic_term_id": "academic_term_id",
                        "planned_size": "planned_size",
                        "credit_provided": "credit_provided",
                        "week_pattern_id": "week_pattern_id",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtModule.objects.filter(id=module_id).update(**update_data)

                if "week" in row:
                    week_map[module_id] = row.get("week") or []

                updated += 1

            week_rel = replace_relation_rows(
                TtModuleWeek,
                "module",
                "week",
                week_map,
            ) if week_map else 0

        return {
            "step": "module_update",
            "updated": updated,
            "week_rel": week_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def module_group_update(response_data):
    try:
        payload = normalize_payload(response_data, "module_group")

        if not payload:
            return {
                "step": "module_group_update",
                "updated": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtModuleGroup.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                module_group_id = row.get("id")
                if not module_group_id or module_group_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtModuleGroup.objects.filter(id=module_group_id).update(**update_data)

                updated += 1

        return {
            "step": "module_group_update",
            "updated": updated,
            "skipped": skipped,
        }

    except Exception:
        raise

def activity_type_update(response_data):
    try:
        payload = normalize_payload(response_data, "activity_type")

        if not payload:
            return {
                "step": "activity_type_update",
                "updated": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtActivityType.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                activity_type_id = row.get("id")
                if not activity_type_id or activity_type_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "department_id": "department_id",
                        "booking": "booking",
                        "online": "online",
                        "contact": "contact",
                        "cover": "cover",
                        "color": "color",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtActivityType.objects.filter(id=activity_type_id).update(**update_data)

                updated += 1

        return {
            "step": "activity_type_update",
            "updated": updated,
            "skipped": skipped,
        }

    except Exception:
        raise

def activity_template_update(response_data):
    try:
        payload = normalize_payload(response_data, "activity_template")

        if not payload:
            return {
                "step": "activity_template_update",
                "updated": 0,
                "week_rel": 0,
                "activity_updated": 0,
                "activity_skipped": 0,
                "activity_week_rel": 0,
                "activity_week_deleted": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtActivityTemplate.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        week_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                activity_template_id = row.get("id")
                if not activity_template_id or activity_template_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "department_id": "department_id",
                        "duration": "duration",
                        "slot_required": "slot_required",
                        "planned_size": "planned_size",
                        "week_pattern_id": "week_pattern_id",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtActivityTemplate.objects.filter(id=activity_template_id).update(**update_data)

                if "week" in row:
                    week_map[activity_template_id] = row.get("week") or []

                updated += 1

            week_rel = replace_relation_rows(
                TtActivityTemplateWeek,
                "activity_template",
                "week",
                week_map,
            ) if week_map else 0

            activity_payload = normalize_payload(response_data, "activity")
            activity_updated = 0
            activity_skipped = 0
            activity_week_rel = 0
            activity_week_deleted = 0

            if activity_payload:
                activity_ids = [row.get("id") for row in activity_payload if row.get("id")]

                existing_activity_map = {
                    obj.id: obj
                    for obj in TtActivity.objects.filter(id__in=activity_ids)
                }

                existing_activity_ids = set(existing_activity_map.keys())

                incoming_week_pattern_ids = {
                    row.get("week_pattern_id")
                    for row in activity_payload
                    if row.get("week_pattern_id")
                }

                existing_week_pattern_ids = {
                    obj.week_pattern_id
                    for obj in existing_activity_map.values()
                    if obj.week_pattern_id
                }

                valid_week_pattern_ids = set(
                    TtWeekPattern.objects
                    .filter(id__in=(incoming_week_pattern_ids | existing_week_pattern_ids))
                    .values_list("id", flat=True)
                )

                activity_week_map = {}
                activities_with_valid_week_pattern = set()

                for row in activity_payload:
                    activity_id = row.get("id")
                    if not activity_id or activity_id not in existing_activity_ids:
                        activity_skipped += 1
                        continue

                    existing_activity = existing_activity_map[activity_id]

                    effective_week_pattern_id = (
                        row.get("week_pattern_id")
                        if "week_pattern_id" in row
                        else existing_activity.week_pattern_id
                    )

                    if effective_week_pattern_id and effective_week_pattern_id in valid_week_pattern_ids:
                        activities_with_valid_week_pattern.add(activity_id)

                    activity_update_data = build_update_data(
                        row,
                        {
                            "code": "code",
                            "name": "name",
                            "desc": "desc",
                            "activity_template_id": "activity_template_id",
                            "activity_type_id": "activity_type_id",
                            "department_id": "department_id",
                            "academic_term_id": "academic_term_id",
                            "module_id": "module_id",
                            "duration": "duration",
                            "slot_required": "slot_required",
                            "planned_size": "planned_size",
                            "scheduled_start_time": "scheduled_start_time",
                            "scheduled_day": "scheduled_day",
                            "scheduled_start_slot": "scheduled_start_slot",
                            "week_pattern_id": "week_pattern_id",
                            "is_jta": "is_jta",
                            "jta_parent": "jta_parent_id",
                            "is_variant": "is_variant",
                            "variant_parent": "variant_parent_id",
                            "scheduled": "scheduled",
                            "status": "status",
                        },
                        single_item=(len(activity_payload) == 1),
                    )

                    if activity_update_data:
                        TtActivity.objects.filter(id=activity_id).update(**activity_update_data)

                    if "week" in row:
                        if effective_week_pattern_id and effective_week_pattern_id in valid_week_pattern_ids:
                            activity_week_map[activity_id] = []
                        else:
                            activity_week_map[activity_id] = row.get("week") or []

                    activity_updated += 1

                if activities_with_valid_week_pattern:
                    activity_week_deleted, _ = TtActivityWeek.objects.filter(
                        activity_id__in=activities_with_valid_week_pattern
                    ).delete()

                activity_week_rel = replace_relation_rows(
                    TtActivityWeek,
                    "activity",
                    "week",
                    activity_week_map,
                ) if activity_week_map else 0

        return {
            "step": "activity_template_update",
            "updated": updated,
            "week_rel": week_rel,
            "activity_updated": activity_updated,
            "activity_skipped": activity_skipped,
            "activity_week_rel": activity_week_rel,
            "activity_week_deleted": activity_week_deleted,
            "skipped": skipped,
        }

    except Exception:
        raise

def student_update(response_data):
    try:
        payload = normalize_payload(response_data, "student")

        if not payload:
            return {
                "step": "student_update",
                "updated": 0,
                "academic_term_rel": 0,
                "pos_rel": 0,
                "pos_module_rel": 0,
                "pos_module_rel_skipped": 0,
                # "student_set_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtStudent.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        academic_term_map = {}
        pos_map = {}
        pos_module_payload_map = {}
        student_set_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                student_id = row.get("id")
                if not student_id or student_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "email": "email",
                        "desc": "desc",
                        "department_id": "department_id",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtStudent.objects.filter(id=student_id).update(**update_data)

                if "academic_term" in row:
                    academic_term_map[student_id] = row.get("academic_term") or []

                if "pos" in row:
                    pos_map[student_id] = [
                        pos.get("pos_id") if isinstance(pos, dict) else pos
                        for pos in (row.get("pos") or [])
                    ]

                if "pos_module" in row:
                    pos_module_payload_map[student_id] = row.get("pos_module") or []

                if "student_set" in row:
                    student_set_map[student_id] = row.get("student_set") or []

                updated += 1

            pos_module_keys = {
                (
                    pos_module.get("pos_id"),
                    pos_module.get("module_group_id"),
                    pos_module.get("module_id"),
                )
                for pos_modules in pos_module_payload_map.values()
                for pos_module in pos_modules
                if isinstance(pos_module, dict)
                and pos_module.get("pos_id")
                and pos_module.get("module_group_id")
                and pos_module.get("module_id")
            }

            pos_module_lookup = {}
            if pos_module_keys:
                pos_ids = {key[0] for key in pos_module_keys}
                module_group_ids = {key[1] for key in pos_module_keys}
                module_ids = {key[2] for key in pos_module_keys}

                pos_module_lookup = {
                    (pos_id, module_group_id, module_id): relation_id
                    for pos_id, module_group_id, module_id, relation_id in (
                        TtPosModuleGroupModule.objects.filter(
                            pos_module_group__pos_id__in=pos_ids,
                            pos_module_group__module_group_id__in=module_group_ids,
                            module_id__in=module_ids,
                        ).values_list(
                            "pos_module_group__pos_id",
                            "pos_module_group__module_group_id",
                            "module_id",
                            "id",
                        )
                    )
                }

            raw_pos_module_ids = {
                pos_module
                for pos_modules in pos_module_payload_map.values()
                for pos_module in pos_modules
                if not isinstance(pos_module, dict) and pos_module
            }
            valid_raw_pos_module_ids = set(
                TtPosModuleGroupModule.objects.filter(
                    id__in=raw_pos_module_ids
                ).values_list("id", flat=True)
            )

            pos_module_map = {}
            pos_module_rel_skipped = 0
            for student_id, pos_modules in pos_module_payload_map.items():
                local_ids = []
                for pos_module in pos_modules:
                    if isinstance(pos_module, dict):
                        local_id = pos_module_lookup.get(
                            (
                                pos_module.get("pos_id"),
                                pos_module.get("module_group_id"),
                                pos_module.get("module_id"),
                            )
                        )
                    else:
                        local_id = (
                            pos_module
                            if pos_module in valid_raw_pos_module_ids
                            else None
                        )

                    if local_id:
                        local_ids.append(local_id)
                    else:
                        pos_module_rel_skipped += 1

                pos_module_map[student_id] = local_ids

            academic_term_rel = replace_relation_rows(
                TtStudentAcademicTerm,
                "student",
                "academic_term",
                academic_term_map,
            ) if academic_term_map else 0

            pos_rel = replace_relation_rows(
                TtStudentPos,
                "student",
                "pos",
                pos_map,
            ) if pos_map else 0

            pos_module_rel = replace_relation_rows(
                TtStudentPosModule,
                "student",
                "pos_module_group_module",
                pos_module_map,
            ) if pos_module_map else 0

            # student_set_rel = replace_relation_rows(
            #     TtStudentSetStudent,
            #     "student",
            #     "student_set",
            #     student_set_map,
            # ) if student_set_map else 0

        return {
            "step": "student_update",
            "updated": updated,
            "academic_term_rel": academic_term_rel,
            "pos_rel": pos_rel,
            "pos_module_rel": pos_module_rel,
            "pos_module_rel_skipped": pos_module_rel_skipped,
            # "student_set_rel": student_set_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def student_set_update(response_data):
    try:
        payload = normalize_payload(response_data, "student_set")

        if not payload:
            return {
                "step": "student_set_update",
                "updated": 0,
                "module_rel": 0,
                "activity_rel": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtStudentSet.objects.filter(id__in=incoming_ids).values_list("id", flat=True)
        )

        module_map = {}
        activity_map = {}
        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                student_set_id = row.get("id")
                if not student_set_id or student_set_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "department_id": "department_id",
                        "academic_term_id": "academic_term_id",
                        "pos_id": "pos_id",
                        "planned_size": "planned_size",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtStudentSet.objects.filter(id=student_set_id).update(**update_data)

                if "module" in row:
                    module_map[student_set_id] = row.get("module") or []

                if "activity" in row:
                    activity_map[student_set_id] = row.get("activity") or []

                updated += 1

            module_rel = replace_relation_rows(
                TtStudentSetModule,
                "student_set",
                "module",
                module_map,
            ) if module_map else 0

            activity_rel = replace_relation_rows(
                TtStudentSetActivity,
                "student_set",
                "activity",
                activity_map,
            ) if activity_map else 0

        return {
            "step": "student_set_update",
            "updated": updated,
            "module_rel": module_rel,
            "activity_rel": activity_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

def activity_update(response_data):
    try:
        payload = normalize_payload(response_data, "activity")

        if not payload:
            return {
                "step": "activity_update",
                "updated": 0,
                "staff_rel": 0,
                "location_rel": 0,
                "student_set_rel": 0,
                "week_rel": 0,
                "activity_week_deleted": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]

        existing_activity_map = {
            obj.id: obj
            for obj in TtActivity.objects.filter(id__in=incoming_ids)
        }

        existing_ids = set(existing_activity_map.keys())

        incoming_week_pattern_ids = {
            row.get("week_pattern_id")
            for row in payload
            if row.get("week_pattern_id")
        }

        existing_week_pattern_ids = {
            obj.week_pattern_id
            for obj in existing_activity_map.values()
            if obj.week_pattern_id
        }

        valid_week_pattern_ids = set(
            TtWeekPattern.objects
            .filter(id__in=(incoming_week_pattern_ids | existing_week_pattern_ids))
            .values_list("id", flat=True)
        )

        staff_map = {}
        location_map = {}
        student_set_map = {}
        week_map = {}
        activities_with_valid_week_pattern = set()

        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                activity_id = row.get("id")
                if not activity_id or activity_id not in existing_ids:
                    skipped += 1
                    continue

                existing_activity = existing_activity_map[activity_id]

                effective_week_pattern_id = (
                    row.get("week_pattern_id")
                    if "week_pattern_id" in row
                    else existing_activity.week_pattern_id
                )

                if effective_week_pattern_id and effective_week_pattern_id in valid_week_pattern_ids:
                    activities_with_valid_week_pattern.add(activity_id)

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "activity_template_id": "activity_template_id",
                        "activity_type_id": "activity_type_id",
                        "department_id": "department_id",
                        "academic_term_id": "academic_term_id",
                        "module_id": "module_id",
                        "duration": "duration",
                        "slot_required": "slot_required",
                        "planned_size": "planned_size",
                        "scheduled_start_time": "scheduled_start_time",
                        "scheduled_day": "scheduled_day",
                        "scheduled_start_slot": "scheduled_start_slot",
                        "week_pattern_id": "week_pattern_id",
                        "is_jta": "is_jta",
                        "jta_parent": "jta_parent_id",
                        "is_variant": "is_variant",
                        "variant_parent": "variant_parent_id",
                        "scheduled": "scheduled",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtActivity.objects.filter(id=activity_id).update(**update_data)

                if "staff" in row:
                    staff_map[activity_id] = row.get("staff") or []

                if "location" in row:
                    location_map[activity_id] = row.get("location") or []

                if "student_set" in row:
                    student_set_map[activity_id] = row.get("student_set") or []

                if "week" in row:
                    if effective_week_pattern_id and effective_week_pattern_id in valid_week_pattern_ids:
                        week_map[activity_id] = []
                    else:
                        week_map[activity_id] = row.get("week") or []

                updated += 1

            staff_rel = replace_relation_rows(
                TtActivityStaff,
                "activity",
                "staff",
                staff_map,
            ) if staff_map else 0

            location_rel = replace_relation_rows(
                TtActivityLocation,
                "activity",
                "location",
                location_map,
            ) if location_map else 0

            student_set_rel = replace_relation_rows(
                TtStudentSetActivity,
                "activity",
                "student_set",
                student_set_map,
            ) if student_set_map else 0

            activity_week_deleted = 0
            if activities_with_valid_week_pattern:
                activity_week_deleted, _ = TtActivityWeek.objects.filter(
                    activity_id__in=activities_with_valid_week_pattern
                ).delete()

            week_rel = replace_relation_rows(
                TtActivityWeek,
                "activity",
                "week",
                week_map,
            ) if week_map else 0

        return {
            "step": "activity_update",
            "updated": updated,
            "staff_rel": staff_rel,
            "location_rel": location_rel,
            "student_set_rel": student_set_rel,
            "week_rel": week_rel,
            "activity_week_deleted": activity_week_deleted,
            "skipped": skipped,
        }

    except Exception:
        raise

def booking_update(response_data):
    try:
        payload = normalize_payload(response_data, "activity")

        if not payload:
            return {
                "step": "booking_update",
                "updated": 0,
                "week_rel": 0,
                "activity_week_deleted": 0,
                "skipped": 0,
            }

        single_item = len(payload) == 1
        incoming_ids = [row.get("id") for row in payload if row.get("id")]

        existing_activity_map = {
            obj.id: obj
            for obj in TtActivity.objects.filter(id__in=incoming_ids)
        }

        existing_ids = set(existing_activity_map.keys())

        week_map = {}

        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                activity_id = row.get("id")
                if not activity_id or activity_id not in existing_ids:
                    skipped += 1
                    continue

                update_data = build_update_data(
                    row,
                    {
                        "code": "code",
                        "name": "name",
                        "desc": "desc",
                        "activity_type_id": "activity_type_id",
                        "department_id": "department_id",
                        "duration": "duration",
                        "slot_required": "slot_required",
                        "planned_size": "planned_size",
                        "status": "status",
                    },
                    single_item=single_item,
                )

                if update_data:
                    TtActivity.objects.filter(id=activity_id).update(**update_data)

                if "week" in row:
                    week_map[activity_id] = row.get("week") or []

                updated += 1
            activity_week_deleted = 0
            week_rel = replace_relation_rows(
                TtActivityWeek,
                "activity",
                "week",
                week_map,
            ) if week_map else 0

        return {
            "step": "activity_update",
            "updated": updated,
            "week_rel": week_rel,
            "activity_week_deleted": activity_week_deleted,
            "skipped": skipped,
        }

    except Exception:
        raise

# ============================
# DELETE FUNCTIONS
# ============================

def academic_term_delete(response_data):
    try:
        ids = response_data.get("academic_term_ids") or []

        if not ids:
            return {
                "step": "academic_term_delete",
                "deleted": 0,
                "skipped": 0,
                "academic_term_week_rel_deleted": 0,
            }

        existing_ids = set(
            TtAcademicTerm.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "academic_term_delete",
                "deleted": 0,
                "skipped": len(ids),
                "academic_term_week_rel_deleted": 0,
            }

        with transaction.atomic():
            setting_removed_ids = 0

            setting_value = Setting.get_setting("academic_term")
            if setting_value not in [None, ""]:
                parsed_setting = setting_value
                if isinstance(setting_value, str):
                    try:
                        parsed_setting = json.loads(setting_value)
                    except Exception:
                        parsed_setting = [setting_value]

                if isinstance(parsed_setting, (int, str)):
                    parsed_setting = [parsed_setting]

                if isinstance(parsed_setting, (list, tuple, set)):
                    normalized_setting_ids = []
                    for item in parsed_setting:
                        try:
                            normalized_setting_ids.append(int(item))
                        except (TypeError, ValueError):
                            continue

                    kept_setting_ids = [
                        term_id for term_id in normalized_setting_ids
                        if term_id not in existing_ids
                    ]

                    setting_removed_ids = len(normalized_setting_ids) - len(kept_setting_ids)
                    if setting_removed_ids > 0:
                        Setting.objects.update_or_create(
                            param="academic_term",
                            defaults={
                                "value": json.dumps(kept_setting_ids) if kept_setting_ids else None
                            },
                        )

            academic_term_week_rel_deleted = TtAcademicTermWeek.objects.filter(
                academic_term_id__in=existing_ids
            ).count()

            deleted, _ = TtAcademicTerm.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "academic_term_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "academic_term_week_rel_deleted": academic_term_week_rel_deleted,
        }

    except Exception:
        raise

def department_delete(response_data):
    try:
        ids = response_data.get("department_ids") or []

        if not ids:
            return {
                "step": "department_delete",
                "deleted": 0,
                "skipped": 0,
                "child_department_deleted": 0,
            }

        existing_ids = set(
            TtDepartment.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "department_delete",
                "deleted": 0,
                "skipped": len(ids),
                "child_department_deleted": 0,
            }

        with transaction.atomic():
            child_department_deleted = TtDepartment.objects.filter(
                department_id__in=existing_ids
            ).exclude(id__in=existing_ids).count()

            deleted, _ = TtDepartment.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "department_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "child_department_deleted": child_department_deleted,
        }

    except Exception:
        raise

def staff_delete(response_data):
    try:
        ids = response_data.get("staff_ids") or []

        if not ids:
            return {
                "step": "staff_delete",
                "deleted": 0,
                "skipped": 0,
                "shared_department_rel_deleted": 0,
                "activity_staff_rel_deleted": 0,
            }

        existing_ids = set(
            TtStaff.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "staff_delete",
                "deleted": 0,
                "skipped": len(ids),
                "shared_department_rel_deleted": 0,
                "activity_staff_rel_deleted": 0,
            }

        with transaction.atomic():
            shared_department_rel_deleted = TtStaffSharedWithDepartment.objects.filter(
                staff_id__in=existing_ids
            ).count()

            activity_staff_rel_deleted = TtActivityStaff.objects.filter(
                staff_id__in=existing_ids
            ).count()

            deleted, _ = TtStaff.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "staff_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "shared_department_rel_deleted": shared_department_rel_deleted,
            "activity_staff_rel_deleted": activity_staff_rel_deleted,
        }

    except Exception:
        raise

def location_delete(response_data):
    try:
        ids = response_data.get("location_ids") or []

        if not ids:
            return {
                "step": "location_delete",
                "deleted": 0,
                "skipped": 0,
                "shared_department_rel_deleted": 0,
                "activity_location_rel_deleted": 0,
            }

        existing_ids = set(
            TtLocation.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "location_delete",
                "deleted": 0,
                "skipped": len(ids),
                "shared_department_rel_deleted": 0,
                "activity_location_rel_deleted": 0,
            }

        with transaction.atomic():
            shared_department_rel_deleted = TtLocationSharedWithDepartment.objects.filter(
                location_id__in=existing_ids
            ).count()

            activity_location_rel_deleted = TtActivityLocation.objects.filter(
                location_id__in=existing_ids
            ).count()

            deleted, _ = TtLocation.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "location_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "shared_department_rel_deleted": shared_department_rel_deleted,
            "activity_location_rel_deleted": activity_location_rel_deleted,
        }

    except Exception:
        raise

def week_pattern_delete(response_data):
    try:
        ids = response_data.get("week_pattern_ids") or []

        if not ids:
            return {
                "step": "week_pattern_delete",
                "deleted": 0,
                "skipped": 0,
                "week_pattern_week_rel_deleted": 0,
            }

        existing_ids = set(
            TtWeekPattern.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "week_pattern_delete",
                "deleted": 0,
                "skipped": len(ids),
                "week_pattern_week_rel_deleted": 0,
            }

        with transaction.atomic():
            week_pattern_week_rel_deleted = TtWeekPattern.week.through.objects.filter(
                week_pattern_id__in=existing_ids
            ).count()

            deleted, _ = TtWeekPattern.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "week_pattern_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "week_pattern_week_rel_deleted": week_pattern_week_rel_deleted,
        }

    except Exception:
        raise

def module_group_delete(response_data):
    try:
        ids = response_data.get("module_group_ids") or []

        if not ids:
            return {
                "step": "module_group_delete",
                "deleted": 0,
                "skipped": 0,
                "pos_module_group_rel_deleted": 0,
                "pos_module_group_module_rel_deleted": 0,
            }

        existing_ids = set(
            TtModuleGroup.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "module_group_delete",
                "deleted": 0,
                "skipped": len(ids),
                "pos_module_group_rel_deleted": 0,
                "pos_module_group_module_rel_deleted": 0,
            }

        with transaction.atomic():
            pos_module_group_ids = list(
                TtPosModuleGroup.objects.filter(
                    module_group_id__in=existing_ids
                ).values_list("id", flat=True)
            )

            pos_module_group_rel_deleted = len(pos_module_group_ids)

            pos_module_group_module_rel_deleted = TtPosModuleGroupModule.objects.filter(
                pos_module_group_id__in=pos_module_group_ids
            ).count() if pos_module_group_ids else 0

            deleted, _ = TtModuleGroup.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "module_group_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "pos_module_group_rel_deleted": pos_module_group_rel_deleted,
            "pos_module_group_module_rel_deleted": pos_module_group_module_rel_deleted,
        }

    except Exception:
        raise

def pos_delete(response_data):
    try:
        ids = response_data.get("pos_ids") or []

        if not ids:
            return {
                "step": "pos_delete",
                "deleted": 0,
                "skipped": 0,
                "pos_module_group_rel_deleted": 0,
                "pos_module_group_module_rel_deleted": 0,
                "student_set_deleted": 0,
            }

        existing_ids = set(
            TtPos.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "pos_delete",
                "deleted": 0,
                "skipped": len(ids),
                "pos_module_group_rel_deleted": 0,
                "pos_module_group_module_rel_deleted": 0,
                "student_set_deleted": 0,
            }

        with transaction.atomic():
            pos_module_group_ids = list(
                TtPosModuleGroup.objects.filter(
                    pos_id__in=existing_ids
                ).values_list("id", flat=True)
            )

            pos_module_group_rel_deleted = len(pos_module_group_ids)

            pos_module_group_module_rel_deleted = TtPosModuleGroupModule.objects.filter(
                pos_module_group_id__in=pos_module_group_ids
            ).count() if pos_module_group_ids else 0

            student_set_deleted = TtStudentSet.objects.filter(
                pos_id__in=existing_ids
            ).count()

            deleted, _ = TtPos.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "pos_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "pos_module_group_rel_deleted": pos_module_group_rel_deleted,
            "pos_module_group_module_rel_deleted": pos_module_group_module_rel_deleted,
            "student_set_deleted": student_set_deleted,
        }

    except Exception:
        raise

def module_delete(response_data):
    try:
        ids = response_data.get("module_ids") or []

        if not ids:
            return {
                "step": "module_delete",
                "deleted": 0,
                "skipped": 0,
                "module_week_rel_deleted": 0,
                "student_set_module_rel_deleted": 0,
                "activity_deleted": 0,
                "activity_template_deleted": 0,
            }

        existing_ids = set(
            TtModule.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "module_delete",
                "deleted": 0,
                "skipped": len(ids),
                "module_week_rel_deleted": 0,
                "student_set_module_rel_deleted": 0,
                "activity_deleted": 0,
                "activity_template_deleted": 0,
            }

        with transaction.atomic():
            module_week_rel_deleted = TtModuleWeek.objects.filter(
                module_id__in=existing_ids
            ).count()

            student_set_module_rel_deleted = TtStudentSetModule.objects.filter(
                module_id__in=existing_ids
            ).count()

            activity_deleted = TtActivity.objects.filter(
                module_id__in=existing_ids
            ).count()

            activity_template_deleted = TtActivityTemplate.objects.filter(
                module_id__in=existing_ids
            ).count()

            deleted, _ = TtModule.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "module_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "module_week_rel_deleted": module_week_rel_deleted,
            "student_set_module_rel_deleted": student_set_module_rel_deleted,
            "activity_deleted": activity_deleted,
            "activity_template_deleted": activity_template_deleted,
        }

    except Exception:
        raise

def student_delete(response_data):
    try:
        ids = response_data.get("student_ids") or []

        if not ids:
            return {
                "step": "student_delete",
                "deleted": 0,
                "skipped": 0,
                "student_academic_term_rel_deleted": 0,
                "student_pos_rel_deleted": 0,
                "student_pos_module_rel_deleted": 0,
                "student_set_student_rel_deleted": 0,
            }

        existing_ids = set(
            TtStudent.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "student_delete",
                "deleted": 0,
                "skipped": len(ids),
                "student_academic_term_rel_deleted": 0,
                "student_pos_rel_deleted": 0,
                "student_pos_module_rel_deleted": 0,
                # "student_set_student_rel_deleted": 0,
            }

        with transaction.atomic():
            student_academic_term_rel_deleted = (
                TtStudentAcademicTerm.objects.filter(
                    student_id__in=existing_ids
                ).count()
            )

            student_pos_rel_deleted = TtStudentPos.objects.filter(
                student_id__in=existing_ids
            ).count()

            student_pos_module_rel_deleted = TtStudentPosModule.objects.filter(
                student_id__in=existing_ids
            ).count()

            # student_set_student_rel_deleted = TtStudentSetStudent.objects.filter(
            #     student_id__in=existing_ids
            # ).count()

            deleted, _ = TtStudent.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "student_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "student_academic_term_rel_deleted": student_academic_term_rel_deleted,
            "student_pos_rel_deleted": student_pos_rel_deleted,
            "student_pos_module_rel_deleted": student_pos_module_rel_deleted,
            # "student_set_student_rel_deleted": student_set_student_rel_deleted,
        }

    except Exception:
        raise

def student_set_delete(response_data):
    try:
        ids = response_data.get("student_set_ids") or []

        if not ids:
            return {
                "step": "student_set_delete",
                "deleted": 0,
                "skipped": 0,
                "student_set_module_rel_deleted": 0,
                "student_set_activity_rel_deleted": 0,
            }

        existing_ids = set(
            TtStudentSet.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "student_set_delete",
                "deleted": 0,
                "skipped": len(ids),
                "student_set_module_rel_deleted": 0,
                "student_set_activity_rel_deleted": 0,
            }

        with transaction.atomic():
            student_set_module_rel_deleted = TtStudentSetModule.objects.filter(
                student_set_id__in=existing_ids
            ).count()

            student_set_activity_rel_deleted = TtStudentSetActivity.objects.filter(
                student_set_id__in=existing_ids
            ).count()

            deleted, _ = TtStudentSet.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "student_set_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "student_set_module_rel_deleted": student_set_module_rel_deleted,
            "student_set_activity_rel_deleted": student_set_activity_rel_deleted,
        }

    except Exception:
        raise

def activity_template_delete(response_data):
    try:
        ids = response_data.get("activity_template_ids") or []

        if not ids:
            return {
                "step": "activity_template_delete",
                "deleted": 0,
                "skipped": 0,
                "activity_template_week_rel_deleted": 0,
                "activity_deleted": 0,
            }

        existing_ids = set(
            TtActivityTemplate.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "activity_template_delete",
                "deleted": 0,
                "skipped": len(ids),
                "activity_template_week_rel_deleted": 0,
                "activity_deleted": 0,
            }

        with transaction.atomic():
            activity_template_week_rel_deleted = TtActivityTemplateWeek.objects.filter(
                activity_template_id__in=existing_ids
            ).count()

            activity_deleted = TtActivity.objects.filter(
                activity_template_id__in=existing_ids
            ).count()

            deleted, _ = TtActivityTemplate.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "activity_template_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "activity_template_week_rel_deleted": activity_template_week_rel_deleted,
            "activity_deleted": activity_deleted,
        }

    except Exception:
        raise

def activity_delete(response_data):
    try:
        ids = response_data.get("activity_ids") or []

        if not ids:
            return {
                "step": "activity_delete",
                "deleted": 0,
                "skipped": 0,
                "activity_staff_rel_deleted": 0,
                "activity_location_rel_deleted": 0,
                "activity_week_rel_deleted": 0,
                "student_set_activity_rel_deleted": 0,
            }

        existing_ids = set(
            TtActivity.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "activity_delete",
                "deleted": 0,
                "skipped": len(ids),
                "activity_staff_rel_deleted": 0,
                "activity_location_rel_deleted": 0,
                "activity_week_rel_deleted": 0,
                "student_set_activity_rel_deleted": 0,
            }

        with transaction.atomic():
            activity_staff_rel_deleted = TtActivityStaff.objects.filter(
                activity_id__in=existing_ids
            ).count()

            activity_location_rel_deleted = TtActivityLocation.objects.filter(
                activity_id__in=existing_ids
            ).count()

            activity_week_rel_deleted = TtActivityWeek.objects.filter(
                activity_id__in=existing_ids
            ).count()

            student_set_activity_rel_deleted = TtStudentSetActivity.objects.filter(
                activity_id__in=existing_ids
            ).count()

            deleted, _ = TtActivity.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "activity_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "activity_staff_rel_deleted": activity_staff_rel_deleted,
            "activity_location_rel_deleted": activity_location_rel_deleted,
            "activity_week_rel_deleted": activity_week_rel_deleted,
            "student_set_activity_rel_deleted": student_set_activity_rel_deleted,
        }

    except Exception:
        raise

def activity_type_delete(response_data):
    try:
        ids = response_data.get("activity_type_ids") or []

        if not ids:
            return {
                "step": "activity_type_delete",
                "deleted": 0,
                "skipped": 0,
                "activity_deleted": 0,
                "activity_template_deleted": 0,
            }

        existing_ids = set(
            TtActivityType.objects.filter(id__in=ids).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "activity_type_delete",
                "deleted": 0,
                "skipped": len(ids),
                "activity_deleted": 0,
                "activity_template_deleted": 0,
            }

        with transaction.atomic():
            activity_deleted = TtActivity.objects.filter(
                activity_type_id__in=existing_ids
            ).count()

            activity_template_deleted = TtActivityTemplate.objects.filter(
                activity_type_id__in=existing_ids
            ).count()

            deleted, _ = TtActivityType.objects.filter(id__in=existing_ids).delete()

        return {
            "step": "activity_type_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "activity_deleted": activity_deleted,
            "activity_template_deleted": activity_template_deleted,
        }

    except Exception:
        raise

def booking_delete(response_data):
    try:
        ids = response_data.get("activity_ids") or []

        if not ids:
            return {
                "step": "booking_delete",
                "deleted": 0,
                "skipped": 0,
                "activity_week_rel_deleted": 0,
            }

        existing_ids = set(
            TtActivity.objects.filter(
                id__in=ids,
                is_booking=1,
            ).values_list("id", flat=True)
        )

        if not existing_ids:
            return {
                "step": "booking_delete",
                "deleted": 0,
                "skipped": len(ids),
                "activity_week_rel_deleted": 0,
            }

        with transaction.atomic():
            activity_week_rel_deleted = TtActivityWeek.objects.filter(
                activity_id__in=existing_ids
            ).count()

            deleted, _ = TtActivity.objects.filter(
                id__in=existing_ids,
                is_booking=1,
            ).delete()

        return {
            "step": "booking_delete",
            "deleted": deleted,
            "skipped": len(ids) - len(existing_ids),
            "activity_week_rel_deleted": activity_week_rel_deleted,
        }

    except Exception:
        raise

# ============================
# OTHERS
# ============================

def schedule(response_data):
    try:
        payload = normalize_payload(response_data, "activity")
        delete_activity_ids = response_data.get("delete_activity") or []

        if not payload and not delete_activity_ids:
            return {
                "step": "schedule_create",
                "updated": 0,
                "skipped": 0,
                "deleted": 0,
                "staff_rel": 0,
                "location_rel": 0,
                "week_rel": 0,
            }

        activity_ids = [
            row.get("id")
            for row in payload
            if row.get("id")
        ]

        updated = 0
        skipped = 0
        deleted = 0
        staff_rel = 0
        location_rel = 0
        week_rel = 0

        with transaction.atomic():
            if delete_activity_ids:
                deleted, _ = TtActivity.objects.filter(
                    id__in=delete_activity_ids
                ).delete()

            activity_map = {
                activity.id: activity
                for activity in TtActivity.objects.filter(id__in=activity_ids)
            }

            for row in payload:
                activity_id = row.get("id")

                if not activity_id or activity_id not in activity_map:
                    skipped += 1
                    continue

                activity = activity_map[activity_id]

                update_fields = []

                activity.scheduled = 1
                update_fields.append("scheduled")

                if "week_pattern_id" in row:
                    activity.week_pattern_id = row.get("week_pattern_id")
                    update_fields.append("week_pattern_id")

                if "scheduled_day" in row:
                    activity.scheduled_day = row.get("scheduled_day")
                    update_fields.append("scheduled_day")

                if "scheduled_start_time" in row:
                    activity.scheduled_start_time = row.get("scheduled_start_time")
                    update_fields.append("scheduled_start_time")

                if "scheduled_start_slot" in row:
                    activity.scheduled_start_slot = row.get("scheduled_start_slot")
                    update_fields.append("scheduled_start_slot")

                if "code" in row:
                    activity.code = row.get("code")
                    update_fields.append("code")

                if "name" in row:
                    activity.name = row.get("name")
                    update_fields.append("name")

                if "is_variant" in row:
                    activity.is_variant = row.get("is_variant")
                    update_fields.append("is_variant")

                if "variant_parent_id" in row:
                    activity.variant_parent_id = row.get("variant_parent_id")
                    update_fields.append("variant_parent_id")

                activity.save(update_fields=list(set(update_fields)))

                if "staff_ids" in row:
                    staff_ids = row.get("staff_ids") or []
                    activity.staff.set(staff_ids)
                    staff_rel += len(staff_ids)

                if "location_ids" in row:
                    location_ids = row.get("location_ids") or []
                    activity.location.set(location_ids)
                    location_rel += len(location_ids)

                if "week_ids" in row:
                    week_ids = row.get("week_ids") or []
                    activity.week.set(week_ids)
                    week_rel += len(week_ids)

                updated += 1

        return {
            "step": "schedule_create",
            "updated": updated,
            "skipped": skipped,
            "deleted": deleted,
            "staff_rel": staff_rel,
            "location_rel": location_rel,
            "week_rel": week_rel,
        }

    except Exception:
        raise

def unschedule(response_data):
    try:
        payload = normalize_payload(response_data, "activity")

        activity_map = {}
        for row in payload:
            activity_id = row.get("activity_id")
            if not activity_id:
                continue

            activity_map[activity_id] = {
                "scheduled": row.get("scheduled", 0),
                "scheduled_day": row.get("scheduled_day"),
                "scheduled_start_time": row.get("scheduled_start_time"),
                "scheduled_start_slot": row.get("scheduled_start_slot"),
                "staff": row.get("staff", []) or [],
                "location": row.get("location", []) or [],
            }

        activity_ids = list(activity_map.keys())

        if not activity_ids:
            return {
                "step": "unschedule",
                "updated": 0,
                "skipped": len(payload),
                "staff_removed": 0,
                "location_removed": 0,
            }

        with transaction.atomic():
            activities = list(TtActivity.objects.filter(id__in=activity_ids))
            valid_ids = [act.id for act in activities]

            for activity in activities:
                row = activity_map.get(activity.id, {})
                activity.scheduled = row.get("scheduled", 0)
                activity.scheduled_day = row.get("scheduled_day")
                activity.scheduled_start_time = row.get("scheduled_start_time")
                activity.scheduled_start_slot = row.get("scheduled_start_slot")

            if activities:
                TtActivity.objects.bulk_update(
                    activities,
                    ["scheduled", "scheduled_day", "scheduled_start_time", "scheduled_start_slot"]
                )

            staff_removed, _ = TtActivityStaff.objects.filter(
                activity_id__in=valid_ids
            ).delete()

            location_removed, _ = TtActivityLocation.objects.filter(
                activity_id__in=valid_ids
            ).delete()

        return {
            "step": "unschedule",
            "updated": len(valid_ids),
            "skipped": len(activity_ids) - len(valid_ids),
            "staff_removed": staff_removed,
            "location_removed": location_removed,
        }

    except Exception:
        raise

def swap(response_data):
    try:
        payload = normalize_payload(response_data, "activity")
        delete_activity_ids = response_data.get("delete_activity") or []

        if not payload and not delete_activity_ids:
            return {
                "step": "swap",
                "updated": 0,
                "skipped": 0,
                "deleted": 0,
                "staff_rel": 0,
                "location_rel": 0,
                "week_rel": 0,
            }

        activity_ids = [
            row.get("id")
            for row in payload
            if row.get("id")
        ]

        updated = 0
        skipped = 0
        deleted = 0
        staff_rel = 0
        location_rel = 0
        week_rel = 0

        with transaction.atomic():
            if delete_activity_ids:
                deleted, _ = TtActivity.objects.filter(
                    id__in=delete_activity_ids
                ).delete()

            activity_map = {
                activity.id: activity
                for activity in TtActivity.objects.filter(id__in=activity_ids)
            }

            for row in payload:
                activity_id = row.get("id")

                if not activity_id or activity_id not in activity_map:
                    skipped += 1
                    continue

                activity = activity_map[activity_id]

                update_fields = []

                if "code" in row:
                    activity.code = row.get("code")
                    update_fields.append("code")

                if "name" in row:
                    activity.name = row.get("name")
                    update_fields.append("name")

                if "is_variant" in row:
                    activity.is_variant = row.get("is_variant")
                    update_fields.append("is_variant")

                if "variant_parent_id" in row:
                    activity.variant_parent_id = row.get("variant_parent_id")
                    update_fields.append("variant_parent_id")

                if update_fields:
                    activity.save(update_fields=list(set(update_fields)))

                if "staff_ids" in row:
                    staff_ids = row.get("staff_ids") or []
                    activity.staff.set(staff_ids)
                    staff_rel += len(staff_ids)

                if "location_ids" in row:
                    location_ids = row.get("location_ids") or []
                    activity.location.set(location_ids)
                    location_rel += len(location_ids)

                if "week_ids" in row:
                    week_ids = row.get("week_ids") or []
                    activity.week.set(week_ids)
                    week_rel += len(week_ids)

                updated += 1

        return {
            "step": "swap",
            "updated": updated,
            "skipped": skipped,
            "deleted": deleted,
            "staff_rel": staff_rel,
            "location_rel": location_rel,
            "week_rel": week_rel,
        }

    except Exception:
        raise

def pos_update_module(response_data):
    try:
        payload = normalize_payload(response_data, "modules")

        if not payload:
            return {
                "step": "pos_update_module",
                "updated_groups": 0,
                "module_rel_created": 0,
                "skipped": 0,
            }

        lookup_pairs = []
        for row in payload:
            pos_id = row.get("pos_id")
            module_group_id = row.get("module_group")
            if pos_id and module_group_id:
                lookup_pairs.append((pos_id, module_group_id))

        if not lookup_pairs:
            return {
                "step": "pos_update_module",
                "updated_groups": 0,
                "module_rel_created": 0,
                "skipped": len(payload),
            }

        relation_rows = []
        valid_pos_module_group_ids = []
        skipped_count = 0

        with transaction.atomic():
            pos_ids = list({pair[0] for pair in lookup_pairs})
            module_group_ids = list({pair[1] for pair in lookup_pairs})

            pos_module_groups = TtPosModuleGroup.objects.filter(
                pos_id__in=pos_ids,
                module_group_id__in=module_group_ids,
            )

            pmg_map = {
                (obj.pos_id, obj.module_group_id): obj.id
                for obj in pos_module_groups
            }

            for row in payload:
                pos_id = row.get("pos_id")
                module_group_id = row.get("module_group")
                module_ids = row.get("module_id") or []

                if not pos_id or not module_group_id:
                    skipped_count += 1
                    continue

                pos_module_group_id = pmg_map.get((pos_id, module_group_id))
                if not pos_module_group_id:
                    skipped_count += 1
                    continue

                valid_pos_module_group_ids.append(pos_module_group_id)

                for module_id in module_ids:
                    if module_id:
                        relation_rows.append(
                            TtPosModuleGroupModule(
                                pos_module_group_id=pos_module_group_id,
                                module_id=module_id,
                            )
                        )

            valid_pos_module_group_ids = list(set(valid_pos_module_group_ids))

            if valid_pos_module_group_ids:
                TtPosModuleGroupModule.objects.filter(
                    pos_module_group_id__in=valid_pos_module_group_ids
                ).delete()

            if relation_rows:
                TtPosModuleGroupModule.objects.bulk_create(
                    relation_rows,
                    batch_size=500,
                    ignore_conflicts=True,
                )

        return {
            "step": "pos_update_module",
            "updated_groups": len(valid_pos_module_group_ids),
            "module_rel_created": len(relation_rows),
            "skipped": skipped_count,
        }

    except Exception:
        raise

def jta_split(response_data):
    try:
        payload = normalize_payload(response_data, "activity_update")
        delete_id = response_data.get("activity_delete")
        update_child_data = response_data.get("child_activity_update")

        if not payload and not delete_id:
            return {
                "step": "jta_split",
                "updated": 0,
                "deleted": 0,
                "skipped": 0,
                "delete_skipped": 0,
            }

        activity_ids = [row.get("id") for row in payload if row.get("id")]
        existing_ids = set(
            TtActivity.objects.filter(id__in=activity_ids).values_list("id", flat=True)
        )

        updated = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                activity_id = row.get("id")
                if not activity_id or activity_id not in existing_ids:
                    skipped += 1
                    continue

                TtActivity.objects.filter(id=activity_id).update(
                    is_jta=row.get("is_jta", 0),
                    jta_parent_id=row.get("jta_parent_id"),
                )
                updated += 1
            if update_child_data:
                child_activity = TtActivity.objects.filter(id=update_child_data["id"]).first()
                if child_activity:
                    child_activity.name = update_child_data["name"]
                    child_activity.desc = update_child_data["desc"]
                    child_activity.module_id = update_child_data["module_id"]
                    child_activity.activity_template_id = update_child_data["activity_template_id"]
                    child_activity.student_set.set(update_child_data["student_set"])
                    child_activity.save()
                    updated += 1
                else:
                    skipped += 1
            deleted = 0
            delete_skipped = 0

            if delete_id:
                deleted, delete_skipped = bulk_delete_by_ids(TtActivity, [delete_id])

        return {
            "step": "jta_split",
            "updated": updated,
            "deleted": deleted,
            "skipped": skipped,
            "delete_skipped": delete_skipped,
        }

    except Exception:
        raise

def activity_template_allocator(response_data):
    try:
        payload = normalize_payload(response_data, "allocator")

        if not payload:
            return {
                "step": "activity_template_allocator",
                "updated_student_sets": 0,
                "activity_rel": 0,
                "skipped": 0,
            }

        student_set_ids = [row.get("student_set") for row in payload if row.get("student_set")]
        existing_student_set_ids = set(
            TtStudentSet.objects.filter(id__in=student_set_ids).values_list("id", flat=True)
        )

        all_activity_ids = []
        for row in payload:
            all_activity_ids.extend(row.get("activity", []) or [])

        existing_activity_ids = set(
            TtActivity.objects.filter(id__in=all_activity_ids).values_list("id", flat=True)
        )

        relation_map = {}
        updated_student_sets = 0
        skipped = 0

        with transaction.atomic():
            for row in payload:
                student_set_id = row.get("student_set")
                if not student_set_id or student_set_id not in existing_student_set_ids:
                    skipped += 1
                    continue

                relation_map[student_set_id] = [
                    activity_id
                    for activity_id in (row.get("activity", []) or [])
                    if activity_id in existing_activity_ids
                ]
                updated_student_sets += 1

            activity_rel = replace_relation_rows(
                TtStudentSetActivity,
                "student_set",
                "activity",
                relation_map,
            ) if relation_map else 0

        return {
            "step": "activity_template_allocator",
            "updated_student_sets": updated_student_sets,
            "activity_rel": activity_rel,
            "skipped": skipped,
        }

    except Exception:
        raise

# def seed_booking_test_records():
#     with transaction.atomic():
#         # activity_type id=3
#         TtActivityType.objects.update_or_create(
#             id=3,
#             defaults={
#                 "code": "booking",
#                 "name": "Booking",
#                 "desc": "Booking Activity Type",
#                 "booking": 1,
#                 "contact": 0,
#                 "cover": 0,
#                 "color": "#000000",
#                 "status": 1,
#             }
#         )

#         # week ids
#         week_rows = [
#             {"id": 74, "week": 1, "start_date": "2026-01-05"},
#             {"id": 75, "week": 2, "start_date": "2026-01-12"},
#             {"id": 76, "week": 3, "start_date": "2026-01-19"},
#             {"id": 77, "week": 4, "start_date": "2026-01-26"},
#             {"id": 78, "week": 5, "start_date": "2026-02-02"},
#         ]

#         for row in week_rows:
#             TtWeek.objects.update_or_create(
#                 id=row["id"],
#                 defaults={
#                     "week": row["week"],
#                     "start_date": row["start_date"],
#                 }
#             )

#         # staff id=19
#         TtStaff.objects.update_or_create(
#             id=19,
#             defaults={
#                 "code": "staff19",
#                 "name": "Andy",
#                 "email": "andy@test.com",
#                 "desc": "Seed staff for booking_schedule",
#                 "department_id": None,
#                 "is_part_time": 0,
#                 "maximum_period": 144,
#                 "contract_period": 144,
#                 "shared_with_all_department": 1,
#                 "status": 1,
#             }
#         )

#         # staff id=14
#         TtStaff.objects.update_or_create(
#             id=14,
#             defaults={
#                 "code": "staff14",
#                 "name": "Swap Staff",
#                 "email": "swapstaff@test.com",
#                 "desc": "Seed staff for booking_swap",
#                 "department_id": None,
#                 "is_part_time": 0,
#                 "maximum_period": 144,
#                 "contract_period": 144,
#                 "shared_with_all_department": 1,
#                 "status": 1,
#             }
#         )

#         # activity id=1143
#         TtActivity.objects.update_or_create(
#             id=1143,
#             defaults={
#                 "code": "BookingOld",
#                 "name": "Booking Old",
#                 "desc": "Before booking update",
#                 "activity_type_id": 3,
#                 "duration": 120,
#                 "slot_required": 4,
#                 "planned_size": 30,
#                 "is_jta": 0,
#                 "is_variant": 0,
#                 "scheduled": 0,
#                 "is_booking": 1,
#                 "status": 1,
#             }
#         )

#         # activity id=1149
#         TtActivity.objects.update_or_create(
#             id=1149,
#             defaults={
#                 "code": "MV0000000325-OLD",
#                 "name": "Booking Schedule Old",
#                 "desc": "Before booking schedule",
#                 "activity_type_id": 3,
#                 "department_id": None,
#                 "duration": 60,
#                 "slot_required": 2,
#                 "planned_size": 1,
#                 "scheduled_start_time": None,
#                 "scheduled_start_slot": None,
#                 "scheduled_day": None,
#                 "is_jta": 0,
#                 "is_variant": 0,
#                 "scheduled": 0,
#                 "is_booking": 1,
#                 "status": 1,
#             }
#         )

#         # deletes activity ids [1150, 1151, 1152]
#         delete_activity_ids = [1150, 1151, 1152]

#         for activity_id in delete_activity_ids:
#             TtActivity.objects.update_or_create(
#                 id=activity_id,
#                 defaults={
#                     "code": f"BookingDelete{activity_id}",
#                     "name": f"Booking Delete {activity_id}",
#                     "desc": "Seed record for booking_delete test",
#                     "activity_type_id": 3,
#                     "duration": 120,
#                     "slot_required": 4,
#                     "planned_size": 30,
#                     "is_jta": 0,
#                     "is_variant": 0,
#                     "scheduled": 0,
#                     "is_booking": 1,
#                     "status": 1,
#                 }
#             )



TEST = False

if TEST:

    TtActivityStaff.objects.update_or_create(
        activity_id=1149,
        staff_id=19,
        defaults={}
    )

    header_data = {
        "method": "booking_swap"
    }

    response_data = {
        "session_id": "dev",
        "activity": {
            "id": 1149,
            "staff_ids": [14],
            "location_ids": [],
        },
        "request_id": "a58706fb-cd03-4d42-b672-ddfb25348c02",
    }

    result = swap(response_data)
    append_log(response_data, header_data, result)

    print("result =", result)
    exit()


while True:
    # poll(1.0) means 1 second timeout, so everytime call poll within 1 second no data, will have another poll call again
    # if put poll() without timeout, will keep wait until got response only call another poll. (NOT RECOMMEND NO TIMEOUT, WILL CAUSE BUSY LOOP)
    response_data = None
    header_data = {}
    msg_str = None
    msg = None

    try:
        msg = consumer.poll(1.0)
        # print(msg)
        if msg is None:
            continue
        if msg.error():
            # print("Error:", msg.error())
            continue

        # print(msg.headers())
        # Returns [('method', b'schedule')]
        raw_headers = msg.headers()

        # Convert headers to a readable dictionary
        if raw_headers:
            # Headers are bytes, so must decode them
            header_data = {key: value.decode('utf-8') for key, value in raw_headers}

        msg_str = msg.value().decode()
        # print(header_data)
        # print("Received:", msg.value().decode())

        # print(datetime.datetime.now())
        response_data = json.loads(msg_str)
        # kafka_log = KafkaLog.update_response(topic, response_data,header_data)

        result = None

        if "method" in header_data:
            method = header_data["method"]
            match method:
                # create methods
                case "department_create":
                    result = department_create(response_data)
                case "location_create":
                    result = location_create(response_data)
                case "week_pattern_create":
                    result = week_pattern_create(response_data)
                case "academic_term_create":
                    result = academic_term_create(response_data)
                case "module_group_create":
                    result = module_group_create(response_data)
                case "pos_create":
                    result = pos_create(response_data)
                case "activity_type_create":
                    result = activity_type_create(response_data)
                case "module_create":
                    result = module_create(response_data)
                case "activity_template_create":
                    result = activity_template_create(response_data)
                case "student_set_create" | "allocation":
                    result = student_set_create(response_data)
                case "student_create":
                    result = student_create(response_data)
                case "staff_create":
                    result = staff_create(response_data)
                case "activity_create" | "activity_generate":
                    result = activity_create(response_data)
                case "variant_create":
                    result = variant_create(response_data)
                case "jta_create":
                    result = jta_create(response_data)
                case "booking_create":
                    result = booking_create(response_data)
                # update methods
                case "department_update":
                    result = department_update(response_data)
                case "location_update":
                    result = location_update(response_data)
                case "staff_update":
                    result = staff_update(response_data)
                case "week_pattern_update":
                    result = week_pattern_update(response_data)
                case "academic_term_update":
                    result = academic_term_update(response_data)
                case "pos_update":
                    result = pos_update(response_data)
                case "module_group_update":
                    result = module_group_update(response_data)
                case "module_update":
                    result = module_update(response_data)
                case "activity_type_update":
                    result = activity_type_update(response_data)
                case "activity_template_update":
                    result = activity_template_update(response_data)
                case "student_set_update":
                    result = student_set_update(response_data)
                case "student_update":
                    result = student_update(response_data)
                case "activity_update":
                    result = activity_update(response_data)
                case "booking_update":
                    result = booking_update(response_data)
                # delete methods
                case "academic_term_delete":
                    result = academic_term_delete(response_data)
                case "department_delete":
                    result = department_delete(response_data)
                case "staff_delete":
                    result = staff_delete(response_data)
                case "location_delete":
                    result = location_delete(response_data)
                case "week_pattern_delete":
                    result = week_pattern_delete(response_data)
                case "module_group_delete":
                    result = module_group_delete(response_data)
                case "pos_delete":
                    result = pos_delete(response_data)
                case "module_delete":
                    result = module_delete(response_data)
                case "student_delete":
                    result = student_delete(response_data)
                case "student_set_delete":
                    result = student_set_delete(response_data)
                case "activity_template_delete":
                    result = activity_template_delete(response_data)
                case "activity_delete":
                    result = activity_delete(response_data)
                case "activity_type_delete":
                    result = activity_type_delete(response_data)
                case "booking_delete":
                    result = booking_delete(response_data)
                # other operations
                case "schedule" | "booking_schedule":
                    result = schedule(response_data)
                case "unschedule":
                    result = unschedule(response_data)
                case "swap" | "booking_swap":
                    result = swap(response_data)
                case "pos_update_module":
                    result = pos_update_module(response_data)
                case "jta_split":
                    result = jta_split(response_data)
                case "activity_template_allocator":
                    result = activity_template_allocator(response_data)
                case _:
                    print(f"Unknown Kafka method: {method}")
                    continue
        else:
            result = {
                "step": "missing_method",
                "success": False,
                "error": "Kafka header method is missing",
            }
            
        if result:
            append_log(response_data, header_data, result)
        # make it manual commit, so if got any connection problem in db, wont skip this record
        consumer.commit(msg)
    except (OperationalError, InterfaceError, TimeoutError) as e:
        # 2026-05-07 try to prevent the error from timeout, if got any timeout/connection error, sleep for 5 second and try again
        e_details = get_exception_detail(e)

        log_critical_error(
            user_id=None,
            descr=e_details["descr"],
            url=e_details["url"],
            trace=e_details["trace"],
        )

        error_result = {
            "step": "kafka_connection_exception",
            "success": False,
            "error": str(e),
            "descr": e_details["descr"],
            "url": e_details["url"],
            "trace": e_details["trace"],
        }

        try:
            if isinstance(response_data, dict):
                append_log(response_data, header_data, error_result)
            else:
                append_log(
                    {
                        "request_id": None,
                        "raw_data": msg_str,
                    },
                    header_data,
                    error_result,
                )
        except Exception as log_error:
            print("failed to insert kafka connection error log:", log_error)
        
        time.sleep(5)
        continue

    except Exception as e:
        e_details = get_exception_detail(e)
        
        log_critical_error(
            user_id=None,
            descr=e_details["descr"],
            url=e_details["url"],
            trace=e_details["trace"],
        )
        # if is error due to code is not working, still need to commit, if not will keep loop in same problem
        error_result = {
            "step": "kafka_exception",
            "success": False,
            "error": str(e),
            "descr": e_details["descr"],
            "url": e_details["url"],
            "trace": e_details["trace"],
        }

        try:
            if isinstance(response_data, dict):
                append_log(response_data, header_data, error_result)
            else:
                append_log(
                    {
                        "request_id": None,
                        "raw_data": msg_str,
                    },
                    header_data,
                    error_result,
                )
        except Exception as log_error:
            print("failed to insert kafka error log:", log_error)

        if msg:
            consumer.commit(msg)
        continue
