Skip to content

Commit b2a1f29

Browse files
committed
telemetry: add timestamp index companion account for latency samples
Add a TimestampIndex account that tracks sample index → timestamp mappings as a companion to device and internet latency sample accounts. - New onchain account type (AccountTypeTimestampIndex = 5) with PDA derived from samples account PK - InitializeTimestampIndex and WriteLatencySamples instructions updated to optionally append index entries - Read-only deserialization in Go, Python, and TypeScript SDKs with fixture generation and cross-language compat tests - Timestamp reconstruction helpers using binary search and single-pass for efficient sample timestamp lookups - Timestamp index full condition is non-fatal in write path
1 parent 5e339e8 commit b2a1f29

44 files changed

Lines changed: 2473 additions & 13 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

CHANGELOG.md

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -326,6 +326,9 @@ All notable changes to this project will be documented in this file.
326326
- Controller
327327
- Retry transient Solana RPC failures when fetching onchain serviceability accounts so controller polls are more resilient to short-lived provider resets
328328

329+
- Telemetry
330+
- Add timestamp index companion account for device and internet latency samples, enabling reliable timestamp reconstruction when agents experience downtime gaps within an epoch
331+
329332
## [v0.12.0](https://github.com/malbeclabs/doublezero/compare/client/v0.11.0...client/v0.12.0) - 2026-03-16
330333

331334
### Breaking

controlplane/internet-latency-collector/internal/exporter/ledger.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,9 @@ type ServiceabilityProgramClient interface {
3737
}
3838

3939
type TelemetryProgramClient interface {
40+
ProgramID() solana.PublicKey
4041
InitializeInternetLatencySamples(ctx context.Context, config telemetry.InitializeInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
42+
InitializeTimestampIndex(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error)
4143
WriteInternetLatencySamples(ctx context.Context, config telemetry.WriteInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
4244
}
4345

controlplane/internet-latency-collector/internal/exporter/main_test.go

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,6 +52,18 @@ type mockTelemetryProgramClient struct {
5252
InitializeInternetLatencySamplesFunc func(ctx context.Context, config telemetry.InitializeInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
5353
WriteInternetLatencySamplesFunc func(ctx context.Context, config telemetry.WriteInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
5454
GetInternetLatencySamplesFunc func(ctx context.Context, dataProviderName string, originExchangePK solana.PublicKey, targetExchangePK solana.PublicKey, epoch uint64) (*telemetry.InternetLatencySamples, error)
55+
InitializeTimestampIndexFunc func(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error)
56+
}
57+
58+
func (c *mockTelemetryProgramClient) ProgramID() solana.PublicKey {
59+
return solana.MustPublicKeyFromBase58("11111111111111111111111111111111")
60+
}
61+
62+
func (c *mockTelemetryProgramClient) InitializeTimestampIndex(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error) {
63+
if c.InitializeTimestampIndexFunc != nil {
64+
return c.InitializeTimestampIndexFunc(ctx, samplesAccountPK)
65+
}
66+
return solana.Signature{}, nil, nil
5567
}
5668

5769
func (c *mockTelemetryProgramClient) InitializeInternetLatencySamples(ctx context.Context, config telemetry.InitializeInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {

controlplane/internet-latency-collector/internal/exporter/submitter.go

Lines changed: 46 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -117,18 +117,56 @@ func (s *Submitter) SubmitSamples(ctx context.Context, partitionKey PartitionKey
117117
}
118118
}
119119

120+
// Derive the samples PDA so we can derive the timestamp index PDA from it.
121+
samplesPDA, _, err := telemetry.DeriveInternetLatencySamplesPDA(
122+
s.cfg.Telemetry.ProgramID(),
123+
s.cfg.OracleAgentPK,
124+
string(partitionKey.DataProvider),
125+
partitionKey.SourceExchangePK,
126+
partitionKey.TargetExchangePK,
127+
partitionKey.Epoch,
128+
)
129+
if err != nil {
130+
return fmt.Errorf("failed to derive internet latency samples PDA: %w", err)
131+
}
132+
timestampIndexPDA, _, err := telemetry.DeriveTimestampIndexPDA(
133+
s.cfg.Telemetry.ProgramID(),
134+
samplesPDA,
135+
)
136+
if err != nil {
137+
return fmt.Errorf("failed to derive timestamp index PDA: %w", err)
138+
}
139+
120140
writeConfig := telemetry.WriteInternetLatencySamplesInstructionConfig{
121141
DataProviderName: string(partitionKey.DataProvider),
122142
OriginExchangePK: partitionKey.SourceExchangePK,
123143
TargetExchangePK: partitionKey.TargetExchangePK,
124144
Epoch: partitionKey.Epoch,
125145
StartTimestampMicroseconds: uint64(minTimestamp.UnixMicro()),
126146
Samples: rtts,
147+
TimestampIndexPK: &timestampIndexPDA,
127148
}
128149

129-
_, _, err := s.cfg.Telemetry.WriteInternetLatencySamples(ctx, writeConfig)
150+
_, _, err = s.cfg.Telemetry.WriteInternetLatencySamples(ctx, writeConfig)
130151
if err != nil {
131-
if errors.Is(err, telemetry.ErrAccountNotFound) {
152+
if errors.Is(err, telemetry.ErrTimestampIndexNotFound) {
153+
log.Info("Timestamp index account not found, initializing")
154+
_, _, err = s.cfg.Telemetry.InitializeTimestampIndex(ctx, samplesPDA)
155+
if err != nil {
156+
log.Warn("Failed to initialize timestamp index, writes will proceed without it", "error", err)
157+
writeConfig.TimestampIndexPK = nil
158+
}
159+
_, _, err = s.cfg.Telemetry.WriteInternetLatencySamples(ctx, writeConfig)
160+
if err != nil {
161+
if errors.Is(err, telemetry.ErrSamplesAccountFull) {
162+
log.Warn("Partition account is full, dropping samples from buffer and moving on", "droppedSamples", len(samples))
163+
metrics.ExporterSubmitterAccountFull.WithLabelValues(string(partitionKey.DataProvider), partitionKey.SourceExchangePK.String(), partitionKey.TargetExchangePK.String(), strconv.FormatUint(partitionKey.Epoch, 10)).Inc()
164+
s.cfg.Buffer.Remove(partitionKey)
165+
return nil
166+
}
167+
return fmt.Errorf("failed to write internet latency samples after timestamp index init: %w", err)
168+
}
169+
} else if errors.Is(err, telemetry.ErrAccountNotFound) {
132170
log.Info("Account not found, initializing new account")
133171
samplingInterval, ok := s.cfg.DataProviderSamplingIntervals[partitionKey.DataProvider]
134172
if !ok {
@@ -144,6 +182,12 @@ func (s *Submitter) SubmitSamples(ctx context.Context, partitionKey PartitionKey
144182
if err != nil {
145183
return fmt.Errorf("failed to initialize internet latency samples: %w", err)
146184
}
185+
// Initialize the companion timestamp index account.
186+
_, _, err = s.cfg.Telemetry.InitializeTimestampIndex(ctx, samplesPDA)
187+
if err != nil {
188+
log.Warn("Failed to initialize timestamp index, writes will proceed without it", "error", err)
189+
writeConfig.TimestampIndexPK = nil
190+
}
147191
_, _, err = s.cfg.Telemetry.WriteInternetLatencySamples(ctx, writeConfig)
148192
if err != nil {
149193
if errors.Is(err, telemetry.ErrSamplesAccountFull) {

controlplane/internet-latency-collector/internal/exporter/submitter_test.go

Lines changed: 150 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -571,6 +571,156 @@ func TestInternetLatency_Submitter(t *testing.T) {
571571
assert.False(t, buffer.Has(key), "partition key should be removed after account full error")
572572
})
573573

574+
t.Run("initializes_only_timestamp_index_when_timestamp_index_not_found", func(t *testing.T) {
575+
t.Parallel()
576+
577+
log := logger.With("test", t.Name())
578+
579+
key := newTestPartitionKey()
580+
sample := newTestSample()
581+
582+
var initSamplesCalled, initTimestampIndexCalled, writeCalled int32
583+
telemetryProgram := &mockTelemetryProgramClient{
584+
WriteInternetLatencySamplesFunc: func(ctx context.Context, config sdktelemetry.WriteInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {
585+
if atomic.AddInt32(&writeCalled, 1) == 1 {
586+
return solana.Signature{}, nil, sdktelemetry.ErrTimestampIndexNotFound
587+
}
588+
return solana.Signature{}, nil, nil
589+
},
590+
InitializeInternetLatencySamplesFunc: func(ctx context.Context, config sdktelemetry.InitializeInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {
591+
atomic.AddInt32(&initSamplesCalled, 1)
592+
return solana.Signature{}, nil, nil
593+
},
594+
InitializeTimestampIndexFunc: func(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error) {
595+
atomic.AddInt32(&initTimestampIndexCalled, 1)
596+
return solana.Signature{}, nil, nil
597+
},
598+
}
599+
600+
buffer := buffer.NewMemoryPartitionedBuffer[exporter.PartitionKey, exporter.Sample](128)
601+
buffer.Add(key, sample)
602+
603+
submitter, err := exporter.NewSubmitter(log, &exporter.SubmitterConfig{
604+
OracleAgentPK: solana.NewWallet().PublicKey(),
605+
Interval: time.Hour,
606+
Buffer: buffer,
607+
Telemetry: telemetryProgram,
608+
MaxAttempts: 3,
609+
BackoffFunc: func(_ int) time.Duration { return 0 },
610+
EpochFinder: &mockEpochFinder{ApproximateAtTimeFunc: func(ctx context.Context, target time.Time) (uint64, error) {
611+
return key.Epoch, nil
612+
}},
613+
DataProviderSamplingIntervals: map[exporter.DataProviderName]time.Duration{
614+
key.DataProvider: time.Second,
615+
},
616+
})
617+
require.NoError(t, err)
618+
619+
submitter.Tick(t.Context())
620+
621+
assert.Equal(t, int32(0), atomic.LoadInt32(&initSamplesCalled), "should not initialize samples account")
622+
assert.Equal(t, int32(1), atomic.LoadInt32(&initTimestampIndexCalled), "should initialize timestamp index")
623+
assert.Equal(t, int32(2), atomic.LoadInt32(&writeCalled), "should try write twice (before and after timestamp index init)")
624+
})
625+
626+
t.Run("clears_timestamp_index_pk_when_init_fails", func(t *testing.T) {
627+
t.Parallel()
628+
629+
log := logger.With("test", t.Name())
630+
631+
key := newTestPartitionKey()
632+
sample := newTestSample()
633+
634+
var writeCalled int32
635+
var retryTimestampIndexPK *solana.PublicKey
636+
telemetryProgram := &mockTelemetryProgramClient{
637+
WriteInternetLatencySamplesFunc: func(ctx context.Context, config sdktelemetry.WriteInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {
638+
n := atomic.AddInt32(&writeCalled, 1)
639+
if n == 1 {
640+
return solana.Signature{}, nil, sdktelemetry.ErrTimestampIndexNotFound
641+
}
642+
retryTimestampIndexPK = config.TimestampIndexPK
643+
return solana.Signature{}, nil, nil
644+
},
645+
InitializeTimestampIndexFunc: func(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error) {
646+
return solana.Signature{}, nil, errors.New("init failed")
647+
},
648+
}
649+
650+
buffer := buffer.NewMemoryPartitionedBuffer[exporter.PartitionKey, exporter.Sample](128)
651+
buffer.Add(key, sample)
652+
653+
submitter, err := exporter.NewSubmitter(log, &exporter.SubmitterConfig{
654+
OracleAgentPK: solana.NewWallet().PublicKey(),
655+
Interval: time.Hour,
656+
Buffer: buffer,
657+
Telemetry: telemetryProgram,
658+
MaxAttempts: 2,
659+
BackoffFunc: func(_ int) time.Duration { return 0 },
660+
EpochFinder: &mockEpochFinder{ApproximateAtTimeFunc: func(ctx context.Context, target time.Time) (uint64, error) {
661+
return key.Epoch, nil
662+
}},
663+
})
664+
require.NoError(t, err)
665+
666+
submitter.Tick(t.Context())
667+
668+
assert.Equal(t, int32(2), atomic.LoadInt32(&writeCalled), "should retry write after failed timestamp index init")
669+
assert.Nil(t, retryTimestampIndexPK, "retry write should have nil TimestampIndexPK after failed init")
670+
})
671+
672+
t.Run("clears_timestamp_index_pk_when_init_fails_on_new_account", func(t *testing.T) {
673+
t.Parallel()
674+
675+
log := logger.With("test", t.Name())
676+
677+
key := newTestPartitionKey()
678+
sample := newTestSample()
679+
680+
var writeCalled int32
681+
var retryTimestampIndexPK *solana.PublicKey
682+
telemetryProgram := &mockTelemetryProgramClient{
683+
WriteInternetLatencySamplesFunc: func(ctx context.Context, config sdktelemetry.WriteInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {
684+
n := atomic.AddInt32(&writeCalled, 1)
685+
if n == 1 {
686+
return solana.Signature{}, nil, sdktelemetry.ErrAccountNotFound
687+
}
688+
retryTimestampIndexPK = config.TimestampIndexPK
689+
return solana.Signature{}, nil, nil
690+
},
691+
InitializeInternetLatencySamplesFunc: func(ctx context.Context, config sdktelemetry.InitializeInternetLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {
692+
return solana.Signature{}, nil, nil
693+
},
694+
InitializeTimestampIndexFunc: func(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error) {
695+
return solana.Signature{}, nil, errors.New("init failed")
696+
},
697+
}
698+
699+
buffer := buffer.NewMemoryPartitionedBuffer[exporter.PartitionKey, exporter.Sample](128)
700+
buffer.Add(key, sample)
701+
702+
submitter, err := exporter.NewSubmitter(log, &exporter.SubmitterConfig{
703+
OracleAgentPK: solana.NewWallet().PublicKey(),
704+
Interval: time.Hour,
705+
Buffer: buffer,
706+
Telemetry: telemetryProgram,
707+
MaxAttempts: 2,
708+
BackoffFunc: func(_ int) time.Duration { return 0 },
709+
EpochFinder: &mockEpochFinder{ApproximateAtTimeFunc: func(ctx context.Context, target time.Time) (uint64, error) {
710+
return key.Epoch, nil
711+
}},
712+
DataProviderSamplingIntervals: map[exporter.DataProviderName]time.Duration{
713+
key.DataProvider: time.Second,
714+
},
715+
})
716+
require.NoError(t, err)
717+
718+
submitter.Tick(t.Context())
719+
720+
assert.Equal(t, int32(2), atomic.LoadInt32(&writeCalled), "should retry write after failed timestamp index init")
721+
assert.Nil(t, retryTimestampIndexPK, "retry write should have nil TimestampIndexPK after failed init")
722+
})
723+
574724
t.Run("failed_retries_reinsert_at_front_preserving_order", func(t *testing.T) {
575725
t.Parallel()
576726

controlplane/telemetry/internal/telemetry/config_test.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,14 @@ func (m *mockPeerDiscovery) GetPeers() []*Peer {
5151

5252
type mockTelemetryProgramClient struct{}
5353

54+
func (m *mockTelemetryProgramClient) ProgramID() solana.PublicKey {
55+
return solana.MustPublicKeyFromBase58("11111111111111111111111111111111")
56+
}
57+
58+
func (m *mockTelemetryProgramClient) InitializeTimestampIndex(_ context.Context, _ solana.PublicKey) (solana.Signature, *rpc.GetTransactionResult, error) {
59+
return solana.Signature{}, nil, nil
60+
}
61+
5462
func (m *mockTelemetryProgramClient) InitializeDeviceLatencySamples(ctx context.Context, config telemetryprog.InitializeDeviceLatencySamplesInstructionConfig) (solana.Signature, *rpc.GetTransactionResult, error) {
5563
return solana.Signature{}, nil, nil
5664
}

controlplane/telemetry/internal/telemetry/ledger.go

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,9 +20,15 @@ type ServiceabilityProgramClient interface {
2020

2121
// TelemetryProgramClient is the client to the telemetry program.
2222
type TelemetryProgramClient interface {
23+
// ProgramID returns the telemetry program ID.
24+
ProgramID() solana.PublicKey
25+
2326
// InitializeDeviceLatencySamples initializes the device latency samples account.
2427
InitializeDeviceLatencySamples(ctx context.Context, config telemetry.InitializeDeviceLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
2528

29+
// InitializeTimestampIndex initializes a timestamp index companion account.
30+
InitializeTimestampIndex(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error)
31+
2632
// WriteDeviceLatencySamples writes the device latency samples to the account.
2733
WriteDeviceLatencySamples(ctx context.Context, config telemetry.WriteDeviceLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
2834
}

controlplane/telemetry/internal/telemetry/main_test.go

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,18 @@ type mockTelemetryProgramClient struct {
7777
InitializeDeviceLatencySamplesFunc func(ctx context.Context, config sdktelemetry.InitializeDeviceLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
7878
WriteDeviceLatencySamplesFunc func(ctx context.Context, config sdktelemetry.WriteDeviceLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error)
7979
GetDeviceLatencySamplesFunc func(ctx context.Context, originDevicePK solana.PublicKey, targetDevicePK solana.PublicKey, linkPK solana.PublicKey, epoch uint64) (*sdktelemetry.DeviceLatencySamples, error)
80+
InitializeTimestampIndexFunc func(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error)
81+
}
82+
83+
func (c *mockTelemetryProgramClient) ProgramID() solana.PublicKey {
84+
return solana.MustPublicKeyFromBase58("11111111111111111111111111111111")
85+
}
86+
87+
func (c *mockTelemetryProgramClient) InitializeTimestampIndex(ctx context.Context, samplesAccountPK solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error) {
88+
if c.InitializeTimestampIndexFunc != nil {
89+
return c.InitializeTimestampIndexFunc(ctx, samplesAccountPK)
90+
}
91+
return solana.Signature{}, nil, nil
8092
}
8193

8294
func (c *mockTelemetryProgramClient) InitializeDeviceLatencySamples(ctx context.Context, config sdktelemetry.InitializeDeviceLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {
@@ -103,6 +115,14 @@ func newMemoryTelemetryProgramClient() *memoryTelemetryProgramClient {
103115
}
104116
}
105117

118+
func (c *memoryTelemetryProgramClient) ProgramID() solana.PublicKey {
119+
return solana.MustPublicKeyFromBase58("11111111111111111111111111111111")
120+
}
121+
122+
func (c *memoryTelemetryProgramClient) InitializeTimestampIndex(_ context.Context, _ solana.PublicKey) (solana.Signature, *solanarpc.GetTransactionResult, error) {
123+
return solana.Signature{}, nil, nil
124+
}
125+
106126
func (c *memoryTelemetryProgramClient) InitializeDeviceLatencySamples(ctx context.Context, config sdktelemetry.InitializeDeviceLatencySamplesInstructionConfig) (solana.Signature, *solanarpc.GetTransactionResult, error) {
107127
c.mu.Lock()
108128
defer c.mu.Unlock()

0 commit comments

Comments
 (0)