- Task 1: Admin endpoints for wallet management (balance, ledger, adjust, reconcile) - Task 2: AI Credits admin UI (pricing.tsx, credit.tsx) - Task 3: Ollama security (NetworkPolicy, prompt validation, audit logging) - Task 4: LiteLLM integration (litellm.rs, migrated AI feature handlers) - Task 5: Refund architecture (ai_refunds table, endpoints) - Task 6: Coupons, promotions, referrals (order creation with coupon validation) - Task 7: Subscription lifecycle (plan upgrades/downgrades/cancellations, cron jobs) - Task 8: Credit expiration enforcement (daily cron task) - Task 9: Token cost engine (ai_model_cost_config, margin view) - Task 10: Observability (metrics tables, aggregation function) New files: - apps/users/src/litellm.rs (LiteLLM client) - apps/users/src/ai_credits.rs (Charging primitives) - apps/users/src/ai_subscription.rs (Plan lifecycle) - apps/users/src/handlers/ai_credits.rs (Admin endpoints) - apps/payments/src/ai_credits.rs (Package purchase) - apps/payments/src/payu.rs (PayU integration) - apps/cron/src/tasks/ai_credits.rs (Cron jobs) - 9 SQL migrations for schema - crates/db/src/models/ai_credits.rs (Repository) - tests for ai_credits
426 lines
12 KiB
Rust
426 lines
12 KiB
Rust
use axum::{
|
|
extract::State,
|
|
http::StatusCode,
|
|
routing::{get, post},
|
|
Json, Router,
|
|
};
|
|
use contracts::auth_middleware::AuthUser;
|
|
use serde::{Deserialize, Serialize};
|
|
use std::net::SocketAddr;
|
|
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
|
|
use uuid::Uuid;
|
|
use sqlx::postgres::PgPool;
|
|
use sqlx::FromRow;
|
|
|
|
pub mod ai_credits;
|
|
pub mod packages;
|
|
pub mod payu;
|
|
|
|
#[derive(Clone)]
|
|
pub struct AppState {
|
|
beeceptor_url: String,
|
|
client: reqwest::Client,
|
|
pool: PgPool,
|
|
payu: payu::PayuConfig,
|
|
}
|
|
|
|
#[derive(Debug, Serialize, Deserialize)]
|
|
struct CreateOrderRequest {
|
|
// Client-supplied amount/currency are accepted for backward
|
|
// compatibility but NOT trusted for pricing -- the real charge is
|
|
// always derived server-side from pricing_packages.price_inr (see
|
|
// create_order). Trusting a client-supplied amount for a real payment
|
|
// gateway would let a client pay any amount it likes for a package.
|
|
#[allow(dead_code)]
|
|
amount: Option<u64>,
|
|
#[allow(dead_code)]
|
|
currency: Option<String>,
|
|
package_id: Option<String>,
|
|
}
|
|
|
|
#[derive(Debug, Serialize)]
|
|
struct CreateOrderResponse {
|
|
key: String,
|
|
txnid: String,
|
|
amount: String,
|
|
productinfo: String,
|
|
firstname: String,
|
|
email: String,
|
|
phone: String,
|
|
surl: String,
|
|
furl: String,
|
|
hash: String,
|
|
payu_base_url: String,
|
|
udf1: String,
|
|
udf2: String,
|
|
// Kept for callers still reading the pre-PayU response shape.
|
|
order_id: String,
|
|
currency: String,
|
|
status: String,
|
|
}
|
|
|
|
#[derive(Debug, Serialize, Deserialize)]
|
|
struct VerifyPaymentRequest {
|
|
txnid: String,
|
|
mihpayid: String,
|
|
status: String,
|
|
hash: String,
|
|
amount: String,
|
|
productinfo: String,
|
|
firstname: String,
|
|
email: String,
|
|
#[serde(default)]
|
|
phone: Option<String>,
|
|
#[serde(default)]
|
|
udf1: Option<String>,
|
|
#[serde(default)]
|
|
udf2: Option<String>,
|
|
}
|
|
|
|
#[derive(Debug, Serialize)]
|
|
struct VerifyPaymentResponse {
|
|
verified: bool,
|
|
payment_id: String,
|
|
status: String,
|
|
message: String,
|
|
}
|
|
|
|
#[derive(Debug, Serialize)]
|
|
struct PaymentStatusResponse {
|
|
payment_id: String,
|
|
status: String,
|
|
amount: u64,
|
|
currency: String,
|
|
}
|
|
|
|
#[derive(Debug, FromRow)]
|
|
struct PricingPackageRow {
|
|
tracecoins_amount: i32,
|
|
price_inr: i32,
|
|
name: String,
|
|
}
|
|
|
|
#[derive(Debug, FromRow)]
|
|
struct UserContactRow {
|
|
email: String,
|
|
full_name: Option<String>,
|
|
phone: Option<String>,
|
|
}
|
|
|
|
#[derive(Debug, FromRow)]
|
|
#[allow(dead_code)]
|
|
struct PaymentRow {
|
|
id: Uuid,
|
|
user_id: Uuid,
|
|
package_id: Option<Uuid>,
|
|
tracecoins_credited: Option<i32>,
|
|
}
|
|
|
|
async fn create_order(
|
|
auth: AuthUser,
|
|
State(state): State<AppState>,
|
|
Json(payload): Json<CreateOrderRequest>,
|
|
) -> Result<Json<CreateOrderResponse>, (StatusCode, String)> {
|
|
let package_id_str = payload.package_id.as_ref().ok_or((StatusCode::BAD_REQUEST, "package_id is required".to_string()))?;
|
|
let package_id = Uuid::parse_str(package_id_str).map_err(|_| (StatusCode::BAD_REQUEST, "Invalid package id".to_string()))?;
|
|
|
|
let package = sqlx::query_as::<_, PricingPackageRow>(
|
|
"SELECT tracecoins_amount, price_inr, name FROM pricing_packages WHERE id = $1 AND is_active = true",
|
|
)
|
|
.bind(package_id)
|
|
.fetch_optional(&state.pool)
|
|
.await
|
|
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("DB error: {e}")))?;
|
|
|
|
let package = package.ok_or((StatusCode::BAD_REQUEST, "Invalid or inactive package".to_string()))?;
|
|
let tracecoins_credited = package.tracecoins_amount;
|
|
|
|
tracing::info!("Creating PayU order for package {} (₹{})", package_id, payu::paise_to_rupee_string(package.price_inr));
|
|
|
|
let contact = sqlx::query_as::<_, UserContactRow>(
|
|
"SELECT email, full_name, phone FROM users WHERE id = $1",
|
|
)
|
|
.bind(auth.user_id)
|
|
.fetch_optional(&state.pool)
|
|
.await
|
|
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("DB error: {e}")))?
|
|
.ok_or((StatusCode::NOT_FOUND, "User not found".to_string()))?;
|
|
|
|
let firstname = contact
|
|
.full_name
|
|
.as_deref()
|
|
.and_then(|n| n.split_whitespace().next())
|
|
.unwrap_or("Customer")
|
|
.to_string();
|
|
let phone = contact.phone.unwrap_or_default();
|
|
|
|
let txnid = payu::generate_txnid();
|
|
// Price is derived from pricing_packages.price_inr (server-side truth),
|
|
// never from the client-supplied `payload.amount` -- see CreateOrderRequest.
|
|
let amount_str = payu::paise_to_rupee_string(package.price_inr);
|
|
let productinfo = package.name.clone();
|
|
let udf1 = package_id_str.clone();
|
|
let udf2 = String::new();
|
|
|
|
let hash = payu::request_hash(
|
|
&state.payu,
|
|
&txnid,
|
|
&amount_str,
|
|
&productinfo,
|
|
&firstname,
|
|
&contact.email,
|
|
&udf1,
|
|
&udf2,
|
|
);
|
|
|
|
sqlx::query(
|
|
r#"
|
|
INSERT INTO payments (user_id, package_id, razorpay_order_id, amount_inr, tracecoins_credited, status)
|
|
VALUES ($1, $2, $3, $4, $5, 'PENDING')
|
|
"#,
|
|
)
|
|
.bind(auth.user_id)
|
|
.bind(package_id)
|
|
.bind(&txnid)
|
|
.bind(package.price_inr)
|
|
.bind(tracecoins_credited)
|
|
.execute(&state.pool)
|
|
.await
|
|
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("DB error: {e}")))?;
|
|
|
|
Ok(Json(CreateOrderResponse {
|
|
key: state.payu.merchant_key.clone(),
|
|
txnid: txnid.clone(),
|
|
amount: amount_str,
|
|
productinfo,
|
|
firstname,
|
|
email: contact.email,
|
|
phone,
|
|
surl: state.payu.surl.clone(),
|
|
furl: state.payu.furl.clone(),
|
|
hash,
|
|
payu_base_url: state.payu.base_url.clone(),
|
|
udf1,
|
|
udf2,
|
|
order_id: txnid,
|
|
currency: "INR".to_string(),
|
|
status: "created".to_string(),
|
|
}))
|
|
}
|
|
|
|
async fn verify_payment(
|
|
auth: AuthUser,
|
|
State(state): State<AppState>,
|
|
Json(payload): Json<VerifyPaymentRequest>,
|
|
) -> Result<Json<VerifyPaymentResponse>, (StatusCode, String)> {
|
|
tracing::info!("Verifying PayU payment: txnid={}", payload.txnid);
|
|
|
|
if !payu::verify_response_hash(
|
|
&state.payu,
|
|
&payload.status,
|
|
&payload.txnid,
|
|
&payload.amount,
|
|
&payload.productinfo,
|
|
&payload.firstname,
|
|
&payload.email,
|
|
payload.udf1.as_deref().unwrap_or(""),
|
|
payload.udf2.as_deref().unwrap_or(""),
|
|
&payload.hash,
|
|
) {
|
|
return Err((StatusCode::BAD_REQUEST, "Payment hash verification failed".to_string()));
|
|
}
|
|
|
|
if !payload.status.eq_ignore_ascii_case("success") {
|
|
return Err((StatusCode::BAD_REQUEST, "Payment was not successful".to_string()));
|
|
}
|
|
|
|
let payment = sqlx::query_as::<_, PaymentRow>(
|
|
r#"
|
|
SELECT id, user_id, package_id, tracecoins_credited
|
|
FROM payments
|
|
WHERE razorpay_order_id = $1 AND status = 'PENDING'
|
|
"#,
|
|
)
|
|
.bind(&payload.txnid)
|
|
.fetch_optional(&state.pool)
|
|
.await
|
|
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("DB error: {e}")))?;
|
|
|
|
let payment = match payment {
|
|
Some(p) => p,
|
|
None => return Err((StatusCode::NOT_FOUND, "Payment not found or already processed".to_string())),
|
|
};
|
|
|
|
if payment.user_id != auth.user_id {
|
|
return Err((StatusCode::FORBIDDEN, "Payment does not belong to user".to_string()));
|
|
}
|
|
|
|
let tracecoins = payment.tracecoins_credited.unwrap_or(0);
|
|
|
|
sqlx::query(
|
|
r#"
|
|
UPDATE payments SET
|
|
status = 'SUCCESS',
|
|
razorpay_payment_id = $1,
|
|
verified_at = NOW()
|
|
WHERE id = $2
|
|
"#,
|
|
)
|
|
.bind(&payload.mihpayid)
|
|
.bind(payment.id)
|
|
.execute(&state.pool)
|
|
.await
|
|
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("DB error: {e}")))?;
|
|
|
|
sqlx::query(
|
|
r#"
|
|
INSERT INTO tracecoin_wallets (user_id, balance, reserved)
|
|
VALUES ($1, $2, 0)
|
|
ON CONFLICT (user_id) DO UPDATE SET
|
|
balance = tracecoin_wallets.balance + excluded.balance
|
|
"#,
|
|
)
|
|
.bind(payment.user_id)
|
|
.bind(tracecoins as i64)
|
|
.execute(&state.pool)
|
|
.await
|
|
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("DB error: {e}")))?;
|
|
|
|
if let Ok(Some(wallet_id)) = sqlx::query_scalar::<_, Uuid>(
|
|
"SELECT id FROM tracecoin_wallets WHERE user_id = $1"
|
|
)
|
|
.bind(payment.user_id)
|
|
.fetch_optional(&state.pool)
|
|
.await
|
|
{
|
|
sqlx::query(
|
|
r#"
|
|
INSERT INTO tracecoin_ledger (wallet_id, transaction_type, amount, balance_after, reference_type, reference_id, description)
|
|
VALUES ($1, 'CREDIT', $2, $2, 'PAYMENT', $3, 'Package purchase')
|
|
"#,
|
|
)
|
|
.bind(wallet_id)
|
|
.bind(tracecoins as i64)
|
|
.bind(payment.id)
|
|
.execute(&state.pool)
|
|
.await
|
|
.ok();
|
|
}
|
|
|
|
let _ = sqlx::query(
|
|
r#"
|
|
INSERT INTO notifications (user_id, title, body, notification_type, reference_id)
|
|
VALUES ($1, $2, $3, $4, $5)
|
|
"#,
|
|
)
|
|
.bind(payment.user_id)
|
|
.bind("Tracecoins Purchased Successfully")
|
|
.bind(format!("Your {} Tracecoin package has been credited to your wallet.", tracecoins))
|
|
.bind("PAYMENT")
|
|
.bind(payment.id)
|
|
.execute(&state.pool)
|
|
.await
|
|
.ok();
|
|
|
|
Ok(Json(VerifyPaymentResponse {
|
|
verified: true,
|
|
payment_id: payload.mihpayid,
|
|
status: "success".to_string(),
|
|
message: "Payment verified successfully".to_string(),
|
|
}))
|
|
}
|
|
|
|
async fn get_payment_status(
|
|
State(state): State<AppState>,
|
|
axum::extract::Path(payment_id): axum::extract::Path<String>,
|
|
) -> Result<Json<PaymentStatusResponse>, (StatusCode, String)> {
|
|
tracing::info!("Getting payment status: payment_id={}", payment_id);
|
|
|
|
let status_url = format!("{}/{}", state.beeceptor_url.trim_end_matches('/'), payment_id);
|
|
let resp = state
|
|
.client
|
|
.get(&status_url)
|
|
.send()
|
|
.await
|
|
.map_err(|e| (StatusCode::BAD_GATEWAY, format!("Beeceptor error: {}", e)))?;
|
|
|
|
let status = resp.status();
|
|
let body: serde_json::Value = resp
|
|
.json()
|
|
.await
|
|
.map_err(|e| (StatusCode::BAD_GATEWAY, format!("Parse error: {}", e)))?;
|
|
|
|
if !status.is_success() {
|
|
return Ok(Json(PaymentStatusResponse {
|
|
payment_id,
|
|
status: "not_found".to_string(),
|
|
amount: 0,
|
|
currency: "INR".to_string(),
|
|
}));
|
|
}
|
|
|
|
let amount = body.get("amount").and_then(|v| v.as_u64()).unwrap_or(0);
|
|
let currency = body
|
|
.get("currency")
|
|
.and_then(|v| v.as_str())
|
|
.unwrap_or("INR")
|
|
.to_string();
|
|
let status_str = body
|
|
.get("status")
|
|
.and_then(|v| v.as_str())
|
|
.unwrap_or("unknown")
|
|
.to_string();
|
|
|
|
Ok(Json(PaymentStatusResponse {
|
|
payment_id,
|
|
status: status_str,
|
|
amount,
|
|
currency,
|
|
}))
|
|
}
|
|
|
|
#[tokio::main]
|
|
async fn main() {
|
|
tracing_subscriber::registry()
|
|
.with(tracing_subscriber::EnvFilter::new(
|
|
std::env::var("RUST_LOG").unwrap_or_else(|_| "info".into()),
|
|
))
|
|
.with(tracing_subscriber::fmt::layer())
|
|
.init();
|
|
|
|
let beeceptor_url = std::env::var("BEECEPTOR_URL")
|
|
.expect("BEECEPTOR_URL must be set");
|
|
|
|
let db_url = std::env::var("DATABASE_URL")
|
|
.expect("DATABASE_URL must be set");
|
|
let pool = PgPool::connect(&db_url)
|
|
.await
|
|
.expect("Failed to connect to database");
|
|
|
|
let state = AppState {
|
|
beeceptor_url,
|
|
client: reqwest::Client::new(),
|
|
pool,
|
|
payu: payu::PayuConfig::from_env(),
|
|
};
|
|
|
|
let app = Router::new()
|
|
.route("/api/payments/create-order", post(create_order))
|
|
.route("/api/payments/verify", post(verify_payment))
|
|
.route("/api/payments/{id}/status", get(get_payment_status))
|
|
.nest("/api/packages", packages::router())
|
|
.nest("/api/ai-credits", ai_credits::router())
|
|
.nest("/api/admin/ai-credits/packages", ai_credits::admin_router())
|
|
.with_state(state);
|
|
|
|
let port: u16 = std::env::var("PORT")
|
|
.unwrap_or_else(|_| "9116".to_string())
|
|
.parse()
|
|
.expect("PORT must be a valid u16");
|
|
let addr = SocketAddr::from(([0, 0, 0, 0], port));
|
|
|
|
tracing::info!("Payments service listening on {}", addr);
|
|
|
|
let listener = tokio::net::TcpListener::bind(&addr).await.unwrap();
|
|
axum::serve(listener, app).await.unwrap();
|
|
}
|