Files
qitongxue-api/internal/logic/event_import.go
2026-09-29 10:57:01 +08:00

557 lines
17 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package logic
import (
"context"
"encoding/json"
"fmt"
"strconv"
"strings"
"github.com/gogf/gf/v2/database/gdb"
"github.com/gogf/gf/v2/errors/gcode"
"github.com/gogf/gf/v2/errors/gerror"
"github.com/gogf/gf/v2/frame/g"
"github.com/gogf/gf/v2/os/gtime"
v1 "tool-api/api/admin/v1"
"tool-api/internal/consts"
"tool-api/internal/model/entity"
)
// ============================================================================
// 年会域:员工导入(两段式)
// 依据:PRD-02 R7-3/R7-4/R7-5、§5.2 异常流;架构 §3.2.8、§4.3 图③、§8.4。
//
// 约定:
// - 管理端(SheetJS)负责解析文件,后端只收结构化行数组,零新增依赖;
// - 预览只写 event_import_logs(status=draft) 快照,**不落 employees**;
// - 防篡改采用 **DB 快照 by id**(非 HMAC):确认时凭 import_log_id 从库读回快照,
// 只接受快照内的行,无需密钥,规避硬编码密钥隐患;
// - 未裁决冲突 → 4008;确认入库全程单事务,含 event_participants UPSERT。
// ============================================================================
// gcodeImportConflict 导入存在未处理冲突(4008)
func gcodeImportConflict() gcode.Code {
return gcode.New(consts.CodeImportConflict, "", nil)
}
// ============================================================================
// 纯类型与纯函数(不碰 DB,单测直接覆盖)
// ============================================================================
// importRow 一行导入数据(来自管理端解析)
type importRow struct {
Name string
Phone string
Dept string
Remark string
}
// existingEmployee 现有员工快照(分类用)
type existingEmployee struct {
Id int64
Name string
Phone string
}
// classifiedRow 分类后的行(会序列化进快照)
type classifiedRow struct {
Idx int `json:"idx"`
Type string `json:"type"` // new / update / conflict / error
Name string `json:"name"`
Phone string `json:"phone"`
Dept string `json:"dept"`
Remark string `json:"remark"`
Message string `json:"message"`
MatchedEmployeeId int64 `json:"matched_employee_id"` // update:同号命中;conflict:同名命中
}
// importSnapshot 预览快照(落 event_import_logs.snapshot)
type importSnapshot struct {
EnterpriseId int64 `json:"enterprise_id"`
EventId int64 `json:"event_id"`
Rows []classifiedRow `json:"rows"`
}
// importOp 一条待执行操作
type importOp struct {
Kind string // insert / update / skip
RowIdx int
Name string
Phone string
Dept string
Remark string
UpdateId int64
}
// classifyImportRows 逐行分类(R7-4/R7-5):
//
// 手机号缺失/格式非法 → error
// 手机号命中现有员工 → update(视为同一人,更新姓名/部门/备注)
// 姓名命中现有员工(手机号不同) → conflict(疑似换了手机号,需人工裁决)
// 其余 → new
func classifyImportRows(rows []importRow, existing []existingEmployee) []classifiedRow {
byPhone := make(map[string]existingEmployee, len(existing))
byName := make(map[string][]existingEmployee)
for _, e := range existing {
byPhone[e.Phone] = e
byName[e.Name] = append(byName[e.Name], e)
}
out := make([]classifiedRow, 0, len(rows))
for i, r := range rows {
name := strings.TrimSpace(r.Name)
phone := normalizePhone(r.Phone)
c := classifiedRow{Idx: i, Name: name, Phone: phone, Dept: r.Dept, Remark: r.Remark}
switch {
case phone == "":
c.Type = consts.ImportTypeError
c.Message = "手机号缺失"
case !validPhone(phone):
c.Type = consts.ImportTypeError
c.Message = "手机号格式不正确"
default:
if e, ok := byPhone[phone]; ok {
c.Type = consts.ImportTypeUpdate
c.MatchedEmployeeId = e.Id
c.Message = "同手机号,更新"
} else if list := byName[name]; len(list) > 0 {
c.Type = consts.ImportTypeConflict
c.MatchedEmployeeId = list[0].Id
c.Message = "同名不同号,疑似换了手机号"
} else {
c.Type = consts.ImportTypeNew
c.Message = "将新增"
}
}
out = append(out, c)
}
return out
}
// planImportOps 结合人工裁决生成待执行操作。
// 返回 unresolved = 仍未裁决(缺少合法裁决)的 conflict 行数;>0 时确认接口应返回 4008。
func planImportOps(rows []classifiedRow, decisions map[int]string) (ops []importOp, unresolved int) {
ops = make([]importOp, 0, len(rows))
for _, c := range rows {
switch c.Type {
case consts.ImportTypeError:
ops = append(ops, importOp{Kind: "skip", RowIdx: c.Idx})
case consts.ImportTypeNew:
ops = append(ops, importOp{Kind: "insert", RowIdx: c.Idx, Name: c.Name, Phone: c.Phone, Dept: c.Dept, Remark: c.Remark})
case consts.ImportTypeUpdate:
ops = append(ops, importOp{Kind: "update", RowIdx: c.Idx, UpdateId: c.MatchedEmployeeId, Name: c.Name, Phone: c.Phone, Dept: c.Dept, Remark: c.Remark})
case consts.ImportTypeConflict:
switch decisions[c.Idx] {
case consts.ImportActionMerge:
// 视为同一人·更新手机号 → 更新现有员工(含手机号)
ops = append(ops, importOp{Kind: "update", RowIdx: c.Idx, UpdateId: c.MatchedEmployeeId, Name: c.Name, Phone: c.Phone, Dept: c.Dept, Remark: c.Remark})
case consts.ImportActionCreate:
// 视为另一个人新增(手机号不同,不违反 uk_ent_phone)
ops = append(ops, importOp{Kind: "insert", RowIdx: c.Idx, Name: c.Name, Phone: c.Phone, Dept: c.Dept, Remark: c.Remark})
case consts.ImportActionIgnore:
ops = append(ops, importOp{Kind: "skip", RowIdx: c.Idx})
default:
unresolved++
}
}
}
return ops, unresolved
}
// validPhone 基础校验:仅数字,长度 6..20。
func validPhone(phone string) bool {
if len(phone) < 6 || len(phone) > 20 {
return false
}
for _, r := range phone {
if r < '0' || r > '9' {
return false
}
}
return true
}
// ============================================================================
// 预览
// ============================================================================
// AdminImportPreview 导入预览:分类 + 写 draft 快照;不落 employees。
func AdminImportPreview(ctx context.Context, req *v1.ImportPreviewReq) (*v1.ImportPreviewRes, error) {
limit := importRowLimit(ctx)
if len(req.Rows) > limit {
return nil, gerror.New(fmt.Sprintf("单次导入最多 %d 行,请拆分后再试", limit))
}
if exists, err := g.Model(consts.TableEnterprises).Where("id", req.EnterpriseId).Count(); err != nil {
return nil, err
} else if exists == 0 {
return nil, gerror.New("企业不存在")
}
if req.EventId > 0 {
if exists, err := g.Model(consts.TableAnnualEvents).Where("id", req.EventId).Count(); err != nil {
return nil, err
} else if exists == 0 {
return nil, gerror.New("活动不存在")
}
}
existing, err := loadExistingEmployees(ctx, req.EnterpriseId)
if err != nil {
return nil, err
}
rows := make([]importRow, 0, len(req.Rows))
for _, r := range req.Rows {
rows = append(rows, importRow{Name: r.Name, Phone: r.Phone, Dept: r.Dept, Remark: r.Remark})
}
classified := classifyImportRows(rows, existing)
snapBytes, err := json.Marshal(importSnapshot{
EnterpriseId: req.EnterpriseId,
EventId: req.EventId,
Rows: classified,
})
if err != nil {
return nil, err
}
counts := countClassified(classified)
now := gtime.Now()
logId, err := g.Model(consts.TableEventImportLogs).Data(g.Map{
"enterprise_id": req.EnterpriseId,
"event_id": req.EventId,
"file_name": req.FileName,
"status": consts.ImportStatusDraft,
"token": "",
"snapshot": string(snapBytes),
"total": len(classified),
"inserted": counts.inserted,
"updated": counts.updated,
"skipped": counts.skipped,
"conflict": counts.conflict,
"conflicts": "",
"operator": adminOperator(ctx),
"created_at": now,
"updated_at": now,
}).InsertAndGetId()
if err != nil {
return nil, err
}
rowsOut := make([]v1.ImportPreviewRowOut, 0, len(classified))
for _, c := range classified {
rowsOut = append(rowsOut, v1.ImportPreviewRowOut{
Idx: c.Idx,
Type: c.Type,
Name: c.Name,
Phone: c.Phone,
Dept: c.Dept,
Message: c.Message,
})
}
return &v1.ImportPreviewRes{
ImportLogId: logId,
Total: len(classified),
Inserted: counts.inserted,
Updated: counts.updated,
Skipped: counts.skipped,
Conflict: counts.conflict,
Rows: rowsOut,
}, nil
}
// ============================================================================
// 确认入库(单事务)
// ============================================================================
// AdminImportConfirm 确认导入:读回快照 → 校验裁决完整性 → 单事务写 employees + UPSERT 候场名单 + 置 committed。
func AdminImportConfirm(ctx context.Context, req *v1.ImportConfirmReq) (*v1.ImportConfirmRes, error) {
logRec, err := g.Model(consts.TableEventImportLogs).Where("id", req.ImportLogId).One()
if err != nil {
return nil, err
}
if logRec.IsEmpty() {
return nil, gerror.New("导入记录不存在")
}
if logRec["status"].String() != consts.ImportStatusDraft {
return nil, gerror.New("该导入已处理,请重新预览")
}
snapshot := &importSnapshot{}
if err = json.Unmarshal([]byte(logRec["snapshot"].String()), snapshot); err != nil {
return nil, gerror.New("导入快照损坏,请重新预览")
}
// 只允许裁决快照中的 conflict 行
conflictIdx := make(map[int]bool)
for _, c := range snapshot.Rows {
if c.Type == consts.ImportTypeConflict {
conflictIdx[c.Idx] = true
}
}
decisions := make(map[int]string, len(req.Decisions))
for _, d := range req.Decisions {
if !conflictIdx[d.Idx] {
return nil, gerror.NewCode(gcode.CodeValidationFailed, "裁决行索引非法")
}
decisions[d.Idx] = d.Action
}
ops, unresolved := planImportOps(snapshot.Rows, decisions)
if unresolved > 0 {
return nil, gerror.NewCode(gcodeImportConflict(), "存在未处理的导入冲突")
}
var (
inserted int
updated int
skipped int
)
err = g.DB().Transaction(ctx, func(ctx context.Context, tx gdb.TX) error {
inserted, updated, skipped = 0, 0, 0
now := gtime.Now()
for _, op := range ops {
switch op.Kind {
case "skip":
skipped++
case "insert":
empId, isInserted, e := insertEmployeeTx(ctx, tx, snapshot.EnterpriseId, op, now)
if e != nil {
return e
}
if isInserted {
inserted++
} else {
updated++ // 并发/重导命中唯一键 → 转更新
}
if snapshot.EventId > 0 {
if e = addParticipantTx(ctx, tx, snapshot.EventId, empId, now); e != nil {
return e
}
}
case "update":
if _, e := tx.Model(consts.TableEmployees).Where("id", op.UpdateId).Data(g.Map{
"name": op.Name,
"phone": op.Phone,
"dept": op.Dept,
"remark": op.Remark,
"updated_at": now,
}).Update(); e != nil {
return e
}
updated++
if snapshot.EventId > 0 {
if e := addParticipantTx(ctx, tx, snapshot.EventId, op.UpdateId, now); e != nil {
return e
}
}
}
}
conflictsJson, _ := json.Marshal(decisions)
if _, e := tx.Model(consts.TableEventImportLogs).Where("id", req.ImportLogId).Data(g.Map{
"status": consts.ImportStatusCommitted,
"inserted": inserted,
"updated": updated,
"skipped": skipped,
"conflict": len(conflictIdx),
"conflicts": string(conflictsJson),
"updated_at": now,
}).Update(); e != nil {
return e
}
return nil
})
if err != nil {
return nil, err
}
// T17 审计:确认导入属敏感操作
WriteAudit(ctx, AuditEntry{
Action: "event.import_confirm",
TargetType: "import_log",
TargetId: auditId(req.ImportLogId),
After: g.Map{
"enterprise_id": snapshot.EnterpriseId,
"event_id": snapshot.EventId,
"inserted": inserted,
"updated": updated,
"skipped": skipped,
},
Result: consts.AuditResultSuccess,
})
return &v1.ImportConfirmRes{Inserted: inserted, Updated: updated, Skipped: skipped}, nil
}
// AdminImportLogList 导入历史
func AdminImportLogList(ctx context.Context, req *v1.ImportLogListReq) (*v1.ImportLogListRes, error) {
model := g.Model(consts.TableEventImportLogs)
if req.EnterpriseId > 0 {
model = model.Where("enterprise_id", req.EnterpriseId)
}
if req.EventId > 0 {
model = model.Where("event_id", req.EventId)
}
total, err := model.Count()
if err != nil {
return nil, err
}
page, size := normalizeAdminPage(req.Page, req.PageSize)
records, err := model.OrderDesc("id").Page(page, size).All()
if err != nil {
return nil, err
}
list := make([]v1.ImportLogItem, 0, len(records))
for _, r := range records {
log := &entity.EventImportLogs{}
if err = r.Struct(log); err != nil {
continue
}
list = append(list, v1.ImportLogItem{
Id: log.Id,
EnterpriseId: log.EnterpriseId,
EventId: log.EventId,
FileName: log.FileName,
Status: log.Status,
Total: log.Total,
Inserted: log.Inserted,
Updated: log.Updated,
Skipped: log.Skipped,
Conflict: log.Conflict,
Operator: log.Operator,
CreatedAt: formatGTime(log.CreatedAt),
})
}
return &v1.ImportLogListRes{List: list, Total: total}, nil
}
// ============================================================================
// 内部辅助
// ============================================================================
type importCounts struct {
inserted int
updated int
skipped int
conflict int
}
// countClassified 预览口径计数:new→inserted,update→updated,error→skipped,conflict→conflict。
func countClassified(rows []classifiedRow) importCounts {
c := importCounts{}
for _, r := range rows {
switch r.Type {
case consts.ImportTypeNew:
c.inserted++
case consts.ImportTypeUpdate:
c.updated++
case consts.ImportTypeError:
c.skipped++
case consts.ImportTypeConflict:
c.conflict++
}
}
return c
}
// loadExistingEmployees 读取企业现有员工(分类用)
func loadExistingEmployees(ctx context.Context, enterpriseId int64) ([]existingEmployee, error) {
// 软删除:已删员工不参与「同号/同名」分类,重新导入同一手机号视为新增
records, err := g.Model(consts.TableEmployees).Where("enterprise_id", enterpriseId).WhereNull("deleted_at").All()
if err != nil {
return nil, err
}
out := make([]existingEmployee, 0, len(records))
for _, r := range records {
out = append(out, existingEmployee{
Id: r["id"].Int64(),
Name: r["name"].String(),
Phone: r["phone"].String(),
})
}
return out, nil
}
// insertEmployeeTx 事务内插入员工;命中唯一键(并发/重导)则转更新,保证 uk_ent_phone 不产生重复。
func insertEmployeeTx(ctx context.Context, tx gdb.TX, enterpriseId int64, op importOp, now *gtime.Time) (int64, bool, error) {
id, err := tx.Model(consts.TableEmployees).Data(g.Map{
"enterprise_id": enterpriseId,
"name": op.Name,
"phone": op.Phone,
"dept": op.Dept,
"remark": op.Remark,
"status": 1,
"created_at": now,
"updated_at": now,
}).InsertAndGetId()
if err != nil {
if isDuplicateErr(err) {
rec, e := tx.Model(consts.TableEmployees).
Where("enterprise_id", enterpriseId).Where("phone", op.Phone).One()
if e != nil {
return 0, false, e
}
if rec.IsEmpty() {
return 0, false, err
}
existingId := rec["id"].Int64()
if _, e = tx.Model(consts.TableEmployees).Where("id", existingId).Data(g.Map{
"name": op.Name,
"dept": op.Dept,
"remark": op.Remark,
"updated_at": now,
}).Update(); e != nil {
return 0, false, e
}
return existingId, false, nil
}
return 0, false, err
}
return id, true, nil
}
// addParticipantTx 事务内幂等纳入候场名单(UPSERT 语义,靠 UNIQUE uk_event_emp 避免重复挂名单)。
func addParticipantTx(ctx context.Context, tx gdb.TX, eventId, employeeId int64, now *gtime.Time) error {
exists, err := tx.Model(consts.TableEventParticipants).
Where("event_id", eventId).Where("employee_id", employeeId).Count()
if err != nil {
return err
}
if exists > 0 {
return nil
}
if _, err = tx.Model(consts.TableEventParticipants).Data(g.Map{
"event_id": eventId,
"employee_id": employeeId,
"created_at": now,
}).Insert(); err != nil {
// 并发下可能已被插入(uk_event_emp):忽略,保持幂等
return nil
}
return nil
}
// isDuplicateErr 判断是否 MySQL 唯一键冲突(1062)。
func isDuplicateErr(err error) bool {
if err == nil {
return false
}
msg := err.Error()
return strings.Contains(msg, "Duplicate") || strings.Contains(msg, "1062") || strings.Contains(msg, "uk_ent_phone")
}
// adminOperator 记录操作人(管理端 id),无上下文时为空串。
func adminOperator(ctx context.Context) string {
id := CtxAdminId(ctx)
if id <= 0 {
return ""
}
return strconv.FormatInt(id, 10)
}
// importRowLimit 单次导入行数上限:settings 覆盖 > 默认 2000。
func importRowLimit(ctx context.Context) int {
raw := SettingValue(ctx, consts.SettingImportRowLimit)
if n, err := strconv.Atoi(raw); err == nil && n > 0 {
return n
}
return consts.DefaultImportRowLimit
}