整合数据

This commit is contained in:
2026-06-16 01:30:39 +08:00
parent 761a5cb69c
commit c0f70823a9
31 changed files with 4385 additions and 893 deletions
+185 -155
View File
@@ -619,115 +619,150 @@ func replenishPoolRow(c *beego.Controller, module string) {
platform := payload.Platform
remark := strings.TrimSpace(payload.Remark)
replenishWithProbe(c, module, payload.Type, platform, remark, now)
}
type poolReplenishCandidate struct {
id uint64
dataType string
token string
isUsed *int8
row interface{}
}
type poolReplenishFetcher func() (*poolReplenishCandidate, error)
// replenishWithProbe 按 id 顺序补号并探测;不可用则标记 is_extracted=2 后继续下一条。
func replenishWithProbe(c *beego.Controller, module, dataType, platform, remark string, now time.Time) {
var fetch poolReplenishFetcher
switch module {
case "cursor":
checkedCount := 0
unavailableCount := 0
for {
fetch = func() (*poolReplenishCandidate, error) {
var row models.PlatformAccountPoolCursor
if err := models.Orm.QueryTable(new(models.PlatformAccountPoolCursor)).
Filter("is_extracted", 0).Filter("data_type", payload.Type).
OrderBy("id").One(&row); err != nil {
msg := "暂无可用账号"
if checkedCount > 0 {
msg = fmt.Sprintf("已检测%d个账号,其中%d个不可用,暂无可用账号", checkedCount, unavailableCount)
}
poolJSONErr(c, 404, 404, msg)
return
err := models.Orm.QueryTable(new(models.PlatformAccountPoolCursor)).
Filter("is_extracted", 0).
Filter("data_type", dataType).
Filter("delete_time__isnull", true).
OrderBy("id").
One(&row)
if err != nil {
return nil, err
}
checkedCount++
isAvailable := poolProbeToken("cursor", row.DataType, row.Token, row.ID)
if !isAvailable {
unavailableCount++
if _, err := models.Orm.QueryTable(new(models.PlatformAccountPoolCursor)).Filter("id", row.ID).Update(map[string]interface{}{
// 补号流程检测出来不可用/已用完的号,仍然归类为“补号”记录。
// 不要写成已提取/已用完状态;只有接口提取后再标记不可用的号才归到已提取侧。
"is_extracted": int8(2),
"is_used": int8(0),
"extracted_time": now,
"extracted_platform": platform,
"remark": remark,
"update_time": now,
}); err != nil {
poolJSONErr(c, 500, 500, "补号检测失败: "+err.Error())
return
}
continue
}
if _, err := models.Orm.QueryTable(new(models.PlatformAccountPoolCursor)).Filter("id", row.ID).Update(map[string]interface{}{
"is_extracted": int8(2),
"is_used": int8(1),
"extracted_time": now,
"extracted_platform": platform,
"remark": remark,
"update_time": now,
}); err != nil {
poolJSONErr(c, 500, 500, "补号失败: "+err.Error())
return
}
row.IsExtracted = 2
isUsed := int8(1)
row.IsUsed = &isUsed
row.ExtractedTime = &now
row.ExtractedPlatform = &platform
row.Remark = remark
c.Data["json"] = map[string]interface{}{
"code": 200,
"msg": "补号成功",
"data": row,
"probe": map[string]interface{}{
"checkedCount": checkedCount,
"unavailableCount": unavailableCount,
},
}
break
return &poolReplenishCandidate{
id: row.ID, dataType: row.DataType, token: row.Token, isUsed: row.IsUsed, row: row,
}, nil
}
case "windsurf":
var row models.PlatformAccountPoolWindsurf
if err := models.Orm.QueryTable(new(models.PlatformAccountPoolWindsurf)).
Filter("is_extracted", 0).Filter("data_type", payload.Type).
OrderBy("id").One(&row); err != nil {
poolJSONErr(c, 404, 404, "暂无可用账号")
return
fetch = func() (*poolReplenishCandidate, error) {
var row models.PlatformAccountPoolWindsurf
err := models.Orm.QueryTable(new(models.PlatformAccountPoolWindsurf)).
Filter("is_extracted", 0).
Filter("data_type", dataType).
Filter("delete_time__isnull", true).
OrderBy("id").
One(&row)
if err != nil {
return nil, err
}
return &poolReplenishCandidate{
id: row.ID, dataType: row.DataType, token: row.Token, row: row,
}, nil
}
if _, err = models.Orm.QueryTable(new(models.PlatformAccountPoolWindsurf)).Filter("id", row.ID).Update(map[string]interface{}{
"is_extracted": int8(2), "extracted_time": now, "extracted_platform": platform, "remark": remark,
}); err != nil {
poolJSONErr(c, 500, 500, "补号失败: "+err.Error())
return
}
row.IsExtracted = 2
row.ExtractedTime = &now
row.ExtractedPlatform = &platform
row.Remark = remark
c.Data["json"] = map[string]interface{}{"code": 200, "msg": "补号成功", "data": row}
case "krio":
var row models.PlatformAccountPoolKiro
if err := models.Orm.QueryTable(new(models.PlatformAccountPoolKiro)).
Filter("is_extracted", 0).Filter("data_type", payload.Type).
OrderBy("id").One(&row); err != nil {
poolJSONErr(c, 404, 404, "暂无可用账号")
return
fetch = func() (*poolReplenishCandidate, error) {
var row models.PlatformAccountPoolKiro
err := models.Orm.QueryTable(new(models.PlatformAccountPoolKiro)).
Filter("is_extracted", 0).
Filter("data_type", dataType).
Filter("delete_time__isnull", true).
OrderBy("id").
One(&row)
if err != nil {
return nil, err
}
return &poolReplenishCandidate{
id: row.ID, dataType: row.DataType, token: row.Token, row: row,
}, nil
}
if _, err = models.Orm.QueryTable(new(models.PlatformAccountPoolKiro)).Filter("id", row.ID).Update(map[string]interface{}{
"is_extracted": int8(2), "extracted_time": now, "extracted_platform": platform, "remark": remark,
}); err != nil {
poolJSONErr(c, 500, 500, "补号失败: "+err.Error())
return
}
row.IsExtracted = 2
row.ExtractedTime = &now
row.ExtractedPlatform = &platform
row.Remark = remark
c.Data["json"] = map[string]interface{}{"code": 200, "msg": "补号成功", "data": row}
default:
poolJSONErr(c, 400, 400, "无效模块")
return
}
_ = c.ServeJSON()
tableName := poolTableName(module)
if tableName == "" {
poolJSONErr(c, 400, 400, "无效模块")
return
}
for {
candidate, err := fetch()
if err != nil {
if err == orm.ErrNoRows {
poolJSONErr(c, 404, 404, "暂无可用账号")
} else {
poolJSONErr(c, 500, 500, "查询失败")
}
return
}
updateFields := map[string]interface{}{
"is_extracted": int8(2),
"extracted_time": now,
"extracted_platform": platform,
"remark": remark,
"update_time": now,
}
if _, err = models.Orm.QueryTable(tableName).
Filter("id", candidate.id).
Update(updateFields); err != nil {
poolJSONErr(c, 500, 500, "补号失败: "+err.Error())
return
}
if known, available := poolIsUsedAvailable(candidate.isUsed); known {
if !available {
continue
}
} else if !poolProbeToken(module, candidate.dataType, candidate.token, candidate.id) {
continue
}
data := replenishApplyResponse(candidate.row, platform, remark, now)
c.Data["json"] = map[string]interface{}{"code": 200, "msg": "补号成功", "data": data}
_ = c.ServeJSON()
return
}
}
func replenishApplyResponse(row interface{}, platform, remark string, now time.Time) interface{} {
pf := platform
switch r := row.(type) {
case models.PlatformAccountPoolCursor:
r.IsExtracted = 2
r.ExtractedTime = &now
r.ExtractedPlatform = &pf
r.Remark = remark
if r.IsUsed == nil || *r.IsUsed != 1 {
used := int8(1)
r.IsUsed = &used
}
return r
case models.PlatformAccountPoolWindsurf:
r.IsExtracted = 2
r.ExtractedTime = &now
r.ExtractedPlatform = &pf
r.Remark = remark
return r
case models.PlatformAccountPoolKiro:
r.IsExtracted = 2
r.ExtractedTime = &now
r.ExtractedPlatform = &pf
r.Remark = remark
return r
default:
return row
}
}
func updatePoolRemark(c *beego.Controller, module string) {
@@ -843,6 +878,54 @@ func setPoolUnavailable(c *beego.Controller, module string) {
_ = c.ServeJSON()
}
func updatePoolUsable(c *beego.Controller, module string) {
if _, err := requirePlatformAuth(c); err != nil {
poolJSONErr(c, 401, 401, err.Error())
return
}
if module != "cursor" {
poolJSONErr(c, 400, 400, "该模块不支持可用状态修改")
return
}
raw, err := io.ReadAll(c.Ctx.Request.Body)
if err != nil {
poolJSONErr(c, 400, 400, "参数错误")
return
}
var payload struct {
ID uint64 `json:"id"`
Usable int `json:"usable"`
}
if err := json.Unmarshal(raw, &payload); err != nil || payload.ID == 0 {
poolJSONErr(c, 400, 400, "参数错误")
return
}
if payload.Usable != 0 && payload.Usable != 1 {
poolJSONErr(c, 400, 400, "可用状态参数错误")
return
}
now := time.Now()
updated, err := models.Orm.QueryTable(new(models.PlatformAccountPoolCursor)).Filter("id", payload.ID).Update(orm.Params{
"is_used": int8(payload.Usable),
"update_time": now,
})
if err != nil {
poolJSONErr(c, 500, 500, "可用状态更新失败: "+err.Error())
return
}
if updated == 0 {
poolJSONErr(c, 404, 404, "记录不存在")
return
}
msg := "已标记不可用"
if payload.Usable == 1 {
msg = "已标记可用"
}
c.Data["json"] = map[string]interface{}{"code": 200, "msg": msg}
_ = c.ServeJSON()
}
func updatePoolPlatform(c *beego.Controller, module string) {
if _, err := requirePlatformAuth(c); err != nil {
poolJSONErr(c, 401, 401, err.Error())
@@ -1019,12 +1102,12 @@ func probePoolToken(c *beego.Controller, module string) {
if r.StreamNote != "" {
data["streamNote"] = r.StreamNote
}
// Cursor 探测状态只按底层探针结论 r.OK 保存。
// 注意:客户端版本过旧只是 warningToken 仍可用时 r.OK=true,不能因此写成已用完。
if module == "cursor" && payload.ID > 0 && r.HTTPStatus == http.StatusOK {
isUsed := int8(0)
var isUsed int8
if r.OK {
isUsed = 1
} else {
isUsed = 0
}
if _, uerr := models.Orm.QueryTable(new(models.PlatformAccountPoolCursor)).Filter("id", payload.ID).Update(orm.Params{
"is_used": isUsed,
@@ -1041,62 +1124,6 @@ func probePoolToken(c *beego.Controller, module string) {
_ = c.ServeJSON()
}
func poolTableName(module string) string {
switch module {
case "cursor":
return (&models.PlatformAccountPoolCursor{}).TableName()
case "windsurf":
return (&models.PlatformAccountPoolWindsurf{}).TableName()
case "krio":
return (&models.PlatformAccountPoolKiro{}).TableName()
default:
return ""
}
}
func poolIsUsedAvailable(isUsed *int8) (known bool, available bool) {
if isUsed == nil {
return false, false
}
switch *isUsed {
case 1:
return true, true
case 0:
return true, false
default:
return false, false
}
}
func poolProbeToken(module, rowDataType, token string, id uint64) bool {
token = strings.TrimSpace(token)
if rowDataType == "account" || token == "" {
return true
}
r := tokenprobe.ProbeOfficial(module, token)
// Cursor 自动探测只按底层探针结论 r.OK 判定。
// 客户端版本过旧是 warning,不代表 Token 已用完;只有 tokenprobe 明确判定额度用尽/不可用时 r.OK 才为 false。
available := r.OK
// 更新数据库中的 is_used 字段
if module == "cursor" && id > 0 {
isUsed := int8(0)
if available {
isUsed = 1
}
_, _ = models.Orm.QueryTable(new(models.PlatformAccountPoolCursor)).
Filter("id", id).
Update(orm.Params{
"is_used": isUsed,
"update_time": time.Now(),
})
}
return available
}
func (c *PlatformAccountPoolCursorController) List() { listPoolRows(&c.Controller, "cursor") }
func (c *PlatformAccountPoolCursorController) Add() { addPoolRow(&c.Controller, "cursor") }
func (c *PlatformAccountPoolCursorController) BatchAdd() { batchAddPoolRows(&c.Controller, "cursor") }
@@ -1109,6 +1136,9 @@ func (c *PlatformAccountPoolCursorController) UpdateRemark() {
func (c *PlatformAccountPoolCursorController) SetUnavailable() {
setPoolUnavailable(&c.Controller, "cursor")
}
func (c *PlatformAccountPoolCursorController) UpdateUsable() {
updatePoolUsable(&c.Controller, "cursor")
}
func (c *PlatformAccountPoolCursorController) UpdatePlatform() {
updatePoolPlatform(&c.Controller, "cursor")
}