Files
xk-hy-transit-go/internal/db/query.go
2026-05-22 09:17:39 +08:00

279 lines
8.3 KiB
Go
Raw Permalink 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 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)
}