import os
import grpc
from grpc_health.v1 import health
from grpc_health.v1 import health_pb2
from grpc_health.v1 import health_pb2_grpc
from concurrent import futures
import threading
from time import sleep
import cloudianConsumer_pb2_grpc as pb2_grpc
import cloudianConsumer_pb2 as pb2
from cloudianConsumer.consumer import CloudianConsumer,ERROR,MSG,DATA


class cloudianConsumer(pb2_grpc.cloudianConsumer):

    consumer = CloudianConsumer()

    def __init__(self, *args, **kwargs):
        pass
    
    def CreateGroup(self, request, context):
        result = self.consumer.createGroup(groupId=request.groupName, groupName=request.groupDesc)
        data = []
        if DATA in result:
            for ele in result[DATA]:
                datValue = pb2.DataValue()
                match result[DATA][ele]:
                    case str():
                        datValue(strval=(result[DATA][ele]))
                    case bool():
                        datValue(boolval=(result[DATA][ele]))
                    case int():
                        datValue(intval=(result[DATA][ele]))
                datElem = pb2.Data(element={ele:datValue})
                data.append(datElem)
            result[DATA] = data
        return pb2.Response(**result)

    def DeleteGroup(self, request, context):
        result = self.consumer.deleteGroup(groupId=request.groupName)
        return pb2.Response(**result)

    def CreateUser(self, request, context):
        result = self.consumer.createUser(groupId=request.groupName, userId=request.userName, fullName=request.userDesc)
        data = []
        if DATA in result:
            for ele in result[DATA]:
                datValue = pb2.DataValue(boolval=(result[DATA][ele])) if isinstance(result[DATA][ele],bool) else pb2.DataValue(strval=(result[DATA][ele]))
                datElem = pb2.Data(element={ele:datValue})
                data.append(datElem)
            result[DATA] = data
        return pb2.Response(**result)

    def DeleteUser(self, request, context):
        result = self.consumer.deleteUser(groupId=request.groupName, userId=request.userName)
        return pb2.Response(**result)
    
    def CreateS3Access(self, request, context):
        result = self.consumer.createS3Access(groupId=request.groupName, userId=request.userName)
        data = []
        if DATA in result:
            for ele in result[DATA]:
                datValue = None;
                match result[DATA][ele]:
                    case str():
                        datValue = pb2.DataValue(strval=(result[DATA][ele]))
                    case bool():
                        datValue = pb2.DataValue(boolval=(result[DATA][ele]))
                    case int():
                        datValue = pb2.DataValue(intval=(result[DATA][ele]))
                datElem = pb2.Data(element={ele:datValue})
                data.append(datElem)
            result[DATA] = data
        return pb2.Response(**result)
    
def _toggle_health(health_servicer: health.HealthServicer, service: str):
    next_status = health_pb2.HealthCheckResponse.SERVING
    while True:
        if next_status == health_pb2.HealthCheckResponse.SERVING:
            next_status = health_pb2.HealthCheckResponse.NOT_SERVING
        else:
            next_status = health_pb2.HealthCheckResponse.SERVING

        health_servicer.set(service, next_status)
        sleep(5)


def _configure_health_server(server: grpc.Server):
    health_servicer = health.HealthServicer(
        experimental_non_blocking=True,
        experimental_thread_pool=futures.ThreadPoolExecutor(max_workers=10),
    )
    health_pb2_grpc.add_HealthServicer_to_server(health_servicer, server)

    # Use a daemon thread to toggle health status
    toggle_health_status_thread = threading.Thread(
        target=_toggle_health,
        args=(health_servicer, "cloudianConsumer.cloudianConsumer"),
        daemon=True,
    )
    toggle_health_status_thread.start()

def serve():
    server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
    pb2_grpc.add_cloudianConsumerServicer_to_server(cloudianConsumer(), server)
    port=os.getenv("CONTAINER_INSECURE_PORT", "50051")
    server.add_insecure_port(f"0.0.0.0:{port}")
    server.start()
    server.wait_for_termination()

def secure_serve():
    server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
    port=os.getenv("CONTAINER_SECURE_PORT", "8445")
    pb2_grpc.add_cloudianConsumerServicer_to_server(
        cloudianConsumer(), server
    )
    with open("server.key", "rb") as fp:
        server_key = fp.read()
    with open("server.pem", "rb") as fp:
        server_cert = fp.read()

    creds = grpc.ssl_server_credentials([(server_key, server_cert)])
    server.add_secure_port(f"0.0.0.0:{port}", creds)
    _configure_health_server(server)
    server.start()
    server.wait_for_termination()

if __name__ == '__main__':
    #serve()
    secure_serve()
