Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,17 @@ WINDUP_CORS_ORIGIN_REGEX=
# 管理员邮箱白名单,逗号分隔;生产环境在私有 .env 中填写真实邮箱
WINDUP_ADMIN_EMAILS=admin@example.com

# ── 计费(易支付网关) ─────────────────────────────────────────────
# 商户 ID 与密钥。api_url 是网关的站点根地址(下单拼 /submit.php、查询拼 /api.php),
# 必须指向真实网关域名——留空会回落占位域名 pay.example.com,下单将不可用。
WINDUP_BILL_EPAY_PID=0
WINDUP_BILL_EPAY_PRIVATE_KEY=
WINDUP_BILL_EPAY_PUBLIC_KEY=
WINDUP_BILL_EPAY_API_URL=https://pay.example.com
# 异步/同步回调地址,指向部署后端的可公网访问域名
WINDUP_BILL_EPAY_NOTIFY_URL=
WINDUP_BILL_EPAY_RETURN_URL=

# ── MQ(Redis Stream 轻量消息队列) ──
# 本地/Compose 须同时起 web + worker,否则生成任务会一直 PENDING、邮件不会发出。
# docker compose up -d backend worker
Expand Down
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -40,3 +40,4 @@ htmlcov/

# 本地数据库初始化脚本
init.sql
.npm-cache/
81 changes: 81 additions & 0 deletions backend/packages/app/src/windup_app/bootstrap/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
from windup_app.server.project.model import Project # noqa: F401
from windup_app.server.project.service import service as project_service
from windup_app.server.quota import model as quota_model # noqa: F401
from windup_app.server.bill import model as bill_model # noqa: F401
from windup_app.server.sensitive_word.model import SensitiveWord # noqa: F401
from windup_app.server.sensitive_word.seed import seed_sensitive_words
from windup_app.server.sensitive_word.service import (
Expand All @@ -46,6 +47,8 @@
from windup_app.web.api.pixel_perfect import router as pixel_perfect_router
from windup_app.web.api.project import router as project_router
from windup_app.web.api.quota import router as quota_router
from windup_app.web.api.bill import router as bill_router
from windup_app.web.api.bill_notify import router as bill_notify_router
from windup_app.web.api.render3d import router as render3d_router
from windup_app.web.api.workflow_run import router as workflow_run_router
from windup_app.web.handler.exception_handlers import register_exception_handlers
Expand Down Expand Up @@ -116,6 +119,79 @@ async def _lifespan(app: FastAPI):
session.close()
print_banner()
from windup_app.web.api.generation import event_bus
from windup_app.server.bill.service import service as bill_service
import asyncio
import time
from windup_framework.db.redis import get_redis

# 关单重试间隔(秒):unknown 结果延后重试的间隔
_CLOSE_RETRY_DELAY = 600

async def _bill_pending_close_worker():
"""Polling bill:pending_close zset for expired orders.

bump-then-process:取出 order_no 后先把 score 推到 now+600s,只有终态
(closed/paid/skipped)才 zrem。崩溃最多延迟 600s 重试不丢单,兼治多
实例重复处理(score>now 不会被再次选出)。每单处理体包 to_thread,
同步 httpx + DB 不再阻塞事件循环。
"""
logger = logging.getLogger("windup.bill.worker.pending_close")
logger.info("Bill pending close worker started.")
while True:
try:
now_ts = int(time.time())
redis_client = get_redis()
orders = redis_client.zrangebyscore("bill:pending_close", 0, now_ts)
for order_no in orders:
if isinstance(order_no, bytes):
order_no = order_no.decode("utf-8")
# 先 bump score,防止崩溃丢单与多实例重复处理
redis_client.zadd(
"bill:pending_close",
{order_no: now_ts + _CLOSE_RETRY_DELAY},
)
try:
result = await asyncio.to_thread(
_close_one_order, order_no
)
if result in ("closed", "paid", "skipped"):
redis_client.zrem("bill:pending_close", order_no)
except Exception as e:
logger.error(f"Error closing order {order_no}: {e}")
except asyncio.CancelledError:
break
except Exception as e:
logger.error(f"Error in bill pending close worker: {e}")
await asyncio.sleep(10)

def _close_one_order(order_no: str) -> str:
"""在线程池中处理单个订单的关单逻辑。"""
with SessionLocal() as db_session:
result = bill_service.close_expired_orders(db_session, order_no)
db_session.commit()
return result

async def _subscription_expired_worker():
"""Polling expired subscriptions."""
logger = logging.getLogger("windup.bill.worker.subscription")
logger.info("Subscription expiration worker started.")
while True:
try:
await asyncio.to_thread(_process_subscriptions_once)
except asyncio.CancelledError:
break
except Exception as e:
logger.error(f"Error processing expired subscriptions: {e}")
await asyncio.sleep(60)

def _process_subscriptions_once() -> None:
"""在线程池中处理一轮订阅过期。"""
with SessionLocal() as db_session:
bill_service.process_expired_subscriptions(db_session)
db_session.commit()

bill_close_task = asyncio.create_task(_bill_pending_close_worker())
sub_expire_task = asyncio.create_task(_subscription_expired_worker())

sse_subscriber = RedisTaskEventSubscriber(
lambda project_id, task_id, event, data: event_bus.publish(
Expand All @@ -131,6 +207,9 @@ async def _lifespan(app: FastAPI):
try:
yield
finally:
bill_close_task.cancel()
sub_expire_task.cancel()
# await asyncio.gather(bill_close_task, sub_expire_task, return_exceptions=True)
sensitive_word_subscriber.stop()
sse_subscriber.stop()

Expand Down Expand Up @@ -177,6 +256,8 @@ def health() -> dict[str, str]:
app.include_router(generation_router)
app.include_router(quota_router)
app.include_router(admin_quota_router)
app.include_router(bill_router)
app.include_router(bill_notify_router)
app.include_router(render3d_router)
app.include_router(agent_router)
app.include_router(action_preset_router)
Expand Down
44 changes: 44 additions & 0 deletions backend/packages/app/src/windup_app/server/bill/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
"""支付与账单领域模块。"""

from windup_app.server.bill.epay import EpayProvider
from windup_app.server.bill.interface import BillService
from windup_app.server.bill.model import (
Order,
PaymentEvent,
Product,
Refund,
Subscription,
)
from windup_app.server.bill.provider import (
PaymentProvider,
ProviderCreateParams,
ProviderCreateResult,
ProviderNotifyResult,
ProviderQueryResult,
get_provider,
register_provider,
)
from windup_app.server.bill.service import (
SqlAlchemyBillService,
service,
)

__all__ = [
"BillService",
"SqlAlchemyBillService",
"service",
"Product",
"Order",
"PaymentEvent",
"Subscription",
"Refund",
"PaymentProvider",
"ProviderCreateParams",
"ProviderCreateResult",
"ProviderNotifyResult",
"ProviderQueryResult",
"EpayProvider",
"register_provider",
"get_provider",
]

Loading
Loading