aboutsummaryrefslogtreecommitdiff
path: root/inventory-service/app/kafka_producer.py
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