2026-07-15 17:33:55 +08:00
|
|
|
// Package handler provides HTTP handlers for various endpoints.
|
|
|
|
|
package handler
|
|
|
|
|
|
|
|
|
|
import (
|
2026-07-20 15:46:31 +08:00
|
|
|
"bytes"
|
|
|
|
|
"context"
|
|
|
|
|
"encoding/json"
|
|
|
|
|
"errors"
|
2026-07-15 17:33:55 +08:00
|
|
|
"fmt"
|
2026-07-20 15:46:31 +08:00
|
|
|
"strconv"
|
2026-07-15 17:33:55 +08:00
|
|
|
"strings"
|
2026-07-20 15:46:31 +08:00
|
|
|
"time"
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-08-04 17:12:17 +08:00
|
|
|
"modelRT/client/manualsync"
|
2026-07-20 15:46:31 +08:00
|
|
|
"modelRT/common"
|
2026-07-15 17:33:55 +08:00
|
|
|
"modelRT/common/errcode"
|
|
|
|
|
"modelRT/constants"
|
|
|
|
|
"modelRT/database"
|
|
|
|
|
"modelRT/diagram"
|
|
|
|
|
"modelRT/logger"
|
2026-07-20 15:46:31 +08:00
|
|
|
"modelRT/model"
|
2026-07-15 17:33:55 +08:00
|
|
|
"modelRT/orm"
|
2026-08-04 17:12:17 +08:00
|
|
|
redisrepository "modelRT/repository/redis"
|
2026-07-15 17:33:55 +08:00
|
|
|
|
|
|
|
|
"github.com/gin-gonic/gin"
|
2026-07-20 15:46:31 +08:00
|
|
|
"gorm.io/gorm"
|
2026-07-15 17:33:55 +08:00
|
|
|
)
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
type dataObjectAttributeUpdateRequest struct {
|
|
|
|
|
Token string `json:"token"`
|
|
|
|
|
Field string `json:"field"`
|
|
|
|
|
Value json.RawMessage `json:"value"`
|
2026-07-21 16:13:23 +08:00
|
|
|
Data json.RawMessage `json:"data,omitempty"`
|
2026-07-20 15:46:31 +08:00
|
|
|
}
|
|
|
|
|
|
2026-08-04 17:12:17 +08:00
|
|
|
const redisChangeRestoreTimeout = 5 * time.Second
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
// DataObjectAttributeUpdateHandler updates the writable field of one data object.
|
2026-07-15 17:33:55 +08:00
|
|
|
func DataObjectAttributeUpdateHandler(c *gin.Context) {
|
2026-07-20 15:46:31 +08:00
|
|
|
ctx := c.Request.Context()
|
|
|
|
|
var request dataObjectAttributeUpdateRequest
|
2026-07-15 17:33:55 +08:00
|
|
|
if err := c.ShouldBindJSON(&request); err != nil {
|
2026-07-20 15:46:31 +08:00
|
|
|
logger.Error(ctx, "unmarshal data-object update request failed", "error", err)
|
2026-07-15 17:33:55 +08:00
|
|
|
renderRespFailure(c, constants.RespCodeInvalidParams, err.Error(), nil)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
dataObjectType, field, value, err := validateDataObjectAttributeUpdate(request)
|
|
|
|
|
if err != nil {
|
|
|
|
|
logger.Warn(ctx, "validate data-object update request failed", "token", request.Token, "field", request.Field, "error", err)
|
|
|
|
|
renderRespFailure(c, constants.RespCodeInvalidParams, err.Error(), nil)
|
|
|
|
|
return
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
tx := database.GetPostgresDBClient().WithContext(ctx).Begin()
|
2026-07-15 17:33:55 +08:00
|
|
|
if tx.Error != nil {
|
2026-07-20 15:46:31 +08:00
|
|
|
logger.Error(ctx, "begin data-object update transaction failed", "error", tx.Error)
|
2026-07-15 17:33:55 +08:00
|
|
|
renderRespFailure(c, constants.RespCodeServerError, "begin postgres transaction failed", nil)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-07-21 16:13:23 +08:00
|
|
|
transactionCompleted := false
|
|
|
|
|
defer func() {
|
|
|
|
|
if !transactionCompleted {
|
|
|
|
|
_ = tx.Rollback().Error
|
|
|
|
|
}
|
|
|
|
|
}()
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-31 13:39:18 +08:00
|
|
|
redisClient := diagram.GetRedisClientInstance()
|
2026-08-04 17:12:17 +08:00
|
|
|
redisChanges := redisrepository.NewRedisChangeSet(redisClient)
|
2026-07-31 13:39:18 +08:00
|
|
|
canonicalRedisKey, err := model.ResolveDataObjectRedisKey(
|
|
|
|
|
ctx,
|
|
|
|
|
redisClient,
|
|
|
|
|
dataObjectType,
|
|
|
|
|
request.Token,
|
|
|
|
|
)
|
2026-07-20 15:46:31 +08:00
|
|
|
message := "data-object attribute update success"
|
|
|
|
|
var measurementResult measurementUpdateResult
|
2026-07-31 13:39:18 +08:00
|
|
|
switch {
|
|
|
|
|
case err != nil:
|
|
|
|
|
// The shared resolver error is handled by the common failure path below.
|
|
|
|
|
case dataObjectType == constants.DataObjectTypeParameter:
|
2026-07-20 15:46:31 +08:00
|
|
|
parameter, queryErr := database.QueryParameterByDataObjectToken(ctx, tx, request.Token)
|
|
|
|
|
if queryErr == nil {
|
|
|
|
|
queryErr = database.UpdateParameterDataObjectValue(ctx, tx, parameter, value)
|
|
|
|
|
}
|
2026-07-30 15:32:22 +08:00
|
|
|
if queryErr == nil {
|
2026-07-31 13:39:18 +08:00
|
|
|
queryErr = redisChanges.AddHashChange(ctx, canonicalRedisKey, field, value)
|
2026-07-30 15:32:22 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
err = queryErr
|
2026-07-31 13:39:18 +08:00
|
|
|
case dataObjectType == constants.DataObjectTypeMeasurement:
|
2026-07-21 16:13:23 +08:00
|
|
|
measurementResult, err = updateMeasurementDataObject(ctx, tx, request.Token, field, value, request.Data, measurementUpdateDependencies{
|
2026-08-04 17:12:17 +08:00
|
|
|
writeManualValueFunc: func(ctx context.Context, measurement *orm.Measurement, value float64, timestamp time.Time) error {
|
|
|
|
|
key, err := model.GenerateMeasureIdentifier(measurement.DataSource)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return fmt.Errorf("generate measurement redis key: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return redisChanges.AddMeasurementValueChange(ctx, key, value, timestamp, false)
|
2026-07-30 15:32:22 +08:00
|
|
|
},
|
2026-08-04 17:12:17 +08:00
|
|
|
syncManualChangeFunc: manualsync.Sync,
|
|
|
|
|
replaceRedisValueFunc: func(ctx context.Context, measurement *orm.Measurement, value float64, timestamp time.Time) error {
|
|
|
|
|
key, err := model.GenerateMeasureIdentifier(measurement.DataSource)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return fmt.Errorf("generate measurement redis key: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return redisChanges.AddMeasurementValueChange(ctx, key, value, timestamp, true)
|
2026-07-30 15:32:22 +08:00
|
|
|
},
|
2026-08-04 17:12:17 +08:00
|
|
|
nowFunc: time.Now,
|
2026-07-21 16:13:23 +08:00
|
|
|
})
|
2026-07-30 15:32:22 +08:00
|
|
|
if err == nil && measurementResult.modeChanged {
|
2026-07-31 13:39:18 +08:00
|
|
|
err = redisChanges.AddHashChange(ctx, canonicalRedisKey, "mode", measurementResult.mode)
|
2026-07-30 15:32:22 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
message = measurementResult.message
|
|
|
|
|
default:
|
|
|
|
|
err = fmt.Errorf("unsupported data object type %q", dataObjectType)
|
|
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
if err != nil {
|
|
|
|
|
_ = tx.Rollback().Error
|
2026-07-30 15:32:22 +08:00
|
|
|
if measurementResult.recordFailureOnError {
|
2026-08-04 17:12:17 +08:00
|
|
|
if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, measurementFailureTime(measurementResult)); logErr != nil {
|
2026-07-20 15:46:31 +08:00
|
|
|
logger.Error(ctx, "append failed measurement value operation failed", "measurement_id", measurementResult.measurementID, "error", logErr)
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
|
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
logger.Warn(ctx, "update data-object attribute failed", "token", request.Token, "field", field, "error", err)
|
|
|
|
|
if isInvalidDataObjectUpdateError(err) {
|
|
|
|
|
renderRespFailure(c, constants.RespCodeInvalidParams, err.Error(), nil)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
renderRespFailure(c, constants.RespCodeFailed, err.Error(), nil)
|
|
|
|
|
return
|
|
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-30 15:32:22 +08:00
|
|
|
if err := redisChanges.Apply(ctx); err != nil {
|
|
|
|
|
_ = tx.Rollback().Error
|
|
|
|
|
if measurementResult.recordFailureOnError {
|
2026-08-04 17:12:17 +08:00
|
|
|
if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, measurementFailureTime(measurementResult)); logErr != nil {
|
2026-07-30 15:32:22 +08:00
|
|
|
logger.Error(ctx, "append failed measurement value operation failed", "measurement_id", measurementResult.measurementID, "error", logErr)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
logger.Error(ctx, "apply redis data-object changes failed", "token", request.Token, "field", field, "error", err)
|
|
|
|
|
renderRespFailure(c, constants.RespCodeFailed, err.Error(), nil)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
if err := tx.Commit().Error; err != nil {
|
2026-07-30 15:32:22 +08:00
|
|
|
revertCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), redisChangeRestoreTimeout)
|
|
|
|
|
defer cancel()
|
|
|
|
|
if redisErr := redisChanges.Revert(revertCtx); redisErr != nil {
|
|
|
|
|
logger.Error(ctx, "revert redis data-object changes failed", "token", request.Token, "field", field, "error", redisErr)
|
|
|
|
|
}
|
|
|
|
|
if measurementResult.recordFailureOnError {
|
2026-08-04 17:12:17 +08:00
|
|
|
if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, measurementFailureTime(measurementResult)); logErr != nil {
|
2026-07-30 15:32:22 +08:00
|
|
|
logger.Error(ctx, "append failed measurement value operation failed", "measurement_id", measurementResult.measurementID, "error", logErr)
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
logger.Error(ctx, "commit data-object update transaction failed", "token", request.Token, "field", field, "error", err)
|
|
|
|
|
renderRespFailure(c, constants.RespCodeServerError, "transaction commit failed", nil)
|
2026-07-15 17:33:55 +08:00
|
|
|
return
|
|
|
|
|
}
|
2026-07-21 16:13:23 +08:00
|
|
|
transactionCompleted = true
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
renderRespSuccess(c, constants.RespCodeSuccess, message, map[string]any{
|
|
|
|
|
"token": request.Token,
|
|
|
|
|
"field": field,
|
|
|
|
|
"value": value,
|
|
|
|
|
})
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func validateDataObjectAttributeUpdate(request dataObjectAttributeUpdateRequest) (constants.DataObjectType, string, any, error) {
|
|
|
|
|
if request.Token == "" {
|
|
|
|
|
return "", "", nil, fmt.Errorf("token is required")
|
|
|
|
|
}
|
|
|
|
|
if len(bytes.TrimSpace(request.Value)) == 0 || bytes.Equal(bytes.TrimSpace(request.Value), []byte("null")) {
|
|
|
|
|
return "", "", nil, fmt.Errorf("value is required")
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
field := strings.ToLower(strings.TrimSpace(request.Field))
|
|
|
|
|
if field == "" {
|
2026-07-21 16:13:23 +08:00
|
|
|
field = "value"
|
2026-07-20 15:46:31 +08:00
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
dataObjectType, err := model.ClassifyDataObjectToken(request.Token)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return "", "", nil, err
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
switch dataObjectType {
|
|
|
|
|
case constants.DataObjectTypeParameter:
|
|
|
|
|
parts := strings.Split(request.Token, ".")
|
|
|
|
|
attributeGroup := parts[len(parts)-2]
|
|
|
|
|
if !isWritableParameterAttributeGroup(attributeGroup) {
|
|
|
|
|
return "", "", nil, fmt.Errorf("parameter updates do not support token6=%s", attributeGroup)
|
|
|
|
|
}
|
|
|
|
|
if field != "value" {
|
|
|
|
|
return "", "", nil, fmt.Errorf("parameter data objects only support updating field value")
|
|
|
|
|
}
|
|
|
|
|
value, err := decodeDataObjectUpdateValue(request.Value)
|
|
|
|
|
return dataObjectType, field, value, err
|
|
|
|
|
case constants.DataObjectTypeMeasurement:
|
|
|
|
|
parts := strings.Split(request.Token, ".")
|
2026-07-21 16:13:23 +08:00
|
|
|
if len(parts) != 2 && parts[len(parts)-2] != "bay" {
|
|
|
|
|
return "", "", nil, fmt.Errorf("measurement updates require token4.token7 or token6=bay")
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
switch field {
|
|
|
|
|
case "value":
|
|
|
|
|
value, err := parseMeasurementUpdateValue(request.Value)
|
|
|
|
|
return dataObjectType, field, value, err
|
|
|
|
|
case "mode":
|
|
|
|
|
value, err := parseMeasurementUpdateMode(request.Value)
|
|
|
|
|
return dataObjectType, field, value, err
|
|
|
|
|
default:
|
|
|
|
|
return "", "", nil, fmt.Errorf("measurement data objects only support updating fields value and mode")
|
|
|
|
|
}
|
|
|
|
|
default:
|
|
|
|
|
return "", "", nil, fmt.Errorf("unsupported data object type %q", dataObjectType)
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
func isWritableParameterAttributeGroup(group string) bool {
|
|
|
|
|
switch group {
|
2026-07-21 16:13:23 +08:00
|
|
|
case "rated", "setup", "model", "stable", "craft", "integrity", "behavior", "base_extend":
|
2026-07-20 15:46:31 +08:00
|
|
|
return true
|
|
|
|
|
default:
|
|
|
|
|
return false
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
func decodeDataObjectUpdateValue(raw json.RawMessage) (any, error) {
|
|
|
|
|
decoder := json.NewDecoder(bytes.NewReader(raw))
|
|
|
|
|
decoder.UseNumber()
|
|
|
|
|
var value any
|
|
|
|
|
if err := decoder.Decode(&value); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("decode update value: %w", err)
|
|
|
|
|
}
|
|
|
|
|
if number, ok := value.(json.Number); ok {
|
|
|
|
|
if integer, err := number.Int64(); err == nil {
|
|
|
|
|
return integer, nil
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
decimal, err := number.Float64()
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, fmt.Errorf("invalid numeric update value %q: %w", number, err)
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
return decimal, nil
|
|
|
|
|
}
|
|
|
|
|
return value, nil
|
|
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
func parseMeasurementUpdateValue(raw json.RawMessage) (float64, error) {
|
|
|
|
|
var number float64
|
|
|
|
|
if err := json.Unmarshal(raw, &number); err == nil {
|
|
|
|
|
return number, nil
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
var text string
|
|
|
|
|
if err := json.Unmarshal(raw, &text); err != nil {
|
|
|
|
|
return 0, fmt.Errorf("measurement value must be a number or numeric string")
|
|
|
|
|
}
|
|
|
|
|
number, err := strconv.ParseFloat(text, 64)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return 0, fmt.Errorf("measurement value %q is not numeric: %w", text, err)
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
return number, nil
|
|
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-21 16:13:23 +08:00
|
|
|
func parseMeasurementUpdateMode(raw json.RawMessage) (int16, error) {
|
|
|
|
|
var mode int16
|
|
|
|
|
if err := json.Unmarshal(raw, &mode); err != nil {
|
|
|
|
|
return 0, fmt.Errorf("measurement mode must be 0 (manual) or 1 (automatic)")
|
2026-07-20 15:46:31 +08:00
|
|
|
}
|
2026-07-21 16:13:23 +08:00
|
|
|
if mode != constants.MeasurementModeManual && mode != constants.MeasurementModeAutomatic {
|
|
|
|
|
return 0, fmt.Errorf("measurement mode must be 0 (manual) or 1 (automatic)")
|
2026-07-20 15:46:31 +08:00
|
|
|
}
|
2026-07-21 16:13:23 +08:00
|
|
|
return mode, nil
|
2026-07-20 15:46:31 +08:00
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-08-04 17:12:17 +08:00
|
|
|
type measurementManualValueWriter func(context.Context, *orm.Measurement, float64, time.Time) error
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-08-04 17:12:17 +08:00
|
|
|
type measurementManualChangeSyncer func(context.Context, orm.JSONMap, int16, *manualsync.SyntheticData) error
|
2026-07-21 16:13:23 +08:00
|
|
|
|
2026-08-04 17:12:17 +08:00
|
|
|
type measurementRedisValueReplacer func(context.Context, *orm.Measurement, float64, time.Time) error
|
2026-07-21 16:13:23 +08:00
|
|
|
|
|
|
|
|
type measurementUpdateDependencies struct {
|
|
|
|
|
writeManualValueFunc measurementManualValueWriter
|
2026-08-04 17:12:17 +08:00
|
|
|
syncManualChangeFunc measurementManualChangeSyncer
|
2026-07-21 16:13:23 +08:00
|
|
|
replaceRedisValueFunc measurementRedisValueReplacer
|
2026-08-04 17:12:17 +08:00
|
|
|
nowFunc func() time.Time
|
2026-07-21 16:13:23 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
|
|
|
|
|
type measurementUpdateResult struct {
|
2026-07-30 15:32:22 +08:00
|
|
|
message string
|
|
|
|
|
measurementID int64
|
|
|
|
|
value float64
|
|
|
|
|
recordFailureOnError bool
|
|
|
|
|
mode int16
|
|
|
|
|
modeChanged bool
|
2026-08-04 17:12:17 +08:00
|
|
|
operationTime time.Time
|
2026-07-20 15:46:31 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func updateMeasurementDataObject(
|
|
|
|
|
ctx context.Context,
|
|
|
|
|
tx *gorm.DB,
|
|
|
|
|
token, field string,
|
|
|
|
|
value any,
|
2026-07-21 16:13:23 +08:00
|
|
|
modeData json.RawMessage,
|
|
|
|
|
dependencies measurementUpdateDependencies,
|
2026-07-20 15:46:31 +08:00
|
|
|
) (measurementUpdateResult, error) {
|
|
|
|
|
measurement, _, err := database.QueryMeasurementByDataObjectToken(ctx, tx, token)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return measurementUpdateResult{}, err
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-21 16:13:23 +08:00
|
|
|
lockedMeasurement, err := database.QueryMeasurementByIDForUpdate(ctx, tx, measurement.ID)
|
2026-07-20 15:46:31 +08:00
|
|
|
if err != nil {
|
|
|
|
|
return measurementUpdateResult{}, fmt.Errorf("lock measurement %d for update: %w", measurement.ID, err)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
switch field {
|
|
|
|
|
case "mode":
|
2026-07-21 16:13:23 +08:00
|
|
|
mode, ok := value.(int16)
|
2026-07-20 15:46:31 +08:00
|
|
|
if !ok {
|
|
|
|
|
return measurementUpdateResult{}, fmt.Errorf("measurement mode has invalid type %T", value)
|
|
|
|
|
}
|
2026-07-21 16:13:23 +08:00
|
|
|
currentMode, err := measurementModeIsAutomatic(lockedMeasurement.Mode)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return measurementUpdateResult{}, err
|
|
|
|
|
}
|
|
|
|
|
targetAutomatic := mode == constants.MeasurementModeAutomatic
|
|
|
|
|
if currentMode == targetAutomatic {
|
2026-07-20 15:46:31 +08:00
|
|
|
return measurementUpdateResult{message: fmt.Sprintf("measurement is already in %s mode", measurementModeName(mode))}, nil
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
2026-08-04 17:12:17 +08:00
|
|
|
operationTime := measurementUpdateNow(dependencies)
|
2026-07-21 16:13:23 +08:00
|
|
|
var manualValue *float64
|
|
|
|
|
if currentMode && mode == constants.MeasurementModeManual {
|
|
|
|
|
manualValue, err = parseOptionalMeasurementModeData(modeData)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return measurementUpdateResult{}, err
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-08-04 17:12:17 +08:00
|
|
|
if err := database.UpdateMeasurementModeWithOperation(ctx, tx, lockedMeasurement.ID, mode, operationTime); err != nil {
|
2026-07-20 15:46:31 +08:00
|
|
|
return measurementUpdateResult{}, err
|
|
|
|
|
}
|
2026-08-04 17:12:17 +08:00
|
|
|
syncManualMeasurementChange(
|
|
|
|
|
ctx,
|
|
|
|
|
dependencies.syncManualChangeFunc,
|
|
|
|
|
lockedMeasurement.ID,
|
|
|
|
|
lockedMeasurement.DataSource,
|
|
|
|
|
mode,
|
|
|
|
|
nil,
|
|
|
|
|
)
|
2026-07-21 16:13:23 +08:00
|
|
|
if currentMode && mode == constants.MeasurementModeManual {
|
|
|
|
|
if manualValue != nil {
|
|
|
|
|
if dependencies.replaceRedisValueFunc == nil {
|
|
|
|
|
return measurementUpdateResult{}, fmt.Errorf("measurement redis value replacer is nil")
|
|
|
|
|
}
|
2026-08-04 17:12:17 +08:00
|
|
|
if err := dependencies.replaceRedisValueFunc(ctx, &lockedMeasurement, *manualValue, operationTime); err != nil {
|
2026-07-21 16:13:23 +08:00
|
|
|
return measurementUpdateResult{}, fmt.Errorf("replace measurement redis value: %w", err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-07-30 15:32:22 +08:00
|
|
|
return measurementUpdateResult{
|
2026-08-04 17:12:17 +08:00
|
|
|
message: fmt.Sprintf("measurement mode changed to %s", measurementModeName(mode)),
|
|
|
|
|
mode: mode,
|
|
|
|
|
modeChanged: true,
|
|
|
|
|
operationTime: operationTime,
|
2026-07-30 15:32:22 +08:00
|
|
|
}, nil
|
2026-07-20 15:46:31 +08:00
|
|
|
case "value":
|
2026-07-21 16:13:23 +08:00
|
|
|
currentMode, err := measurementModeIsAutomatic(lockedMeasurement.Mode)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return measurementUpdateResult{}, err
|
|
|
|
|
}
|
|
|
|
|
if currentMode {
|
2026-07-20 15:46:31 +08:00
|
|
|
return measurementUpdateResult{}, fmt.Errorf("measurement value is read-only while mode is automatic")
|
|
|
|
|
}
|
|
|
|
|
manualValue, ok := value.(float64)
|
|
|
|
|
if !ok {
|
|
|
|
|
return measurementUpdateResult{}, fmt.Errorf("measurement value has invalid type %T", value)
|
|
|
|
|
}
|
2026-08-04 17:12:17 +08:00
|
|
|
operationTime := measurementUpdateNow(dependencies)
|
2026-07-20 15:46:31 +08:00
|
|
|
failureResult := measurementUpdateResult{
|
2026-07-30 15:32:22 +08:00
|
|
|
measurementID: lockedMeasurement.ID,
|
|
|
|
|
value: manualValue,
|
|
|
|
|
recordFailureOnError: true,
|
2026-08-04 17:12:17 +08:00
|
|
|
operationTime: operationTime,
|
2026-07-20 15:46:31 +08:00
|
|
|
}
|
2026-07-21 16:13:23 +08:00
|
|
|
if dependencies.writeManualValueFunc == nil {
|
2026-07-20 15:46:31 +08:00
|
|
|
return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(fmt.Errorf("measurement manual value writer is nil"))
|
|
|
|
|
}
|
2026-08-04 17:12:17 +08:00
|
|
|
if err := dependencies.writeManualValueFunc(ctx, &lockedMeasurement, manualValue, operationTime); err != nil {
|
2026-07-20 15:46:31 +08:00
|
|
|
return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(err)
|
|
|
|
|
}
|
2026-08-04 17:12:17 +08:00
|
|
|
sample := manualsync.SyntheticData{Time: operationTime.UnixNano(), Value: manualValue}
|
|
|
|
|
syncManualMeasurementChange(
|
|
|
|
|
ctx,
|
|
|
|
|
dependencies.syncManualChangeFunc,
|
|
|
|
|
lockedMeasurement.ID,
|
|
|
|
|
lockedMeasurement.DataSource,
|
|
|
|
|
constants.MeasurementModeManual,
|
|
|
|
|
&sample,
|
|
|
|
|
)
|
|
|
|
|
if err := database.AppendMeasurementValueOperation(ctx, tx, lockedMeasurement.ID, 0, manualValue, operationTime); err != nil {
|
2026-07-20 15:46:31 +08:00
|
|
|
return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(err)
|
|
|
|
|
}
|
2026-07-30 15:32:22 +08:00
|
|
|
failureResult.message = "measurement manual value updated"
|
|
|
|
|
return failureResult, nil
|
2026-07-20 15:46:31 +08:00
|
|
|
default:
|
|
|
|
|
return measurementUpdateResult{}, fmt.Errorf("unsupported measurement update field %q", field)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-04 17:12:17 +08:00
|
|
|
// TODO: This synchronous best-effort manual synchronization is not necessarily
|
|
|
|
|
// the final implementation. Its current behavior is to keep the local update
|
|
|
|
|
// successful when the downstream request fails (or the client is unavailable),
|
|
|
|
|
// log the error, and permanently drop that synchronization event without retry.
|
|
|
|
|
// The HTTP timeout still adds latency to the update request, and downstream state
|
|
|
|
|
// may diverge from local state. Revisit asynchronous delivery or a durable outbox
|
|
|
|
|
// if delivery reliability or request latency becomes important.
|
|
|
|
|
func syncManualMeasurementChange(
|
|
|
|
|
ctx context.Context,
|
|
|
|
|
syncer measurementManualChangeSyncer,
|
|
|
|
|
measurementID int64,
|
|
|
|
|
dataSource orm.JSONMap,
|
|
|
|
|
mode int16,
|
|
|
|
|
sample *manualsync.SyntheticData,
|
|
|
|
|
) {
|
|
|
|
|
if syncer == nil {
|
|
|
|
|
logManualMeasurementSyncError(ctx, "manual measurement synchronization skipped because sync client is unavailable",
|
|
|
|
|
"measurement_id", measurementID,
|
|
|
|
|
"mode", mode,
|
|
|
|
|
"has_data", sample != nil,
|
|
|
|
|
)
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
if err := syncer(ctx, dataSource, mode, sample); err != nil {
|
|
|
|
|
logManualMeasurementSyncError(ctx, "manual measurement synchronization failed; local update will continue",
|
|
|
|
|
"measurement_id", measurementID,
|
|
|
|
|
"mode", mode,
|
|
|
|
|
"has_data", sample != nil,
|
|
|
|
|
"error", err,
|
|
|
|
|
)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func logManualMeasurementSyncError(ctx context.Context, message string, fields ...any) {
|
|
|
|
|
// The application initializes logging before serving requests. This guard keeps
|
|
|
|
|
// the best-effort path safe in isolated unit tests and other pre-init callers.
|
|
|
|
|
if logger.GetLoggerInstance() == nil {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
logger.Error(ctx, message, fields...)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func measurementUpdateNow(dependencies measurementUpdateDependencies) time.Time {
|
|
|
|
|
if dependencies.nowFunc != nil {
|
|
|
|
|
return dependencies.nowFunc().UTC()
|
|
|
|
|
}
|
|
|
|
|
return time.Now().UTC()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func measurementFailureTime(result measurementUpdateResult) time.Time {
|
|
|
|
|
if !result.operationTime.IsZero() {
|
|
|
|
|
return result.operationTime
|
|
|
|
|
}
|
|
|
|
|
return time.Now().UTC()
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-21 16:13:23 +08:00
|
|
|
func parseOptionalMeasurementModeData(raw json.RawMessage) (*float64, error) {
|
|
|
|
|
trimmed := bytes.TrimSpace(raw)
|
|
|
|
|
if len(trimmed) == 0 || bytes.Equal(trimmed, []byte("null")) {
|
|
|
|
|
return nil, nil
|
|
|
|
|
}
|
|
|
|
|
value, err := parseMeasurementUpdateValue(trimmed)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, fmt.Errorf("invalid measurement mode data: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return &value, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func measurementModeName(mode int16) string {
|
|
|
|
|
if mode == constants.MeasurementModeAutomatic {
|
2026-07-20 15:46:31 +08:00
|
|
|
return "automatic"
|
2026-07-15 17:33:55 +08:00
|
|
|
}
|
2026-07-20 15:46:31 +08:00
|
|
|
return "manual"
|
|
|
|
|
}
|
2026-07-15 17:33:55 +08:00
|
|
|
|
2026-07-21 16:13:23 +08:00
|
|
|
func measurementModeIsAutomatic(mode int16) (bool, error) {
|
|
|
|
|
switch mode {
|
|
|
|
|
case constants.MeasurementModeManual:
|
|
|
|
|
return false, nil
|
|
|
|
|
case constants.MeasurementModeAutomatic:
|
|
|
|
|
return true, nil
|
|
|
|
|
default:
|
|
|
|
|
return false, fmt.Errorf("measurement has invalid mode %d", mode)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-07-20 15:46:31 +08:00
|
|
|
func isInvalidDataObjectUpdateError(err error) bool {
|
|
|
|
|
return errors.Is(err, common.ErrInvalidParameterToken) ||
|
|
|
|
|
errors.Is(err, common.ErrParameterTokenNotFound) ||
|
|
|
|
|
errors.Is(err, common.ErrAmbiguousParameterToken) ||
|
|
|
|
|
errors.Is(err, common.ErrInvalidMeasurementToken) ||
|
|
|
|
|
errors.Is(err, common.ErrMeasurementTokenNotFound) ||
|
|
|
|
|
errors.Is(err, common.ErrAmbiguousMeasurementToken)
|
|
|
|
|
}
|