Skip to content

Commit fd0ae2a

Browse files
elffjsclaude
andauthored
Presign externalized payloads via data_index_key (#79)
* feat: presign externalized payloads via data_index_key Adopt cloudevent v0.2.8: events whose payload was split into the blob bucket carry a data_index_key column in ClickHouse. Switch the GraphQL presign trigger from the cloudevent/blobs/ key prefix to a non-empty data_index_key, and presign that key against a new BLOB_BUCKET setting instead of the parquet bucket. CloudEvent responses for externalized events are now header-only at the eventrepo layer (the GraphQL resolver fills in dataUrl); other paths remain unchanged. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * Proper blob bucket --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
1 parent 69a9f6c commit fd0ae2a

12 files changed

Lines changed: 65 additions & 44 deletions

File tree

charts/fetch-api/values-prod.yaml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ env:
1212
CLOUDEVENT_BUCKET: dimo-ingest-cloudevent-prod
1313
EPHEMERAL_BUCKET: dimo-ingest-ephemeral-prod
1414
PARQUET_BUCKET: dimo-storage-prod
15+
BLOB_BUCKET: dimo-blob-storage-prod
1516
VEHICLE_NFT_ADDRESS: '0xbA5738a18d83D41847dfFbDC6101d37C69c9B0cF'
1617
IDENTITY_API_URL: https://identity-api.dimo.zone/query
1718
ingress:

charts/fetch-api/values.yaml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ env:
3434
CLOUDEVENT_BUCKET: dimo-ingest-cloudevent-dev
3535
EPHEMERAL_BUCKET: dimo-ingest-ephemeral-dev
3636
PARQUET_BUCKET: dimo-storage-dev
37+
BLOB_BUCKET: dimo-blob-storage-dev
3738
S3_AWS_REGION: us-east-2
3839
IDENTITY_API_URL: https://identity-api.dev.dimo.zone/query
3940
service:

go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ require (
66
github.com/99designs/gqlgen v0.17.89
77
github.com/ClickHouse/clickhouse-go/v2 v2.43.0
88
github.com/DIMO-Network/clickhouse-infra v0.0.7
9-
github.com/DIMO-Network/cloudevent v0.2.7
9+
github.com/DIMO-Network/cloudevent v0.2.8
1010
github.com/DIMO-Network/server-garage v0.1.1
1111
github.com/DIMO-Network/shared v1.1.7
1212
github.com/DIMO-Network/token-exchange-api v0.4.0

go.sum

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,8 @@ github.com/DATA-DOG/go-sqlmock v1.5.2 h1:OcvFkGmslmlZibjAjaHm3L//6LiuBgolP7Oputl
1616
github.com/DATA-DOG/go-sqlmock v1.5.2/go.mod h1:88MAG/4G7SMwSE3CeA0ZKzrT5CiOU3OJ+JlNzwDqpNU=
1717
github.com/DIMO-Network/clickhouse-infra v0.0.7 h1:TAsjkFFKu3D5Xg6dwBcRBryjCVSlXsNjVbTwJ4UDlTg=
1818
github.com/DIMO-Network/clickhouse-infra v0.0.7/go.mod h1:XS80lhSJNWBWGgZ+m4j7++zFj1wAXfmtV2gJfhGlabQ=
19-
github.com/DIMO-Network/cloudevent v0.2.7 h1:/cgFhUcWcliZYrmITkB8oIZb+zDhZvYNxWVGS2D3894=
20-
github.com/DIMO-Network/cloudevent v0.2.7/go.mod h1:I/9NcpMozV5Fw194WimhbkAsJtKVZf5UKYJ9hgc8Cdg=
19+
github.com/DIMO-Network/cloudevent v0.2.8 h1:Q0xGQVPlOshF2LSX/m15Qzi2n4BI0EQDgOM71gEbNsY=
20+
github.com/DIMO-Network/cloudevent v0.2.8/go.mod h1:I/9NcpMozV5Fw194WimhbkAsJtKVZf5UKYJ9hgc8Cdg=
2121
github.com/DIMO-Network/server-garage v0.1.1 h1:EYmyy+Fgi2BNW0Bufn04BViDtb8BCWaN7C7BbEuoI5s=
2222
github.com/DIMO-Network/server-garage v0.1.1/go.mod h1:Z3A1KDUsXey+XhrPhmw/wyCidfrQvmEdWp7nShno7ZM=
2323
github.com/DIMO-Network/shared v1.1.7 h1:5Ex8bZ6BpOjcLj4u7n5Kih1Ho6b9BVJsKpKn4iU2EaM=

internal/app/app.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ func New(settings config.Settings) (*App, error) {
4646
}
4747
s3Client := s3ClientFromSettings(&settings)
4848
buckets := []string{settings.CloudEventBucket, settings.EphemeralBucket, settings.ParquetBucket}
49-
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket)
49+
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket, settings.BlobBucket)
5050

5151
var identityClient identity.Client
5252
if settings.IdentityAPIURL != "" {
@@ -139,7 +139,7 @@ func CreateGRPCServer(logger *zerolog.Logger, settings *config.Settings) (*grpc.
139139
}
140140

141141
s3Client := s3ClientFromSettings(settings)
142-
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket)
142+
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket, settings.BlobBucket)
143143

144144
rpcServer := rpc.NewServer([]string{settings.CloudEventBucket, settings.EphemeralBucket, settings.ParquetBucket}, eventService)
145145

internal/config/settings.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ type Settings struct {
1717
CloudEventBucket string `yaml:"CLOUDEVENT_BUCKET"`
1818
EphemeralBucket string `yaml:"EPHEMERAL_BUCKET"`
1919
ParquetBucket string `yaml:"PARQUET_BUCKET"`
20+
BlobBucket string `yaml:"BLOB_BUCKET"`
2021
S3AWSRegion string `yaml:"S3_AWS_REGION"`
2122
S3AWSAccessKeyID string `yaml:"S3_AWS_ACCESS_KEY_ID"`
2223
S3AWSSecretAccessKey string `yaml:"S3_AWS_SECRET_ACCESS_KEY"`

internal/graph/base.resolvers.go

Lines changed: 4 additions & 5 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

pkg/eventrepo/db_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -60,11 +60,11 @@ func setupClickHouseContainer(t *testing.T) *container.Container {
6060
return globalTestContainer.container
6161
}
6262

63-
// insertTestData inserts test data into ClickHouse.
63+
// insertTestData inserts test data into ClickHouse and returns the index_key.
6464
func insertTestData(t *testing.T, ctx context.Context, conn clickhouse.Conn, index *cloudevent.CloudEventHeader) string {
6565
values := chindexer.CloudEventToSlice(index)
6666

6767
err := conn.Exec(ctx, chindexer.InsertStmt, values...)
6868
require.NoError(t, err)
69-
return values[len(values)-1].(string)
69+
return chindexer.CloudEventToObjectKey(index)
7070
}

pkg/eventrepo/event_repo_test.go

Lines changed: 14 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -93,7 +93,7 @@ func TestGetLatestIndexKey(t *testing.T) {
9393
},
9494
}
9595

96-
indexService := eventrepo.New(conn, nil, nil, "")
96+
indexService := eventrepo.New(conn, nil, nil, "", "")
9797

9898
for _, tt := range tests {
9999
t.Run(tt.name, func(t *testing.T) {
@@ -172,7 +172,7 @@ func TestGetDataFromIndex(t *testing.T) {
172172
ContentLength: ref(int64(len(content))),
173173
}, nil).AnyTimes()
174174

175-
indexService := eventrepo.New(conn, mockS3Client, nil, "")
175+
indexService := eventrepo.New(conn, mockS3Client, nil, "", "")
176176

177177
for _, tt := range tests {
178178
t.Run(tt.name, func(t *testing.T) {
@@ -204,7 +204,7 @@ func TestStoreObject(t *testing.T) {
204204
mockS3Client := NewMockObjectGetter(ctrl)
205205
mockS3Client.EXPECT().PutObject(gomock.Any(), gomock.Any(), gomock.Any()).Return(&s3.PutObjectOutput{}, nil).AnyTimes()
206206

207-
indexService := eventrepo.New(conn, mockS3Client, nil, "")
207+
indexService := eventrepo.New(conn, mockS3Client, nil, "", "")
208208

209209
content := []byte(`{"vin": "1HGCM82633A123456"}`)
210210
did := cloudevent.ERC721DID{
@@ -333,7 +333,7 @@ func TestGetData(t *testing.T) {
333333
ctrl := gomock.NewController(t)
334334
mockS3Client := NewMockObjectGetter(ctrl)
335335

336-
indexService := eventrepo.New(conn, mockS3Client, nil, "")
336+
indexService := eventrepo.New(conn, mockS3Client, nil, "", "")
337337
// Allow GetObject calls in any order since fetches are concurrent.
338338
if len(tt.expectedIndexKeys) > 0 {
339339
mockS3Client.EXPECT().GetObject(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn(func(ctx context.Context, params *s3.GetObjectInput, optFns ...func(*s3.Options)) (*s3.GetObjectOutput, error) {
@@ -425,7 +425,7 @@ func TestGetEventWithAllHeaderFields(t *testing.T) {
425425
eventDataEnvelope := []byte(`{"data":` + string(eventData) + `}`)
426426

427427
// Create service
428-
indexService := eventrepo.New(conn, mockS3Client, nil, "")
428+
indexService := eventrepo.New(conn, mockS3Client, nil, "", "")
429429

430430
// Test retrieving the event
431431
t.Run("retrieve event with full headers", func(t *testing.T) {
@@ -535,8 +535,12 @@ func ref[T any](x T) *T {
535535
// returns the raw bytes and the map of event index to index key.
536536
func encodeTestParquet(t *testing.T, events []cloudevent.RawEvent, objectKey string) ([]byte, map[int]string) {
537537
t.Helper()
538+
stored := make([]cloudevent.StoredEvent, len(events))
539+
for i, ev := range events {
540+
stored[i] = cloudevent.StoredEvent{RawEvent: ev}
541+
}
538542
var buf bytes.Buffer
539-
indexKeys, err := parquet.Encode(&buf, events, objectKey)
543+
indexKeys, err := parquet.Encode(&buf, stored, objectKey)
540544
require.NoError(t, err)
541545
return buf.Bytes(), indexKeys
542546
}
@@ -609,7 +613,7 @@ func TestGetCloudEventFromIndex_ParquetRef(t *testing.T) {
609613
ctrl := gomock.NewController(t)
610614
mockS3 := mockS3ParquetReader(t, ctrl, parquetBytes)
611615

612-
indexService := eventrepo.New(conn, mockS3, nil, "test-parquet-bucket")
616+
indexService := eventrepo.New(conn, mockS3, nil, "test-parquet-bucket", "")
613617

614618
// Build the index object as GetCloudEventFromIndex expects
615619
index := cloudevent.CloudEvent[eventrepo.ObjectInfo]{
@@ -688,7 +692,7 @@ func TestListCloudEventsFromIndexes_ParquetCaching(t *testing.T) {
688692
ctrl := gomock.NewController(t)
689693
mockS3 := mockS3ParquetReader(t, ctrl, parquetBytes)
690694

691-
indexService := eventrepo.New(conn, mockS3, nil, "test-parquet-bucket")
695+
indexService := eventrepo.New(conn, mockS3, nil, "test-parquet-bucket", "")
692696

693697
indexes := []cloudevent.CloudEvent[eventrepo.ObjectInfo]{
694698
{CloudEventHeader: hdr0, Data: eventrepo.ObjectInfo{Key: indexKeys[0]}},
@@ -780,7 +784,7 @@ func TestListIndexesAdvanced(t *testing.T) {
780784
keyTypeStatusSource1Producer3 := insertTestData(t, ctx, conn, eventIdx3)
781785
keyTypeStatusSource3Producer4 := insertTestData(t, ctx, conn, eventIdx4)
782786

783-
indexService := eventrepo.New(conn, nil, nil, "")
787+
indexService := eventrepo.New(conn, nil, nil, "", "")
784788

785789
tests := []struct {
786790
name string
@@ -1065,7 +1069,7 @@ func TestGetCloudEventTypeSummaries(t *testing.T) {
10651069
insertTestData(t, ctx, conn, status3)
10661070
insertTestData(t, ctx, conn, fp1)
10671071

1068-
indexService := eventrepo.New(conn, nil, nil, "")
1072+
indexService := eventrepo.New(conn, nil, nil, "", "")
10691073

10701074
t.Run("no filter returns all types", func(t *testing.T) {
10711075
opts := &grpc.SearchOptions{

pkg/eventrepo/eventrepo.go

Lines changed: 31 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -36,11 +36,19 @@ type Service struct {
3636
chConn clickhouse.Conn
3737
// parquetBucket is the object storage bucket for Iceberg Parquet files.
3838
parquetBucket string
39+
// blobBucket is the object storage bucket for externalized event payloads
40+
// (referenced by data_index_key).
41+
blobBucket string
3942
}
4043

4144
// ObjectInfo is the information about the object in S3.
4245
type ObjectInfo struct {
46+
// Key is the index_key — a pointer into parquet (key#row) or a legacy JSON object.
4347
Key string
48+
// DataIndexKey, when non-empty, is the key in the blob bucket holding the
49+
// event's externalized payload. Set by the producer when the payload was
50+
// split out of the inline event.
51+
DataIndexKey string
4452
}
4553

4654
// ObjectGetter is an interface for getting an object from S3.
@@ -54,37 +62,34 @@ type Presigner interface {
5462
PresignGetObject(ctx context.Context, params *s3.GetObjectInput, optFns ...func(*s3.PresignOptions)) (*v4.PresignedHTTPRequest, error)
5563
}
5664

57-
// BlobKeyPrefix is the S3 key prefix used for large binary blob objects.
58-
// Keys with this prefix are served via presigned URL instead of inline in the response.
59-
const BlobKeyPrefix = "cloudevent/blobs/"
60-
6165
// presignTTL is the lifetime of generated presigned S3 URLs.
6266
const presignTTL = 15 * time.Minute
6367

6468
// New creates a new instance of Service.
65-
func New(chConn clickhouse.Conn, objGetter ObjectGetter, presigner Presigner, parquetBucket string) *Service {
69+
func New(chConn clickhouse.Conn, objGetter ObjectGetter, presigner Presigner, parquetBucket, blobBucket string) *Service {
6670
return &Service{
6771
objGetter: objGetter,
6872
presigner: presigner,
6973
chConn: chConn,
7074
parquetBucket: parquetBucket,
75+
blobBucket: blobBucket,
7176
}
7277
}
7378

74-
// PresignBlobURL returns a short-lived presigned GET URL for the given S3 key in the parquet bucket.
79+
// PresignBlobURL returns a short-lived presigned GET URL for the given S3 key in the blob bucket.
7580
func (s *Service) PresignBlobURL(ctx context.Context, key string) (string, error) {
7681
if s.presigner == nil {
7782
return "", fmt.Errorf("presigner not configured")
7883
}
79-
if s.parquetBucket == "" {
80-
return "", fmt.Errorf("parquet bucket not configured")
84+
if s.blobBucket == "" {
85+
return "", fmt.Errorf("blob bucket not configured")
8186
}
8287
req, err := s.presigner.PresignGetObject(ctx, &s3.GetObjectInput{
83-
Bucket: aws.String(s.parquetBucket),
88+
Bucket: aws.String(s.blobBucket),
8489
Key: aws.String(key),
8590
}, s3.WithPresignExpires(presignTTL))
8691
if err != nil {
87-
return "", fmt.Errorf("presign %s/%s: %w", s.parquetBucket, key, err)
92+
return "", fmt.Errorf("presign %s/%s: %w", s.blobBucket, key, err)
8893
}
8994
return req.URL, nil
9095
}
@@ -143,6 +148,7 @@ func (s *Service) ListIndexesAdvanced(ctx context.Context, limit int, advancedOp
143148
chindexer.DataVersionColumn,
144149
chindexer.ExtrasColumn,
145150
chindexer.IndexKeyColumn,
151+
chindexer.DataIndexKeyColumn,
146152
),
147153
qm.From(chindexer.TableName),
148154
qm.OrderBy(chindexer.TimestampColumn + order),
@@ -164,7 +170,7 @@ func (s *Service) ListIndexesAdvanced(ctx context.Context, limit int, advancedOp
164170
var extras string
165171
for rows.Next() {
166172
var event cloudevent.CloudEvent[ObjectInfo]
167-
err = rows.Scan(&event.Subject, &event.Time, &event.Type, &event.ID, &event.Source, &event.Producer, &event.DataContentType, &event.DataVersion, &extras, &event.Data.Key)
173+
err = rows.Scan(&event.Subject, &event.Time, &event.Type, &event.ID, &event.Source, &event.Producer, &event.DataContentType, &event.DataVersion, &extras, &event.Data.Key, &event.Data.DataIndexKey)
168174
if err != nil {
169175
_ = rows.Close()
170176
return nil, fmt.Errorf("failed to scan cloud event: %w", err)
@@ -349,12 +355,12 @@ func (s *Service) ListCloudEventsFromIndexes(ctx context.Context, indexes []clou
349355
}
350356
defer func() { _ = pr.Close() }()
351357
for _, item := range items {
352-
event, err := pr.SeekToRow(item.rowOffset)
358+
stored, err := pr.SeekToRow(item.rowOffset)
353359
if err != nil {
354360
return fmt.Errorf("seek to row %d in %s: %w", item.rowOffset, item.objectKey, err)
355361
}
356-
event.Tags = grpc.TagsOrEmpty(event.Tags)
357-
events[item.idx] = event
362+
stored.Tags = grpc.TagsOrEmpty(stored.Tags)
363+
events[item.idx] = stored.RawEvent
358364
}
359365
return nil
360366
})
@@ -387,7 +393,15 @@ func (s *Service) ListCloudEventsFromIndexes(ctx context.Context, indexes []clou
387393
}
388394

389395
// GetCloudEventFromIndex fetches and returns the cloud event for the given index.
396+
// Events with a non-empty DataIndexKey have their payload externalized to the
397+
// blob bucket; this method returns header-only since the data is meant to be
398+
// served via presigned URL by the GraphQL layer (see PresignBlobURL).
390399
func (s *Service) GetCloudEventFromIndex(ctx context.Context, index *cloudevent.CloudEvent[ObjectInfo], bucketName string) (cloudevent.RawEvent, error) {
400+
if index.Data.DataIndexKey != "" {
401+
hdr := index.CloudEventHeader
402+
hdr.Tags = grpc.TagsOrEmpty(hdr.Tags)
403+
return cloudevent.RawEvent{CloudEventHeader: hdr}, nil
404+
}
391405
if parquet.IsParquetRef(index.Data.Key) {
392406
return s.getCloudEventFromParquet(ctx, index.Data.Key)
393407
}
@@ -454,12 +468,12 @@ func (s *Service) getCloudEventFromParquet(ctx context.Context, key string) (clo
454468
return cloudevent.RawEvent{}, fmt.Errorf("create s3 reader for %s: %w", objectKey, err)
455469
}
456470

457-
event, err := parquet.SeekToRow(reader, reader.Size(), rowOffset)
471+
stored, err := parquet.SeekToRow(reader, reader.Size(), rowOffset)
458472
if err != nil {
459473
return cloudevent.RawEvent{}, fmt.Errorf("seek to row %d in %s: %w", rowOffset, objectKey, err)
460474
}
461-
event.Tags = grpc.TagsOrEmpty(event.Tags)
462-
return event, nil
475+
stored.Tags = grpc.TagsOrEmpty(stored.Tags)
476+
return stored.RawEvent, nil
463477
}
464478

465479
// parseParquetRef parses a parquet index key into bucket, object key, and row offset.

0 commit comments

Comments
 (0)