Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions core/cmd/shell.go
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,7 @@ func newBeholderClient(
ChipIngressEmitterGRPCEndpoint: cfgTelemetry.ChipIngressEndpoint(),
ChipIngressInsecureConnection: cfgTelemetry.ChipIngressInsecureConnection(),
ChipIngressBatchEmitterEnabled: cfgTelemetry.ChipIngressBatchEmitterEnabled(),
ChipIngressMaxMessageBufferBytes: cfgTelemetry.ChipIngressMaxMessageBufferBytes(),
ChipIngressLogger: lggr,
LogStreamingEnabled: cfgTelemetry.LogStreamingEnabled(),
LogLevel: cfgTelemetry.LogLevel(),
Expand Down
2 changes: 2 additions & 0 deletions core/config/docs/core.toml
Original file line number Diff line number Diff line change
Expand Up @@ -889,6 +889,8 @@ ChipIngressInsecureConnection = false # Default
# ChipIngressBatchEmitterEnabled enables batching for chip-ingress events.
# When false, events are sent individually (legacy behavior).
ChipIngressBatchEmitterEnabled = true # Default
# ChipIngressMaxMessageBufferBytes is the max live byte size of messages buffered in the chip-ingress batch client before drop (default 1 GiB in batch client when unset).
ChipIngressMaxMessageBufferBytes = 1073741824 # Default
# DurableEmitterEnabled enables persisting outbound CHIP events to Postgres for at-least-once delivery.
DurableEmitterEnabled = false # Default
# HeartbeatInterval is the interval at which a the application heartbeat is sent to telemetry backends.
Expand Down
1 change: 1 addition & 0 deletions core/config/telemetry_config.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ type Telemetry interface {
ChipIngressEndpoint() string
ChipIngressInsecureConnection() bool
ChipIngressBatchEmitterEnabled() bool
ChipIngressMaxMessageBufferBytes() uint
DurableEmitterEnabled() bool
HeartbeatInterval() time.Duration
LogStreamingEnabled() bool
Expand Down
7 changes: 7 additions & 0 deletions core/config/toml/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -2985,6 +2985,7 @@ type Telemetry struct {
ChipIngressEndpoint *string
ChipIngressInsecureConnection *bool
ChipIngressBatchEmitterEnabled *bool
ChipIngressMaxMessageBufferBytes *uint
DurableEmitterEnabled *bool
HeartbeatInterval *commonconfig.Duration
LogLevel *string
Expand Down Expand Up @@ -3035,6 +3036,9 @@ func (b *Telemetry) setFrom(f *Telemetry) {
if v := f.ChipIngressBatchEmitterEnabled; v != nil {
b.ChipIngressBatchEmitterEnabled = v
}
if v := f.ChipIngressMaxMessageBufferBytes; v != nil {
b.ChipIngressMaxMessageBufferBytes = v
}
if v := f.DurableEmitterEnabled; v != nil {
b.DurableEmitterEnabled = v
}
Expand Down Expand Up @@ -3081,6 +3085,9 @@ func (b *Telemetry) ValidateConfig() (err error) {
if ratio := b.TraceSampleRatio; ratio != nil && (*ratio < 0 || *ratio > 1) {
err = errors.Join(err, configutils.ErrInvalid{Name: "TraceSampleRatio", Value: *ratio, Msg: "must be between 0 and 1"})
}
if v := b.ChipIngressMaxMessageBufferBytes; v != nil && *v == 0 {
err = errors.Join(err, configutils.ErrInvalid{Name: "ChipIngressMaxMessageBufferBytes", Value: *v, Msg: "must be greater than 0"})
}
return err
}

Expand Down
3 changes: 3 additions & 0 deletions core/services/chainlink/application.go
Original file line number Diff line number Diff line change
Expand Up @@ -423,6 +423,9 @@ func NewApplication(ctx context.Context, opts ApplicationOpts) (Application, err
MaxConcurrentSends: 8,
MessageBufferSize: 50_000,
}
if maxBuf := int(cfg.Telemetry().ChipIngressMaxMessageBufferBytes()); maxBuf > 0 {
durableCfg.MaxMessageBufferBytes = maxBuf
}
pgStore := durableemitter.NewPgDurableEventStore(opts.DS)
durableEmitter, setupErr := durableemitter.Setup(pgStore, durableCfg, globalLogger)
if setupErr != nil {
Expand Down
7 changes: 7 additions & 0 deletions core/services/chainlink/config_telemetry.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,13 @@ func (b *telemetryConfig) ChipIngressBatchEmitterEnabled() bool {
return *b.s.ChipIngressBatchEmitterEnabled
}

func (b *telemetryConfig) ChipIngressMaxMessageBufferBytes() uint {
if b.s.ChipIngressMaxMessageBufferBytes == nil {
return 0
}
return *b.s.ChipIngressMaxMessageBufferBytes
}

func (b *telemetryConfig) DurableEmitterEnabled() bool {
if b.s.DurableEmitterEnabled == nil {
return false
Expand Down
18 changes: 18 additions & 0 deletions core/services/chainlink/config_telemetry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -237,6 +237,24 @@ func TestTelemetryConfig_ChipIngressBatchEmitterEnabled(t *testing.T) {
}
}

func TestTelemetryConfig_ChipIngressMaxMessageBufferBytes(t *testing.T) {
tests := []struct {
name string
telemetry toml.Telemetry
expected uint
}{
{"ChipIngressMaxMessageBufferBytesSet", toml.Telemetry{ChipIngressMaxMessageBufferBytes: ptr[uint](1073741824)}, 1073741824},
{"ChipIngressMaxMessageBufferBytesNil", toml.Telemetry{ChipIngressMaxMessageBufferBytes: nil}, 0},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
tc := telemetryConfig{s: tt.telemetry}
assert.Equal(t, tt.expected, tc.ChipIngressMaxMessageBufferBytes())
})
}
}

func ptrDuration(d time.Duration) *config.Duration {
return config.MustNewDuration(d)
}
Expand Down
7 changes: 7 additions & 0 deletions docs/CONFIG.md
Original file line number Diff line number Diff line change
Expand Up @@ -2393,6 +2393,7 @@ AuthHeadersTTL = '0s' # Default
ChipIngressEndpoint = '' # Default
ChipIngressInsecureConnection = false # Default
ChipIngressBatchEmitterEnabled = true # Default
ChipIngressMaxMessageBufferBytes = 1073741824 # Default
DurableEmitterEnabled = false # Default
HeartbeatInterval = '1s' # Default
LogLevel = "info" # Default
Expand Down Expand Up @@ -2477,6 +2478,12 @@ ChipIngressBatchEmitterEnabled = true # Default
ChipIngressBatchEmitterEnabled enables batching for chip-ingress events.
When false, events are sent individually (legacy behavior).

### ChipIngressMaxMessageBufferBytes
```toml
ChipIngressMaxMessageBufferBytes = 1073741824 # Default
```
ChipIngressMaxMessageBufferBytes is the max live byte size of messages buffered in the chip-ingress batch client before drop (default 1 GiB in batch client when unset).

### DurableEmitterEnabled
```toml
DurableEmitterEnabled = false # Default
Expand Down
4 changes: 4 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -437,3 +437,7 @@ tool github.com/smartcontractkit/chainlink-common/pkg/loop/cmd/loopinstall
tool github.com/smartcontractkit/chainlink-common/script/cmd/dependabot

replace github.com/smartcontractkit/chainlink-sui => github.com/smartcontractkit/chainlink-sui v0.0.0-20260707125635-abec997b6eae

replace github.com/smartcontractkit/chainlink-common => ../../chainlink-common/chipingress-buffer-byte-cap

replace github.com/smartcontractkit/chainlink-common/pkg/chipingress => ../../chainlink-common/chipingress-buffer-byte-cap/pkg/chipingress
4 changes: 0 additions & 4 deletions go.sum

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions plugins/loop_registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,7 @@ func (m *LoopRegistry) Register(id string) (*RegisteredLoop, error) {
envCfg.ChipIngressEndpoint = m.cfgTelemetry.ChipIngressEndpoint()
envCfg.ChipIngressInsecureConnection = m.cfgTelemetry.ChipIngressInsecureConnection()
envCfg.ChipIngressBatchEmitterEnabled = m.cfgTelemetry.ChipIngressBatchEmitterEnabled()
envCfg.ChipIngressMaxMessageBufferBytes = m.cfgTelemetry.ChipIngressMaxMessageBufferBytes()
envCfg.ChipIngressDurableEmitterEnabled = m.cfgTelemetry.DurableEmitterEnabled()
envCfg.TelemetryLogStreamingEnabled = m.cfgTelemetry.LogStreamingEnabled()
envCfg.TelemetryLogLevel = m.cfgTelemetry.LogLevel()
Expand Down
Loading