-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathapi.py
More file actions
124 lines (98 loc) · 3.63 KB
/
Copy pathapi.py
File metadata and controls
124 lines (98 loc) · 3.63 KB
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
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
import logging
from typing import Optional
from fastapi import Depends, FastAPI, HTTPException
from pydantic import BaseModel, Field, condecimal, conint, constr, root_validator
from sqlalchemy import insert, text, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.db import get_session
from app.models import OrderDB, OutboxEvent
from app.order_types import Order, OrderSide, OrderType
logger = logging.getLogger(__name__)
app = FastAPI()
class CreateOrderModel(BaseModel):
"""Request model for creating an order."""
type_: OrderType = Field(..., alias="type")
side: OrderSide
instrument: constr(min_length=12, max_length=12)
limit_price: Optional[condecimal(decimal_places=2)]
quantity: conint(gt=0)
@root_validator
def validator(cls, values: dict):
type_ = values.get("type_")
limit_price = values.get("limit_price")
if type_ == OrderType.MARKET and limit_price:
raise ValueError(
"Providing a `limit_price` is prohibited for type `market`"
)
if type_ == OrderType.LIMIT and not limit_price:
raise ValueError("Attribute `limit_price` is required for type `limit`")
return values
class CreateOrderResponseModel(Order):
"""Response model for created order."""
@app.post("/requeue-failed-orders/{order_id}")
async def requeue_failed_order(
order_id: str, session: AsyncSession = Depends(get_session)
):
"""
Admin endpoint: manually requeue failed orders by setting status to PENDING
and re-inserting into outbox.
"""
await session.execute(
update(OrderDB).where(OrderDB.id == order_id).values(status="PENDING")
)
stmt = insert(OutboxEvent).values(order_id=order_id, event_type="ORDER_CREATED")
await session.execute(stmt)
await session.commit()
logger.info(f"Order {order_id} manually requeued")
return {"message": f"Order {order_id} requeued"}
@app.get("/health")
async def health_check():
return {"status": "ok"}
@app.post(
"/orders",
status_code=201,
response_model=CreateOrderResponseModel,
response_model_by_alias=True,
)
async def create_order(
model: CreateOrderModel, session: AsyncSession = Depends(get_session)
):
"""
Create a new order and queue it for placement via outbox pattern.
Returns 201 if persisted and guaranteed for async processing.
"""
try:
# 1. Persist order
order = OrderDB(
type=model.type_.value,
side=model.side.value,
instrument=model.instrument,
limit_price=model.limit_price,
quantity=model.quantity,
)
session.add(order)
await session.flush()
# 2. Insert outbox event
stmt = insert(OutboxEvent).values(order_id=order.id, event_type="ORDER_CREATED")
await session.execute(stmt)
# 3. Trigger NOTIFY (fires only if commit succeeds)
await session.execute(text("NOTIFY outbox_channel, 'new_order_created'"))
# 4. Commit atomically
await session.commit()
logger.info(f"Order {order.id} created and enqueued for placement")
return Order(
id=order.id,
created_at=order.created_at,
type=model.type_,
side=model.side,
instrument=model.instrument,
limit_price=model.limit_price,
quantity=model.quantity,
)
except Exception:
# Rollback on any error
await session.rollback()
logger.exception("Error creating order")
raise HTTPException(
status_code=500, detail="Internal server error while placing the order"
)