#!/usr/bin/python3
# -*- coding: utf-8 -*-
# Version: 1.01
# Author: wfr
# last edited: 14.07.2025
# by: wfr
import os
import json
import mysql.connector
from mysql.connector import Error
import sys

DB_CONFIG = {
    "host": "chge.at",
    "user": "monta_test_db_user",
    "password": "Ycze9_733",
    "database": "monta_test_db",
}

def get_nested_value(data_dict, keys, default=None):
    """Safely get a value from a nested dictionary."""
    temp_dict = data_dict
    for key in keys:
        if isinstance(temp_dict, dict):
            temp_dict = temp_dict.get(key)
        else:
            return default
        if temp_dict is None:
            return default
    return temp_dict

def insert_update_team(cursor, team_payload, event_time):
    """Inserts or updates a team record in the database."""
    sql = """
    INSERT INTO teams (
        id, name, legalEmail, externalId, partnerExternalId, partnerCustomPayload,
        joinCode, recipientCode, companyName, taxIdentificationNumber, vatNumber,
        operator_id, operator_name, operator_identifier, operator_partnerId, operator_vatNumber,
        address_address1, address_address2, address_address3, address_city, address_province,
        address_zipCode, address_countryId,
        currency_identifier, currency_name, currency_decimals,
        type, category, operatorId, userId, isFrozen, frozenReason, frozenAt,
        blockedAt, createdAt, updatedAt, deletedAt, last_event_time
    ) VALUES (
        %(id)s, %(name)s, %(legalEmail)s, %(externalId)s, %(partnerExternalId)s, %(partnerCustomPayload)s,
        %(joinCode)s, %(recipientCode)s, %(companyName)s, %(taxIdentificationNumber)s, %(vatNumber)s,
        %(operator_id)s, %(operator_name)s, %(operator_identifier)s, %(operator_partnerId)s, %(operator_vatNumber)s,
        %(address_address1)s, %(address_address2)s, %(address_address3)s, %(address_city)s, %(address_province)s,
        %(address_zipCode)s, %(address_countryId)s,
        %(currency_identifier)s, %(currency_name)s, %(currency_decimals)s,
        %(type)s, %(category)s, %(operatorId)s, %(userId)s, %(isFrozen)s, %(frozenReason)s, %(frozenAt)s,
        %(blockedAt)s, %(createdAt)s, %(updatedAt)s, %(deletedAt)s, %(last_event_time)s
    )
    ON DUPLICATE KEY UPDATE
        name = VALUES(name), legalEmail = VALUES(legalEmail), externalId = VALUES(externalId),
        partnerExternalId = VALUES(partnerExternalId), partnerCustomPayload = VALUES(partnerCustomPayload),
        joinCode = VALUES(joinCode), recipientCode = VALUES(recipientCode), companyName = VALUES(companyName),
        taxIdentificationNumber = VALUES(taxIdentificationNumber), vatNumber = VALUES(vatNumber),
        operator_id = VALUES(operator_id), operator_name = VALUES(operator_name),
        operator_identifier = VALUES(operator_identifier), operator_partnerId = VALUES(operator_partnerId),
        operator_vatNumber = VALUES(operator_vatNumber), address_address1 = VALUES(address_address1),
        address_address2 = VALUES(address_address2), address_address3 = VALUES(address_address3),
        address_city = VALUES(address_city), address_province = VALUES(address_province),
        address_zipCode = VALUES(address_zipCode), address_countryId = VALUES(address_countryId),
        currency_identifier = VALUES(currency_identifier), currency_name = VALUES(currency_name),
        currency_decimals = VALUES(currency_decimals), type = VALUES(type), category = VALUES(category),
        operatorId = VALUES(operatorId), userId = VALUES(userId), isFrozen = VALUES(isFrozen),
        frozenReason = VALUES(frozenReason), frozenAt = VALUES(frozenAt), blockedAt = VALUES(blockedAt),
        createdAt = VALUES(createdAt), updatedAt = VALUES(updatedAt), deletedAt = VALUES(deletedAt),
        last_event_time = VALUES(last_event_time)
    """
    # Prepare data dictionary, handling None and nested structures
    data = {
        'id': team_payload.get('id'),
        'name': team_payload.get('name'),
        'legalEmail': team_payload.get('legalEmail'),
        'externalId': team_payload.get('externalId'),
        'partnerExternalId': team_payload.get('partnerExternalId'),
        'partnerCustomPayload': team_payload.get('partnerCustomPayload'), # Assuming TEXT or similar
        'joinCode': team_payload.get('joinCode'),
        'recipientCode': team_payload.get('recipientCode'),
        'companyName': team_payload.get('companyName'),
        'taxIdentificationNumber': team_payload.get('taxIdentificationNumber'),
        'vatNumber': team_payload.get('vatNumber'),
        'operator_id': get_nested_value(team_payload, ['operator', 'id']),
        'operator_name': get_nested_value(team_payload, ['operator', 'name']),
        'operator_identifier': get_nested_value(team_payload, ['operator', 'identifier']),
        'operator_partnerId': get_nested_value(team_payload, ['operator', 'partnerId']),
        'operator_vatNumber': get_nested_value(team_payload, ['operator', 'vatNumber']),
        'address_address1': get_nested_value(team_payload, ['address', 'address1']),
        'address_address2': get_nested_value(team_payload, ['address', 'address2']),
        'address_address3': get_nested_value(team_payload, ['address', 'address3']),
        'address_city': get_nested_value(team_payload, ['address', 'city']),
        'address_province': get_nested_value(team_payload, ['address', 'province']),
        'address_zipCode': get_nested_value(team_payload, ['address', 'zipCode']),
        'address_countryId': get_nested_value(team_payload, ['address', 'countryId']),
        'currency_identifier': get_nested_value(team_payload, ['currency', 'identifier']),
        'currency_name': get_nested_value(team_payload, ['currency', 'name']),
        'currency_decimals': get_nested_value(team_payload, ['currency', 'decimals']),
        'type': team_payload.get('type'),
        'category': team_payload.get('category'),
        'operatorId': team_payload.get('operatorId'),
        'userId': team_payload.get('userId'),
        'isFrozen': team_payload.get('isFrozen'),
        'frozenReason': team_payload.get('frozenReason'),
        'frozenAt': team_payload.get('frozenAt'),
        'blockedAt': team_payload.get('blockedAt'),
        'createdAt': team_payload.get('createdAt'),
        'updatedAt': team_payload.get('updatedAt'),
        'deletedAt': team_payload.get('deletedAt'),
        'last_event_time': event_time
    }
    cursor.execute(sql, data)
    print(cursor.statement)
    print(f"Processed team ID: {data['id']}")


def insert_update_team_member(cursor, member_payload, event_time):
    """Inserts or updates a team member record in the database."""
    sql = """
    INSERT INTO team_members (
        id, teamId, displayName,
        operator_id, operator_name, operator_identifier, operator_partnerId, operator_vatNumber,
        partnerExternalTeamId, userId, role, access, state, note,
        invitedAt, acceptedAt, rejectedAt, createdAt, updatedAt, deletedAt,
        chargePointIds, priceGroupId, costGroupId, teamMemberProfileId,
        canConfigureChargePoints, canPayWithTeamWallet, canManageTeamWallet,
        canRequestSponsoring, canManageTeamMembers, canPayForChargesCountryIds,
        teamWalletChargePaymentType, partnerExternalId, partnerCustomPayload,
        last_event_time
    ) VALUES (
        %(id)s, %(teamId)s, %(displayName)s,
        %(operator_id)s, %(operator_name)s, %(operator_identifier)s, %(operator_partnerId)s, %(operator_vatNumber)s,
        %(partnerExternalTeamId)s, %(userId)s, %(role)s, %(access)s, %(state)s, %(note)s,
        %(invitedAt)s, %(acceptedAt)s, %(rejectedAt)s, %(createdAt)s, %(updatedAt)s, %(deletedAt)s,
        %(chargePointIds)s, %(priceGroupId)s, %(costGroupId)s, %(teamMemberProfileId)s,
        %(canConfigureChargePoints)s, %(canPayWithTeamWallet)s, %(canManageTeamWallet)s,
        %(canRequestSponsoring)s, %(canManageTeamMembers)s, %(canPayForChargesCountryIds)s,
        %(teamWalletChargePaymentType)s, %(partnerExternalId)s, %(partnerCustomPayload)s,
        %(last_event_time)s
    )
    ON DUPLICATE KEY UPDATE
        teamId = VALUES(teamId), displayName = VALUES(displayName),
        operator_id = VALUES(operator_id), operator_name = VALUES(operator_name),
        operator_identifier = VALUES(operator_identifier), operator_partnerId = VALUES(operator_partnerId),
        operator_vatNumber = VALUES(operator_vatNumber), partnerExternalTeamId = VALUES(partnerExternalTeamId),
        userId = VALUES(userId), role = VALUES(role), access = VALUES(access), state = VALUES(state),
        note = VALUES(note), invitedAt = VALUES(invitedAt), acceptedAt = VALUES(acceptedAt),
        rejectedAt = VALUES(rejectedAt), createdAt = VALUES(createdAt), updatedAt = VALUES(updatedAt),
        deletedAt = VALUES(deletedAt), chargePointIds = VALUES(chargePointIds),
        priceGroupId = VALUES(priceGroupId), costGroupId = VALUES(costGroupId),
        teamMemberProfileId = VALUES(teamMemberProfileId), canConfigureChargePoints = VALUES(canConfigureChargePoints),
        canPayWithTeamWallet = VALUES(canPayWithTeamWallet), canManageTeamWallet = VALUES(canManageTeamWallet),
        canRequestSponsoring = VALUES(canRequestSponsoring), canManageTeamMembers = VALUES(canManageTeamMembers),
        canPayForChargesCountryIds = VALUES(canPayForChargesCountryIds),
        teamWalletChargePaymentType = VALUES(teamWalletChargePaymentType),
        partnerExternalId = VALUES(partnerExternalId), partnerCustomPayload = VALUES(partnerCustomPayload),
        last_event_time = VALUES(last_event_time)
    """
    # Prepare data dictionary
    data = {
        'id': member_payload.get('id'),
        'teamId': member_payload.get('teamId'),
        'displayName': member_payload.get('displayName'),
        'operator_id': get_nested_value(member_payload, ['operator', 'id']),
        'operator_name': get_nested_value(member_payload, ['operator', 'name']),
        'operator_identifier': get_nested_value(member_payload, ['operator', 'identifier']),
        'operator_partnerId': get_nested_value(member_payload, ['operator', 'partnerId']),
        'operator_vatNumber': get_nested_value(member_payload, ['operator', 'vatNumber']),
        'partnerExternalTeamId': member_payload.get('partnerExternalTeamId'),
        'userId': member_payload.get('userId'),
        'role': member_payload.get('role'),
        'access': member_payload.get('access'),
        'state': member_payload.get('state'),
        'note': member_payload.get('note'),
        'invitedAt': member_payload.get('invitedAt'),
        'acceptedAt': member_payload.get('acceptedAt'),
        'rejectedAt': member_payload.get('rejectedAt'),
        'createdAt': member_payload.get('createdAt'),
        'updatedAt': member_payload.get('updatedAt'),
        'deletedAt': member_payload.get('deletedAt'),
        # Convert lists to JSON strings for storage
        'chargePointIds': json.dumps(member_payload.get('chargePointIds', [])),
        'priceGroupId': member_payload.get('priceGroupId'),
        'costGroupId': member_payload.get('costGroupId'),
        'teamMemberProfileId': member_payload.get('teamMemberProfileId'),
        'canConfigureChargePoints': member_payload.get('canConfigureChargePoints'),
        'canPayWithTeamWallet': member_payload.get('canPayWithTeamWallet'),
        'canManageTeamWallet': member_payload.get('canManageTeamWallet'),
        'canRequestSponsoring': member_payload.get('canRequestSponsoring'),
        'canManageTeamMembers': member_payload.get('canManageTeamMembers'),
        # Convert lists to JSON strings for storage
        'canPayForChargesCountryIds': json.dumps(member_payload.get('canPayForChargesCountryIds', [])),
        'teamWalletChargePaymentType': member_payload.get('teamWalletChargePaymentType'),
        'partnerExternalId': member_payload.get('partnerExternalId'),
        'partnerCustomPayload': member_payload.get('partnerCustomPayload'),
        'last_event_time': event_time
    }
    cursor.execute(sql, data)
    print(f"Processed team member ID: {data['id']}")

def insert_update_charge_point(cursor, cp_payload, event_time):
    """Inserts or updates a charge point record in the database."""
    sql = """
    INSERT INTO charge_points (
        id, teamId, partnerExternalId, partnerCustomPayload, siteId,
        operator_id, operator_name, operator_identifier, operator_partnerId, operator_vatNumber,
        createdAt, updatedAt, deletedAt, activeAt, serialNumber, name, visibility, maxKw, type,
        note, operatorNote, state,
        location_latitude, location_longitude, location_address1, location_address2, location_address3,
        location_zip, location_city, location_country,
        connectors, deeplinks_app, deeplinks_web,
        lastMeterReadingKwh, cablePluggedIn, brandName, modelName, firmwareVersion, integrationType, evseId,
        priceGroupId, costGroupId, roamingPriceGroupId, sponsoredPriceGroupId, last_event_time
    ) VALUES (
        %(id)s, %(teamId)s, %(partnerExternalId)s, %(partnerCustomPayload)s, %(siteId)s,
        %(operator_id)s, %(operator_name)s, %(operator_identifier)s, %(operator_partnerId)s, %(operator_vatNumber)s,
        %(createdAt)s, %(updatedAt)s, %(deletedAt)s, %(activeAt)s, %(serialNumber)s, %(name)s, %(visibility)s, %(maxKW)s, %(type)s,
        %(note)s, %(operatorNote)s, %(state)s,
        %(location_latitude)s, %(location_longitude)s, %(location_address1)s, %(location_address2)s, %(location_address3)s,
        %(location_zip)s, %(location_city)s, %(location_country)s,
        %(connectors)s, %(deeplinks_app)s, %(deeplinks_web)s,
        %(lastMeterReadingKwh)s, %(cablePluggedIn)s, %(brandName)s, %(modelName)s, %(firmwareVersion)s, %(integrationType)s, %(evseId)s,
        %(priceGroupId)s, %(costGroupId)s, %(roamingPriceGroupId)s, %(sponsoredPriceGroupId)s, %(last_event_time)s
    )
    ON DUPLICATE KEY UPDATE
        teamId = VALUES(teamId), partnerExternalId = VALUES(partnerExternalId), partnerCustomPayload = VALUES(partnerCustomPayload), siteId = VALUES(siteId),
        operator_id = VALUES(operator_id), operator_name = VALUES(operator_name), operator_identifier = VALUES(operator_identifier), 
        operator_partnerId = VALUES(operator_partnerId), operator_vatNumber = VALUES(operator_vatNumber),
        updatedAt = VALUES(updatedAt), deletedAt = VALUES(deletedAt), activeAt = VALUES(activeAt), serialNumber = VALUES(serialNumber), 
        name = VALUES(name), visibility = VALUES(visibility), maxKw = VALUES(maxKw), type = VALUES(type),
        note = VALUES(note), operatorNote = VALUES(operatorNote), state = VALUES(state),
        location_latitude = VALUES(location_latitude), location_longitude = VALUES(location_longitude), 
        location_address1 = VALUES(location_address1), location_address2 = VALUES(location_address2), location_address3 = VALUES(location_address3),
        location_zip = VALUES(location_zip), location_city = VALUES(location_city), location_country = VALUES(location_country),
        connectors = VALUES(connectors), deeplinks_app = VALUES(deeplinks_app), deeplinks_web = VALUES(deeplinks_web),
        lastMeterReadingKwh = VALUES(lastMeterReadingKwh), cablePluggedIn = VALUES(cablePluggedIn), brandName = VALUES(brandName), 
        modelName = VALUES(modelName), firmwareVersion = VALUES(firmwareVersion), integrationType = VALUES(integrationType), evseId = VALUES(evseId),
        priceGroupId = VALUES(priceGroupId), costGroupId = VALUES(costGroupId), roamingPriceGroupId = VALUES(roamingPriceGroupId), 
        sponsoredPriceGroupId = VALUES(sponsoredPriceGroupId), last_event_time = VALUES(last_event_time)
    """
    data = {
        'id': cp_payload.get('id'), 'teamId': cp_payload.get('teamId'), 'partnerExternalId': cp_payload.get('partnerExternalId'),
        'partnerCustomPayload': cp_payload.get('partnerCustomPayload'), 'siteId': cp_payload.get('siteId'),
        'operator_id': get_nested_value(cp_payload, ['operator', 'id']), 'operator_name': get_nested_value(cp_payload, ['operator', 'name']),
        'operator_identifier': get_nested_value(cp_payload, ['operator', 'identifier']), 'operator_partnerId': get_nested_value(cp_payload, ['operator', 'partnerId']),
        'operator_vatNumber': get_nested_value(cp_payload, ['operator', 'vatNumber']), 'createdAt': cp_payload.get('createdAt'),
        'updatedAt': cp_payload.get('updatedAt'), 'deletedAt': cp_payload.get('deletedAt'), 'activeAt': cp_payload.get('activeAt'),
        'serialNumber': cp_payload.get('serialNumber'), 'name': cp_payload.get('name'), 'visibility': cp_payload.get('visibility'),
        'maxKW': cp_payload.get('maxKW'), 'maxKw': cp_payload.get('maxKw'), 'type': cp_payload.get('type'), 'note': cp_payload.get('note'),
        'operatorNote': cp_payload.get('operatorNote'), 'state': cp_payload.get('state'),
        'location_latitude': get_nested_value(cp_payload, ['location', 'coordinates', 'latitude']), 'location_longitude': get_nested_value(cp_payload, ['location', 'coordinates', 'longitude']),
        'location_address1': get_nested_value(cp_payload, ['location', 'address', 'address1']), 'location_address2': get_nested_value(cp_payload, ['location', 'address', 'address2']),
        'location_address3': get_nested_value(cp_payload, ['location', 'address', 'address3']), 'location_zip': get_nested_value(cp_payload, ['location', 'address', 'zip']),
        'location_city': get_nested_value(cp_payload, ['location', 'address', 'city']), 'location_country': get_nested_value(cp_payload, ['location', 'address', 'country']),
        'connectors': json.dumps(cp_payload.get('connectors', [])), 'deeplinks_app': get_nested_value(cp_payload, ['deeplinks', 'app']),
        'deeplinks_web': get_nested_value(cp_payload, ['deeplinks', 'web']), 'lastMeterReadingKwh': cp_payload.get('lastMeterReadingKwh'),
        'cablePluggedIn': cp_payload.get('cablePluggedIn'), 'brandName': cp_payload.get('brandName'), 'modelName': cp_payload.get('modelName'),
        'firmwareVersion': cp_payload.get('firmwareVersion'), 'integrationType': cp_payload.get('integrationType'),
        'evseId': cp_payload.get('evseId'), 'priceGroupId': cp_payload.get('priceGroupId'), 'costGroupId': cp_payload.get('costGroupId'),
        'roamingPriceGroupId': cp_payload.get('roamingPriceGroupId'), 'sponsoredPriceGroupId': cp_payload.get('sponsoredPriceGroupId'),
        'last_event_time': event_time
    }
    cursor.execute(sql, data)
    print(f"Processed charge point ID: {data['id']}")

def insert_update_site(cursor, site_payload, event_time):
    """Inserts or updates a site record in the database."""
    sql = """
    INSERT INTO sites (
        id, teamId, name, partnerExternalId, partnerCustomPayload, siteType, visibility, note,
        operator_id, operator_name, operator_identifier, operator_partnerId, operator_vatNumber,
        location_latitude, location_longitude, location_address1, location_address2, location_address3,
        location_zip, location_city, location_country,
        createdAt, updatedAt, deletedAt, last_event_time
    ) VALUES (
        %(id)s, %(teamId)s, %(name)s, %(partnerExternalId)s, %(partnerCustomPayload)s, %(siteType)s, %(visibility)s, %(note)s,
        %(operator_id)s, %(operator_name)s, %(operator_identifier)s, %(operator_partnerId)s, %(operator_vatNumber)s,
        %(location_latitude)s, %(location_longitude)s, %(location_address1)s, %(location_address2)s, %(location_address3)s,
        %(location_zip)s, %(location_city)s, %(location_country)s,
        %(createdAt)s, %(updatedAt)s, %(deletedAt)s, %(last_event_time)s
    )
    ON DUPLICATE KEY UPDATE
        teamId = VALUES(teamId), name = VALUES(name), partnerExternalId = VALUES(partnerExternalId),
        partnerCustomPayload = VALUES(partnerCustomPayload), siteType = VALUES(siteType),
        visibility = VALUES(visibility), note = VALUES(note), operator_id = VALUES(operator_id),
        operator_name = VALUES(operator_name), operator_identifier = VALUES(operator_identifier),
        operator_partnerId = VALUES(operator_partnerId), operator_vatNumber = VALUES(operator_vatNumber),
        location_latitude = VALUES(location_latitude), location_longitude = VALUES(location_longitude),
        location_address1 = VALUES(location_address1), location_address2 = VALUES(location_address2),
        location_address3 = VALUES(location_address3), location_zip = VALUES(location_zip),
        location_city = VALUES(location_city), location_country = VALUES(location_country),
        updatedAt = VALUES(updatedAt), deletedAt = VALUES(deletedAt), last_event_time = VALUES(last_event_time)
    """
    data = {
        'id': site_payload.get('id'), 'teamId': site_payload.get('teamId'), 'name': site_payload.get('name'),
        'partnerExternalId': site_payload.get('partnerExternalId'), 'partnerCustomPayload': site_payload.get('partnerCustomPayload'),
        'siteType': site_payload.get('siteType'), 'visibility': site_payload.get('visibility'), 'note': site_payload.get('note'),
        'operator_id': get_nested_value(site_payload, ['operator', 'id']), 'operator_name': get_nested_value(site_payload, ['operator', 'name']),
        'operator_identifier': get_nested_value(site_payload, ['operator', 'identifier']), 'operator_partnerId': get_nested_value(site_payload, ['operator', 'partnerId']),
        'operator_vatNumber': get_nested_value(site_payload, ['operator', 'vatNumber']),
        'location_latitude': get_nested_value(site_payload, ['location', 'coordinates', 'latitude']), 'location_longitude': get_nested_value(site_payload, ['location', 'coordinates', 'longitude']),
        'location_address1': get_nested_value(site_payload, ['location', 'address', 'address1']), 'location_address2': get_nested_value(site_payload, ['location', 'address', 'address2']),
        'location_address3': get_nested_value(site_payload, ['location', 'address', 'address3']), 'location_zip': get_nested_value(site_payload, ['location', 'address', 'zip']),
        'location_city': get_nested_value(site_payload, ['location', 'address', 'city']), 'location_country': get_nested_value(site_payload, ['location', 'address', 'country']),
        'createdAt': site_payload.get('createdAt'), 'updatedAt': site_payload.get('updatedAt'), 'deletedAt': site_payload.get('deletedAt'),
        'last_event_time': event_time
    }
    cursor.execute(sql, data)
    print(f"Processed site ID: {data['id']}")

def insert_update_charge_auth_token(cursor, token_payload, event_time):
    """Inserts or updates a charge auth token record in the database."""
    sql = """
    INSERT INTO charge_auth_tokens (
        id, uid, teamId, userId, name, type, visualId, issuer, state,
        isRoaming, lastUsedAt, createdAt, updatedAt, deletedAt, last_event_time
    ) VALUES (
        %(id)s, %(uid)s, %(teamId)s, %(userId)s, %(name)s, %(type)s, %(visualId)s, %(issuer)s, %(state)s,
        %(isRoaming)s, %(lastUsedAt)s, %(createdAt)s, %(updatedAt)s, %(deletedAt)s, %(last_event_time)s
    )
    ON DUPLICATE KEY UPDATE
        uid = VALUES(uid), teamId = VALUES(teamId), userId = VALUES(userId), name = VALUES(name),
        type = VALUES(type), visualId = VALUES(visualId), issuer = VALUES(issuer), state = VALUES(state),
        isRoaming = VALUES(isRoaming), lastUsedAt = VALUES(lastUsedAt), updatedAt = VALUES(updatedAt),
        deletedAt = VALUES(deletedAt), last_event_time = VALUES(last_event_time)
    """
    data = {
        'id': token_payload.get('id'), 'uid': token_payload.get('uid'), 'teamId': token_payload.get('teamId'),
        'userId': token_payload.get('userId'), 'name': token_payload.get('name'), 'type': token_payload.get('type'),
        'visualId': token_payload.get('visualId'), 'issuer': token_payload.get('issuer'), 'state': token_payload.get('state'),
        'isRoaming': token_payload.get('isRoaming'), 'lastUsedAt': token_payload.get('lastUsedAt'),
        'createdAt': token_payload.get('createdAt'), 'updatedAt': token_payload.get('updatedAt'),
        'deletedAt': token_payload.get('deletedAt'), 'last_event_time': event_time
    }
    cursor.execute(sql, data)
    print(f"Processed charge auth token ID: {data['id']}")

def insert_update_charge(cursor, charge_payload, event_time):
    """Inserts or updates a charge record in the database based on the provided example."""
    sql = """
    INSERT INTO charges (
        id, humanReadableId, chargePointId, userId, teamId, consumedKwh,
        price, currency, state, note, startedAt, endedAt,
        fullyChargedAt, stoppedAt, completedAt, paymentMethod,
        chargeAuthTokenId, createdAt, updatedAt, deletedAt,
        cablePluggedInAt, priceGroupId, costPriceGroupId, memberCostPriceGroupId, siteId,
        startMeterKwh, endMeterKwh,
        last_event_time
    ) VALUES (
        %(id)s, %(humanReadableId)s, %(chargePointId)s, %(userId)s, %(teamId)s, %(consumedKwh)s,
        %(price)s, %(currency)s, %(state)s, %(note)s, %(startedAt)s, %(endedAt)s,
        %(fullyChargedAt)s, %(stoppedAt)s, %(completedAt)s, %(paymentMethod)s,
        %(chargeAuthTokenId)s, %(createdAt)s, %(updatedAt)s, %(deletedAt)s,
        %(cablePluggedInAt)s, %(priceGroupId)s, %(costPriceGroupId)s, %(memberCostPriceGroupId)s, %(siteId)s,
        %(startMeterKwh)s, %(endMeterKwh)s,
        %(last_event_time)s
    )
    ON DUPLICATE KEY UPDATE
        chargePointId = VALUES(chargePointId), userId = VALUES(userId), teamId = VALUES(teamId),
        consumedKwh = VALUES(consumedKwh), price = VALUES(price), currency = VALUES(currency),
        state = VALUES(state), note = VALUES(note), startedAt = VALUES(startedAt), endedAt = VALUES(endedAt),
        fullyChargedAt = VALUES(fullyChargedAt), stoppedAt = VALUES(stoppedAt), completedAt = VALUES(completedAt),
        paymentMethod = VALUES(paymentMethod), chargeAuthTokenId = VALUES(chargeAuthTokenId),
        updatedAt = VALUES(updatedAt), deletedAt = VALUES(deletedAt),
        cablePluggedInAt = VALUES(cablePluggedInAt), priceGroupId = VALUES(priceGroupId),
        costPriceGroupId = VALUES(costPriceGroupId), memberCostPriceGroupId = VALUES(memberCostPriceGroupId),
        siteId = VALUES(siteId), startMeterKwh = VALUES(startMeterKwh), endMeterKwh = VALUES(endMeterKwh),
        last_event_time = VALUES(last_event_time)
    """
    data = {
        'id': charge_payload.get('id'),
        'humanReadableId': charge_payload.get('humanReadableId'),
        'chargePointId': charge_payload.get('chargePointId'),
        'userId': get_nested_value(charge_payload, ['user', 'id']),
        'teamId': get_nested_value(charge_payload, ['payingTeam', 'id']),
        'consumedKwh': charge_payload.get('consumedKwh'),
        'price': charge_payload.get('price'),
        'currency': get_nested_value(charge_payload, ['currency', 'identifier']),
        'state': charge_payload.get('state'),
        'note': charge_payload.get('note'),
        'startedAt': charge_payload.get('startedAt'),
        'endedAt': charge_payload.get('endedAt'),
        'fullyChargedAt': charge_payload.get('fullyChargedAt'),
        'stoppedAt': charge_payload.get('stoppedAt'),
        'completedAt': charge_payload.get('completedAt'),
        'paymentMethod': charge_payload.get('paymentMethod'),
        'chargeAuthTokenId': get_nested_value(charge_payload, ['chargeAuth', 'id']),
        'createdAt': charge_payload.get('createdAt'),
        'updatedAt': charge_payload.get('updatedAt'),
        'deletedAt': charge_payload.get('deletedAt'),
        'cablePluggedInAt': charge_payload.get('cablePluggedInAt'),
        'priceGroupId': charge_payload.get('priceGroupId'),
        'costPriceGroupId': charge_payload.get('costPriceGroupId'),
        'memberCostPriceGroupId': charge_payload.get('memberCostPriceGroupId'),
        'siteId': charge_payload.get('siteId'),
        'startMeterKwh': charge_payload.get('startMeterKwh'),
        'endMeterKwh': charge_payload.get('endMeterKwh'),
        'last_event_time': event_time
    }
    cursor.execute(sql, data)
    print(cursor.statement)
    print(f"Processed charge ID: {data['id']}")


connection = mysql.connector.connect(**DB_CONFIG)
connection2 = mysql.connector.connect(**DB_CONFIG)
if connection.is_connected():
    print("Successfully connected to database.")
    cursor = connection.cursor(dictionary=True)
    cursor2 = connection2.cursor(dictionary=True)
    sql = "SELECT * FROM Webhooks WHERE id >= 648983"
#    sql = "SELECT * FROM Webhooks"
    cursor.execute(sql)
    for res in cursor:
        print(res['json'])
        print('id: ',res['id'])
        print('timestamp: ',res['timestamp'])
        # Load JSON data
        try:
            events = json.loads(res['json'],strict=False)
        except json.JSONDecodeError as e:
            print(res['json'].rpartition(",{\"entityType\"")[0]+"]")
            try:
                events = json.loads(res['json'].rpartition(",{\"entityType\"")[0]+"]",strict=False)
            except json.JSONDecodeError as e:
                print (e)
        if events:
            for event in events:
                entity_type = event.get('entityType')
                payload = event.get('payload')
                event_time = event.get('eventTime') # Get event time
                if not payload or not entity_type:
                    print(f"Skipping event due to missing payload or entityType: {event}")
                    continue

                if entity_type == 'team':
                    insert_update_team(cursor2, payload, event_time)
                elif entity_type == 'team-member':
                    insert_update_team_member(cursor2, payload, event_time)
                elif entity_type == 'charge-point':
                    insert_update_charge_point(cursor2, payload, event_time)
                elif entity_type == 'site':
                    insert_update_site(cursor2, payload, event_time)
                elif entity_type == 'charge-auth-token':
                    insert_update_charge_auth_token(cursor2, payload, event_time)
                elif entity_type == 'charge':
                    insert_update_charge(cursor2, payload, event_time)
                else:
                    print(f"Skipping unknown entityType: {entity_type}")
connection.commit()
connection2.commit()

if cursor:
    cursor.close()
if cursor2:
    cursor2.close()
if connection and connection.is_connected():
    connection.close()
if connection2 and connection2.is_connected():
    connection2.close()
