From 2fbef3e1faa6226b35a7927d622f3712f749a61f Mon Sep 17 00:00:00 2001 From: douxu Date: Tue, 4 Aug 2026 17:12:17 +0800 Subject: [PATCH] feat: add manual measurement synchronization workflow - add configurable CL3611 and IEC104 manual sync clients - synchronize measurement mode and value updates with consistent timestamps - move Redis change sets into the repository layer - nest manual sync configuration under dataRT - improve startup transaction error handling and JSONB type mapping - update Kubernetes and MongoDB deployment documentation --- client/manualsync/client.go | 255 ++++++++++++++++++ client/manualsync/client_test.go | 240 +++++++++++++++++ client/manualsync/init.go | 33 +++ client/manualsync/init_test.go | 46 ++++ config/config.go | 22 +- config/config_test.go | 31 +++ constants/redis.go | 8 +- deploy/deploy.md | 96 ++++++- deploy/k8s/modelrt-configmap.yaml | 10 +- handler/data_object_attribute_update.go | 165 ++++++++---- handler/data_object_attribute_update_test.go | 97 +++++-- main.go | 65 ++--- model/measurement_data_object_init.go | 6 +- .../redis}/data_object_redis_change.go | 59 ++-- .../redis}/data_object_redis_change_test.go | 2 +- 15 files changed, 965 insertions(+), 170 deletions(-) create mode 100644 client/manualsync/client.go create mode 100644 client/manualsync/client_test.go create mode 100644 client/manualsync/init.go create mode 100644 client/manualsync/init_test.go create mode 100644 config/config_test.go rename {handler => repository/redis}/data_object_redis_change.go (86%) rename {handler => repository/redis}/data_object_redis_change_test.go (98%) diff --git a/client/manualsync/client.go b/client/manualsync/client.go new file mode 100644 index 0000000..cecf2ae --- /dev/null +++ b/client/manualsync/client.go @@ -0,0 +1,255 @@ +// Package manualsync synchronizes measurement manual-mode changes with the +// protocol service responsible for the measurement's data source +package manualsync + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "math" + "net/http" + "net/url" + "strconv" + "strings" + + "modelRT/config" + "modelRT/constants" + "modelRT/orm" +) + +const maxErrorResponseBody = 4 << 10 + +// SyntheticData is one manually supplied measurement value +type SyntheticData struct { + Time int64 `json:"time"` + Value float64 `json:"value"` +} + +// Target identifies a measurement in a downstream protocol service +type Target struct { + Type int `json:"type"` + Station string `json:"station"` + MainPos string `json:"main_pos"` + SubPos string `json:"sub_pos"` + Option string `json:"option"` +} + +// Request is the payload accepted by POST /api/manual +type Request struct { + Mode int16 `json:"mode"` + Data []SyntheticData `json:"data,omitempty"` + Target Target `json:"target"` +} + +// Syncer synchronizes a measurement mode or manual-value change +type Syncer interface { + Sync(context.Context, orm.JSONMap, int16, *SyntheticData) error +} + +// Client calls the protocol-specific manual synchronization endpoint +type Client struct { + httpClient *http.Client + protocolCL3611URL string + protocol104URL string +} + +// NewClient validates the configuration and constructs a reusable client +func NewClient(cfg config.ManualSyncConfig) (*Client, error) { + if cfg.Timeout <= 0 { + return nil, fmt.Errorf("manual sync timeout must be greater than zero") + } + protocolCL3611URL, err := endpointURL(cfg.ProtocolCL3611URL, cfg.APIPath) + if err != nil { + return nil, fmt.Errorf("invalid protocol CL3611 URL: %w", err) + } + protocol104URL, err := endpointURL(cfg.Protocol104URL, cfg.APIPath) + if err != nil { + return nil, fmt.Errorf("invalid protocol 104 URL: %w", err) + } + + return &Client{ + httpClient: &http.Client{Timeout: cfg.Timeout}, + protocolCL3611URL: protocolCL3611URL, + protocol104URL: protocol104URL, + }, nil +} + +// Sync posts one mode transition or manual-value update. Data is omitted for +// mode transitions and included only when sample is non-nil in manual mode +func (c *Client) Sync(ctx context.Context, dataSource orm.JSONMap, mode int16, data *SyntheticData) error { + if c == nil || c.httpClient == nil { + return fmt.Errorf("manual sync client is not initialized") + } + if mode != constants.MeasurementModeManual && mode != constants.MeasurementModeAutomatic { + return fmt.Errorf("manual sync mode must be 0 or 1, got %d", mode) + } + if mode == constants.MeasurementModeAutomatic && data != nil { + return fmt.Errorf("automatic mode manual sync request cannot contain data") + } + + endpoint, target, err := c.resolveTarget(dataSource) + if err != nil { + return err + } + requestPayload := Request{Mode: mode, Target: target} + if data != nil { + requestPayload.Data = []SyntheticData{*data} + } + body, err := json.Marshal(requestPayload) + if err != nil { + return fmt.Errorf("encode manual sync request: %w", err) + } + + request, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("create manual sync request: %w", err) + } + request.Header.Set("Content-Type", "application/json") + + response, err := c.httpClient.Do(request) + if err != nil { + return fmt.Errorf("call manual sync endpoint: %w", err) + } + defer response.Body.Close() + if response.StatusCode >= http.StatusOK && response.StatusCode < http.StatusMultipleChoices { + _, _ = io.Copy(io.Discard, response.Body) + return nil + } + + responseBody, readErr := io.ReadAll(io.LimitReader(response.Body, maxErrorResponseBody)) + if readErr != nil { + return fmt.Errorf("manual sync endpoint returned %s and response body could not be read: %w", response.Status, readErr) + } + message := strings.TrimSpace(string(responseBody)) + if message == "" { + return fmt.Errorf("manual sync endpoint returned %s", response.Status) + } + return fmt.Errorf("manual sync endpoint returned %s: %s", response.Status, message) +} + +type rawDataSource struct { + Type int `json:"type"` + IOAddress rawIOAddress `json:"io_address"` +} + +type rawIOAddress struct { + DType int `json:"dtype"` + Station string `json:"station"` + Device string `json:"device"` + Channel string `json:"channel"` + Option string `json:"option"` + Packet any `json:"packet"` + Offset any `json:"offset"` +} + +func (c *Client) resolveTarget(dataSource orm.JSONMap) (string, Target, error) { + if dataSource == nil { + return "", Target{}, fmt.Errorf("measurement data_source is null") + } + encoded, err := json.Marshal(dataSource) + if err != nil { + return "", Target{}, fmt.Errorf("encode measurement data_source: %w", err) + } + var source rawDataSource + if err := json.Unmarshal(encoded, &source); err != nil { + return "", Target{}, fmt.Errorf("decode measurement data_source: %w", err) + } + station := strings.TrimSpace(source.IOAddress.Station) + if station == "" { + return "", Target{}, fmt.Errorf("measurement data_source io_address.station is required") + } + + switch source.Type { + case 1: + device := strings.TrimSpace(source.IOAddress.Device) + channel := strings.TrimSpace(source.IOAddress.Channel) + if device == "" { + return "", Target{}, fmt.Errorf("CL3611 data_source io_address.device is required") + } + if channel == "" { + return "", Target{}, fmt.Errorf("CL3611 data_source io_address.channel is required") + } + target := Target{ + Station: station, + MainPos: device, + SubPos: channel, + } + switch source.IOAddress.DType { + case 1: + target.Type = 1 + target.Option = strings.TrimSpace(source.IOAddress.Option) + case 2: + target.Type = 2 + default: + return "", Target{}, fmt.Errorf("CL3611 data_source dtype must be 1 or 2, got %d", source.IOAddress.DType) + } + return c.protocolCL3611URL, target, nil + case 2: + packet, err := integerString(source.IOAddress.Packet) + if err != nil { + return "", Target{}, fmt.Errorf("104 data_source io_address.packet: %w", err) + } + offset, err := integerString(source.IOAddress.Offset) + if err != nil { + return "", Target{}, fmt.Errorf("104 data_source io_address.offset: %w", err) + } + return c.protocol104URL, Target{ + Type: 3, + Station: station, + MainPos: packet, + SubPos: offset, + Option: "", + }, nil + default: + return "", Target{}, fmt.Errorf("unsupported measurement data_source type %d", source.Type) + } +} + +func endpointURL(baseURL, apiPath string) (string, error) { + baseURL = strings.TrimSpace(baseURL) + if baseURL == "" { + return "", fmt.Errorf("base URL is required") + } + apiPath = strings.TrimSpace(apiPath) + if apiPath == "" { + return "", fmt.Errorf("API path is required") + } + endpoint := strings.TrimRight(baseURL, "/") + "/" + strings.TrimLeft(apiPath, "/") + parsed, err := url.Parse(endpoint) + if err != nil { + return "", err + } + if parsed.Scheme != "http" && parsed.Scheme != "https" { + return "", fmt.Errorf("URL scheme must be http or https") + } + if parsed.Host == "" { + return "", fmt.Errorf("URL host is required") + } + return parsed.String(), nil +} + +func integerString(value any) (string, error) { + switch typed := value.(type) { + case nil: + return "", fmt.Errorf("is required") + case float64: + if math.IsNaN(typed) || math.IsInf(typed, 0) || math.Trunc(typed) != typed { + return "", fmt.Errorf("must be an integer") + } + return strconv.FormatInt(int64(typed), 10), nil + case string: + trimmed := strings.TrimSpace(typed) + if trimmed == "" { + return "", fmt.Errorf("is required") + } + integer, err := strconv.ParseInt(trimmed, 10, 64) + if err != nil { + return "", fmt.Errorf("must be an integer: %w", err) + } + return strconv.FormatInt(integer, 10), nil + default: + return "", fmt.Errorf("has unsupported type %T", value) + } +} diff --git a/client/manualsync/client_test.go b/client/manualsync/client_test.go new file mode 100644 index 0000000..17afee9 --- /dev/null +++ b/client/manualsync/client_test.go @@ -0,0 +1,240 @@ +package manualsync + +import ( + "context" + "encoding/json" + "io" + "net/http" + "strings" + "testing" + "time" + + "modelRT/config" + "modelRT/constants" + "modelRT/orm" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type capturedRequest struct { + Path string + Method string + ContentType string + Payload Request +} + +func TestClientSyncRoutesAndMapsDataSources(t *testing.T) { + cl3611Requests := make(chan capturedRequest, 2) + protocol104Requests := make(chan capturedRequest, 1) + client, err := NewClient(config.ManualSyncConfig{ + ProtocolCL3611URL: "http://cl3611.test", + Protocol104URL: "http://protocol104.test", + APIPath: "/api/manual", + Timeout: time.Second, + }) + require.NoError(t, err) + client.httpClient.Transport = roundTripFunc(func(request *http.Request) (*http.Response, error) { + var payload Request + require.NoError(t, json.NewDecoder(request.Body).Decode(&payload)) + captured := capturedRequest{ + Path: request.URL.Path, + Method: request.Method, + ContentType: request.Header.Get("Content-Type"), + Payload: payload, + } + if request.URL.Host == "cl3611.test" { + cl3611Requests <- captured + } else { + protocol104Requests <- captured + } + return httpResponse(http.StatusNoContent, ""), nil + }) + + tests := []struct { + name string + dataSource orm.JSONMap + requests <-chan capturedRequest + wantTarget Target + }{ + { + name: "CL3611 phasor", + dataSource: orm.JSONMap{ + "type": 1, + "io_address": map[string]any{ + "dtype": 1, "station": "001", "device": "ssu001", + "channel": "TM1", "option": "RMS", + }, + }, + requests: cl3611Requests, + wantTarget: Target{ + Type: 1, Station: "001", MainPos: "ssu001", SubPos: "TM1", Option: "RMS", + }, + }, + { + name: "CL3611 sample", + dataSource: orm.JSONMap{ + "type": 1, + "io_address": map[string]any{ + "dtype": 2, "station": "002", "device": "ssu002", + "channel": "TS01", "option": "ignored", + }, + }, + requests: cl3611Requests, + wantTarget: Target{ + Type: 2, Station: "002", MainPos: "ssu002", SubPos: "TS01", Option: "", + }, + }, + { + name: "104", + dataSource: orm.JSONMap{ + "type": 2, + "io_address": map[string]any{ + "station": "station000", "packet": 10, "offset": 35, + }, + }, + requests: protocol104Requests, + wantTarget: Target{ + Type: 3, Station: "station000", MainPos: "10", SubPos: "35", Option: "", + }, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + err := client.Sync(context.Background(), test.dataSource, constants.MeasurementModeAutomatic, nil) + require.NoError(t, err) + captured := <-test.requests + assert.Equal(t, http.MethodPost, captured.Method) + assert.Equal(t, "/api/manual", captured.Path) + assert.Equal(t, "application/json", captured.ContentType) + assert.Equal(t, constants.MeasurementModeAutomatic, captured.Payload.Mode) + assert.Nil(t, captured.Payload.Data) + assert.Equal(t, test.wantTarget, captured.Payload.Target) + }) + } +} + +func TestClientSyncIncludesManualValueData(t *testing.T) { + requests := make(chan capturedRequest, 1) + client, err := NewClient(config.ManualSyncConfig{ + ProtocolCL3611URL: "http://cl3611.test", + Protocol104URL: "http://protocol104.test", + APIPath: "api/manual", + Timeout: time.Second, + }) + require.NoError(t, err) + client.httpClient.Transport = captureTransport(t, requests) + + sample := SyntheticData{Time: 1736305467506000000, Value: 1.25} + err = client.Sync(context.Background(), orm.JSONMap{ + "type": 2, + "io_address": map[string]any{ + "station": "station000", "packet": 1, "offset": 2, + }, + }, constants.MeasurementModeManual, &sample) + require.NoError(t, err) + captured := <-requests + assert.Equal(t, constants.MeasurementModeManual, captured.Payload.Mode) + assert.Equal(t, []SyntheticData{sample}, captured.Payload.Data) +} + +func TestClientSyncRejectsInvalidDataSources(t *testing.T) { + client, err := NewClient(config.ManualSyncConfig{ + ProtocolCL3611URL: "http://cl3611.test", + Protocol104URL: "http://protocol104.test", + APIPath: "/api/manual", + Timeout: time.Second, + }) + require.NoError(t, err) + client.httpClient.Transport = roundTripFunc(func(*http.Request) (*http.Response, error) { + t.Fatal("invalid data source must not call endpoint") + return nil, nil + }) + + tests := []struct { + name string + dataSource orm.JSONMap + wantError string + }{ + {name: "unsupported type", dataSource: orm.JSONMap{"type": 3, "io_address": map[string]any{"station": "s"}}, wantError: "unsupported"}, + {name: "invalid dtype", dataSource: orm.JSONMap{"type": 1, "io_address": map[string]any{"dtype": 3, "station": "s", "device": "d", "channel": "c"}}, wantError: "dtype"}, + {name: "missing station", dataSource: orm.JSONMap{"type": 2, "io_address": map[string]any{"packet": 1, "offset": 2}}, wantError: "station"}, + {name: "fractional packet", dataSource: orm.JSONMap{"type": 2, "io_address": map[string]any{"station": "s", "packet": 1.5, "offset": 2}}, wantError: "packet"}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + err := client.Sync(context.Background(), test.dataSource, constants.MeasurementModeManual, nil) + require.Error(t, err) + assert.Contains(t, err.Error(), test.wantError) + }) + } +} + +func TestClientSyncReturnsNonSuccessResponse(t *testing.T) { + client, err := NewClient(config.ManualSyncConfig{ + ProtocolCL3611URL: "http://cl3611.test", + Protocol104URL: "http://protocol104.test", + APIPath: "/api/manual", + Timeout: time.Second, + }) + require.NoError(t, err) + client.httpClient.Transport = roundTripFunc(func(*http.Request) (*http.Response, error) { + return httpResponse(http.StatusServiceUnavailable, "downstream unavailable"), nil + }) + + err = client.Sync(context.Background(), orm.JSONMap{ + "type": 2, + "io_address": map[string]any{ + "station": "s", "packet": 1, "offset": 2, + }, + }, constants.MeasurementModeManual, nil) + require.Error(t, err) + assert.Contains(t, err.Error(), "Service Unavailable") + assert.Contains(t, err.Error(), "downstream unavailable") +} + +func TestNewClientValidatesConfiguration(t *testing.T) { + _, err := NewClient(config.ManualSyncConfig{}) + require.Error(t, err) + assert.Contains(t, err.Error(), "timeout") + + _, err = NewClient(config.ManualSyncConfig{ + ProtocolCL3611URL: "127.0.0.1:9001", + Protocol104URL: "http://127.0.0.1:9002", + APIPath: "/api/manual", + Timeout: time.Second, + }) + require.Error(t, err) + assert.Contains(t, err.Error(), "invalid protocol CL3611 URL") +} + +type roundTripFunc func(*http.Request) (*http.Response, error) + +func (roundTrip roundTripFunc) RoundTrip(request *http.Request) (*http.Response, error) { + return roundTrip(request) +} + +func captureTransport(t *testing.T, requests chan<- capturedRequest) http.RoundTripper { + t.Helper() + return roundTripFunc(func(request *http.Request) (*http.Response, error) { + var payload Request + require.NoError(t, json.NewDecoder(request.Body).Decode(&payload)) + requests <- capturedRequest{ + Path: request.URL.Path, + Method: request.Method, + ContentType: request.Header.Get("Content-Type"), + Payload: payload, + } + return httpResponse(http.StatusNoContent, ""), nil + }) +} + +func httpResponse(status int, body string) *http.Response { + return &http.Response{ + StatusCode: status, + Status: http.StatusText(status), + Body: io.NopCloser(strings.NewReader(body)), + Header: make(http.Header), + } +} diff --git a/client/manualsync/init.go b/client/manualsync/init.go new file mode 100644 index 0000000..facb4bd --- /dev/null +++ b/client/manualsync/init.go @@ -0,0 +1,33 @@ +package manualsync + +import ( + "context" + "fmt" + "sync" + + "modelRT/orm" +) + +var ( + defaultSyncerMu sync.RWMutex + defaultSyncer Syncer +) + +// SetDefaultSyncer sets the process-wide manual measurement syncer. +// Passing nil clears the current default. +func SetDefaultSyncer(syncer Syncer) { + defaultSyncerMu.Lock() + defer defaultSyncerMu.Unlock() + defaultSyncer = syncer +} + +// Sync uses the process-wide manual measurement syncer. +func Sync(ctx context.Context, dataSource orm.JSONMap, mode int16, sample *SyntheticData) error { + defaultSyncerMu.RLock() + syncer := defaultSyncer + defaultSyncerMu.RUnlock() + if syncer == nil { + return fmt.Errorf("manual measurement sync client is not initialized") + } + return syncer.Sync(ctx, dataSource, mode, sample) +} diff --git a/client/manualsync/init_test.go b/client/manualsync/init_test.go new file mode 100644 index 0000000..98f0996 --- /dev/null +++ b/client/manualsync/init_test.go @@ -0,0 +1,46 @@ +package manualsync + +import ( + "context" + "testing" + + "modelRT/constants" + "modelRT/orm" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +type syncerFunc func(context.Context, orm.JSONMap, int16, *SyntheticData) error + +func (syncer syncerFunc) Sync(ctx context.Context, dataSource orm.JSONMap, mode int16, sample *SyntheticData) error { + return syncer(ctx, dataSource, mode, sample) +} + +func TestSyncRequiresDefaultSyncer(t *testing.T) { + SetDefaultSyncer(nil) + t.Cleanup(func() { SetDefaultSyncer(nil) }) + + err := Sync(context.Background(), nil, constants.MeasurementModeManual, nil) + require.Error(t, err) + assert.Contains(t, err.Error(), "not initialized") +} + +func TestSyncUsesDefaultSyncer(t *testing.T) { + SetDefaultSyncer(nil) + t.Cleanup(func() { SetDefaultSyncer(nil) }) + + dataSource := orm.JSONMap{"type": 2} + sample := &SyntheticData{Time: 123, Value: 4.5} + called := false + SetDefaultSyncer(syncerFunc(func(_ context.Context, gotDataSource orm.JSONMap, gotMode int16, gotSample *SyntheticData) error { + called = true + assert.Equal(t, dataSource, gotDataSource) + assert.Equal(t, constants.MeasurementModeManual, gotMode) + assert.Same(t, sample, gotSample) + return nil + })) + + require.NoError(t, Sync(context.Background(), dataSource, constants.MeasurementModeManual, sample)) + assert.True(t, called) +} diff --git a/config/config.go b/config/config.go index a0ddc7f..3f6d0d6 100644 --- a/config/config.go +++ b/config/config.go @@ -91,12 +91,18 @@ type AntsConfig struct { RTDReceiveConcurrentQuantity int `mapstructure:"rtd_receive_concurrent_quantity"` // polling real time data concurrent quantity } -// DataRTConfig define config struct of data runtime server api config +// ManualSyncConfig defines protocol endpoints used to synchronize manual +// measurement mode and value changes. +type ManualSyncConfig struct { + ProtocolCL3611URL string `mapstructure:"protocol_cl3611_url"` + Protocol104URL string `mapstructure:"protocol_104_url"` + APIPath string `mapstructure:"api_path"` + Timeout time.Duration `mapstructure:"timeout"` +} + +// DataRTConfig defines APIs provided by dataRT. type DataRTConfig struct { - Host string `mapstructure:"host"` - Port int64 `mapstructure:"port"` - PollingAPI string `mapstructure:"polling_api"` - Method string `mapstructure:"polling_api_method"` + ManualSync ManualSyncConfig `mapstructure:"manual_sync"` } // OtelConfig define config struct of OpenTelemetry tracing @@ -124,9 +130,9 @@ type ModelRTConfig struct { KafkaConfig `mapstructure:"kafka"` LoggerConfig `mapstructure:"logger"` AntsConfig `mapstructure:"ants"` - DataRTConfig `mapstructure:"dataRT"` - LockerRedisConfig RedisConfig `mapstructure:"locker_redis"` - StorageRedisConfig RedisConfig `mapstructure:"storage_redis"` + DataRTConfig DataRTConfig `mapstructure:"dataRT"` + LockerRedisConfig RedisConfig `mapstructure:"locker_redis"` + StorageRedisConfig RedisConfig `mapstructure:"storage_redis"` AsyncTaskConfig AsyncTaskConfig `mapstructure:"async_task"` OtelConfig OtelConfig `mapstructure:"otel"` PostgresDBURI string `mapstructure:"-"` diff --git a/config/config_test.go b/config/config_test.go new file mode 100644 index 0000000..c1e4559 --- /dev/null +++ b/config/config_test.go @@ -0,0 +1,31 @@ +package config + +import ( + "os" + "path/filepath" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestReadAndInitConfigReadsNestedDataRTManualSync(t *testing.T) { + configPath := filepath.Join(t.TempDir(), "config.yaml") + contents := []byte(` +dataRT: + manual_sync: + protocol_cl3611_url: "http://127.0.0.1:9001" + protocol_104_url: "http://127.0.0.1:9002" + api_path: "/api/manual" + timeout: 3s +`) + require.NoError(t, os.WriteFile(configPath, contents, 0o600)) + + cfg := ReadAndInitConfig(filepath.Dir(configPath), "config", "yaml") + + assert.Equal(t, "http://127.0.0.1:9001", cfg.DataRTConfig.ManualSync.ProtocolCL3611URL) + assert.Equal(t, "http://127.0.0.1:9002", cfg.DataRTConfig.ManualSync.Protocol104URL) + assert.Equal(t, "/api/manual", cfg.DataRTConfig.ManualSync.APIPath) + assert.Equal(t, 3*time.Second, cfg.DataRTConfig.ManualSync.Timeout) +} diff --git a/constants/redis.go b/constants/redis.go index ee317af..b4c3a33 100644 --- a/constants/redis.go +++ b/constants/redis.go @@ -9,16 +9,16 @@ const ( // startup so stale parameter data-object keys can be removed safely. RedisParameterDataObjectKeySet = "modelrt:parameter-data-object:keys" + // RedisMeasurementDataObjectKeySet tracks measurement hashes created during + // startup so stale measurement data-object keys can be removed safely. + RedisMeasurementDataObjectKeySet = "modelrt:measurement-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" diff --git a/deploy/deploy.md b/deploy/deploy.md index e8f2e15..2e1678a 100644 --- a/deploy/deploy.md +++ b/deploy/deploy.md @@ -439,10 +439,10 @@ go run deploy/redis-test-data/measurments-recommend/measurement_injection.go | | `station_id` | 项目所操作的默认变电站 `ID`。 | `1` | | **Service Config** | `service_name` | 服务名称,用于日志、监控等标识。 | `"modelRT"` | | | `secret_key` | 服务内部使用的秘钥,用于签名或认证。 | `"modelrt_key"` | -| **DataRT API** | `host` | 外部 `DataRT` 服务的主机地址。 | `"http://127.0.0.1"` | -| | `port` | `DataRT` 服务的端口号。 | `8888` | -| | `polling_api` | 轮询数据的 `API` 路径。 | `"datart/getPointData"` | -| | `polling_api_method` | 调用该 `API` 使用的 `HTTP` 方法。 | `"GET"` | +| **DataRT Manual Sync** | `manual_sync.protocol_cl3611_url` | CL3611 协议服务地址。 | `"http://127.0.0.1:9001"` | +| | `manual_sync.protocol_104_url` | IEC 60870-5-104 协议服务地址。 | `"http://127.0.0.1:9002"` | +| | `manual_sync.api_path` | 手动测量值及模式同步 API 路径。 | `"/api/manual"` | +| | `manual_sync.timeout` | 同步请求超时时间。 | `"3s"` | #### 3.2 编译 ModelRT 服务 @@ -759,6 +759,94 @@ kubectl delete -f deploy/k8s/pg-service.yaml \ -f deploy/k8s/pg-configmap.yaml ``` +#### 4.5 部署 MongoDB 并创建应用用户 + +使用以下清单部署 MongoDB: + +```bash +kubectl apply -f deploy/k8s/mongodb-secret.yaml +kubectl apply -f deploy/k8s/mongodb-pvc.yaml +kubectl apply -f deploy/k8s/mongodb-statefulset.yaml +kubectl apply -f deploy/k8s/mongodb-service.yaml +``` + +等待 MongoDB Pod 就绪: + +```bash +kubectl wait --for=condition=ready pod/mongodb-0 --timeout=180s +``` + +MongoDB 首次初始化时会根据 `mongodb-secret.yaml` 创建 `admin` 管理员。Pod 就绪后,以管理员身份在 `admin` 认证库中创建应用用户 `coslight`,并授予其对 `eventdb` 的读写和数据库管理权限: + +```bash +kubectl exec mongodb-0 -- mongosh \ + -u admin \ + -p coslight \ + --authenticationDatabase admin \ + --quiet \ + --eval ' +const adminDb = db.getSiblingDB("admin"); +adminDb.createUser({ + user: "coslight", + pwd: "coslight", + roles: [ + { role: "readWrite", db: "eventdb" }, + { role: "dbAdmin", db: "eventdb" } + ] +}); +' +``` + +> **注意:** `use admin` 是 `mongosh` 的交互式命令,不应在 `--eval` 脚本中使用。这里通过 `db.getSiblingDB("admin")` 明确指定用户所属的认证库。用户存储在 `admin` 库中,因此应用连接时必须将认证库配置为 `admin`。 + +检查用户是否创建成功: + +```bash +kubectl exec mongodb-0 -- mongosh \ + -u admin \ + -p coslight \ + --authenticationDatabase admin \ + --quiet \ + --eval 'printjson(db.getSiblingDB("admin").getUser("coslight"));' +``` + +使用 `coslight` 用户连接并验证 `eventdb` 权限: + +```bash +kubectl exec mongodb-0 -- mongosh \ + -u coslight \ + -p coslight \ + --authenticationDatabase admin \ + --quiet \ + --eval ' +const eventDb = db.getSiblingDB("eventdb"); +eventDb.__permission_check.insertOne({ checkedAt: new Date() }); +eventDb.__permission_check.deleteMany({}); +printjson({ ok: 1, database: eventDb.getName() }); +' +``` + +如果 `coslight` 用户已经存在,重复执行 `createUser` 会返回 `User already exists`。需要重置密码或修正角色时,使用: + +```bash +kubectl exec mongodb-0 -- mongosh \ + -u admin \ + -p coslight \ + --authenticationDatabase admin \ + --quiet \ + --eval ' +db.getSiblingDB("admin").updateUser("coslight", { + pwd: "coslight", + roles: [ + { role: "readWrite", db: "eventdb" }, + { role: "dbAdmin", db: "eventdb" } + ] +}); +' +``` + +> **安全提示:** 示例使用仓库当前的测试密码。生产环境应修改管理员和应用用户密码,并避免在命令行或版本库中保存明文凭据。 + ### 5\. 部署 ModelRT(Kubernetes) 所有资源部署在 `default` 命名空间,YAML 文件位于 `deploy/k8s/`。 diff --git a/deploy/k8s/modelrt-configmap.yaml b/deploy/k8s/modelrt-configmap.yaml index 39f2b59..c32ddaf 100644 --- a/deploy/k8s/modelrt-configmap.yaml +++ b/deploy/k8s/modelrt-configmap.yaml @@ -80,7 +80,9 @@ data: deploy_env: "development" dataRT: - host: "http://127.0.0.1" - port: 8888 - polling_api: "datart/getPointData" - polling_api_method: "GET" + # manual measurement synchronization endpoints + manual_sync: + protocol_cl3611_url: "http://protocol-cl3611-service:9001" + protocol_104_url: "http://protocol-104-service:9002" + api_path: "/api/manual" + timeout: 3s diff --git a/handler/data_object_attribute_update.go b/handler/data_object_attribute_update.go index 71b3c64..36df8d4 100644 --- a/handler/data_object_attribute_update.go +++ b/handler/data_object_attribute_update.go @@ -11,6 +11,7 @@ import ( "strings" "time" + "modelRT/client/manualsync" "modelRT/common" "modelRT/common/errcode" "modelRT/constants" @@ -19,6 +20,7 @@ import ( "modelRT/logger" "modelRT/model" "modelRT/orm" + redisrepository "modelRT/repository/redis" "github.com/gin-gonic/gin" "gorm.io/gorm" @@ -31,6 +33,8 @@ type dataObjectAttributeUpdateRequest struct { Data json.RawMessage `json:"data,omitempty"` } +const redisChangeRestoreTimeout = 5 * time.Second + // DataObjectAttributeUpdateHandler updates the writable field of one data object. func DataObjectAttributeUpdateHandler(c *gin.Context) { ctx := c.Request.Context() @@ -62,7 +66,7 @@ func DataObjectAttributeUpdateHandler(c *gin.Context) { }() redisClient := diagram.GetRedisClientInstance() - redisChanges := NewRedisChangeSet(redisClient) + redisChanges := redisrepository.NewRedisChangeSet(redisClient) canonicalRedisKey, err := model.ResolveDataObjectRedisKey( ctx, redisClient, @@ -85,14 +89,22 @@ func DataObjectAttributeUpdateHandler(c *gin.Context) { err = queryErr 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) + 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) }, - updateDataRTFunc: callRealTimeDataWriteStopInterface, - startDataRTFunc: callRealTimeDataWriteStartInterface, - replaceRedisValueFunc: func(ctx context.Context, measurement *orm.Measurement, value float64) error { - return redisChanges.AddMeasurementValueChange(ctx, measurement, value, true) + 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) }, + nowFunc: time.Now, }) if err == nil && measurementResult.modeChanged { err = redisChanges.AddHashChange(ctx, canonicalRedisKey, "mode", measurementResult.mode) @@ -105,7 +117,7 @@ func DataObjectAttributeUpdateHandler(c *gin.Context) { if err != nil { _ = tx.Rollback().Error if measurementResult.recordFailureOnError { - if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, time.Now().UTC()); logErr != nil { + if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, measurementFailureTime(measurementResult)); logErr != nil { logger.Error(ctx, "append failed measurement value operation failed", "measurement_id", measurementResult.measurementID, "error", logErr) } } @@ -121,7 +133,7 @@ func DataObjectAttributeUpdateHandler(c *gin.Context) { if err := redisChanges.Apply(ctx); err != nil { _ = tx.Rollback().Error if measurementResult.recordFailureOnError { - if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, time.Now().UTC()); logErr != nil { + if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, measurementFailureTime(measurementResult)); logErr != nil { logger.Error(ctx, "append failed measurement value operation failed", "measurement_id", measurementResult.measurementID, "error", logErr) } } @@ -137,7 +149,7 @@ func DataObjectAttributeUpdateHandler(c *gin.Context) { logger.Error(ctx, "revert redis data-object changes failed", "token", request.Token, "field", field, "error", redisErr) } if measurementResult.recordFailureOnError { - if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, time.Now().UTC()); logErr != nil { + if logErr := database.AppendMeasurementValueOperation(ctx, database.GetPostgresDBClient(), measurementResult.measurementID, 1, measurementResult.value, measurementFailureTime(measurementResult)); logErr != nil { logger.Error(ctx, "append failed measurement value operation failed", "measurement_id", measurementResult.measurementID, "error", logErr) } } @@ -261,17 +273,17 @@ func parseMeasurementUpdateMode(raw json.RawMessage) (int16, error) { return mode, nil } -type measurementManualValueWriter func(context.Context, *orm.Measurement, float64) error +type measurementManualValueWriter func(context.Context, *orm.Measurement, float64, time.Time) error -type measurementDataRTUpdater func(context.Context, orm.JSONMap, *float64) error +type measurementManualChangeSyncer func(context.Context, orm.JSONMap, int16, *manualsync.SyntheticData) error -type measurementRedisValueReplacer func(context.Context, *orm.Measurement, float64) error +type measurementRedisValueReplacer func(context.Context, *orm.Measurement, float64, time.Time) error type measurementUpdateDependencies struct { writeManualValueFunc measurementManualValueWriter - updateDataRTFunc measurementDataRTUpdater - startDataRTFunc measurementDataRTUpdater + syncManualChangeFunc measurementManualChangeSyncer replaceRedisValueFunc measurementRedisValueReplacer + nowFunc func() time.Time } type measurementUpdateResult struct { @@ -281,6 +293,7 @@ type measurementUpdateResult struct { recordFailureOnError bool mode int16 modeChanged bool + operationTime time.Time } func updateMeasurementDataObject( @@ -315,6 +328,7 @@ func updateMeasurementDataObject( if currentMode == targetAutomatic { return measurementUpdateResult{message: fmt.Sprintf("measurement is already in %s mode", measurementModeName(mode))}, nil } + operationTime := measurementUpdateNow(dependencies) var manualValue *float64 if currentMode && mode == constants.MeasurementModeManual { manualValue, err = parseOptionalMeasurementModeData(modeData) @@ -322,37 +336,32 @@ func updateMeasurementDataObject( return measurementUpdateResult{}, err } } - if err := database.UpdateMeasurementModeWithOperation(ctx, tx, lockedMeasurement.ID, mode, time.Now().UTC()); err != nil { + if err := database.UpdateMeasurementModeWithOperation(ctx, tx, lockedMeasurement.ID, mode, operationTime); err != nil { return measurementUpdateResult{}, err } + syncManualMeasurementChange( + ctx, + dependencies.syncManualChangeFunc, + lockedMeasurement.ID, + lockedMeasurement.DataSource, + mode, + nil, + ) if currentMode && mode == constants.MeasurementModeManual { - if dependencies.updateDataRTFunc == nil { - return measurementUpdateResult{}, fmt.Errorf("measurement dataRT updater is nil") - } - if err := dependencies.updateDataRTFunc(ctx, lockedMeasurement.DataSource, nil); err != nil { - return measurementUpdateResult{}, fmt.Errorf("stop automatic measurement write to dataRT: %w", err) - } if manualValue != nil { if dependencies.replaceRedisValueFunc == nil { return measurementUpdateResult{}, fmt.Errorf("measurement redis value replacer is nil") } - if err := dependencies.replaceRedisValueFunc(ctx, &lockedMeasurement, *manualValue); err != nil { + if err := dependencies.replaceRedisValueFunc(ctx, &lockedMeasurement, *manualValue, operationTime); err != nil { return measurementUpdateResult{}, fmt.Errorf("replace measurement redis value: %w", err) } } } - if !currentMode && mode == constants.MeasurementModeAutomatic { - if dependencies.startDataRTFunc == nil { - return measurementUpdateResult{}, fmt.Errorf("measurement dataRT starter is nil") - } - if err := dependencies.startDataRTFunc(ctx, lockedMeasurement.DataSource, nil); err != nil { - return measurementUpdateResult{}, fmt.Errorf("start automatic measurement write to dataRT: %w", err) - } - } return measurementUpdateResult{ - message: fmt.Sprintf("measurement mode changed to %s", measurementModeName(mode)), - mode: mode, - modeChanged: true, + message: fmt.Sprintf("measurement mode changed to %s", measurementModeName(mode)), + mode: mode, + modeChanged: true, + operationTime: operationTime, }, nil case "value": currentMode, err := measurementModeIsAutomatic(lockedMeasurement.Mode) @@ -366,24 +375,29 @@ func updateMeasurementDataObject( if !ok { return measurementUpdateResult{}, fmt.Errorf("measurement value has invalid type %T", value) } + operationTime := measurementUpdateNow(dependencies) failureResult := measurementUpdateResult{ measurementID: lockedMeasurement.ID, value: manualValue, recordFailureOnError: true, + operationTime: operationTime, } if dependencies.writeManualValueFunc == nil { return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(fmt.Errorf("measurement manual value writer is nil")) } - if err := dependencies.writeManualValueFunc(ctx, &lockedMeasurement, manualValue); err != nil { + if err := dependencies.writeManualValueFunc(ctx, &lockedMeasurement, manualValue, operationTime); err != nil { return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(err) } - if dependencies.updateDataRTFunc == nil { - return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(fmt.Errorf("measurement dataRT updater is nil")) - } - if err := dependencies.updateDataRTFunc(ctx, lockedMeasurement.DataSource, &manualValue); err != nil { - return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(err) - } - if err := database.AppendMeasurementValueOperation(ctx, tx, lockedMeasurement.ID, 0, manualValue, time.Now().UTC()); err != nil { + 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 { return failureResult, errcode.ErrMeasurementValueUpdateFailed.WithCause(err) } failureResult.message = "measurement manual value updated" @@ -393,6 +407,62 @@ func updateMeasurementDataObject( } } +// 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() +} + func parseOptionalMeasurementModeData(raw json.RawMessage) (*float64, error) { trimmed := bytes.TrimSpace(raw) if len(trimmed) == 0 || bytes.Equal(trimmed, []byte("null")) { @@ -431,14 +501,3 @@ func isInvalidDataObjectUpdateError(err error) bool { errors.Is(err, common.ErrMeasurementTokenNotFound) || errors.Is(err, common.ErrAmbiguousMeasurementToken) } - -func callRealTimeDataWriteStopInterface(_ context.Context, _ orm.JSONMap, _ *float64) error { - // TODO: call the dataRT HTTP API. A nil value stops automatic writes; - // a non-nil value writes the supplied manual measurement value. - return nil -} - -func callRealTimeDataWriteStartInterface(_ context.Context, _ orm.JSONMap, _ *float64) error { - // TODO: call the dataRT HTTP API to start automatic measurement writes. - return nil -} diff --git a/handler/data_object_attribute_update_test.go b/handler/data_object_attribute_update_test.go index 9f1e2b9..e196f4d 100644 --- a/handler/data_object_attribute_update_test.go +++ b/handler/data_object_attribute_update_test.go @@ -5,7 +5,9 @@ import ( "encoding/json" "fmt" "testing" + "time" + "modelRT/client/manualsync" "modelRT/common/errcode" "modelRT/constants" "modelRT/orm" @@ -193,10 +195,11 @@ func TestUpdateMeasurementDataObjectLocksRowAndUpdatesMode(t *testing.T) { startCalled := false result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "mode", constants.MeasurementModeAutomatic, nil, measurementUpdateDependencies{ - startDataRTFunc: func(_ context.Context, dataSource orm.JSONMap, value *float64) error { + syncManualChangeFunc: func(_ context.Context, dataSource orm.JSONMap, mode int16, sample *manualsync.SyntheticData) error { startCalled = true assert.Equal(t, float64(1), dataSource["type"]) - assert.Nil(t, value) + assert.Equal(t, constants.MeasurementModeAutomatic, mode) + assert.Nil(t, sample) return nil }, }) @@ -207,7 +210,7 @@ func TestUpdateMeasurementDataObjectLocksRowAndUpdatesMode(t *testing.T) { require.NoError(t, mock.ExpectationsWereMet()) } -func TestUpdateMeasurementModeToAutomaticReturnsErrorWhenDataRTStartFails(t *testing.T) { +func TestUpdateMeasurementModeToAutomaticContinuesWhenManualSyncFails(t *testing.T) { db, mock, closeDB := newDataObjectUpdateTestDB(t) defer closeDB() @@ -220,13 +223,14 @@ func TestUpdateMeasurementModeToAutomaticReturnsErrorWhenDataRTStartFails(t *tes WillReturnResult(sqlmock.NewResult(0, 1)) mock.ExpectRollback() - _, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "mode", constants.MeasurementModeAutomatic, nil, measurementUpdateDependencies{ - startDataRTFunc: func(context.Context, orm.JSONMap, *float64) error { + result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "mode", constants.MeasurementModeAutomatic, nil, measurementUpdateDependencies{ + syncManualChangeFunc: func(context.Context, orm.JSONMap, int16, *manualsync.SyntheticData) error { return fmt.Errorf("dataRT unavailable") }, }) - require.Error(t, err) - assert.Contains(t, err.Error(), "start automatic measurement write") + require.NoError(t, err) + assert.True(t, result.modeChanged) + assert.Equal(t, constants.MeasurementModeAutomatic, result.mode) require.NoError(t, tx.Rollback().Error) require.NoError(t, mock.ExpectationsWereMet()) } @@ -264,13 +268,14 @@ func TestUpdateMeasurementModeToManualWithoutDataOnlyStopsDataRT(t *testing.T) { stopCalled := false replaceCalled := false result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "mode", constants.MeasurementModeManual, nil, measurementUpdateDependencies{ - updateDataRTFunc: func(_ context.Context, dataSource orm.JSONMap, value *float64) error { + syncManualChangeFunc: func(_ context.Context, dataSource orm.JSONMap, mode int16, sample *manualsync.SyntheticData) error { stopCalled = true assert.Equal(t, float64(1), dataSource["type"]) - assert.Nil(t, value) + assert.Equal(t, constants.MeasurementModeManual, mode) + assert.Nil(t, sample) return nil }, - replaceRedisValueFunc: func(context.Context, *orm.Measurement, float64) error { + replaceRedisValueFunc: func(context.Context, *orm.Measurement, float64, time.Time) error { replaceCalled = true return nil }, @@ -298,12 +303,13 @@ func TestUpdateMeasurementModeToManualReplacesRedisValueWhenDataProvided(t *test callOrder := make([]string, 0, 2) result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "mode", constants.MeasurementModeManual, json.RawMessage(`0`), measurementUpdateDependencies{ - updateDataRTFunc: func(_ context.Context, _ orm.JSONMap, value *float64) error { - callOrder = append(callOrder, "stop-dataRT") - assert.Nil(t, value) + syncManualChangeFunc: func(_ context.Context, _ orm.JSONMap, mode int16, sample *manualsync.SyntheticData) error { + callOrder = append(callOrder, "sync-mode") + assert.Equal(t, constants.MeasurementModeManual, mode) + assert.Nil(t, sample) return nil }, - replaceRedisValueFunc: func(_ context.Context, measurement *orm.Measurement, value float64) error { + replaceRedisValueFunc: func(_ context.Context, measurement *orm.Measurement, value float64, _ time.Time) error { callOrder = append(callOrder, "replace-redis") assert.Equal(t, int64(10), measurement.ID) assert.Equal(t, float64(0), value) @@ -311,13 +317,13 @@ func TestUpdateMeasurementModeToManualReplacesRedisValueWhenDataProvided(t *test }, }) require.NoError(t, err) - assert.Equal(t, []string{"stop-dataRT", "replace-redis"}, callOrder) + assert.Equal(t, []string{"sync-mode", "replace-redis"}, callOrder) assert.Contains(t, result.message, "manual") require.NoError(t, tx.Rollback().Error) require.NoError(t, mock.ExpectationsWereMet()) } -func TestUpdateMeasurementModeToManualDoesNotTouchRedisWhenDataRTStopFails(t *testing.T) { +func TestUpdateMeasurementModeToManualStillUpdatesRedisWhenManualSyncFails(t *testing.T) { db, mock, closeDB := newDataObjectUpdateTestDB(t) defer closeDB() @@ -331,19 +337,20 @@ func TestUpdateMeasurementModeToManualDoesNotTouchRedisWhenDataRTStopFails(t *te mock.ExpectRollback() replaceCalled := false - _, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "mode", constants.MeasurementModeManual, json.RawMessage(`15.2`), measurementUpdateDependencies{ - updateDataRTFunc: func(_ context.Context, _ orm.JSONMap, value *float64) error { - assert.Nil(t, value) + result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "mode", constants.MeasurementModeManual, json.RawMessage(`15.2`), measurementUpdateDependencies{ + syncManualChangeFunc: func(_ context.Context, _ orm.JSONMap, mode int16, sample *manualsync.SyntheticData) error { + assert.Equal(t, constants.MeasurementModeManual, mode) + assert.Nil(t, sample) return fmt.Errorf("dataRT unavailable") }, - replaceRedisValueFunc: func(context.Context, *orm.Measurement, float64) error { + replaceRedisValueFunc: func(context.Context, *orm.Measurement, float64, time.Time) error { replaceCalled = true return nil }, }) - require.Error(t, err) - assert.Contains(t, err.Error(), "stop automatic measurement write") - assert.False(t, replaceCalled) + require.NoError(t, err) + assert.True(t, replaceCalled) + assert.True(t, result.modeChanged) require.NoError(t, tx.Rollback().Error) require.NoError(t, mock.ExpectationsWereMet()) } @@ -395,23 +402,28 @@ func TestUpdateMeasurementDataObjectWritesValueInManualMode(t *testing.T) { mock.ExpectRollback() called := false - writer := func(_ context.Context, measurement *orm.Measurement, value float64) error { + operationTime := time.Date(2026, time.August, 4, 10, 0, 0, 123, time.UTC) + writer := func(_ context.Context, measurement *orm.Measurement, value float64, timestamp time.Time) error { called = true assert.Equal(t, int64(10), measurement.ID) assert.Equal(t, float64(15.2), value) + assert.Equal(t, operationTime, timestamp) return nil } dataRTCalled := false - dataRTWriter := func(_ context.Context, dataSource orm.JSONMap, value *float64) error { + dataRTWriter := func(_ context.Context, dataSource orm.JSONMap, mode int16, sample *manualsync.SyntheticData) error { dataRTCalled = true - require.NotNil(t, value) - assert.Equal(t, float64(15.2), *value) + assert.Equal(t, constants.MeasurementModeManual, mode) + require.NotNil(t, sample) + assert.Equal(t, float64(15.2), sample.Value) + assert.Equal(t, operationTime.UnixNano(), sample.Time) assert.Equal(t, float64(1), dataSource["type"]) return nil } result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "value", float64(15.2), nil, measurementUpdateDependencies{ writeManualValueFunc: writer, - updateDataRTFunc: dataRTWriter, + syncManualChangeFunc: dataRTWriter, + nowFunc: func() time.Time { return operationTime }, }) require.NoError(t, err) assert.True(t, called) @@ -424,6 +436,33 @@ func TestUpdateMeasurementDataObjectWritesValueInManualMode(t *testing.T) { require.NoError(t, mock.ExpectationsWereMet()) } +func TestUpdateMeasurementDataObjectContinuesValueUpdateWhenManualSyncFails(t *testing.T) { + db, mock, closeDB := newDataObjectUpdateTestDB(t) + defer closeDB() + + mock.ExpectBegin() + tx := db.Begin() + require.NoError(t, tx.Error) + expectMeasurementResolution(mock, constants.MeasurementModeManual) + mock.ExpectExec(`UPDATE "measurement" SET "operations"=.*WHERE id = \$3`). + WithArgs(sqlmock.AnyArg(), 500, int64(10)). + WillReturnResult(sqlmock.NewResult(0, 1)) + mock.ExpectRollback() + + result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "value", float64(15.2), nil, measurementUpdateDependencies{ + writeManualValueFunc: func(context.Context, *orm.Measurement, float64, time.Time) error { + return nil + }, + syncManualChangeFunc: func(context.Context, orm.JSONMap, int16, *manualsync.SyntheticData) error { + return fmt.Errorf("manual sync unavailable") + }, + }) + require.NoError(t, err) + assert.Contains(t, result.message, "updated") + require.NoError(t, tx.Rollback().Error) + require.NoError(t, mock.ExpectationsWereMet()) +} + func TestUpdateMeasurementDataObjectReturnsFailureResultAndAppError(t *testing.T) { db, mock, closeDB := newDataObjectUpdateTestDB(t) defer closeDB() @@ -435,7 +474,7 @@ func TestUpdateMeasurementDataObjectReturnsFailureResultAndAppError(t *testing.T mock.ExpectRollback() writeErr := fmt.Errorf("write value failed") - writer := func(context.Context, *orm.Measurement, float64) error { return writeErr } + writer := func(context.Context, *orm.Measurement, float64, time.Time) error { return writeErr } result, err := updateMeasurementDataObject(context.Background(), tx, "nspath.measurement", "value", float64(15.2), nil, measurementUpdateDependencies{ writeManualValueFunc: writer, }) diff --git a/main.go b/main.go index 32affa1..20d3dad 100644 --- a/main.go +++ b/main.go @@ -14,6 +14,7 @@ import ( "syscall" "time" + "modelRT/client/manualsync" "modelRT/config" "modelRT/constants" "modelRT/database" @@ -98,10 +99,11 @@ func main() { logger.InitLoggerInstance(modelRTConfig.LoggerConfig) defer logger.GetLoggerInstance().Sync() + baseCtx := context.Background() // init OTel TracerProvider tp, tpErr := middleware.InitTracerProvider(context.Background(), modelRTConfig) if tpErr != nil { - log.Printf("warn: OTLP tracer init failed, tracing disabled: %v", tpErr) + logger.Error(baseCtx, "init OTLP tracer provider failed, tracing disabled", "error", tpErr) } if tp != nil { defer func() { @@ -111,7 +113,7 @@ func main() { }() } - ctx, startupSpan := otel.Tracer("modelRT/main").Start(context.Background(), "startup") + ctx, startupSpan := otel.Tracer("modelRT/main").Start(baseCtx, "startup") defer startupSpan.End() hostName, err := os.Hostname() @@ -126,6 +128,16 @@ func main() { panic(err) } + manualSyncClient, err := manualsync.NewClient(modelRTConfig.DataRTConfig.ManualSync) + if err != nil { + // TODO: This is a best-effort integration and may not be the final design. + // The current behavior lets modelRT start without a manual-sync client; + // affected updates still succeed and log an error while sync events are lost. + // Revisit fail-fast validation if this downstream service becomes mandatory. + logger.Error(ctx, "init manual measurement sync client failed", "error", err) + } + manualsync.SetDefaultSyncer(manualSyncClient) + // init postgresDBClient postgresDBClient = database.InitPostgresDBInstance(ctx, modelRTConfig.PostgresDBURI) @@ -191,7 +203,7 @@ func main() { // async push task message to rabbitMQ go task.PushTaskToRabbitMQ(ctx, modelRTConfig.RabbitMQConfig, task.TaskMsgChan) - postgresDBClient.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if err := postgresDBClient.WithContext(ctx).Transaction(func(tx *gorm.DB) error { // load circuit diagram from postgres // componentTypeMap, err := database.QueryCircuitDiagramComponentFromDB(cancelCtx, tx, parsePool) // if err != nil { @@ -201,91 +213,80 @@ func main() { cacheMap, err := model.GetNSpathToIsLocalMap(ctx, postgresDBClient) if err != nil { - logger.Error(ctx, "get nspath to is_local map failed", "error", err) - panic(err) + return fmt.Errorf("get nspath to is_local map: %w", err) } model.NSPathToIsLocalMap = cacheMap err = model.CleanupRecommendRedisCache(ctx) if err != nil { - logger.Error(ctx, "clean up component measurement and attribute group failed", "error", err) - panic(err) + return fmt.Errorf("clean up component measurement and attribute group: %w", err) } measurementSet, err := database.GetFullMeasurementSet(ctx, postgresDBClient) if err != nil { - logger.Error(ctx, "generate component measurement group failed", "error", err) - panic(err) + return fmt.Errorf("generate component measurement group: %w", err) } fullParentPath, isLocalParentPath, err := model.TraverseMeasurementGroupTables(ctx, *measurementSet) if err != nil { - logger.Error(ctx, "store component measurement group into redis failed", "error", err) - panic(err) + return fmt.Errorf("store component measurement group into redis: %w", err) } compAttrSet, err := database.GenAllAttributeMap(tx) if err != nil { - logger.Error(ctx, "generate component attribute group failed", "error", err) - panic(err) + return fmt.Errorf("generate component attribute group: %w", err) } err = model.TraverseAttributeGroupTables(ctx, tx, fullParentPath, isLocalParentPath, compAttrSet) if err != nil { - logger.Error(ctx, "store component attribute group into redis failed", "error", err) - panic(err) + return fmt.Errorf("store component attribute group into redis: %w", err) } componentColumnNames, err := database.QueryComponentColumnNames(ctx, tx) if err != nil { - logger.Error(ctx, "query component table column names failed", "error", err) - panic(err) + return fmt.Errorf("query component table column names: %w", err) } err = model.StoreComponentColumnRecommend(ctx, fullParentPath, isLocalParentPath, componentColumnNames) if err != nil { - logger.Error(ctx, "store component column recommend content failed", "error", err) - panic(err) + return fmt.Errorf("store component column recommend content: %w", err) } parameterRecords, err := database.QueryParameterInitializationRecords(ctx, tx) if err != nil { - logger.Error(ctx, "load parameter data objects from postgres failed", "error", err) - panic(err) + return fmt.Errorf("load parameter data objects from postgres: %w", err) } err = model.InitializeParameterDataObjects(ctx, parameterRecords) if err != nil { - logger.Error(ctx, "initialize parameter data objects failed", "error", err) - panic(err) + return fmt.Errorf("initialize parameter data objects: %w", err) } measurementRecords, err := database.QueryMeasurementInitializationRecords(ctx, tx) if err != nil { - logger.Error(ctx, "load measurement data objects from postgres failed", "error", err) - panic(err) + return fmt.Errorf("load measurement data objects from postgres: %w", err) } err = model.InitializeMeasurementDataObjects(ctx, measurementRecords) if err != nil { - logger.Error(ctx, "initialize measurement data objects failed", "error", err) - panic(err) + return fmt.Errorf("initialize measurement data objects: %w", err) } allMeasurement, err := database.GetAllMeasurements(ctx, tx) if err != nil { - logger.Error(ctx, "load topologic info from postgres failed", "error", err) - panic(err) + return fmt.Errorf("load measurements from postgres: %w", err) } go realtimedata.StartComputingRealTimeDataLimit(ctx, allMeasurement) topologics, err := database.QueryTopologic(ctx, tx) if err != nil { - logger.Error(ctx, "load topologic info from postgres failed", "error", err) - panic(err) + return fmt.Errorf("load topologic info from postgres: %w", err) } diagram.SetGlobalTopologyGraph(diagram.NewTopologyGraph(topologics)) return nil - }) + }); err != nil { + logger.Error(ctx, "initialize modelRT startup data failed", "error", err) + panic(err) + } // use release mode in production if modelRTConfig.DeployEnv == constants.ProductionDeployMode { diff --git a/model/measurement_data_object_init.go b/model/measurement_data_object_init.go index 768c98c..db6a112 100644 --- a/model/measurement_data_object_init.go +++ b/model/measurement_data_object_init.go @@ -37,9 +37,9 @@ type MeasurementInitializationRecord struct { 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"` + MeasurementDataSource orm.JSONMap `gorm:"column:measurement_data_source;type:jsonb"` + MeasurementEventPlan orm.JSONMap `gorm:"column:measurement_event_plan;type:jsonb"` + MeasurementBinding orm.JSONMap `gorm:"column:measurement_binding;type:jsonb"` } // InitializeMeasurementDataObjects creates one seven-part Redis hash and diff --git a/handler/data_object_redis_change.go b/repository/redis/data_object_redis_change.go similarity index 86% rename from handler/data_object_redis_change.go rename to repository/redis/data_object_redis_change.go index 5ee2e43..64ea4ef 100644 --- a/handler/data_object_redis_change.go +++ b/repository/redis/data_object_redis_change.go @@ -1,4 +1,5 @@ -package handler +// Package redis provides Redis persistence helpers. +package redis import ( "context" @@ -8,10 +9,7 @@ import ( "strconv" "time" - "modelRT/model" - "modelRT/orm" - - "github.com/redis/go-redis/v9" + redisclient "github.com/redis/go-redis/v9" ) const redisChangeRestoreTimeout = 5 * time.Second @@ -25,21 +23,21 @@ type redisHashChange struct { type redisZSetChange struct { key string - oldValues []redis.Z - newValues []redis.Z + oldValues []redisclient.Z + newValues []redisclient.Z } // RedisChangeSet keeps the Redis changes belonging to one PostgreSQL // transaction. Changes are prepared first and applied together immediately // before the PostgreSQL transaction is committed. type RedisChangeSet struct { - client *redis.Client + client *redisclient.Client hashChanges []redisHashChange zsetChanges []redisZSetChange applied bool } -func NewRedisChangeSet(client *redis.Client) *RedisChangeSet { +func NewRedisChangeSet(client *redisclient.Client) *RedisChangeSet { return &RedisChangeSet{client: client} } @@ -74,19 +72,16 @@ func (changes *RedisChangeSet) AddHashChange( func (changes *RedisChangeSet) AddMeasurementValueChange( ctx context.Context, - measurement *orm.Measurement, + key string, value float64, + timestamp time.Time, replace bool, ) error { if changes == nil || changes.client == nil { return fmt.Errorf("redis client is not initialized") } - if measurement == nil { - return fmt.Errorf("measurement is nil") - } - key, err := model.GenerateMeasureIdentifier(measurement.DataSource) - if err != nil { - return fmt.Errorf("generate measurement redis key: %w", err) + if key == "" { + return fmt.Errorf("measurement redis key is empty") } keyType, err := changes.client.Type(ctx, key).Result() if err != nil { @@ -100,8 +95,8 @@ func (changes *RedisChangeSet) AddMeasurementValueChange( return fmt.Errorf("query measurement redis values for %q: %w", key, err) } - newMember := strconv.FormatInt(time.Now().UnixNano(), 10) - newValues := []redis.Z{{Score: value, Member: newMember}} + newMember := strconv.FormatInt(timestamp.UnixNano(), 10) + newValues := []redisclient.Z{{Score: value, Member: newMember}} if !replace { newValues = mergeRedisZValues(oldValues, newValues...) } @@ -122,11 +117,11 @@ func (changes *RedisChangeSet) Apply(ctx context.Context) error { } keys := changes.keys() - err := changes.client.Watch(ctx, func(tx *redis.Tx) error { + err := changes.client.Watch(ctx, func(tx *redisclient.Tx) error { if err := changes.verify(ctx, tx, false); err != nil { return err } - _, err := tx.TxPipelined(ctx, func(pipe redis.Pipeliner) error { + _, err := tx.TxPipelined(ctx, func(pipe redisclient.Pipeliner) error { for _, change := range changes.hashChanges { pipe.HSet(ctx, change.key, change.field, change.newValue) } @@ -165,13 +160,13 @@ func (changes *RedisChangeSet) Revert(ctx context.Context) error { func (changes *RedisChangeSet) restoreOldValues(ctx context.Context, compareNew bool) error { keys := changes.keys() - return changes.client.Watch(ctx, func(tx *redis.Tx) error { + return changes.client.Watch(ctx, func(tx *redisclient.Tx) error { if compareNew { if err := changes.verify(ctx, tx, true); err != nil { return err } } - _, err := tx.TxPipelined(ctx, func(pipe redis.Pipeliner) error { + _, err := tx.TxPipelined(ctx, func(pipe redisclient.Pipeliner) error { for _, change := range changes.hashChanges { pipe.HSet(ctx, change.key, change.field, change.oldValue) } @@ -189,7 +184,7 @@ func (changes *RedisChangeSet) restoreOldValues(ctx context.Context, compareNew func (changes *RedisChangeSet) restoreAfterApplyFailure(ctx context.Context) error { keys := changes.keys() - return changes.client.Watch(ctx, func(tx *redis.Tx) error { + return changes.client.Watch(ctx, func(tx *redisclient.Tx) error { for _, change := range changes.hashChanges { actual, err := tx.HGet(ctx, change.key, change.field).Result() if err != nil { @@ -208,7 +203,7 @@ func (changes *RedisChangeSet) restoreAfterApplyFailure(ctx context.Context) err return fmt.Errorf("redis zset %q changed concurrently", change.key) } } - _, err := tx.TxPipelined(ctx, func(pipe redis.Pipeliner) error { + _, err := tx.TxPipelined(ctx, func(pipe redisclient.Pipeliner) error { for _, change := range changes.hashChanges { pipe.HSet(ctx, change.key, change.field, change.oldValue) } @@ -224,7 +219,7 @@ func (changes *RedisChangeSet) restoreAfterApplyFailure(ctx context.Context) err }, keys...) } -func (changes *RedisChangeSet) verify(ctx context.Context, tx *redis.Tx, expectNew bool) error { +func (changes *RedisChangeSet) verify(ctx context.Context, tx *redisclient.Tx, expectNew bool) error { for _, change := range changes.hashChanges { expected := change.oldValue if expectNew { @@ -300,28 +295,28 @@ func redisChangeString(value any) (string, error) { } } -func cloneRedisZValues(values []redis.Z) []redis.Z { - cloned := make([]redis.Z, len(values)) +func cloneRedisZValues(values []redisclient.Z) []redisclient.Z { + cloned := make([]redisclient.Z, len(values)) copy(cloned, values) return cloned } -func mergeRedisZValues(current []redis.Z, additions ...redis.Z) []redis.Z { - valuesByMember := make(map[string]redis.Z, len(current)+len(additions)) +func mergeRedisZValues(current []redisclient.Z, additions ...redisclient.Z) []redisclient.Z { + valuesByMember := make(map[string]redisclient.Z, len(current)+len(additions)) for _, value := range current { valuesByMember[fmt.Sprint(value.Member)] = value } for _, value := range additions { valuesByMember[fmt.Sprint(value.Member)] = value } - values := make([]redis.Z, 0, len(valuesByMember)) + values := make([]redisclient.Z, 0, len(valuesByMember)) for _, value := range valuesByMember { values = append(values, value) } return normalizeRedisZValues(values) } -func normalizeRedisZValues(values []redis.Z) []redis.Z { +func normalizeRedisZValues(values []redisclient.Z) []redisclient.Z { normalized := cloneRedisZValues(values) sort.Slice(normalized, func(i, j int) bool { if normalized[i].Score != normalized[j].Score { @@ -332,7 +327,7 @@ func normalizeRedisZValues(values []redis.Z) []redis.Z { return normalized } -func equalRedisZValues(left, right []redis.Z) bool { +func equalRedisZValues(left, right []redisclient.Z) bool { left = normalizeRedisZValues(left) right = normalizeRedisZValues(right) if len(left) != len(right) { diff --git a/handler/data_object_redis_change_test.go b/repository/redis/data_object_redis_change_test.go similarity index 98% rename from handler/data_object_redis_change_test.go rename to repository/redis/data_object_redis_change_test.go index a0c117e..20a2319 100644 --- a/handler/data_object_redis_change_test.go +++ b/repository/redis/data_object_redis_change_test.go @@ -1,4 +1,4 @@ -package handler +package redis import ( "testing"