279 lines
8.3 KiB
Go
279 lines
8.3 KiB
Go
package db
|
||
|
||
import (
|
||
"database/sql"
|
||
"fmt"
|
||
"strings"
|
||
"time"
|
||
)
|
||
|
||
// SyncJobRow hy_sync_job 查询行。
|
||
type SyncJobRow struct {
|
||
ID int64
|
||
AnchorDate string
|
||
Step string
|
||
Status string
|
||
TotalCount int
|
||
SuccessCount int
|
||
FailedCount int
|
||
ErrorMessage sql.NullString
|
||
StartedAt sql.NullTime
|
||
FinishedAt sql.NullTime
|
||
UpdatedAt time.Time
|
||
}
|
||
|
||
// PushLogRow hy_push_log 查询行。
|
||
type PushLogRow struct {
|
||
ID int64
|
||
RecordID int64
|
||
HTTPCode int
|
||
MsgCode sql.NullInt64
|
||
Msg string
|
||
TraceID string
|
||
RequestBodyLen int
|
||
RequestPlainJSON sql.NullString
|
||
RequestHeadersJSON sql.NullString
|
||
ResponseBody sql.NullString
|
||
DurationMs int
|
||
CreatedAt time.Time
|
||
}
|
||
|
||
// RecordWithLog 业务记录 + 最新一条推送日志。
|
||
type RecordWithLog struct {
|
||
ID int64
|
||
JobID int64
|
||
AnchorDate string
|
||
Method string
|
||
ServiceMethod string
|
||
BussID string
|
||
BizKey string
|
||
PayloadHash string
|
||
PayloadJSON string
|
||
ValidationErrors sql.NullString
|
||
PushStatus string
|
||
RetryCount int
|
||
LastError sql.NullString
|
||
UpdatedAt time.Time
|
||
LastLog *PushLogRow
|
||
}
|
||
|
||
// RecordDetail 记录详情(含全部推送日志)。
|
||
type RecordDetail struct {
|
||
Record RecordWithLog
|
||
Logs []PushLogRow
|
||
}
|
||
|
||
// ListJobs 查询同步任务,按 started_at 降序。
|
||
func (s *Store) ListJobs(limit int, anchorDate, step string) ([]SyncJobRow, error) {
|
||
if limit <= 0 || limit > 200 {
|
||
limit = 50
|
||
}
|
||
q := `SELECT id, anchor_date, step, status, total_count, success_count, failed_count,
|
||
error_message, started_at, finished_at, updated_at
|
||
FROM hy_sync_job WHERE 1=1`
|
||
var args []any
|
||
if anchorDate != "" {
|
||
q += ` AND anchor_date = ?`
|
||
args = append(args, anchorDate)
|
||
}
|
||
if step != "" {
|
||
q += ` AND step = ?`
|
||
args = append(args, step)
|
||
}
|
||
q += ` ORDER BY COALESCE(started_at, updated_at) DESC LIMIT ?`
|
||
args = append(args, limit)
|
||
|
||
rows, err := s.db.Query(q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
var out []SyncJobRow
|
||
for rows.Next() {
|
||
var j SyncJobRow
|
||
var ad time.Time
|
||
if err := rows.Scan(&j.ID, &ad, &j.Step, &j.Status, &j.TotalCount, &j.SuccessCount, &j.FailedCount,
|
||
&j.ErrorMessage, &j.StartedAt, &j.FinishedAt, &j.UpdatedAt); err != nil {
|
||
return nil, err
|
||
}
|
||
j.AnchorDate = ad.Format("2006-01-02")
|
||
out = append(out, j)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
// ListRecordsByAnchorMethod 按锚定日与监管 method 查业务记录及最新日志。
|
||
func (s *Store) ListRecordsByAnchorMethod(anchorDate, method string, limit int) ([]RecordWithLog, error) {
|
||
if limit <= 0 || limit > 500 {
|
||
limit = 100
|
||
}
|
||
rows, err := s.db.Query(`
|
||
SELECT r.id, r.job_id, r.anchor_date, r.method, r.service_method, r.buss_id, r.biz_key,
|
||
r.payload_hash, r.payload_json, r.validation_errors, r.push_status, r.retry_count, r.last_error, r.updated_at,
|
||
l.id, l.record_id, l.http_code, l.msg_code, l.msg, l.trace_id, l.request_body_len, l.request_plain_json, l.request_headers_json, l.response_body, l.duration_ms, l.created_at
|
||
FROM hy_push_record r
|
||
LEFT JOIN hy_push_log l ON l.id = (
|
||
SELECT id FROM hy_push_log WHERE record_id = r.id ORDER BY id DESC LIMIT 1
|
||
)
|
||
WHERE r.anchor_date = ? AND r.method = ?
|
||
ORDER BY r.updated_at DESC
|
||
LIMIT ?`, anchorDate, method, limit)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
return scanRecordsWithOptionalLog(rows)
|
||
}
|
||
|
||
// ListRecordsRecent 最近更新的业务记录(跨 step),含最新日志。
|
||
func (s *Store) ListRecordsRecent(limit int, anchorDate string) ([]RecordWithLog, error) {
|
||
if limit <= 0 || limit > 500 {
|
||
limit = 100
|
||
}
|
||
q := `
|
||
SELECT r.id, r.job_id, r.anchor_date, r.method, r.service_method, r.buss_id, r.biz_key,
|
||
r.payload_hash, r.payload_json, r.validation_errors, r.push_status, r.retry_count, r.last_error, r.updated_at,
|
||
l.id, l.record_id, l.http_code, l.msg_code, l.msg, l.trace_id, l.request_body_len, l.request_plain_json, l.request_headers_json, l.response_body, l.duration_ms, l.created_at
|
||
FROM hy_push_record r
|
||
LEFT JOIN hy_push_log l ON l.id = (
|
||
SELECT id FROM hy_push_log WHERE record_id = r.id ORDER BY id DESC LIMIT 1
|
||
)
|
||
WHERE 1=1`
|
||
var args []any
|
||
if anchorDate != "" {
|
||
q += ` AND r.anchor_date = ?`
|
||
args = append(args, anchorDate)
|
||
}
|
||
q += ` ORDER BY r.updated_at DESC LIMIT ?`
|
||
args = append(args, limit)
|
||
|
||
rows, err := s.db.Query(q, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
return scanRecordsWithOptionalLog(rows)
|
||
}
|
||
|
||
func scanRecordsWithOptionalLog(rows *sql.Rows) ([]RecordWithLog, error) {
|
||
var out []RecordWithLog
|
||
for rows.Next() {
|
||
var r RecordWithLog
|
||
var ad time.Time
|
||
var logID sql.NullInt64
|
||
var logRecID sql.NullInt64
|
||
var httpCode sql.NullInt64
|
||
var msgCode sql.NullInt64
|
||
var msg, traceID sql.NullString
|
||
var bodyLen sql.NullInt64
|
||
var plainJSON, headersJSON sql.NullString
|
||
var respBody sql.NullString
|
||
var dur sql.NullInt64
|
||
var logCreated sql.NullTime
|
||
|
||
if err := rows.Scan(
|
||
&r.ID, &r.JobID, &ad, &r.Method, &r.ServiceMethod, &r.BussID, &r.BizKey,
|
||
&r.PayloadHash, &r.PayloadJSON, &r.ValidationErrors, &r.PushStatus, &r.RetryCount, &r.LastError, &r.UpdatedAt,
|
||
&logID, &logRecID, &httpCode, &msgCode, &msg, &traceID, &bodyLen, &plainJSON, &headersJSON, &respBody, &dur, &logCreated,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
r.AnchorDate = ad.Format("2006-01-02")
|
||
if logID.Valid {
|
||
r.LastLog = &PushLogRow{
|
||
ID: logID.Int64,
|
||
RecordID: logRecID.Int64,
|
||
HTTPCode: int(httpCode.Int64),
|
||
Msg: msg.String,
|
||
TraceID: traceID.String,
|
||
RequestBodyLen: int(bodyLen.Int64),
|
||
RequestPlainJSON: plainJSON,
|
||
RequestHeadersJSON: headersJSON,
|
||
ResponseBody: respBody,
|
||
DurationMs: int(dur.Int64),
|
||
}
|
||
if msgCode.Valid {
|
||
r.LastLog.MsgCode = sql.NullInt64{Int64: msgCode.Int64, Valid: true}
|
||
}
|
||
if logCreated.Valid {
|
||
r.LastLog.CreatedAt = logCreated.Time
|
||
}
|
||
}
|
||
out = append(out, r)
|
||
}
|
||
return out, rows.Err()
|
||
}
|
||
|
||
// GetRecordDetail 按 id 查记录及全部推送日志。
|
||
func (s *Store) GetRecordDetail(recordID int64) (*RecordDetail, error) {
|
||
var r RecordWithLog
|
||
var ad time.Time
|
||
err := s.db.QueryRow(`
|
||
SELECT id, job_id, anchor_date, method, service_method, buss_id, biz_key,
|
||
payload_hash, payload_json, validation_errors, push_status, retry_count, last_error, updated_at
|
||
FROM hy_push_record WHERE id = ?`, recordID).Scan(
|
||
&r.ID, &r.JobID, &ad, &r.Method, &r.ServiceMethod, &r.BussID, &r.BizKey,
|
||
&r.PayloadHash, &r.PayloadJSON, &r.ValidationErrors, &r.PushStatus, &r.RetryCount, &r.LastError, &r.UpdatedAt,
|
||
)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
r.AnchorDate = ad.Format("2006-01-02")
|
||
|
||
logRows, err := s.db.Query(`
|
||
SELECT id, record_id, http_code, msg_code, msg, trace_id, request_body_len, request_plain_json, request_headers_json, response_body, duration_ms, created_at
|
||
FROM hy_push_log WHERE record_id = ? ORDER BY id DESC`, recordID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer logRows.Close()
|
||
|
||
var logs []PushLogRow
|
||
for logRows.Next() {
|
||
var l PushLogRow
|
||
if err := logRows.Scan(&l.ID, &l.RecordID, &l.HTTPCode, &l.MsgCode, &l.Msg, &l.TraceID,
|
||
&l.RequestBodyLen, &l.RequestPlainJSON, &l.RequestHeadersJSON, &l.ResponseBody, &l.DurationMs, &l.CreatedAt); err != nil {
|
||
return nil, err
|
||
}
|
||
logs = append(logs, l)
|
||
}
|
||
if len(logs) > 0 {
|
||
r.LastLog = &logs[0]
|
||
}
|
||
return &RecordDetail{Record: r, Logs: logs}, logRows.Err()
|
||
}
|
||
|
||
// Ping 检测数据库可用。
|
||
func (s *Store) Ping() error {
|
||
return s.db.Ping()
|
||
}
|
||
|
||
// StepForMethod 由 method 反查 step(用于 UI)。
|
||
func StepForMethod(method string) string {
|
||
m := map[string]string{
|
||
"uploadConsultIndicators": "consult",
|
||
"uploadReferralIndicators": "referral",
|
||
"uploadRecipeIndicators": "recipe",
|
||
"uploadRecipeVerificationIndicators": "verification",
|
||
}
|
||
if step, ok := m[method]; ok {
|
||
return step
|
||
}
|
||
return strings.TrimSpace(method)
|
||
}
|
||
|
||
// MethodForStep 由 step 查 method。
|
||
func MethodForStep(step string) (string, error) {
|
||
m := map[string]string{
|
||
"consult": "uploadConsultIndicators",
|
||
"referral": "uploadReferralIndicators",
|
||
"recipe": "uploadRecipeIndicators",
|
||
"verification": "uploadRecipeVerificationIndicators",
|
||
}
|
||
if method, ok := m[step]; ok {
|
||
return method, nil
|
||
}
|
||
return "", fmt.Errorf("unknown step: %s", step)
|
||
}
|