modelRT/model/measurement_data_object_ini...

239 lines
7.9 KiB
Go

package model
import (
"context"
"encoding/json"
"fmt"
"strings"
"modelRT/constants"
"modelRT/diagram"
"modelRT/logger"
"modelRT/orm"
"github.com/redis/go-redis/v9"
)
const measurementDataObjectPipelineSize = 500
type measurementDataObjectHash struct {
Key string
Fields map[string]any
}
// MeasurementInitializationRecord contains a measurement and the hierarchy
// needed to create all supported Redis data-object token aliases.
type MeasurementInitializationRecord struct {
GridTag string `gorm:"column:grid_tag"`
ZoneTag string `gorm:"column:zone_tag"`
StationTag string `gorm:"column:station_tag"`
ComponentUUID string `gorm:"column:component_uuid"`
ComponentNSPath string `gorm:"column:component_nspath"`
ComponentTag string `gorm:"column:component_tag"`
MeasurementID int64 `gorm:"column:measurement_id"`
MeasurementTag string `gorm:"column:measurement_tag"`
MeasurementName string `gorm:"column:measurement_name"`
MeasurementType int16 `gorm:"column:measurement_type"`
MeasurementMode int16 `gorm:"column:measurement_mode"`
MeasurementSize int `gorm:"column:measurement_size"`
MeasurementDataSource orm.JSONMap `gorm:"column:measurement_data_source"`
MeasurementEventPlan orm.JSONMap `gorm:"column:measurement_event_plan"`
MeasurementBinding orm.JSONMap `gorm:"column:measurement_binding"`
}
// InitializeMeasurementDataObjects creates the seven-part, four-part, and
// token4.token7 Redis hashes for measurements loaded from PostgreSQL.
func InitializeMeasurementDataObjects(ctx context.Context, records []MeasurementInitializationRecord) error {
hashes, err := buildMeasurementDataObjectHashes(records)
if err != nil {
return fmt.Errorf("build measurement data-object hashes: %w", err)
}
if err := storeMeasurementDataObjectHashes(ctx, diagram.GetRedisClientInstance(), hashes); err != nil {
return fmt.Errorf("store measurement data-object hashes in redis: %w", err)
}
logger.Info(ctx, "initialize measurement data objects completed",
"postgres_record_count", len(records),
"redis_hash_count", len(hashes),
)
return nil
}
func buildMeasurementDataObjectHashes(records []MeasurementInitializationRecord) ([]measurementDataObjectHash, error) {
hashes := make([]measurementDataObjectHash, 0, len(records)*3)
seenKeys := make(map[string]string, len(records)*3)
for _, record := range records {
if record.MeasurementMode != constants.MeasurementModeManual &&
record.MeasurementMode != constants.MeasurementModeAutomatic {
return nil, fmt.Errorf(
"measurement %q mode must be %d or %d, got %d",
record.MeasurementTag,
constants.MeasurementModeManual,
constants.MeasurementModeAutomatic,
record.MeasurementMode,
)
}
if record.MeasurementDataSource == nil ||
record.MeasurementEventPlan == nil ||
record.MeasurementBinding == nil {
return nil, fmt.Errorf("measurement %q contains a null JSONB field", record.MeasurementTag)
}
fullToken := strings.Join([]string{
record.GridTag,
record.ZoneTag,
record.StationTag,
record.ComponentNSPath,
record.ComponentTag,
"bay",
record.MeasurementTag,
}, ".")
fourPartToken := strings.Join([]string{
record.ComponentNSPath,
record.ComponentTag,
"bay",
record.MeasurementTag,
}, ".")
twoPartToken := record.ComponentNSPath + "." + record.MeasurementTag
for _, token := range []string{fullToken, fourPartToken, twoPartToken} {
if err := validateInitializedMeasurementToken(token); err != nil {
return nil, err
}
}
measurementType, err := MeasurementTypeString(record.MeasurementType)
if err != nil {
return nil, fmt.Errorf("derive type for measurement %q: %w", fullToken, err)
}
dataSource, err := measurementInitializationJSON(record.MeasurementDataSource)
if err != nil {
return nil, fmt.Errorf("encode data_source for measurement %q: %w", fullToken, err)
}
eventPlan, err := measurementInitializationJSON(record.MeasurementEventPlan)
if err != nil {
return nil, fmt.Errorf("encode event_plan for measurement %q: %w", fullToken, err)
}
binding, err := measurementInitializationJSON(record.MeasurementBinding)
if err != nil {
return nil, fmt.Errorf("encode binding for measurement %q: %w", fullToken, err)
}
fields := map[string]any{
"mode": record.MeasurementMode,
"meta": "MEASUREMENT",
"type": measurementType,
"name": twoPartToken,
"description": record.MeasurementName,
"id": fullToken,
"size": record.MeasurementSize,
"data_source": dataSource,
"event_plan": eventPlan,
"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
}
}
}
return hashes, nil
}
func validateInitializedMeasurementToken(token string) error {
dataObjectType, err := ClassifyDataObjectToken(token)
if err != nil {
return fmt.Errorf("generated invalid measurement token %q: %w", token, err)
}
if dataObjectType != constants.DataObjectTypeMeasurement {
return fmt.Errorf("generated token %q is not a measurement", token)
}
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 {
return "", err
}
return string(encoded), nil
}
func storeMeasurementDataObjectHashes(
ctx context.Context,
rdb *redis.Client,
hashes []measurementDataObjectHash,
) error {
if rdb == nil {
return fmt.Errorf("redis client is nil")
}
oldKeys, err := rdb.SMembers(ctx, constants.RedisMeasurementDataObjectKeySet).Result()
if err != nil {
return fmt.Errorf("query previously initialized measurement keys: %w", err)
}
currentKeys := make(map[string]struct{}, len(hashes))
for start := 0; start < len(hashes); start += measurementDataObjectPipelineSize {
end := min(start+measurementDataObjectPipelineSize, len(hashes))
pipeline := rdb.TxPipeline()
keyMembers := make([]any, 0, end-start)
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{}{}
}
if len(keyMembers) > 0 {
pipeline.SAdd(ctx, constants.RedisMeasurementDataObjectKeySet, keyMembers...)
}
if _, err := pipeline.Exec(ctx); err != nil {
return fmt.Errorf("write measurement data-object hash batch starting at %d: %w", start, err)
}
}
staleKeys := make([]string, 0)
for _, key := range oldKeys {
if _, exists := currentKeys[key]; !exists {
staleKeys = append(staleKeys, key)
}
}
cleanupPipeline := rdb.TxPipeline()
for start := 0; start < len(staleKeys); start += measurementDataObjectPipelineSize {
end := min(start+measurementDataObjectPipelineSize, len(staleKeys))
cleanupPipeline.Del(ctx, staleKeys[start:end]...)
members := make([]any, 0, end-start)
for _, key := range staleKeys[start:end] {
members = append(members, key)
}
cleanupPipeline.SRem(ctx, constants.RedisMeasurementDataObjectKeySet, members...)
}
if len(hashes) == 0 {
cleanupPipeline.Del(ctx, constants.RedisMeasurementDataObjectKeySet)
}
if _, err := cleanupPipeline.Exec(ctx); err != nil {
return fmt.Errorf("remove stale measurement data-object hashes: %w", err)
}
return nil
}