feat(oneos): add signed scope API client

This commit is contained in:
lingniu
2026-07-16 18:49:01 +08:00
parent 1243efc7dd
commit 46f2026c03
10 changed files with 555 additions and 41 deletions

View File

@@ -23,10 +23,6 @@ func main() {
}
func run() error {
sourceDSN, err := normalizedDSN("ONEOS_MYSQL_DSN", os.Getenv("ONEOS_MYSQL_DSN"), true)
if err != nil {
return err
}
targetDSN, err := normalizedDSN("MYSQL_DSN", os.Getenv("MYSQL_DSN"), false)
if err != nil {
return err
@@ -34,21 +30,12 @@ func run() error {
timeout := time.Duration(envInt("ONEOS_SCOPE_SYNC_TIMEOUT_SEC", 60)) * time.Second
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
source, err := openDB(ctx, sourceDSN, 2)
if err != nil {
return fmt.Errorf("open OneOS read-only database: %w", err)
}
defer source.Close()
target, err := openDB(ctx, targetDSN, 4)
if err != nil {
return fmt.Errorf("open vehicle platform database: %w", err)
}
defer target.Close()
candidates, err := businessscope.ReadCandidates(ctx, source)
if err != nil {
return err
}
snapshot, err := businessscope.BuildSnapshot(candidates, time.Now())
snapshot, sourceKind, err := readSnapshot(ctx)
if err != nil {
return err
}
@@ -64,11 +51,50 @@ func run() error {
if err != nil {
return err
}
log.Printf("OneOS scope sync complete changed=%t candidates=%d accepted=%d rejected=%d version=%s",
result.Changed, snapshot.Candidates, len(snapshot.Items), len(snapshot.Rejections), snapshot.SourceVersion)
log.Printf("OneOS scope sync complete source=%s changed=%t candidates=%d accepted=%d rejected=%d version=%s",
sourceKind, result.Changed, snapshot.Candidates, len(snapshot.Items), len(snapshot.Rejections), snapshot.SourceVersion)
return nil
}
func readSnapshot(ctx context.Context) (businessscope.Snapshot, string, error) {
mode := strings.ToLower(strings.TrimSpace(os.Getenv("ONEOS_SCOPE_SOURCE")))
if mode == "" {
if strings.TrimSpace(os.Getenv("ONEOS_SCOPE_API_URL")) != "" {
mode = "api"
} else {
mode = "database"
}
}
switch mode {
case "api":
snapshot, err := businessscope.ReadAPISnapshot(ctx, businessscope.APIConfig{
URL: os.Getenv("ONEOS_SCOPE_API_URL"),
ServiceToken: os.Getenv("ONEOS_SCOPE_API_SERVICE_TOKEN"),
SigningSecret: os.Getenv("ONEOS_SCOPE_API_SIGNING_SECRET"),
MaxAttempts: envInt("ONEOS_SCOPE_API_MAX_ATTEMPTS", 3),
})
return snapshot, mode, err
case "database":
sourceDSN, err := normalizedDSN("ONEOS_MYSQL_DSN", os.Getenv("ONEOS_MYSQL_DSN"), true)
if err != nil {
return businessscope.Snapshot{}, mode, err
}
source, err := openDB(ctx, sourceDSN, 2)
if err != nil {
return businessscope.Snapshot{}, mode, fmt.Errorf("open OneOS read-only database: %w", err)
}
defer source.Close()
candidates, err := businessscope.ReadCandidates(ctx, source)
if err != nil {
return businessscope.Snapshot{}, mode, err
}
snapshot, err := businessscope.BuildSnapshot(candidates, time.Now())
return snapshot, mode, err
default:
return businessscope.Snapshot{}, mode, fmt.Errorf("ONEOS_SCOPE_SOURCE must be api or database")
}
}
func normalizedDSN(name, raw string, requireOneOSDatabase bool) (string, error) {
raw = strings.TrimSpace(raw)
if raw == "" {

View File

@@ -0,0 +1,343 @@
package businessscope
import (
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"io"
"net"
"net/http"
"net/url"
"strconv"
"strings"
"time"
)
const (
apiPageLimit = 500
apiMaxPages = 100
apiMaxRows = 50000
apiMaxBodyBytes = 8 << 20
apiSourcePrefix = "oneos-api-v1:"
)
type APIConfig struct {
URL string
ServiceToken string
SigningSecret string
MaxAttempts int
HTTPClient *http.Client
Now func() time.Time
}
type apiEnvelope struct {
Code int `json:"code"`
Message string `json:"message"`
Data apiSnapshotPage `json:"data"`
RequestID string `json:"requestId"`
}
type apiSnapshotPage struct {
ScopeVersion string `json:"scopeVersion"`
GeneratedAt string `json:"generatedAt"`
Complete bool `json:"complete"`
NextCursor string `json:"nextCursor"`
Items []apiScopeItem `json:"items"`
Rejected []apiRejected `json:"rejected"`
}
type apiScopeItem struct {
VehicleID json.RawMessage `json:"vehicleId"`
VIN string `json:"vin"`
PlateNumber string `json:"plateNumber"`
CustomerID json.RawMessage `json:"customerId"`
CustomerName string `json:"customerName"`
ContractID json.RawMessage `json:"contractId"`
ContractCode string `json:"contractCode"`
ProjectName string `json:"projectName"`
DepartmentID json.RawMessage `json:"departmentId"`
DepartmentName string `json:"departmentName"`
ResponsibleUserID json.RawMessage `json:"responsibleUserId"`
ResponsibleName string `json:"responsibleUserName"`
OperationStatus string `json:"operationStatus"`
ScopeStartAt string `json:"scopeStartAt"`
SourceUpdatedAt string `json:"sourceUpdatedAt"`
}
type apiRejected struct {
VehicleID json.RawMessage `json:"vehicleId"`
VIN string `json:"vin"`
CustomerID json.RawMessage `json:"customerId"`
ContractID json.RawMessage `json:"contractId"`
Reason string `json:"reason"`
ReasonCode string `json:"reasonCode"`
}
func ReadAPISnapshot(ctx context.Context, config APIConfig) (Snapshot, error) {
endpoint, err := url.Parse(strings.TrimSpace(config.URL))
if err != nil || endpoint.Scheme == "" || endpoint.Host == "" {
return Snapshot{}, fmt.Errorf("ONEOS_SCOPE_API_URL must be an absolute HTTP(S) URL")
}
if endpoint.Scheme != "https" && !isPrivateAPIHost(endpoint.Hostname()) {
return Snapshot{}, fmt.Errorf("OneOS scope API must use HTTPS outside private networks")
}
if strings.TrimSpace(config.ServiceToken) == "" || strings.TrimSpace(config.SigningSecret) == "" {
return Snapshot{}, fmt.Errorf("OneOS scope API service token and signing secret are required")
}
if config.MaxAttempts <= 0 {
config.MaxAttempts = 3
}
if config.MaxAttempts > 5 {
config.MaxAttempts = 5
}
if config.HTTPClient == nil {
config.HTTPClient = &http.Client{Timeout: 10 * time.Second}
}
if config.Now == nil {
config.Now = time.Now
}
candidates := make([]Candidate, 0, 1024)
upstreamRejections := make([]Rejection, 0)
cursor := ""
seenCursors := map[string]bool{}
scopeVersion := ""
var generatedAt time.Time
for pageNumber := 0; pageNumber < apiMaxPages; pageNumber++ {
pageURL := *endpoint
query := pageURL.Query()
query.Set("limit", strconv.Itoa(apiPageLimit))
if cursor != "" {
query.Set("cursor", cursor)
} else {
query.Del("cursor")
}
pageURL.RawQuery = query.Encode()
page, err := fetchAPIPage(ctx, config, &pageURL)
if err != nil {
return Snapshot{}, err
}
if !page.Complete {
return Snapshot{}, fmt.Errorf("OneOS scope API returned an incomplete snapshot")
}
if strings.TrimSpace(page.ScopeVersion) == "" {
return Snapshot{}, fmt.Errorf("OneOS scope API omitted scopeVersion")
}
pageGeneratedAt, err := time.Parse(time.RFC3339Nano, page.GeneratedAt)
if err != nil {
return Snapshot{}, fmt.Errorf("OneOS scope API generatedAt is invalid: %w", err)
}
if scopeVersion == "" {
scopeVersion = page.ScopeVersion
generatedAt = pageGeneratedAt
} else if page.ScopeVersion != scopeVersion || !pageGeneratedAt.Equal(generatedAt) {
return Snapshot{}, fmt.Errorf("OneOS scope API pagination changed snapshot version")
}
for _, item := range page.Items {
candidate, err := candidateFromAPI(item, len(candidates)+len(upstreamRejections)+1)
if err != nil {
return Snapshot{}, err
}
candidates = append(candidates, candidate)
}
for _, rejected := range page.Rejected {
upstreamRejections = append(upstreamRejections, rejectionFromAPI(rejected, len(candidates)+len(upstreamRejections)+1))
}
if len(candidates)+len(upstreamRejections) > apiMaxRows {
return Snapshot{}, fmt.Errorf("OneOS scope API exceeded %d rows", apiMaxRows)
}
next := strings.TrimSpace(page.NextCursor)
if next == "" {
break
}
if seenCursors[next] {
return Snapshot{}, fmt.Errorf("OneOS scope API repeated pagination cursor")
}
seenCursors[next] = true
cursor = next
if pageNumber == apiMaxPages-1 {
return Snapshot{}, fmt.Errorf("OneOS scope API exceeded %d pages", apiMaxPages)
}
}
snapshot, err := BuildSnapshot(candidates, generatedAt)
if err != nil {
return Snapshot{}, err
}
snapshot.Candidates += len(upstreamRejections)
snapshot.Rejections = append(snapshot.Rejections, upstreamRejections...)
versionHash := sha256.Sum256([]byte(scopeVersion + "\n" + snapshot.Checksum))
snapshot.SourceVersion = apiSourcePrefix + hex.EncodeToString(versionHash[:])
return snapshot, nil
}
func fetchAPIPage(ctx context.Context, config APIConfig, pageURL *url.URL) (apiSnapshotPage, error) {
var lastErr error
for attempt := 1; attempt <= config.MaxAttempts; attempt++ {
requestID, err := randomHex(16)
if err != nil {
return apiSnapshotPage{}, err
}
timestamp := strconv.FormatInt(config.Now().UTC().Unix(), 10)
canonical := http.MethodGet + "\n" + pageURL.RequestURI() + "\n" + timestamp + "\n" + requestID
mac := hmac.New(sha256.New, []byte(config.SigningSecret))
_, _ = mac.Write([]byte(canonical))
request, err := http.NewRequestWithContext(ctx, http.MethodGet, pageURL.String(), nil)
if err != nil {
return apiSnapshotPage{}, err
}
request.Header.Set("Accept", "application/json")
request.Header.Set("Authorization", "Service "+config.ServiceToken)
request.Header.Set("X-Request-Id", requestID)
request.Header.Set("X-Request-Timestamp", timestamp)
request.Header.Set("X-Request-Signature", hex.EncodeToString(mac.Sum(nil)))
response, err := config.HTTPClient.Do(request)
if err != nil {
lastErr = err
} else {
page, retry, responseErr := decodeAPIResponse(response)
if responseErr == nil {
return page, nil
}
lastErr = responseErr
if !retry {
return apiSnapshotPage{}, responseErr
}
}
if attempt < config.MaxAttempts {
select {
case <-ctx.Done():
return apiSnapshotPage{}, ctx.Err()
case <-time.After(time.Duration(attempt*100) * time.Millisecond):
}
}
}
return apiSnapshotPage{}, fmt.Errorf("OneOS scope API failed after %d attempts: %w", config.MaxAttempts, lastErr)
}
func decodeAPIResponse(response *http.Response) (apiSnapshotPage, bool, error) {
defer response.Body.Close()
retry := response.StatusCode == http.StatusTooManyRequests || response.StatusCode >= 500
if response.StatusCode != http.StatusOK {
_, _ = io.Copy(io.Discard, io.LimitReader(response.Body, 4096))
return apiSnapshotPage{}, retry, fmt.Errorf("OneOS scope API returned HTTP %d", response.StatusCode)
}
var envelope apiEnvelope
decoder := json.NewDecoder(io.LimitReader(response.Body, apiMaxBodyBytes+1))
if err := decoder.Decode(&envelope); err != nil {
return apiSnapshotPage{}, false, fmt.Errorf("decode OneOS scope API response: %w", err)
}
if envelope.Code != 0 {
return apiSnapshotPage{}, false, fmt.Errorf("OneOS scope API returned business code %d", envelope.Code)
}
return envelope.Data, false, nil
}
func candidateFromAPI(item apiScopeItem, row int) (Candidate, error) {
vehicleID, err := requiredInt64ID(item.VehicleID, "vehicleId")
if err != nil {
return Candidate{}, err
}
customerID, err := requiredInt64ID(item.CustomerID, "customerId")
if err != nil {
return Candidate{}, err
}
contractID, err := requiredInt64ID(item.ContractID, "contractId")
if err != nil {
return Candidate{}, err
}
scopeStart, err := time.Parse(time.RFC3339Nano, strings.TrimSpace(item.ScopeStartAt))
if err != nil {
return Candidate{}, fmt.Errorf("OneOS scope API item %d has invalid scopeStartAt", row)
}
var sourceUpdatedAt *time.Time
if value := strings.TrimSpace(item.SourceUpdatedAt); value != "" {
parsed, err := time.Parse(time.RFC3339Nano, value)
if err != nil {
return Candidate{}, fmt.Errorf("OneOS scope API item %d has invalid sourceUpdatedAt", row)
}
sourceUpdatedAt = &parsed
}
return Candidate{
RowNumber: row, VehicleID: vehicleID, VehiclePresent: true,
VIN: item.VIN, PlateNumber: item.PlateNumber,
CustomerID: customerID, CustomerName: item.CustomerName, CustomerPresent: true, CustomerProfileExists: true,
ContractID: contractID, ContractPresent: true, ContractProfileExists: true,
EffectiveCustomerID: customerID, EffectiveCustomerSet: true,
ContractCode: item.ContractCode, ProjectName: item.ProjectName,
DepartmentID: rawIDString(item.DepartmentID), DepartmentName: item.DepartmentName,
ResponsibleUserID: rawIDString(item.ResponsibleUserID), ResponsibleUserName: item.ResponsibleName,
OperationStatus: item.OperationStatus, ScopeStartAt: scopeStart, SourceUpdatedAt: sourceUpdatedAt,
}, nil
}
func rejectionFromAPI(item apiRejected, row int) Rejection {
reason := strings.ToUpper(strings.TrimSpace(firstNonEmptyScope(item.ReasonCode, item.Reason)))
if reason == "" || len(reason) > 64 {
reason = "UPSTREAM_REJECTED"
}
rejection := Rejection{RowNumber: row, VIN: normalizeVIN(item.VIN), ReasonCode: reason}
if value, err := optionalInt64ID(item.VehicleID); err == nil && value > 0 {
rejection.VehicleID = int64Pointer(value)
}
if value, err := optionalInt64ID(item.CustomerID); err == nil && value > 0 {
rejection.CustomerID = int64Pointer(value)
}
if value, err := optionalInt64ID(item.ContractID); err == nil && value > 0 {
rejection.ContractID = int64Pointer(value)
}
return rejection
}
func requiredInt64ID(raw json.RawMessage, name string) (int64, error) {
value, err := optionalInt64ID(raw)
if err != nil || value <= 0 {
return 0, fmt.Errorf("OneOS scope API %s must be a positive integer string", name)
}
return value, nil
}
func optionalInt64ID(raw json.RawMessage) (int64, error) {
value := rawIDString(raw)
if value == "" {
return 0, nil
}
return strconv.ParseInt(value, 10, 64)
}
func rawIDString(raw json.RawMessage) string {
value := strings.TrimSpace(string(raw))
if value == "" || value == "null" {
return ""
}
if strings.HasPrefix(value, `"`) {
var decoded string
if json.Unmarshal(raw, &decoded) == nil {
return strings.TrimSpace(decoded)
}
return ""
}
return value
}
func isPrivateAPIHost(host string) bool {
host = strings.ToLower(strings.TrimSpace(host))
if host == "localhost" {
return true
}
address := net.ParseIP(host)
return address != nil && (address.IsPrivate() || address.IsLoopback())
}
func firstNonEmptyScope(values ...string) string {
for _, value := range values {
if strings.TrimSpace(value) != "" {
return value
}
}
return ""
}

View File

@@ -0,0 +1,99 @@
package businessscope
import (
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/hex"
"fmt"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"time"
)
func TestReadAPISnapshotAuthenticatesPaginatesAndPreservesDimensions(t *testing.T) {
const token = "service-token"
const secret = "signing-secret"
now := time.Date(2026, 7, 16, 10, 0, 0, 0, time.UTC)
var requests atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
requests.Add(1)
if r.Header.Get("Authorization") != "Service "+token {
t.Fatalf("missing service token")
}
canonical := r.Method + "\n" + r.URL.RequestURI() + "\n" + r.Header.Get("X-Request-Timestamp") + "\n" + r.Header.Get("X-Request-Id")
mac := hmac.New(sha256.New, []byte(secret))
_, _ = mac.Write([]byte(canonical))
if r.Header.Get("X-Request-Signature") != hex.EncodeToString(mac.Sum(nil)) {
t.Fatalf("invalid request signature")
}
w.Header().Set("Content-Type", "application/json")
if r.URL.Query().Get("cursor") == "" {
fmt.Fprint(w, `{"code":0,"data":{"scopeVersion":"scope-42","generatedAt":"2026-07-16T10:00:00Z","complete":true,"nextCursor":"page-2","items":[{"vehicleId":"10","vin":" lvin0001 ","plateNumber":"沪A00001","customerId":"100","customerName":"客户甲","contractId":"1010","contractCode":"HT-1","projectName":"项目甲","departmentId":"20","departmentName":"运营一部","responsibleUserId":"30","responsibleUserName":"张三","operationStatus":"active","scopeStartAt":"2026-07-01T08:00:00+08:00","sourceUpdatedAt":"2026-07-16T09:59:00Z"}],"rejected":[]}}`)
return
}
fmt.Fprint(w, `{"code":0,"data":{"scopeVersion":"scope-42","generatedAt":"2026-07-16T10:00:00Z","complete":true,"nextCursor":"","items":[{"vehicleId":"11","vin":"LVIN0002","plateNumber":"沪A00002","customerId":"100","customerName":"客户甲","contractId":"1011","contractCode":"HT-2","projectName":"项目甲","departmentId":"20","departmentName":"运营一部","responsibleUserId":"31","responsibleUserName":"李四","operationStatus":"active","scopeStartAt":"2026-07-02T08:00:00+08:00"}],"rejected":[{"vehicleId":"12","vin":"","customerId":"100","contractId":"1012","reasonCode":"VIN_MISSING"}]}}`)
}))
defer server.Close()
snapshot, err := ReadAPISnapshot(context.Background(), APIConfig{
URL: server.URL, ServiceToken: token, SigningSecret: secret,
MaxAttempts: 1, HTTPClient: server.Client(), Now: func() time.Time { return now },
})
if err != nil {
t.Fatal(err)
}
if requests.Load() != 2 || snapshot.Candidates != 3 || len(snapshot.Items) != 2 || len(snapshot.Rejections) != 1 {
t.Fatalf("unexpected API snapshot: requests=%d snapshot=%+v", requests.Load(), snapshot)
}
if !strings.HasPrefix(snapshot.SourceVersion, apiSourcePrefix) || len(snapshot.SourceVersion) > 96 {
t.Fatalf("unexpected source version %q", snapshot.SourceVersion)
}
item := snapshot.Items[0]
if item.VIN != "LVIN0001" || item.CustomerName != "客户甲" || item.DepartmentName != "运营一部" ||
item.ResponsibleUserName != "张三" {
t.Fatalf("business dimensions lost: %+v", item)
}
}
func TestReadAPISnapshotFailsClosedOnVersionDrift(t *testing.T) {
var request int
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
request++
version := "scope-1"
cursor := "next"
if request == 2 {
version = "scope-2"
cursor = ""
}
fmt.Fprintf(w, `{"code":0,"data":{"scopeVersion":%q,"generatedAt":"2026-07-16T10:00:00Z","complete":true,"nextCursor":%q,"items":[],"rejected":[]}}`, version, cursor)
}))
defer server.Close()
_, err := ReadAPISnapshot(context.Background(), APIConfig{
URL: server.URL, ServiceToken: "token", SigningSecret: "secret", MaxAttempts: 1, HTTPClient: server.Client(),
})
if err == nil || !strings.Contains(err.Error(), "changed snapshot version") {
t.Fatalf("version drift must fail closed: %v", err)
}
}
func TestReadAPISnapshotRetriesOnlyTransientHTTPFailure(t *testing.T) {
var requests atomic.Int32
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if requests.Add(1) == 1 {
http.Error(w, "temporary", http.StatusServiceUnavailable)
return
}
fmt.Fprint(w, `{"code":0,"data":{"scopeVersion":"scope-1","generatedAt":"2026-07-16T10:00:00Z","complete":true,"nextCursor":"","items":[{"vehicleId":"10","vin":"LVIN0001","customerId":"100","contractId":"1010","scopeStartAt":"2026-07-01T00:00:00+08:00"}],"rejected":[]}}`)
}))
defer server.Close()
snapshot, err := ReadAPISnapshot(context.Background(), APIConfig{
URL: server.URL, ServiceToken: "token", SigningSecret: "secret", MaxAttempts: 2, HTTPClient: server.Client(),
})
if err != nil || requests.Load() != 2 || len(snapshot.Items) != 1 {
t.Fatalf("transient retry failed: requests=%d snapshot=%+v err=%v", requests.Load(), snapshot, err)
}
}

View File

@@ -33,6 +33,7 @@ type Candidate struct {
VIN string
PlateNumber string
CustomerID int64
CustomerName string
CustomerPresent bool
CustomerProfileExists bool
ContractID int64
@@ -42,22 +43,31 @@ type Candidate struct {
EffectiveCustomerSet bool
ContractCode string
ProjectName string
DepartmentID string
DepartmentName string
ResponsibleUserID string
ResponsibleUserName string
OperationStatus string
ScopeStartAt time.Time
SourceUpdatedAt *time.Time
}
type ScopeItem struct {
VehicleID int64
VIN string
PlateNumber string
CustomerID int64
ContractID int64
ContractCode string
ProjectName string
OperationStatus string
ScopeStartAt time.Time
SourceUpdatedAt *time.Time
VehicleID int64
VIN string
PlateNumber string
CustomerID int64
CustomerName string
ContractID int64
ContractCode string
ProjectName string
DepartmentID string
DepartmentName string
ResponsibleUserID string
ResponsibleUserName string
OperationStatus string
ScopeStartAt time.Time
SourceUpdatedAt *time.Time
}
type Rejection struct {
@@ -119,8 +129,10 @@ func BuildSnapshot(candidates []Candidate, generatedAt time.Time) (Snapshot, err
}
items = append(items, ScopeItem{
VehicleID: candidate.VehicleID, VIN: candidate.VIN, PlateNumber: strings.TrimSpace(candidate.PlateNumber),
CustomerID: candidate.CustomerID, ContractID: candidate.ContractID,
CustomerID: candidate.CustomerID, CustomerName: strings.TrimSpace(candidate.CustomerName), ContractID: candidate.ContractID,
ContractCode: strings.TrimSpace(candidate.ContractCode), ProjectName: strings.TrimSpace(candidate.ProjectName),
DepartmentID: strings.TrimSpace(candidate.DepartmentID), DepartmentName: strings.TrimSpace(candidate.DepartmentName),
ResponsibleUserID: strings.TrimSpace(candidate.ResponsibleUserID), ResponsibleUserName: strings.TrimSpace(candidate.ResponsibleUserName),
OperationStatus: strings.TrimSpace(candidate.OperationStatus), ScopeStartAt: candidate.ScopeStartAt,
SourceUpdatedAt: candidate.SourceUpdatedAt,
})

View File

@@ -61,8 +61,9 @@ run_id,source_system,source_version,status,candidate_count,accepted_count,reject
return PublishResult{}, fmt.Errorf("clear inactive business scope version: %w", err)
}
statement, err := tx.PrepareContext(ctx, `INSERT INTO business_customer_vehicle_scope(
source_version,customer_id,vin,vehicle_id,contract_id,contract_code,plate_number,project_name,operation_status,scope_start_at,source_updated_at,published_at
) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)`)
source_version,customer_id,customer_name,vin,vehicle_id,contract_id,contract_code,plate_number,project_name,
department_id,department_name,responsible_user_id,responsible_user_name,operation_status,scope_start_at,source_updated_at,published_at
) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`)
if err != nil {
return PublishResult{}, fmt.Errorf("prepare business scope insert: %w", err)
}
@@ -72,15 +73,16 @@ source_version,customer_id,vin,vehicle_id,contract_id,contract_code,plate_number
if item.ContractID > 0 {
contractID = item.ContractID
}
if _, err := statement.ExecContext(ctx, snapshot.SourceVersion, item.CustomerID, item.VIN, item.VehicleID, contractID,
item.ContractCode, item.PlateNumber, item.ProjectName, item.OperationStatus, item.ScopeStartAt, item.SourceUpdatedAt, now); err != nil {
if _, err := statement.ExecContext(ctx, snapshot.SourceVersion, item.CustomerID, item.CustomerName, item.VIN, item.VehicleID, contractID,
item.ContractCode, item.PlateNumber, item.ProjectName, item.DepartmentID, item.DepartmentName,
item.ResponsibleUserID, item.ResponsibleUserName, item.OperationStatus, item.ScopeStartAt, item.SourceUpdatedAt, now); err != nil {
return PublishResult{}, fmt.Errorf("insert business scope item: %w", err)
}
}
}
if len(snapshot.Rejections) > 0 {
statement, err := tx.PrepareContext(ctx, `INSERT INTO business_scope_rejection(
run_id,row_number,vehicle_id,vin,customer_id,contract_id,reason_code,created_at
run_id,`+"`row_number`"+`,vehicle_id,vin,customer_id,contract_id,reason_code,created_at
) VALUES(?,?,?,?,?,?,?,?)`)
if err != nil {
return PublishResult{}, fmt.Errorf("prepare scope rejection insert: %w", err)

View File

@@ -11,8 +11,9 @@ import (
)
const scopeInsertSQL = `INSERT INTO business_customer_vehicle_scope(
source_version,customer_id,vin,vehicle_id,contract_id,contract_code,plate_number,project_name,operation_status,scope_start_at,source_updated_at,published_at
) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)`
source_version,customer_id,customer_name,vin,vehicle_id,contract_id,contract_code,plate_number,project_name,
department_id,department_name,responsible_user_id,responsible_user_name,operation_status,scope_start_at,source_updated_at,published_at
) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`
func publishSnapshot(version string) Snapshot {
generatedAt := time.Date(2026, 7, 14, 8, 0, 0, 0, time.UTC)
@@ -21,8 +22,9 @@ func publishSnapshot(version string) Snapshot {
Checksum: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
GeneratedAt: generatedAt, Candidates: 1,
Items: []ScopeItem{{
VehicleID: 10, VIN: "LVIN0001", PlateNumber: "沪A00001", CustomerID: 100,
VehicleID: 10, VIN: "LVIN0001", PlateNumber: "沪A00001", CustomerID: 100, CustomerName: "示例客户",
ContractID: 1010, ContractCode: "HT-1", ProjectName: "项目", OperationStatus: "1",
DepartmentID: "20", DepartmentName: "运营一部", ResponsibleUserID: "30", ResponsibleUserName: "负责人",
ScopeStartAt: generatedAt.Add(-time.Hour),
}},
}
@@ -47,8 +49,9 @@ func TestPublishRebuildsPreviouslyStoredVersionBeforeActivation(t *testing.T) {
prepared := mock.ExpectPrepare(regexp.QuoteMeta(scopeInsertSQL))
item := snapshot.Items[0]
prepared.ExpectExec().WithArgs(
snapshot.SourceVersion, item.CustomerID, item.VIN, item.VehicleID, item.ContractID,
item.ContractCode, item.PlateNumber, item.ProjectName, item.OperationStatus, item.ScopeStartAt, nil, sqlmock.AnyArg(),
snapshot.SourceVersion, item.CustomerID, item.CustomerName, item.VIN, item.VehicleID, item.ContractID,
item.ContractCode, item.PlateNumber, item.ProjectName, item.DepartmentID, item.DepartmentName,
item.ResponsibleUserID, item.ResponsibleUserName, item.OperationStatus, item.ScopeStartAt, nil, sqlmock.AnyArg(),
).WillReturnResult(sqlmock.NewResult(1, 1))
mock.ExpectExec(`UPDATE business_scope_state SET`).
WithArgs(snapshot.SourceVersion, snapshot.Checksum, 1, 1, 0, snapshot.GeneratedAt, sqlmock.AnyArg(), sqlmock.AnyArg()).