Skip to content

Commit dca9a19

Browse files
Update main.py
1 parent 9a9006e commit dca9a19

1 file changed

Lines changed: 90 additions & 32 deletions

File tree

‎main.py‎

Lines changed: 90 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -1,48 +1,106 @@
11
import asyncio
2-
import os
2+
from fastapi import FastAPI
33

4-
from src.event-bus.kafka_bus import KafkaEventBus
5-
from src.event-bus.redis_bus import RedisEventBus
6-
from src.event-bus.rabbitmq_bus import RabbitMQEventBus
7-
8-
from src.matching-service.consumer import MatchingConsumer
4+
from src.common.event_bus import EventBus
95
from src.common.utils import get_logger
6+
from src.common.models import MatchResultEvent, DriverLocationEvent
7+
8+
from src.dispatch-service.producer import DispatchProducer
9+
from src.dispatch-service.consumer import DispatchConsumer
10+
from src.dispatch-service.api_router import router, DRIVER_STORE, MATCH_RESULTS_STORE
11+
12+
# --------------------------------------------------------------------
13+
# CORE STATE STORES
14+
# --------------------------------------------------------------------
15+
16+
# In-memory driver store populated by Driver Location Service
17+
_driver_locations: list[DriverLocationEvent] = []
18+
19+
def driver_store():
20+
return list(_driver_locations)
21+
22+
# In-memory match results
23+
_match_results: list[MatchResultEvent] = []
24+
25+
def match_results_store():
26+
return _match_results
27+
28+
29+
# Surge lookup (injected later; stub for now)
30+
def surge_lookup(zone_id: str) -> float:
31+
# In production → query pricing service cache
32+
return 1.0
33+
34+
35+
# --------------------------------------------------------------------
36+
# SERVICE INITIALIZATION
37+
# --------------------------------------------------------------------
38+
39+
logger = get_logger("DispatchMain")
40+
event_bus = EventBus()
41+
42+
app = FastAPI(
43+
title="Dispatch Service",
44+
description="Driver-rider matching microservice.",
45+
version="1.0.0",
46+
)
47+
48+
49+
# --------------------------------------------------------------------
50+
# STARTUP SEQUENCE
51+
# --------------------------------------------------------------------
52+
53+
@app.on_event("startup")
54+
async def startup_event():
55+
56+
logger.info("Starting Dispatch Service...")
57+
58+
# Inject stores into API router
59+
global DRIVER_STORE, MATCH_RESULTS_STORE
60+
DRIVER_STORE = driver_store
61+
MATCH_RESULTS_STORE = _match_results
1062

63+
# Initialize producer and consumer
64+
producer = DispatchProducer(event_bus)
65+
consumer = DispatchConsumer(
66+
event_bus=event_bus,
67+
driver_store=driver_store,
68+
surge_lookup=surge_lookup,
69+
)
1170

12-
def get_event_bus():
13-
backend = os.getenv("EVENT_BUS_BACKEND", "kafka").lower()
71+
# Subscribe to trip requests
72+
await event_bus.subscribe("trip_requests", consumer.handle_trip_request)
1473

15-
if backend == "kafka":
16-
return KafkaEventBus()
17-
elif backend == "redis":
18-
return RedisEventBus()
19-
elif backend == "rabbitmq":
20-
return RabbitMQEventBus()
21-
else:
22-
raise ValueError(f"Unsupported EVENT_BUS_BACKEND: {backend}")
74+
# Subscribe to match results
75+
async def match_result_listener(data):
76+
event = MatchResultEvent(**data)
77+
_match_results.append(event)
78+
logger.info(f"[DISPATCH] Stored match result for rider {event.rider_id}")
2379

80+
await event_bus.subscribe("match_results", match_result_listener)
2481

25-
async def main():
26-
logger = get_logger("MatchingService")
82+
# Start background async tasks
83+
asyncio.create_task(producer.start())
2784

28-
event_bus = get_event_bus()
85+
logger.info("Dispatch Service started successfully.")
2986

30-
logger.info(f"Using event bus backend: {event_bus.__class__.__name__}")
3187

32-
# Connect to Kafka / Redis / RabbitMQ
33-
await event_bus.connect()
88+
# --------------------------------------------------------------------
89+
# ROUTER
90+
# --------------------------------------------------------------------
3491

35-
consumer = MatchingConsumer(event_bus)
92+
app.include_router(router, prefix="/dispatch")
3693

37-
try:
38-
logger.info("Starting Matching Service...")
39-
await consumer.start()
40-
except Exception as e:
41-
logger.error(f"Fatal error in Matching Service: {e}")
42-
finally:
43-
logger.info("Shutting down Matching Service...")
44-
await event_bus.close()
4594

95+
# --------------------------------------------------------------------
96+
# LOCAL RUN
97+
# --------------------------------------------------------------------
4698

4799
if __name__ == "__main__":
48-
asyncio.run(main())
100+
import uvicorn
101+
uvicorn.run(
102+
"src.dispatch-service.main:app",
103+
host="0.0.0.0",
104+
port=8002,
105+
reload=True
106+
)

0 commit comments

Comments
 (0)