Compare commits

..

10 Commits

Author SHA1 Message Date
shanshanzhong147 af22430101 fix: 修复订单 claim 机制状态冲突与 stuck 订单恢复
Build docker and publish / build (20.15.1) (pull_request) Failing after 8m15s
- 将 OrderStatusClaimed 从 4 改为 6,与 OrderStatusFailed(4) 区分
- releaseClaim 改为返回 error,失败时不再静默吞掉
- claimAndGetOrder 对 status=claimed 返回可重试错误而非静默跳过
- ProcessTask 区分 claimed stuck 错误(触发重试)和其他非 paid 状态(跳过)
- finalizeCouponAndOrder 使用直接 DB 更新 claimed→finished,绕过
  UpdateOrderStatus 的 status<target 守卫(6 > 5 无法通过该条件)
  并显式删除 Redis 缓存避免缓存脏读
- 新增 StuckOrderRecoveryLogic:每 10 分钟扫描超时 claimed 订单,
  重置为 paid 并重新入队激活任务

Co-authored-by: multica-agent <github@multica.ai>
2026-05-24 23:42:37 -07:00
shanshanzhong147 fd522b6c71 path
Build docker and publish / build (20.15.1) (push) Successful in 8m56s
2026-05-19 19:35:44 -07:00
shanshanzhong147 a1184ef5ed feat: add userinfo bind-email trial use status
Build docker and publish / build (20.15.1) (push) Successful in 9m2s
2026-05-18 03:02:35 -07:00
shanshanzhong147 7c6efe9dfe x
Build docker and publish / build (20.15.1) (push) Failing after 9m23s
2026-05-16 23:58:53 -07:00
shanshanzhong147 c0ece054a0 x
Build docker and publish / build (20.15.1) (push) Failing after 9m44s
2026-05-16 04:13:05 -07:00
shanshanzhong147 0d57450283 x
Build docker and publish / build (20.15.1) (push) Has been cancelled
2026-05-15 09:53:15 -07:00
shanshanzhong147 f2033fd4b9 feat: add rustfs direct upload endpoint
Build docker and publish / build (20.15.1) (push) Failing after 8m36s
2026-05-14 23:33:50 -07:00
shanshanzhong147 3284cb45f0 ci: trigger internal branch workflow
Build docker and publish / build (20.15.1) (push) Failing after 8m36s
2026-05-14 05:58:05 -07:00
shanshanzhong147 6041bc3419 x
Build docker and publish / build (20.15.1) (push) Has been cancelled
2026-05-14 05:50:44 -07:00
shanshanzhong147 4581a6fc17 ci: use ssh key for main deploy
Build docker and publish / build (20.15.1) (push) Failing after 8m20s
2026-05-13 11:43:01 -07:00
54 changed files with 1663 additions and 222 deletions
@@ -0,0 +1,25 @@
package main
import (
"fmt"
"github.com/perfect-panel/server/internal/config"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/pkg/conf"
)
func main() {
var c config.Config
conf.MustLoad("/private/tmp/ppanel-local-upload.yaml", &c)
ctx := svc.NewServiceContext(c)
const key = "cache:auth:method:device"
before, _ := ctx.Redis.Get(ctx.DB.Statement.Context, key).Result()
fmt.Printf("cache before=%q\n", before)
m, err := ctx.AuthModel.FindOneByMethod(ctx.DB.Statement.Context, "device")
fmt.Printf("model err=%v enabled_nil=%v", err, m == nil || m.Enabled == nil)
if m != nil && m.Enabled != nil { fmt.Printf(" enabled=%v", *m.Enabled) }
if m != nil { fmt.Printf(" config=%s", m.Config) }
fmt.Println()
after, _ := ctx.Redis.Get(ctx.DB.Statement.Context, key).Result()
fmt.Printf("cache after=%q\n", after)
}
+23
View File
@@ -0,0 +1,23 @@
package main
import (
"fmt"
initpkg "github.com/perfect-panel/server/initialize"
"github.com/perfect-panel/server/internal/config"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/pkg/conf"
)
func main() {
var c config.Config
conf.MustLoad("/private/tmp/ppanel-local-upload.yaml", &c)
ctx := svc.NewServiceContext(c)
method, err := ctx.AuthModel.FindOneByMethod(ctx.DB.Statement.Context, "device")
if err != nil {
panic(err)
}
fmt.Printf("db auth_method.enabled=%v config=%s\n", *method.Enabled, method.Config)
initpkg.Device(ctx)
fmt.Printf("ctx.Config.Device.Enable=%v SecuritySecret=%q EnableSecurity=%v\n", ctx.Config.Device.Enable, ctx.Config.Device.SecuritySecret, ctx.Config.Device.EnableSecurity)
}
+36
View File
@@ -0,0 +1,36 @@
package main
import (
"fmt"
"github.com/perfect-panel/server/internal/config"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/pkg/conf"
)
type row struct {
ID int64
Method string
Enabled int
Config string
}
func main() {
var c config.Config
conf.MustLoad("/private/tmp/ppanel-local-upload.yaml", &c)
ctx := svc.NewServiceContext(c)
var rows []row
if err := ctx.DB.Raw("SELECT id, method, enabled, config FROM auth_method WHERE method = ?", "device").Scan(&rows).Error; err != nil {
panic(err)
}
fmt.Printf("raw rows: %+v\n", rows)
m, err := ctx.AuthModel.FindOneByMethod(ctx.DB.Statement.Context, "device")
fmt.Printf("model err=%v\n", err)
if err == nil && m != nil && m.Enabled != nil {
fmt.Printf("model row: id=%d method=%s enabled=%v config=%s\n", m.Id, m.Method, *m.Enabled, m.Config)
} else {
fmt.Printf("model row nil or no enabled ptr: %#v\n", m)
}
}
+28
View File
@@ -0,0 +1,28 @@
package main
import (
"fmt"
"github.com/perfect-panel/server/internal/config"
authmodel "github.com/perfect-panel/server/internal/model/auth"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/pkg/conf"
)
func main() {
var c config.Config
conf.MustLoad("/private/tmp/ppanel-local-upload.yaml", &c)
ctx := svc.NewServiceContext(c)
var a1 authmodel.Auth
err1 := ctx.DB.Model(&authmodel.Auth{}).Where("method = ?", "device").First(&a1).Error
fmt.Printf("gorm direct err=%v enabled_nil=%v", err1, a1.Enabled == nil)
if a1.Enabled != nil { fmt.Printf(" enabled=%v", *a1.Enabled) }
fmt.Printf(" config=%s\n", a1.Config)
var a2 authmodel.Auth
err2 := ctx.DB.Table("auth_method").Where("method = ?", "device").First(&a2).Error
fmt.Printf("gorm table err=%v enabled_nil=%v", err2, a2.Enabled == nil)
if a2.Enabled != nil { fmt.Printf(" enabled=%v", *a2.Enabled) }
fmt.Printf(" config=%s\n", a2.Config)
}
+27 -24
View File
@@ -14,11 +14,12 @@ on:
env: env:
# Docker镜像仓库 # Docker镜像仓库
REPO: ${{ vars.REPO || 'registry.kxsw.us/vpn-server' }} REPO: ${{ vars.REPO || 'registry.kxsw.us/vpn-server' }}
# SSH连接信息 (根据分支自动选择) # SSH连接信息 (根据分支自动选择服务器和用户)
SSH_HOST: ${{ github.ref_name == 'main' && vars.SSH_HOST || vars.DEV_SSH_HOST }} SSH_HOST: ${{ github.ref_name == 'main' && vars.SSH_HOST || vars.DEV_SSH_HOST }}
SSH_PORT: ${{ vars.SSH_PORT }} SSH_PORT: ${{ vars.SSH_PORT }}
SSH_USER: ${{ vars.SSH_USER }} SSH_USER: ${{ github.ref_name == 'main' && 'ubuntu' || 'root' }}
SSH_PASSWORD: ${{ github.ref_name == 'main' && vars.SSH_PASSWORD || vars.DEV_SSH_PASSWORD }} # SSH私钥(Gitea Secret 名称:AWS
SSH_KEY: ${{ secrets.AWS }}
# TG通知 # TG通知
TG_BOT_TOKEN: 8114337882:AAHkEx03HSu7RxN4IHBJJEnsK9aPPzNLIk0 TG_BOT_TOKEN: 8114337882:AAHkEx03HSu7RxN4IHBJJEnsK9aPPzNLIk0
TG_CHAT_ID: "-4940243803" TG_CHAT_ID: "-4940243803"
@@ -49,12 +50,12 @@ jobs:
if [ "${{ github.ref_name }}" = "main" ]; then if [ "${{ github.ref_name }}" = "main" ]; then
echo "DOCKER_TAG_SUFFIX=latest" >> $GITHUB_ENV echo "DOCKER_TAG_SUFFIX=latest" >> $GITHUB_ENV
echo "CONTAINER_NAME=ppanel-server" >> $GITHUB_ENV echo "CONTAINER_NAME=ppanel-server" >> $GITHUB_ENV
echo "DEPLOY_PATH=/root/hifast" >> $GITHUB_ENV echo "DEPLOY_PATH=/opt/ppanel" >> $GITHUB_ENV
echo "为 main 分支设置生产环境变量" echo "为 main 分支设置生产环境变量"
elif [ "${{ github.ref_name }}" = "internal" ]; then elif [ "${{ github.ref_name }}" = "internal" ]; then
echo "DOCKER_TAG_SUFFIX=internal" >> $GITHUB_ENV echo "DOCKER_TAG_SUFFIX=internal" >> $GITHUB_ENV
echo "CONTAINER_NAME=ppanel-server-internal" >> $GITHUB_ENV echo "CONTAINER_NAME=ppanel-server-internal" >> $GITHUB_ENV
echo "DEPLOY_PATH=/root/hifast" >> $GITHUB_ENV echo "DEPLOY_PATH=/root/bindbox" >> $GITHUB_ENV
echo "为 internal 分支设置开发环境变量" echo "为 internal 分支设置开发环境变量"
else else
echo "DOCKER_TAG_SUFFIX=${{ github.ref_name }}" >> $GITHUB_ENV echo "DOCKER_TAG_SUFFIX=${{ github.ref_name }}" >> $GITHUB_ENV
@@ -137,17 +138,15 @@ jobs:
echo "镜像推送完成" echo "镜像推送完成"
# 调试: 打印 SSH 连接信息 # 调试: 打印部署目标(不输出敏感信息
- name: 🔍 调试 - 打印 SSH 连接信息 - name: 🔍 调试 - 打印部署目标
run: | run: |
echo "========== SSH 连接信息调试 ==========" echo "========== 部署目标调试 =========="
echo "当前分支: ${{ github.ref_name }}" echo "当前分支: ${{ github.ref_name }}"
echo "SSH_HOST: ${{ env.SSH_HOST }}" echo "SSH_HOST: ${{ env.SSH_HOST }}"
echo "SSH_PORT: ${{ env.SSH_PORT }}" echo "SSH_PORT: ${{ env.SSH_PORT }}"
echo "SSH_USER: ${{ env.SSH_USER }}" echo "SSH_USER: ${{ env.SSH_USER }}"
echo "SSH_PASSWORD 长度: ${#SSH_PASSWORD}" echo "SSH认证方式: 私钥 (AWS)"
echo "SSH_PASSWORD 前3位: $(echo "$SSH_PASSWORD" | cut -c1-3)***"
echo "SSH_PASSWORD 完整值: ${{ env.SSH_PASSWORD }}"
echo "DEPLOY_PATH: ${{ env.DEPLOY_PATH }}" echo "DEPLOY_PATH: ${{ env.DEPLOY_PATH }}"
echo "=====================================" echo "====================================="
@@ -157,10 +156,10 @@ jobs:
with: with:
host: ${{ env.SSH_HOST }} host: ${{ env.SSH_HOST }}
username: ${{ env.SSH_USER }} username: ${{ env.SSH_USER }}
password: ${{ env.SSH_PASSWORD }} key: ${{ env.SSH_KEY }}
port: ${{ env.SSH_PORT }} port: ${{ env.SSH_PORT }}
source: "docker-compose.cloud.yml" source: "docker-compose.cloud.yml"
target: "${{ env.DEPLOY_PATH }}/" target: "/tmp/ppanel-deploy/"
# 步骤6: 连接服务器更新并启动 # 步骤6: 连接服务器更新并启动
- name: 🚀 连接服务器更新并启动 - name: 🚀 连接服务器更新并启动
@@ -168,7 +167,7 @@ jobs:
with: with:
host: ${{ env.SSH_HOST }} host: ${{ env.SSH_HOST }}
username: ${{ env.SSH_USER }} username: ${{ env.SSH_USER }}
password: ${{ env.SSH_PASSWORD }} key: ${{ env.SSH_KEY }}
port: ${{ env.SSH_PORT }} port: ${{ env.SSH_PORT }}
timeout: 300s timeout: 300s
command_timeout: 600s command_timeout: 600s
@@ -176,23 +175,27 @@ jobs:
echo "连接服务器成功,开始部署..." echo "连接服务器成功,开始部署..."
echo "部署目录: ${{ env.DEPLOY_PATH }}" echo "部署目录: ${{ env.DEPLOY_PATH }}"
echo "部署标签: ${{ env.DOCKER_TAG_SUFFIX }}" echo "部署标签: ${{ env.DOCKER_TAG_SUFFIX }}"
echo "登录用户: ${{ env.SSH_USER }}"
# 进入部署目录 if [ "${{ github.ref_name }}" = "main" ]; then
sudo mkdir -p ${{ env.DEPLOY_PATH }}
sudo cp /tmp/ppanel-deploy/docker-compose.cloud.yml ${{ env.DEPLOY_PATH }}/docker-compose.cloud.yml
cd ${{ env.DEPLOY_PATH }}
echo "📥 拉取镜像..."
sudo docker-compose -f docker-compose.cloud.yml pull ppanel-server
echo "🚀 启动服务..."
sudo docker-compose -f docker-compose.cloud.yml up -d ppanel-server
sudo docker image prune -f || true
else
mkdir -p ${{ env.DEPLOY_PATH }}
cp /tmp/ppanel-deploy/docker-compose.cloud.yml ${{ env.DEPLOY_PATH }}/docker-compose.cloud.yml
cd ${{ env.DEPLOY_PATH }} cd ${{ env.DEPLOY_PATH }}
# 创建/更新环境变量文件
# echo "PPANEL_SERVER_TAG=${{ env.DOCKER_TAG_SUFFIX }}" > .env
# 拉取最新镜像
echo "📥 拉取镜像..." echo "📥 拉取镜像..."
docker-compose -f docker-compose.cloud.yml pull ppanel-server docker-compose -f docker-compose.cloud.yml pull ppanel-server
# 启动服务
echo "🚀 启动服务..." echo "🚀 启动服务..."
docker-compose -f docker-compose.cloud.yml up -d ppanel-server docker-compose -f docker-compose.cloud.yml up -d ppanel-server
# 清理未使用的镜像
docker image prune -f || true docker image prune -f || true
fi
echo "✅ 部署命令执行完成" echo "✅ 部署命令执行完成"
+76
View File
@@ -0,0 +1,76 @@
upstream api_backend {
server 127.0.0.1:8080;
}
server {
listen 80;
server_name api.hifast.biz 4d3vsw888xgaen.hifast.biz;
location /.well-known/acme-challenge/ {
root /var/www/html;
allow all;
}
location / {
return 301 https://$host$request_uri;
}
}
server {
listen 443 ssl http2;
server_name api.hifast.biz;
client_max_body_size 150M;
ssl_certificate /etc/nginx/ssl/hifast.biz/_.hifast.biz.pem; # managed by Certbot
ssl_certificate_key /etc/nginx/ssl/hifast.biz/_.hifast.biz.key; # managed by Certbot
add_header Strict-Transport-Security "max-age=31536000; includeSubDomains" always;
add_header X-Frame-Options "DENY";
add_header X-Content-Type-Options nosniff;
if ($http_user_agent ~* '(9999|91\.78)') {
return 444;
}
location ~ ^/v1/common/client/download/file/(?<download_name>Hi快VPN-(?<os>windows|mac|android)-1.0.0-ic-.*\.(?<ext>exe|dmg|apk))$ {
alias /var/www/download/Hi快VPN-$os-1.0.0.$ext;
charset utf-8;
add_header Content-Disposition 'attachment; filename="$download_name"';
add_header Content-Type application/octet-stream;
}
location /v1/common/client/download/file/ {
alias /var/www/download/;
}
location / {
proxy_pass http://api_backend;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
proxy_http_version 1.1;
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
}
}
server {
listen 443 ssl http2;
server_name 4d3vsw888xgaen.hifast.biz;
client_max_body_size 150M;
ssl_certificate /etc/nginx/ssl/hifast.biz/_.hifast.biz.pem; # managed by Certbot
ssl_certificate_key /etc/nginx/ssl/hifast.biz/_.hifast.biz.key; # managed by Certbot
add_header Strict-Transport-Security "max-age=31536000; includeSubDomains" always;
add_header X-Frame-Options DENY;
add_header X-Content-Type-Options nosniff;
gzip on;
gzip_vary on;
gzip_min_length 1024;
gzip_types text/plain text/css text/xml text/javascript application/javascript application/xml+rss application/json image/svg+xml;
root /var/www/admin;
location / {
try_files $uri $uri/ /index.html;
}
}
+29 -1
View File
@@ -230,6 +230,23 @@ type (
FamilyId int64 `json:"family_id" validate:"required,gt=0"` FamilyId int64 `json:"family_id" validate:"required,gt=0"`
Reason string `json:"reason,omitempty"` Reason string `json:"reason,omitempty"`
} }
GetWithdrawalListRequest {
Page int `form:"page"`
Size int `form:"size"`
UserId *int64 `form:"user_id,omitempty"`
Status *uint8 `form:"status,omitempty"`
}
GetWithdrawalListResponse {
List []WithdrawalLog `json:"list"`
Total int64 `json:"total"`
}
ApproveWithdrawalRequest {
WithdrawalId int64 `json:"withdrawal_id" validate:"required,gt=0"`
}
RejectWithdrawalRequest {
WithdrawalId int64 `json:"withdrawal_id" validate:"required,gt=0"`
Reason string `json:"reason" validate:"required,max=500"`
}
) )
@server ( @server (
@@ -370,5 +387,16 @@ service ppanel {
@doc "Dissolve family" @doc "Dissolve family"
@handler DissolveFamily @handler DissolveFamily
put /family/dissolve (DissolveFamilyRequest) put /family/dissolve (DissolveFamilyRequest)
}
@doc "Get withdrawal list"
@handler GetWithdrawalList
get /withdrawal/list (GetWithdrawalListRequest) returns (GetWithdrawalListResponse)
@doc "Approve withdrawal"
@handler ApproveWithdrawal
post /withdrawal/approve (ApproveWithdrawalRequest)
@doc "Reject withdrawal"
@handler RejectWithdrawal
post /withdrawal/reject (RejectWithdrawalRequest)
}
+18
View File
@@ -11,6 +11,20 @@ info (
import "../types.api" import "../types.api"
type ( type (
FileUploadRequest {
BizType string `form:"biz_type" validate:"required"`
}
FileUploadResponse {
FileId string `json:"file_id"`
FileName string `json:"file_name"`
ObjectKey string `json:"object_key"`
Size int64 `json:"size"`
ContentType string `json:"content_type"`
Etag string `json:"etag"`
Status string `json:"status"`
}
FileUploadInitRequest { FileUploadInitRequest {
BizType string `json:"biz_type" validate:"required"` BizType string `json:"biz_type" validate:"required"`
FileName string `json:"file_name" validate:"required"` FileName string `json:"file_name" validate:"required"`
@@ -48,6 +62,10 @@ type (
middleware: AuthMiddleware,DeviceMiddleware middleware: AuthMiddleware,DeviceMiddleware
) )
service ppanel { service ppanel {
@doc "Upload file to RustFS"
@handler FileUpload
post /upload (FileUploadRequest) returns (FileUploadResponse)
@doc "Init file upload" @doc "Init file upload"
@handler FileUploadInit @handler FileUploadInit
post /upload/init (FileUploadInitRequest) returns (FileUploadInitResponse) post /upload/init (FileUploadInitRequest) returns (FileUploadInitResponse)
+1 -1
View File
@@ -27,6 +27,7 @@ type (
EnableLoginNotify bool `json:"enable_login_notify"` EnableLoginNotify bool `json:"enable_login_notify"`
EnableSubscribeNotify bool `json:"enable_subscribe_notify"` EnableSubscribeNotify bool `json:"enable_subscribe_notify"`
EnableTradeNotify bool `json:"enable_trade_notify"` EnableTradeNotify bool `json:"enable_trade_notify"`
UseStatus bool `json:"use_status"` // Whether to show the "bind email to get free trial" prompt
AuthMethods []UserAuthMethod `json:"auth_methods"` AuthMethods []UserAuthMethod `json:"auth_methods"`
UserDevices []UserDevice `json:"user_devices"` UserDevices []UserDevice `json:"user_devices"`
Rules []string `json:"rules"` Rules []string `json:"rules"`
@@ -1004,4 +1005,3 @@ type (
ConfigSnapshot map[string]interface{} `json:"config_snapshot,omitempty"` ConfigSnapshot map[string]interface{} `json:"config_snapshot,omitempty"`
} }
) )
@@ -17,7 +17,7 @@ REDIS_HOST=127.0.0.1
REDIS_PORT=6379 REDIS_PORT=6379
REDIS_PASSWORD=CHANGE_ME REDIS_PASSWORD=CHANGE_ME
REDIS_SOURCE_HOST=43.198.248.161 REDIS_SOURCE_HOST=18.163.33.75
REDIS_SOURCE_PORT=6379 REDIS_SOURCE_PORT=6379
REDIS_SOURCE_USER= REDIS_SOURCE_USER=
REDIS_SOURCE_PASSWORD=CHANGE_ME REDIS_SOURCE_PASSWORD=CHANGE_ME
+184
View File
@@ -0,0 +1,184 @@
# TAPI 文件上传接入说明
本文档说明 `https://tapi.hifast.biz/v1/public/file/upload` 相关上传接口的推荐接入方式、签名规则与常见排查方式。
## 总览
上传能力包含两类接入方式:
- 推荐方式:`init -> S3 PUT -> complete`
- 兼容方式:`/upload` multipart 直传
推荐优先使用预签名三段式,因为:
- 现有签名串包含 `BODY_SHA256`
- `/upload``multipart/form-data`
- multipart 原始 body 的签名和调试成本更高
- `init``complete` 是 JSON,更适合客户端和 Apifox 调试
## 签名生效逻辑
项目保持现有旧逻辑,不做强制签名改造:
- `Signature.EnableSignature = false` 时:不校验签名
- `Signature.EnableSignature = true` 且未携带 `X-App-Id` 时:不校验签名,兼容老客户端
- `Signature.EnableSignature = true` 且携带 `X-App-Id` 时:必须同时携带并校验
- `X-Timestamp`
- `X-Nonce`
- `X-Signature`
这意味着:
- 新客户端建议始终带完整签名头
- 老客户端如果没有 `X-App-Id`,仍可按旧逻辑访问
## 签名头定义
- `X-App-Id`: 客户端标识,例如 `ios-client`
- `X-Timestamp`: Unix 秒级时间戳
- `X-Nonce`: 每次请求唯一随机串
- `X-Signature`: `HMAC-SHA256` 结果的十六进制小写字符串
## StringToSign 规则
StringToSign 由下面 7 段按换行符 `\n` 拼接:
```text
METHOD
PATH
CANONICAL_QUERY
BODY_SHA256
X-App-Id
X-Timestamp
X-Nonce
```
说明:
- `METHOD`HTTP 方法大写,例如 `POST`
- `PATH`:请求路径,例如 `/v1/public/file/upload/init`
- `CANONICAL_QUERY`:按 key 排序后的 query string,没有 query 则为空字符串
- `BODY_SHA256`:请求体原始字节的 SHA-256 十六进制小写
- 其余三项直接使用请求头值
签名计算方式:
```text
signature = hex_lower(HMAC_SHA256(app_secret, string_to_sign))
```
时间窗与防重放:
- `X-Timestamp` 默认有效时间窗是 300 秒
- `X-Nonce` 在有效时间窗内不能重复使用
## 推荐接入:预签名三段式
### 1. 初始化上传
请求:
```bash
curl -X POST 'https://tapi.hifast.biz/v1/public/file/upload/init' \
-H 'Accept: application/json, text/plain, */*' \
-H 'Content-Type: application/json' \
-H 'authorization: your-token' \
-H 'X-App-Id: ios-client' \
-H 'X-Timestamp: 1778776400' \
-H 'X-Nonce: nonce-001' \
-H 'X-Signature: your-signature' \
-d '{
"biz_type": "app-package",
"file_name": "demo.zip",
"content_type": "application/zip",
"size": 123456,
"sha256": ""
}'
```
典型返回:
```json
{
"code": 200,
"msg": "success",
"data": {
"file_id": "c29274ee26ab5aa211e0396e",
"object_key": "app-upload/app-package/519/2026/05/c29274ee26ab5aa211e0396e_demo.zip",
"upload_url": "https://bucket.s3.ap-east-1.amazonaws.com/...",
"method": "PUT",
"headers": {
"Content-Type": "application/zip"
},
"expired_at": 1778776715
}
}
```
### 2. 直传 S3
这一步是直接上传二进制文件到 S3,不走业务签名中间件。
```bash
curl -X PUT 'https://bucket.s3.ap-east-1.amazonaws.com/...' \
-H 'Content-Type: application/zip' \
--upload-file '/tmp/demo.zip'
```
说明:
- `Content-Type` 需和 `init` 返回的 `headers.Content-Type` 一致
- `upload_url` 有过期时间,通常 300 秒
- 成功时 S3 常见返回 `200``204`
### 3. 完成上传
```bash
curl -X POST 'https://tapi.hifast.biz/v1/public/file/upload/complete' \
-H 'Accept: application/json, text/plain, */*' \
-H 'Content-Type: application/json' \
-H 'authorization: your-token' \
-H 'X-App-Id: ios-client' \
-H 'X-Timestamp: 1778776405' \
-H 'X-Nonce: nonce-002' \
-H 'X-Signature: your-signature' \
-d '{
"file_id": "c29274ee26ab5aa211e0396e"
}'
```
## 兼容接入:单接口 multipart 直传
接口:
- `POST /v1/public/file/upload`
表单字段:
- `biz_type`
- `file`
说明:
- 该接口继续保留,兼容旧客户端
- 如果请求带了 `X-App-Id`,就按现有逻辑验签
- 如果没有 `X-App-Id`,仍按旧逻辑放行
- 如果要给该接口加签,签名时必须对原始 multipart body 计算 `BODY_SHA256`
## 常见错误码
- `200`: 成功
- `400`: 参数错误
- `40008`: 缺少签名头
- `40009`: 签名已过期
- `40010`: 签名无效
- `40011`: nonce 重放
- `10001`: 上传元数据不存在或对象不存在
## 排查建议
- `40008`:确认带了 `X-App-Id` 后,也同时带上 `X-Timestamp / X-Nonce / X-Signature`
- `40009`:检查客户端时间是否偏差过大
- `40010`:确认 `PATH`、query 排序、body 原始字节、secret 是否完全一致
- `40011`:确保每次请求都生成新的 `X-Nonce`
- `complete` 失败:确认 S3 `PUT` 已成功,且上传大小与 `init.size` 一致
@@ -34,7 +34,7 @@
}, },
"id": 1, "id": 1,
"options": { "options": {
"content": "<b>AWS CloudWatch overview</b><br/>Region: ap-east-1 (Hong Kong)<br/>RDS DBInstanceIdentifier: database-1<br/>Redis ReplicationGroupId: hifastapp-redis<br/><br/>Note: Redis panels are configured using <code>ReplicationGroupId</code>. If a panel shows no data in your account, switch that dimension to the concrete <code>CacheClusterId</code> in Grafana query editor.", "content": "<b>AWS CloudWatch overview</b><br/>Region: ap-east-1 (Hong Kong)<br/>RDS DBInstanceIdentifier: hifast-mysql-prod-v2<br/>Redis: current production uses a local Docker Redis container (<code>hifast-redis</code>) on EC2 rather than AWS ElastiCache.<br/><br/>This dashboard keeps the RDS CloudWatch panels. Redis should be observed from the local ops dashboard via Prometheus/cAdvisor instead of ElastiCache metrics.",
"mode": "html" "mode": "html"
}, },
"pluginVersion": "11.0.0", "pluginVersion": "11.0.0",
@@ -75,7 +75,7 @@
"uid": "cloudwatch" "uid": "cloudwatch"
}, },
"dimensions": { "dimensions": {
"DBInstanceIdentifier": "database-1" "DBInstanceIdentifier": "hifast-mysql-prod-v2"
}, },
"metricName": "CPUUtilization", "metricName": "CPUUtilization",
"namespace": "AWS/RDS", "namespace": "AWS/RDS",
@@ -122,7 +122,7 @@
"uid": "cloudwatch" "uid": "cloudwatch"
}, },
"dimensions": { "dimensions": {
"DBInstanceIdentifier": "database-1" "DBInstanceIdentifier": "hifast-mysql-prod-v2"
}, },
"metricName": "DatabaseConnections", "metricName": "DatabaseConnections",
"namespace": "AWS/RDS", "namespace": "AWS/RDS",
@@ -169,7 +169,7 @@
"uid": "cloudwatch" "uid": "cloudwatch"
}, },
"dimensions": { "dimensions": {
"DBInstanceIdentifier": "database-1" "DBInstanceIdentifier": "hifast-mysql-prod-v2"
}, },
"metricName": "FreeStorageSpace", "metricName": "FreeStorageSpace",
"namespace": "AWS/RDS", "namespace": "AWS/RDS",
@@ -216,7 +216,7 @@
"uid": "cloudwatch" "uid": "cloudwatch"
}, },
"dimensions": { "dimensions": {
"DBInstanceIdentifier": "database-1" "DBInstanceIdentifier": "hifast-mysql-prod-v2"
}, },
"metricName": "ReadLatency", "metricName": "ReadLatency",
"namespace": "AWS/RDS", "namespace": "AWS/RDS",
@@ -231,7 +231,7 @@
"uid": "cloudwatch" "uid": "cloudwatch"
}, },
"dimensions": { "dimensions": {
"DBInstanceIdentifier": "database-1" "DBInstanceIdentifier": "hifast-mysql-prod-v2"
}, },
"metricName": "WriteLatency", "metricName": "WriteLatency",
"namespace": "AWS/RDS", "namespace": "AWS/RDS",
@@ -278,7 +278,7 @@
"uid": "cloudwatch" "uid": "cloudwatch"
}, },
"dimensions": { "dimensions": {
"DBInstanceIdentifier": "database-1" "DBInstanceIdentifier": "hifast-mysql-prod-v2"
}, },
"metricName": "ReadIOPS", "metricName": "ReadIOPS",
"namespace": "AWS/RDS", "namespace": "AWS/RDS",
@@ -293,7 +293,7 @@
"uid": "cloudwatch" "uid": "cloudwatch"
}, },
"dimensions": { "dimensions": {
"DBInstanceIdentifier": "database-1" "DBInstanceIdentifier": "hifast-mysql-prod-v2"
}, },
"metricName": "WriteIOPS", "metricName": "WriteIOPS",
"namespace": "AWS/RDS", "namespace": "AWS/RDS",
@@ -350,7 +350,7 @@
"statistic": "Average" "statistic": "Average"
} }
], ],
"title": "Redis Host CPU", "title": "Redis Host CPU (Legacy ElastiCache)",
"type": "timeseries" "type": "timeseries"
}, },
{ {
@@ -397,7 +397,7 @@
"statistic": "Average" "statistic": "Average"
} }
], ],
"title": "Redis Engine CPU", "title": "Redis Engine CPU (Legacy ElastiCache)",
"type": "timeseries" "type": "timeseries"
}, },
{ {
@@ -444,7 +444,7 @@
"statistic": "Average" "statistic": "Average"
} }
], ],
"title": "Redis Connections", "title": "Redis Connections (Legacy ElastiCache)",
"type": "timeseries" "type": "timeseries"
}, },
{ {
@@ -491,7 +491,7 @@
"statistic": "Average" "statistic": "Average"
} }
], ],
"title": "Redis Memory Usage %", "title": "Redis Memory Usage % (Legacy ElastiCache)",
"type": "timeseries" "type": "timeseries"
} }
], ],
@@ -58,7 +58,7 @@
"uid": "P8E80F9AEF21F6940" "uid": "P8E80F9AEF21F6940"
}, },
"editorMode": "code", "editorMode": "code",
"expr": "sum(count_over_time({container=\"ppanel-server\"} |~ \"${user_id:raw}\" |~ \"${email:raw}\" |~ \"${order:raw}\" |~ \"${keyword:raw}\" |~ \"${level:raw}\" [5m]))", "expr": "sum(count_over_time({compose_service=\"ppanel-server\"}[5m]))",
"queryType": "range", "queryType": "range",
"refId": "A" "refId": "A"
} }
@@ -103,7 +103,7 @@
"uid": "P8E80F9AEF21F6940" "uid": "P8E80F9AEF21F6940"
}, },
"editorMode": "code", "editorMode": "code",
"expr": "sum(count_over_time({container=\"ppanel-server\"} |~ \"${user_id:raw}\" |~ \"${email:raw}\" |~ \"${order:raw}\" |~ \"${keyword:raw}\" |~ \"(?i)(error|panic|fatal)\" [5m]))", "expr": "sum(count_over_time({compose_service=\"ppanel-server\"} |~ \"(?i)(error|panic|fatal)\" [5m]))",
"queryType": "range", "queryType": "range",
"refId": "A" "refId": "A"
} }
@@ -140,7 +140,7 @@
"uid": "P8E80F9AEF21F6940" "uid": "P8E80F9AEF21F6940"
}, },
"editorMode": "code", "editorMode": "code",
"expr": "{container=\"ppanel-server\"} |~ \"${user_id:raw}\" |~ \"${email:raw}\" |~ \"${order:raw}\" |~ \"${keyword:raw}\" |~ \"${level:raw}\"", "expr": "{compose_service=\"ppanel-server\"}",
"queryType": "range", "queryType": "range",
"refId": "A" "refId": "A"
} }
@@ -177,7 +177,7 @@
"uid": "P8E80F9AEF21F6940" "uid": "P8E80F9AEF21F6940"
}, },
"editorMode": "code", "editorMode": "code",
"expr": "{container=\"ppanel-server\"} |~ \"${user_id:raw}\" |~ \"${email:raw}\" |~ \"${order:raw}\" |~ \"${keyword:raw}\" |~ \"(?i)(error|panic|fatal)\"", "expr": "{compose_service=\"ppanel-server\"} |~ \"(?i)(error|panic|fatal)\"",
"queryType": "range", "queryType": "range",
"refId": "A" "refId": "A"
} }
@@ -318,7 +318,7 @@
] ]
}, },
"time": { "time": {
"from": "now-7d", "from": "now-1h",
"to": "now" "to": "now"
}, },
"timepicker": {}, "timepicker": {},
@@ -449,7 +449,7 @@ CREATE TABLE IF NOT EXISTS `user_device`
`subscribe_id` bigint DEFAULT NULL COMMENT 'Subscribe ID', `subscribe_id` bigint DEFAULT NULL COMMENT 'Subscribe ID',
`ip` varchar(191) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT 'Device Ip.', `ip` varchar(191) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT 'Device Ip.',
`Identifier` varchar(191) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT 'Device Identifier.', `Identifier` varchar(191) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT 'Device Identifier.',
`user_agent` varchar(64) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT 'Device User Agent.', `user_agent` varchar(255) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci DEFAULT NULL COMMENT 'Device User Agent.',
`online` tinyint(1) NOT NULL DEFAULT '0' COMMENT 'Online', `online` tinyint(1) NOT NULL DEFAULT '0' COMMENT 'Online',
`enabled` tinyint(1) NOT NULL DEFAULT '1' COMMENT 'EnableDeviceNumber', `enabled` tinyint(1) NOT NULL DEFAULT '1' COMMENT 'EnableDeviceNumber',
`created_at` datetime(3) DEFAULT NULL COMMENT 'Creation Time', `created_at` datetime(3) DEFAULT NULL COMMENT 'Creation Time',
@@ -0,0 +1,7 @@
ALTER TABLE `log_message`
MODIFY COLUMN `app_version` VARCHAR(32) NULL,
MODIFY COLUMN `os_name` VARCHAR(32) NULL,
MODIFY COLUMN `os_version` VARCHAR(32) NULL,
MODIFY COLUMN `device_id` VARCHAR(64) NULL,
MODIFY COLUMN `session_id` VARCHAR(64) NULL,
MODIFY COLUMN `error_code` VARCHAR(64) NULL;
@@ -0,0 +1,7 @@
ALTER TABLE `log_message`
MODIFY COLUMN `app_version` VARCHAR(64) NULL,
MODIFY COLUMN `os_name` VARCHAR(64) NULL,
MODIFY COLUMN `os_version` VARCHAR(64) NULL,
MODIFY COLUMN `device_id` VARCHAR(255) NULL,
MODIFY COLUMN `session_id` VARCHAR(255) NULL,
MODIFY COLUMN `error_code` VARCHAR(128) NULL;
@@ -0,0 +1,2 @@
ALTER TABLE `user_device`
MODIFY COLUMN `user_agent` VARCHAR(64) NULL COMMENT 'Device User Agent.';
@@ -0,0 +1,2 @@
ALTER TABLE `user_device`
MODIFY COLUMN `user_agent` VARCHAR(255) NULL COMMENT 'Device User Agent.';
+72
View File
@@ -20,6 +20,14 @@ type schemaColumnPatch struct {
ddl string ddl string
} }
type schemaColumnDefinitionPatch struct {
table string
column string
dataType string
characterMaxLen *int64
ddl string
}
func EnsureSchemaCompatibility(ctx *svc.ServiceContext) error { func EnsureSchemaCompatibility(ctx *svc.ServiceContext) error {
tablePatches := []schemaTablePatch{ tablePatches := []schemaTablePatch{
{ {
@@ -142,6 +150,17 @@ func EnsureSchemaCompatibility(ctx *svc.ServiceContext) error {
}, },
} }
varchar255 := int64(255)
columnDefinitionPatches := []schemaColumnDefinitionPatch{
{
table: "user_device",
column: "user_agent",
dataType: "varchar",
characterMaxLen: &varchar255,
ddl: "ALTER TABLE `user_device` MODIFY COLUMN `user_agent` VARCHAR(255) NULL COMMENT 'Device User Agent.';",
},
}
for _, patch := range tablePatches { for _, patch := range tablePatches {
exists, err := tableExists(ctx.DB, patch.table) exists, err := tableExists(ctx.DB, patch.table)
if err != nil { if err != nil {
@@ -199,6 +218,27 @@ func EnsureSchemaCompatibility(ctx *svc.ServiceContext) error {
logger.Infof("[SchemaCompat] created missing index: %s.%s", patch.table, patch.index) logger.Infof("[SchemaCompat] created missing index: %s.%s", patch.table, patch.index)
} }
for _, patch := range columnDefinitionPatches {
tblExists, err := tableExists(ctx.DB, patch.table)
if err != nil {
return errors.Wrapf(err, "check table %s failed", patch.table)
}
if !tblExists {
continue
}
matches, err := columnDefinitionMatches(ctx.DB, patch.table, patch.column, patch.dataType, patch.characterMaxLen)
if err != nil {
return errors.Wrapf(err, "check column definition %s.%s failed", patch.table, patch.column)
}
if matches {
continue
}
if err = ctx.DB.Exec(patch.ddl).Error; err != nil {
return errors.Wrapf(err, "modify column %s.%s failed", patch.table, patch.column)
}
logger.Infof("[SchemaCompat] repaired column definition: %s.%s", patch.table, patch.column)
}
return nil return nil
} }
@@ -237,6 +277,38 @@ func indexExists(db *gorm.DB, table, index string) (bool, error) {
return count > 0, nil return count > 0, nil
} }
func columnDefinitionMatches(db *gorm.DB, table, column, dataType string, characterMaxLen *int64) (bool, error) {
type columnMeta struct {
DataType string
CharacterMaximumLen *int64
}
var meta columnMeta
err := db.Raw(
`SELECT DATA_TYPE AS data_type, CHARACTER_MAXIMUM_LENGTH AS character_maximum_len
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME = ? AND COLUMN_NAME = ?`,
table,
column,
).Scan(&meta).Error
if err != nil {
return false, err
}
if meta.DataType == "" {
return false, nil
}
if meta.DataType != dataType {
return false, nil
}
if characterMaxLen == nil {
return true, nil
}
if meta.CharacterMaximumLen == nil {
return false, nil
}
return *meta.CharacterMaximumLen == *characterMaxLen, nil
}
func _schemaCompatDebug(table, column string) string { func _schemaCompatDebug(table, column string) string {
if column == "" { if column == "" {
return table return table
@@ -0,0 +1,26 @@
package user
import (
"github.com/gin-gonic/gin"
"github.com/perfect-panel/server/internal/logic/admin/user"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/result"
)
// Approve withdrawal
func ApproveWithdrawalHandler(svcCtx *svc.ServiceContext) func(c *gin.Context) {
return func(c *gin.Context) {
var req types.ApproveWithdrawalRequest
_ = c.ShouldBind(&req)
validateErr := svcCtx.Validate(&req)
if validateErr != nil {
result.ParamErrorResult(c, validateErr)
return
}
l := user.NewApproveWithdrawalLogic(c.Request.Context(), svcCtx)
err := l.ApproveWithdrawal(&req)
result.HttpResult(c, nil, err)
}
}
@@ -0,0 +1,26 @@
package user
import (
"github.com/gin-gonic/gin"
"github.com/perfect-panel/server/internal/logic/admin/user"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/result"
)
// Get withdrawal list
func GetWithdrawalListHandler(svcCtx *svc.ServiceContext) func(c *gin.Context) {
return func(c *gin.Context) {
var req types.GetWithdrawalListRequest
_ = c.ShouldBind(&req)
validateErr := svcCtx.Validate(&req)
if validateErr != nil {
result.ParamErrorResult(c, validateErr)
return
}
l := user.NewGetWithdrawalListLogic(c.Request.Context(), svcCtx)
resp, err := l.GetWithdrawalList(&req)
result.HttpResult(c, resp, err)
}
}
@@ -0,0 +1,26 @@
package user
import (
"github.com/gin-gonic/gin"
"github.com/perfect-panel/server/internal/logic/admin/user"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/result"
)
// Reject withdrawal
func RejectWithdrawalHandler(svcCtx *svc.ServiceContext) func(c *gin.Context) {
return func(c *gin.Context) {
var req types.RejectWithdrawalRequest
_ = c.ShouldBind(&req)
validateErr := svcCtx.Validate(&req)
if validateErr != nil {
result.ParamErrorResult(c, validateErr)
return
}
l := user.NewRejectWithdrawalLogic(c.Request.Context(), svcCtx)
err := l.RejectWithdrawal(&req)
result.HttpResult(c, nil, err)
}
}
@@ -0,0 +1,43 @@
package file
import (
"mime/multipart"
"github.com/gin-gonic/gin"
"github.com/perfect-panel/server/internal/logic/public/file"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/result"
)
// Upload file to RustFS
func FileUploadHandler(svcCtx *svc.ServiceContext) func(c *gin.Context) {
return func(c *gin.Context) {
var req types.FileUploadRequest
_ = c.ShouldBind(&req)
validateErr := svcCtx.Validate(&req)
if validateErr != nil {
result.ParamErrorResult(c, validateErr)
return
}
fileHeader, err := c.FormFile("file")
if err != nil {
result.ParamErrorResult(c, err)
return
}
fileReader, err := fileHeader.Open()
if err != nil {
result.HttpResult(c, nil, err)
return
}
defer func(file multipart.File) {
_ = file.Close()
}(fileReader)
l := file.NewFileUploadLogic(c.Request.Context(), svcCtx)
resp, err := l.FileUpload(&req, fileHeader, fileReader)
result.HttpResult(c, resp, err)
}
}
+12
View File
@@ -707,6 +707,15 @@ func RegisterHandlers(router *gin.Engine, serverCtx *svc.ServiceContext) {
// Get admin user invite list // Get admin user invite list
adminUserGroupRouter.GET("/invite/list", adminUser.GetAdminUserInviteListHandler(serverCtx)) adminUserGroupRouter.GET("/invite/list", adminUser.GetAdminUserInviteListHandler(serverCtx))
// Get withdrawal list
adminUserGroupRouter.GET("/withdrawal/list", adminUser.GetWithdrawalListHandler(serverCtx))
// Approve withdrawal
adminUserGroupRouter.POST("/withdrawal/approve", adminUser.ApproveWithdrawalHandler(serverCtx))
// Reject withdrawal
adminUserGroupRouter.POST("/withdrawal/reject", adminUser.RejectWithdrawalHandler(serverCtx))
} }
authGroupRouter := router.Group("/v1/auth") authGroupRouter := router.Group("/v1/auth")
@@ -870,6 +879,9 @@ func RegisterHandlers(router *gin.Engine, serverCtx *svc.ServiceContext) {
publicFileGroupRouter.Use(middleware.AuthMiddleware(serverCtx), middleware.DeviceMiddleware(serverCtx)) publicFileGroupRouter.Use(middleware.AuthMiddleware(serverCtx), middleware.DeviceMiddleware(serverCtx))
{ {
// Upload file to RustFS
publicFileGroupRouter.POST("/upload", publicFile.FileUploadHandler(serverCtx))
// Init file upload // Init file upload
publicFileGroupRouter.POST("/upload/init", publicFile.FileUploadInitHandler(serverCtx)) publicFileGroupRouter.POST("/upload/init", publicFile.FileUploadInitHandler(serverCtx))
@@ -0,0 +1,27 @@
package user
import (
"context"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/logger"
)
type ApproveWithdrawalLogic struct {
logger.Logger
ctx context.Context
svcCtx *svc.ServiceContext
}
func NewApproveWithdrawalLogic(ctx context.Context, svcCtx *svc.ServiceContext) *ApproveWithdrawalLogic {
return &ApproveWithdrawalLogic{
Logger: logger.WithContext(ctx),
ctx: ctx,
svcCtx: svcCtx,
}
}
func (l *ApproveWithdrawalLogic) ApproveWithdrawal(req *types.ApproveWithdrawalRequest) error {
return approveWithdrawal(l.ctx, l.svcCtx, req.WithdrawalId)
}
@@ -0,0 +1,74 @@
package user
import (
"context"
usermodel "github.com/perfect-panel/server/internal/model/user"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/logger"
"github.com/perfect-panel/server/pkg/xerr"
"github.com/pkg/errors"
)
type GetWithdrawalListLogic struct {
logger.Logger
ctx context.Context
svcCtx *svc.ServiceContext
}
func NewGetWithdrawalListLogic(ctx context.Context, svcCtx *svc.ServiceContext) *GetWithdrawalListLogic {
return &GetWithdrawalListLogic{
Logger: logger.WithContext(ctx),
ctx: ctx,
svcCtx: svcCtx,
}
}
func (l *GetWithdrawalListLogic) GetWithdrawalList(req *types.GetWithdrawalListRequest) (*types.GetWithdrawalListResponse, error) {
page := req.Page
size := req.Size
if page <= 0 {
page = 1
}
if size <= 0 {
size = 10
}
query := l.svcCtx.DB.WithContext(l.ctx).Model(&usermodel.Withdrawal{})
if req.UserId != nil {
query = query.Where("user_id = ?", *req.UserId)
}
if req.Status != nil {
query = query.Where("status = ?", *req.Status)
}
var total int64
if err := query.Count(&total).Error; err != nil {
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "count withdrawals failed: %v", err)
}
var rows []usermodel.Withdrawal
if err := query.Order("id DESC").Limit(size).Offset((page - 1) * size).Find(&rows).Error; err != nil {
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query withdrawals failed: %v", err)
}
list := make([]types.WithdrawalLog, 0, len(rows))
for _, row := range rows {
list = append(list, types.WithdrawalLog{
Id: row.Id,
UserId: row.UserId,
Amount: row.Amount,
Content: row.Content,
Status: row.Status,
Reason: row.Reason,
CreatedAt: row.CreatedAt.UnixMilli(),
UpdatedAt: row.UpdatedAt.UnixMilli(),
})
}
return &types.GetWithdrawalListResponse{
List: list,
Total: total,
}, nil
}
@@ -0,0 +1,27 @@
package user
import (
"context"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/logger"
)
type RejectWithdrawalLogic struct {
logger.Logger
ctx context.Context
svcCtx *svc.ServiceContext
}
func NewRejectWithdrawalLogic(ctx context.Context, svcCtx *svc.ServiceContext) *RejectWithdrawalLogic {
return &RejectWithdrawalLogic{
Logger: logger.WithContext(ctx),
ctx: ctx,
svcCtx: svcCtx,
}
}
func (l *RejectWithdrawalLogic) RejectWithdrawal(req *types.RejectWithdrawalRequest) error {
return rejectWithdrawal(l.ctx, l.svcCtx, req.WithdrawalId, req.Reason)
}
@@ -6,6 +6,7 @@ import (
"strings" "strings"
"time" "time"
logicCommon "github.com/perfect-panel/server/internal/logic/common"
"github.com/perfect-panel/server/internal/model/log" "github.com/perfect-panel/server/internal/model/log"
"github.com/perfect-panel/server/internal/svc" "github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types" "github.com/perfect-panel/server/internal/types"
@@ -100,21 +101,15 @@ func (l *UpdateUserBasicInfoLogic) UpdateUserBasicInfo(req *types.UpdateUserBasi
} }
if req.Commission != userInfo.Commission { if req.Commission != userInfo.Commission {
if isWithdrawalScene(req.Remark) {
commentLog := log.Commission{ logWithdrawalGuard(l.Logger, userInfo.Id)
Type: log.CommissionTypeAdjust, return errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "commission overwrite is blocked in withdrawal scene")
Amount: req.Commission - userInfo.Commission,
Timestamp: time.Now().UnixMilli(),
} }
change := req.Commission - userInfo.Commission
content, _ := commentLog.Marshal() if err = l.svcCtx.UserModel.UpdateCommission(l.ctx, userInfo.Id, change, tx); err != nil {
err = tx.Create(&log.SystemLog{ return err
Type: log.TypeCommission.Uint8(), }
Date: time.Now().Format(time.DateOnly), if err = logicCommon.WriteCommissionLog(tx, userInfo.Id, log.CommissionTypeAdjust, change, ""); err != nil {
ObjectID: userInfo.Id,
Content: string(content),
}).Error
if err != nil {
return err return err
} }
userInfo.Commission = req.Commission userInfo.Commission = req.Commission
@@ -0,0 +1,81 @@
package user
import (
"context"
"strings"
logicCommon "github.com/perfect-panel/server/internal/logic/common"
"github.com/perfect-panel/server/internal/model/log"
usermodel "github.com/perfect-panel/server/internal/model/user"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/pkg/logger"
"github.com/perfect-panel/server/pkg/xerr"
"github.com/pkg/errors"
"gorm.io/gorm"
)
func approveWithdrawal(ctx context.Context, svcCtx *svc.ServiceContext, withdrawalID int64) error {
return svcCtx.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
withdrawal, err := logicCommon.LoadPendingWithdrawalForUpdate(ctx, tx, withdrawalID)
if err != nil {
if err.Error() == "withdrawal status invalid" {
return errors.Wrapf(xerr.NewErrCode(xerr.WithdrawalStatusInvalid), "withdrawal %d already processed", withdrawalID)
}
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "load withdrawal failed: %v", err)
}
if err := tx.Model(&usermodel.Withdrawal{}).
Where("id = ? AND status = 0", withdrawalID).
Updates(map[string]interface{}{
"status": 1,
"reason": "",
}).Error; err != nil {
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "approve withdrawal failed: %v", err)
}
if err := logicCommon.WriteCommissionLog(tx, withdrawal.UserId, log.CommissionTypeWithdraw, withdrawal.Amount, ""); err != nil {
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "write commission log failed: %v", err)
}
return nil
})
}
func rejectWithdrawal(ctx context.Context, svcCtx *svc.ServiceContext, withdrawalID int64, reason string) error {
reason = strings.TrimSpace(reason)
return svcCtx.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
withdrawal, err := logicCommon.LoadPendingWithdrawalForUpdate(ctx, tx, withdrawalID)
if err != nil {
if err.Error() == "withdrawal status invalid" {
return errors.Wrapf(xerr.NewErrCode(xerr.WithdrawalStatusInvalid), "withdrawal %d already processed", withdrawalID)
}
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "load withdrawal failed: %v", err)
}
if err := tx.Model(&usermodel.Withdrawal{}).
Where("id = ? AND status = 0", withdrawalID).
Updates(map[string]interface{}{
"status": 2,
"reason": reason,
}).Error; err != nil {
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "reject withdrawal failed: %v", err)
}
if err := svcCtx.UserModel.UpdateCommission(ctx, withdrawal.UserId, withdrawal.Amount, tx); err != nil {
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "refund commission failed: %v", err)
}
if err := logicCommon.WriteCommissionLog(tx, withdrawal.UserId, log.CommissionTypeWithdrawReject, withdrawal.Amount, ""); err != nil {
return errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "write commission log failed: %v", err)
}
return nil
})
}
func isWithdrawalScene(remark string) bool {
normalized := strings.ToLower(strings.TrimSpace(remark))
return strings.Contains(normalized, "withdraw") || strings.Contains(normalized, "提现")
}
func logWithdrawalGuard(l logger.Logger, userID int64) {
l.Errorw("blocked commission overwrite in withdrawal scene", logger.Field("user_id", userID))
}
+25 -13
View File
@@ -35,9 +35,13 @@ func (l *ReportLogMessageLogic) ReportLogMessage(req *types.ReportLogMessageRequ
// 简单限流:设备ID优先,其次IP // 简单限流:设备ID优先,其次IP
limitKey := "logmsg:" + strings.TrimSpace(req.DeviceId) limitKey := "logmsg:" + strings.TrimSpace(req.DeviceId)
if limitKey == "logmsg:" { limitKey = "logmsg:" + ip } if limitKey == "logmsg:" {
limitKey = "logmsg:" + ip
}
count, _ := l.svcCtx.Redis.Incr(l.ctx, limitKey).Result() count, _ := l.svcCtx.Redis.Incr(l.ctx, limitKey).Result()
if count == 1 { _ = l.svcCtx.Redis.Expire(l.ctx, limitKey, 60*time.Second).Err() } if count == 1 {
_ = l.svcCtx.Redis.Expire(l.ctx, limitKey, 60*time.Second).Err()
}
if count > 120 { // 每分钟最多120条 if count > 120 { // 每分钟最多120条
return nil, errors.Wrapf(xerr.NewErrCode(xerr.TooManyRequests), "too many reports") return nil, errors.Wrapf(xerr.NewErrCode(xerr.TooManyRequests), "too many reports")
} }
@@ -60,18 +64,20 @@ func (l *ReportLogMessageLogic) ReportLogMessage(req *types.ReportLogMessageRequ
} }
var userIdPtr *int64 var userIdPtr *int64
if req.UserId > 0 { userIdPtr = &req.UserId } if req.UserId > 0 {
userIdPtr = &req.UserId
}
row := &logmessage.LogMessage{ row := &logmessage.LogMessage{
Platform: req.Platform, Platform: safeTruncate(req.Platform, 32),
AppVersion: req.AppVersion, AppVersion: safeTruncate(req.AppVersion, 64),
OsName: req.OsName, OsName: safeTruncate(req.OsName, 64),
OsVersion: req.OsVersion, OsVersion: safeTruncate(req.OsVersion, 64),
DeviceId: req.DeviceId, DeviceId: safeTruncate(req.DeviceId, 255),
UserId: userIdPtr, UserId: userIdPtr,
SessionId: req.SessionId, SessionId: safeTruncate(req.SessionId, 255),
Level: req.Level, Level: req.Level,
ErrorCode: req.ErrorCode, ErrorCode: safeTruncate(req.ErrorCode, 128),
Message: safeTruncate(req.Message, 1024*64), Message: safeTruncate(req.Message, 1024*64),
Stack: safeTruncate(req.Stack, 1024*1024), Stack: safeTruncate(req.Stack, 1024*1024),
Context: ctxStr, Context: ctxStr,
@@ -95,14 +101,20 @@ func (l *ReportLogMessageLogic) ReportLogMessage(req *types.ReportLogMessageRequ
} }
func safeTruncate(s string, n int) string { func safeTruncate(s string, n int) string {
if len(s) <= n { return s } if len(s) <= n {
return s
}
return s[:n] return s[:n]
} }
func clientIP(c *gin.Context) string { func clientIP(c *gin.Context) string {
ip := c.ClientIP() ip := c.ClientIP()
if ip != "" { return ip } if ip != "" {
return ip
}
host, _, err := net.SplitHostPort(strings.TrimSpace(c.Request.RemoteAddr)) host, _, err := net.SplitHostPort(strings.TrimSpace(c.Request.RemoteAddr))
if err == nil && host != "" { return host } if err == nil && host != "" {
return host
}
return "" return ""
} }
+48
View File
@@ -0,0 +1,48 @@
package common
import (
"context"
"time"
"github.com/perfect-panel/server/internal/model/log"
usermodel "github.com/perfect-panel/server/internal/model/user"
"github.com/pkg/errors"
"gorm.io/gorm"
"gorm.io/gorm/clause"
)
func WriteCommissionLog(tx *gorm.DB, objectID int64, logType uint16, amount int64, orderNo string) error {
logInfo := log.Commission{
Type: logType,
Amount: amount,
OrderNo: orderNo,
Timestamp: time.Now().UnixMilli(),
}
content, err := logInfo.Marshal()
if err != nil {
return err
}
return tx.Model(log.SystemLog{}).Create(&log.SystemLog{
Type: log.TypeCommission.Uint8(),
Date: time.Now().Format(time.DateOnly),
ObjectID: objectID,
Content: string(content),
CreatedAt: time.Now(),
}).Error
}
func LoadPendingWithdrawalForUpdate(ctx context.Context, tx *gorm.DB, withdrawalID int64) (*usermodel.Withdrawal, error) {
var withdrawal usermodel.Withdrawal
if err := tx.WithContext(ctx).
Clauses(clause.Locking{Strength: "UPDATE"}).
Where("id = ?", withdrawalID).
First(&withdrawal).Error; err != nil {
return nil, err
}
if withdrawal.Status != 0 {
return nil, errors.New("withdrawal status invalid")
}
return &withdrawal, nil
}
+4 -4
View File
@@ -109,14 +109,14 @@ func buildObjectKey(prefix string, userID int64, bizType, fileID, fileName strin
if prefix == "" { if prefix == "" {
prefix = "app-upload" prefix = "app-upload"
} }
return fmt.Sprintf("%s/%s/%d/%04d/%02d/%s_%s", return fmt.Sprintf("%s/%04d/%02d/%02d/%d/%s__%s",
prefix, prefix,
strings.Trim(safeFileNameRegexp.ReplaceAllString(strings.ToLower(strings.TrimSpace(bizType)), "-"), "-"),
userID,
now.Year(), now.Year(),
int(now.Month()), int(now.Month()),
fileID, now.Day(),
userID,
safeFileName(fileName), safeFileName(fileName),
fileID,
) )
} }
@@ -0,0 +1,88 @@
package file
import (
"bytes"
"context"
"io"
"mime/multipart"
"net/http"
"time"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/logger"
"github.com/perfect-panel/server/pkg/xerr"
"github.com/pkg/errors"
)
type FileUploadLogic struct {
logger.Logger
ctx context.Context
svcCtx *svc.ServiceContext
}
// Upload file to RustFS
func NewFileUploadLogic(ctx context.Context, svcCtx *svc.ServiceContext) *FileUploadLogic {
return &FileUploadLogic{
Logger: logger.WithContext(ctx),
ctx: ctx,
svcCtx: svcCtx,
}
}
func (l *FileUploadLogic) FileUpload(req *types.FileUploadRequest, fileHeader *multipart.FileHeader, file multipart.File) (resp *types.FileUploadResponse, err error) {
u, err := currentUserFromContext(l.ctx)
if err != nil {
return nil, err
}
if fileHeader == nil {
return nil, errors.Wrapf(xerr.NewErrCode(xerr.InvalidParams), "file is required")
}
contentType, err := sniffContentType(fileHeader, file)
if err != nil {
return nil, err
}
if err := validateInitRequest(l.svcCtx, req.BizType, fileHeader.Filename, contentType, fileHeader.Size); err != nil {
return nil, err
}
now := time.Now()
fileID := buildFileID(u.Id, req.BizType, fileHeader.Filename)
objectKey := buildObjectKey(l.svcCtx.Config.S3.Prefix, u.Id, req.BizType, fileID, fileHeader.Filename, now)
putResult, err := l.svcCtx.S3Store.PutObject(l.ctx, objectKey, file, fileHeader.Size, contentType)
if err != nil {
l.Errorw("put object failed", logger.Field("error", err.Error()), logger.Field("user_id", u.Id), logger.Field("file_id", fileID))
return nil, err
}
return &types.FileUploadResponse{
FileId: fileID,
FileName: fileHeader.Filename,
ObjectKey: objectKey,
Size: fileHeader.Size,
ContentType: contentType,
Etag: putResult.ETag,
Status: fileUploadCompleteStatus,
}, nil
}
func sniffContentType(fileHeader *multipart.FileHeader, file multipart.File) (string, error) {
if headerType := fileHeader.Header.Get("Content-Type"); headerType != "" {
return headerType, nil
}
if seeker, ok := file.(io.ReadSeeker); ok {
buf := make([]byte, 512)
n, readErr := seeker.Read(buf)
if readErr != nil && readErr != io.EOF {
return "", readErr
}
if _, err := seeker.Seek(0, io.SeekStart); err != nil {
return "", err
}
return http.DetectContentType(bytes.TrimRight(buf[:n], "\x00")), nil
}
return "application/octet-stream", nil
}
@@ -4,6 +4,7 @@ import (
"context" "context"
"time" "time"
logicCommon "github.com/perfect-panel/server/internal/logic/common"
"github.com/perfect-panel/server/internal/model/log" "github.com/perfect-panel/server/internal/model/log"
"github.com/perfect-panel/server/internal/model/user" "github.com/perfect-panel/server/internal/model/user"
"github.com/perfect-panel/server/internal/svc" "github.com/perfect-panel/server/internal/svc"
@@ -12,6 +13,7 @@ import (
"github.com/perfect-panel/server/pkg/logger" "github.com/perfect-panel/server/pkg/logger"
"github.com/perfect-panel/server/pkg/xerr" "github.com/perfect-panel/server/pkg/xerr"
"github.com/pkg/errors" "github.com/pkg/errors"
"gorm.io/gorm"
) )
type CommissionWithdrawLogic struct { type CommissionWithdrawLogic struct {
@@ -42,38 +44,21 @@ func (l *CommissionWithdrawLogic) CommissionWithdraw(req *types.CommissionWithdr
} }
tx := l.svcCtx.DB.WithContext(l.ctx).Begin() tx := l.svcCtx.DB.WithContext(l.ctx).Begin()
now := time.Now()
// update user commission balance // Atomically deduct the requested amount so concurrent commission growth is preserved.
u.Commission -= req.Amount if err = l.svcCtx.DB.WithContext(l.ctx).
if err = l.svcCtx.UserModel.Update(l.ctx, u, tx); err != nil { Model(&user.User{}).
Where("id = ? AND commission >= ?", u.Id, req.Amount).
UpdateColumn("commission", gorm.Expr("commission - ?", req.Amount)).Error; err != nil {
tx.Rollback() tx.Rollback()
l.Errorf("Failed to update user %d commission balance: %v", u.Id, err) l.Errorf("Failed to update user %d commission balance: %v", u.Id, err)
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "Failed to update user %d commission balance: %v", u.Id, err) return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseUpdateError), "Failed to update user %d commission balance: %v", u.Id, err)
} }
_ = l.svcCtx.UserModel.ClearUserCache(l.ctx, u)
// create withdrawal log // create withdrawal log
logInfo := log.Commission{ if err = logicCommon.WriteCommissionLog(tx, u.Id, log.CommissionTypeConvertBalance, req.Amount, ""); err != nil {
Type: log.CommissionTypeConvertBalance,
Amount: req.Amount,
Timestamp: time.Now().UnixMilli(),
}
b, err := logInfo.Marshal()
if err != nil {
tx.Rollback()
l.Errorf("Failed to marshal commission log for user %d: %v", u.Id, err)
return nil, errors.Wrapf(xerr.NewErrCode(xerr.ERROR), "Failed to marshal commission log for user %d: %v", u.Id, err)
}
err = tx.Model(log.SystemLog{}).Create(&log.SystemLog{
Type: log.TypeCommission.Uint8(),
Date: time.Now().Format("2006-01-02"),
ObjectID: u.Id,
Content: string(b),
CreatedAt: time.Now(),
}).Error
if err != nil {
tx.Rollback() tx.Rollback()
l.Errorf("Failed to create commission log for user %d: %v", u.Id, err) l.Errorf("Failed to create commission log for user %d: %v", u.Id, err)
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "Failed to create commission log for user %d: %v", u.Id, err) return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseInsertError), "Failed to create commission log for user %d: %v", u.Id, err)
@@ -103,6 +88,7 @@ func (l *CommissionWithdrawLogic) CommissionWithdraw(req *types.CommissionWithdr
Content: req.Content, Content: req.Content,
Status: 0, Status: 0,
Reason: "", Reason: "",
CreatedAt: time.Now().UnixMilli(), CreatedAt: now.UnixMilli(),
UpdatedAt: now.UnixMilli(),
}, nil }, nil
} }
@@ -6,7 +6,9 @@ import (
"sort" "sort"
"strings" "strings"
authlogic "github.com/perfect-panel/server/internal/logic/auth"
logicCommon "github.com/perfect-panel/server/internal/logic/common" logicCommon "github.com/perfect-panel/server/internal/logic/common"
modelOrder "github.com/perfect-panel/server/internal/model/order"
"github.com/perfect-panel/server/pkg/constant" "github.com/perfect-panel/server/pkg/constant"
"github.com/perfect-panel/server/pkg/uuidx" "github.com/perfect-panel/server/pkg/uuidx"
"github.com/perfect-panel/server/pkg/xerr" "github.com/perfect-panel/server/pkg/xerr"
@@ -44,6 +46,7 @@ func (l *QueryUserInfoLogic) QueryUserInfo() (resp *types.User, err error) {
return nil, errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "Invalid Access") return nil, errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "Invalid Access")
} }
tool.DeepCopy(resp, u) tool.DeepCopy(resp, u)
resp.UseStatus = true
// 用家庭范围查设备,而不是只看当前用户自己的 UserDevices // 用家庭范围查设备,而不是只看当前用户自己的 UserDevices
scopeHelper := newFamilyScopeHelper(l.ctx, l.svcCtx) scopeHelper := newFamilyScopeHelper(l.ctx, l.svcCtx)
@@ -65,6 +68,14 @@ func (l *QueryUserInfoLogic) QueryUserInfo() (resp *types.User, err error) {
} }
resp.UserDevices = userDevices resp.UserDevices = userDevices
} }
useStatus, useStatusErr := l.resolveBindEmailTrialUseStatus(u.Id, scopeUserIds)
if useStatusErr != nil {
l.Errorw("resolve bind email trial use status failed", logger.Field("user_id", u.Id), logger.Field("error", useStatusErr.Error()))
} else {
resp.UseStatus = useStatus
}
// refer_code 为空时自动生成 // refer_code 为空时自动生成
if resp.ReferCode == "" { if resp.ReferCode == "" {
resp.ReferCode = uuidx.UserInviteCode(u.Id) resp.ReferCode = uuidx.UserInviteCode(u.Id)
@@ -108,6 +119,49 @@ func (l *QueryUserInfoLogic) QueryUserInfo() (resp *types.User, err error) {
return resp, nil return resp, nil
} }
// resolveBindEmailTrialUseStatus determines whether userinfo should show the
// "bind email to get free trial" prompt. `true` means show the prompt.
func (l *QueryUserInfoLogic) resolveBindEmailTrialUseStatus(currentUserId int64, scopeUserIds []int64) (bool, error) {
if len(scopeUserIds) == 0 {
scopeUserIds = []int64{currentUserId}
}
var hasBoundEmailCount int64
if err := l.svcCtx.DB.WithContext(l.ctx).
Model(&user.AuthMethods{}).
Where("user_id IN ? AND auth_type = ? AND auth_identifier != ''", scopeUserIds, "email").
Count(&hasBoundEmailCount).Error; err != nil {
return false, err
}
var hasPurchaseCount int64
if err := l.svcCtx.DB.WithContext(l.ctx).
Model(&modelOrder.Order{}).
Where("user_id IN ? AND type IN ? AND status IN ?", scopeUserIds, []int64{1, 2}, []int64{2, 5}).
Count(&hasPurchaseCount).Error; err != nil {
return false, err
}
hasTrial := false
registerCfg := l.svcCtx.Config.Register
if authlogic.IsTrialConfigReady(registerCfg) && registerCfg.TrialSubscribe > 0 {
var hasTrialCount int64
if err := l.svcCtx.DB.WithContext(l.ctx).
Model(&user.Subscribe{}).
Where("user_id IN ? AND subscribe_id = ?", scopeUserIds, registerCfg.TrialSubscribe).
Count(&hasTrialCount).Error; err != nil {
return false, err
}
hasTrial = hasTrialCount > 0
}
return shouldShowBindEmailTrialPrompt(hasBoundEmailCount > 0, hasPurchaseCount > 0, hasTrial), nil
}
func shouldShowBindEmailTrialPrompt(hasBoundEmail, hasPurchased, hasTrial bool) bool {
return !hasBoundEmail && !hasPurchased && !hasTrial
}
func (l *QueryUserInfoLogic) fillFamilyContext(resp *types.User, userId int64) *user.AuthMethods { func (l *QueryUserInfoLogic) fillFamilyContext(resp *types.User, userId int64) *user.AuthMethods {
type familyRelation struct { type familyRelation struct {
FamilyId int64 FamilyId int64
@@ -0,0 +1,67 @@
package user
import "testing"
func TestShouldShowBindEmailTrialPrompt(t *testing.T) {
tests := []struct {
name string
hasBoundEmail bool
hasPurchased bool
hasTrial bool
want bool
}{
{
name: "new user should see prompt",
want: true,
},
{
name: "bound email should hide prompt",
hasBoundEmail: true,
want: false,
},
{
name: "paid purchase should hide prompt",
hasPurchased: true,
want: false,
},
{
name: "trial claimed should hide prompt",
hasTrial: true,
want: false,
},
{
name: "bound email and purchase should hide prompt",
hasBoundEmail: true,
hasPurchased: true,
want: false,
},
{
name: "bound email and trial should hide prompt",
hasBoundEmail: true,
hasTrial: true,
want: false,
},
{
name: "purchase and trial should hide prompt",
hasPurchased: true,
hasTrial: true,
want: false,
},
{
name: "all blockers should hide prompt",
hasBoundEmail: true,
hasPurchased: true,
hasTrial: true,
want: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got := shouldShowBindEmailTrialPrompt(tt.hasBoundEmail, tt.hasPurchased, tt.hasTrial)
if got != tt.want {
t.Fatalf("shouldShowBindEmailTrialPrompt(%v, %v, %v) = %v, want %v", tt.hasBoundEmail, tt.hasPurchased, tt.hasTrial, got, tt.want)
}
})
}
}
@@ -3,9 +3,13 @@ package user
import ( import (
"context" "context"
"github.com/perfect-panel/server/internal/model/user"
"github.com/perfect-panel/server/internal/svc" "github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/internal/types" "github.com/perfect-panel/server/internal/types"
"github.com/perfect-panel/server/pkg/constant"
"github.com/perfect-panel/server/pkg/logger" "github.com/perfect-panel/server/pkg/logger"
"github.com/perfect-panel/server/pkg/xerr"
"github.com/pkg/errors"
) )
type QueryWithdrawalLogLogic struct { type QueryWithdrawalLogLogic struct {
@@ -24,7 +28,49 @@ func NewQueryWithdrawalLogLogic(ctx context.Context, svcCtx *svc.ServiceContext)
} }
func (l *QueryWithdrawalLogLogic) QueryWithdrawalLog(req *types.QueryWithdrawalLogListRequest) (resp *types.QueryWithdrawalLogListResponse, err error) { func (l *QueryWithdrawalLogLogic) QueryWithdrawalLog(req *types.QueryWithdrawalLogListRequest) (resp *types.QueryWithdrawalLogListResponse, err error) {
// todo: add your logic here and delete this line u, ok := l.ctx.Value(constant.CtxKeyUser).(*user.User)
if !ok {
return l.Error("current user is not found in context")
return nil, errors.Wrapf(xerr.NewErrCode(xerr.InvalidAccess), "Invalid Access")
}
page := req.Page
size := req.Size
if page <= 0 {
page = 1
}
if size <= 0 {
size = 10
}
query := l.svcCtx.DB.WithContext(l.ctx).Model(&user.Withdrawal{}).Where("user_id = ?", u.Id)
var total int64
if err = query.Count(&total).Error; err != nil {
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "count withdrawal logs failed: %v", err)
}
var rows []user.Withdrawal
if err = query.Order("id DESC").Limit(size).Offset((page - 1) * size).Find(&rows).Error; err != nil {
return nil, errors.Wrapf(xerr.NewErrCode(xerr.DatabaseQueryError), "query withdrawal logs failed: %v", err)
}
list := make([]types.WithdrawalLog, 0, len(rows))
for _, row := range rows {
list = append(list, types.WithdrawalLog{
Id: row.Id,
UserId: row.UserId,
Amount: row.Amount,
Content: row.Content,
Status: row.Status,
Reason: row.Reason,
CreatedAt: row.CreatedAt.UnixMilli(),
UpdatedAt: row.UpdatedAt.UnixMilli(),
})
}
return &types.QueryWithdrawalLogListResponse{
List: list,
Total: total,
}, nil
} }
+17
View File
@@ -31,6 +31,8 @@ const (
ctxDecryptedQueryKey = "decrypted_query" ctxDecryptedQueryKey = "decrypted_query"
ctxEncryptedBodyKey = "encrypted_request_body" ctxEncryptedBodyKey = "encrypted_request_body"
ctxDecryptedBodyKey = "decrypted_request_body" ctxDecryptedBodyKey = "decrypted_request_body"
deviceDecryptSkipPathPublicFileUpload = "/v1/public/file/upload"
) )
func DeviceMiddleware(srvCtx *svc.ServiceContext) func(c *gin.Context) { func DeviceMiddleware(srvCtx *svc.ServiceContext) func(c *gin.Context) {
@@ -70,6 +72,14 @@ func DeviceMiddleware(srvCtx *svc.ServiceContext) func(c *gin.Context) {
} }
rw := NewResponseWriter(c, srvCtx) rw := NewResponseWriter(c, srvCtx)
if shouldSkipDeviceRequestDecrypt(c) {
c.Set(ctxDeviceDecryptStatusKey, "skipped")
c.Set(ctxDeviceDecryptReasonKey, "multipart_upload_passthrough")
c.Writer = rw
c.Next()
rw.FlushAbort()
return
}
if !rw.Decrypt() { if !rw.Decrypt() {
c.Set(ctxDeviceDecryptStatusKey, "failed") c.Set(ctxDeviceDecryptStatusKey, "failed")
if _, exists := c.Get(ctxDeviceDecryptReasonKey); !exists { if _, exists := c.Get(ctxDeviceDecryptReasonKey); !exists {
@@ -85,6 +95,13 @@ func DeviceMiddleware(srvCtx *svc.ServiceContext) func(c *gin.Context) {
} }
} }
func shouldSkipDeviceRequestDecrypt(c *gin.Context) bool {
if c.Request.URL.Path != deviceDecryptSkipPathPublicFileUpload {
return false
}
return strings.HasPrefix(strings.ToLower(strings.TrimSpace(c.GetHeader("Content-Type"))), "multipart/form-data")
}
func NewResponseWriter(c *gin.Context, srvCtx *svc.ServiceContext) (rw *ResponseWriter) { func NewResponseWriter(c *gin.Context, srvCtx *svc.ServiceContext) (rw *ResponseWriter) {
rw = &ResponseWriter{ rw = &ResponseWriter{
c: c, c: c,
+1
View File
@@ -49,6 +49,7 @@ const (
CommissionTypeWithdraw uint16 = 334 // withdraw CommissionTypeWithdraw uint16 = 334 // withdraw
CommissionTypeAdjust uint16 = 335 // Admin Adjust CommissionTypeAdjust uint16 = 335 // Admin Adjust
CommissionTypeConvertBalance uint16 = 336 // Convert to Balance CommissionTypeConvertBalance uint16 = 336 // Convert to Balance
CommissionTypeWithdrawReject uint16 = 337 // Withdraw rejected refund
GiftTypeIncrease uint16 = 341 // Increase GiftTypeIncrease uint16 = 341 // Increase
GiftTypeReduce uint16 = 342 // Reduce GiftTypeReduce uint16 = 342 // Reduce
) )
+6 -6
View File
@@ -5,14 +5,14 @@ import "time"
type LogMessage struct { type LogMessage struct {
Id int64 `gorm:"primaryKey;AUTO_INCREMENT"` Id int64 `gorm:"primaryKey;AUTO_INCREMENT"`
Platform string `gorm:"type:varchar(32);not null"` Platform string `gorm:"type:varchar(32);not null"`
AppVersion string `gorm:"type:varchar(32);default:null"` AppVersion string `gorm:"type:varchar(64);default:null"`
OsName string `gorm:"type:varchar(32);default:null"` OsName string `gorm:"type:varchar(64);default:null"`
OsVersion string `gorm:"type:varchar(32);default:null"` OsVersion string `gorm:"type:varchar(64);default:null"`
DeviceId string `gorm:"type:varchar(64);default:null"` DeviceId string `gorm:"type:varchar(255);default:null"`
UserId *int64 `gorm:"type:bigint;default:null"` UserId *int64 `gorm:"type:bigint;default:null"`
SessionId string `gorm:"type:varchar(64);default:null"` SessionId string `gorm:"type:varchar(255);default:null"`
Level uint8 `gorm:"type:tinyint(1);not null;default:3"` Level uint8 `gorm:"type:tinyint(1);not null;default:3"`
ErrorCode string `gorm:"type:varchar(64);default:null"` ErrorCode string `gorm:"type:varchar(128);default:null"`
Message string `gorm:"type:text;not null"` Message string `gorm:"type:text;not null"`
Stack string `gorm:"type:mediumtext;default:null"` Stack string `gorm:"type:mediumtext;default:null"`
Context string `gorm:"type:json;default:null"` Context string `gorm:"type:json;default:null"`
+17
View File
@@ -26,6 +26,7 @@ type (
Insert(ctx context.Context, data *User, tx ...*gorm.DB) error Insert(ctx context.Context, data *User, tx ...*gorm.DB) error
FindOne(ctx context.Context, id int64) (*User, error) FindOne(ctx context.Context, id int64) (*User, error)
Update(ctx context.Context, data *User, tx ...*gorm.DB) error Update(ctx context.Context, data *User, tx ...*gorm.DB) error
UpdateCommission(ctx context.Context, userId int64, delta int64, tx ...*gorm.DB) error
Delete(ctx context.Context, id int64, tx ...*gorm.DB) error Delete(ctx context.Context, id int64, tx ...*gorm.DB) error
Transaction(ctx context.Context, fn func(db *gorm.DB) error) error Transaction(ctx context.Context, fn func(db *gorm.DB) error) error
} }
@@ -111,6 +112,22 @@ func (m *defaultUserModel) Update(ctx context.Context, data *User, tx ...*gorm.D
return err return err
} }
func (m *defaultUserModel) UpdateCommission(ctx context.Context, userId int64, delta int64, tx ...*gorm.DB) error {
old, err := m.FindOne(ctx, userId)
if err != nil {
return err
}
return m.ExecCtx(ctx, func(conn *gorm.DB) error {
if len(tx) > 0 {
conn = tx[0]
}
return conn.Model(&User{}).
Where("id = ?", userId).
UpdateColumn("commission", gorm.Expr("commission + ?", delta)).Error
}, m.getCacheKeys(old)...)
}
func (m *defaultUserModel) Delete(ctx context.Context, id int64, tx ...*gorm.DB) error { func (m *defaultUserModel) Delete(ctx context.Context, id int64, tx ...*gorm.DB) error {
data, err := m.FindOne(ctx, id) data, err := m.FindOne(ctx, id)
if err != nil { if err != nil {
+18
View File
@@ -8,6 +8,22 @@ import (
"gorm.io/gorm" "gorm.io/gorm"
) )
const userDeviceUserAgentMaxLength = 255
func normalizeDeviceForStorage(data *Device) {
if data == nil {
return
}
data.UserAgent = truncateForColumn(data.UserAgent, userDeviceUserAgentMaxLength)
}
func truncateForColumn(s string, max int) string {
if len(s) <= max {
return s
}
return s[:max]
}
func (m *customUserModel) FindOneDevice(ctx context.Context, id int64) (*Device, error) { func (m *customUserModel) FindOneDevice(ctx context.Context, id int64) (*Device, error) {
deviceIdKey := fmt.Sprintf("%s%v", cacheUserDeviceIdPrefix, id) deviceIdKey := fmt.Sprintf("%s%v", cacheUserDeviceIdPrefix, id)
var resp Device var resp Device
@@ -69,6 +85,7 @@ func (m *customUserModel) QueryDeviceListByUserIds(ctx context.Context, userIds
} }
func (m *customUserModel) UpdateDevice(ctx context.Context, data *Device, tx ...*gorm.DB) error { func (m *customUserModel) UpdateDevice(ctx context.Context, data *Device, tx ...*gorm.DB) error {
normalizeDeviceForStorage(data)
old, err := m.FindOneDevice(ctx, data.Id) old, err := m.FindOneDevice(ctx, data.Id)
if err != nil { if err != nil {
return err return err
@@ -100,6 +117,7 @@ func (m *customUserModel) DeleteDevice(ctx context.Context, id int64, tx ...*gor
} }
func (m *customUserModel) InsertDevice(ctx context.Context, data *Device, tx ...*gorm.DB) error { func (m *customUserModel) InsertDevice(ctx context.Context, data *Device, tx ...*gorm.DB) error {
normalizeDeviceForStorage(data)
defer func() { defer func() {
if clearErr := m.ClearDeviceCache(ctx, data); clearErr != nil { if clearErr := m.ClearDeviceCache(ctx, data); clearErr != nil {
// log cache clear error // log cache clear error
+36
View File
@@ -762,6 +762,20 @@ type FamilySummary struct {
UpdatedAt int64 `json:"updated_at"` UpdatedAt int64 `json:"updated_at"`
} }
type FileUploadRequest struct {
BizType string `form:"biz_type" validate:"required"`
}
type FileUploadResponse struct {
FileId string `json:"file_id"`
FileName string `json:"file_name"`
ObjectKey string `json:"object_key"`
Size int64 `json:"size"`
ContentType string `json:"content_type"`
Etag string `json:"etag"`
Status string `json:"status"`
}
type FileUploadCompleteRequest struct { type FileUploadCompleteRequest struct {
FileId string `json:"file_id" validate:"required"` FileId string `json:"file_id" validate:"required"`
} }
@@ -2270,6 +2284,27 @@ type QueryWithdrawalLogListResponse struct {
Total int64 `json:"total"` Total int64 `json:"total"`
} }
type GetWithdrawalListRequest struct {
Page int `form:"page"`
Size int `form:"size"`
UserId *int64 `form:"user_id,omitempty"`
Status *uint8 `form:"status,omitempty"`
}
type GetWithdrawalListResponse struct {
List []WithdrawalLog `json:"list"`
Total int64 `json:"total"`
}
type ApproveWithdrawalRequest struct {
WithdrawalId int64 `json:"withdrawal_id" validate:"required,gt=0"`
}
type RejectWithdrawalRequest struct {
WithdrawalId int64 `json:"withdrawal_id" validate:"required,gt=0"`
Reason string `json:"reason" validate:"required,max=500"`
}
type QuotaTask struct { type QuotaTask struct {
Id int64 `json:"id"` Id int64 `json:"id"`
Subscribers []int64 `json:"subscribers"` Subscribers []int64 `json:"subscribers"`
@@ -3264,6 +3299,7 @@ type User struct {
EnableLoginNotify bool `json:"enable_login_notify"` EnableLoginNotify bool `json:"enable_login_notify"`
EnableSubscribeNotify bool `json:"enable_subscribe_notify"` EnableSubscribeNotify bool `json:"enable_subscribe_notify"`
EnableTradeNotify bool `json:"enable_trade_notify"` EnableTradeNotify bool `json:"enable_trade_notify"`
UseStatus bool `json:"use_status"`
AuthMethods []UserAuthMethod `json:"auth_methods"` AuthMethods []UserAuthMethod `json:"auth_methods"`
UserDevices []UserDevice `json:"user_devices"` UserDevices []UserDevice `json:"user_devices"`
Rules []string `json:"rules"` Rules []string `json:"rules"`
+4 -4
View File
@@ -22,7 +22,7 @@
- Region: `ap-east-1` - Region: `ap-east-1`
- AWS app EC2: - AWS app EC2:
- Name: `hifast-hk-app-01` - Name: `hifast-hk-app-01`
- Public IP: `43.198.248.161` - Public IP: `18.163.33.75`
- Private IP: `10.0.1.201` - Private IP: `10.0.1.201`
- AWS MySQL: - AWS MySQL:
- Type: `RDS MySQL` - Type: `RDS MySQL`
@@ -183,7 +183,7 @@ redis-cli INFO replication
重点看: 重点看:
- `role:slave` - `role:slave`
- `master_host:43.198.248.161` - `master_host:18.163.33.75`
- `master_port:6379` - `master_port:6379`
- `master_link_status:up` - `master_link_status:up`
@@ -416,7 +416,7 @@ SHOW REPLICA STATUS\G
```bash ```bash
redis-cli CONFIG SET masterauth '0BVz9XOHf7KUfEuoFJRK-dURdKUGFiZ8QeaHpysHnKeKhLskZb55HPK121lFsKtr' redis-cli CONFIG SET masterauth '0BVz9XOHf7KUfEuoFJRK-dURdKUGFiZ8QeaHpysHnKeKhLskZb55HPK121lFsKtr'
redis-cli REPLICAOF 43.198.248.161 6379 redis-cli REPLICAOF 18.163.33.75 6379
redis-cli CONFIG SET replica-read-only yes redis-cli CONFIG SET replica-read-only yes
redis-cli INFO replication redis-cli INFO replication
``` ```
@@ -426,7 +426,7 @@ redis-cli INFO replication
检查 `/etc/redis/redis.conf` 至少包含: 检查 `/etc/redis/redis.conf` 至少包含:
```conf ```conf
replicaof 43.198.248.161 6379 replicaof 18.163.33.75 6379
masterauth 0BVz9XOHf7KUfEuoFJRK-dURdKUGFiZ8QeaHpysHnKeKhLskZb55HPK121lFsKtr masterauth 0BVz9XOHf7KUfEuoFJRK-dURdKUGFiZ8QeaHpysHnKeKhLskZb55HPK121lFsKtr
replica-read-only yes replica-read-only yes
``` ```
+9 -9
View File
@@ -66,7 +66,7 @@
- Instance ID: `i-079cd9d3ef3748714` - Instance ID: `i-079cd9d3ef3748714`
- 角色:当前实际生产入口 / Nginx / 业务服务 / AWS 侧 Redis 主库宿主机 - 角色:当前实际生产入口 / Nginx / 业务服务 / AWS 侧 Redis 主库宿主机
- 私网 IP: `10.0.1.201` - 私网 IP: `10.0.1.201`
- 公网 IP: `43.198.248.161` - 公网 IP: `18.163.33.75`
- 业务服务运行方式:`Docker Compose` - 业务服务运行方式:`Docker Compose`
- 业务容器:`ppanel-server` - 业务容器:`ppanel-server`
- 部署目录:`/opt/ppanel` - 部署目录:`/opt/ppanel`
@@ -103,7 +103,7 @@
- 容器名:`hifast-redis` - 容器名:`hifast-redis`
- 版本:`redis:8.2.1` - 版本:`redis:8.2.1`
- 访问端口:`6379` - 访问端口:`6379`
- 主库出口地址:`43.198.248.161:6379` - 主库出口地址:`18.163.33.75:6379`
- 应用当前实际连接:`127.0.0.1:6379` - 应用当前实际连接:`127.0.0.1:6379`
- 认证方式:已启用密码认证 - 认证方式:已启用密码认证
@@ -130,10 +130,10 @@
```mermaid ```mermaid
flowchart TB flowchart TB
USER["用户 / 客户端"] --> DNS["域名 / DNS / 入口层"] USER["用户 / 客户端"] --> DNS["域名 / DNS / 入口层"]
DNS --> APP["AWS EC2\nhifast-hk-app-01\n43.198.248.161\n10.0.1.201"] DNS --> APP["AWS EC2\nhifast-hk-app-01\n18.163.33.75\n10.0.1.201"]
APP --> RDS["AWS RDS MySQL\nhifast-mysql-prod-v2\n主库"] APP --> RDS["AWS RDS MySQL\nhifast-mysql-prod-v2\n主库"]
APP --> REDISM["AWS Redis 主库\nDocker redis:8.2.1\n43.198.248.161:6379"] APP --> REDISM["AWS Redis 主库\nDocker redis:8.2.1\n18.163.33.75:6379"]
RDS -. MySQL 备用 / 同步 .-> MYSQLS["104.238.220.230\nMySQL 备用库"] RDS -. MySQL 备用 / 同步 .-> MYSQLS["104.238.220.230\nMySQL 备用库"]
REDISM -. Redis 主从复制 .-> REDISS["104.238.220.230\n原生 Redis 8.6.3\n从库"] REDISM -. Redis 主从复制 .-> REDISS["104.238.220.230\n原生 Redis 8.6.3\n从库"]
@@ -238,7 +238,7 @@ App / Nginx
```mermaid ```mermaid
flowchart LR flowchart LR
SG["hifast-hk-app-core-sg"] --> REDIS["AWS Redis 主库\n43.198.248.161:6379"] SG["hifast-hk-app-core-sg"] --> REDIS["AWS Redis 主库\n18.163.33.75:6379"]
STANDBY["104.238.220.230/32"] --> SG STANDBY["104.238.220.230/32"] --> SG
``` ```
@@ -311,7 +311,7 @@ RDS 当前状态已经比之前干净很多:
`104.238.220.230` 连接 AWS Redis,不是通过 PEM 证书,也不是通过 SSH 登录 AWS 机器,而是直接作为 Redis 从库去访问 AWS Redis 主库: `104.238.220.230` 连接 AWS Redis,不是通过 PEM 证书,也不是通过 SSH 登录 AWS 机器,而是直接作为 Redis 从库去访问 AWS Redis 主库:
- 目标地址:`43.198.248.161:6379` - 目标地址:`18.163.33.75:6379`
- 连接方式:`TCP` - 连接方式:`TCP`
- 认证方式:`Redis 密码` - 认证方式:`Redis 密码`
- 网络前提:AWS EC2 安全组已放行 `104.238.220.230/32 -> 6379` - 网络前提:AWS EC2 安全组已放行 `104.238.220.230/32 -> 6379`
@@ -325,7 +325,7 @@ RDS 当前状态已经比之前干净很多:
示意命令: 示意命令:
```bash ```bash
redis-cli -h 43.198.248.161 -p 6379 -a '<REDIS_PASSWORD>' redis-cli -h 18.163.33.75 -p 6379 -a '<REDIS_PASSWORD>'
``` ```
### 4.6.2 104 连接 AWS MySQL RDS 的方式 ### 4.6.2 104 连接 AWS MySQL RDS 的方式
@@ -380,7 +380,7 @@ mysql -h hifast-mysql-prod-v2.cd6aey40m6ag.ap-east-1.rds.amazonaws.com -u admin
- 部署方式:Docker - 部署方式:Docker
- 版本:`8.2.1` - 版本:`8.2.1`
- 主库地址:`43.198.248.161:6379` - 主库地址:`18.163.33.75:6379`
- 运行容器:`hifast-redis` - 运行容器:`hifast-redis`
### 5.2 104 Redis 从库 ### 5.2 104 Redis 从库
@@ -417,7 +417,7 @@ mysql -h hifast-mysql-prod-v2.cd6aey40m6ag.ap-east-1.rds.amazonaws.com -u admin
```mermaid ```mermaid
flowchart LR flowchart LR
REDISMASTER["AWS Redis 主库\n43.198.248.161:6379\nDocker redis:8.2.1"] REDISMASTER["AWS Redis 主库\n18.163.33.75:6379\nDocker redis:8.2.1"]
REDISSLAVE["104.238.220.230\n原生 Redis 8.6.3\nrole: slave"] REDISSLAVE["104.238.220.230\n原生 Redis 8.6.3\nrole: slave"]
REDISMASTER --> REDISSLAVE REDISMASTER --> REDISSLAVE
``` ```
+19
View File
@@ -3,6 +3,7 @@ package storage
import ( import (
"context" "context"
"fmt" "fmt"
"io"
"net/url" "net/url"
"strings" "strings"
"time" "time"
@@ -27,6 +28,10 @@ type ObjectMeta struct {
ContentLength int64 ContentLength int64
} }
type PutObjectResult struct {
ETag string
}
type S3Store struct { type S3Store struct {
cfg config.S3Config cfg config.S3Config
client *s3.Client client *s3.Client
@@ -96,6 +101,20 @@ func (s *S3Store) HeadObject(ctx context.Context, objectKey string) (*ObjectMeta
}, nil }, nil
} }
func (s *S3Store) PutObject(ctx context.Context, objectKey string, body io.Reader, size int64, contentType string) (*PutObjectResult, error) {
resp, err := s.client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(s.cfg.Bucket),
Key: aws.String(objectKey),
Body: body,
ContentLength: aws.Int64(size),
ContentType: aws.String(contentType),
})
if err != nil {
return nil, err
}
return &PutObjectResult{ETag: strings.Trim(aws.ToString(resp.ETag), "\"")}, nil
}
func (s *S3Store) BuildObjectURL(objectKey string) string { func (s *S3Store) BuildObjectURL(objectKey string) string {
if s.cfg.PublicBaseURL != "" { if s.cfg.PublicBaseURL != "" {
return strings.TrimRight(s.cfg.PublicBaseURL, "/") + "/" + strings.TrimLeft(objectKey, "/") return strings.TrimRight(s.cfg.PublicBaseURL, "/") + "/" + strings.TrimLeft(objectKey, "/")
+1
View File
@@ -37,6 +37,7 @@ const (
FamilyNotExist uint32 = 20017 FamilyNotExist uint32 = 20017
FamilyStatusInvalid uint32 = 20018 FamilyStatusInvalid uint32 = 20018
FamilyOwnerOperationForbidden uint32 = 20019 FamilyOwnerOperationForbidden uint32 = 20019
WithdrawalStatusInvalid uint32 = 20020
) )
// Node error // Node error
+1
View File
@@ -46,6 +46,7 @@ func init() {
FamilyNotExist: "家庭组不存在", FamilyNotExist: "家庭组不存在",
FamilyStatusInvalid: "家庭组状态无效", FamilyStatusInvalid: "家庭组状态无效",
FamilyOwnerOperationForbidden: "家庭组所有者不允许此操作", FamilyOwnerOperationForbidden: "家庭组所有者不允许此操作",
WithdrawalStatusInvalid: "提现状态无效",
// Node error // Node error
NodeExist: "Node already exists", NodeExist: "Node already exists",
+3
View File
@@ -48,4 +48,7 @@ func RegisterHandlers(mux *asynq.ServeMux, serverCtx *svc.ServiceContext) {
// Apple IAP 对账(第二层:5min 扫描 + 第三层:日终全量) // Apple IAP 对账(第二层:5min 扫描 + 第三层:日终全量)
mux.Handle(types.SchedulerIAPReconcile, iapLogic.NewReconcileLogic(serverCtx)) mux.Handle(types.SchedulerIAPReconcile, iapLogic.NewReconcileLogic(serverCtx))
mux.Handle(types.SchedulerIAPDailyReconcile, iapLogic.NewDailyReconcileLogic(serverCtx)) mux.Handle(types.SchedulerIAPDailyReconcile, iapLogic.NewDailyReconcileLogic(serverCtx))
// Stuck order recovery: reset claimed orders that timed out back to paid
mux.Handle(types.SchedulerStuckOrderRecovery, orderLogic.NewStuckOrderRecoveryLogic(serverCtx))
} }
+50 -8
View File
@@ -7,6 +7,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt" "fmt"
"strings"
"time" "time"
"github.com/perfect-panel/server/internal/logic/admin/group" "github.com/perfect-panel/server/internal/logic/admin/group"
@@ -46,7 +47,7 @@ const (
OrderStatusPaid = 2 // Order paid and ready for processing OrderStatusPaid = 2 // Order paid and ready for processing
OrderStatusClose = 3 // Order closed/cancelled OrderStatusClose = 3 // Order closed/cancelled
OrderStatusFailed = 4 // Order processing failed OrderStatusFailed = 4 // Order processing failed
OrderStatusClaimed = 4 // Internal transient claim while a worker processes the order OrderStatusClaimed = 6 // Internal transient claim while a worker processes the order
OrderStatusFinished = 5 // Order successfully completed OrderStatusFinished = 5 // Order successfully completed
) )
@@ -94,6 +95,11 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
if err != nil { if err != nil {
// 如果订单不存在或状态不对,不重试 // 如果订单不存在或状态不对,不重试
if errors.Is(err, ErrInvalidOrderStatus) { if errors.Is(err, ErrInvalidOrderStatus) {
if strings.Contains(err.Error(), "stuck in claimed") {
logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单卡在 claimed,将重试",
logger.Field("order_no", payload.OrderNo))
return err
}
logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过", logger.WithContext(ctx).Info("[ActivateOrderLogic] 订单状态不是已支付,跳过",
logger.Field("order_no", payload.OrderNo)) logger.Field("order_no", payload.OrderNo))
return nil return nil
@@ -116,7 +122,12 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
) )
if err = l.processOrderByType(ctx, orderInfo, payload.IAPExpireAt); err != nil { if err = l.processOrderByType(ctx, orderInfo, payload.IAPExpireAt); err != nil {
l.releaseClaim(ctx, orderInfo.OrderNo) if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil {
logger.WithContext(ctx).Error("[ActivateOrderLogic] releaseClaim also failed, stuck recovery will handle",
logger.Field("order_no", orderInfo.OrderNo),
logger.Field("release_error", releaseErr.Error()),
)
}
logger.WithContext(ctx).Error("[ActivateOrderLogic] 处理订单失败,将重试", logger.WithContext(ctx).Error("[ActivateOrderLogic] 处理订单失败,将重试",
logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_no", orderInfo.OrderNo),
logger.Field("order_type", orderInfo.Type), logger.Field("order_type", orderInfo.Type),
@@ -125,7 +136,12 @@ func (l *ActivateOrderLogic) ProcessTask(ctx context.Context, task *asynq.Task)
} }
if err = l.reconcilePostOrderSubscriptions(ctx, orderInfo); err != nil { if err = l.reconcilePostOrderSubscriptions(ctx, orderInfo); err != nil {
l.releaseClaim(ctx, orderInfo.OrderNo) if releaseErr := l.releaseClaim(ctx, orderInfo.OrderNo); releaseErr != nil {
logger.WithContext(ctx).Error("[ActivateOrderLogic] releaseClaim also failed, stuck recovery will handle",
logger.Field("order_no", orderInfo.OrderNo),
logger.Field("release_error", releaseErr.Error()),
)
}
logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单订阅兜底合并失败,将重试", logger.WithContext(ctx).Error("[ActivateOrderLogic] 订单订阅兜底合并失败,将重试",
logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_no", orderInfo.OrderNo),
logger.Field("order_type", orderInfo.Type), logger.Field("order_type", orderInfo.Type),
@@ -176,6 +192,14 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
return nil, nil return nil, nil
} }
if orderInfo.Status == OrderStatusClaimed {
logger.WithContext(ctx).Error("Order stuck in claimed status",
logger.Field("order_no", orderInfo.OrderNo),
logger.Field("status", orderInfo.Status),
)
return nil, fmt.Errorf("order %s stuck in claimed status: %w", orderNo, ErrInvalidOrderStatus)
}
if orderInfo.Status != OrderStatusPaid { if orderInfo.Status != OrderStatusPaid {
logger.WithContext(ctx).Error("Order status error", logger.WithContext(ctx).Error("Order status error",
logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_no", orderInfo.OrderNo),
@@ -201,7 +225,7 @@ func (l *ActivateOrderLogic) claimAndGetOrder(ctx context.Context, orderNo strin
return &orderInfo, nil return &orderInfo, nil
} }
func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) { func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) error {
if err := l.svc.DB.WithContext(ctx). if err := l.svc.DB.WithContext(ctx).
Model(&order.Order{}). Model(&order.Order{}).
Where("order_no = ? AND status = ?", orderNo, OrderStatusClaimed). Where("order_no = ? AND status = ?", orderNo, OrderStatusClaimed).
@@ -210,7 +234,9 @@ func (l *ActivateOrderLogic) releaseClaim(ctx context.Context, orderNo string) {
logger.Field("error", err.Error()), logger.Field("error", err.Error()),
logger.Field("order_no", orderNo), logger.Field("order_no", orderNo),
) )
return fmt.Errorf("release claim failed for order %s: %w", orderNo, err)
} }
return nil
} }
// processOrderByType routes order processing based on the order type // processOrderByType routes order processing based on the order type
@@ -563,10 +589,26 @@ func (l *ActivateOrderLogic) finalizeCouponAndOrder(ctx context.Context, orderIn
} }
} }
// Update order status using state-guarded UpdateOrderStatus to prevent double finalization // Direct update from claimed(6) → finished(5), bypassing the model's status<target guard.
if err := l.svc.OrderModel.UpdateOrderStatus(ctx, orderInfo.OrderNo, OrderStatusFinished); err != nil { // UpdateOrderStatus uses WHERE status < ?, so status=6 > 5 would never match.
logger.WithContext(ctx).Error("Update order status failed", result := l.svc.DB.WithContext(ctx).
logger.Field("error", err.Error()), Model(&order.Order{}).
Where("order_no = ? AND status = ?", orderInfo.OrderNo, OrderStatusClaimed).
Update("status", OrderStatusFinished)
if result.Error != nil {
logger.WithContext(ctx).Error("Update order status from claimed to finished failed",
logger.Field("error", result.Error.Error()),
logger.Field("order_no", orderInfo.OrderNo),
)
}
// Invalidate cache entries; key format matches order/default.go cacheOrderIdPrefix / cacheOrderNoPrefix.
cacheKeys := []string{
fmt.Sprintf("cache:order:id:%d", orderInfo.Id),
fmt.Sprintf("cache:order:no:%s", orderInfo.OrderNo),
}
if delErr := l.svc.Redis.Del(ctx, cacheKeys...).Err(); delErr != nil {
logger.WithContext(ctx).Error("Failed to invalidate order cache after status update",
logger.Field("error", delErr.Error()),
logger.Field("order_no", orderInfo.OrderNo), logger.Field("order_no", orderInfo.OrderNo),
) )
} }
@@ -0,0 +1,86 @@
package orderLogic
import (
"context"
"encoding/json"
"time"
"github.com/hibiken/asynq"
"github.com/perfect-panel/server/internal/model/order"
"github.com/perfect-panel/server/internal/svc"
"github.com/perfect-panel/server/pkg/logger"
queueTypes "github.com/perfect-panel/server/queue/types"
)
const stuckClaimTimeout = 10 * time.Minute
// StuckOrderRecoveryLogic scans orders stuck in claimed(6) status and resets them to paid(2)
// so that asynq retry can re-claim and process them.
type StuckOrderRecoveryLogic struct {
svc *svc.ServiceContext
}
func NewStuckOrderRecoveryLogic(svc *svc.ServiceContext) *StuckOrderRecoveryLogic {
return &StuckOrderRecoveryLogic{svc: svc}
}
func (l *StuckOrderRecoveryLogic) ProcessTask(ctx context.Context, _ *asynq.Task) error {
cutoff := time.Now().Add(-stuckClaimTimeout)
var stuckOrders []order.Order
if err := l.svc.DB.WithContext(ctx).
Model(&order.Order{}).
Where("status = ? AND updated_at < ?", OrderStatusClaimed, cutoff).
Find(&stuckOrders).Error; err != nil {
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to query stuck orders",
logger.Field("error", err.Error()),
)
return err
}
if len(stuckOrders) == 0 {
return nil
}
orderNos := make([]string, 0, len(stuckOrders))
for i := range stuckOrders {
orderNos = append(orderNos, stuckOrders[i].OrderNo)
}
logger.WithContext(ctx).Error("[StuckOrderRecovery] Found stuck claimed orders, resetting to paid",
logger.Field("count", len(stuckOrders)),
logger.Field("order_nos", orderNos),
)
result := l.svc.DB.WithContext(ctx).
Model(&order.Order{}).
Where("status = ? AND updated_at < ?", OrderStatusClaimed, cutoff).
Update("status", OrderStatusPaid)
if result.Error != nil {
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to reset stuck orders",
logger.Field("error", result.Error.Error()),
)
return result.Error
}
for i := range stuckOrders {
ord := &stuckOrders[i]
payload, err := json.Marshal(queueTypes.ForthwithActivateOrderPayload{OrderNo: ord.OrderNo})
if err != nil {
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to marshal payload",
logger.Field("order_no", ord.OrderNo),
logger.Field("error", err.Error()),
)
continue
}
task := asynq.NewTask(queueTypes.ForthwithActivateOrder, payload, asynq.MaxRetry(5))
if _, err := l.svc.Queue.EnqueueContext(ctx, task); err != nil {
logger.WithContext(ctx).Error("[StuckOrderRecovery] Failed to re-enqueue order",
logger.Field("order_no", ord.OrderNo),
logger.Field("error", err.Error()),
)
}
}
return nil
}
+1
View File
@@ -7,4 +7,5 @@ const (
SchedulerTrafficStat = "scheduler:traffic:stat" SchedulerTrafficStat = "scheduler:traffic:stat"
SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单 SchedulerIAPReconcile = "scheduler:iap:reconcile" // 第二层:每 5 分钟扫描待支付 IAP 订单
SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账 SchedulerIAPDailyReconcile = "scheduler:iap:daily:reconcile" // 第三层:日终全量对账
SchedulerStuckOrderRecovery = "scheduler:stuck:order:recovery" // 定时恢复卡在 claimed 状态的订单
) )
+6
View File
@@ -64,6 +64,12 @@ func (m *Service) Start() {
logger.Errorf("register iap daily reconcile task failed: %s", err.Error()) logger.Errorf("register iap daily reconcile task failed: %s", err.Error())
} }
// schedule stuck order recovery: every 10 minutes
stuckOrderTask := asynq.NewTask(types.SchedulerStuckOrderRecovery, nil)
if _, err := m.server.Register("@every 10m", stuckOrderTask, asynq.MaxRetry(1)); err != nil {
logger.Errorf("register stuck order recovery task failed: %s", err.Error())
}
if err := m.server.Run(); err != nil { if err := m.server.Run(); err != nil {
logger.Errorf("run scheduler failed: %s", err.Error()) logger.Errorf("run scheduler failed: %s", err.Error())
} }
+44
View File
@@ -0,0 +1,44 @@
导出数据库:
mysqldump --socket=/var/run/mysqld/mysqld.sock -uroot -p --single-transaction --routines --triggers --events --set-gtid-purged=OFF --source-data=2 --no-tablespaces hifast > /root/data/hifast-final-0515.sql
jpcV41ppanel
gzip -1 /root/hifast-final-0515.sql
导出redis
redis-cli -a '0BVz9XOHf7KUfEuoFJRK-dURdKUGFiZ8QeaHpysHnKeKhLskZb55HPK121lFsKtr' BGSAVE
sleep 5
cp /var/lib/redis/dump.rdb /root/data/redis-backup-0515.rdb
导入数据库:
gunzip -c /root/hifast-final-0515.sql.gz | docker exec -i ppanel-mysql mysql -uroot -p'jpcV41ppanel' hifast
导入redis
docker stop ppanel-redis
docker cp /root/data/redis-backup-0515.rdb ppanel-redis:/data/dump.rdb
docker start ppanel-redis
docker exec ppanel-redis redis-cli -a 'hifast67yj' DBSIZE
docker stop ppanel-redis || true
docker rm -f ppanel-redis
rm -f /root/data/dump.rdb
rm -rf /root/data/appendonlydir
rm -f /root/data/appendonly.aof*
mkdir -p /root/data
docker pull redis:8.6.3
docker run -d \
--name ppanel-redis \
--restart always \
-p 6379:6379 \
redis:8.6.3 \
redis-server --requirepass 'hifast67yj'