aboutsummaryrefslogtreecommitdiff
path: root/order-service/app
diff options
context:
space:
mode:
Diffstat (limited to 'order-service/app')
-rw-r--r--order-service/app/api/v1/endpoints/orders.py8
-rw-r--r--order-service/app/init_db.py2
-rw-r--r--order-service/app/kafka_producer.py20
3 files changed, 28 insertions, 2 deletions
diff --git a/order-service/app/api/v1/endpoints/orders.py b/order-service/app/api/v1/endpoints/orders.py
index 4530958..dc81835 100644
--- a/order-service/app/api/v1/endpoints/orders.py
+++ b/order-service/app/api/v1/endpoints/orders.py
@@ -5,6 +5,8 @@ from sqlalchemy.orm import Session
from app.api import deps
from app.models.order import Order, OrderItem
from app.schemas.order import OrderCreate, OrderOut
+from app.kafka_producer import send_order_created_event
+
router = APIRouter()
@@ -38,10 +40,14 @@ def create_order(order_in: OrderCreate, db: Session = Depends(deps.get_db)):
db.commit()
db.refresh(db_order)
+ event_payload = {"order_id": db_order.id, "customer_id": order_in.customer_id}
+ send_order_created_event(event_payload)
+
return db_order
- except Exception:
+ except Exception as e:
db.rollback()
+ print(e)
raise HTTPException(
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
detail="Could not initialize checkout process.",
diff --git a/order-service/app/init_db.py b/order-service/app/init_db.py
index 7ee1465..88a3149 100644
--- a/order-service/app/init_db.py
+++ b/order-service/app/init_db.py
@@ -1,6 +1,6 @@
from app.core.database import engine, Base
-from app.models.order import Order, OrderItem
+from app.models.order import Order, OrderItem # noqa: F401
def init_database():
diff --git a/order-service/app/kafka_producer.py b/order-service/app/kafka_producer.py
new file mode 100644
index 0000000..6b24395
--- /dev/null
+++ b/order-service/app/kafka_producer.py
@@ -0,0 +1,20 @@
+import json
+from confluent_kafka import Producer
+
+conf = {"bootstrap.servers": "localhost:9092"}
+producer = Producer(conf)
+
+
+def delivery_report(err, msg):
+ if err is not None:
+ print(f"Message delivery failed: {err}")
+ else:
+ print(f"Message delivered to {msg.topic()} [{msg.partition()}]")
+
+
+def send_order_created_event(order_data: dict):
+ data_str = json.dumps(order_data)
+
+ producer.produce("orders", value=data_str.encode("utf-8"), callback=delivery_report)
+
+ producer.flush()