diff options
Diffstat (limited to 'inventory-service/app/kafka_producer.py')
| -rw-r--r-- | inventory-service/app/kafka_producer.py | 33 |
1 files changed, 33 insertions, 0 deletions
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 |
