@@ -0,0 +1,24 @@
|
||||
package router
|
||||
|
||||
import (
|
||||
"github.com/gin-gonic/gin"
|
||||
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/provisioning"
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
|
||||
)
|
||||
|
||||
func init() { routerCheckRole = append(routerCheckRole, registerSenseProvisioningRouter) }
|
||||
|
||||
func registerSenseProvisioningRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
|
||||
api := &provisioning.API{}
|
||||
r := v1.Group("/provisioning/batches").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
|
||||
r.GET("", api.List)
|
||||
r.POST("", api.Create)
|
||||
r.GET("/:id", api.Get)
|
||||
r.POST("/:id/execute", api.Execute)
|
||||
r.POST("/:id/retry-failed", api.RetryFailed)
|
||||
r.POST("/:id/items/:itemId/retry", api.RetryItem)
|
||||
r.GET("/:id/export", api.Export)
|
||||
}
|
||||
@@ -0,0 +1,170 @@
|
||||
package provisioning
|
||||
|
||||
import (
|
||||
"encoding/csv"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/go-admin-team/go-admin-core/sdk/api"
|
||||
"github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth/user"
|
||||
)
|
||||
|
||||
type API struct{ api.Api }
|
||||
|
||||
func (e *API) service(c *gin.Context) (*Service, error) {
|
||||
service := &Service{}
|
||||
if err := e.MakeContext(c).MakeOrm().MakeService(&service.Service).Errors; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return service, nil
|
||||
}
|
||||
|
||||
func (e *API) List(c *gin.Context) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
request := BatchPageRequest{}
|
||||
if err = e.MakeContext(c).Bind(&request).Errors; err != nil {
|
||||
e.Error(http.StatusBadRequest, err, "查询条件格式不正确")
|
||||
return
|
||||
}
|
||||
list, count, err := service.ListBatches(&request)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
e.PageOK(list, int(count), request.GetPageIndex(), request.GetPageSize(), "查询成功")
|
||||
}
|
||||
|
||||
func (e *API) Get(c *gin.Context) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
response, err := service.GetBatch(c.Param("id"))
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
e.OK(response, "查询成功")
|
||||
}
|
||||
|
||||
func (e *API) Create(c *gin.Context) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
request := CreateBatchRequest{CreateBy: user.GetUserId(c)}
|
||||
if err = bindStrictJSON(c, &request); err != nil {
|
||||
e.Error(http.StatusBadRequest, err, "导入内容格式不正确")
|
||||
return
|
||||
}
|
||||
response, err := service.CreateBatch(request)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
e.OK(response, "批次已导入并完成预校验")
|
||||
}
|
||||
|
||||
func (e *API) Execute(c *gin.Context) { e.execute(c, false, "") }
|
||||
func (e *API) RetryFailed(c *gin.Context) { e.execute(c, true, "") }
|
||||
func (e *API) RetryItem(c *gin.Context) { e.execute(c, true, c.Param("itemId")) }
|
||||
|
||||
func (e *API) execute(c *gin.Context, retryFailed bool, itemID string) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
request := ExecuteRequest{UpdateBy: user.GetUserId(c)}
|
||||
defer func() { clearCredentials(request.Credentials) }()
|
||||
if err = bindStrictJSON(c, &request); err != nil {
|
||||
e.Error(http.StatusBadRequest, err, "执行内容格式不正确")
|
||||
return
|
||||
}
|
||||
response, err := service.Execute(c.Request.Context(), c.Param("id"), itemID, retryFailed, request)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
e.OK(response, "批量开通处理完成")
|
||||
}
|
||||
|
||||
func (e *API) Export(c *gin.Context) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
batch, err := service.GetBatch(c.Param("id"))
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
c.Header("Content-Type", "text/csv; charset=utf-8")
|
||||
c.Header("Content-Disposition", "attachment; filename=provisioning-"+batch.ID+".csv")
|
||||
c.Status(http.StatusOK)
|
||||
writer := csv.NewWriter(c.Writer)
|
||||
if err = writer.Write([]string{"line_number", "name", "location", "address", "status", "failure_code", "detail", "device_id"}); err != nil {
|
||||
e.Logger.Error(err)
|
||||
return
|
||||
}
|
||||
for _, item := range batch.Items {
|
||||
if err = writer.Write([]string{strconv.Itoa(item.LineNumber), item.Name, item.Location, item.Address, item.Status, item.FailureCode, item.Detail, item.DeviceID}); err != nil {
|
||||
e.Logger.Error(err)
|
||||
return
|
||||
}
|
||||
}
|
||||
writer.Flush()
|
||||
if err = writer.Error(); err != nil {
|
||||
e.Logger.Error(err)
|
||||
}
|
||||
}
|
||||
|
||||
func (e *API) writeError(err error) {
|
||||
switch {
|
||||
case errors.Is(err, ErrInvalidRequest):
|
||||
e.Error(http.StatusBadRequest, err, err.Error())
|
||||
case errors.Is(err, ErrBatchNotFound), errors.Is(err, ErrItemNotFound):
|
||||
e.Error(http.StatusNotFound, err, err.Error())
|
||||
default:
|
||||
e.Error(http.StatusInternalServerError, err, "批量开通操作失败")
|
||||
}
|
||||
}
|
||||
|
||||
func bindStrictJSON(c *gin.Context, target any) error {
|
||||
if !strings.HasPrefix(strings.ToLower(strings.TrimSpace(c.GetHeader("Content-Type"))), "application/json") {
|
||||
return errors.New("content type must be application/json")
|
||||
}
|
||||
decoder := json.NewDecoder(http.MaxBytesReader(c.Writer, c.Request.Body, 2<<20))
|
||||
decoder.DisallowUnknownFields()
|
||||
if err := decoder.Decode(target); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := decoder.Decode(&struct{}{}); !errors.Is(err, io.EOF) {
|
||||
if err == nil {
|
||||
return errors.New("request body must contain one JSON object")
|
||||
}
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func clearCredentials(values []CredentialInput) {
|
||||
for index := range values {
|
||||
values[index].ONVIFUsername = ""
|
||||
values[index].ONVIFPassword = ""
|
||||
values[index].RTSPUsername = ""
|
||||
values[index].RTSPPassword = ""
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,32 @@
|
||||
package provisioning
|
||||
|
||||
import (
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
func TestCreateRequestRejectsCredentialColumns(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
recorder := httptest.NewRecorder()
|
||||
context, _ := gin.CreateTestContext(recorder)
|
||||
context.Request = httptest.NewRequest("POST", "/api/v1/provisioning/batches", strings.NewReader(`{
|
||||
"idempotencyKey":"import-1",
|
||||
"rows":[{"lineNumber":1,"name":"东门摄像机","location":"东门","address":"http://192.0.2.10/onvif","password":"must-not-be-accepted"}]
|
||||
}`))
|
||||
context.Request.Header.Set("Content-Type", "application/json")
|
||||
var request CreateBatchRequest
|
||||
if err := bindStrictJSON(context, &request); err == nil || !strings.Contains(err.Error(), "unknown field") {
|
||||
t.Fatalf("expected unknown credential field to be rejected, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestClearCredentialsOverwritesTransientValues(t *testing.T) {
|
||||
values := []CredentialInput{{ONVIFUsername: "installer", ONVIFPassword: "temporary-secret", RTSPUsername: "stream", RTSPPassword: "stream-secret"}}
|
||||
clearCredentials(values)
|
||||
if values[0].ONVIFUsername != "" || values[0].ONVIFPassword != "" || values[0].RTSPUsername != "" || values[0].RTSPPassword != "" {
|
||||
t.Fatalf("credentials were not cleared: %#v", values[0])
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
package provisioning
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
commonDTO "git.ilapage.cn/ila/yovision/Sense/server/common/dto"
|
||||
)
|
||||
|
||||
type BatchPageRequest struct {
|
||||
commonDTO.Pagination `search:"-"`
|
||||
Status string `form:"status"`
|
||||
}
|
||||
|
||||
type ImportRow struct {
|
||||
LineNumber int `json:"lineNumber"`
|
||||
Name string `json:"name"`
|
||||
Location string `json:"location"`
|
||||
Address string `json:"address"`
|
||||
}
|
||||
|
||||
type CreateBatchRequest struct {
|
||||
IdempotencyKey string `json:"idempotencyKey"`
|
||||
Rows []ImportRow `json:"rows"`
|
||||
CreateBy int `json:"-"`
|
||||
}
|
||||
|
||||
type CredentialInput struct {
|
||||
ItemID string `json:"itemId"`
|
||||
ONVIFUsername string `json:"onvifUsername"`
|
||||
ONVIFPassword string `json:"onvifPassword"`
|
||||
RTSPSameAsONVIF bool `json:"rtspSameAsOnvif"`
|
||||
RTSPUsername string `json:"rtspUsername"`
|
||||
RTSPPassword string `json:"rtspPassword"`
|
||||
}
|
||||
|
||||
type ExecuteRequest struct {
|
||||
Credentials []CredentialInput `json:"credentials"`
|
||||
UpdateBy int `json:"-"`
|
||||
}
|
||||
|
||||
type ItemResponse struct {
|
||||
ID string `json:"id"`
|
||||
LineNumber int `json:"lineNumber"`
|
||||
Name string `json:"name"`
|
||||
Location string `json:"location"`
|
||||
Address string `json:"address"`
|
||||
Status string `json:"status"`
|
||||
FailureCode string `json:"failureCode,omitempty"`
|
||||
Detail string `json:"detail,omitempty"`
|
||||
DeviceID string `json:"deviceId,omitempty"`
|
||||
Attempts int `json:"attempts"`
|
||||
LastTriedAt *time.Time `json:"lastTriedAt,omitempty"`
|
||||
}
|
||||
|
||||
type BatchResponse struct {
|
||||
ID string `json:"id"`
|
||||
IdempotencyKey string `json:"idempotencyKey"`
|
||||
Status string `json:"status"`
|
||||
QuotaLimit int `json:"quotaLimit"`
|
||||
ExistingCount int `json:"existingCount"`
|
||||
TotalCount int `json:"totalCount"`
|
||||
ReadyCount int `json:"readyCount"`
|
||||
SuccessCount int `json:"successCount"`
|
||||
FailureCount int `json:"failureCount"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
Items []ItemResponse `json:"items,omitempty"`
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
package provisioning
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
|
||||
)
|
||||
|
||||
const (
|
||||
BatchReady = "ready"
|
||||
BatchRunning = "running"
|
||||
BatchSucceeded = "succeeded"
|
||||
BatchPartial = "partial"
|
||||
BatchFailed = "failed"
|
||||
|
||||
ItemInvalid = "invalid"
|
||||
ItemReady = "ready"
|
||||
ItemQuotaExceeded = "quota_exceeded"
|
||||
ItemRunning = "running"
|
||||
ItemSucceeded = "succeeded"
|
||||
ItemFailed = "failed"
|
||||
)
|
||||
|
||||
type Batch struct {
|
||||
ID string `gorm:"size:36;primaryKey" json:"id"`
|
||||
IdempotencyKey string `gorm:"size:128;not null;uniqueIndex" json:"idempotencyKey"`
|
||||
Status string `gorm:"size:32;not null;index" json:"status"`
|
||||
QuotaLimit int `gorm:"not null" json:"quotaLimit"`
|
||||
ExistingCount int `gorm:"not null" json:"existingCount"`
|
||||
TotalCount int `gorm:"not null" json:"totalCount"`
|
||||
ReadyCount int `gorm:"not null" json:"readyCount"`
|
||||
SuccessCount int `gorm:"not null" json:"successCount"`
|
||||
FailureCount int `gorm:"not null" json:"failureCount"`
|
||||
common.ControlBy
|
||||
common.ModelTime
|
||||
Items []Item `gorm:"foreignKey:BatchID" json:"items,omitempty"`
|
||||
}
|
||||
|
||||
func (Batch) TableName() string { return "sense_provisioning_batches" }
|
||||
|
||||
type Item struct {
|
||||
ID string `gorm:"size:36;primaryKey" json:"id"`
|
||||
BatchID string `gorm:"size:36;not null;uniqueIndex:batch_line;index" json:"batchId"`
|
||||
LineNumber int `gorm:"not null;uniqueIndex:batch_line" json:"lineNumber"`
|
||||
Name string `gorm:"size:128;not null" json:"name"`
|
||||
Location string `gorm:"size:255;not null;default:''" json:"location"`
|
||||
Address string `gorm:"size:1024;not null" json:"address"`
|
||||
Status string `gorm:"size:32;not null;index" json:"status"`
|
||||
FailureCode string `gorm:"size:64;not null;default:''" json:"failureCode,omitempty"`
|
||||
Detail string `gorm:"size:512;not null;default:''" json:"detail,omitempty"`
|
||||
DeviceID string `gorm:"size:36;not null;default:'';index" json:"deviceId,omitempty"`
|
||||
Attempts int `gorm:"not null;default:0" json:"attempts"`
|
||||
LastTriedAt *time.Time `json:"lastTriedAt,omitempty"`
|
||||
common.ControlBy
|
||||
common.ModelTime
|
||||
}
|
||||
|
||||
func (Item) TableName() string { return "sense_provisioning_items" }
|
||||
@@ -0,0 +1,366 @@
|
||||
package provisioning
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/url"
|
||||
"os"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgconn"
|
||||
"gorm.io/gorm"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/admission"
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/credential"
|
||||
deviceService "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service"
|
||||
deviceDTO "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/service/dto"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrInvalidRequest = errors.New("批量开通请求不符合要求")
|
||||
ErrBatchNotFound = errors.New("批量开通批次不存在")
|
||||
ErrItemNotFound = errors.New("批量开通条目不存在")
|
||||
)
|
||||
|
||||
type Activator func(context.Context, *Service, *Item, CredentialInput, int) (string, string, string, error)
|
||||
|
||||
type Service struct {
|
||||
coreService.Service
|
||||
Quota int
|
||||
Activator Activator
|
||||
}
|
||||
|
||||
func QuotaFromEnvironment() int {
|
||||
value := strings.TrimSpace(os.Getenv("SENSE_PROVISIONING_QUOTA"))
|
||||
if value == "" {
|
||||
return 16
|
||||
}
|
||||
parsed, err := strconv.Atoi(value)
|
||||
if err != nil || parsed < 1 || parsed > 100000 {
|
||||
return 16
|
||||
}
|
||||
return parsed
|
||||
}
|
||||
|
||||
func (s *Service) quota() int {
|
||||
if s.Quota > 0 {
|
||||
return s.Quota
|
||||
}
|
||||
return QuotaFromEnvironment()
|
||||
}
|
||||
|
||||
func (s *Service) CreateBatch(request CreateBatchRequest) (BatchResponse, error) {
|
||||
request.IdempotencyKey = strings.TrimSpace(request.IdempotencyKey)
|
||||
if request.IdempotencyKey == "" || len(request.IdempotencyKey) > 128 || len(request.Rows) == 0 || len(request.Rows) > 1000 {
|
||||
return BatchResponse{}, ErrInvalidRequest
|
||||
}
|
||||
var existing Batch
|
||||
if err := s.Orm.Preload("Items", func(db *gorm.DB) *gorm.DB { return db.Order("line_number") }).First(&existing, "idempotency_key = ?", request.IdempotencyKey).Error; err == nil {
|
||||
return batchResponse(existing), nil
|
||||
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return BatchResponse{}, fmt.Errorf("read provisioning idempotency key: %w", err)
|
||||
}
|
||||
|
||||
var existingDevices int64
|
||||
if err := s.Orm.Table("sense_devices").Where("status <> ?", "disabled").Count(&existingDevices).Error; err != nil {
|
||||
return BatchResponse{}, fmt.Errorf("count provisioned devices: %w", err)
|
||||
}
|
||||
quota := s.quota()
|
||||
available := quota - int(existingDevices)
|
||||
if available < 0 {
|
||||
available = 0
|
||||
}
|
||||
batch := Batch{ID: uuid.NewString(), IdempotencyKey: request.IdempotencyKey, Status: BatchReady, QuotaLimit: quota, ExistingCount: int(existingDevices), TotalCount: len(request.Rows)}
|
||||
batch.CreateBy, batch.UpdateBy = request.CreateBy, request.CreateBy
|
||||
seenLines := map[int]bool{}
|
||||
seenAddresses := map[string]bool{}
|
||||
ready := 0
|
||||
for _, row := range request.Rows {
|
||||
item := Item{ID: uuid.NewString(), BatchID: batch.ID, LineNumber: row.LineNumber, Name: strings.TrimSpace(row.Name), Location: strings.TrimSpace(row.Location), Address: strings.TrimSpace(row.Address), Status: ItemReady}
|
||||
item.CreateBy, item.UpdateBy = request.CreateBy, request.CreateBy
|
||||
code, detail := validateImportRow(item, seenLines, seenAddresses)
|
||||
if code != "" {
|
||||
item.Status, item.FailureCode, item.Detail = ItemInvalid, code, detail
|
||||
} else if ready >= available {
|
||||
item.Status, item.FailureCode, item.Detail = ItemQuotaExceeded, "quota_exceeded", "超出当前可用配额,请调整配额或批次后重试"
|
||||
} else {
|
||||
ready++
|
||||
}
|
||||
batch.Items = append(batch.Items, item)
|
||||
}
|
||||
batch.ReadyCount = ready
|
||||
if ready == 0 {
|
||||
batch.Status = BatchFailed
|
||||
}
|
||||
if err := s.Orm.Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Omit("Items").Create(&batch).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(&batch.Items).Error
|
||||
}); err != nil {
|
||||
if isDuplicateKey(err) {
|
||||
return s.GetBatchByKey(request.IdempotencyKey)
|
||||
}
|
||||
return BatchResponse{}, fmt.Errorf("create provisioning batch: %w", err)
|
||||
}
|
||||
return s.GetBatch(batch.ID)
|
||||
}
|
||||
|
||||
func validateImportRow(item Item, seenLines map[int]bool, seenAddresses map[string]bool) (string, string) {
|
||||
if item.LineNumber < 1 || seenLines[item.LineNumber] {
|
||||
return "invalid_line_number", "行号必须为正整数且批次内唯一"
|
||||
}
|
||||
seenLines[item.LineNumber] = true
|
||||
if item.Name == "" || len([]rune(item.Name)) > 128 || len([]rune(item.Location)) > 255 {
|
||||
return "invalid_device", "设备名称不能为空,名称或安装位置长度不能超过限制"
|
||||
}
|
||||
parsed, err := url.Parse(item.Address)
|
||||
if err != nil || (parsed.Scheme != "http" && parsed.Scheme != "https") || parsed.Hostname() == "" || parsed.User != nil || parsed.RawQuery != "" || parsed.Fragment != "" {
|
||||
return "invalid_address", "设备地址必须是无账号、查询参数和片段的 HTTP(S) 地址"
|
||||
}
|
||||
key := strings.ToLower(parsed.String())
|
||||
if seenAddresses[key] {
|
||||
return "duplicate_address", "同一批次中设备地址不能重复"
|
||||
}
|
||||
seenAddresses[key] = true
|
||||
return "", ""
|
||||
}
|
||||
|
||||
func (s *Service) ListBatches(request *BatchPageRequest) ([]BatchResponse, int64, error) {
|
||||
query := s.Orm.Model(&Batch{})
|
||||
if request.Status != "" {
|
||||
query = query.Where("status = ?", request.Status)
|
||||
}
|
||||
var count int64
|
||||
if err := query.Count(&count).Error; err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
pageSize := request.GetPageSize()
|
||||
if pageSize > 100 {
|
||||
pageSize = 100
|
||||
}
|
||||
var batches []Batch
|
||||
if err := query.Order("created_at DESC").Limit(pageSize).Offset((request.GetPageIndex() - 1) * pageSize).Find(&batches).Error; err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
result := make([]BatchResponse, 0, len(batches))
|
||||
for _, batch := range batches {
|
||||
result = append(result, batchResponse(batch))
|
||||
}
|
||||
return result, count, nil
|
||||
}
|
||||
|
||||
func (s *Service) GetBatch(id string) (BatchResponse, error) {
|
||||
var batch Batch
|
||||
if err := s.Orm.Preload("Items", func(db *gorm.DB) *gorm.DB { return db.Order("line_number") }).First(&batch, "id = ?", id).Error; err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return BatchResponse{}, ErrBatchNotFound
|
||||
}
|
||||
return BatchResponse{}, err
|
||||
}
|
||||
return batchResponse(batch), nil
|
||||
}
|
||||
|
||||
func (s *Service) GetBatchByKey(key string) (BatchResponse, error) {
|
||||
var batch Batch
|
||||
if err := s.Orm.Preload("Items", func(db *gorm.DB) *gorm.DB { return db.Order("line_number") }).First(&batch, "idempotency_key = ?", key).Error; err != nil {
|
||||
return BatchResponse{}, err
|
||||
}
|
||||
return batchResponse(batch), nil
|
||||
}
|
||||
|
||||
func (s *Service) Execute(ctx context.Context, batchID string, itemID string, retryFailed bool, request ExecuteRequest) (BatchResponse, error) {
|
||||
credentialByItem := make(map[string]CredentialInput, len(request.Credentials))
|
||||
for _, input := range request.Credentials {
|
||||
if input.ItemID == "" || credentialByItem[input.ItemID].ItemID != "" {
|
||||
return BatchResponse{}, ErrInvalidRequest
|
||||
}
|
||||
credentialByItem[input.ItemID] = input
|
||||
}
|
||||
var batch Batch
|
||||
if err := s.Orm.First(&batch, "id = ?", batchID).Error; err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return BatchResponse{}, ErrBatchNotFound
|
||||
}
|
||||
return BatchResponse{}, err
|
||||
}
|
||||
query := s.Orm.Where("batch_id = ?", batchID)
|
||||
if itemID != "" {
|
||||
query = query.Where("id = ?", itemID)
|
||||
}
|
||||
if retryFailed {
|
||||
query = query.Where("status = ?", ItemFailed)
|
||||
} else {
|
||||
query = query.Where("status = ?", ItemReady)
|
||||
}
|
||||
var items []Item
|
||||
if err := query.Order("line_number").Find(&items).Error; err != nil {
|
||||
return BatchResponse{}, err
|
||||
}
|
||||
if itemID != "" && len(items) == 0 {
|
||||
return BatchResponse{}, ErrItemNotFound
|
||||
}
|
||||
if len(items) == 0 {
|
||||
return s.GetBatch(batchID)
|
||||
}
|
||||
if err := s.Orm.Model(&Batch{}).Where("id = ?", batchID).Updates(map[string]any{"status": BatchRunning, "update_by": request.UpdateBy}).Error; err != nil {
|
||||
return BatchResponse{}, fmt.Errorf("mark provisioning batch running: %w", err)
|
||||
}
|
||||
activate := s.Activator
|
||||
if activate == nil {
|
||||
activate = defaultActivate
|
||||
}
|
||||
for index := range items {
|
||||
item := &items[index]
|
||||
now := time.Now().UTC()
|
||||
expectedStatus := ItemReady
|
||||
if retryFailed {
|
||||
expectedStatus = ItemFailed
|
||||
}
|
||||
claim := s.Orm.Model(&Item{}).Where("id = ? AND status = ?", item.ID, expectedStatus).Updates(map[string]any{"status": ItemRunning, "attempts": gorm.Expr("attempts + 1"), "last_tried_at": now, "update_by": request.UpdateBy})
|
||||
if claim.Error != nil {
|
||||
return BatchResponse{}, fmt.Errorf("mark provisioning item running: %w", claim.Error)
|
||||
}
|
||||
if claim.RowsAffected == 0 {
|
||||
continue
|
||||
}
|
||||
input, ok := credentialByItem[item.ID]
|
||||
if !ok || strings.TrimSpace(input.ONVIFUsername) == "" || input.ONVIFPassword == "" {
|
||||
if err := s.failItem(item.ID, request.UpdateBy, "credentials_required", "请安全填写该设备的账号密码后重试"); err != nil {
|
||||
return BatchResponse{}, err
|
||||
}
|
||||
continue
|
||||
}
|
||||
deviceID, status, detail, err := activate(ctx, s, item, input, request.UpdateBy)
|
||||
if deviceID != "" && deviceID != item.DeviceID {
|
||||
item.DeviceID = deviceID
|
||||
if err := s.Orm.Model(&Item{}).Where("id = ?", item.ID).Update("device_id", deviceID).Error; err != nil {
|
||||
return BatchResponse{}, fmt.Errorf("link provisioning item to device: %w", err)
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
code, safeDetail := safeFailure(err)
|
||||
if updateErr := s.failItem(item.ID, request.UpdateBy, code, safeDetail); updateErr != nil {
|
||||
return BatchResponse{}, updateErr
|
||||
}
|
||||
continue
|
||||
}
|
||||
if status != "ready" {
|
||||
if detail == "" {
|
||||
detail = "设备或视频流尚未通过验证"
|
||||
}
|
||||
if err := s.failItem(item.ID, request.UpdateBy, status, detail); err != nil {
|
||||
return BatchResponse{}, err
|
||||
}
|
||||
continue
|
||||
}
|
||||
if err := s.Orm.Model(&Item{}).Where("id = ?", item.ID).Updates(map[string]any{"status": ItemSucceeded, "failure_code": "", "detail": "设备已开通并通过视频验证", "update_by": request.UpdateBy}).Error; err != nil {
|
||||
return BatchResponse{}, fmt.Errorf("mark provisioning item succeeded: %w", err)
|
||||
}
|
||||
}
|
||||
if err := s.refreshBatch(batchID, request.UpdateBy); err != nil {
|
||||
return BatchResponse{}, err
|
||||
}
|
||||
return s.GetBatch(batchID)
|
||||
}
|
||||
|
||||
func isDuplicateKey(err error) bool {
|
||||
if errors.Is(err, gorm.ErrDuplicatedKey) {
|
||||
return true
|
||||
}
|
||||
var postgresError *pgconn.PgError
|
||||
return errors.As(err, &postgresError) && postgresError.Code == "23505"
|
||||
}
|
||||
|
||||
func (s *Service) failItem(id string, updateBy int, code, detail string) error {
|
||||
if err := s.Orm.Model(&Item{}).Where("id = ?", id).Updates(map[string]any{"status": ItemFailed, "failure_code": code, "detail": detail, "update_by": updateBy}).Error; err != nil {
|
||||
return fmt.Errorf("mark provisioning item failed: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Service) refreshBatch(batchID string, updateBy int) error {
|
||||
var items []Item
|
||||
if err := s.Orm.Where("batch_id = ?", batchID).Find(&items).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
ready, success, failure := 0, 0, 0
|
||||
for _, item := range items {
|
||||
switch item.Status {
|
||||
case ItemReady, ItemRunning:
|
||||
ready++
|
||||
case ItemSucceeded:
|
||||
success++
|
||||
case ItemFailed:
|
||||
failure++
|
||||
}
|
||||
}
|
||||
status := BatchFailed
|
||||
if success == len(items) {
|
||||
status = BatchSucceeded
|
||||
} else if success > 0 {
|
||||
status = BatchPartial
|
||||
} else if ready > 0 {
|
||||
status = BatchReady
|
||||
}
|
||||
return s.Orm.Model(&Batch{}).Where("id = ?", batchID).Updates(map[string]any{"status": status, "ready_count": ready, "success_count": success, "failure_count": failure, "update_by": updateBy}).Error
|
||||
}
|
||||
|
||||
func defaultActivate(ctx context.Context, service *Service, item *Item, input CredentialInput, updateBy int) (string, string, string, error) {
|
||||
device := deviceService.Device{Service: service.Service}
|
||||
deviceID := item.DeviceID
|
||||
var response deviceDTO.DeviceResponse
|
||||
if deviceID == "" {
|
||||
if err := device.Insert(&deviceDTO.CreateReq{Name: item.Name, Location: item.Location, Modality: "video", Capabilities: []string{"video"}, CreateBy: updateBy}, &response); err != nil {
|
||||
return "", "", "", err
|
||||
}
|
||||
deviceID = response.ID
|
||||
if err := service.Orm.Model(&Item{}).Where("id = ?", item.ID).Update("device_id", deviceID).Error; err != nil {
|
||||
return deviceID, "", "", err
|
||||
}
|
||||
} else if err := device.Get(deviceID, &response); err != nil {
|
||||
return deviceID, "", "", err
|
||||
}
|
||||
if err := device.UpdateCredentials(&deviceDTO.CredentialUpdateReq{ID: deviceID, ONVIFUsername: input.ONVIFUsername, ONVIFPassword: input.ONVIFPassword, RTSPSameAsONVIF: input.RTSPSameAsONVIF, RTSPUsername: input.RTSPUsername, RTSPPassword: input.RTSPPassword, Version: response.Version, UpdateBy: updateBy}, &response); err != nil {
|
||||
return deviceID, "", "", err
|
||||
}
|
||||
probeService, err := admission.NewRuntime(service.Service)
|
||||
if err != nil {
|
||||
return deviceID, "", "", err
|
||||
}
|
||||
result, err := probeService.Probe(ctx, admission.ProbeRequest{DeviceID: deviceID, Address: item.Address, Version: response.Version, UpdateBy: updateBy})
|
||||
if err != nil {
|
||||
return deviceID, "", "", err
|
||||
}
|
||||
return deviceID, result.Status, result.Detail, nil
|
||||
}
|
||||
|
||||
func safeFailure(err error) (string, string) {
|
||||
switch {
|
||||
case errors.Is(err, credential.ErrKeyUnavailable):
|
||||
return "credential_key_unavailable", "摄像头凭据安全配置不可用"
|
||||
case errors.Is(err, credential.ErrCredentialNotConfigured):
|
||||
return "credentials_required", "请安全填写摄像头账号密码后重试"
|
||||
case errors.Is(err, deviceService.ErrInvalidDevice), errors.Is(err, admission.ErrInvalid):
|
||||
return "invalid_device", "设备信息或凭据不符合要求"
|
||||
case errors.Is(err, deviceService.ErrVersionConflict), errors.Is(err, admission.ErrConflict):
|
||||
return "version_conflict", "设备已被其他操作更新,请重试"
|
||||
default:
|
||||
return "activation_failed", "设备开通失败,请检查网络、地址和凭据后重试"
|
||||
}
|
||||
}
|
||||
|
||||
func batchResponse(batch Batch) BatchResponse {
|
||||
response := BatchResponse{ID: batch.ID, IdempotencyKey: batch.IdempotencyKey, Status: batch.Status, QuotaLimit: batch.QuotaLimit, ExistingCount: batch.ExistingCount, TotalCount: batch.TotalCount, ReadyCount: batch.ReadyCount, SuccessCount: batch.SuccessCount, FailureCount: batch.FailureCount, CreatedAt: batch.CreatedAt, UpdatedAt: batch.UpdatedAt, Items: make([]ItemResponse, 0, len(batch.Items))}
|
||||
sort.Slice(batch.Items, func(i, j int) bool { return batch.Items[i].LineNumber < batch.Items[j].LineNumber })
|
||||
for _, item := range batch.Items {
|
||||
response.Items = append(response.Items, ItemResponse{ID: item.ID, LineNumber: item.LineNumber, Name: item.Name, Location: item.Location, Address: item.Address, Status: item.Status, FailureCode: item.FailureCode, Detail: item.Detail, DeviceID: item.DeviceID, Attempts: item.Attempts, LastTriedAt: item.LastTriedAt})
|
||||
}
|
||||
return response
|
||||
}
|
||||
@@ -0,0 +1,278 @@
|
||||
package provisioning
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgconn"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
|
||||
deviceModels "git.ilapage.cn/ila/yovision/Sense/server/app/sense/device/models"
|
||||
)
|
||||
|
||||
func provisioningService(t *testing.T, quota int) *Service {
|
||||
t.Helper()
|
||||
db, err := gorm.Open(sqlite.Open("file:"+uuid.NewString()+"?mode=memory&cache=shared"), &gorm.Config{TranslateError: true})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
if err = db.AutoMigrate(&deviceModels.Device{}, &Batch{}, &Item{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return &Service{Service: coreService.Service{Orm: db}, Quota: quota}
|
||||
}
|
||||
|
||||
func TestCreateBatchValidatesQuotaAndIsIdempotent(t *testing.T) {
|
||||
service := provisioningService(t, 2)
|
||||
if err := service.Orm.Create(&deviceModels.Device{ID: "existing", Name: "已接入摄像机", Modality: "video", Status: "active", AdapterStatus: "ready", Version: 1}).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
request := CreateBatchRequest{
|
||||
IdempotencyKey: "import-001",
|
||||
CreateBy: 7,
|
||||
Rows: []ImportRow{
|
||||
{LineNumber: 1, Name: "东门摄像机", Location: "东门", Address: "http://192.0.2.10/onvif"},
|
||||
{LineNumber: 2, Name: "重复地址", Location: "东门", Address: "http://192.0.2.10/onvif"},
|
||||
{LineNumber: 3, Name: "西门摄像机", Location: "西门", Address: "https://192.0.2.11/onvif"},
|
||||
},
|
||||
}
|
||||
created, err := service.CreateBatch(request)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if created.QuotaLimit != 2 || created.ExistingCount != 1 || created.ReadyCount != 1 || len(created.Items) != 3 {
|
||||
t.Fatalf("unexpected batch: %#v", created)
|
||||
}
|
||||
if created.Items[0].Status != ItemReady || created.Items[1].Status != ItemInvalid || created.Items[2].Status != ItemQuotaExceeded {
|
||||
t.Fatalf("unexpected item statuses: %#v", created.Items)
|
||||
}
|
||||
|
||||
request.Rows = []ImportRow{{LineNumber: 1, Name: "不应覆盖", Address: "http://192.0.2.99/onvif"}}
|
||||
repeated, err := service.CreateBatch(request)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if repeated.ID != created.ID || repeated.TotalCount != 3 || repeated.Items[0].Name != "东门摄像机" {
|
||||
t.Fatalf("idempotent replay changed the batch: %#v", repeated)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExecuteSupportsPartialSuccessAndFailedOnlyRetry(t *testing.T) {
|
||||
service := provisioningService(t, 16)
|
||||
created, err := service.CreateBatch(CreateBatchRequest{
|
||||
IdempotencyKey: "execute-001",
|
||||
Rows: []ImportRow{
|
||||
{LineNumber: 1, Name: "东门摄像机", Address: "http://192.0.2.10/onvif"},
|
||||
{LineNumber: 2, Name: "西门摄像机", Address: "http://192.0.2.11/onvif"},
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
calls := map[int]int{}
|
||||
service.Activator = func(_ context.Context, _ *Service, item *Item, input CredentialInput, _ int) (string, string, string, error) {
|
||||
calls[item.LineNumber]++
|
||||
if input.ONVIFPassword != "temporary-secret" {
|
||||
t.Fatalf("activator did not receive the transient credential")
|
||||
}
|
||||
if item.LineNumber == 2 && calls[item.LineNumber] == 1 {
|
||||
return "device-2", "", "", errors.New("synthetic network failure containing temporary-secret")
|
||||
}
|
||||
return "device-" + string(rune('0'+item.LineNumber)), "ready", "验证通过", nil
|
||||
}
|
||||
credentials := make([]CredentialInput, 0, len(created.Items))
|
||||
for _, item := range created.Items {
|
||||
credentials = append(credentials, CredentialInput{ItemID: item.ID, ONVIFUsername: "installer", ONVIFPassword: "temporary-secret", RTSPSameAsONVIF: true})
|
||||
}
|
||||
partial, err := service.Execute(context.Background(), created.ID, "", false, ExecuteRequest{Credentials: credentials, UpdateBy: 8})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if partial.Status != BatchPartial || partial.SuccessCount != 1 || partial.FailureCount != 1 {
|
||||
t.Fatalf("expected partial success, got %#v", partial)
|
||||
}
|
||||
if partial.Items[1].Detail == "" || strings.Contains(partial.Items[1].Detail, "temporary-secret") {
|
||||
t.Fatalf("unsafe failure detail: %q", partial.Items[1].Detail)
|
||||
}
|
||||
|
||||
retried, err := service.Execute(context.Background(), created.ID, "", true, ExecuteRequest{Credentials: []CredentialInput{{ItemID: created.Items[1].ID, ONVIFUsername: "installer", ONVIFPassword: "temporary-secret", RTSPSameAsONVIF: true}}, UpdateBy: 8})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if retried.Status != BatchSucceeded || retried.SuccessCount != 2 || calls[1] != 1 || calls[2] != 2 {
|
||||
t.Fatalf("failed-only retry was not idempotent: response=%#v calls=%#v", retried, calls)
|
||||
}
|
||||
encoded, err := json.Marshal(retried)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if strings.Contains(string(encoded), "temporary-secret") || strings.Contains(string(encoded), "installer") {
|
||||
t.Fatalf("response contains credentials: %s", encoded)
|
||||
}
|
||||
var stored []Item
|
||||
if err = service.Orm.Find(&stored).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
storedJSON, _ := json.Marshal(stored)
|
||||
if strings.Contains(string(storedJSON), "temporary-secret") || strings.Contains(string(storedJSON), "installer") {
|
||||
t.Fatalf("provisioning rows contain credentials: %s", storedJSON)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExecuteRequiresCredentialsWithoutCallingActivator(t *testing.T) {
|
||||
service := provisioningService(t, 16)
|
||||
created, err := service.CreateBatch(CreateBatchRequest{IdempotencyKey: "missing-credentials", Rows: []ImportRow{{LineNumber: 1, Name: "东门摄像机", Address: "http://192.0.2.10/onvif"}}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
called := false
|
||||
service.Activator = func(context.Context, *Service, *Item, CredentialInput, int) (string, string, string, error) {
|
||||
called = true
|
||||
return "", "", "", nil
|
||||
}
|
||||
result, err := service.Execute(context.Background(), created.ID, "", false, ExecuteRequest{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if called || result.Status != BatchFailed || result.Items[0].FailureCode != "credentials_required" {
|
||||
t.Fatalf("missing credentials were not handled safely: %#v called=%v", result, called)
|
||||
}
|
||||
}
|
||||
|
||||
func TestQuotaFromEnvironment(t *testing.T) {
|
||||
t.Setenv("SENSE_PROVISIONING_QUOTA", "24")
|
||||
if got := QuotaFromEnvironment(); got != 24 {
|
||||
t.Fatalf("quota=%d", got)
|
||||
}
|
||||
t.Setenv("SENSE_PROVISIONING_QUOTA", "invalid")
|
||||
if got := QuotaFromEnvironment(); got != 16 {
|
||||
t.Fatalf("fallback quota=%d", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestConcurrentCreateUsesOneIdempotentBatch(t *testing.T) {
|
||||
service := provisioningService(t, 16)
|
||||
sqlDB, err := service.Orm.DB()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
request := CreateBatchRequest{IdempotencyKey: "concurrent-import", Rows: []ImportRow{{LineNumber: 1, Name: "东门摄像机", Address: "http://192.0.2.10/onvif"}}}
|
||||
const workers = 8
|
||||
ids := make(chan string, workers)
|
||||
errs := make(chan error, workers)
|
||||
var wait sync.WaitGroup
|
||||
for index := 0; index < workers; index++ {
|
||||
wait.Add(1)
|
||||
go func() {
|
||||
defer wait.Done()
|
||||
batch, createErr := service.CreateBatch(request)
|
||||
if createErr != nil {
|
||||
errs <- createErr
|
||||
return
|
||||
}
|
||||
ids <- batch.ID
|
||||
}()
|
||||
}
|
||||
wait.Wait()
|
||||
close(ids)
|
||||
close(errs)
|
||||
for createErr := range errs {
|
||||
t.Fatalf("concurrent create failed: %v", createErr)
|
||||
}
|
||||
var first string
|
||||
for id := range ids {
|
||||
if first == "" {
|
||||
first = id
|
||||
} else if id != first {
|
||||
t.Fatalf("idempotent creates returned different batches: %q and %q", first, id)
|
||||
}
|
||||
}
|
||||
var count int64
|
||||
if err = service.Orm.Model(&Batch{}).Count(&count).Error; err != nil || count != 1 {
|
||||
t.Fatalf("batch count=%d err=%v", count, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSixteenDeviceBaselineSmoke(t *testing.T) {
|
||||
service := provisioningService(t, 16)
|
||||
rows := make([]ImportRow, 0, 16)
|
||||
for index := 1; index <= 16; index++ {
|
||||
rows = append(rows, ImportRow{LineNumber: index, Name: fmt.Sprintf("摄像机-%02d", index), Address: fmt.Sprintf("http://192.0.2.%d/onvif", index)})
|
||||
}
|
||||
created, err := service.CreateBatch(CreateBatchRequest{IdempotencyKey: "sixteen-device-smoke", Rows: rows})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if created.ReadyCount != 16 {
|
||||
t.Fatalf("ready count=%d", created.ReadyCount)
|
||||
}
|
||||
service.Activator = func(_ context.Context, _ *Service, item *Item, _ CredentialInput, _ int) (string, string, string, error) {
|
||||
return "device-" + item.ID, "ready", "验证通过", nil
|
||||
}
|
||||
credentials := make([]CredentialInput, 0, 16)
|
||||
for _, item := range created.Items {
|
||||
credentials = append(credentials, CredentialInput{ItemID: item.ID, ONVIFUsername: "installer", ONVIFPassword: "temporary-secret", RTSPSameAsONVIF: true})
|
||||
}
|
||||
completed, err := service.Execute(context.Background(), created.ID, "", false, ExecuteRequest{Credentials: credentials})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if completed.Status != BatchSucceeded || completed.SuccessCount != 16 || completed.FailureCount != 0 {
|
||||
t.Fatalf("unexpected 16-device result: %#v", completed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestConcurrentExecuteClaimsAnItemOnce(t *testing.T) {
|
||||
service := provisioningService(t, 16)
|
||||
created, err := service.CreateBatch(CreateBatchRequest{IdempotencyKey: "concurrent-execute", Rows: []ImportRow{{LineNumber: 1, Name: "东门摄像机", Address: "http://192.0.2.10/onvif"}}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var calls atomic.Int32
|
||||
service.Activator = func(_ context.Context, _ *Service, item *Item, _ CredentialInput, _ int) (string, string, string, error) {
|
||||
calls.Add(1)
|
||||
return "device-" + item.ID, "ready", "验证通过", nil
|
||||
}
|
||||
request := ExecuteRequest{Credentials: []CredentialInput{{ItemID: created.Items[0].ID, ONVIFUsername: "installer", ONVIFPassword: "temporary-secret", RTSPSameAsONVIF: true}}}
|
||||
var wait sync.WaitGroup
|
||||
errs := make(chan error, 2)
|
||||
for index := 0; index < 2; index++ {
|
||||
wait.Add(1)
|
||||
go func() {
|
||||
defer wait.Done()
|
||||
_, executeErr := service.Execute(context.Background(), created.ID, "", false, request)
|
||||
errs <- executeErr
|
||||
}()
|
||||
}
|
||||
wait.Wait()
|
||||
close(errs)
|
||||
for executeErr := range errs {
|
||||
if executeErr != nil {
|
||||
t.Fatal(executeErr)
|
||||
}
|
||||
}
|
||||
if calls.Load() != 1 {
|
||||
t.Fatalf("activator calls=%d", calls.Load())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPostgresDuplicateKeyDetection(t *testing.T) {
|
||||
if !isDuplicateKey(&pgconn.PgError{Code: "23505"}) || isDuplicateKey(&pgconn.PgError{Code: "23503"}) {
|
||||
t.Fatal("PostgreSQL duplicate-key detection is incorrect")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
package version
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"runtime"
|
||||
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/provisioning"
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration"
|
||||
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
|
||||
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
|
||||
)
|
||||
|
||||
func init() {
|
||||
_, fileName, _, _ := runtime.Caller(0)
|
||||
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSenseProvisioning)
|
||||
}
|
||||
|
||||
func migrateSenseProvisioning(db *gorm.DB, version string) error {
|
||||
return db.Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.AutoMigrate(&provisioning.Batch{}, &provisioning.Item{}, &deviceCasbinRule{}); err != nil {
|
||||
return err
|
||||
}
|
||||
root, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: senseLayoutMenuName, Title: "视频感知", Icon: "video-camera", Path: "/sense", MenuType: "M", Component: "Layout", Sort: 5, Visible: "0", IsFrame: "1"})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
page, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseProvisioning", Title: "批量开通", Icon: "upload", Path: "provisioning", Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:provisioning:list", ParentId: root.MenuId, Component: "/sense/provisioning/index", Sort: 2, Visible: "0", IsFrame: "1"})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
definitions := []struct{ name, title, action, permission string }{
|
||||
{"SenseProvisioningImport", "导入批次", "POST", "sense:provisioning:import"},
|
||||
{"SenseProvisioningExecute", "执行开通", "POST", "sense:provisioning:execute"},
|
||||
{"SenseProvisioningRetry", "重试失败项", "POST", "sense:provisioning:retry"},
|
||||
{"SenseProvisioningExport", "导出结果", "GET", "sense:provisioning:export"},
|
||||
}
|
||||
buttons := make([]migrationModels.SysMenu, 0, len(definitions))
|
||||
for index, definition := range definitions {
|
||||
button, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: definition.name, Title: definition.title, MenuType: "F", Action: definition.action, Permission: definition.permission, ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: index + 1, Visible: "1", IsFrame: "1"})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
buttons = append(buttons, button)
|
||||
}
|
||||
allMenus := append([]migrationModels.SysMenu{page}, buttons...)
|
||||
for _, role := range []string{"implementation_operator", "site_admin"} {
|
||||
if err = attachDeviceRole(tx, role, allMenus); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if err = attachDeviceRole(tx, "viewer", []migrationModels.SysMenu{page, buttons[3]}); err != nil {
|
||||
return err
|
||||
}
|
||||
read := [][2]string{{"/api/v1/provisioning/batches", "GET"}, {"/api/v1/provisioning/batches/:id", "GET"}, {"/api/v1/provisioning/batches/:id/export", "GET"}}
|
||||
write := [][2]string{{"/api/v1/provisioning/batches", "POST"}, {"/api/v1/provisioning/batches/:id/execute", "POST"}, {"/api/v1/provisioning/batches/:id/retry-failed", "POST"}, {"/api/v1/provisioning/batches/:id/items/:itemId/retry", "POST"}}
|
||||
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
|
||||
policies := append([][2]string{}, read...)
|
||||
if role != "viewer" {
|
||||
policies = append(policies, write...)
|
||||
}
|
||||
for _, policy := range policies {
|
||||
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: policy[0], V2: policy[1]}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
if err = rebuildSenseMenuPaths(tx, root.MenuId, "/0"); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(&common.Migration{Version: version}).Error
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,61 @@
|
||||
package version
|
||||
|
||||
import (
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"gorm.io/driver/postgres"
|
||||
"gorm.io/gorm"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/provisioning"
|
||||
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
|
||||
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
|
||||
)
|
||||
|
||||
func TestProvisioningMigrationOnPostgres(t *testing.T) {
|
||||
dsn := os.Getenv("SENSE_PROVISIONING_MIGRATION_TEST_DATABASE_URL")
|
||||
if dsn == "" {
|
||||
t.Skip("set SENSE_PROVISIONING_MIGRATION_TEST_DATABASE_URL to run the PostgreSQL migration test")
|
||||
}
|
||||
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
const schema = "sense_provisioning_72_test"
|
||||
if err = db.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err = db.Exec("CREATE SCHEMA " + schema).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = db.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error })
|
||||
sqlDB, err := db.DB()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
if err = db.Exec("SET search_path TO " + schema).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
|
||||
if err = db.Create(&migrationModels.SysRole{RoleName: role, RoleKey: role, Status: "2"}).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
const version = "2026082809000_provisioning.go"
|
||||
if err = migrateSenseProvisioning(db, version); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var menus, policies, batches, items, applied int64
|
||||
db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseProvisioning%").Count(&menus)
|
||||
db.Model(&deviceCasbinRule{}).Where("v1 LIKE ?", "/api/v1/provisioning/%").Count(&policies)
|
||||
db.Model(&provisioning.Batch{}).Count(&batches)
|
||||
db.Model(&provisioning.Item{}).Count(&items)
|
||||
db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied)
|
||||
if menus != 5 || policies != 17 || batches != 0 || items != 0 || applied != 1 {
|
||||
t.Fatalf("menus=%d policies=%d batches=%d items=%d applied=%d", menus, policies, batches, items, applied)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,36 @@
|
||||
import request from '@/utils/request'
|
||||
import axios from 'axios'
|
||||
import { getToken } from '@/utils/auth'
|
||||
|
||||
export function listProvisioningBatches(query) {
|
||||
return request({ url: '/api/v1/provisioning/batches', method: 'get', params: query })
|
||||
}
|
||||
|
||||
export function getProvisioningBatch(id) {
|
||||
return request({ url: `/api/v1/provisioning/batches/${id}`, method: 'get' })
|
||||
}
|
||||
|
||||
export function createProvisioningBatch(data) {
|
||||
return request({ url: '/api/v1/provisioning/batches', method: 'post', data })
|
||||
}
|
||||
|
||||
export function executeProvisioningBatch(id, data) {
|
||||
return request({ url: `/api/v1/provisioning/batches/${id}/execute`, method: 'post', data })
|
||||
}
|
||||
|
||||
export function retryProvisioningFailures(id, data) {
|
||||
return request({ url: `/api/v1/provisioning/batches/${id}/retry-failed`, method: 'post', data })
|
||||
}
|
||||
|
||||
export function retryProvisioningItem(batchId, itemId, data) {
|
||||
return request({ url: `/api/v1/provisioning/batches/${batchId}/items/${itemId}/retry`, method: 'post', data })
|
||||
}
|
||||
|
||||
export async function downloadProvisioningBatch(id) {
|
||||
const baseURL = String(process.env.VUE_APP_BASE_API || '').replace(/\/$/, '')
|
||||
const response = await axios.get(`${baseURL}/api/v1/provisioning/batches/${id}/export`, {
|
||||
responseType: 'blob',
|
||||
headers: { Authorization: `Bearer ${getToken()}` }
|
||||
})
|
||||
return response.data
|
||||
}
|
||||
@@ -0,0 +1,200 @@
|
||||
<template>
|
||||
<BasicLayout>
|
||||
<template #wrapper>
|
||||
<el-card class="box-card">
|
||||
<div class="page-header">
|
||||
<div>
|
||||
<h3>摄像头批量开通</h3>
|
||||
<p>按模板导入、检查配额、临时填写凭据并执行开通;密码不会保存在批次记录中。</p>
|
||||
</div>
|
||||
<el-button :icon="Refresh" @click="loadBatches">刷新记录</el-button>
|
||||
</div>
|
||||
|
||||
<el-steps :active="activeStep" finish-status="success" simple class="workflow-steps">
|
||||
<el-step title="导入清单" />
|
||||
<el-step title="预校验" />
|
||||
<el-step title="填写凭据" />
|
||||
<el-step title="查看结果" />
|
||||
</el-steps>
|
||||
|
||||
<el-alert title="默认交付配额为 16 路,实际可开通数量以服务端预校验结果为准。模板不得包含账号或密码列。" type="info" :closable="false" show-icon />
|
||||
|
||||
<div class="toolbar">
|
||||
<el-upload v-permisaction="['sense:provisioning:import']" :auto-upload="false" :show-file-list="false" accept=".xlsx,.xls,.csv" :on-change="handleFile">
|
||||
<el-button type="primary" :icon="Upload">选择清单</el-button>
|
||||
</el-upload>
|
||||
<el-button :icon="Download" @click="downloadTemplate">下载模板</el-button>
|
||||
<span class="file-name">{{ fileName || '尚未选择文件' }}</span>
|
||||
</div>
|
||||
|
||||
<el-table v-if="importRows.length" :data="importRows" border stripe max-height="360">
|
||||
<el-table-column prop="lineNumber" label="行号" width="80" />
|
||||
<el-table-column prop="name" label="设备名称" min-width="150" />
|
||||
<el-table-column prop="location" label="安装位置" min-width="160" />
|
||||
<el-table-column prop="address" label="ONVIF 地址" min-width="250" show-overflow-tooltip />
|
||||
</el-table>
|
||||
<div v-if="importRows.length" class="import-actions">
|
||||
<el-button v-permisaction="['sense:provisioning:import']" type="primary" :loading="importing" @click="submitImport">提交预校验</el-button>
|
||||
<span>共 {{ importRows.length }} 条;提交后服务端会检查格式、重复地址和剩余配额。</span>
|
||||
</div>
|
||||
|
||||
<el-divider />
|
||||
<div class="section-title">
|
||||
<h4>开通批次</h4>
|
||||
<el-select v-model="statusFilter" placeholder="全部状态" clearable @change="loadBatches">
|
||||
<el-option v-for="item in statusOptions" :key="item.value" :label="item.label" :value="item.value" />
|
||||
</el-select>
|
||||
</div>
|
||||
<el-table v-loading="loading" :data="batches" border stripe>
|
||||
<el-table-column prop="createdAt" label="创建时间" width="180">
|
||||
<template #default="scope">{{ parseTime(scope.row.createdAt) }}</template>
|
||||
</el-table-column>
|
||||
<el-table-column label="进度" min-width="190">
|
||||
<template #default="scope">总数 {{ scope.row.totalCount }} / 可执行 {{ scope.row.readyCount }} / 成功 {{ scope.row.successCount }} / 失败 {{ scope.row.failureCount }}</template>
|
||||
</el-table-column>
|
||||
<el-table-column label="配额" width="140">
|
||||
<template #default="scope">已用 {{ scope.row.existingCount }} / {{ scope.row.quotaLimit }}</template>
|
||||
</el-table-column>
|
||||
<el-table-column label="状态" width="110" align="center">
|
||||
<template #default="scope"><el-tag :type="statusType(scope.row.status)">{{ statusLabel(scope.row.status) }}</el-tag></template>
|
||||
</el-table-column>
|
||||
<el-table-column label="操作" width="230" fixed="right">
|
||||
<template #default="scope">
|
||||
<el-button type="primary" link @click="openBatch(scope.row.id)">查看</el-button>
|
||||
<el-button v-if="scope.row.readyCount" v-permisaction="['sense:provisioning:execute']" type="primary" link @click="openCredentials(scope.row.id, false)">填写凭据</el-button>
|
||||
<el-button v-if="scope.row.failureCount" v-permisaction="['sense:provisioning:retry']" type="warning" link @click="openCredentials(scope.row.id, true)">重试失败</el-button>
|
||||
</template>
|
||||
</el-table-column>
|
||||
</el-table>
|
||||
<pagination v-show="total > 0" v-model:current-page="query.pageIndex" v-model:page-size="query.pageSize" :total="total" @pagination="loadBatches" />
|
||||
</el-card>
|
||||
|
||||
<el-dialog v-model="detailOpen" title="批次详情" width="min(980px, calc(100vw - 24px))" :close-on-click-modal="false">
|
||||
<div v-if="selectedBatch" class="summary-cards" role="status" aria-live="polite">
|
||||
<el-tag>{{ statusLabel(selectedBatch.status) }}</el-tag>
|
||||
<span>配额 {{ selectedBatch.existingCount }} / {{ selectedBatch.quotaLimit }}</span>
|
||||
<span>成功 {{ selectedBatch.successCount }}</span>
|
||||
<span>失败 {{ selectedBatch.failureCount }}</span>
|
||||
</div>
|
||||
<el-table v-if="selectedBatch" :data="selectedBatch.items" border stripe max-height="470">
|
||||
<el-table-column prop="lineNumber" label="行号" width="70" />
|
||||
<el-table-column prop="name" label="设备名称" min-width="130" />
|
||||
<el-table-column prop="location" label="安装位置" min-width="130" />
|
||||
<el-table-column prop="address" label="地址" min-width="210" show-overflow-tooltip />
|
||||
<el-table-column label="状态" width="110"><template #default="scope"><el-tag :type="statusType(scope.row.status)">{{ statusLabel(scope.row.status) }}</el-tag></template></el-table-column>
|
||||
<el-table-column prop="detail" label="说明" min-width="200" show-overflow-tooltip />
|
||||
<el-table-column label="操作" width="90" fixed="right"><template #default="scope"><el-button v-if="scope.row.status === 'failed'" v-permisaction="['sense:provisioning:retry']" type="warning" link @click="openSingleRetry(scope.row)">重试</el-button></template></el-table-column>
|
||||
</el-table>
|
||||
<template #footer>
|
||||
<el-button v-permisaction="['sense:provisioning:export']" :icon="Download" @click="exportSelected">导出结果</el-button>
|
||||
<el-button @click="detailOpen = false">关闭</el-button>
|
||||
</template>
|
||||
</el-dialog>
|
||||
|
||||
<el-dialog v-model="credentialOpen" :title="singleItemId ? '重试单个失败项' : retryMode ? '重试失败项' : '填写临时凭据并开通'" width="min(1180px, calc(100vw - 24px))" :close-on-click-modal="false" @closed="clearCredentials">
|
||||
<el-alert title="凭据仅用于本次请求,提交成功或关闭窗口后立即从页面内存清除。" type="warning" :closable="false" show-icon />
|
||||
<el-table :data="credentialItems" border stripe max-height="500" class="credential-table">
|
||||
<el-table-column prop="name" label="设备" min-width="130" />
|
||||
<el-table-column prop="address" label="地址" min-width="190" show-overflow-tooltip />
|
||||
<el-table-column label="ONVIF 用户名" min-width="145"><template #default="scope"><el-input v-model="credentials[scope.row.id].onvifUsername" autocomplete="off" /></template></el-table-column>
|
||||
<el-table-column label="ONVIF 密码" min-width="160"><template #default="scope"><el-input v-model="credentials[scope.row.id].onvifPassword" type="password" show-password autocomplete="new-password" /></template></el-table-column>
|
||||
<el-table-column label="RTSP" width="120"><template #default="scope"><el-checkbox v-model="credentials[scope.row.id].rtspSameAsOnvif">同上</el-checkbox></template></el-table-column>
|
||||
<el-table-column label="RTSP 用户名" min-width="145"><template #default="scope"><el-input v-model="credentials[scope.row.id].rtspUsername" :disabled="credentials[scope.row.id].rtspSameAsOnvif" autocomplete="off" /></template></el-table-column>
|
||||
<el-table-column label="RTSP 密码" min-width="160"><template #default="scope"><el-input v-model="credentials[scope.row.id].rtspPassword" :disabled="credentials[scope.row.id].rtspSameAsOnvif" type="password" show-password autocomplete="new-password" /></template></el-table-column>
|
||||
</el-table>
|
||||
<template #footer>
|
||||
<el-button type="primary" :loading="executing" @click="executeBatch">确认执行</el-button>
|
||||
<el-button @click="credentialOpen = false">取消</el-button>
|
||||
</template>
|
||||
</el-dialog>
|
||||
</template>
|
||||
</BasicLayout>
|
||||
</template>
|
||||
|
||||
<script setup>
|
||||
import { computed, onMounted, reactive, ref } from 'vue'
|
||||
import { Download, Refresh, Upload } from '@element-plus/icons-vue'
|
||||
import { ElMessage } from 'element-plus'
|
||||
import * as XLSX from 'xlsx'
|
||||
import { createProvisioningBatch, downloadProvisioningBatch, executeProvisioningBatch, getProvisioningBatch, listProvisioningBatches, retryProvisioningFailures, retryProvisioningItem } from '@/api/sense/provisioning'
|
||||
import { clearCredentialMap, createBatchPayload, credentialPayload, normalizeImportRows } from './provisioningPayload'
|
||||
|
||||
defineOptions({ name: 'SenseProvisioning' })
|
||||
|
||||
const loading = ref(false); const importing = ref(false); const executing = ref(false)
|
||||
const batches = ref([]); const total = ref(0); const statusFilter = ref(''); const fileName = ref(''); const importRows = ref([])
|
||||
const selectedBatch = ref(null); const detailOpen = ref(false); const credentialOpen = ref(false); const retryMode = ref(false)
|
||||
const singleItemId = ref('')
|
||||
const credentials = reactive({}); const query = reactive({ pageIndex: 1, pageSize: 10 })
|
||||
const statusOptions = [{ value: 'ready', label: '待执行' }, { value: 'partial', label: '部分成功' }, { value: 'succeeded', label: '已完成' }, { value: 'failed', label: '失败' }]
|
||||
const activeStep = computed(() => credentialOpen.value ? 2 : selectedBatch.value?.status === 'succeeded' || selectedBatch.value?.status === 'partial' ? 3 : importRows.value.length ? 1 : 0)
|
||||
const credentialItems = computed(() => (selectedBatch.value?.items || []).filter(item => singleItemId.value ? item.id === singleItemId.value : retryMode.value ? item.status === 'failed' : item.status === 'ready'))
|
||||
|
||||
function unwrap(response) { return response?.data?.data ?? response?.data ?? response }
|
||||
function statusLabel(status) { return ({ ready: '待执行', running: '执行中', succeeded: '已完成', partial: '部分成功', failed: '失败', invalid: '格式错误', quota_exceeded: '超出配额' })[status] || status }
|
||||
function statusType(status) { return ({ succeeded: 'success', partial: 'warning', failed: 'danger', invalid: 'danger', quota_exceeded: 'info', ready: 'primary' })[status] || 'info' }
|
||||
function saveText(content, name, type = 'text/csv;charset=utf-8') { const url = URL.createObjectURL(new Blob([content], { type })); const link = document.createElement('a'); link.href = url; link.download = name; link.click(); URL.revokeObjectURL(url) }
|
||||
function saveBlob(content, name) { const url = URL.createObjectURL(content); const link = document.createElement('a'); link.href = url; link.download = name; link.click(); URL.revokeObjectURL(url) }
|
||||
|
||||
async function loadBatches() {
|
||||
loading.value = true
|
||||
try {
|
||||
const payload = unwrap(await listProvisioningBatches({ ...query, status: statusFilter.value }))
|
||||
batches.value = payload?.list || payload?.data || []
|
||||
total.value = payload?.count || 0
|
||||
} finally { loading.value = false }
|
||||
}
|
||||
async function handleFile(file) {
|
||||
try {
|
||||
const workbook = XLSX.read(await file.raw.arrayBuffer(), { type: 'array' })
|
||||
const sheet = workbook.Sheets[workbook.SheetNames[0]]
|
||||
const rows = XLSX.utils.sheet_to_json(sheet, { defval: '' })
|
||||
importRows.value = normalizeImportRows(rows)
|
||||
fileName.value = file.name
|
||||
if (!importRows.value.length) ElMessage.warning('清单中没有可导入的数据')
|
||||
} catch (error) { ElMessage.error(error.message || '无法读取清单文件') }
|
||||
}
|
||||
async function submitImport() {
|
||||
if (!importRows.value.length) return ElMessage.warning('请先选择清单')
|
||||
importing.value = true
|
||||
try {
|
||||
const key = `sense-${Date.now()}-${Math.random().toString(16).slice(2)}`
|
||||
selectedBatch.value = unwrap(await createProvisioningBatch(createBatchPayload(importRows.value, key)))
|
||||
detailOpen.value = true
|
||||
ElMessage.success('预校验完成')
|
||||
await loadBatches()
|
||||
} catch (error) { ElMessage.error(error.message || '清单预校验失败') } finally { importing.value = false }
|
||||
}
|
||||
async function openBatch(id) { selectedBatch.value = unwrap(await getProvisioningBatch(id)); detailOpen.value = true }
|
||||
async function openCredentials(id, retry) {
|
||||
selectedBatch.value = unwrap(await getProvisioningBatch(id)); retryMode.value = retry; singleItemId.value = ''
|
||||
credentialItems.value.forEach(item => { credentials[item.id] = { onvifUsername: '', onvifPassword: '', rtspSameAsOnvif: true, rtspUsername: '', rtspPassword: '' } })
|
||||
credentialOpen.value = true
|
||||
}
|
||||
function openSingleRetry(item) { retryMode.value = true; singleItemId.value = item.id; credentials[item.id] = { onvifUsername: '', onvifPassword: '', rtspSameAsOnvif: true, rtspUsername: '', rtspPassword: '' }; detailOpen.value = false; credentialOpen.value = true }
|
||||
function clearCredentials() { clearCredentialMap(credentials); Object.keys(credentials).forEach(key => delete credentials[key]); singleItemId.value = '' }
|
||||
async function executeBatch() {
|
||||
if (credentialItems.value.some(item => !credentials[item.id]?.onvifUsername || !credentials[item.id]?.onvifPassword)) return ElMessage.warning('请填写所有待处理设备的 ONVIF 用户名和密码')
|
||||
if (credentialItems.value.some(item => credentials[item.id]?.rtspSameAsOnvif === false && (!credentials[item.id]?.rtspUsername || !credentials[item.id]?.rtspPassword))) return ElMessage.warning('请填写所有独立 RTSP 凭据')
|
||||
executing.value = true
|
||||
try {
|
||||
const payload = credentialPayload(credentialItems.value, credentials)
|
||||
const action = singleItemId.value
|
||||
? retryProvisioningItem(selectedBatch.value.id, singleItemId.value, payload)
|
||||
: retryMode.value
|
||||
? retryProvisioningFailures(selectedBatch.value.id, payload)
|
||||
: executeProvisioningBatch(selectedBatch.value.id, payload)
|
||||
selectedBatch.value = unwrap(await action)
|
||||
clearCredentials(); credentialOpen.value = false; detailOpen.value = true
|
||||
ElMessage.success('批量开通处理完成')
|
||||
await loadBatches()
|
||||
} catch (error) { clearCredentials(); ElMessage.error(error.message || '批量开通失败') } finally { executing.value = false }
|
||||
}
|
||||
async function exportSelected() { if (!selectedBatch.value) return; try { saveBlob(await downloadProvisioningBatch(selectedBatch.value.id), `sense-provisioning-${selectedBatch.value.id}.csv`) } catch (error) { ElMessage.error(error.message || '结果导出失败') } }
|
||||
function downloadTemplate() { saveText('\uFEFFline_number,name,location,address\r\n1,东门摄像机,教学楼一楼东门,http://192.0.2.10/onvif', 'sense-provisioning-template.csv') }
|
||||
onMounted(loadBatches)
|
||||
</script>
|
||||
|
||||
<style scoped>
|
||||
.page-header,.section-title,.toolbar,.import-actions,.summary-cards{display:flex;align-items:center;gap:12px}.page-header,.section-title{justify-content:space-between}.page-header h3,.section-title h4{margin:0 0 6px}.page-header p{margin:0;color:#909399}.workflow-steps{margin:18px 0}.toolbar{margin:18px 0}.file-name,.import-actions span{color:#606266;font-size:13px}.import-actions{justify-content:flex-end;margin-top:16px}.section-title{margin-bottom:12px}.section-title .el-select{width:160px}.summary-cards{margin-bottom:14px}.credential-table{margin-top:16px}
|
||||
</style>
|
||||
@@ -0,0 +1,58 @@
|
||||
const aliases = {
|
||||
lineNumber: ['lineNumber', 'line_number', '序号', '行号'],
|
||||
name: ['name', '设备名称', '摄像头名称'],
|
||||
location: ['location', '安装位置', '位置'],
|
||||
address: ['address', '设备地址', 'ONVIF地址', 'onvif_address']
|
||||
}
|
||||
|
||||
function readAlias(row, names) {
|
||||
const key = names.find(name => Object.prototype.hasOwnProperty.call(row, name))
|
||||
return key ? row[key] : ''
|
||||
}
|
||||
|
||||
export function normalizeImportRows(rows) {
|
||||
return rows.map((row, index) => ({
|
||||
lineNumber: Number(readAlias(row, aliases.lineNumber)) || index + 1,
|
||||
name: String(readAlias(row, aliases.name) || '').trim(),
|
||||
location: String(readAlias(row, aliases.location) || '').trim(),
|
||||
address: String(readAlias(row, aliases.address) || '').trim()
|
||||
}))
|
||||
}
|
||||
|
||||
export function createBatchPayload(rows, idempotencyKey) {
|
||||
return {
|
||||
idempotencyKey: String(idempotencyKey || '').trim(),
|
||||
rows: normalizeImportRows(rows).map(row => ({
|
||||
lineNumber: row.lineNumber,
|
||||
name: row.name,
|
||||
location: row.location,
|
||||
address: row.address
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
export function credentialPayload(items, credentialByItem) {
|
||||
return {
|
||||
credentials: items.map(item => {
|
||||
const value = credentialByItem[item.id] || {}
|
||||
const same = value.rtspSameAsOnvif !== false
|
||||
return {
|
||||
itemId: item.id,
|
||||
onvifUsername: String(value.onvifUsername || ''),
|
||||
onvifPassword: String(value.onvifPassword || ''),
|
||||
rtspSameAsOnvif: same,
|
||||
rtspUsername: same ? '' : String(value.rtspUsername || ''),
|
||||
rtspPassword: same ? '' : String(value.rtspPassword || '')
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
export function clearCredentialMap(credentialByItem) {
|
||||
Object.keys(credentialByItem).forEach(key => {
|
||||
credentialByItem[key].onvifUsername = ''
|
||||
credentialByItem[key].onvifPassword = ''
|
||||
credentialByItem[key].rtspUsername = ''
|
||||
credentialByItem[key].rtspPassword = ''
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
import axios from 'axios'
|
||||
import { downloadProvisioningBatch } from '@/api/sense/provisioning'
|
||||
import { getToken } from '@/utils/auth'
|
||||
|
||||
jest.mock('axios', () => ({ get: jest.fn() }))
|
||||
jest.mock('@/utils/auth', () => ({ getToken: jest.fn() }))
|
||||
jest.mock('@/utils/request', () => jest.fn())
|
||||
|
||||
describe('Sense provisioning result export', () => {
|
||||
test('downloads the audited server export as a blob', async() => {
|
||||
const blob = new Blob(['status'], { type: 'text/csv' })
|
||||
getToken.mockReturnValue('synthetic-token')
|
||||
axios.get.mockResolvedValue({ data: blob })
|
||||
await expect(downloadProvisioningBatch('batch-1')).resolves.toBe(blob)
|
||||
expect(axios.get).toHaveBeenCalledWith(expect.stringContaining('/api/v1/provisioning/batches/batch-1/export'), {
|
||||
responseType: 'blob',
|
||||
headers: { Authorization: 'Bearer synthetic-token' }
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,19 @@
|
||||
import { clearCredentialMap, createBatchPayload, credentialPayload, normalizeImportRows } from '@/views/sense/provisioning/provisioningPayload'
|
||||
|
||||
describe('Sense provisioning payloads', () => {
|
||||
test('normalizes Chinese template headers and excludes unknown columns', () => {
|
||||
const rows = normalizeImportRows([{ 序号: 3, 设备名称: ' 东门摄像机 ', 安装位置: ' 东门 ', ONVIF地址: 'http://192.0.2.10/onvif', 密码: 'must-not-be-sent' }])
|
||||
expect(createBatchPayload(rows, ' import-1 ')).toEqual({
|
||||
idempotencyKey: 'import-1',
|
||||
rows: [{ lineNumber: 3, name: '东门摄像机', location: '东门', address: 'http://192.0.2.10/onvif' }]
|
||||
})
|
||||
})
|
||||
|
||||
test('credential payload is allowlisted and can be cleared in place', () => {
|
||||
const values = { item1: { onvifUsername: 'installer', onvifPassword: 'temporary-secret', rtspSameAsOnvif: true, rtspUsername: 'ignored', rtspPassword: 'ignored', unexpected: 'ignored' }}
|
||||
expect(credentialPayload([{ id: 'item1' }], values)).toEqual({ credentials: [{ itemId: 'item1', onvifUsername: 'installer', onvifPassword: 'temporary-secret', rtspSameAsOnvif: true, rtspUsername: '', rtspPassword: '' }] })
|
||||
clearCredentialMap(values)
|
||||
expect(values.item1.onvifPassword).toBe('')
|
||||
expect(values.item1.onvifUsername).toBe('')
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user