From 8e796c9cfcd65f6225a6ae3ec2a4419f265b0ebc Mon Sep 17 00:00:00 2001 From: alex Date: Fri, 17 Jul 2026 18:16:35 +0200 Subject: inventory-service init --- inventory-service/app/kafka_producer.py | 33 +++++++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) create mode 100644 inventory-service/app/kafka_producer.py (limited to 'inventory-service/app/kafka_producer.py') diff --git a/inventory-service/app/kafka_producer.py b/inventory-service/app/kafka_producer.py new file mode 100644 index 0000000..ede6012 --- /dev/null +++ b/inventory-service/app/kafka_producer.py @@ -0,0 +1,33 @@ +import json +from confluent_kafka import Producer +from app.core.config import settings + +_conf = {"bootstrap.servers": settings.KAFKA_BOOTSTRAP_SERVERS} +_producer = Producer(_conf) + + +def _delivery_report(err, msg): + if err is not None: + print(f"[Kafka] Delivery failed: {err}") + else: + print(f"[Kafka] Message delivered → {msg.topic()} [partition {msg.partition()}]") + + +def send_inventory_event(event_payload: dict) -> None: + """ + Publish an inventory result event to the inventory topic. + Called after the consumer decides to emit StockReserved or StockReservationFailed. + + Args: + event_payload: A dict that conforms to either StockReservedEvent or + StockReservationFailedEvent schema. + """ + pass # TODO: implement + + +def flush_producer() -> None: + """ + Block until all buffered messages are delivered. + Call this on application shutdown. + """ + pass # TODO: implement -- cgit v1.2.3