diff --git a/alembic/versions/0080_usage_cost_subsidy_fields.py b/alembic/versions/0080_usage_cost_subsidy_fields.py new file mode 100644 index 0000000..fb73be9 --- /dev/null +++ b/alembic/versions/0080_usage_cost_subsidy_fields.py @@ -0,0 +1,46 @@ +"""0080 compute_usage_records 加成本/毛利/补贴字段 + +为补贴体系提供数据底座:每笔消费记录采购成本、毛利、是否参与补贴。 + +Revision ID: 0080 +Revises: 0079 +Create Date: 2026-09-13 +""" +from alembic import op +import sqlalchemy as sa + +revision = "0080" +down_revision = "0079" +branch_labels = None +depends_on = None + + +def _column_exists(conn, table_name: str, column_name: str) -> bool: + try: + r = conn.execute(sa.text( + "SELECT COUNT(*) FROM INFORMATION_SCHEMA.COLUMNS " + "WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = :t AND COLUMN_NAME = :c" + ), {"t": table_name, "c": column_name}) + return r.scalar() > 0 + except Exception: + return False + + +def _add(table, col): + conn = op.get_bind() + if not _column_exists(conn, table, col.name): + op.add_column(table, col) + + +def upgrade() -> None: + _add("compute_usage_records", sa.Column("cost_amount", sa.BigInteger, server_default="0")) + _add("compute_usage_records", sa.Column("gross_profit", sa.BigInteger, server_default="0")) + _add("compute_usage_records", sa.Column("model_cost_ratio", sa.Numeric(5, 4), server_default="1.0000")) + _add("compute_usage_records", sa.Column("subsidy_eligible", sa.Integer, server_default="1")) + + +def downgrade() -> None: + op.drop_column("compute_usage_records", "subsidy_eligible") + op.drop_column("compute_usage_records", "model_cost_ratio") + op.drop_column("compute_usage_records", "gross_profit") + op.drop_column("compute_usage_records", "cost_amount") diff --git a/app/api/routers/compute_internal.py b/app/api/routers/compute_internal.py index e801699..26d8952 100644 --- a/app/api/routers/compute_internal.py +++ b/app/api/routers/compute_internal.py @@ -28,6 +28,9 @@ class EngineDeductRequest(BaseModel): engine_log_id: int model_name: str = "" token_count: int = 0 + cost_micro: int = 0 # 采购成本(微元),补贴体系数据底座 + model_cost_ratio: float = 1.0 # 模型成本比例快照 + subsidy_eligible: int = 1 # 是否参与补贴 def _verify_internal_token(authorization: str = Header(default="")) -> None: @@ -38,6 +41,40 @@ def _verify_internal_token(authorization: str = Header(default="")) -> None: raise HTTPException(status_code=401, detail="invalid internal token") +class SubsidyGrantRequest(BaseModel): + username: str + amount_micro: int + period: str = "" + record_id: int = 0 + + +@router.post("/subsidy-grant", summary="compute 补贴发放(增加用户算力余额)") +async def subsidy_grant( + req: SubsidyGrantRequest, + db: Database = Depends(get_db), + _auth: None = Depends(_verify_internal_token), +): + """补贴发放:amount_micro(微元)→ 增加用户个人算力余额 + 记录审计。""" + if req.amount_micro <= 0: + return {"ok": True, "amount": 0, "reason": "zero_amount"} + amount_fen = max(0, int(req.amount_micro) // 10000) # 微元→分 + if amount_fen <= 0: + return {"ok": True, "amount": 0, "reason": "zero_fen"} + + user = await db.users.get_by_username(req.username) + if not user: + return {"ok": False, "amount": 0, "reason": "user_not_found"} + user_id = user["id"] + + # 增加个人算力余额 + from ...services.compute_pricing_service import adjust_user_quota + await adjust_user_quota(db, user_id, amount_fen, reason=f"subsidy_grant period={req.period} record_id={req.record_id}") + + logger.info("[compute-internal] subsidy granted username=%s amount_fen=%s period=%s record_id=%s", + req.username, amount_fen, req.period, req.record_id) + return {"ok": True, "amount": amount_fen, "username": req.username, "period": req.period} + + @router.post("/deduct", summary="compute 引擎扣费回调(按实际费用扣平台账本)") async def engine_deduct( req: EngineDeductRequest, @@ -53,6 +90,9 @@ async def engine_deduct( engine_log_id=req.engine_log_id, model_name=req.model_name, token_count=req.token_count, + cost_micro=req.cost_micro, + model_cost_ratio=req.model_cost_ratio, + subsidy_eligible=req.subsidy_eligible, ) return {"ok": result.get("ok", True), **result} except Exception as exc: # noqa: BLE001 diff --git a/app/infrastructure/models.py b/app/infrastructure/models.py index 43d6daa..6ecfd26 100644 --- a/app/infrastructure/models.py +++ b/app/infrastructure/models.py @@ -1431,6 +1431,10 @@ class ComputeUsageRecord(Base): actual_amount: Mapped[int] = mapped_column(BigInteger, default=0) # 实际扣费(分) balance_source: Mapped[str] = mapped_column(String, default="personal") # company/personal engine_log_id: Mapped[int] = mapped_column(Integer, default=0, index=True) # compute 引擎日志 id(幂等键) + cost_amount: Mapped[int] = mapped_column(BigInteger, default=0) # 采购成本(分) + gross_profit: Mapped[int] = mapped_column(BigInteger, default=0) # 毛利(分) + model_cost_ratio: Mapped[float] = mapped_column(Float, default=1.0) # 模型成本比例快照 + subsidy_eligible: Mapped[int] = mapped_column(Integer, default=1) # 是否参与补贴 created_at: Mapped[str] = mapped_column(String, default="") diff --git a/app/services/compute_pricing_service.py b/app/services/compute_pricing_service.py index ac10c09..e9cebc7 100644 --- a/app/services/compute_pricing_service.py +++ b/app/services/compute_pricing_service.py @@ -398,12 +398,14 @@ async def check_user_balance_available(db: Database, user_id: str) -> bool: async def deduct_by_engine_cost( db: Database, username: str, actual_cost_micro: int, engine_log_id: int, model_name: str = "", token_count: int = 0, + cost_micro: int = 0, model_cost_ratio: float = 1.0, subsidy_eligible: int = 1, ) -> dict: """compute 引擎扣费后回调:按引擎实际费用(微元)扣平台来源账本。 幂等:engine_log_id 已存在记录则直接返回成功(防网络重试重复扣)。 扣费优先级:企业分配余额(折扣从低到高)→ 个人余额。 actual_cost_micro 单位微元(1元=1e6),换算为分(1元=100分):amount_fen = micro // 10000。 + cost_micro 为采购成本,model_cost_ratio 为成本比例快照,subsidy_eligible 是否参与补贴。 """ # 幂等检查 from sqlalchemy import select as _sa_select @@ -423,6 +425,10 @@ async def deduct_by_engine_cost( return {"ok": False, "amount": 0, "source": "", "reason": "user_not_found"} user_id = user["id"] + # 成本/毛利(微元→分) + cost_fen = max(0, int(cost_micro) // MICRO_PER_FEN) + gross_fen = max(0, amount_fen - cost_fen) + # 构造 price_info(引擎已按折扣扣费,平台侧记录 actual_amount 即可,standard_price 用 actual 反推) price_info = { "standard_price": amount_fen, @@ -450,7 +456,8 @@ async def deduct_by_engine_cost( charged_from = cb["company_id"] if remaining <= 0: await _record_usage_with_engine_id(db, user_id, model_name, token_count, price_info, - amount_fen, "company", charged_from, int(engine_log_id)) + amount_fen, "company", charged_from, int(engine_log_id), + cost_fen, gross_fen, model_cost_ratio, subsidy_eligible) return {"ok": True, "amount": amount_fen, "source": "company", "company_id": charged_from} # 2. 企业不足部分从个人余额扣 @@ -458,20 +465,24 @@ async def deduct_by_engine_cost( ok = await _deduct_from_user_balance(db, user_id, remaining) if ok: await _record_usage_with_engine_id(db, user_id, model_name, token_count, price_info, - amount_fen, "personal", "", int(engine_log_id)) + amount_fen, "personal", "", int(engine_log_id), + cost_fen, gross_fen, model_cost_ratio, subsidy_eligible) return {"ok": True, "amount": amount_fen, "source": "personal", "reason": ""} # 3. 全部不足:记为欠费(引擎已放行),不阻断 await _record_usage_with_engine_id(db, user_id, model_name, token_count, price_info, - amount_fen, "personal", "", int(engine_log_id)) + amount_fen, "personal", "", int(engine_log_id), + cost_fen, gross_fen, model_cost_ratio, subsidy_eligible) return {"ok": False, "amount": amount_fen, "source": "personal", "reason": "insufficient_balance"} async def _record_usage_with_engine_id( db: Database, user_id: str, model: str, token_count: int, price_info: dict, amount: int, balance_source: str, company_id: str, engine_log_id: int, + cost_amount: int = 0, gross_profit: int = 0, model_cost_ratio: float = 1.0, + subsidy_eligible: int = 1, ): - """记录算力使用(含引擎日志 id 幂等键)。""" + """记录算力使用(含引擎日志 id 幂等键 + 成本/毛利字段)。""" record = ComputeUsageRecord( id=new_id(), user_id=user_id, @@ -484,6 +495,10 @@ async def _record_usage_with_engine_id( actual_amount=amount, balance_source=balance_source, engine_log_id=engine_log_id, + cost_amount=cost_amount, + gross_profit=gross_profit, + model_cost_ratio=model_cost_ratio, + subsidy_eligible=subsidy_eligible, created_at=now_str(), ) db.session.add(record)