Skip to content

Commit 9e2fbca

Browse files
committed
fix: bill agent backtests atomically and correct plugin gateway contracts
1 parent 4d47498 commit 9e2fbca

24 files changed

Lines changed: 1096 additions & 125 deletions

‎backend_api_python/app/celery_app.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,13 +50,18 @@ def __call__(self, *args, **kwargs):
5050
task_routes={
5151
"quantdinger.tasks.fast_analysis": {"queue": "ai"},
5252
"quantdinger.tasks.agent_job": {"queue": "jobs"},
53+
"quantdinger.tasks.expire_agent_jobs": {"queue": "maintenance"},
5354
"quantdinger.tasks.reflection": {"queue": "maintenance"},
5455
"quantdinger.tasks.ai_calibration": {"queue": "maintenance"},
5556
"quantdinger.tasks.market_catalog_sync": {"queue": "maintenance"},
5657
"quantdinger.tasks.worker_heartbeat": {"queue": "maintenance"},
5758
"quantdinger.tasks.cleanup_runtime_metadata": {"queue": "maintenance"},
5859
},
5960
beat_schedule={
61+
"expire-billed-agent-jobs": {
62+
"task": "quantdinger.tasks.expire_agent_jobs",
63+
"schedule": 60.0,
64+
},
6065
"reflection-cycle": {
6166
"task": "quantdinger.tasks.reflection",
6267
"schedule": max(300, int(os.getenv("REFLECTION_WORKER_INTERVAL_SEC", "86400"))),

‎backend_api_python/app/data_sources/futures.py‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -171,6 +171,8 @@ def _get_ticker_yfinance(self, symbol: str) -> Dict[str, Any]:
171171
yf_symbol = yf_symbol + "=F"
172172
t = yf.Ticker(yf_symbol)
173173
last = None
174+
source = "yfinance"
175+
timestamp = None
174176
try:
175177
last = getattr(t, "fast_info", {}).get("last_price")
176178
except Exception:
@@ -179,7 +181,9 @@ def _get_ticker_yfinance(self, symbol: str) -> Dict[str, Any]:
179181
hist = t.history(period="2d", interval="1d")
180182
if hist is not None and not hist.empty:
181183
last = float(hist["Close"].iloc[-1])
182-
return {"symbol": yf_symbol, "last": float(last or 0.0)}
184+
source = "kline_1d"
185+
timestamp = int(hist.index[-1].timestamp())
186+
return {"symbol": yf_symbol, "last": float(last or 0.0), "source": source, "timestamp": timestamp}
183187
except Exception:
184188
return {"symbol": symbol, "last": 0.0}
185189

‎backend_api_python/app/routes/agent_v1/backtests.py‎

Lines changed: 23 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -9,11 +9,12 @@
99
from typing import Any
1010

1111
from app.services.strategy_v2 import StrategyV2BacktestService
12+
from app.services.billing_service import BillingError
1213
from app.utils.agent_auth import (
1314
SCOPE_B, agent_required, current_token, current_user_id,
1415
instrument_allowed, market_allowed, with_idempotency,
1516
)
16-
from app.utils.agent_jobs import count_active_jobs, submit_job
17+
from app.utils.agent_jobs import _job_receipt, count_active_jobs, submit_job
1718
from app.utils.logger import get_logger
1819
from flask import request
1920

@@ -187,6 +188,10 @@ def create_backtest():
187188
if validation_error:
188189
return validation_error
189190

191+
with with_idempotency("backtest") as existing:
192+
if existing:
193+
return envelope(_job_receipt(existing, duplicate=True), message="idempotent replay")
194+
190195
token_id = int(current_token().get("id") or 0)
191196
tenant_cap = max(1, int(os.getenv("AGENT_MAX_CONCURRENT_JOBS_PER_TENANT", "4")))
192197
token_cap = max(1, int(os.getenv("AGENT_MAX_CONCURRENT_JOBS_PER_TOKEN", "2")))
@@ -207,22 +212,23 @@ def create_backtest():
207212
http=429,
208213
)
209214

210-
with with_idempotency("backtest") as existing:
211-
if existing:
212-
return envelope({
213-
"job_id": existing["job_id"],
214-
"status": existing["status"],
215-
"duplicate": True,
216-
}, message="idempotent replay")
217-
218215
payload = dict(body)
219216
payload["__user_id"] = current_user_id()
220-
job = submit_job(
221-
user_id=current_user_id(),
222-
agent_token_id=token_id,
223-
kind="backtest",
224-
request_payload=payload,
225-
runner=_run_backtest,
226-
idempotency_key=request.headers.get("Idempotency-Key"),
227-
)
217+
try:
218+
job = submit_job(
219+
user_id=current_user_id(),
220+
agent_token_id=token_id,
221+
kind="backtest",
222+
request_payload=payload,
223+
runner=_run_backtest,
224+
idempotency_key=request.headers.get("Idempotency-Key"),
225+
)
226+
except BillingError as exc:
227+
return error(exc.status, exc.code, details={"error_type": exc.code, **exc.details},
228+
retriable=exc.status >= 500, http=exc.status)
229+
except Exception:
230+
logger.exception("Agent backtest submission failed")
231+
return error(503, "BILLING_OR_JOB_UNAVAILABLE", retriable=True, http=503)
232+
if job["status"] == "failed":
233+
return error(503, "AGENT_JOB_DISPATCH_FAILED", details=job, http=503)
228234
return envelope(job, message="queued", status=202)

‎backend_api_python/app/routes/agent_v1/markets.py‎

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
11
"""Read-class market data endpoints."""
22
from __future__ import annotations
33

4+
import math
5+
46
from app.data.market_symbols_seed import (
57
get_hot_symbols as seed_get_hot_symbols,
68
search_symbols as seed_search_symbols,
@@ -136,19 +138,19 @@ def price():
136138
if not instrument_allowed(symbol):
137139
return error(403, f"Instrument not allowed: {symbol}", http=403)
138140
try:
139-
rows = _kline_service.get_kline(market=market, symbol=symbol, timeframe="1m", limit=1) or []
140-
if not rows:
141-
return envelope({"market": market, "symbol": symbol, "price": None})
142-
last = rows[-1]
143-
# KlineService rows are typically dicts with 'close'/'c' keys.
144-
close = (
145-
last.get("close") if isinstance(last, dict) else None
146-
) or (last.get("c") if isinstance(last, dict) else None)
141+
quote = _kline_service.get_realtime_price(market=market, symbol=symbol) or {}
142+
close = float(quote.get("price") or 0)
143+
if not math.isfinite(close) or close <= 0:
144+
return error(503, "brokerAccounts.quoteUnavailable", retriable=True, http=503)
145+
source = str(quote.get("source") or "unknown")
147146
return envelope({
148147
"market": market,
149148
"symbol": symbol,
150149
"price": close,
151-
"raw": last,
150+
"source": source,
151+
"quote_type": "historical" if source == "kline_1d" else "latest_available",
152+
"timestamp": quote.get("timestamp"),
153+
"raw": quote,
152154
})
153155
except Exception as exc:
154156
logger.error(f"agent_v1/price failed: {exc}", exc_info=True)

‎backend_api_python/app/services/alpaca_trading/client.py‎

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -580,7 +580,7 @@ def get_account_summary(self) -> Dict[str, Any]:
580580
logger.error(f"Alpaca get_account_summary failed: {e}")
581581
return {"success": False, "error": str(e)}
582582

583-
def get_positions(self) -> List[Dict[str, Any]]:
583+
def get_positions(self, *, raise_on_error: bool = False) -> List[Dict[str, Any]]:
584584
"""Get current positions."""
585585
try:
586586
self._ensure_connected()
@@ -606,9 +606,11 @@ def get_positions(self) -> List[Dict[str, Any]]:
606606
]
607607
except Exception as e:
608608
logger.error(f"Alpaca get_positions failed: {e}")
609+
if raise_on_error:
610+
raise
609611
return []
610612

611-
def get_orders(self, status: str = "all", limit: int = 100) -> List[Dict[str, Any]]:
613+
def get_orders(self, status: str = "all", limit: int = 100, *, raise_on_error: bool = False) -> List[Dict[str, Any]]:
612614
"""Get recent orders, including filled orders by default."""
613615
try:
614616
self._ensure_connected()
@@ -628,6 +630,7 @@ def get_orders(self, status: str = "all", limit: int = 100) -> List[Dict[str, An
628630
"id": str(o.id),
629631
"orderId": str(o.id),
630632
"symbol": o.symbol,
633+
"asset_class": _enum_value(getattr(o, "asset_class", "")),
631634
"side": _enum_value(getattr(o, "side", "")).lower(),
632635
"action": _enum_value(getattr(o, "side", "")).upper(),
633636
"quantity": _num(getattr(o, "qty", 0)),
@@ -652,6 +655,8 @@ def get_orders(self, status: str = "all", limit: int = 100) -> List[Dict[str, An
652655
]
653656
except Exception as e:
654657
logger.error(f"Alpaca get_orders failed: {e}")
658+
if raise_on_error:
659+
raise
655660
return []
656661

657662
def get_open_orders(self) -> List[Dict[str, Any]]:

‎backend_api_python/app/services/billing_service.py‎

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,16 @@
1818
logger = get_logger(__name__)
1919

2020

21+
class BillingError(Exception):
22+
"""A charge could not be accepted; callers must roll back their transaction."""
23+
24+
def __init__(self, code: str, *, status: int = 503, details: Optional[dict] = None):
25+
super().__init__(code)
26+
self.code = code
27+
self.status = status
28+
self.details = details or {}
29+
30+
2131
class BillingService:
2232
"""Billing and credit accounting service."""
2333

@@ -595,6 +605,78 @@ def _grant_lifetime_monthly_credits_best_effort(self, cur, user_id: int):
595605
# Best-effort; never break caller
596606
pass
597607

608+
def consume_in_transaction(self, cur, user_id: int, feature: str, reference_id: str) -> dict:
609+
"""Deduct and record credits using the caller's transaction, without committing."""
610+
enabled = self.is_billing_enabled()
611+
cost = max(0, int(self.get_feature_cost(feature) or 0))
612+
cur.execute("SELECT credits FROM qd_users WHERE id = ? FOR UPDATE", (user_id,))
613+
account = cur.fetchone()
614+
if not account:
615+
raise BillingError("BILLING_ACCOUNT_NOT_FOUND")
616+
balance = Decimal(str(account.get("credits") or 0))
617+
charge = {"enabled": enabled, "feature": feature, "cost": cost, "charged": 0,
618+
"refunded": 0, "remaining": float(balance), "referenceId": reference_id,
619+
"transactionId": None, "refundTransactionId": None, "status": "free"}
620+
if reference_id:
621+
cur.execute("""SELECT id, amount, balance_after FROM qd_credits_log
622+
WHERE user_id = ? AND action = 'consume' AND feature = ? AND reference_id = ?
623+
ORDER BY id DESC LIMIT 1""", (user_id, feature, reference_id))
624+
existing = cur.fetchone()
625+
if existing:
626+
amount = abs(int(existing["amount"]))
627+
return {**charge, "enabled": True, "cost": amount, "charged": amount,
628+
"transactionId": existing["id"], "status": "charged"}
629+
if not enabled or cost == 0:
630+
return charge
631+
cur.execute("""UPDATE qd_users SET credits = credits - ?, updated_at = NOW()
632+
WHERE id = ? AND credits >= ? RETURNING credits""", (cost, user_id, cost))
633+
updated = cur.fetchone()
634+
if not updated:
635+
raise BillingError("INSUFFICIENT_CREDITS", status=402, details={
636+
"feature": feature, "current": float(balance), "required": cost,
637+
"shortage": max(0, float(Decimal(cost) - balance)),
638+
})
639+
remaining = float(updated["credits"])
640+
cur.execute("""INSERT INTO qd_credits_log
641+
(user_id, action, amount, balance_after, feature, reference_id, remark, created_at)
642+
VALUES (?, 'consume', ?, ?, ?, ?, ?, ?) RETURNING id""",
643+
(user_id, -cost, remaining, feature, reference_id,
644+
f"Consume: {FEATURE_NAMES.get(feature, feature)}", datetime.now(timezone.utc)))
645+
transaction_id = cur.fetchone()["id"]
646+
return {**charge, "charged": cost, "remaining": remaining,
647+
"transactionId": transaction_id, "status": "charged"}
648+
649+
def refund_in_transaction(self, cur, user_id: int, charge: dict) -> dict:
650+
"""Refund the recorded debit once, atomically with the caller's job transition."""
651+
if int(charge.get("charged") or 0) <= 0:
652+
return charge
653+
reference = str(charge["referenceId"])
654+
feature = str(charge["feature"])
655+
cur.execute("SELECT credits FROM qd_users WHERE id = ? FOR UPDATE", (user_id,))
656+
if not cur.fetchone():
657+
raise BillingError("BILLING_ACCOUNT_NOT_FOUND")
658+
cur.execute("""SELECT amount FROM qd_credits_log WHERE user_id = ?
659+
AND action = 'consume' AND feature = ? AND reference_id = ?
660+
ORDER BY id DESC LIMIT 1""", (user_id, feature, reference))
661+
debit = cur.fetchone()
662+
if not debit:
663+
raise BillingError("BILLING_DEBIT_NOT_FOUND")
664+
amount = abs(int(debit["amount"]))
665+
cur.execute("""SELECT id, balance_after FROM qd_credits_log WHERE user_id = ?
666+
AND action = 'refund' AND reference_id = ? ORDER BY id DESC LIMIT 1""", (user_id, reference))
667+
refund = cur.fetchone()
668+
if not refund:
669+
cur.execute("""UPDATE qd_users SET credits = credits + ?, updated_at = NOW()
670+
WHERE id = ? RETURNING credits""", (amount, user_id))
671+
remaining = float(cur.fetchone()["credits"])
672+
cur.execute("""INSERT INTO qd_credits_log
673+
(user_id, action, amount, balance_after, feature, reference_id, remark, created_at)
674+
VALUES (?, 'refund', ?, ?, ?, ?, ?, NOW()) RETURNING id""",
675+
(user_id, amount, remaining, feature, reference, "Automatic refund: agent job did not complete"))
676+
refund = {"id": cur.fetchone()["id"], "balance_after": remaining}
677+
return {**charge, "refunded": amount, "remaining": float(refund["balance_after"]),
678+
"refundTransactionId": refund["id"], "status": "refunded"}
679+
598680
def check_and_consume(self, user_id: int, feature: str, reference_id: str = '') -> Tuple[bool, str]:
599681
"""
600682
Check and consume credits for a feature.

‎backend_api_python/app/services/kline.py‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -157,7 +157,8 @@ def get_realtime_price(
157157
'low': ticker.get('low', 0),
158158
'open': ticker.get('open', 0),
159159
'previousClose': ticker.get('previousClose', 0),
160-
'source': 'ticker'
160+
'source': ticker.get('source') or 'ticker',
161+
'timestamp': ticker.get('timestamp'),
161162
}
162163
self.cache.set(cache_key, result, 30)
163164
return result
@@ -184,7 +185,8 @@ def get_realtime_price(
184185
'low': latest.get('low', 0),
185186
'open': latest.get('open', 0),
186187
'previousClose': prev_close,
187-
'source': 'kline_1m'
188+
'source': 'kline_1m',
189+
'timestamp': latest.get('time'),
188190
}
189191
self.cache.set(cache_key, result, 30)
190192
return result
@@ -216,7 +218,8 @@ def get_realtime_price(
216218
'low': latest.get('low', 0),
217219
'open': latest.get('open', 0),
218220
'previousClose': prev_close,
219-
'source': 'kline_1d'
221+
'source': 'kline_1d',
222+
'timestamp': latest.get('time'),
220223
}
221224
self.cache.set(cache_key, result, 300)
222225
return result

‎backend_api_python/app/services/live_trading/account_snapshot.py‎

Lines changed: 61 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -620,6 +620,61 @@ def _fetch_binance_snapshot(
620620
return swap_pos, spot_pos, orders
621621

622622

623+
def _fetch_alpaca_snapshot(exchange_config: Dict[str, Any], errors: List[str]) -> Tuple[List[Dict], List[Dict], List[Dict]]:
624+
from app.services.alpaca_trading.symbols import parse_symbol
625+
626+
positions: List[Dict[str, Any]] = []
627+
orders: List[Dict[str, Any]] = []
628+
try:
629+
client = create_client(exchange_config, market_type="spot")
630+
except Exception:
631+
logger.warning("Alpaca snapshot connection failed", exc_info=True)
632+
errors.append("brokerAccounts.snapshotConnectionFailed")
633+
return [], positions, orders
634+
635+
def instrument(item):
636+
asset_class = str(item.get("asset_class") or "").lower()
637+
hint = "Crypto" if asset_class == "crypto" else "USStock" if asset_class == "us_equity" else None
638+
symbol, asset_class = parse_symbol(str(item.get("symbol") or ""), market_hint=hint)
639+
return {"symbol": symbol, "inst_id": str(item.get("symbol") or ""), "market_type": "spot",
640+
"market": "Crypto" if asset_class == "crypto" else "USStock", "asset_class": asset_class}
641+
642+
try:
643+
for item in client.get_positions(raise_on_error=True):
644+
quantity = float(item.get("qty") or item.get("quantity") or 0)
645+
if not quantity:
646+
continue
647+
positions.append({
648+
**instrument(item),
649+
"side": "short" if quantity < 0 or str(item.get("side")).lower() == "short" else "long",
650+
"size": abs(quantity),
651+
"entry_price": float(item.get("avg_entry_price") or item.get("avgCost") or 0),
652+
"mark_price": float(item.get("current_price") or item.get("currentPrice") or 0),
653+
"market_value": float(item.get("market_value") or item.get("marketValue") or 0),
654+
"unrealized_pnl": float(item.get("unrealized_pnl") or item.get("unrealizedPnL") or 0),
655+
})
656+
except Exception:
657+
logger.warning("Alpaca snapshot positions failed", exc_info=True)
658+
errors.append("brokerAccounts.snapshotPositionsFailed")
659+
try:
660+
for item in client.get_orders(status="open", limit=500, raise_on_error=True):
661+
orders.append({
662+
**instrument(item),
663+
"exchange_order_id": str(item.get("id") or item.get("orderId") or ""),
664+
"side": str(item.get("side") or "").lower(),
665+
"order_type": str(item.get("order_type") or item.get("orderType") or "").lower(),
666+
"price": float(item.get("limit_price") or item.get("limitPrice") or 0),
667+
"amount": abs(float(item.get("qty") or item.get("quantity") or 0)),
668+
"filled": abs(float(item.get("filled_qty") or item.get("filled") or 0)),
669+
"notional": float(item.get("notional") or 0),
670+
"status": str(item.get("status") or ""),
671+
})
672+
except Exception:
673+
logger.warning("Alpaca snapshot orders failed", exc_info=True)
674+
errors.append("brokerAccounts.snapshotOrdersFailed")
675+
return [], positions, orders
676+
677+
623678
def fetch_account_snapshot(*, user_id: int, credential_id: int) -> Dict[str, Any]:
624679
"""Live fetch swap/spot legs + open orders for one credential."""
625680
cred = int(credential_id or 0)
@@ -649,7 +704,12 @@ def fetch_account_snapshot(*, user_id: int, credential_id: int) -> Dict[str, Any
649704
spot_all: List[Dict[str, Any]] = []
650705
orders_all: List[Dict[str, Any]] = []
651706

652-
if exchange_id in ("okx", "okex"):
707+
if exchange_id == "alpaca":
708+
sp, st, od = _fetch_alpaca_snapshot(exchange_config, errors)
709+
swap_all.extend(sp)
710+
spot_all.extend(st)
711+
orders_all.extend(od)
712+
elif exchange_id in ("okx", "okex"):
653713
try:
654714
client = create_client(exchange_config, market_type="swap")
655715
sp, st, od = _fetch_okx_snapshot(client, exchange_id, errors)

0 commit comments

Comments
 (0)