import datetime
import json
import logging
import os
import re
import uuid
from collections import defaultdict
from contextlib import contextmanager
from json import JSONDecodeError

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

# 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, PostCommitDelivery, TtActivity, TtSetting, TtActivityStaff, TtStaffResourceMap, TtStaff, TtLocation, \
    TtActivityLocation, TtLocationResourceMap, TtActivityWeek, TtStudentSet, TtStudentSetResourceMap, \
    TtStudentSetActivity, TtAcademicTermWeek, AuditTrail, AuditTrailDetails, User
from django.conf import settings
from api.utils import log_critical_error, get_exception_detail, bitwise_calculator_from_string, bulk_sync_to_redis, \
    push_websocket_notification, remove_variant_activity_name_range, \
    format_variant_activity_name_week_pattern_to_ranges, convert_to_redis_week, recalculate_resource_map, \
    get_resource_map_redis_data
from backend.redis_client import redis_client
from api.translation import __
from backend.kafka import send_request
from api.helper.activity_sequencing_helper import helper_build_update_redis_sequencing_data
from api.helper.resource_map_helper import helper_recalculate_and_get_redis_resource_map
from api.services.integration.quarantine import quarantine_engine_response
from api.services.integration.receipts import EngineResponseConflict, claim_engine_response
from api.services.integration.context import IntegrationContext, reset_context, set_context
from api.services.integration.tracking import mutation_tracking, track_activity_ids

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': 'earliest', # 'earliest': consume from beginning, 'latest': consume only new messages
    'enable.auto.commit': False,
    'enable.auto.offset.store': False,
    'max.poll.interval.ms': settings.KAFKA_CONFIG["MAX_POLL_INTERVAL_MS"],
    'session.timeout.ms': settings.KAFKA_CONFIG["SESSION_TIMEOUT_MS"],
    'heartbeat.interval.ms': settings.KAFKA_CONFIG["HEARTBEAT_INTERVAL_MS"],
    'fetch.message.max.bytes': settings.KAFKA_CONFIG["FETCH_MESSAGE_MAX_BYTES"],
}

logger = logging.getLogger("kafka_consumer.tt_response")

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_swap_get_resources_endpoint = str(settings.WEBSOCKET_CONFIG["SWAP_GET_RESOURCES_ENDPOINT"])
websocket_swap_endpoint = str(settings.WEBSOCKET_CONFIG["SWAP_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
websocket_swap_get_resources_url = websocket_host + ":" + websocket_port + "/" + websocket_swap_get_resources_endpoint
websocket_swap_url = websocket_host + ":" + websocket_port + "/" + websocket_swap_endpoint

topic = settings.KAFKA_CONFIG["RESPONSE_TOPIC"]
timeout = settings.KAFKA_CONFIG["TIMEOUT"]
microservices_tt_topic = settings.KAFKA_CONFIG["MICROSERVICES_TT_TOPIC"]
request_topic = settings.KAFKA_CONFIG["REQUEST_TOPIC"]
# Consumer construction is intentionally lazy so importing the atomic apply
# service (for simulation/tests) does not start network activity.
consumer = None


def response_consumer():
    global consumer
    if consumer is None:
        consumer = Consumer(conf)
        consumer.subscribe([topic], on_assign=_on_assign, on_revoke=_on_revoke)
    return consumer


def _on_assign(active_consumer, partitions):
    logger.info("engine_response_partitions_assigned", extra={"partitions": [str(p) for p in partitions]})
    active_consumer.assign(partitions)


def _on_revoke(active_consumer, partitions):
    # Uncommitted work is intentionally not committed here. The new owner will
    # redeliver it and the applied-response receipt prevents double mutation.
    logger.warning("engine_response_partitions_revoked", extra={"partitions": [str(p) for p in partitions]})
    active_consumer.unassign()

setting_params = {
    "minute_per_slot",
    "slot_per_week",
    "slot_per_day"
}
tt_settings = TtSetting.get_multiple_setting(setting_params)
slot_per_week = int(tt_settings["slot_per_week"])
slot_per_day = int(tt_settings["slot_per_day"])
minute_per_slot = int(tt_settings["minute_per_slot"])

def get_match_key(activity):
    """
    helper function to get variant staff,location,and slot for checking
    """
    staff_ids = tuple(sorted([s.id for s in activity.staff.all()]))
    location_ids = tuple(sorted([l.id for l in activity.location.all()]))
    start_slot = activity.scheduled_start_slot
    slot_required = activity.slot_required
    return (location_ids, staff_ids, start_slot, slot_required)

def _apply_schedule(response_data,kafka_log):
    """
    Sample Response
    {
        'session_id': 'user123-abc456',
        'request_id': 'req789',
        'status': 'success',
        'schedule': '[
            {"activity":"5","teaching_staff":"[1]","start_slot":"352","location":"[1,2]"},
            {"activity":"2","teaching_staff":"[1]","start_slot":"22","location":"[2]"},
            {"activity":"9","teaching_staff":"[1]","start_slot":"364","location":"[2]"},
            {"activity":"3","teaching_staff":"[1]","start_slot":"64","location":"[2]"},
            {"activity":"4","teaching_staff":"[1]","start_slot":"406","location":"[2]"},
            {"activity":"7","teaching_staff":"[1]","start_slot":"412","location":"[2]"},
            {"activity":"8","teaching_staff":"[1]","start_slot":"448","location":"[2]"},
            {"activity":"1","teaching_staff":"[1]","start_slot":"454","location":"[2]"},
            {"activity":"10","teaching_staff":"[1]","start_slot":"458","location":"[2]"},
            {"activity":"6","teaching_staff":"[]","start_slot":"370","location":"[2]"}
        ]'
    }
    2025-12-12 exact response
    {
        "session_id": "dev",
        "request_id": "4998764b-1555-4faa-a41e-dd7f91b204c7",
        "status": "success",
        "schedule": "[{\"activity\":\"6\",\"teaching_staff\":\"[1]\",\"start_slot\":\"582\",\"location\":\"[1]\"}]"
    }
    """
    try:
        update_activity_objs = []
        insert_activity_staff_objs = []
        insert_activity_location_objs = []
        need_delete_variant_activity_ids = []
        sequencing_activity_ids = []
        staff_with_affected_week = {}
        location_with_affected_week = {}
        student_set_with_affected_week = {}

        # key will be activity id and value will be week_ids
        activity_week_microservices_relation = {}

        # param used to update redis resource map
        scheduled_staff_ids = set()
        scheduled_location_ids = set()
        scheduled_student_set_ids = set()
        scheduled_variant_activity_ids = []
        activity_ids = []
        success_activity_ids = [] # got start_slot only add in
        all_staff_ids = set()
        all_location_ids = set()

        redis_data = {
            "update": {},
            "delete": {}
        }
        redis_staff_table = "staff"
        redis_location_table = "location"
        redis_student_set_table = "student_set"
        redis_activity_table = "activity"
        kafka_data = []

        old_new_data = {
            "old_data": {},
            "new_data": {}
        }
        # got kafka_log only do action, else need double check why no request but receive response
        if response_data['status'] == "success" and kafka_log:
            # 2025-12-12 the exact response for the "schedule" is json_encode 1 more time, so need decode again for the "schedule"
            schedule = json.loads(response_data['schedule'])
            for item in schedule:
                activity_ids.append(item["activity"])
                if item.get("teaching_staff"):
                    all_staff_ids.update(json.loads(item["teaching_staff"]))
                if item.get("location"):
                    all_location_ids.update(json.loads(item["location"]))

            # used 1 query get all activity data
            activities = (
                TtActivity.objects.filter(id__in=activity_ids)
                .select_related("week_pattern")
                .prefetch_related("week", "week_pattern__week", "staff", "location", "student_set")
            )

            student_set_ids = list(activities.values_list("student_set__id", flat=True).distinct())
            activities_by_id = {act.id: act for act in activities}

            # loop the schedule
            success_count = 0
            fail_count = 0
            for val in schedule:
                # at the end will get those affected variant ids, to check code and name to change
                activity_id = int(val['activity'])
                activity = activities_by_id[activity_id]
                # if got week_pattern_id, will get from week_pattern table, else will used self week_pattern
                if activity.week_pattern_id:
                    week_ids = list(aw.id for aw in activity.week_pattern.week.all())
                else:
                    week_ids = list(aw.id for aw in activity.week.all())

                # 2025-12-12 the data response in teaching_staff and location still have 1 more level of json encode, so still need decode to get the array
                """
                2025-12-12 start_slot is in format include previous week, for example
                start_slot return value is 582, after modulo by slot_per_week(336), will get 246, so this activity will schedule to slot 246, and follow the week pattern
    
                2026-02-10
                only no slot wont do schedule
                """
                if val.get('start_slot'):
                    success_activity_ids.append(activity_id)
                    
                    old_new_data["old_data"].setdefault(activity_id, {})
                    old_new_data["old_data"][activity_id]["scheduled_start_slot"] = activity.scheduled_start_slot
                    old_new_data["old_data"][activity_id]["scheduled_day"] = activity.scheduled_day
                    old_new_data["old_data"][activity_id]["scheduled_start_time"] = activity.scheduled_start_time.isoformat() if activity.scheduled_start_time else None
                    old_new_data["old_data"][activity_id]["staff_ids"] = sorted([s.id for s in activity.staff.all()])
                    old_new_data["old_data"][activity_id]["location_ids"] = sorted([l.id for l in activity.location.all()])
                    # if is scheduled activity, need delete current staff and location first, cause engine possible give different resources when is scheduled activity
                    # 2026-07-23 normally will have scheduled activity pass in is doing change slot, so will always schedule for 1 activity only
                    if activity.scheduled:
                        delete_activity_staff_ids = []
                        delete_activity_location_ids = []
                        current_scheduled_staff_ids = [scheduled_staff.id for scheduled_staff in activity.staff.all()]
                        for current_scheduled_staff_id in current_scheduled_staff_ids:
                            delete_activity_staff_ids.append(current_scheduled_staff_id)
                            staff_with_affected_week.setdefault(current_scheduled_staff_id, set()).update(week_ids)
                            # even is delete staff also need add to this scheduled_staff_ids, so no need add extra handle for redis
                            scheduled_staff_ids.add(current_scheduled_staff_id)

                        current_scheduled_location_ids = [scheduled_location.id for scheduled_location in activity.location.all()]
                        for current_scheduled_location_id in current_scheduled_location_ids:
                            delete_activity_location_ids.append(current_scheduled_location_id)
                            location_with_affected_week.setdefault(current_scheduled_location_id, set()).update(week_ids)
                            # even is delete location also need add to this scheduled_location_ids, so no need add extra handle for redis
                            scheduled_location_ids.add(current_scheduled_location_id)

                        if delete_activity_staff_ids:
                            TtActivityStaff.objects.filter(activity_id=activity_id,staff_id__in=delete_activity_staff_ids).delete()

                        if delete_activity_location_ids:
                            TtActivityLocation.objects.filter(activity_id=activity_id,location_id__in=delete_activity_location_ids).delete()

                    # group variant activity later will used to check variant merge
                    if activity.is_variant == 1:
                        scheduled_variant_activity_ids.append(activity_id)
                    raw_staff = val.get('teaching_staff')
                    staff_ids = json.loads(raw_staff) if raw_staff else []

                    raw_location = val.get('location')
                    location_ids = json.loads(raw_location) if raw_location else []

                    student_set_ids = list(ss.id for ss in activity.student_set.all())

                    response_start_slot = int(val['start_slot'])

                    start_slot = response_start_slot % slot_per_week
                    # calculate how many slot taken for this activity
                    slot_taken = activity.duration // minute_per_slot
                    # 0 is Monday 6 is Sunday
                    day = start_slot // slot_per_day
                    # * 60 means need convert to second
                    time_in_second = (start_slot % slot_per_day) * minute_per_slot * 60
                    hours = time_in_second // 3600
                    minutes = (time_in_second % 3600) // 60
                    time = datetime.time(hour=hours, minute=minutes)

                    # data need to update put in a object first, at the end will used bulk update to perform, so can reduce the query
                    activity_obj = TtActivity(id=activity_id)
                    activity_obj.scheduled_start_time = time
                    activity_obj.scheduled_day = day
                    activity_obj.scheduled_start_slot = start_slot
                    activity_obj.scheduled = 1
                    update_activity_objs.append(activity_obj)
                    
                    old_new_data["new_data"].setdefault(activity_id, {})
                    old_new_data["new_data"][activity_id]["scheduled_start_slot"] = start_slot
                    old_new_data["new_data"][activity_id]["scheduled_day"] = day
                    old_new_data["new_data"][activity_id]["scheduled_start_time"] = time.isoformat()
                    old_new_data["new_data"][activity_id]["staff_ids"] = sorted(staff_ids)
                    old_new_data["new_data"][activity_id]["location_ids"] = sorted(location_ids)

                    redis_data["update"].setdefault(redis_activity_table, []).append({
                        "id": activity_id,
                        "scheduled": 1,
                        "scheduled_start_slot": start_slot,
                    })

                    if staff_ids:
                        for staff_id in staff_ids:
                            activity_staff_obj = TtActivityStaff(
                                activity_id=activity_id,
                                staff_id=staff_id
                            )
                            insert_activity_staff_objs.append(activity_staff_obj)
                            staff_with_affected_week.setdefault(staff_id, set()).update(week_ids)
                        scheduled_staff_ids.update(staff_ids)
                        # staff_with_affected_week = {k: v for k, v in staff_with_affected_week.items() if v}

                    if location_ids:
                        for location_id in location_ids:
                            activity_location_obj = TtActivityLocation(
                                activity_id=activity_id,
                                location_id=location_id
                            )
                            insert_activity_location_objs.append(activity_location_obj)
                            location_with_affected_week.setdefault(location_id, set()).update(week_ids)
                        scheduled_location_ids.update(location_ids)

                    # if this activity have student_set, will recalculate the resource map for the student set
                    if student_set_ids:
                        for student_set_id in student_set_ids:
                            student_set_with_affected_week.setdefault(student_set_id, set()).update(week_ids)
                        scheduled_student_set_ids.update(student_set_ids)

                    # success count + 1
                    success_count += 1
                else:
                    # failed count + 1 and continue next
                    fail_count += 1
                    continue

            # if not empty, do bulk update and bulk create, if got duplicate key, will ignore and continue other
            if insert_activity_staff_objs:
                TtActivityStaff.objects.bulk_create(insert_activity_staff_objs, ignore_conflicts=True)
            if insert_activity_location_objs:
                TtActivityLocation.objects.bulk_create(insert_activity_location_objs, ignore_conflicts=True)
            # for bulk_update, must have id only can work
            if update_activity_objs:
                TtActivity.objects.bulk_update(update_activity_objs, ["scheduled_start_time", "scheduled_day", "scheduled",
                                                                      "scheduled_start_slot"])
            # check have variant activity or not, if got variant activity, need check need to merge or not
            if scheduled_variant_activity_ids:
                parent_ids = list(
                    TtActivity.objects.filter(id__in=scheduled_variant_activity_ids, is_variant=1)
                    .values_list('variant_parent_id', flat=True)
                    .distinct()
                )
                parent_ids = [pid for pid in parent_ids if pid]
                variant_activities = (
                    TtActivity.objects.filter(
                        Q(id__in=scheduled_variant_activity_ids) |
                        Q(variant_parent_id__in=scheduled_variant_activity_ids) |
                        Q(variant_parent_id__in=parent_ids) |
                        Q(id__in=parent_ids),
                        is_variant=1,
                        scheduled=1
                    )
                    .select_related("week_pattern")
                    .prefetch_related("week", "week_pattern__week", "staff", "location", "student_set","sequencing")
                )
                family_tree = defaultdict(list)
                parents_by_id = {}
                variant_activities_by_id = {v_act.id: v_act for v_act in variant_activities}
                for act in variant_activities:
                    if act.variant_parent_id is None:
                        parents_by_id[act.id] = act
                    else:
                        family_tree[act.variant_parent_id].append(act)

                m2m_weeks_to_add = defaultdict(list)  # Key: Target Activity ID, Value: List of Week IDs
                activity_ids_to_delete = set()

                for scheduled_variant_activity_id in scheduled_variant_activity_ids:
                    variant_activity = variant_activities_by_id.get(scheduled_variant_activity_id)
                    if not variant_activity or variant_activity.id in activity_ids_to_delete:
                        continue  # Skip if already merged/deleted in a previous loop pass

                    # Determine the parent ID and the actual parent object
                    p_id = variant_activity.id if variant_activity.variant_parent_id is None else variant_activity.variant_parent_id
                    parent_obj = parents_by_id.get(p_id)
                    children = family_tree.get(p_id, [])

                    # this scenario is handle when variant activity is parent
                    if variant_activity.variant_parent_id is None:
                        parent_key = get_match_key(variant_activity)

                        for child in children:
                            if child.id not in activity_ids_to_delete and get_match_key(child) == parent_key:
                                # Merge child into parent
                                child_weeks = [w.id for w in child.week.all()]
                                m2m_weeks_to_add[variant_activity.id].extend(child_weeks)
                                activity_ids_to_delete.add(child.id)
                    else:
                        # when variant activity is child
                        # when parent is scheduled only need get parent_obj
                        if parent_obj:
                            parent_key = get_match_key(parent_obj)
                        child_key = get_match_key(variant_activity)

                        if parent_obj and child_key == parent_key:
                            # Merge target child into parent
                            target_weeks = [w.id for w in variant_activity.week.all()]
                            m2m_weeks_to_add[parent_obj.id].extend(target_weeks)
                            activity_ids_to_delete.add(variant_activity.id)
                            sequencing_activity_ids.extend([va_s.id for va_s in variant_activity.sequencing.all()])
                            if parent_obj.id not in success_activity_ids:
                                success_activity_ids.append(parent_obj.id)
                        else:
                            # Compare against siblings to see if any match the target child
                            for sibling in children:
                                if sibling.id == variant_activity.id or sibling.id in activity_ids_to_delete:
                                    continue

                                if get_match_key(sibling) == child_key:
                                    # Merge target child into this matching sibling
                                    target_weeks = [w.id for w in variant_activity.week.all()]
                                    m2m_weeks_to_add[sibling.id].extend(target_weeks)
                                    activity_ids_to_delete.add(variant_activity.id)
                                    sequencing_activity_ids.extend([va_s.id for va_s in variant_activity.sequencing.all()])

                                    if sibling.id not in success_activity_ids:
                                        success_activity_ids.append(sibling.id)
                                    break  # Found a match, stop looking at other siblings

                merge_variant_activity_ids = []
                # update activity week to merge variant
                for target_activity_id, week_ids in m2m_weeks_to_add.items():
                    if target_activity_id not in activity_ids_to_delete and week_ids:
                        target_act = variant_activities_by_id.get(target_activity_id)
                        if target_act:
                            merge_variant_activity_ids.append(target_act.id)
                            target_act.week.add(*week_ids)

                # delete variant child
                if activity_ids_to_delete:
                    TtActivity.objects.filter(id__in=activity_ids_to_delete).delete()
                    need_delete_variant_activity_ids = list(activity_ids_to_delete)

                # get again merge variant activity to check need update data or not
                merge_variant_activities = (
                    TtActivity.objects.filter(id__in=merge_variant_activity_ids)
                    .select_related("week_pattern")
                    .prefetch_related("week", "week_pattern__week", )
                )
                for merge_variant_activity in merge_variant_activities:
                    # this query mostly just want to confirm this variant still have how many variant, if only 1 record, will remove variant tag
                    # for where statement, id and variant_parent_id need check for both self and parent id
                    related_variant_objs = (
                        TtActivity.objects.filter(
                            Q(id__in=[merge_variant_activity.id, merge_variant_activity.variant_parent_id]) |
                            Q(variant_parent_id__in=[merge_variant_activity.id, merge_variant_activity.variant_parent_id])
                        )
                    )
                    if merge_variant_activity.week_pattern:
                        current_week_pattern = list({"id":aw.id, "week":aw.week} for aw in merge_variant_activity.week_pattern.week.all())
                    else:
                        current_week_pattern = list({"id":aw.id, "week":aw.week} for aw in merge_variant_activity.week.all())

                    academic_term_week_pattern = list(
                        TtAcademicTermWeek.objects
                        .filter(academic_term_id=merge_variant_activity.academic_term_id)
                        .select_related("week")
                        .order_by("week__week")
                        .values("week_id", "week__week")
                    )

                    academic_term_week_pattern = [
                        {
                            "id": w["week_id"],
                            "week": w["week__week"],
                        }
                        for w in academic_term_week_pattern
                    ]
                    # if only have 1 record, means already not a variant
                    merge_variant_activity.name = remove_variant_activity_name_range(merge_variant_activity.name)
                    if related_variant_objs.count() == 1:
                        # if not a variant, remove code " 00x" and update is_variant to 0 and variant_parent_id = None
                        merge_variant_activity.code = re.sub(r"\s\d{3}$", "", merge_variant_activity.code)
                        merge_variant_activity.is_variant = 0
                        merge_variant_activity.variant_parent_id = None
                    else:
                        merge_variant_activity.name = merge_variant_activity.name + " " + format_variant_activity_name_week_pattern_to_ranges(
                            current_week_pattern, academic_term_week_pattern)
                    main_week = sorted({w["week"] for w in current_week_pattern})
                    merge_variant_activity.save()
                    redis_data["update"].setdefault(redis_activity_table, []).append({
                        "id": merge_variant_activity.id,
                        "week_pattern": convert_to_redis_week(main_week),
                        "code": merge_variant_activity.code,
                        "name": merge_variant_activity.name,
                        "is_variant": merge_variant_activity.is_variant,
                        "variant_parent_id": merge_variant_activity.variant_parent_id,
                    })
                    # add set() for activity with main week to handle microservices week, if same parent will keep merge for 2 time, will replace latest
                    activity_week_microservices_relation[merge_variant_activity.id] = sorted({w["id"] for w in current_week_pattern})

            # need after merging of variant only start to do recalculate resource map part
            # update resource map
            staff_with_affected_week = {k: v for k, v in staff_with_affected_week.items() if v}
            location_with_affected_week = {k: v for k, v in location_with_affected_week.items() if v}
            student_set_with_affected_week = {k: v for k, v in student_set_with_affected_week.items() if v}
            # since got update student set, this helper already cater if got update student set, no need pass in student details also will update student resource map
            redis_data = helper_recalculate_and_get_redis_resource_map(
                redis_data=redis_data,
                staff_with_affected_week=staff_with_affected_week,
                location_with_affected_week=location_with_affected_week,
                student_set_with_affected_week=student_set_with_affected_week,
                slot_per_week=slot_per_week,
            )
            # get again latest activities data from db
            # only get success schedule activity from db
            activities = (
                TtActivity.objects.filter(id__in=success_activity_ids)
                .select_related("activity_template", "department", "academic_term", "module", "week_pattern",
                                "activity_type", "zone", "availability", "start_preference", "usage_preference")
                .prefetch_related("staff", "location", "week", "week_pattern__week", )
            )

            # format the data need put to FE
            # dd request if no failed, dun show "failed: 0" to FE, cause will confuse user
            if fail_count > 0:
                message = __("message.scheduled_success_failed_count", success_count=success_count, fail_count=fail_count)
            else:
                message = __("message.scheduled_success_only", success_count=success_count)
            websocket_data = {
                "message": message,
                "success": success_count,
                "failed": fail_count,
                "data": [],
            }
            for activity in activities:
                microservices_week_ids = activity_week_microservices_relation.get(activity.id) or None
                final_data = {
                    "id": activity.id,
                    "code": activity.code,
                    "name": activity.name,
                    "desc": activity.desc,
                    "department": activity.department.name if activity.department else None,
                    "department_id": activity.department_id,
                    "planned_size": activity.planned_size,
                    "real_size": 0,
                    "activity_type_id": activity.activity_type_id,
                    "activity_type": activity.activity_type.name if activity.activity_type else None,
                    "activity_type_color": activity.activity_type.color if activity.activity_type else None,
                    "scheduled_start_time": str(activity.scheduled_start_time),
                    "scheduled_start_slot": activity.scheduled_start_slot,
                    "scheduled_day": __(
                        "attr.days_name." + str(activity.scheduled_day)) if activity.scheduled_day is not None else None,
                    "duration": activity.duration,
                    "slot_required": activity.slot_required,
                    "module": activity.module.name if activity.module else None,
                    "module_id": activity.module_id,
                    "location": list(activity.location.values_list("name", flat=True)) if activity.location else None,
                    "staff": list(activity.staff.values_list("name", flat=True)) if activity.staff else None,
                    "is_variant": activity.is_variant,
                    "variant_parent_id": activity.variant_parent_id,
                }
                # need used .copy() if not changes in final_data will change in this websocket_data also
                websocket_data['data'].append(final_data.copy())
                final_data['staff_ids'] = [s.id for s in activity.staff.all()]
                final_data['location_ids'] = [loc.id for loc in activity.location.all()]
                final_data['scheduled_day'] = activity.scheduled_day if activity.scheduled_day is not None else None

                # since engine not yet done medical schedule, now will force update the week_ids first, after engine done, need change it
                final_data["week_ids"] = [w.id for w in activity.week.all()]
                final_data["week_pattern_id"] = activity.week_pattern_id
                if microservices_week_ids:
                    final_data['week_ids'] = microservices_week_ids
                kafka_data.append(final_data)

            if sequencing_activity_ids:
                sequencing_related_activities = TtActivity.objects.filter(id__in=sequencing_activity_ids).prefetch_related("sequencing_from")
                for related_sequencing_activity in sequencing_related_activities:
                    redis_data["update"].setdefault(redis_activity_table, []).append(
                        helper_build_update_redis_sequencing_data(related_sequencing_activity)
                    )

            # send scheduled activity ids and slots to kafka
            # 2026-07-30 ENGINE NOT YET DO FOR CONFIRMED SCHEDULE, TEMPORARY COMMENT IT FIRST
            request_method = "confirmed_schedule"
            kafka_request_topic = request_topic
            scheduled_slots = [
                {
                    "activity": item["activity"],
                    "start_slot": item["start_slot"],
                    "teaching_staff": item.get("teaching_staff"),
                    "location": item.get("location"),
                }
                for item in schedule
                if int(item["activity"]) in success_activity_ids
            ]
            if (
                scheduled_slots
                and kafka_request_topic
            ):
                kafka_request_data = {
                    "session_id": response_data.get("session_id") or "",
                    "status": response_data["status"],
                    "schedule": scheduled_slots,
                }
                send_request(kafka_request_topic, kafka_request_data, None, request_method)

            # after get activities, will need get latest week_pattern for update to redis
            if redis_data['update'] or redis_data['delete']:
                bulk_sync_to_redis(redis_client, redis_data)

            # push websocket for FE
            if kafka_log.socket_id:
                push_websocket_notification(websocket_schedule_url, websocket_data, kafka_log.socket_id)

            # prepare for audit trail
            user_id = None
            
            if response_data.get("session_id"):
                user_id = User.objects.filter(name=response_data.get("session_id")).values_list("id", flat=True).first()

            audit_trail = AuditTrail.objects.create(
                user_id=user_id,
                type=AuditTrail.TYPE["admin"],
                ip_address=None,
            )

            remark_name = ", ".join(list(activities.values_list("name", flat=True)))
            AuditTrailDetails.custom_insert(
                audit_trail=audit_trail,
                action="schedule",
                remark_param={"name": remark_name},
                old_data=old_new_data["old_data"],
                new_data=old_new_data["new_data"],
            )
            
            # push kfaka microservice, add main week
            method = "schedule"
            kafka_topic = microservices_tt_topic
            if kafka_data and kafka_topic:
                kafka_request_data = {
                    "session_id": "",
                    "activity": kafka_data,
                    "delete_activity": need_delete_variant_activity_ids,
                }
                send_request(kafka_topic, kafka_request_data, None, method)

        else:
            # 2026-01-14 when not return success, just return a general msg first, cause they no return reason also
            err = {
                "code": status.HTTP_400_BAD_REQUEST,
                "error": __("validation.unavailable_to_schedule"),
                "errors": None,
            }
            if kafka_log.socket_id:
                push_websocket_notification(websocket_schedule_url, err, kafka_log.socket_id)
    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'])
        err = {
            "code": status.HTTP_400_BAD_REQUEST,
            "error": __("validation.unavailable_to_schedule"),
            "errors": None,
        }
        if kafka_log and kafka_log.socket_id:
            push_websocket_notification(websocket_schedule_url, err,kafka_log.socket_id)
        raise


class EngineProtocolError(ValueError):
    pass


def _parse_final_assignments(response_data):
    if not isinstance(response_data, dict):
        raise EngineProtocolError("Engine response must be an object")
    request_id = response_data.get("request_id")
    if not request_id:
        raise EngineProtocolError("Engine response is missing request_id")
    raw_schedule = response_data.get("schedule", "[]")
    assignments = json.loads(raw_schedule) if isinstance(raw_schedule, str) else raw_schedule
    if not isinstance(assignments, list):
        raise EngineProtocolError("Engine schedule must be a list")
    for index, item in enumerate(assignments):
        if not isinstance(item, dict) or "activity" not in item:
            raise EngineProtocolError(f"Engine assignment {index} is missing activity")
        try:
            int(item["activity"])
            if item.get("start_slot") not in (None, ""):
                int(item["start_slot"])
            for key in ("teaching_staff", "location"):
                value = item.get(key)
                decoded = json.loads(value) if isinstance(value, str) and value else (value or [])
                if not isinstance(decoded, list):
                    raise EngineProtocolError(f"Engine assignment {index} {key} must be a list")
                [int(resource_id) for resource_id in decoded]
        except (TypeError, ValueError, JSONDecodeError) as error:
            raise EngineProtocolError(f"Engine assignment {index} is invalid: {error}") from error
    return request_id, assignments


def _lock_engine_scope(assignments):
    activity_ids = sorted({int(item["activity"]) for item in assignments})
    activities = list(
        TtActivity.objects.filter(id__in=activity_ids)
        .values("id", "variant_parent_id")
        .order_by("id")
    )
    if len(activities) != len(activity_ids):
        found = {item["id"] for item in activities}
        raise EngineProtocolError(f"Unknown activity IDs in engine response: {sorted(set(activity_ids) - found)}")
    family_ids = {item["id"] for item in activities}
    family_ids.update(item["variant_parent_id"] for item in activities if item["variant_parent_id"])
    family_ids.update(
        TtActivity.objects.filter(variant_parent_id__in=family_ids).values_list("id", flat=True)
    )
    staff_ids = set()
    location_ids = set()
    for item in assignments:
        if item.get("teaching_staff"):
            staff_ids.update(json.loads(item["teaching_staff"]) if isinstance(item["teaching_staff"], str) else item["teaching_staff"])
        if item.get("location"):
            location_ids.update(json.loads(item["location"]) if isinstance(item["location"], str) else item["location"])
    list(TtActivity.objects.select_for_update().filter(id__in=sorted(family_ids)).order_by("id"))
    list(TtStaff.objects.select_for_update().filter(id__in=sorted(staff_ids)).order_by("id"))
    list(TtLocation.objects.select_for_update().filter(id__in=sorted(location_ids)).order_by("id"))
    return sorted(family_ids)


def _set_transaction_timeouts():
    if connection.vendor != "postgresql":
        return
    config = settings.RESOURCE_BOOKING_INTEGRATION
    with connection.cursor() as cursor:
        cursor.execute("SET LOCAL lock_timeout = %s", [f"{config['LOCK_TIMEOUT_MS']}ms"])
        cursor.execute("SET LOCAL statement_timeout = %s", [f"{config['STATEMENT_TIMEOUT_MS']}ms"])


@contextmanager
def _engine_response_integration_context(request_id):
    """Keep the initiating request correlation on an engine-applied change set."""

    delivery = None
    try:
        delivery_id = uuid.UUID(str(request_id))
    except (TypeError, ValueError, AttributeError):
        delivery_id = None
    if delivery_id is not None:
        delivery = (
            PostCommitDelivery.objects.filter(delivery_id=delivery_id)
            .only("correlation_id")
            .first()
        )
    correlation_id = (
        str(delivery.correlation_id)
        if delivery is not None and delivery.correlation_id
        else str(request_id)
    )
    token = set_context(
        IntegrationContext(
            correlation_id=correlation_id,
            causation_id=str(request_id),
            request_id=str(request_id),
            origin="scheduling_engine",
        )
    )
    try:
        yield
    finally:
        reset_context(token)


def schedule(response_data, kafka_log):
    """Validate outside the transaction, then atomically apply one final engine result."""
    request_id, assignments = _parse_final_assignments(response_data)
    with _engine_response_integration_context(request_id):
        with transaction.atomic():
            _set_transaction_timeouts()
            with mutation_tracking(request_id=request_id, origin="scheduling_engine") as tracker:
                affected_ids = _lock_engine_scope(assignments)
                track_activity_ids(affected_ids)
                claim = claim_engine_response(
                    request_id=request_id,
                    response_data=response_data,
                    change_set_id=tracker.change_set_id,
                    response_received_at=timezone.now(),
                )
                if claim.duplicate:
                    return {"duplicate": True, "change_set_id": str(claim.receipt.change_set_id)}
                _apply_schedule(response_data, kafka_log)
                return {"duplicate": False, "change_set_id": str(tracker.change_set_id)}

def preschedule(response_data,kafka_log):
    """
    2026-01-13 exact response
    {
        "session_id": "dev",
        "request_id": "ef4b11ff-9941-493e-8945-d12ed4d6b40a",
        "status": "success",
        "schedule": "[
            {
                \"activity\":\"19\",
                \"preference\":[0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,9,9,9,9,9,0,0,0,0,0,0,0,9,9,9,9,9,9,9,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,9,0,0,0,0,0,0,0,0,0,0,0]
            }
        ]"
    }
    """
    try:
        # got kafka_log only do action, else need double check why no request but receive response
        if response_data['status'] == "success" and kafka_log:
            # 2025-12-12 the exact response for the "schedule" is json_encode 1 more time, so need decode again for the "schedule"
            schedule = json.loads(response_data['schedule'])
            activity_ids = [int(item["activity"]) for item in schedule]

            # activities = (
            #     TtActivity.objects.filter(id__in=activity_ids)
            #     .select_related("activity_template", "department", "academic_term", "module", "week_pattern",
            #                     "activity_type", "zone", "availability", "start_preference", "usage_preference")
            #     .prefetch_related("staff","location")
            # )
            # format the data need put to FE
            # push websocket for FE
            if kafka_log.socket_id:
                push_websocket_notification(websocket_preschedule_url, schedule, kafka_log.socket_id)
    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'])

def preschedule_details(response_data,kafka_log):
    try:
        # got kafka_log only do action, else need double check why no request but receive response
        if response_data['status'] == "success" and kafka_log:
            # 2025-12-12 the exact response for the "schedule" is json_encode 1 more time, so need decode again for the "schedule"
            schedule = json.loads(response_data['schedule'])

            # push websocket for FE
            if kafka_log.socket_id:
                push_websocket_notification(websocket_preschedule_details_url, schedule, kafka_log.socket_id)
    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'])

def swap_get_available_resources(response_data,kafka_log):
    """
    sample response
    {
        'session_id': 'user123-abc456',
        'request_id': 'req789',
        'status': 'success',
        'schedule': '[{
            "activity":"405",
            "teaching_staff":"[25,34,28,46,36,18,32,27,26,30,40,42,31,39,17,35,23,24,41,38,33]",
            "start_slot":20,
            "location":"[18]"
        }]'
    }
    """
    try:
        # got kafka_log only do action, else need double check why no request but receive response
        if response_data['status'] == "success" and kafka_log:
            schedule = json.loads(response_data['schedule'])
            staff_ids = set()
            location_ids = set()
            for item in schedule:
                if item.get("teaching_staff"):
                    staff_ids.update(json.loads(item["teaching_staff"]))
                if item.get("location"):
                    location_ids.update(json.loads(item["location"]))

            staff = list(TtStaff.objects.filter(id__in=staff_ids).values_list("id","code","name").order_by("name"))
            location = list(TtLocation.objects.filter(id__in=location_ids).values_list("id","code","name").order_by("name"))
            websocket_data = {
                "staff": staff,
                "location": location
            }
            # push websocket for FE
            if kafka_log.socket_id:
                push_websocket_notification(websocket_swap_get_resources_url, websocket_data, kafka_log.socket_id)
    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'])

def _apply_swap(response_data,kafka_log):
    """
    Swap examples:
    1 fail swap (location can´t be swapped)
    Req {'session_id': 'user123-abc456', 'request_id': 'req789', 'slot': 20, 'activities': [405], 'teaching_staff': [25], 'location': [20]}
    Res {'session_id': 'user123-abc456', 'request_id': 'req789', 'status': 'failed', 'schedule': '[{"activity":"405","teaching_staff":"[]","start_slot":"","location":"[]"}]'}

    2 success swap, swap location and staff
    Req {'session_id': 'user123-abc456', 'request_id': 'req789', 'slot': 20, 'activities': [405], 'teaching_staff': [25], 'location': [18]}
    Res {'session_id': 'user123-abc456', 'request_id': 'req789', 'status': 'success', 'schedule': '[{"activity":"405","teaching_staff":"[25]","start_slot":"19508","location":"[18]"}]'}

    3 success swap, only swap staff
    Req {'session_id': 'user123-abc456', 'request_id': 'req789', 'slot': 20, 'activities': [405], 'teaching_staff': [25], 'location': []}
    Res {'session_id': 'user123-abc456', 'request_id': 'req789', 'status': 'success', 'schedule': '[{"activity":"405","teaching_staff":"[25]","start_slot":"19508","location":"[]"}]'}
    """
    if response_data['status'] == "success" and kafka_log:
        schedule = json.loads(response_data['schedule'])
        old_new_data = {
            "old_data": {},
            "new_data": {},
        }
        delete_activity_ids = []
        staff_ids = set()
        location_ids = set()
        sequencing_activity_ids = []
        kafka_data = []
        microservices_data = {}
        merge_to_variant_id = None

        response_slot = None
        for item in schedule:
            if item.get("teaching_staff"):
                staff_ids.update(json.loads(item["teaching_staff"]))
            if item.get("location"):
                location_ids.update(json.loads(item["location"]))
            if item.get("start_slot"):
                response_slot = int(item.get("start_slot"))
        staff_with_affected_week = {}
        location_with_affected_week = {}
        redis_data = {
            "update": {},
            "delete": {},
        }
        redis_staff_table = "staff"
        redis_location_table = "location"
        redis_activity_table = "activity"
        activity_id = schedule[0].get("activity") if schedule else None
        if response_slot and activity_id:
            slot = response_slot % slot_per_week
            activity = (
                TtActivity.objects.filter(id=activity_id)
                .select_related("week_pattern")
                .prefetch_related("week", "week_pattern__week", "staff", "location")
                .first()
            )
            redis_activity_data = {}
            old_new_data["old_data"].setdefault(activity_id, {})
            old_new_data["old_data"][activity_id]["staff_ids"] = sorted([s.id for s in activity.staff.all()])
            old_new_data["old_data"][activity_id]["location_ids"] = sorted([l.id for l in activity.location.all()])

            # check activity is variant or not, if is variant, have 1 more extra check is check it have parent or child in same staff and location or not
            if activity.is_variant:
                selected_variant_activity=None
                if activity.variant_parent_id:
                    variant_parent_id = activity.variant_parent_id
                else:
                    # if variant_parent_id is null, means self is parent
                    variant_parent_id = activity.id
                variant_activities = (
                    TtActivity.objects.filter(
                        Q(id=variant_parent_id) |
                        Q(variant_parent_id=variant_parent_id),
                        scheduled=1,
                        scheduled_start_slot=slot #only same slot have possible to merge
                    )
                    .exclude(id=activity.id)
                    .select_related("week_pattern")
                    .prefetch_related("week", "week_pattern__week", "staff", "location","sequencing")
                )
                # if got do merging, need get the delete activity staff or location id first, only can delete, cause need update resource map
                for variant_activity in variant_activities:
                    variant_staff_ids = set(variant_staff.id for variant_staff in variant_activity.staff.all())
                    variant_location_ids = set(variant_location.id for variant_location in variant_activity.location.all())
                    if variant_staff_ids == staff_ids and variant_location_ids == location_ids:
                        selected_variant_activity = variant_activity
                        break
                if selected_variant_activity:
                    merge_to_activity = activity
                    merge_to_variant_id = merge_to_activity.id
                    delete_activity = selected_variant_activity
                    # only selected variant_acitivty is parent, only will merge to selected_variant_activity, other scenario all merge to activity
                    if not selected_variant_activity.variant_parent_id:
                        merge_to_activity = selected_variant_activity
                        delete_activity = activity

                    # get the week_ids need merge to "merge_to_activity"
                    if delete_activity.week_pattern:
                        variant_merge_week_ids = list(aw.id for aw in delete_activity.week_pattern.week.all())
                    else:
                        variant_merge_week_ids = list(aw.id for aw in delete_activity.week.all())
                    merge_to_activity.week.add(*variant_merge_week_ids)
                    delete_activity_ids.append(delete_activity.id)
                    sequencing_activity_ids.extend([da_s.id for da_s in delete_activity.sequencing.all()])

                    delete_activity.delete()
                    activity = merge_to_activity

                    related_variant_objs = (
                        TtActivity.objects.filter(
                            Q(id__in=[activity.id, activity.variant_parent_id]) |
                            Q(variant_parent_id__in=[activity.id, activity.variant_parent_id])
                        )
                    )
                    related_variant_list = list(related_variant_objs)
                    # just used for activity after merge still is variant, to update the name
                    if activity.week_pattern:
                        current_week_pattern = list({"id":aw.id, "week":aw.week} for aw in activity.week_pattern.week.all())
                    else:
                        current_week_pattern = list({"id":aw.id, "week":aw.week} for aw in activity.week.all())

                    academic_term_week_pattern = list(
                        TtAcademicTermWeek.objects
                        .filter(academic_term_id=activity.academic_term_id)
                        .select_related("week")
                        .order_by("week__week")
                        .values("week_id", "week__week")
                    )

                    academic_term_week_pattern = [
                        {
                            "id": w["week_id"],
                            "week": w["week__week"],
                        }
                        for w in academic_term_week_pattern
                    ]
                    # if only have 1 record, means already not a variant
                    activity.name = remove_variant_activity_name_range(activity.name)
                    if len(related_variant_list) == 1:
                        # if not a variant, remove code " 00x" and update is_variant to 0 and variant_parent_id = None
                        activity.code = re.sub(r"\s\d{3}$", "", activity.code)
                        activity.is_variant = 0
                        activity.variant_parent_id = None
                        redis_activity_data["code"] = activity.code
                        redis_activity_data["is_variant"] = activity.is_variant
                        redis_activity_data["variant_parent_id"] = activity.variant_parent_id
                    else:
                        activity.name = activity.name + " " + format_variant_activity_name_week_pattern_to_ranges(current_week_pattern,academic_term_week_pattern)
                    redis_activity_data["name"] = activity.name
                    main_week = sorted({w["week"] for w in current_week_pattern})
                    redis_activity_data["week_pattern"] = convert_to_redis_week(main_week)
                    activity.save()

            redis_activity_data["id"] = activity.id
            # get the week ids for update resources map
            if activity.week_pattern:
                activity_week_ids = list(aw.id for aw in activity.week_pattern.week.all())
            else:
                activity_week_ids = list(aw.id for aw in activity.week.all())

            if staff_ids is not None:
                old_staff_ids = set(activity_staff.id for activity_staff in activity.staff.all())
                all_staff_ids = old_staff_ids.union(staff_ids)
                # set() will replace for the relation
                activity.staff.set(staff_ids)
                activity.staff_requirement = len(staff_ids)
                # all staff that need update resource map and redis data
                for staff_id in all_staff_ids:
                    staff_with_affected_week.setdefault(staff_id, set()).update(activity_week_ids)
                staff_with_affected_week = {k: v for k, v in staff_with_affected_week.items() if v}
                redis_data = helper_recalculate_and_get_redis_resource_map(
                    redis_data=redis_data,
                    staff_with_affected_week=staff_with_affected_week,
                    slot_per_week=slot_per_week,
                )

            if location_ids is not None:
                old_location_ids = set(activity_location.id for activity_location in activity.location.all())
                all_location_ids = old_location_ids.union(location_ids)
                # set() will replace for the relation
                activity.location.set(location_ids)
                activity.location_requirement = len(location_ids)
                # all location that need update resource map and redis data
                for location_id in all_location_ids:
                    location_with_affected_week.setdefault(location_id, set()).update(activity_week_ids)
                location_with_affected_week = {k: v for k, v in location_with_affected_week.items() if v}
                redis_data = helper_recalculate_and_get_redis_resource_map(
                    redis_data=redis_data,
                    location_with_affected_week=location_with_affected_week,
                    slot_per_week=slot_per_week,
                )
            activity.save()

            old_new_data["new_data"].setdefault(activity_id, {})
            old_new_data["new_data"][activity_id]["staff_ids"] = sorted([s.id for s in activity.staff.all()])
            old_new_data["new_data"][activity_id]["location_ids"] = sorted([l.id for l in activity.location.all()])

            redis_activity_data["location_required_no"] = activity.location_requirement
            redis_activity_data["staff_required_no"] = activity.staff_requirement
            redis_data["update"].setdefault(redis_activity_table, []).append(redis_activity_data)
            if delete_activity_ids:
                redis_data["delete"][redis_activity_table] = delete_activity_ids

            if sequencing_activity_ids:
                sequencing_related_activities = TtActivity.objects.filter(id__in=sequencing_activity_ids).prefetch_related("sequencing_from")
                for related_sequencing_activity in sequencing_related_activities:
                    redis_data["update"].setdefault(redis_activity_table, []).append(
                        helper_build_update_redis_sequencing_data(related_sequencing_activity)
                    )

            bulk_sync_to_redis(redis_client, redis_data)
            message = __("message.success_swap")
            websocket_data = {
                "message": message,
                "data": [],
            }
            final_data = {
                "id": activity.id,
                "code": activity.code,
                "name": activity.name,
                "desc": activity.desc,
                "department": activity.department.name if activity.department else None,
                "department_id": activity.department_id,
                "planned_size": activity.planned_size,
                "real_size": 0,
                "activity_type_id": activity.activity_type_id,
                "activity_type": activity.activity_type.name if activity.activity_type else None,
                "activity_type_color": activity.activity_type.color if activity.activity_type else None,
                "scheduled_start_time": str(activity.scheduled_start_time),
                "scheduled_start_slot": activity.scheduled_start_slot,
                "scheduled_day": __(
                    "attr.days_name." + str(activity.scheduled_day)) if activity.scheduled_day is not None else None,
                "duration": activity.duration,
                "slot_required": activity.slot_required,
                "module": activity.module.name if activity.module else None,
                "module_id": activity.module_id,
                "location": [loc.name for loc in activity.location.all()] if activity.location else None,
                "staff": [s.name for s in activity.staff.all()] if activity.staff else None,
            }
            # need used .copy() if not changes in final_data will change in this websocket_data also
            websocket_data['data'].append(final_data.copy())
            # copy from redis_activity_data, if got week_pattern, just remove it and custom week_ids to it
            microservices_data = redis_activity_data.copy()
            del microservices_data["location_required_no"]
            del microservices_data["staff_required_no"]
            if "week_pattern" in microservices_data:
                del microservices_data["week_pattern"]
                microservices_data["week_ids"] = activity_week_ids

            microservices_data['id'] = activity.id
            microservices_data['staff_ids'] = [s.id for s in activity.staff.all()]
            microservices_data['location_ids'] = [loc.id for loc in activity.location.all()]

            if kafka_log.socket_id:
                push_websocket_notification(websocket_swap_url, websocket_data, kafka_log.socket_id)

            # prepare for audit trail
            user_id = None
            
            if response_data.get("session_id"):
                user_id = User.objects.filter(name=response_data.get("session_id")).values_list("id", flat=True).first()

            audit_trail = AuditTrail.objects.create(
                user_id=user_id,
                type=AuditTrail.TYPE["admin"],
                ip_address=None,
            )
            AuditTrailDetails.custom_insert(
                audit_trail=audit_trail,
                action="swap",
                remark_param={"name": activity.name},
                old_data=old_new_data["old_data"],
                new_data=old_new_data["new_data"],
            )

            # push kafka microservice
            method = "swap"
            kafka_topic = microservices_tt_topic
            if microservices_data and kafka_topic:
                kafka_request_data = {
                    "session_id": "",
                    "activity": microservices_data,
                    "delete_activity": delete_activity_ids,
                }
                send_request(kafka_topic, kafka_request_data, None, method)
        else:
            # no return slot means failed, return websocket to FE let them know
            err = {
                "code": status.HTTP_400_BAD_REQUEST,
                "error": __("validation.unavailable_to_swap"),
                "errors": None,
            }
            if kafka_log.socket_id:
                push_websocket_notification(websocket_swap_url, err, kafka_log.socket_id)
    else:
        # not status success also return same fail message
        err = {
            "code": status.HTTP_400_BAD_REQUEST,
            "error": __("validation.unavailable_to_swap"),
            "errors": None,
        }
        if kafka_log and kafka_log.socket_id:
            push_websocket_notification(websocket_swap_url, err, kafka_log.socket_id)

def swap(response_data, kafka_log):
    request_id, assignments = _parse_final_assignments(response_data)
    with _engine_response_integration_context(request_id):
        with transaction.atomic():
            _set_transaction_timeouts()
            with mutation_tracking(request_id=request_id, origin="scheduling_engine") as tracker:
                affected_ids = _lock_engine_scope(assignments)
                track_activity_ids(affected_ids)
                claim = claim_engine_response(
                    request_id=request_id,
                    response_data=response_data,
                    change_set_id=tracker.change_set_id,
                    response_received_at=timezone.now(),
                )
                if claim.duplicate:
                    return {"duplicate": True, "change_set_id": str(claim.receipt.change_set_id)}
                _apply_swap(response_data, kafka_log)
                return {"duplicate": False, "change_set_id": str(tracker.change_set_id)}


def consume_forever():
    active_consumer = response_consumer()
    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
        msg_str = ""
        header_data = {}
        try:
            msg = active_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'preschedule')]
            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)
        # when no kafka_log return from update_response, means the request id is not found, need skip handle this response
            kafka_log = KafkaLog.update_response(topic, response_data,header_data)
            if not kafka_log:
                quarantine_engine_response(
                    payload=response_data,
                    headers=header_data,
                    reason="unknown_request_id",
                )
                active_consumer.store_offsets(msg)
                active_consumer.commit(asynchronous=False)
                continue
            if timezone.now() > kafka_log.request_at + datetime.timedelta(seconds=timeout):
                # Age is observable but no longer a discard condition: an old
                # matching response may be a valid redelivery after rollback.
                logger.warning(
                    "late_engine_response",
                    extra={"request_id": response_data.get("request_id")},
                )

            if "method" in header_data:
                method = header_data["method"]
                assigned = active_consumer.assignment()
                if method in {"schedule", "swap"} and assigned:
                    active_consumer.pause(assigned)
                try:
                    match method:
                        case "schedule":
                            schedule(response_data, kafka_log)
                        case "preschedule":
                            preschedule(response_data, kafka_log)
                        case "preschedule_details":
                            preschedule_details(response_data, kafka_log)
                        case "swap_get_available_resources":
                            swap_get_available_resources(response_data, kafka_log)
                        case "swap":
                            swap(response_data, kafka_log)
                        case _:
                            quarantine_engine_response(
                                payload=response_data,
                                headers=header_data,
                                reason="unsupported_method",
                            )
                finally:
                    if method in {"schedule", "swap"} and assigned:
                        active_consumer.resume(assigned)
            else:
                quarantine_engine_response(
                    payload=response_data,
                    headers=header_data,
                    reason="missing_method",
                )
            active_consumer.store_offsets(msg)
            active_consumer.commit(asynchronous=False)

        except (EngineProtocolError, EngineResponseConflict, JSONDecodeError) as e:
            payload = response_data if response_data is not None else {"raw": msg_str[:10000]}
            quarantine_engine_response(
                payload=payload,
                headers=header_data,
                reason=e.__class__.__name__,
            )
            logger.exception("engine_response_quarantined")
            active_consumer.store_offsets(msg)
            active_consumer.commit(asynchronous=False)
            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'])
        # since merge to 1 topic only, if need return to socket, will return at the function
        # if kafka_log.socket_id:
        #     push_websocket_notification(websocket_url, err,kafka_log.socket_id)
            continue

    # print(datetime.datetime.now())


if __name__ == "__main__":
    consume_forever()
