Skip to content
Merged
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 charts/fetch-api/values-prod.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ env:
CLOUDEVENT_BUCKET: dimo-ingest-cloudevent-prod
EPHEMERAL_BUCKET: dimo-ingest-ephemeral-prod
PARQUET_BUCKET: dimo-storage-prod
BLOB_BUCKET: dimo-blob-storage-prod
VEHICLE_NFT_ADDRESS: '0xbA5738a18d83D41847dfFbDC6101d37C69c9B0cF'
IDENTITY_API_URL: https://identity-api.dimo.zone/query
ingress:
Expand Down
1 change: 1 addition & 0 deletions charts/fetch-api/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ env:
CLOUDEVENT_BUCKET: dimo-ingest-cloudevent-dev
EPHEMERAL_BUCKET: dimo-ingest-ephemeral-dev
PARQUET_BUCKET: dimo-storage-dev
BLOB_BUCKET: dimo-blob-storage-dev
S3_AWS_REGION: us-east-2
IDENTITY_API_URL: https://identity-api.dev.dimo.zone/query
service:
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ require (
github.com/99designs/gqlgen v0.17.89
github.com/ClickHouse/clickhouse-go/v2 v2.43.0
github.com/DIMO-Network/clickhouse-infra v0.0.7
github.com/DIMO-Network/cloudevent v0.2.7
github.com/DIMO-Network/cloudevent v0.2.8
github.com/DIMO-Network/server-garage v0.1.1
github.com/DIMO-Network/shared v1.1.7
github.com/DIMO-Network/token-exchange-api v0.4.0
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@ github.com/DATA-DOG/go-sqlmock v1.5.2 h1:OcvFkGmslmlZibjAjaHm3L//6LiuBgolP7Oputl
github.com/DATA-DOG/go-sqlmock v1.5.2/go.mod h1:88MAG/4G7SMwSE3CeA0ZKzrT5CiOU3OJ+JlNzwDqpNU=
github.com/DIMO-Network/clickhouse-infra v0.0.7 h1:TAsjkFFKu3D5Xg6dwBcRBryjCVSlXsNjVbTwJ4UDlTg=
github.com/DIMO-Network/clickhouse-infra v0.0.7/go.mod h1:XS80lhSJNWBWGgZ+m4j7++zFj1wAXfmtV2gJfhGlabQ=
github.com/DIMO-Network/cloudevent v0.2.7 h1:/cgFhUcWcliZYrmITkB8oIZb+zDhZvYNxWVGS2D3894=
github.com/DIMO-Network/cloudevent v0.2.7/go.mod h1:I/9NcpMozV5Fw194WimhbkAsJtKVZf5UKYJ9hgc8Cdg=
github.com/DIMO-Network/cloudevent v0.2.8 h1:Q0xGQVPlOshF2LSX/m15Qzi2n4BI0EQDgOM71gEbNsY=
github.com/DIMO-Network/cloudevent v0.2.8/go.mod h1:I/9NcpMozV5Fw194WimhbkAsJtKVZf5UKYJ9hgc8Cdg=
github.com/DIMO-Network/server-garage v0.1.1 h1:EYmyy+Fgi2BNW0Bufn04BViDtb8BCWaN7C7BbEuoI5s=
github.com/DIMO-Network/server-garage v0.1.1/go.mod h1:Z3A1KDUsXey+XhrPhmw/wyCidfrQvmEdWp7nShno7ZM=
github.com/DIMO-Network/shared v1.1.7 h1:5Ex8bZ6BpOjcLj4u7n5Kih1Ho6b9BVJsKpKn4iU2EaM=
Expand Down
4 changes: 2 additions & 2 deletions internal/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ func New(settings config.Settings) (*App, error) {
}
s3Client := s3ClientFromSettings(&settings)
buckets := []string{settings.CloudEventBucket, settings.EphemeralBucket, settings.ParquetBucket}
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket)
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket, settings.BlobBucket)

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

s3Client := s3ClientFromSettings(settings)
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket)
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket, settings.BlobBucket)

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

Expand Down
1 change: 1 addition & 0 deletions internal/config/settings.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ type Settings struct {
CloudEventBucket string `yaml:"CLOUDEVENT_BUCKET"`
EphemeralBucket string `yaml:"EPHEMERAL_BUCKET"`
ParquetBucket string `yaml:"PARQUET_BUCKET"`
BlobBucket string `yaml:"BLOB_BUCKET"`
S3AWSRegion string `yaml:"S3_AWS_REGION"`
S3AWSAccessKeyID string `yaml:"S3_AWS_ACCESS_KEY_ID"`
S3AWSSecretAccessKey string `yaml:"S3_AWS_SECRET_ACCESS_KEY"`
Expand Down
9 changes: 4 additions & 5 deletions internal/graph/base.resolvers.go

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

4 changes: 2 additions & 2 deletions pkg/eventrepo/db_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,11 +60,11 @@ func setupClickHouseContainer(t *testing.T) *container.Container {
return globalTestContainer.container
}

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

err := conn.Exec(ctx, chindexer.InsertStmt, values...)
require.NoError(t, err)
return values[len(values)-1].(string)
return chindexer.CloudEventToObjectKey(index)
}
24 changes: 14 additions & 10 deletions pkg/eventrepo/event_repo_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ func TestGetLatestIndexKey(t *testing.T) {
},
}

indexService := eventrepo.New(conn, nil, nil, "")
indexService := eventrepo.New(conn, nil, nil, "", "")

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

indexService := eventrepo.New(conn, mockS3Client, nil, "")
indexService := eventrepo.New(conn, mockS3Client, nil, "", "")

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

indexService := eventrepo.New(conn, mockS3Client, nil, "")
indexService := eventrepo.New(conn, mockS3Client, nil, "", "")

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

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

// Create service
indexService := eventrepo.New(conn, mockS3Client, nil, "")
indexService := eventrepo.New(conn, mockS3Client, nil, "", "")

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

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

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

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

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

indexService := eventrepo.New(conn, nil, nil, "")
indexService := eventrepo.New(conn, nil, nil, "", "")

tests := []struct {
name string
Expand Down Expand Up @@ -1065,7 +1069,7 @@ func TestGetCloudEventTypeSummaries(t *testing.T) {
insertTestData(t, ctx, conn, status3)
insertTestData(t, ctx, conn, fp1)

indexService := eventrepo.New(conn, nil, nil, "")
indexService := eventrepo.New(conn, nil, nil, "", "")

t.Run("no filter returns all types", func(t *testing.T) {
opts := &grpc.SearchOptions{
Expand Down
48 changes: 31 additions & 17 deletions pkg/eventrepo/eventrepo.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,11 +36,19 @@ type Service struct {
chConn clickhouse.Conn
// parquetBucket is the object storage bucket for Iceberg Parquet files.
parquetBucket string
// blobBucket is the object storage bucket for externalized event payloads
// (referenced by data_index_key).
blobBucket string
}

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

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

// BlobKeyPrefix is the S3 key prefix used for large binary blob objects.
// Keys with this prefix are served via presigned URL instead of inline in the response.
const BlobKeyPrefix = "cloudevent/blobs/"

// presignTTL is the lifetime of generated presigned S3 URLs.
const presignTTL = 15 * time.Minute

// New creates a new instance of Service.
func New(chConn clickhouse.Conn, objGetter ObjectGetter, presigner Presigner, parquetBucket string) *Service {
func New(chConn clickhouse.Conn, objGetter ObjectGetter, presigner Presigner, parquetBucket, blobBucket string) *Service {
return &Service{
objGetter: objGetter,
presigner: presigner,
chConn: chConn,
parquetBucket: parquetBucket,
blobBucket: blobBucket,
}
}

// PresignBlobURL returns a short-lived presigned GET URL for the given S3 key in the parquet bucket.
// PresignBlobURL returns a short-lived presigned GET URL for the given S3 key in the blob bucket.
func (s *Service) PresignBlobURL(ctx context.Context, key string) (string, error) {
if s.presigner == nil {
return "", fmt.Errorf("presigner not configured")
}
if s.parquetBucket == "" {
return "", fmt.Errorf("parquet bucket not configured")
if s.blobBucket == "" {
return "", fmt.Errorf("blob bucket not configured")
}
req, err := s.presigner.PresignGetObject(ctx, &s3.GetObjectInput{
Bucket: aws.String(s.parquetBucket),
Bucket: aws.String(s.blobBucket),
Key: aws.String(key),
}, s3.WithPresignExpires(presignTTL))
if err != nil {
return "", fmt.Errorf("presign %s/%s: %w", s.parquetBucket, key, err)
return "", fmt.Errorf("presign %s/%s: %w", s.blobBucket, key, err)
}
return req.URL, nil
}
Expand Down Expand Up @@ -143,6 +148,7 @@ func (s *Service) ListIndexesAdvanced(ctx context.Context, limit int, advancedOp
chindexer.DataVersionColumn,
chindexer.ExtrasColumn,
chindexer.IndexKeyColumn,
chindexer.DataIndexKeyColumn,
),
qm.From(chindexer.TableName),
qm.OrderBy(chindexer.TimestampColumn + order),
Expand All @@ -164,7 +170,7 @@ func (s *Service) ListIndexesAdvanced(ctx context.Context, limit int, advancedOp
var extras string
for rows.Next() {
var event cloudevent.CloudEvent[ObjectInfo]
err = rows.Scan(&event.Subject, &event.Time, &event.Type, &event.ID, &event.Source, &event.Producer, &event.DataContentType, &event.DataVersion, &extras, &event.Data.Key)
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)
if err != nil {
_ = rows.Close()
return nil, fmt.Errorf("failed to scan cloud event: %w", err)
Expand Down Expand Up @@ -349,12 +355,12 @@ func (s *Service) ListCloudEventsFromIndexes(ctx context.Context, indexes []clou
}
defer func() { _ = pr.Close() }()
for _, item := range items {
event, err := pr.SeekToRow(item.rowOffset)
stored, err := pr.SeekToRow(item.rowOffset)
if err != nil {
return fmt.Errorf("seek to row %d in %s: %w", item.rowOffset, item.objectKey, err)
}
event.Tags = grpc.TagsOrEmpty(event.Tags)
events[item.idx] = event
stored.Tags = grpc.TagsOrEmpty(stored.Tags)
events[item.idx] = stored.RawEvent
}
return nil
})
Expand Down Expand Up @@ -387,7 +393,15 @@ func (s *Service) ListCloudEventsFromIndexes(ctx context.Context, indexes []clou
}

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

event, err := parquet.SeekToRow(reader, reader.Size(), rowOffset)
stored, err := parquet.SeekToRow(reader, reader.Size(), rowOffset)
if err != nil {
return cloudevent.RawEvent{}, fmt.Errorf("seek to row %d in %s: %w", rowOffset, objectKey, err)
}
event.Tags = grpc.TagsOrEmpty(event.Tags)
return event, nil
stored.Tags = grpc.TagsOrEmpty(stored.Tags)
return stored.RawEvent, nil
}

// parseParquetRef parses a parquet index key into bucket, object key, and row offset.
Expand Down
10 changes: 5 additions & 5 deletions pkg/eventrepo/presign_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ func TestPresignBlobURL(t *testing.T) {
expectedURL = "https://s3.amazonaws.com/test-bucket/cloudevent/blobs/some-scan.bin?X-Amz-Signature=abc123"
)

svc := eventrepo.New(nil, nil, mockPresigner, bucket)
svc := eventrepo.New(nil, nil, mockPresigner, "", bucket)

mockPresigner.EXPECT().
PresignGetObject(gomock.Any(), gomock.Any(), gomock.Any()).
Expand Down Expand Up @@ -54,7 +54,7 @@ func TestPresignBlobURL_PresignerError(t *testing.T) {
ctrl := gomock.NewController(t)
mockPresigner := NewMockPresigner(ctrl)

svc := eventrepo.New(nil, nil, mockPresigner, "test-bucket")
svc := eventrepo.New(nil, nil, mockPresigner, "", "test-bucket")

mockPresigner.EXPECT().
PresignGetObject(gomock.Any(), gomock.Any(), gomock.Any()).
Expand All @@ -67,7 +67,7 @@ func TestPresignBlobURL_PresignerError(t *testing.T) {

func TestPresignBlobURL_NilPresigner(t *testing.T) {
t.Parallel()
svc := eventrepo.New(nil, nil, nil, "test-bucket")
svc := eventrepo.New(nil, nil, nil, "", "test-bucket")

_, err := svc.PresignBlobURL(context.Background(), "cloudevent/blobs/test.bin")
require.Error(t, err)
Expand All @@ -78,9 +78,9 @@ func TestPresignBlobURL_NoBucket(t *testing.T) {
ctrl := gomock.NewController(t)
mockPresigner := NewMockPresigner(ctrl)

svc := eventrepo.New(nil, nil, mockPresigner, "")
svc := eventrepo.New(nil, nil, mockPresigner, "", "")

_, err := svc.PresignBlobURL(context.Background(), "cloudevent/blobs/test.bin")
require.Error(t, err)
assert.Contains(t, err.Error(), "parquet bucket not configured")
assert.Contains(t, err.Error(), "blob bucket not configured")
}
1 change: 1 addition & 0 deletions settings.sample.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ TOKEN_EXCHANGE_ISSUER_URL: http://127.0.0.1:5556/dex
CLOUDEVENT_BUCKET: ""
EPHEMERAL_BUCKET: ""
PARQUET_BUCKET: ""
BLOB_BUCKET: ""
S3_AWS_REGION: ""
S3_AWS_ACCESS_KEY_ID: ""
S3_AWS_SECRET_ACCESS_KEY: ""
Expand Down
Loading