blob: ede6012cc54000a7a7d5b523cd25937c5ffa5262 (
plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
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
|