From d67096385465421f58b71a3a7e1c50c8d9660d16 Mon Sep 17 00:00:00 2001 From: douxu Date: Fri, 31 Jul 2026 13:39:18 +0800 Subject: [PATCH] refactor(redis): normalize data object storage with aliases - store one canonical hash for each parameter and measurement - map supported token formats to canonical Redis keys through aliases - resolve aliases in query and update handlers - clean stale hashes and aliases during initialization - update Redis change logic and related tests --- constants/redis.go | 12 ++ handler/data_object_attribute_query.go | 53 +++----- handler/data_object_attribute_update.go | 21 +++- handler/data_object_redis_change.go | 107 ++-------------- handler/data_object_redis_change_test.go | 35 ------ model/data_object_redis_alias.go | 78 ++++++++++++ model/data_object_redis_alias_test.go | 25 ++++ model/measurement_data_object_init.go | 81 ++++++++----- model/measurement_data_object_init_test.go | 9 +- model/parameter_data_object_init.go | 134 +++++++++++++-------- model/parameter_data_object_init_test.go | 8 +- 11 files changed, 308 insertions(+), 255 deletions(-) create mode 100644 model/data_object_redis_alias.go create mode 100644 model/data_object_redis_alias_test.go diff --git a/constants/redis.go b/constants/redis.go index 3d4d347..ee317af 100644 --- a/constants/redis.go +++ b/constants/redis.go @@ -9,7 +9,19 @@ const ( // startup so stale parameter data-object keys can be removed safely. RedisParameterDataObjectKeySet = "modelrt:parameter-data-object:keys" + // RedisParameterDataObjectAliasKeySet tracks parameter alias string keys. + RedisParameterDataObjectAliasKeySet = "modelrt:parameter-data-object:alias-keys" + + // RedisParameterDataObjectAliasPrefix prefixes parameter token aliases. + RedisParameterDataObjectAliasPrefix = "modelrt:data-object:alias:parameter:" + // RedisMeasurementDataObjectKeySet tracks measurement hashes created during // startup so stale measurement data-object keys can be removed safely. RedisMeasurementDataObjectKeySet = "modelrt:measurement-data-object:keys" + + // RedisMeasurementDataObjectAliasKeySet tracks measurement alias string keys. + RedisMeasurementDataObjectAliasKeySet = "modelrt:measurement-data-object:alias-keys" + + // RedisMeasurementDataObjectAliasPrefix prefixes measurement token aliases. + RedisMeasurementDataObjectAliasPrefix = "modelrt:data-object:alias:measurement:" ) diff --git a/handler/data_object_attribute_query.go b/handler/data_object_attribute_query.go index 84b0e4f..d100226 100644 --- a/handler/data_object_attribute_query.go +++ b/handler/data_object_attribute_query.go @@ -196,26 +196,16 @@ func loadDataObjectHashField( if rdb == nil { return "", fmt.Errorf("redis client is not initialized") } - value, err := rdb.HGet(ctx, token, field).Result() + canonicalKey, err := model.ResolveDataObjectRedisKey(ctx, rdb, dataObjectType, token) + if err != nil { + return "", err + } + value, err := rdb.HGet(ctx, canonicalKey, field).Result() if errors.Is(err, redis.Nil) { - exists, existsErr := rdb.Exists(ctx, token).Result() - if existsErr != nil { - return "", fmt.Errorf("check redis data-object hash %q: %w", token, existsErr) - } - if exists > 0 { - return "", fmt.Errorf("redis data-object hash %q does not contain field %q", token, field) - } - switch dataObjectType { - case constants.DataObjectTypeParameter: - return "", fmt.Errorf("%w: %q", common.ErrParameterTokenNotFound, token) - case constants.DataObjectTypeMeasurement: - return "", fmt.Errorf("%w: %q", common.ErrMeasurementTokenNotFound, token) - default: - return "", fmt.Errorf("invalid data object type %q", dataObjectType) - } + return "", fmt.Errorf("canonical redis data-object hash %q does not contain field %q", canonicalKey, field) } if err != nil { - return "", fmt.Errorf("query redis hash %q field %q: %w", token, field, err) + return "", fmt.Errorf("query canonical redis hash %q field %q: %w", canonicalKey, field, err) } return value, nil } @@ -225,23 +215,18 @@ func loadMeasurementValueMetadata(ctx context.Context, token string) (orm.JSONMa if rdb == nil { return nil, 0, fmt.Errorf("redis client is not initialized") } - - values, err := rdb.HMGet(ctx, token, "data_source", "size").Result() + canonicalKey, err := model.ResolveDataObjectRedisKey(ctx, rdb, constants.DataObjectTypeMeasurement, token) if err != nil { - return nil, 0, fmt.Errorf("query redis hash %q measurement value metadata: %w", token, err) + return nil, 0, err + } + values, err := rdb.HMGet(ctx, canonicalKey, "data_source", "size").Result() + if err != nil { + return nil, 0, fmt.Errorf("query canonical redis hash %q measurement value metadata: %w", canonicalKey, err) } if len(values) != 2 { - return nil, 0, fmt.Errorf("redis hash %q returned %d measurement metadata fields", token, len(values)) + return nil, 0, fmt.Errorf("canonical redis hash %q returned %d measurement metadata fields", canonicalKey, len(values)) } if values[0] == nil || values[1] == nil { - exists, existsErr := rdb.Exists(ctx, token).Result() - if existsErr != nil { - return nil, 0, fmt.Errorf("check redis data-object hash %q: %w", token, existsErr) - } - if exists == 0 { - return nil, 0, fmt.Errorf("%w: %q", common.ErrMeasurementTokenNotFound, token) - } - missingFields := make([]string, 0, 2) if values[0] == nil { missingFields = append(missingFields, "data_source") @@ -250,24 +235,24 @@ func loadMeasurementValueMetadata(ctx context.Context, token string) (orm.JSONMa missingFields = append(missingFields, "size") } return nil, 0, fmt.Errorf( - "redis measurement hash %q does not contain field(s) %s", - token, + "canonical redis measurement hash %q does not contain field(s) %s", + canonicalKey, strings.Join(missingFields, ", "), ) } rawDataSource, ok := values[0].(string) if !ok { - return nil, 0, fmt.Errorf("redis measurement hash %q data_source has type %T", token, values[0]) + return nil, 0, fmt.Errorf("canonical redis measurement hash %q data_source has type %T", canonicalKey, values[0]) } var dataSource orm.JSONMap if err := json.Unmarshal([]byte(rawDataSource), &dataSource); err != nil { - return nil, 0, fmt.Errorf("decode measurement data_source from redis hash %q: %w", token, err) + return nil, 0, fmt.Errorf("decode measurement data_source from canonical redis hash %q: %w", canonicalKey, err) } rawSize, ok := values[1].(string) if !ok { - return nil, 0, fmt.Errorf("redis measurement hash %q size has type %T", token, values[1]) + return nil, 0, fmt.Errorf("canonical redis measurement hash %q size has type %T", canonicalKey, values[1]) } size, err := strconv.Atoi(rawSize) if err != nil { diff --git a/handler/data_object_attribute_update.go b/handler/data_object_attribute_update.go index feb1038..71b3c64 100644 --- a/handler/data_object_attribute_update.go +++ b/handler/data_object_attribute_update.go @@ -61,20 +61,29 @@ func DataObjectAttributeUpdateHandler(c *gin.Context) { } }() - redisChanges := NewRedisChangeSet(diagram.GetRedisClientInstance()) + redisClient := diagram.GetRedisClientInstance() + redisChanges := NewRedisChangeSet(redisClient) + canonicalRedisKey, err := model.ResolveDataObjectRedisKey( + ctx, + redisClient, + dataObjectType, + request.Token, + ) message := "data-object attribute update success" var measurementResult measurementUpdateResult - switch dataObjectType { - case constants.DataObjectTypeParameter: + switch { + case err != nil: + // The shared resolver error is handled by the common failure path below. + case dataObjectType == constants.DataObjectTypeParameter: parameter, queryErr := database.QueryParameterByDataObjectToken(ctx, tx, request.Token) if queryErr == nil { queryErr = database.UpdateParameterDataObjectValue(ctx, tx, parameter, value) } if queryErr == nil { - queryErr = redisChanges.AddDataObjectHashChange(ctx, dataObjectType, request.Token, field, value) + queryErr = redisChanges.AddHashChange(ctx, canonicalRedisKey, field, value) } err = queryErr - case constants.DataObjectTypeMeasurement: + case dataObjectType == constants.DataObjectTypeMeasurement: measurementResult, err = updateMeasurementDataObject(ctx, tx, request.Token, field, value, request.Data, measurementUpdateDependencies{ writeManualValueFunc: func(ctx context.Context, measurement *orm.Measurement, value float64) error { return redisChanges.AddMeasurementValueChange(ctx, measurement, value, false) @@ -86,7 +95,7 @@ func DataObjectAttributeUpdateHandler(c *gin.Context) { }, }) if err == nil && measurementResult.modeChanged { - err = redisChanges.AddDataObjectHashChange(ctx, dataObjectType, request.Token, "mode", measurementResult.mode) + err = redisChanges.AddHashChange(ctx, canonicalRedisKey, "mode", measurementResult.mode) } message = measurementResult.message default: diff --git a/handler/data_object_redis_change.go b/handler/data_object_redis_change.go index f8d1f3b..5ee2e43 100644 --- a/handler/data_object_redis_change.go +++ b/handler/data_object_redis_change.go @@ -3,14 +3,11 @@ package handler import ( "context" "encoding/json" - "errors" "fmt" "sort" "strconv" - "strings" "time" - "modelRT/constants" "modelRT/model" "modelRT/orm" @@ -46,70 +43,32 @@ func NewRedisChangeSet(client *redis.Client) *RedisChangeSet { return &RedisChangeSet{client: client} } -func (changes *RedisChangeSet) AddDataObjectHashChange( +func (changes *RedisChangeSet) AddHashChange( ctx context.Context, - dataObjectType constants.DataObjectType, - token, field string, + canonicalKey, field string, value any, ) error { if changes == nil || changes.client == nil { return fmt.Errorf("redis client is not initialized") } - - metadata, err := changes.client.HMGet(ctx, token, "id", "name").Result() - if err != nil { - return fmt.Errorf("query redis data-object aliases for %q: %w", token, err) - } - if len(metadata) != 2 || metadata[0] == nil || metadata[1] == nil { - return fmt.Errorf("redis data-object hash %q does not contain id and name", token) - } - - id, ok := metadata[0].(string) - if !ok { - return fmt.Errorf("redis data-object hash %q id has type %T", token, metadata[0]) - } - name, ok := metadata[1].(string) - if !ok { - return fmt.Errorf("redis data-object hash %q name has type %T", token, metadata[1]) - } - keys, err := dataObjectRedisAliasKeys(dataObjectType, id, name) - if err != nil { - return err - } - if !containsString(keys, token) { - return fmt.Errorf("redis data-object hash %q metadata points to different aliases", token) + if canonicalKey == "" { + return fmt.Errorf("canonical redis key is empty") } newValue, err := redisChangeString(value) if err != nil { return fmt.Errorf("encode redis data-object value: %w", err) } - for _, key := range keys { - keyType, err := changes.client.Type(ctx, key).Result() - if err != nil { - return fmt.Errorf("query redis key type for %q: %w", key, err) - } - if keyType == "none" && dataObjectType == constants.DataObjectTypeParameter && key == name && key != token { - // A parameter short alias is only initialized for a local station. - continue - } - if keyType != "hash" { - return fmt.Errorf("redis data-object key %q has type %q, expected hash", key, keyType) - } - oldValue, err := changes.client.HGet(ctx, key, field).Result() - if errors.Is(err, redis.Nil) { - return fmt.Errorf("redis data-object hash %q does not contain field %q", key, field) - } - if err != nil { - return fmt.Errorf("query redis hash %q field %q: %w", key, field, err) - } - changes.hashChanges = append(changes.hashChanges, redisHashChange{ - key: key, - field: field, - oldValue: oldValue, - newValue: newValue, - }) + oldValue, err := changes.client.HGet(ctx, canonicalKey, field).Result() + if err != nil { + return fmt.Errorf("query canonical redis hash %q field %q: %w", canonicalKey, field, err) } + changes.hashChanges = append(changes.hashChanges, redisHashChange{ + key: canonicalKey, + field: field, + oldValue: oldValue, + newValue: newValue, + }) return nil } @@ -314,24 +273,6 @@ func (changes *RedisChangeSet) keys() []string { return keys } -func dataObjectRedisAliasKeys(dataObjectType constants.DataObjectType, id, name string) ([]string, error) { - switch dataObjectType { - case constants.DataObjectTypeParameter: - if len(strings.Split(id, ".")) != 7 || len(strings.Split(name, ".")) != 4 { - return nil, fmt.Errorf("invalid parameter redis aliases id=%q name=%q", id, name) - } - return uniqueStrings(id, name), nil - case constants.DataObjectTypeMeasurement: - parts := strings.Split(id, ".") - if len(parts) != 7 || len(strings.Split(name, ".")) != 2 { - return nil, fmt.Errorf("invalid measurement redis aliases id=%q name=%q", id, name) - } - return uniqueStrings(id, strings.Join(parts[3:], "."), name), nil - default: - return nil, fmt.Errorf("unsupported data object type %q", dataObjectType) - } -} - func redisChangeString(value any) (string, error) { switch typedValue := value.(type) { case string: @@ -405,25 +346,3 @@ func equalRedisZValues(left, right []redis.Z) bool { } return true } - -func uniqueStrings(values ...string) []string { - seen := make(map[string]struct{}, len(values)) - result := make([]string, 0, len(values)) - for _, value := range values { - if _, ok := seen[value]; ok { - continue - } - seen[value] = struct{}{} - result = append(result, value) - } - return result -} - -func containsString(values []string, target string) bool { - for _, value := range values { - if value == target { - return true - } - } - return false -} diff --git a/handler/data_object_redis_change_test.go b/handler/data_object_redis_change_test.go index 3961e2e..a0c117e 100644 --- a/handler/data_object_redis_change_test.go +++ b/handler/data_object_redis_change_test.go @@ -3,46 +3,11 @@ package handler import ( "testing" - "modelRT/constants" - "github.com/redis/go-redis/v9" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) -func TestDataObjectRedisAliasKeys(t *testing.T) { - parameterKeys, err := dataObjectRedisAliasKeys( - constants.DataObjectTypeParameter, - "grid.zone.station.nspath.component.rated.voltage", - "nspath.component.rated.voltage", - ) - require.NoError(t, err) - assert.Equal(t, []string{ - "grid.zone.station.nspath.component.rated.voltage", - "nspath.component.rated.voltage", - }, parameterKeys) - - measurementKeys, err := dataObjectRedisAliasKeys( - constants.DataObjectTypeMeasurement, - "grid.zone.station.nspath.component.bay.current", - "nspath.current", - ) - require.NoError(t, err) - assert.Equal(t, []string{ - "grid.zone.station.nspath.component.bay.current", - "nspath.component.bay.current", - "nspath.current", - }, measurementKeys) -} - -func TestDataObjectRedisAliasKeysRejectsInvalidMetadata(t *testing.T) { - _, err := dataObjectRedisAliasKeys(constants.DataObjectTypeParameter, "short.id", "short.name") - require.Error(t, err) - - _, err = dataObjectRedisAliasKeys(constants.DataObjectTypeMeasurement, "short.id", "short.name") - require.Error(t, err) -} - func TestRedisChangeStringUsesCacheRepresentations(t *testing.T) { tests := []struct { name string diff --git a/model/data_object_redis_alias.go b/model/data_object_redis_alias.go new file mode 100644 index 0000000..3afd1a4 --- /dev/null +++ b/model/data_object_redis_alias.go @@ -0,0 +1,78 @@ +package model + +import ( + "context" + "errors" + "fmt" + + "modelRT/common" + "modelRT/constants" + + "github.com/redis/go-redis/v9" +) + +// DataObjectRedisAliasKey returns the Redis string key used to resolve any +// supported token form to the one canonical full-token hash. +func DataObjectRedisAliasKey(dataObjectType constants.DataObjectType, token string) (string, error) { + switch dataObjectType { + case constants.DataObjectTypeParameter: + return constants.RedisParameterDataObjectAliasPrefix + token, nil + case constants.DataObjectTypeMeasurement: + return constants.RedisMeasurementDataObjectAliasPrefix + token, nil + default: + return "", fmt.Errorf("unsupported data object type %q", dataObjectType) + } +} + +// ResolveDataObjectRedisKey resolves a full or short token to the canonical +// full-token Redis hash created during data-object initialization. +func ResolveDataObjectRedisKey( + ctx context.Context, + rdb *redis.Client, + dataObjectType constants.DataObjectType, + token string, +) (string, error) { + if rdb == nil { + return "", fmt.Errorf("redis client is not initialized") + } + classifiedType, err := ClassifyDataObjectToken(token) + if err != nil { + return "", err + } + if classifiedType != dataObjectType { + return "", fmt.Errorf("token %q is %q, expected %q", token, classifiedType, dataObjectType) + } + + aliasKey, err := DataObjectRedisAliasKey(dataObjectType, token) + if err != nil { + return "", err + } + canonicalKey, err := rdb.Get(ctx, aliasKey).Result() + if errors.Is(err, redis.Nil) { + switch dataObjectType { + case constants.DataObjectTypeParameter: + return "", fmt.Errorf("%w: %q", common.ErrParameterTokenNotFound, token) + case constants.DataObjectTypeMeasurement: + return "", fmt.Errorf("%w: %q", common.ErrMeasurementTokenNotFound, token) + } + } + if err != nil { + return "", fmt.Errorf("resolve redis data-object alias %q: %w", token, err) + } + if canonicalKey == "" { + return "", fmt.Errorf("redis data-object alias %q points to an empty key", token) + } + keyType, err := rdb.Type(ctx, canonicalKey).Result() + if err != nil { + return "", fmt.Errorf("query canonical redis key type for %q: %w", token, err) + } + if keyType != "hash" { + return "", fmt.Errorf( + "redis data-object alias %q points to key %q with type %q, expected hash", + token, + canonicalKey, + keyType, + ) + } + return canonicalKey, nil +} diff --git a/model/data_object_redis_alias_test.go b/model/data_object_redis_alias_test.go new file mode 100644 index 0000000..4a213c4 --- /dev/null +++ b/model/data_object_redis_alias_test.go @@ -0,0 +1,25 @@ +package model + +import ( + "testing" + + "modelRT/constants" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestDataObjectRedisAliasKey(t *testing.T) { + parameterKey, err := DataObjectRedisAliasKey(constants.DataObjectTypeParameter, "nspath.component.rated.voltage") + require.NoError(t, err) + assert.Equal(t, constants.RedisParameterDataObjectAliasPrefix+"nspath.component.rated.voltage", parameterKey) + + measurementKey, err := DataObjectRedisAliasKey(constants.DataObjectTypeMeasurement, "nspath.current") + require.NoError(t, err) + assert.Equal(t, constants.RedisMeasurementDataObjectAliasPrefix+"nspath.current", measurementKey) +} + +func TestDataObjectRedisAliasKeyRejectsUnsupportedType(t *testing.T) { + _, err := DataObjectRedisAliasKey(constants.DataObjectType("unknown"), "token") + require.Error(t, err) +} diff --git a/model/measurement_data_object_init.go b/model/measurement_data_object_init.go index e0823bb..768c98c 100644 --- a/model/measurement_data_object_init.go +++ b/model/measurement_data_object_init.go @@ -17,8 +17,9 @@ import ( const measurementDataObjectPipelineSize = 500 type measurementDataObjectHash struct { - Key string - Fields map[string]any + Key string + Aliases []string + Fields map[string]any } // MeasurementInitializationRecord contains a measurement and the hierarchy @@ -41,8 +42,8 @@ type MeasurementInitializationRecord struct { MeasurementBinding orm.JSONMap `gorm:"column:measurement_binding"` } -// InitializeMeasurementDataObjects creates the seven-part, four-part, and -// token4.token7 Redis hashes for measurements loaded from PostgreSQL. +// InitializeMeasurementDataObjects creates one seven-part Redis hash and +// aliases for all supported measurement token forms. func InitializeMeasurementDataObjects(ctx context.Context, records []MeasurementInitializationRecord) error { hashes, err := buildMeasurementDataObjectHashes(records) if err != nil { @@ -60,7 +61,7 @@ func InitializeMeasurementDataObjects(ctx context.Context, records []Measurement } func buildMeasurementDataObjectHashes(records []MeasurementInitializationRecord) ([]measurementDataObjectHash, error) { - hashes := make([]measurementDataObjectHash, 0, len(records)*3) + hashes := make([]measurementDataObjectHash, 0, len(records)) seenKeys := make(map[string]string, len(records)*3) for _, record := range records { if record.MeasurementMode != constants.MeasurementModeManual && @@ -103,7 +104,8 @@ func buildMeasurementDataObjectHashes(records []MeasurementInitializationRecord) }, ".") twoPartToken := record.ComponentNSPath + "." + record.MeasurementTag - for _, token := range []string{fullToken, fourPartToken, twoPartToken} { + aliases := []string{fullToken, fourPartToken, twoPartToken} + for _, token := range aliases { if err := validateInitializedMeasurementToken(token); err != nil { return nil, err } @@ -139,11 +141,18 @@ func buildMeasurementDataObjectHashes(records []MeasurementInitializationRecord) "binding": binding, } owner := fmt.Sprintf("%d/%s", record.MeasurementID, record.ComponentUUID) - for _, token := range []string{fullToken, fourPartToken, twoPartToken} { - if err := appendMeasurementDataObjectHash(&hashes, seenKeys, token, owner, fields); err != nil { - return nil, err + for _, alias := range aliases { + if existingOwner, exists := seenKeys[alias]; exists { + return nil, fmt.Errorf( + "ambiguous measurement token %q is produced by %q and %q", + alias, + existingOwner, + owner, + ) } + seenKeys[alias] = owner } + hashes = append(hashes, measurementDataObjectHash{Key: fullToken, Aliases: aliases, Fields: fields}) } return hashes, nil } @@ -159,26 +168,6 @@ func validateInitializedMeasurementToken(token string) error { return nil } -func appendMeasurementDataObjectHash( - hashes *[]measurementDataObjectHash, - seenKeys map[string]string, - key string, - owner string, - fields map[string]any, -) error { - if existingOwner, exists := seenKeys[key]; exists { - return fmt.Errorf( - "ambiguous measurement token %q is produced by %q and %q", - key, - existingOwner, - owner, - ) - } - seenKeys[key] = owner - *hashes = append(*hashes, measurementDataObjectHash{Key: key, Fields: fields}) - return nil -} - func measurementInitializationJSON(value orm.JSONMap) (string, error) { encoded, err := json.Marshal(value) if err != nil { @@ -200,20 +189,38 @@ func storeMeasurementDataObjectHashes( if err != nil { return fmt.Errorf("query previously initialized measurement keys: %w", err) } + oldAliasKeys, err := rdb.SMembers(ctx, constants.RedisMeasurementDataObjectAliasKeySet).Result() + if err != nil { + return fmt.Errorf("query previously initialized measurement alias keys: %w", err) + } currentKeys := make(map[string]struct{}, len(hashes)) + currentAliasKeys := make(map[string]struct{}, len(hashes)*3) for start := 0; start < len(hashes); start += measurementDataObjectPipelineSize { end := min(start+measurementDataObjectPipelineSize, len(hashes)) pipeline := rdb.TxPipeline() keyMembers := make([]any, 0, end-start) + aliasKeyMembers := make([]any, 0, (end-start)*3) for _, hash := range hashes[start:end] { pipeline.Del(ctx, hash.Key) pipeline.HSet(ctx, hash.Key, hash.Fields) keyMembers = append(keyMembers, hash.Key) currentKeys[hash.Key] = struct{}{} + for _, alias := range hash.Aliases { + aliasKey, err := DataObjectRedisAliasKey(constants.DataObjectTypeMeasurement, alias) + if err != nil { + return err + } + pipeline.Set(ctx, aliasKey, hash.Key, 0) + aliasKeyMembers = append(aliasKeyMembers, aliasKey) + currentAliasKeys[aliasKey] = struct{}{} + } } if len(keyMembers) > 0 { pipeline.SAdd(ctx, constants.RedisMeasurementDataObjectKeySet, keyMembers...) } + if len(aliasKeyMembers) > 0 { + pipeline.SAdd(ctx, constants.RedisMeasurementDataObjectAliasKeySet, aliasKeyMembers...) + } if _, err := pipeline.Exec(ctx); err != nil { return fmt.Errorf("write measurement data-object hash batch starting at %d: %w", start, err) } @@ -235,8 +242,24 @@ func storeMeasurementDataObjectHashes( } cleanupPipeline.SRem(ctx, constants.RedisMeasurementDataObjectKeySet, members...) } + staleAliasKeys := make([]string, 0) + for _, aliasKey := range oldAliasKeys { + if _, exists := currentAliasKeys[aliasKey]; !exists { + staleAliasKeys = append(staleAliasKeys, aliasKey) + } + } + for start := 0; start < len(staleAliasKeys); start += measurementDataObjectPipelineSize { + end := min(start+measurementDataObjectPipelineSize, len(staleAliasKeys)) + cleanupPipeline.Del(ctx, staleAliasKeys[start:end]...) + members := make([]any, 0, end-start) + for _, aliasKey := range staleAliasKeys[start:end] { + members = append(members, aliasKey) + } + cleanupPipeline.SRem(ctx, constants.RedisMeasurementDataObjectAliasKeySet, members...) + } if len(hashes) == 0 { cleanupPipeline.Del(ctx, constants.RedisMeasurementDataObjectKeySet) + cleanupPipeline.Del(ctx, constants.RedisMeasurementDataObjectAliasKeySet) } if _, err := cleanupPipeline.Exec(ctx); err != nil { return fmt.Errorf("remove stale measurement data-object hashes: %w", err) diff --git a/model/measurement_data_object_init_test.go b/model/measurement_data_object_init_test.go index e236c30..c5ffc5b 100644 --- a/model/measurement_data_object_init_test.go +++ b/model/measurement_data_object_init_test.go @@ -9,19 +9,18 @@ import ( "github.com/stretchr/testify/require" ) -func TestBuildMeasurementDataObjectHashesCreatesAllTokenForms(t *testing.T) { +func TestBuildMeasurementDataObjectHashesCreatesCanonicalHashAndAllAliases(t *testing.T) { record := measurementInitializationRecordForTest() hashes, err := buildMeasurementDataObjectHashes([]MeasurementInitializationRecord{record}) require.NoError(t, err) - require.Len(t, hashes, 3) + require.Len(t, hashes, 1) fullToken := "grid000.zone000.station000.220kV_xuefulu1.CTA.bay.IA_rms" fourPartToken := "220kV_xuefulu1.CTA.bay.IA_rms" twoPartToken := "220kV_xuefulu1.IA_rms" assert.Equal(t, fullToken, hashes[0].Key) - assert.Equal(t, fourPartToken, hashes[1].Key) - assert.Equal(t, twoPartToken, hashes[2].Key) + assert.Equal(t, []string{fullToken, fourPartToken, twoPartToken}, hashes[0].Aliases) fields := hashes[0].Fields assert.NotContains(t, fields, "value") @@ -35,8 +34,6 @@ func TestBuildMeasurementDataObjectHashesCreatesAllTokenForms(t *testing.T) { assert.Equal(t, `{"io_address":{"channel":"TM1","device":"CTA","dtype":1,"option":"rms","station":"001"},"type":1}`, fields["data_source"]) assert.Equal(t, `{}`, fields["event_plan"]) assert.Equal(t, `{"ct":{"index":0,"polarity":1,"ratio":1250}}`, fields["binding"]) - assert.Equal(t, fields, hashes[1].Fields) - assert.Equal(t, fields, hashes[2].Fields) } func TestBuildMeasurementDataObjectHashesRejectsAmbiguousShortToken(t *testing.T) { diff --git a/model/parameter_data_object_init.go b/model/parameter_data_object_init.go index 4819685..e27a7c9 100644 --- a/model/parameter_data_object_init.go +++ b/model/parameter_data_object_init.go @@ -18,8 +18,9 @@ import ( const parameterDataObjectPipelineSize = 500 type parameterDataObjectHash struct { - Key string - Fields map[string]any + Key string + Aliases []string + Fields map[string]any } // ParameterInitializationRecord contains one parameter attribute together with @@ -41,8 +42,8 @@ type ParameterInitializationRecord struct { DynamicRecordCount int64 `gorm:"column:dynamic_record_count"` } -// InitializeParameterDataObjects creates full and local-short Redis hashes for -// parameter attributes previously loaded from PostgreSQL. +// InitializeParameterDataObjects creates one full-token Redis hash and full or +// local-short token aliases for each parameter loaded from PostgreSQL. func InitializeParameterDataObjects(ctx context.Context, records []ParameterInitializationRecord) error { hashes, err := buildParameterDataObjectHashes(records) if err != nil { @@ -60,7 +61,7 @@ func InitializeParameterDataObjects(ctx context.Context, records []ParameterInit } func buildParameterDataObjectHashes(records []ParameterInitializationRecord) ([]parameterDataObjectHash, error) { - hashes := make([]parameterDataObjectHash, 0, len(records)*2) + hashes := make([]parameterDataObjectHash, 0, len(records)) seenKeys := make(map[string]string, len(records)*2) for _, record := range records { @@ -91,8 +92,9 @@ func buildParameterDataObjectHashes(records []ParameterInitializationRecord) ([] record.AttributeName, }, ".") - if err := validateInitializedParameterToken(fullToken); err != nil { - return nil, err + aliases := []string{fullToken} + if record.StationIsLocal { + aliases = append(aliases, shortToken) } fields := map[string]any{ "value": value, @@ -107,19 +109,21 @@ func buildParameterDataObjectHashes(records []ParameterInitializationRecord) ([] record.AttributeGroup, record.AttributeName, }, "/") - if err := appendParameterDataObjectHash(&hashes, seenKeys, fullToken, owner, fields); err != nil { - return nil, err - } - - if !record.StationIsLocal { - continue - } - if err := validateInitializedParameterToken(shortToken); err != nil { - return nil, err - } - if err := appendParameterDataObjectHash(&hashes, seenKeys, shortToken, owner, fields); err != nil { - return nil, err + for _, alias := range aliases { + if err := validateInitializedParameterToken(alias); err != nil { + return nil, err + } + if existingOwner, exists := seenKeys[alias]; exists { + return nil, fmt.Errorf( + "ambiguous parameter token %q is produced by %q and %q", + alias, + existingOwner, + owner, + ) + } + seenKeys[alias] = owner } + hashes = append(hashes, parameterDataObjectHash{Key: fullToken, Aliases: aliases, Fields: fields}) } return hashes, nil } @@ -135,26 +139,6 @@ func validateInitializedParameterToken(token string) error { return nil } -func appendParameterDataObjectHash( - hashes *[]parameterDataObjectHash, - seenKeys map[string]string, - key string, - owner string, - fields map[string]any, -) error { - if existingOwner, exists := seenKeys[key]; exists { - return fmt.Errorf( - "ambiguous parameter token %q is produced by %q and %q", - key, - existingOwner, - owner, - ) - } - seenKeys[key] = owner - *hashes = append(*hashes, parameterDataObjectHash{Key: key, Fields: fields}) - return nil -} - func parameterRedisValue(rawJSON string) (any, error) { decoder := json.NewDecoder(bytes.NewBufferString(rawJSON)) decoder.UseNumber() @@ -194,30 +178,86 @@ func storeParameterDataObjectHashes( if err != nil { return fmt.Errorf("query previously initialized parameter keys: %w", err) } - cleanupPipeline := rdb.TxPipeline() - for start := 0; start < len(oldKeys); start += parameterDataObjectPipelineSize { - end := min(start+parameterDataObjectPipelineSize, len(oldKeys)) - cleanupPipeline.Del(ctx, oldKeys[start:end]...) - } - cleanupPipeline.Del(ctx, constants.RedisParameterDataObjectKeySet) - if _, err := cleanupPipeline.Exec(ctx); err != nil { - return fmt.Errorf("remove stale parameter data-object hashes: %w", err) + oldAliasKeys, err := rdb.SMembers(ctx, constants.RedisParameterDataObjectAliasKeySet).Result() + if err != nil { + return fmt.Errorf("query previously initialized parameter alias keys: %w", err) } + currentKeys := make(map[string]struct{}, len(hashes)) + currentAliasKeys := make(map[string]struct{}, len(hashes)*2) for start := 0; start < len(hashes); start += parameterDataObjectPipelineSize { end := min(start+parameterDataObjectPipelineSize, len(hashes)) pipeline := rdb.TxPipeline() keyMembers := make([]any, 0, end-start) + aliasKeyMembers := make([]any, 0, (end-start)*2) for _, hash := range hashes[start:end] { + pipeline.Del(ctx, hash.Key) pipeline.HSet(ctx, hash.Key, hash.Fields) keyMembers = append(keyMembers, hash.Key) + currentKeys[hash.Key] = struct{}{} + for _, alias := range hash.Aliases { + aliasKey, err := DataObjectRedisAliasKey(constants.DataObjectTypeParameter, alias) + if err != nil { + return err + } + pipeline.Set(ctx, aliasKey, hash.Key, 0) + aliasKeyMembers = append(aliasKeyMembers, aliasKey) + currentAliasKeys[aliasKey] = struct{}{} + } } if len(keyMembers) > 0 { pipeline.SAdd(ctx, constants.RedisParameterDataObjectKeySet, keyMembers...) } + if len(aliasKeyMembers) > 0 { + pipeline.SAdd(ctx, constants.RedisParameterDataObjectAliasKeySet, aliasKeyMembers...) + } if _, err := pipeline.Exec(ctx); err != nil { return fmt.Errorf("write parameter data-object hash batch starting at %d: %w", start, err) } } + + if err := cleanupStaleParameterDataObjectKeys( + ctx, + rdb, + oldKeys, + oldAliasKeys, + currentKeys, + currentAliasKeys, + ); err != nil { + return err + } + return nil +} + +func cleanupStaleParameterDataObjectKeys( + ctx context.Context, + rdb *redis.Client, + oldKeys, oldAliasKeys []string, + currentKeys, currentAliasKeys map[string]struct{}, +) error { + pipeline := rdb.TxPipeline() + for _, key := range oldKeys { + if _, exists := currentKeys[key]; exists { + continue + } + pipeline.Del(ctx, key) + pipeline.SRem(ctx, constants.RedisParameterDataObjectKeySet, key) + } + for _, aliasKey := range oldAliasKeys { + if _, exists := currentAliasKeys[aliasKey]; exists { + continue + } + pipeline.Del(ctx, aliasKey) + pipeline.SRem(ctx, constants.RedisParameterDataObjectAliasKeySet, aliasKey) + } + if len(currentKeys) == 0 { + pipeline.Del(ctx, constants.RedisParameterDataObjectKeySet) + } + if len(currentAliasKeys) == 0 { + pipeline.Del(ctx, constants.RedisParameterDataObjectAliasKeySet) + } + if _, err := pipeline.Exec(ctx); err != nil { + return fmt.Errorf("remove stale parameter data-object keys: %w", err) + } return nil } diff --git a/model/parameter_data_object_init_test.go b/model/parameter_data_object_init_test.go index 0f3a600..4c9ac43 100644 --- a/model/parameter_data_object_init_test.go +++ b/model/parameter_data_object_init_test.go @@ -8,7 +8,7 @@ import ( "github.com/stretchr/testify/require" ) -func TestBuildParameterDataObjectHashesCreatesFullAndLocalKeys(t *testing.T) { +func TestBuildParameterDataObjectHashesCreatesCanonicalHashAndLocalAliases(t *testing.T) { records := []ParameterInitializationRecord{ { GridTag: "grid000", @@ -30,19 +30,18 @@ func TestBuildParameterDataObjectHashesCreatesFullAndLocalKeys(t *testing.T) { hashes, err := buildParameterDataObjectHashes(records) require.NoError(t, err) - require.Len(t, hashes, 2) + require.Len(t, hashes, 1) fullToken := "grid000.zone000.station000.220kV_xuefulu1.cable_22.base_extend.vnom_kv" shortToken := "220kV_xuefulu1.cable_22.base_extend.vnom_kv" assert.Equal(t, fullToken, hashes[0].Key) - assert.Equal(t, shortToken, hashes[1].Key) + assert.Equal(t, []string{fullToken, shortToken}, hashes[0].Aliases) assert.Equal(t, "7800.00", hashes[0].Fields["value"]) assert.Equal(t, "PARAM", hashes[0].Fields["meta"]) assert.Equal(t, "DOUBLE PRECISION", hashes[0].Fields["type"]) assert.Equal(t, shortToken, hashes[0].Fields["name"]) assert.Equal(t, "额定电压", hashes[0].Fields["description"]) assert.Equal(t, fullToken, hashes[0].Fields["id"]) - assert.Equal(t, hashes[0].Fields, hashes[1].Fields) } func TestBuildParameterDataObjectHashesSkipsShortKeyForNonLocalStation(t *testing.T) { @@ -69,6 +68,7 @@ func TestBuildParameterDataObjectHashesSkipsShortKeyForNonLocalStation(t *testin require.NoError(t, err) require.Len(t, hashes, 1) assert.Equal(t, "grid.zone.station.nspath.component.stable.attribute", hashes[0].Key) + assert.Equal(t, []string{"grid.zone.station.nspath.component.stable.attribute"}, hashes[0].Aliases) assert.Equal(t, true, hashes[0].Fields["value"]) }