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
This commit is contained in:
douxu 2026-08-04 17:12:17 +08:00
parent d670963854
commit 2fbef3e1fa
15 changed files with 965 additions and 170 deletions

255
client/manualsync/client.go Normal file
View File

@ -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)
}
}

View File

@ -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),
}
}

33
client/manualsync/init.go Normal file
View File

@ -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)
}

View File

@ -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)
}

View File

@ -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:"-"`

31
config/config_test.go Normal file
View File

@ -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)
}

View File

@ -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"

View File

@ -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\. 部署 ModelRTKubernetes
所有资源部署在 `default` 命名空间YAML 文件位于 `deploy/k8s/`

View File

@ -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

View File

@ -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
}

View File

@ -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,
})

65
main.go
View File

@ -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 {

View File

@ -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

View File

@ -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) {

View File

@ -1,4 +1,4 @@
package handler
package redis
import (
"testing"