From ae19e59703f2a4913ae6717e869de80f838a3a92 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Fri, 28 Aug 2026 09:53:00 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=AE=9E=E7=8E=B0=20Sense=20=E6=89=B9?= =?UTF-8?q?=E9=87=8F=E5=BC=80=E9=80=9A=E4=B8=8E=E5=A4=B1=E8=B4=A5=E9=87=8D?= =?UTF-8?q?=E8=AF=95=20(#72)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../app/admin/router/sense_provisioning.go | 24 ++ Sense/server/app/sense/provisioning/apis.go | 170 ++++++++ .../app/sense/provisioning/apis_test.go | 32 ++ Sense/server/app/sense/provisioning/dto.go | 68 ++++ Sense/server/app/sense/provisioning/models.go | 58 +++ .../server/app/sense/provisioning/service.go | 366 ++++++++++++++++++ .../app/sense/provisioning/service_test.go | 278 +++++++++++++ .../version/2026082809000_provisioning.go | 75 ++++ .../2026082809000_provisioning_test.go | 61 +++ Sense/ui/src/api/sense/provisioning.js | 36 ++ .../ui/src/views/sense/provisioning/index.vue | 200 ++++++++++ .../sense/provisioning/provisioningPayload.js | 58 +++ .../tests/unit/sense/provisioningApi.spec.js | 20 + .../unit/sense/provisioningPayload.spec.js | 19 + 14 files changed, 1465 insertions(+) create mode 100644 Sense/server/app/admin/router/sense_provisioning.go create mode 100644 Sense/server/app/sense/provisioning/apis.go create mode 100644 Sense/server/app/sense/provisioning/apis_test.go create mode 100644 Sense/server/app/sense/provisioning/dto.go create mode 100644 Sense/server/app/sense/provisioning/models.go create mode 100644 Sense/server/app/sense/provisioning/service.go create mode 100644 Sense/server/app/sense/provisioning/service_test.go create mode 100644 Sense/server/cmd/migrate/migration/version/2026082809000_provisioning.go create mode 100644 Sense/server/cmd/migrate/migration/version/2026082809000_provisioning_test.go create mode 100644 Sense/ui/src/api/sense/provisioning.js create mode 100644 Sense/ui/src/views/sense/provisioning/index.vue create mode 100644 Sense/ui/src/views/sense/provisioning/provisioningPayload.js create mode 100644 Sense/ui/tests/unit/sense/provisioningApi.spec.js create mode 100644 Sense/ui/tests/unit/sense/provisioningPayload.spec.js diff --git a/Sense/server/app/admin/router/sense_provisioning.go b/Sense/server/app/admin/router/sense_provisioning.go new file mode 100644 index 0000000..87de6b8 --- /dev/null +++ b/Sense/server/app/admin/router/sense_provisioning.go @@ -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) +} diff --git a/Sense/server/app/sense/provisioning/apis.go b/Sense/server/app/sense/provisioning/apis.go new file mode 100644 index 0000000..5c6589a --- /dev/null +++ b/Sense/server/app/sense/provisioning/apis.go @@ -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 = "" + } +} diff --git a/Sense/server/app/sense/provisioning/apis_test.go b/Sense/server/app/sense/provisioning/apis_test.go new file mode 100644 index 0000000..c95e07a --- /dev/null +++ b/Sense/server/app/sense/provisioning/apis_test.go @@ -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]) + } +} diff --git a/Sense/server/app/sense/provisioning/dto.go b/Sense/server/app/sense/provisioning/dto.go new file mode 100644 index 0000000..b32cbb4 --- /dev/null +++ b/Sense/server/app/sense/provisioning/dto.go @@ -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"` +} diff --git a/Sense/server/app/sense/provisioning/models.go b/Sense/server/app/sense/provisioning/models.go new file mode 100644 index 0000000..85dabb1 --- /dev/null +++ b/Sense/server/app/sense/provisioning/models.go @@ -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" } diff --git a/Sense/server/app/sense/provisioning/service.go b/Sense/server/app/sense/provisioning/service.go new file mode 100644 index 0000000..d957ecf --- /dev/null +++ b/Sense/server/app/sense/provisioning/service.go @@ -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 +} diff --git a/Sense/server/app/sense/provisioning/service_test.go b/Sense/server/app/sense/provisioning/service_test.go new file mode 100644 index 0000000..ed78bae --- /dev/null +++ b/Sense/server/app/sense/provisioning/service_test.go @@ -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") + } +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082809000_provisioning.go b/Sense/server/cmd/migrate/migration/version/2026082809000_provisioning.go new file mode 100644 index 0000000..9e9417f --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082809000_provisioning.go @@ -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 + }) +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082809000_provisioning_test.go b/Sense/server/cmd/migrate/migration/version/2026082809000_provisioning_test.go new file mode 100644 index 0000000..2ab2b15 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082809000_provisioning_test.go @@ -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) + } +} diff --git a/Sense/ui/src/api/sense/provisioning.js b/Sense/ui/src/api/sense/provisioning.js new file mode 100644 index 0000000..b31463e --- /dev/null +++ b/Sense/ui/src/api/sense/provisioning.js @@ -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 +} diff --git a/Sense/ui/src/views/sense/provisioning/index.vue b/Sense/ui/src/views/sense/provisioning/index.vue new file mode 100644 index 0000000..5904070 --- /dev/null +++ b/Sense/ui/src/views/sense/provisioning/index.vue @@ -0,0 +1,200 @@ + + + + + diff --git a/Sense/ui/src/views/sense/provisioning/provisioningPayload.js b/Sense/ui/src/views/sense/provisioning/provisioningPayload.js new file mode 100644 index 0000000..8c1d64a --- /dev/null +++ b/Sense/ui/src/views/sense/provisioning/provisioningPayload.js @@ -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 = '' + }) +} diff --git a/Sense/ui/tests/unit/sense/provisioningApi.spec.js b/Sense/ui/tests/unit/sense/provisioningApi.spec.js new file mode 100644 index 0000000..d66436f --- /dev/null +++ b/Sense/ui/tests/unit/sense/provisioningApi.spec.js @@ -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' } + }) + }) +}) diff --git a/Sense/ui/tests/unit/sense/provisioningPayload.spec.js b/Sense/ui/tests/unit/sense/provisioningPayload.spec.js new file mode 100644 index 0000000..b4b0ee4 --- /dev/null +++ b/Sense/ui/tests/unit/sense/provisioningPayload.spec.js @@ -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('') + }) +}) -- 2.34.1