"""Local aiocoap sensor and stdlib HTTP bridge with validation and replay."""
import argparse
import asyncio
import hashlib
import http.client
from http.server import BaseHTTPRequestHandler, HTTPServer
import json
import threading

import aiocoap
from aiocoap import resource

COAP_PORT = 17682
HTTP_PORT = 17683
QUEUE_LIMIT = 2


class Sensor(resource.Resource):
    payload = b'{"event_id":"reading-001","celsius":21.5,"timestamp":"2026-10-08T12:00:00Z"}'

    async def render_get(self, request):
        reply = aiocoap.Message(payload=self.payload)
        reply.opt.content_format = 50
        return reply


class Receiver(BaseHTTPRequestHandler):
    receipts = {}
    wire = []

    def do_POST(self):
        body = self.rfile.read(int(self.headers["Content-Length"]))
        key = self.headers.get("Idempotency-Key")
        self.__class__.wire.append((self.raw_requestline, body, key))
        if self.path != "/measurements" or not key:
            self.send_error(400)
            return
        if key in self.__class__.receipts:
            status = 200
        else:
            self.__class__.receipts[key] = body
            status = 201
        self.send_response(status)
        self.send_header("Content-Length", "0")
        self.end_headers()

    def log_message(self, format, *args):
        pass


def start_http():
    server = HTTPServer(("127.0.0.1", HTTP_PORT), Receiver)
    thread = threading.Thread(target=server.serve_forever, daemon=True)
    thread.start()
    return server, thread


def stop_http(server, thread):
    server.shutdown()
    thread.join()
    server.server_close()


def map_payload(payload):
    record = json.loads(payload)
    if (set(record) != {"event_id", "celsius", "timestamp"}
            or not isinstance(record["event_id"], str)
            or not isinstance(record["celsius"], (int, float))
            or isinstance(record["celsius"], bool)
            or not -40 <= record["celsius"] <= 85
            or not isinstance(record["timestamp"], str)
            or not record["timestamp"].endswith("Z")):
        raise ValueError("sensor contract rejected: event_id, numeric celsius, UTC timestamp required")
    mapped = {"resource": "/temperature", "value": record["celsius"],
              "unit": "Cel", "timestamp": record["timestamp"]}
    return record["event_id"], json.dumps(mapped, separators=(",", ":")).encode()


def post(body, key):
    conn = http.client.HTTPConnection("127.0.0.1", HTTP_PORT, timeout=1)
    try:
        conn.request("POST", "/measurements", body=body, headers={
            "Content-Type": "application/json", "Idempotency-Key": key})
        response = conn.getresponse()
        response.read()
        return response.status
    finally:
        conn.close()


async def run(step):
    site = resource.Site()
    sensor = Sensor()
    site.add_resource(("temperature",), sensor)
    server = await aiocoap.Context.create_server_context(site, bind=("127.0.0.1", COAP_PORT))
    client = await aiocoap.Context.create_client_context()
    http_server, http_thread = start_http()
    outputs = {}
    try:
        async def read_sensor():
            request = aiocoap.Message(code=aiocoap.GET, uri=f"coap://127.0.0.1:{COAP_PORT}/temperature")
            return await asyncio.wait_for(client.request(request).response, 3)

        response = await read_sensor()
        event_id, body = map_payload(response.payload)
        key = hashlib.sha256(event_id.encode()).hexdigest()[:16]
        outputs[1] = [
            "STEP 1: GET /temperature from the local aiocoap sensor",
            f"CoAP response={response.code} content-format={response.opt.content_format}",
            f"CoAP payload bytes hex={response.payload.hex()}",
            f"CoAP payload UTF-8={response.payload.decode()}",
            "These payload bytes came from the local CoAP server in this run.",
        ]
        status = await asyncio.to_thread(post, body, key)
        line, received_body, received_key = Receiver.wire[-1]
        outputs[2] = [
            "STEP 2: map and POST the sensor record",
            f"HTTP request line bytes hex={line.hex()}",
            f"HTTP request line={line.decode().strip()}",
            f"HTTP Idempotency-Key={received_key}",
            f"HTTP body bytes hex={received_body.hex()}",
            f"HTTP body UTF-8={received_body.decode()}",
            f"HTTP status={status}; stored receipts={len(Receiver.receipts)}",
            "HTTP bytes above were received by the local stdlib server.",
        ]
        sensor.payload = b'{"event_id":"reading-002","celsius":"bad","timestamp":"2026-10-08T12:01:00Z"}'
        bad = await read_sensor()
        try:
            map_payload(bad.payload)
            raise AssertionError("malformed value passed validation")
        except ValueError as error:
            reason = str(error)
        outputs[3] = [
            "STEP 3: reject a malformed CoAP reading",
            f"CoAP response={bad.code} payload bytes hex={bad.payload.hex()}",
            f"decoded sensor reading={bad.payload.decode()}",
            f"validation={reason}",
            f"HTTP receipts remain={len(Receiver.receipts)}",
            "Malformed input was sent by the real local CoAP server, then rejected before POST.",
        ]
        sensor.payload = b'{"event_id":"reading-003","celsius":22.0,"timestamp":"2026-10-08T12:02:00Z"}'
        outage_response = await read_sensor()
        outage_id, outage_body = map_payload(outage_response.payload)
        outage_key = hashlib.sha256(outage_id.encode()).hexdigest()[:16]
        stop_http(http_server, http_thread)
        http_server = None
        queue = []
        try:
            await asyncio.to_thread(post, outage_body, outage_key)
            raise AssertionError("outage request unexpectedly succeeded")
        except OSError as error:
            if len(queue) >= QUEUE_LIMIT:
                raise RuntimeError("retry queue full")
            queue.append((outage_body, outage_key))
            failure = error.__class__.__name__
        outputs[4] = [
            "STEP 4: HTTP endpoint outage after a valid CoAP GET",
            f"CoAP payload UTF-8={outage_response.payload.decode()}",
            f"HTTP POST result={failure}",
            f"queue depth={len(queue)}/{QUEUE_LIMIT} key={outage_key}",
            f"queued HTTP body={outage_body.decode()}",
            "The HTTP listener was closed before POST; this is a local failure, not a cloud outage.",
        ]
        http_server, http_thread = start_http()
        queued_body, queued_key = queue.pop(0)
        retry_status = await asyncio.to_thread(post, queued_body, queued_key)
        line, received_body, received_key = Receiver.wire[-1]
        outputs[5] = [
            "STEP 5: replay the queued HTTP POST",
            f"HTTP request line bytes hex={line.hex()}",
            f"HTTP Idempotency-Key={received_key}",
            f"HTTP body UTF-8={received_body.decode()}",
            f"HTTP retry status={retry_status}; queue depth={len(queue)}",
            f"stored receipts={len(Receiver.receipts)}",
            "The request line and body were received by the restarted local HTTP server.",
        ]
        duplicate_status = await asyncio.to_thread(post, queued_body, queued_key)
        outputs[6] = [
            "STEP 6: submit the same idempotency key again",
            f"HTTP request line bytes hex={Receiver.wire[-1][0].hex()}",
            f"HTTP Idempotency-Key={Receiver.wire[-1][2]}",
            f"HTTP duplicate status={duplicate_status}",
            f"stored receipts={len(Receiver.receipts)}",
            f"same key stored once={sum(1 for item in Receiver.receipts if item == queued_key) == 1}",
            "The second POST reached the real local HTTP server; its dedupe table did not grow.",
        ]
        for number in ([step] if step else sorted(outputs)):
            print("\n".join(outputs[number]), flush=True)
    finally:
        if http_server:
            stop_http(http_server, http_thread)
        await client.shutdown()
        await server.shutdown()


if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument("--step", type=int, choices=range(1, 7))
    args = parser.parse_args()
    asyncio.run(run(args.step))
