Files
lingniu-vehicle-ingest/vehicle-data-platform/apps/api/internal/businessscope/source.go

201 lines
7.1 KiB
Go

package businessscope
import (
"context"
"database/sql"
"fmt"
"strings"
"time"
)
// candidateQuery deliberately derives the active lease lifecycle from delivery and
// return task facts. vehicle_lease_order_record remains an independent customer
// ownership check, but its last_return_time is not authoritative: saving a return
// draft currently updates that aggregate field before the return is completed.
const candidateQuery = `SELECT
dv.vehicle_id,
v.id,
COALESCE(v.vin, ''),
COALESCE(v.plate_number, dv.plate_number, r.plate_number, ''),
r.customer_id,
ci.id,
dv.contract_id,
co.id,
COALESCE(co.other_customer_id, co.customer_id),
COALESCE(co.contract_code, r.contract_code, ''),
COALESCE(co.project_name, r.project_name, ''),
COALESCE(vs.operation_status, ''),
dv.delivery_time,
dv.update_time
FROM delivery_vehicle dv
LEFT JOIN vehicle_info v ON v.id = dv.vehicle_id AND v.del_flag = '0'
LEFT JOIN vehicle_lease_order_record r ON r.vehicle_id = dv.vehicle_id AND r.del_flag = '0'
LEFT JOIN customer_info ci ON ci.id = r.customer_id AND ci.del_flag = '0'
LEFT JOIN vehicle_lease_contract_info co ON co.id = dv.contract_id AND co.del_flag = '0'
LEFT JOIN vehicle_status vs ON vs.vehicle_id = dv.vehicle_id AND vs.del_flag = '0'
WHERE dv.del_flag = '0'
AND dv.delivery_status IN (2, 3)
AND dv.vehicle_id IS NOT NULL
AND dv.delivery_time IS NOT NULL
AND NOT EXISTS (
SELECT 1
FROM delivery_vehicle newer
WHERE newer.del_flag = '0'
AND newer.vehicle_id = dv.vehicle_id
AND newer.vehicle_id IS NOT NULL
AND (
COALESCE(newer.delivery_time, '1000-01-01') > COALESCE(dv.delivery_time, '1000-01-01')
OR (
COALESCE(newer.delivery_time, '1000-01-01') = COALESCE(dv.delivery_time, '1000-01-01')
AND newer.id > dv.id
)
)
)
AND NOT EXISTS (
SELECT 1
FROM return_vehicle_task rt
WHERE rt.delivery_vehicle_id = dv.id
AND rt.del_flag = '0'
AND rt.status IN (2, 3, 5)
)
ORDER BY dv.vehicle_id, dv.id, r.id`
func ReadCandidates(ctx context.Context, db *sql.DB) ([]Candidate, error) {
connection, err := db.Conn(ctx)
if err != nil {
return nil, fmt.Errorf("acquire OneOS connection: %w", err)
}
defer connection.Close()
if err := VerifyReadOnlyGrants(ctx, connection); err != nil {
return nil, err
}
if _, err := connection.ExecContext(ctx, `SET SESSION TRANSACTION READ ONLY`); err != nil {
return nil, fmt.Errorf("force OneOS session read only: %w", err)
}
if _, err := connection.ExecContext(ctx, `SET SESSION MAX_EXECUTION_TIME = 10000`); err != nil {
return nil, fmt.Errorf("set OneOS query timeout: %w", err)
}
tx, err := connection.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelRepeatableRead, ReadOnly: true})
if err != nil {
return nil, fmt.Errorf("begin OneOS read-only snapshot: %w", err)
}
defer tx.Rollback()
rows, err := tx.QueryContext(ctx, candidateQuery)
if err != nil {
return nil, fmt.Errorf("query OneOS customer vehicle candidates: %w", err)
}
defer rows.Close()
candidates := make([]Candidate, 0, 1024)
for rows.Next() {
var vehicleID, vehicleProfileID, customerID, customerProfileID sql.NullInt64
var contractID, contractProfileID, effectiveCustomerID sql.NullInt64
var vin, plate, contractCode, projectName, operationStatus string
var scopeStart time.Time
var updatedAt sql.NullTime
if err := rows.Scan(
&vehicleID, &vehicleProfileID, &vin, &plate, &customerID, &customerProfileID,
&contractID, &contractProfileID, &effectiveCustomerID, &contractCode, &projectName,
&operationStatus, &scopeStart, &updatedAt,
); err != nil {
return nil, fmt.Errorf("scan OneOS scope candidate: %w", err)
}
candidate := Candidate{
RowNumber: len(candidates) + 1,
VehicleID: vehicleID.Int64, VehiclePresent: vehicleID.Valid && vehicleProfileID.Valid,
VIN: vin, PlateNumber: plate,
CustomerID: customerID.Int64, CustomerPresent: customerID.Valid,
CustomerProfileExists: customerProfileID.Valid,
ContractID: contractID.Int64, ContractPresent: contractID.Valid,
ContractProfileExists: contractProfileID.Valid,
EffectiveCustomerID: effectiveCustomerID.Int64, EffectiveCustomerSet: effectiveCustomerID.Valid,
ContractCode: contractCode, ProjectName: projectName, OperationStatus: operationStatus,
ScopeStartAt: scopeStart,
}
if updatedAt.Valid {
value := updatedAt.Time
candidate.SourceUpdatedAt = &value
}
candidates = append(candidates, candidate)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate OneOS scope candidates: %w", err)
}
if err := tx.Commit(); err != nil {
return nil, fmt.Errorf("finish OneOS read-only snapshot: %w", err)
}
return candidates, nil
}
func VerifyReadOnlyGrants(ctx context.Context, connection *sql.Conn) error {
rows, err := connection.QueryContext(ctx, `SHOW GRANTS FOR CURRENT_USER`)
if err != nil {
return fmt.Errorf("inspect OneOS database grants: %w", err)
}
defer rows.Close()
grants := make([]string, 0)
for rows.Next() {
var grant string
if err := rows.Scan(&grant); err != nil {
return fmt.Errorf("scan OneOS database grant: %w", err)
}
grants = append(grants, grant)
}
if err := rows.Err(); err != nil {
return fmt.Errorf("iterate OneOS database grants: %w", err)
}
if err := ValidateReadOnlyGrants(grants); err != nil {
return fmt.Errorf("OneOS database account is not read-only: %w", err)
}
return nil
}
func ValidateReadOnlyGrants(grants []string) error {
if len(grants) == 0 {
return fmt.Errorf("SHOW GRANTS returned no rows")
}
allowedSelectScopes := map[string]bool{
"LN_ASSET_MANAGEMENT.VEHICLE_LEASE_ORDER_RECORD": true,
"LN_ASSET_MANAGEMENT.DELIVERY_VEHICLE": true,
"LN_ASSET_MANAGEMENT.RETURN_VEHICLE_TASK": true,
"LN_ASSET_MANAGEMENT.VEHICLE_INFO": true,
"LN_ASSET_MANAGEMENT.CUSTOMER_INFO": true,
"LN_ASSET_MANAGEMENT.VEHICLE_LEASE_CONTRACT_INFO": true,
"LN_ASSET_MANAGEMENT.VEHICLE_STATUS": true,
}
for _, raw := range grants {
grant := strings.ToUpper(strings.TrimSpace(raw))
if strings.HasPrefix(grant, "SET DEFAULT ROLE ") {
return fmt.Errorf("role-based grants are not accepted; grant SELECT directly to the sync account")
}
if !strings.HasPrefix(grant, "GRANT ") {
return fmt.Errorf("unsupported grant statement")
}
onIndex := strings.Index(grant, " ON ")
if onIndex < 0 {
return fmt.Errorf("unsupported role or dynamic grant")
}
privileges := strings.TrimSpace(strings.TrimPrefix(grant[:onIndex], "GRANT "))
toIndex := strings.Index(grant[onIndex+4:], " TO ")
if toIndex < 0 {
return fmt.Errorf("unsupported grant scope")
}
scope := strings.ReplaceAll(strings.TrimSpace(grant[onIndex+4:onIndex+4+toIndex]), "`", "")
for _, privilege := range strings.Split(privileges, ",") {
privilege = strings.TrimSpace(privilege)
switch privilege {
case "USAGE":
if scope != "*.*" {
return fmt.Errorf("USAGE has unsupported scope %q", scope)
}
case "SELECT":
if !allowedSelectScopes[scope] {
return fmt.Errorf("SELECT scope %q is not required by the sync query", scope)
}
default:
return fmt.Errorf("disallowed privilege %q", privilege)
}
}
}
return nil
}