aboutsummaryrefslogtreecommitdiff
path: root/inventory-service/app/kafka_producer.py
diff options
context:
space:
mode:
Diffstat (limited to 'inventory-service/app/kafka_producer.py')
-rw-r--r--inventory-service/app/kafka_producer.py33
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