import json

from pulsar import Function, SerDe

from pulsar.schema import *

class ProcessedCustomerMessage(object):
    def __init__(self):
        self.id = ''
        self.status = ''
        self.status_metadata = ''
        self.update_metadata = ''
        self.create_metadata = ''

class CustomSerDe(SerDe):
    def __init__(self):
        pass

    def serialize(self, object):
        return bytes(object, encoding='utf-8')

    def deserialize(self, input_bytes):
        input_str = str(input_bytes.decode())
        return input_str

class ParseFunction(Function):
    def __init__(self):
        pass

    def process(self, input, context):
        payload = json.loads(input)
        operation = payload["op"]

        output = ProcessedCustomerMessage()

        # Ensure REPLICA IDENTITY for PostGRES is set to FULL
        if payload["source"]["connector"] == "postgresql":
            if operation in ['c', 'u', 'r']:
                # TODO: Get schema from Schema Registry??
                output.id = payload["after"]["id"]
                output.status = payload["after"]["status"]
                output.status_metadata = payload["after"]["status_metadata"]
                output.create_metadata = payload["after"]["create_metadata"]
                output.update_metadata = payload["after"]["update_metadata"]

            elif operation == "d":
                # TODO: Get primary key from schema
                output.id = payload["before"]["id"]

            retval = {
                "id": output.id,
                "status": output.status,
                "status_metadata": output.status_metadata,
                "create_metadata": output.create_metadata,
                "update_metadata": output.update_metadata
            }
            return json.dumps(retval)

