Compare commits

..

5 Commits

Author SHA1 Message Date
Swapnil Nakade
b419d4ab04 Merge branch 'main' into issue-2977 2026-09-29 19:07:32 +05:30
swapnil-signoz
ccf72a274b refactor: addressing review comments 2026-09-29 19:06:38 +05:30
swapnil-signoz
3ae7d4441f feat: adding integration tests 2026-09-28 15:51:38 +05:30
Swapnil Nakade
cbfe328936 Merge branch 'main' into issue-2977 2026-09-26 06:53:21 +05:30
swapnil-signoz
9708e89d8c feat: adding sync state in cloud integration 2026-09-26 06:48:19 +05:30
22 changed files with 768 additions and 192 deletions

View File

@@ -1763,12 +1763,15 @@ components:
additionalProperties: {}
nullable: true
type: object
syncState:
$ref: '#/components/schemas/CloudintegrationtypesSyncState'
timestampMillis:
format: int64
type: integer
required:
- timestampMillis
- data
- syncState
type: object
CloudintegrationtypesAzureAccountConfig:
properties:
@@ -2013,6 +2016,8 @@ components:
format: date-time
nullable: true
type: string
syncState:
$ref: '#/components/schemas/CloudintegrationtypesSyncState'
required:
- account_id
- cloud_account_id
@@ -2022,6 +2027,7 @@ components:
- providerAccountId
- integrationConfig
- removedAt
- syncState
type: object
CloudintegrationtypesGettableServicesMetadata:
properties:
@@ -2121,6 +2127,9 @@ components:
type: object
providerAccountId:
type: string
syncedVersion:
nullable: true
type: integer
required:
- data
type: object
@@ -2133,6 +2142,18 @@ components:
gcp:
$ref: '#/components/schemas/CloudintegrationtypesGCPIntegrationConfig'
type: object
CloudintegrationtypesRegionState:
enum:
- enabled
- disabled
type: string
CloudintegrationtypesRegionSyncState:
properties:
state:
$ref: '#/components/schemas/CloudintegrationtypesRegionState'
required:
- state
type: object
CloudintegrationtypesService:
properties:
assets:
@@ -2274,6 +2295,23 @@ components:
metrics:
type: boolean
type: object
CloudintegrationtypesSyncState:
nullable: true
properties:
inSync:
type: boolean
regions:
additionalProperties:
$ref: '#/components/schemas/CloudintegrationtypesRegionSyncState'
type: object
version:
format: int64
type: integer
required:
- version
- inSync
- regions
type: object
CloudintegrationtypesUpdatableAccount:
properties:
config:

View File

@@ -183,32 +183,52 @@ func (module *module) AgentCheckIn(ctx context.Context, orgID valuer.UUID, provi
return nil, errors.New(errors.TypeAlreadyExists, cloudintegrationtypes.ErrCodeCloudIntegrationAlreadyConnected, errMessage)
}
account, err := module.store.GetAccountByID(ctx, orgID, req.CloudIntegrationID, provider)
storableAccount, err := module.store.GetAccountByID(ctx, orgID, req.CloudIntegrationID, provider)
if err != nil {
return nil, err
}
account, err := cloudintegrationtypes.NewAccountFromStorable(storableAccount)
if err != nil {
return nil, err
}
syncState := account.NextSyncState(req.SyncedVersion)
// If account has been removed (disconnected), return a minimal response with empty integration config.
// The agent uses this response to clean up resources
if account.RemovedAt != nil {
// Heartbeat stays frozen after removal, only the sync state is updated.
if account.AgentReport != nil && syncState != nil {
account.UpdateSyncState(syncState)
storableAccount, err = cloudintegrationtypes.NewStorableCloudIntegration(account)
if err != nil {
return nil, err
}
err = module.store.UpdateAgentReport(ctx, storableAccount)
if err != nil {
return nil, err
}
}
return cloudintegrationtypes.NewAgentCheckInResponse(
req.ProviderAccountID,
account.ID.StringValue(),
new(cloudintegrationtypes.ProviderIntegrationConfig),
account.RemovedAt,
syncState,
), nil
}
// update account with cloud provider account id and agent report (heartbeat)
account.Update(&req.ProviderAccountID, cloudintegrationtypes.NewAgentReport(req.Data))
account.UpdateAgentReport(&req.ProviderAccountID, cloudintegrationtypes.NewAgentReport(req.Data, syncState))
err = module.store.UpdateAccount(ctx, account)
storableAccount, err = cloudintegrationtypes.NewStorableCloudIntegration(account)
if err != nil {
return nil, err
}
// Get account as domain object for config access (enabled regions, etc.)
domainAccount, err := cloudintegrationtypes.NewAccountFromStorable(account)
err = module.store.UpdateAgentReport(ctx, storableAccount)
if err != nil {
return nil, err
}
@@ -223,8 +243,7 @@ func (module *module) AgentCheckIn(ctx context.Context, orgID valuer.UUID, provi
return nil, err
}
// Delegate integration config building entirely to the provider module
integrationConfig, err := cloudProvider.BuildIntegrationConfig(ctx, domainAccount, storedServices)
integrationConfig, err := cloudProvider.BuildIntegrationConfig(ctx, account, storedServices)
if err != nil {
return nil, err
}
@@ -234,6 +253,7 @@ func (module *module) AgentCheckIn(ctx context.Context, orgID valuer.UUID, provi
account.ID.StringValue(),
integrationConfig,
account.RemovedAt,
syncState,
), nil
}

View File

@@ -3366,6 +3366,37 @@ export interface CloudintegrationtypesAWSServiceConfigDTO {
metrics?: CloudintegrationtypesAWSServiceMetricsConfigDTO;
}
export enum CloudintegrationtypesRegionStateDTO {
enabled = 'enabled',
disabled = 'disabled',
}
export interface CloudintegrationtypesRegionSyncStateDTO {
state: CloudintegrationtypesRegionStateDTO;
}
export type CloudintegrationtypesSyncStateDTORegions = {
[key: string]: CloudintegrationtypesRegionSyncStateDTO;
};
/**
* @nullable
*/
export type CloudintegrationtypesSyncStateDTO = {
/**
* @type boolean
*/
inSync: boolean;
/**
* @type object
*/
regions: CloudintegrationtypesSyncStateDTORegions;
/**
* @type integer
* @format int64
*/
version: number;
} | null;
export type CloudintegrationtypesAgentReportDTODataAnyOf = {
[key: string]: unknown;
};
@@ -3384,6 +3415,7 @@ export type CloudintegrationtypesAgentReportDTO = {
* @type object,null
*/
data: CloudintegrationtypesAgentReportDTOData;
syncState: CloudintegrationtypesSyncStateDTO | null;
/**
* @type integer
* @format int64
@@ -3812,6 +3844,7 @@ export interface CloudintegrationtypesGettableAgentCheckInDTO {
* @format date-time
*/
removedAt: string | null;
syncState: CloudintegrationtypesSyncStateDTO | null;
}
export interface CloudintegrationtypesServiceMetadataDTO {
@@ -3882,6 +3915,10 @@ export interface CloudintegrationtypesPostableAgentCheckInDTO {
* @type string
*/
providerAccountId?: string;
/**
* @type integer,null
*/
syncedVersion?: number | null;
}
export interface CloudintegrationtypesStorableIntegrationDashboardDTO {

View File

@@ -24,6 +24,7 @@ const accountsResponse: ListAccounts200 = {
agentReport: {
timestampMillis: 1747114366214,
data: null,
syncState: null,
},
providerAccountId: PROVIDER_ACCOUNT_ID,
removedAt: null,

View File

@@ -295,7 +295,11 @@ const account = (
provider,
providerAccountId: ACCOUNTS[provider][index],
config: accountConfig(provider),
agentReport: { timestampMillis: Date.now() - 45 * 1000, data: null },
agentReport: {
timestampMillis: Date.now() - 45 * 1000,
data: null,
syncState: null,
},
createdAt: new Date(Date.now() - 21 * 24 * 60 * 60 * 1000).toISOString(),
updatedAt: new Date(Date.now() - 60 * 60 * 1000).toISOString(),
removedAt: null,

View File

@@ -134,6 +134,24 @@ func (store *store) UpdateAccount(ctx context.Context, account *cloudintegration
BunDBCtx(ctx).
NewUpdate().
Model(account).
Column("config").
Column("updated_at").
WherePK().
Where("org_id = ?", account.OrgID).
Where("provider = ?", account.Provider).
Exec(ctx)
return err
}
func (store *store) UpdateAgentReport(ctx context.Context, account *cloudintegrationtypes.StorableCloudIntegration) error {
_, err := store.
store.
BunDBCtx(ctx).
NewUpdate().
Model(account).
Column("account_id").
Column("last_agent_report").
WherePK().
Where("org_id = ?", account.OrgID).
Where("provider = ?", account.Provider).

View File

@@ -40,15 +40,11 @@ func stripKeyAlias(name string) string {
return keyAliasRe.ReplaceAllString(name, "")
}
// unwrapVariant returns the concrete value inside the driver's scan envelopes: chcol.Variant for a
// Dynamic column (a JSON path such as body_v2.level), and chcol.JSON for a whole JSON column, decoded
// into a nested document.
// unwrapVariant returns the concrete value inside the chcol.Variant envelope the driver scans a
// Dynamic column — a JSON path such as body_v2.level — into.
func unwrapVariant(val any) any {
switch v := val.(type) {
case chcol.Variant:
if v, ok := val.(chcol.Variant); ok {
return v.Any()
case chcol.JSON:
return telemetrystoretypes.NestedJSON(v)
}
return val
}
@@ -62,7 +58,7 @@ func labelValue(val any) string {
if val == nil {
return ""
}
if v, ok := val.(map[string]any); ok {
if v, ok := val.(telemetrystoretypes.JSONValue); ok {
if raw, err := json.Marshal(v); err == nil {
return string(raw)
}
@@ -208,7 +204,7 @@ func readAsTimeSeries(rows driver.Rows, queryWindow *qbtypes.TimeRange, step qbt
Value: *val,
})
case *chcol.JSON, *chcol.Variant:
case *telemetrystoretypes.JSONValue, *chcol.Variant:
val := labelValue(derefValue(ptr))
lblVals = append(lblVals, val)
lblObjs = append(lblObjs, &qbtypes.Label{
@@ -540,14 +536,7 @@ func readAsRaw(rows driver.Rows, queryName string) (*qbtypes.RawData, error) {
name := stripKeyAlias(colNames[i])
// de-reference the typed pointer to any
raw := reflect.ValueOf(cellPtr).Elem().Interface()
// the attributes bag is flattened to dotted keys downstream; decode it flat so a key stored as both a scalar and an object is not collapsed into a mislabeled key.
var val any
if j, ok := raw.(chcol.JSON); ok && name == "attributes" {
val = telemetrystoretypes.FlattenJSON(j)
} else {
val = unwrapVariant(raw)
}
val := unwrapVariant(reflect.ValueOf(cellPtr).Elem().Interface())
// special-case: timestamp column
if name == "timestamp" || name == "timestamp_datetime" {
@@ -587,6 +576,8 @@ func flattenJSONPaths(prefix string, m map[string]any, out map[string]any) {
switch child := v.(type) {
case map[string]any:
flattenJSONPaths(key, child, out)
case telemetrystoretypes.JSONValue:
flattenJSONPaths(key, child, out)
default:
out[key] = v
}
@@ -602,7 +593,7 @@ func mergeSpanAttributeColumns(data map[string]any) {
attrStr, hasStr := data["attributes_string"]
attrNum, hasNum := data["attributes_number"]
attrBool, hasBool := data["attributes_bool"]
attrJSON, _ := data["attributes"].(map[string]any)
attrJSON, _ := data["attributes"].(telemetrystoretypes.JSONValue)
// todo(nitya): move to resource json
resStr, hasRes := data["resources_string"]
if hasStr || hasNum || hasBool || attrJSON != nil || hasRes {

View File

@@ -3,9 +3,14 @@ package querier
import (
"reflect"
"testing"
"time"
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
cmock "github.com/SigNoz/clickhouse-go-mock"
"github.com/SigNoz/signoz/pkg/telemetrystore"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/spantypes"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
@@ -78,31 +83,101 @@ func TestMergeSpanAttributeColumns_ParsesEventsAndLinks(t *testing.T) {
}
}
// A JSON path (e.g. `body_v2.level`) comes back as a Dynamic column, which the driver scans into a
// chcol.Variant envelope; a whole JSON column comes back as chcol.JSON, decoded into a nested document.
// A ClickHouse query can put a JSON column in the result of any request type — e.g.
// `select * from signoz_logs.logs_v2` on a body_v2 stack, where `*` covers body_v2.
func TestConsume_JSONColumn(t *testing.T) {
ts := time.Date(2026, 8, 14, 10, 0, 0, 0, time.UTC)
body := `{"level":"error","attrs":{"code":500}}`
wantBody := telemetrystoretypes.JSONValue{
"level": "error",
"attrs": map[string]any{"code": float64(500)},
}
// the scalar reader reuses its scan slots across rows, so each row must still carry its own body
t.Run("scalar", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "body_v2", Type: "JSON"},
{Name: "__result_0", Type: "UInt64"},
}, [][]any{{body, uint64(3)}, {`{"level":"warn"}`, uint64(1)}}))
payload, err := consume(rows, qbtypes.RequestTypeScalar, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.ScalarData)
require.Len(t, data.Data, 2)
assert.Equal(t, wantBody, data.Data[0][0])
assert.Equal(t, uint64(3), data.Data[0][1])
assert.Equal(t, telemetrystoretypes.JSONValue{"level": "warn"}, data.Data[1][0])
assert.Equal(t, uint64(1), data.Data[1][1])
})
t.Run("time series", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "ts", Type: "DateTime"},
{Name: "body_v2", Type: "JSON"},
{Name: "__result_0", Type: "UInt64"},
}, [][]any{{ts, body, uint64(3)}}))
payload, err := consume(rows, qbtypes.RequestTypeTimeSeries, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.TimeSeriesData)
require.Len(t, data.Aggregations, 1)
require.Len(t, data.Aggregations[0].Series, 1)
require.Len(t, data.Aggregations[0].Series[0].Values, 1)
assert.Equal(t, float64(3), data.Aggregations[0].Series[0].Values[0].Value)
})
// grouping by a JSON column is legal in ClickHouse, so each document has to label its own
// series rather than being dropped, which would merge every group into one
t.Run("time series grouped by the JSON column", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "ts", Type: "DateTime"},
{Name: "body_v2", Type: "JSON"},
{Name: "__result_0", Type: "UInt64"},
}, [][]any{
{ts, `{"level":"error"}`, uint64(7)},
{ts, `{"level":"warn"}`, uint64(2)},
}))
payload, err := consume(rows, qbtypes.RequestTypeTimeSeries, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.TimeSeriesData)
require.Len(t, data.Aggregations, 1)
require.Len(t, data.Aggregations[0].Series, 2)
got := map[string]float64{}
for _, series := range data.Aggregations[0].Series {
require.Len(t, series.Labels, 1)
require.Len(t, series.Values, 1)
got[series.Labels[0].Value.(string)] = series.Values[0].Value
}
assert.Equal(t, map[string]float64{`{"level":"error"}`: 7, `{"level":"warn"}`: 2}, got)
})
t.Run("raw", func(t *testing.T) {
rows := telemetrystore.WrapRows(cmock.NewRows([]cmock.ColumnType{
{Name: "timestamp", Type: "DateTime"},
{Name: "body_v2", Type: "JSON"},
}, [][]any{{ts, body}}))
payload, err := consume(rows, qbtypes.RequestTypeRaw, nil, qbtypes.Step{}, "A")
require.NoError(t, err)
data := payload.(*qbtypes.RawData)
require.Len(t, data.Rows, 1)
assert.Equal(t, ts, data.Rows[0].Timestamp.UTC())
assert.Equal(t, wantBody, data.Rows[0].Data["body_v2"])
})
}
// A JSON path (e.g. `body_v2.level`) comes back as a Dynamic column, which the driver scans
// into a chcol.Variant envelope rather than the value itself.
func TestUnwrapVariant(t *testing.T) {
assert.Equal(t, "error", unwrapVariant(chcol.NewDynamicWithType("error", "String")))
assert.Nil(t, unwrapVariant(chcol.Dynamic{}))
assert.Equal(t, uint64(3), unwrapVariant(uint64(3)))
j := chcol.NewJSON()
j.SetValueAtPath("level", "error")
j.SetValueAtPath("attrs.code", int64(500))
assert.Equal(t, map[string]any{
"level": "error",
"attrs": map[string]any{"code": float64(500)},
}, unwrapVariant(*j))
}
// labelValue renders a JSON group-by value as a stable, sorted-key string so structurally equal
// documents share a series.
func TestLabelValue(t *testing.T) {
assert.Equal(t, "", labelValue(nil))
assert.Equal(t, "error", labelValue(chcol.NewDynamicWithType("error", "String")))
assert.Equal(t, `{"attrs":{"code":500},"level":"error"}`, labelValue(map[string]any{
"level": "error",
"attrs": map[string]any{"code": 500},
}))
}
func TestMergeSpanAttributeColumns_EmptyEventsAndLinks(t *testing.T) {
@@ -132,7 +207,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
{
name: "JSONOnly_FlattensNestedPaths_PreservesTypes",
data: map[string]any{
"attributes": map[string]any{
"attributes": telemetrystoretypes.JSONValue{
"http": map[string]any{"route": "/api/pay", "retry": map[string]any{"count": float64(3)}},
"cache.hit": true,
},
@@ -144,7 +219,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
data: map[string]any{
"attributes_string": map[string]string{"http.route": "/old", "only.map": "m"},
"attributes_number": map[string]float64{"http.status": 500},
"attributes": map[string]any{"http": map[string]any{"route": "/new"}, "only.json": "j"},
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"route": "/new"}, "only.json": "j"},
},
want: map[string]any{"http.route": "/old", "only.map": "m", "http.status": float64(500), "only.json": "j"},
},
@@ -154,7 +229,7 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
"attributes_string": map[string]string{"http.route": "/map"},
"attributes_number": map[string]float64{"http.status": 200},
"attributes_bool": map[string]bool{"cache.hit": true},
"attributes": map[string]any{},
"attributes": telemetrystoretypes.JSONValue{},
},
want: map[string]any{"http.route": "/map", "http.status": float64(200), "cache.hit": true},
},
@@ -162,28 +237,28 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
name: "MapOnly_NilJSON_BehavesAsAbsent",
data: map[string]any{
"attributes_string": map[string]string{"http.route": "/map"},
"attributes": map[string]any(nil),
"attributes": telemetrystoretypes.JSONValue(nil),
},
want: map[string]any{"http.route": "/map"},
},
{
name: "Arrays_StayLeafValues",
data: map[string]any{
"attributes": map[string]any{"http": map[string]any{"tags": []any{"a", "b"}, "codes": []any{float64(1), float64(2)}}},
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"tags": []any{"a", "b"}, "codes": []any{float64(1), float64(2)}}},
},
want: map[string]any{"http.tags": []any{"a", "b"}, "http.codes": []any{float64(1), float64(2)}},
},
{
name: "TopLevelArrayOfMaps_StaysNativeLeaf",
data: map[string]any{
"attributes": map[string]any{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
"attributes": telemetrystoretypes.JSONValue{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
},
want: map[string]any{"key": []any{map[string]any{"a": float64(1)}, map[string]any{"b": float64(2)}}},
},
{
name: "NestedArrayOfMaps_StaysNativeLeaf_NoIndexPaths",
data: map[string]any{
"attributes": map[string]any{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
},
want: map[string]any{"http.items": []any{map[string]any{"a": float64(1)}}},
},
@@ -191,28 +266,28 @@ func TestMergeSpanAttributeColumns_JSONColumn(t *testing.T) {
name: "DualWritten_NestedArray_IndexKeysAndJSONArrayCoexist",
data: map[string]any{
"attributes_number": map[string]float64{"http.items.0.a": 1},
"attributes": map[string]any{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"items": []any{map[string]any{"a": float64(1)}}}},
},
want: map[string]any{"http.items.0.a": float64(1), "http.items": []any{map[string]any{"a": float64(1)}}},
},
{
name: "JSONNull_KeptAsNil",
data: map[string]any{
"attributes": map[string]any{"k": nil},
"attributes": telemetrystoretypes.JSONValue{"k": nil},
},
want: map[string]any{"k": nil},
},
{
name: "KeyIsLeafValue_NotFlattened",
data: map[string]any{
"attributes": map[string]any{"http": "plaintext"},
"attributes": telemetrystoretypes.JSONValue{"http": "plaintext"},
},
want: map[string]any{"http": "plaintext"},
},
{
name: "KeyIsParent_FlattensToDottedPath",
data: map[string]any{
"attributes": map[string]any{"http": map[string]any{"route": "/a"}},
"attributes": telemetrystoretypes.JSONValue{"http": map[string]any{"route": "/a"}},
},
want: map[string]any{"http.route": "/a"},
},

View File

@@ -16,6 +16,7 @@ import (
"github.com/SigNoz/signoz/pkg/querybuilder"
"github.com/SigNoz/signoz/pkg/types/featuretypes"
qbtypes "github.com/SigNoz/signoz/pkg/types/querybuildertypes/querybuildertypesv5"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
"github.com/SigNoz/signoz/pkg/types/telemetrytypes"
"github.com/SigNoz/signoz/pkg/valuer"
)
@@ -1201,7 +1202,7 @@ func (q *querier) postProcessLogBody(ctx context.Context, orgID valuer.UUID, res
// carried one. Anything that is not a decoded document — the legacy string body, a NULL cell —
// is legal under these names and left alone.
func stripEmptyBodyMessage(val any) {
bodyMap, ok := val.(map[string]any)
bodyMap, ok := val.(telemetrystoretypes.JSONValue)
if !ok {
return
}

View File

@@ -97,8 +97,8 @@ func New(ctx context.Context, providerSettings factory.ProviderSettings, config
options.MaxIdleConns = config.Connection.MaxIdleConns
options.MaxOpenConns = config.Connection.MaxOpenConns
options.DialTimeout = config.Connection.DialTimeout
// Decode JSON columns via the flattened native serialization (CH 25.6+); without it clickhouse-go mis-decodes the SharedData layout of JSON(max_dynamic_paths=0) columns and desyncs the native protocol.
options.Settings["output_format_native_use_flattened_dynamic_and_json_serialization"] = 1
// This is to avoid the driver decoding issues with JSON columns
options.Settings["output_format_native_write_json_as_string"] = 1
chConn, err := clickhouse.Open(options)
if err != nil {
@@ -184,7 +184,7 @@ func (p *provider) Query(ctx context.Context, query string, args ...interface{})
}
return &rowsWithHooks{
Rows: rows,
Rows: telemetrystore.WrapRows(rows),
ctx: ctx,
event: event,
onClose: func() { telemetrystore.WrapAfterQuery(p.hooks, ctx, event) },

View File

@@ -0,0 +1,39 @@
package telemetrystore
import (
"reflect"
"strings"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
"github.com/SigNoz/signoz/pkg/types/telemetrystoretypes"
)
// WrapRows reports JSONValue as the scan type of every JSON column. Nested JSON — Array(JSON),
// Map(String, JSON) — is not covered.
func WrapRows(rows driver.Rows) driver.Rows {
return &rowsWithJSONScanType{Rows: rows}
}
type rowsWithJSONScanType struct {
driver.Rows
}
func (r *rowsWithJSONScanType) ColumnTypes() []driver.ColumnType {
colTypes := r.Rows.ColumnTypes()
wrapped := make([]driver.ColumnType, len(colTypes))
for i, colType := range colTypes {
wrapped[i] = colType
if strings.HasPrefix(strings.ToUpper(colType.DatabaseTypeName()), "JSON") {
wrapped[i] = jsonColumnType{ColumnType: colType}
}
}
return wrapped
}
type jsonColumnType struct {
driver.ColumnType
}
func (jsonColumnType) ScanType() reflect.Type {
return reflect.TypeFor[telemetrystoretypes.JSONValue]()
}

View File

@@ -0,0 +1,23 @@
package telemetrystoretest
import (
"context"
"github.com/ClickHouse/clickhouse-go/v2"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
"github.com/SigNoz/signoz/pkg/telemetrystore"
)
// conn wraps rows the way the clickhouse provider does, so mocked JSON columns report the scan
// type they do in production.
type conn struct {
clickhouse.Conn
}
func (c conn) Query(ctx context.Context, query string, args ...any) (driver.Rows, error) {
rows, err := c.Conn.Query(ctx, query, args...)
if err != nil {
return nil, err
}
return telemetrystore.WrapRows(rows), nil
}

View File

@@ -32,7 +32,7 @@ func New(_ telemetrystore.Config, matcher sqlmock.QueryMatcher) *Provider {
// ClickhouseDB returns the mock Clickhouse connection.
func (p *Provider) ClickhouseDB() clickhouse.Conn {
return p.clickhouseDB.(clickhouse.Conn)
return conn{Conn: p.clickhouseDB.(clickhouse.Conn)}
}
// Cluster returns the cluster name.

View File

@@ -3,6 +3,7 @@ package cloudintegrationtypes
import (
"encoding/json"
"fmt"
"maps"
"time"
"github.com/SigNoz/signoz/pkg/errors"
@@ -26,6 +27,17 @@ type Account struct {
type AgentReport struct {
TimestampMillis int64 `json:"timestampMillis" required:"true"`
Data map[string]any `json:"data" required:"true" nullable:"true"`
SyncState *SyncState `json:"syncState" required:"true" nullable:"true"`
}
type SyncState struct {
Version int64 `json:"version" required:"true"`
InSync bool `json:"inSync" required:"true"`
Regions map[string]*RegionSyncState `json:"regions" required:"true" nullable:"false"`
}
type RegionSyncState struct {
State RegionState `json:"state" required:"true"`
}
type AccountConfig struct {
@@ -150,6 +162,7 @@ func NewAccountFromStorable(storableAccount *StorableCloudIntegration) (*Account
account.AgentReport = &AgentReport{
TimestampMillis: storableAccount.LastAgentReport.TimestampMillis,
Data: storableAccount.LastAgentReport.Data,
SyncState: NewSyncStateFromStorable(storableAccount.LastAgentReport.SyncState),
}
}
@@ -308,10 +321,28 @@ func NewAccountConfigFromUpdatable(provider CloudProviderType, config *Updatable
}
}
func NewAgentReport(data map[string]any) *AgentReport {
func NewAgentReport(data map[string]any, syncState *SyncState) *AgentReport {
return &AgentReport{
TimestampMillis: time.Now().UnixMilli(),
Data: data,
SyncState: syncState,
}
}
func NewSyncStateFromStorable(storableSyncState *StorableSyncState) *SyncState {
if storableSyncState == nil {
return nil
}
regions := make(map[string]*RegionSyncState, len(storableSyncState.Regions))
for region, regionSyncState := range storableSyncState.Regions {
regions[region] = &RegionSyncState{State: regionSyncState.State}
}
return &SyncState{
Version: storableSyncState.Version,
InSync: storableSyncState.InSync,
Regions: regions,
}
}
@@ -335,6 +366,40 @@ func (account *Account) Update(provider CloudProviderType, config *AccountConfig
return nil
}
func (account *Account) UpdateAgentReport(providerAccountID *string, agentReport *AgentReport) {
account.ProviderAccountID = providerAccountID
account.AgentReport = agentReport
}
// UpdateSyncState keeps the rest of the agent report, and is a no-op when the agent has never checked in.
func (account *Account) UpdateSyncState(syncState *SyncState) {
if account.AgentReport == nil {
return
}
account.AgentReport.SyncState = syncState
}
// NextSyncState returns the sync state for this check-in, or nil for providers without one.
func (account *Account) NextSyncState(syncedVersion *int64) *SyncState {
if account.Provider != CloudProviderTypeAWS {
return nil
}
var previous *SyncState
if account.AgentReport != nil {
previous = account.AgentReport.SyncState
}
regions := account.Config.AWS.Regions
// Removed before the agent ever checked in: no region was sent to it, so there is nothing to clean up.
if account.AgentReport == nil && account.RemovedAt != nil {
regions = nil
}
return newSyncState(previous, regions, account.RemovedAt != nil, syncedVersion)
}
func (postableAccount *PostableAccount) UnmarshalJSON(data []byte) error {
type Alias PostableAccount
@@ -406,3 +471,79 @@ func (config *AccountConfig) ToJSON() ([]byte, error) {
func NewIngestionKeyName(provider CloudProviderType) string {
return fmt.Sprintf("%s-integration", provider.StringValue())
}
// newSyncState returns the sync state after a check-in without mutating previous.
func newSyncState(previous *SyncState, regions []string, removed bool, syncedVersion *int64) *SyncState {
if previous == nil {
previous = newSyncStateFromRegions(regions)
}
next := previous.copy()
// The agent synced this version, so its disabled regions are cleaned up and can be dropped.
if syncedVersion != nil && *syncedVersion == next.Version {
next.InSync = true
maps.DeleteFunc(next.Regions, func(_ string, regionSyncState *RegionSyncState) bool {
return regionSyncState.State == RegionStateDisabled
})
}
// Once the integration is removed, every region is disabled.
if removed {
regions = nil
}
changed := false
desiredRegionsMap := make(map[string]struct{}, len(regions))
for _, region := range regions {
desiredRegionsMap[region] = struct{}{}
if regionSyncState, ok := next.Regions[region]; ok && regionSyncState.State == RegionStateEnabled {
continue
}
next.Regions[region] = &RegionSyncState{State: RegionStateEnabled}
changed = true
}
for region, regionSyncState := range next.Regions {
_, ok := desiredRegionsMap[region]
if ok && regionSyncState.State == RegionStateEnabled {
continue
}
if !ok && regionSyncState.State == RegionStateDisabled {
continue
}
regionSyncState.State = RegionStateDisabled
changed = true
}
if changed {
next.Version++
next.InSync = false
}
return next
}
// newSyncStateFromRegions is used on the first check-in, when the agent has already deployed regions, so it starts in sync.
func newSyncStateFromRegions(regions []string) *SyncState {
syncState := &SyncState{Version: 1, InSync: true, Regions: make(map[string]*RegionSyncState, len(regions))}
for _, region := range regions {
syncState.Regions[region] = &RegionSyncState{State: RegionStateEnabled}
}
return syncState
}
func (syncState *SyncState) copy() *SyncState {
regions := make(map[string]*RegionSyncState, len(syncState.Regions))
for region, regionSyncState := range syncState.Regions {
regions[region] = &RegionSyncState{State: regionSyncState.State}
}
return &SyncState{Version: syncState.Version, InSync: syncState.InSync, Regions: regions}
}

View File

@@ -12,7 +12,8 @@ type AgentCheckInRequest struct {
ProviderAccountID string `json:"providerAccountId" required:"false"`
CloudIntegrationID valuer.UUID `json:"cloudIntegrationId" required:"false"`
Data map[string]any `json:"data" required:"true" nullable:"true"`
Data map[string]any `json:"data" required:"true" nullable:"true"`
SyncedVersion *int64 `json:"syncedVersion" required:"false" nullable:"true"`
}
type PostableAgentCheckIn struct {
@@ -28,6 +29,7 @@ type AgentCheckInResponse struct {
ProviderAccountID string `json:"providerAccountId" required:"true"`
IntegrationConfig *ProviderIntegrationConfig `json:"integrationConfig" required:"true"`
RemovedAt *time.Time `json:"removedAt" required:"true" nullable:"true"`
SyncState *SyncState `json:"syncState" required:"true" nullable:"true"`
}
type GettableAgentCheckIn struct {
@@ -73,12 +75,13 @@ func NewGettableAgentCheckIn(provider CloudProviderType, resp *AgentCheckInRespo
return gettable
}
func NewAgentCheckInResponse(providerAccountID, cloudIntegrationID string, integrationConfig *ProviderIntegrationConfig, removedAt *time.Time) *AgentCheckInResponse {
func NewAgentCheckInResponse(providerAccountID, cloudIntegrationID string, integrationConfig *ProviderIntegrationConfig, removedAt *time.Time, syncState *SyncState) *AgentCheckInResponse {
return &AgentCheckInResponse{
CloudIntegrationID: cloudIntegrationID,
ProviderAccountID: providerAccountID,
IntegrationConfig: integrationConfig,
RemovedAt: removedAt,
SyncState: syncState,
}
}

View File

@@ -25,6 +25,17 @@ var (
ErrCodeServiceDefinitionNotFound = errors.MustNewCode("service_definition_not_found")
)
var (
RegionStateEnabled = RegionState{valuer.NewString("enabled")}
RegionStateDisabled = RegionState{valuer.NewString("disabled")}
)
type RegionState struct{ valuer.String }
func (RegionState) Enum() []any {
return []any{RegionStateEnabled, RegionStateDisabled}
}
// StorableCloudIntegration represents a cloud integration stored in the database.
// This is also referred as "Account" in the context of cloud integrations.
type StorableCloudIntegration struct {
@@ -43,8 +54,16 @@ type StorableCloudIntegration struct {
// StorableAgentReport represents the last heartbeat and arbitrary data sent by the agent
// as of now there is no use case for Data field, but keeping it for backwards compatibility with older structure.
type StorableAgentReport struct {
TimestampMillis int64 `json:"timestamp_millis"` // backward compatibility
Data map[string]any `json:"data"`
TimestampMillis int64 `json:"timestamp_millis"` // backward compatibility
Data map[string]any `json:"data"`
SyncState *StorableSyncState `json:"sync_state,omitempty"`
}
// StorableSyncState holds every region sent to the agent. A disabled region is dropped only after the agent acks Version.
type StorableSyncState struct {
Version int64 `json:"version"`
InSync bool `json:"in_sync"`
Regions map[string]*RegionSyncState `json:"regions"`
}
// StorableCloudIntegrationService is to store service config for a cloud integration, which is a cloud provider specific configuration.
@@ -148,12 +167,30 @@ func NewStorableCloudIntegration(account *Account) (*StorableCloudIntegration, e
storableAccount.LastAgentReport = &StorableAgentReport{
TimestampMillis: account.AgentReport.TimestampMillis,
Data: account.AgentReport.Data,
SyncState: NewStorableSyncState(account.AgentReport.SyncState),
}
}
return storableAccount, nil
}
func NewStorableSyncState(syncState *SyncState) *StorableSyncState {
if syncState == nil {
return nil
}
regions := make(map[string]*RegionSyncState, len(syncState.Regions))
for region, regionSyncState := range syncState.Regions {
regions[region] = &RegionSyncState{State: regionSyncState.State}
}
return &StorableSyncState{
Version: syncState.Version,
InSync: syncState.InSync,
Regions: regions,
}
}
// NewStorableCloudIntegrationService creates a new StorableCloudIntegrationService with
// generated ID and timestamps from a CloudIntegrationService and its serialized config JSON.
func NewStorableCloudIntegrationService(svc *CloudIntegrationService, configJSON string) *StorableCloudIntegrationService {
@@ -172,6 +209,7 @@ func (account *StorableCloudIntegration) Update(providerAccountID *string, agent
account.LastAgentReport = &StorableAgentReport{
TimestampMillis: agentReport.TimestampMillis,
Data: agentReport.Data,
SyncState: NewStorableSyncState(agentReport.SyncState),
}
}
}

View File

@@ -25,9 +25,12 @@ type Store interface {
// CreateAccount creates a new cloud integration account
CreateAccount(ctx context.Context, account *StorableCloudIntegration) error
// UpdateAccount updates an existing cloud integration account
// UpdateAccount updates the user updatable fields (config) of an existing cloud integration account
UpdateAccount(ctx context.Context, account *StorableCloudIntegration) error
// UpdateAgentReport updates the provider account id and last agent report of an existing cloud integration account
UpdateAgentReport(ctx context.Context, account *StorableCloudIntegration) error
// RemoveAccount marks a cloud integration account as removed by setting the RemovedAt field
RemoveAccount(ctx context.Context, orgID, id valuer.UUID, provider CloudProviderType) error

View File

@@ -1,51 +1,37 @@
package telemetrystoretypes
import (
"encoding/json"
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
"github.com/SigNoz/signoz/pkg/errors"
"github.com/bytedance/sonic"
)
// NestedJSON decodes a native JSON column into a nested document via the driver's own marshaler, so
// arrays of objects and typed sub-paths survive and Dynamic values arrive unwrapped. A key stored as
// both a scalar and an object collapses, as the nested form cannot hold both.
func NestedJSON(j chcol.JSON) map[string]any {
raw, err := j.MarshalJSON()
if err != nil {
return nil
}
var out map[string]any
if err := sonic.Unmarshal(raw, &out); err != nil {
return nil
}
return out
}
var ErrCodeUnmarshalJSONColumn = errors.MustNewCode("fail_unmarshal_json_column")
// FlattenJSON decodes a native JSON column into its leaf paths as dotted keys, so a key stored as
// both a scalar and an object survives as two distinct keys — unlike the nested form, which cannot
// hold both. Dynamic values arrive unwrapped, arrays of objects intact.
func FlattenJSON(j chcol.JSON) map[string]any {
paths := j.ValuesByPath()
out := make(map[string]any, len(paths))
for path, value := range paths {
out[path] = decodePathValue(value)
}
return out
}
// JSONValue is the scan target for a ClickHouse JSON column: the connection sets
// output_format_native_write_json_as_string, so the column arrives as a raw document rather than
// the chcol.JSON the driver reports as its scan type.
type JSONValue map[string]any
func decodePathValue(value any) any {
variant, ok := value.(chcol.Variant)
if !ok {
return value
}
raw, err := json.Marshal(variant)
if err != nil {
// Scan decodes into a fresh map every time: a scan target is reused across rows, and unmarshalling
// into the map already there would both keep its keys and hand every row the same map.
func (v *JSONValue) Scan(src any) error {
var raw []byte
switch value := src.(type) {
case nil:
*v = nil
return nil
case string:
raw = []byte(value)
case []byte:
raw = value
default:
return errors.NewInternalf(ErrCodeUnmarshalJSONColumn, "cannot decode %T as a JSON column", src)
}
var out any
if err := sonic.Unmarshal(raw, &out); err != nil {
return nil
decoded := JSONValue{}
if err := sonic.Unmarshal(raw, &decoded); err != nil {
return errors.WrapInternalf(err, ErrCodeUnmarshalJSONColumn, "failed to unmarshal JSON column")
}
return out
*v = decoded
return nil
}

View File

@@ -1,82 +0,0 @@
package telemetrystoretypes
import (
"testing"
"github.com/ClickHouse/clickhouse-go/v2/lib/chcol"
"github.com/stretchr/testify/assert"
)
func TestNestedJSON(t *testing.T) {
testCases := []struct {
name string
paths map[string]any
want map[string]any
}{
{
name: "Empty",
paths: nil,
want: map[string]any{},
},
{
name: "FlatScalars",
paths: map[string]any{"level": "error", "status": int64(500)},
want: map[string]any{"level": "error", "status": float64(500)},
},
{
name: "DottedPathsBecomeNested",
paths: map[string]any{"attrs.code": int64(500), "attrs.path": "/checkout"},
want: map[string]any{"attrs": map[string]any{"code": float64(500), "path": "/checkout"}},
},
{
name: "ArrayOfObjectsPreserved",
paths: map[string]any{"education": []any{map[string]any{"name": "IIT"}}},
want: map[string]any{"education": []any{map[string]any{"name": "IIT"}}},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
j := chcol.NewJSON()
for path, value := range testCase.paths {
j.SetValueAtPath(path, value)
}
assert.Equal(t, testCase.want, NestedJSON(*j))
})
}
}
func TestFlattenJSON(t *testing.T) {
testCases := []struct {
name string
paths map[string]any
want map[string]any
}{
{
name: "Empty",
paths: nil,
want: map[string]any{},
},
{
name: "DottedPathsStayFlat",
paths: map[string]any{"http.method": "GET", "level": "error"},
want: map[string]any{"http.method": "GET", "level": "error"},
},
{
// A scalar and an object under the same prefix survive as two distinct dotted keys.
name: "ScalarAndObjectKey_BothSurvive",
paths: map[string]any{"scope": "x", "scope.attributes.name": "y"},
want: map[string]any{"scope": "x", "scope.attributes.name": "y"},
},
}
for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
j := chcol.NewJSON()
for path, value := range testCase.paths {
j.SetValueAtPath(path, value)
}
assert.Equal(t, testCase.want, FlattenJSON(*j))
})
}
}

View File

@@ -34,6 +34,8 @@ class ProviderAccountSpec:
expected_config: Callable[[dict], dict]
# only the suites that exercise updates need to supply it.
updated_params: dict = field(default_factory=dict)
# params -> the agentReport.syncState the API is expected to return after the first check-in.
expected_sync_state: Callable[[dict], dict | None] = lambda p: None
# id shown in parametrized test names; defaults to the provider slug.
id: str = field(default="")
@@ -315,6 +317,7 @@ def simulate_agent_checkin(
account_id: str,
cloud_account_id: str,
data: dict | None = None,
synced_version: int | None = None,
) -> requests.Response:
endpoint = f"/api/v1/cloud_integrations/{cloud_provider}/accounts/check_in"
@@ -323,6 +326,8 @@ def simulate_agent_checkin(
"providerAccountId": cloud_account_id,
"data": data or {},
}
if synced_version is not None:
checkin_payload["syncedVersion"] = synced_version
response = requests.post(
signoz.self.host_configs["8080"].get(endpoint),

View File

@@ -3,6 +3,7 @@ from collections.abc import Callable
from http import HTTPStatus
import pytest
import requests
from fixtures import types
from fixtures.auth import USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD, add_license
@@ -152,3 +153,230 @@ def test_duplicate_cloud_account_checkins(
# Second check-in: account2 tries to claim the same provider account ID → 409
response = simulate_agent_checkin(signoz, admin_token, spec.provider, account2["id"], same_provider_account_id)
assert response.status_code == HTTPStatus.CONFLICT, f"Expected 409 for duplicate providerAccountId, got {response.status_code}: {response.text}"
def test_sync_state_drops_removed_region_after_ack(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_cloud_integration_account: Callable,
) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
provider_account_id = str(uuid.uuid4())
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
headers={"Authorization": f"Bearer {admin_token}"},
json={"config": {"aws": {"regions": ["us-east-1"]}}},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {
"version": 2,
"inSync": False,
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}},
}, "removed region should be marked disabled and the version bumped"
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=2)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {
"version": 2,
"inSync": True,
"regions": {"us-east-1": {"state": "enabled"}},
}, "acked removed region should be dropped"
def test_sync_state_keeps_removed_region_without_ack(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_cloud_integration_account: Callable,
) -> None:
"""The agent failed to clean up or crashed, so it never acks: the removed region stays and the version stays put."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
provider_account_id = str(uuid.uuid4())
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
headers={"Authorization": f"Bearer {admin_token}"},
json={"config": {"aws": {"regions": ["us-east-1"]}}},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
for _ in range(3):
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {
"version": 2,
"inSync": False,
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}},
}, "unacked removed region should stay without bumping the version"
@pytest.mark.parametrize("synced_version", [2, 9], ids=["stale", "ahead"])
def test_sync_state_ignores_mismatched_ack(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_cloud_integration_account: Callable,
synced_version: int,
) -> None:
"""An ack for any version other than the current one (v3) is ignored,
so us-west-2, removed at v2 and still unacked, is not dropped.
"""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
provider_account_id = str(uuid.uuid4())
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
for regions in (["us-east-1"], ["us-east-1", "eu-west-1"]):
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
headers={"Authorization": f"Bearer {admin_token}"},
json={"config": {"aws": {"regions": regions}}},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
expected_sync_state = {
"version": 3,
"inSync": False,
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}, "eu-west-1": {"state": "enabled"}},
}
assert response.json()["data"]["syncState"] == expected_sync_state
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=synced_version)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == expected_sync_state, "an ack for another version should be ignored"
def test_sync_state_applies_ack_before_config_change(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_cloud_integration_account: Callable,
) -> None:
"""The user changes regions while the agent syncs: the ack for the version it synced still lands."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
provider_account_id = str(uuid.uuid4())
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
for regions, synced_version in ((["us-east-1"], None), (["us-east-1", "eu-west-1"], 2)):
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
headers={"Authorization": f"Bearer {admin_token}"},
json={"config": {"aws": {"regions": regions}}},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=synced_version)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {
"version": 3,
"inSync": False,
"regions": {"us-east-1": {"state": "enabled"}, "eu-west-1": {"state": "enabled"}},
}, "ack should drop the removed region before the new region bumps the version"
@pytest.mark.parametrize("synced_version", [1, None], ids=["agent_acks_synced_version", "agent_crashed"])
def test_sync_state_region_removed_during_sync(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_cloud_integration_account: Callable,
synced_version: int | None,
) -> None:
"""The user removes a region while the agent syncs v1; whether the agent acks v1 or crashed, the region must not be lost."""
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
provider_account_id = str(uuid.uuid4())
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {
"version": 1,
"inSync": True,
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "enabled"}},
}
response = requests.put(
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
headers={"Authorization": f"Bearer {admin_token}"},
json={"config": {"aws": {"regions": ["us-east-1"]}}},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=synced_version)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {
"version": 2,
"inSync": False,
"regions": {"us-east-1": {"state": "enabled"}, "us-west-2": {"state": "disabled"}},
}, "region removed mid-sync should be marked disabled"
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=2)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {
"version": 2,
"inSync": True,
"regions": {"us-east-1": {"state": "enabled"}},
}
def test_sync_state_after_disconnect(
signoz: types.SigNoz,
create_user_admin: types.Operation, # pylint: disable=unused-argument
get_token: Callable[[str, str], str],
create_cloud_integration_account: Callable,
) -> None:
admin_token = get_token(USER_ADMIN_EMAIL, USER_ADMIN_PASSWORD)
account_id = create_cloud_integration_account(admin_token, "aws", regions=["us-east-1", "us-west-2"])["id"]
provider_account_id = str(uuid.uuid4())
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
response = requests.delete(
signoz.self.host_configs["8080"].get(f"/api/v1/cloud_integrations/aws/accounts/{account_id}"),
headers={"Authorization": f"Bearer {admin_token}"},
timeout=10,
)
assert response.status_code == HTTPStatus.NO_CONTENT, response.text
for _ in range(2):
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["removedAt"] is not None, "removedAt should be set after disconnect"
assert response.json()["data"]["syncState"] == {
"version": 2,
"inSync": False,
"regions": {"us-east-1": {"state": "disabled"}, "us-west-2": {"state": "disabled"}},
}, "every region should be disabled once, without bumping the version on later check-ins"
for _ in range(2):
response = simulate_agent_checkin(signoz, admin_token, "aws", account_id, provider_account_id, synced_version=2)
assert response.status_code == HTTPStatus.OK, response.text
assert response.json()["data"]["syncState"] == {"version": 2, "inSync": True, "regions": {}}, "acked removal should leave no regions"

View File

@@ -21,6 +21,11 @@ AWS_ACCOUNT_SPEC = ProviderAccountSpec(
updated_params={"deployment_region": "us-east-1", "regions": ["us-east-1", "us-west-2", "eu-west-1"]},
build_config=lambda p: {"aws": {"deploymentRegion": p["deployment_region"], "regions": p["regions"]}},
expected_config=lambda p: {"regions": p["regions"]},
expected_sync_state=lambda p: {
"version": 1,
"inSync": True,
"regions": {region: {"state": "enabled"} for region in p["regions"]},
},
)
GCP_ACCOUNT_SPEC = ProviderAccountSpec(
@@ -128,6 +133,7 @@ def test_list_accounts_after_checkin(
assert found["providerAccountId"] == provider_account_id, "providerAccountId should match"
assert found["config"][spec.provider] == spec.expected_config(spec.initial_params), "config should match account config"
assert found["agentReport"] is not None, "agentReport should be present after check-in"
assert found["agentReport"]["syncState"] == spec.expected_sync_state(spec.initial_params), "syncState should be seeded from the account regions on first check-in"
assert found["removedAt"] is None, "removedAt should be null for a live account"
@@ -282,6 +288,7 @@ def test_update_account_after_checkin_preserves_connected_status(
assert found_after is not None, "Account must still be listed after config update (account_id should not be reset)"
assert found_after["providerAccountId"] == provider_account_id, "providerAccountId should be preserved after update"
assert found_after["agentReport"] is not None, "agentReport should be preserved after update"
assert found_after["agentReport"]["syncState"] == found_before["agentReport"]["syncState"], "config update must not change syncState"
assert found_after["config"][spec.provider] == spec.expected_config(spec.updated_params), "Config should reflect the update"
assert found_after["removedAt"] is None, "removedAt should still be null"