Skip to content

Commit 0779bd1

Browse files
elffjsclaude
andauthored
Add presigned S3 URL support for large blob payloads (#66)
* Add S3 presigned URL support for blob CloudEvents Objects stored under the cloudevent/blobs/ key prefix are now served via a short-lived presigned S3 GET URL (dataUrl field) instead of being downloaded and embedded inline in the GraphQL response. This avoids ballooning the JSON payload for large binary objects (e.g. scans). - Add Presigner interface + PresignBlobURL method to eventrepo.Service - Export BlobKeyPrefix constant ("cloudevent/blobs/") - Add dataUrl: String field to GraphQL CloudEvent type - Detect blob prefix in LatestCloudEvent / CloudEvents resolvers; skip GetObject and presign against the primary bucket instead - Wire s3.NewPresignClient into both GraphQL and gRPC app paths - Update eventrepo.New signature; pass nil in tests that don't need presigning - Add MockPresigner to generated mock file https://claude.ai/code/session_015ReeLGeCywJfYkkrrng5wU * Improve dataUrl field description in GraphQL schema Describe the field in terms of behavior (large files) rather than the internal key prefix convention, which is an implementation detail that may change. https://claude.ai/code/session_015ReeLGeCywJfYkkrrng5wU * Add tests for presigned URL generation and blob resolver fields - pkg/eventrepo/presign_test.go: unit tests for PresignBlobURL covering correct bucket/key routing, 15-minute TTL, presigner error propagation, and nil-presigner guard - internal/graph/blob_resolver_test.go: unit tests for the DataUrl resolver (nil wrapper, empty DataURL, populated DataURL) and a composite test that a blob wrapper returns nil for data/dataBase64 and a URL for dataUrl https://claude.ai/code/session_015ReeLGeCywJfYkkrrng5wU * Get rid of these imagined handlers * These tests don't do much --------- Co-authored-by: Claude <noreply@anthropic.com>
1 parent cfb4038 commit 0779bd1

10 files changed

Lines changed: 263 additions & 27 deletions

File tree

internal/app/app.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import (
1414
"github.com/DIMO-Network/fetch-api/internal/limits"
1515
"github.com/DIMO-Network/fetch-api/pkg/eventrepo"
1616
fetchgrpc "github.com/DIMO-Network/fetch-api/pkg/grpc"
17+
"github.com/aws/aws-sdk-go-v2/service/s3"
1718
"github.com/DIMO-Network/server-garage/pkg/gql/errorhandler"
1819
gqlmetrics "github.com/DIMO-Network/server-garage/pkg/gql/metrics"
1920
"github.com/DIMO-Network/shared/pkg/middleware/metrics"
@@ -42,7 +43,7 @@ func New(settings config.Settings) (*App, error) {
4243
}
4344
s3Client := s3ClientFromSettings(&settings)
4445
buckets := []string{settings.CloudEventBucket, settings.EphemeralBucket, settings.ParquetBucket}
45-
eventService := eventrepo.New(chConn, s3Client, settings.ParquetBucket)
46+
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket)
4647

4748
var identityClient identity.Client
4849
if settings.IdentityAPIURL != "" {
@@ -118,7 +119,7 @@ func CreateGRPCServer(logger *zerolog.Logger, settings *config.Settings) (*grpc.
118119
}
119120

120121
s3Client := s3ClientFromSettings(settings)
121-
eventService := eventrepo.New(chConn, s3Client, settings.ParquetBucket)
122+
eventService := eventrepo.New(chConn, s3Client, s3.NewPresignClient(s3Client), settings.ParquetBucket)
122123

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

internal/config/settings.go

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,13 +7,13 @@ import (
77

88
// Settings contains the application config.
99
type Settings struct {
10-
Port int `yaml:"PORT"`
11-
MonPort int `yaml:"MON_PORT"`
12-
GRPCPort int `yaml:"GRPC_PORT"`
13-
EnablePprof bool `yaml:"ENABLE_PPROF"`
14-
MaxRequestDuration string `yaml:"MAX_REQUEST_DURATION"`
15-
TokenExchangeJWTKeySetURL string `yaml:"TOKEN_EXCHANGE_JWK_KEY_SET_URL"`
16-
TokenExchangeIssuer string `yaml:"TOKEN_EXCHANGE_ISSUER_URL"`
10+
Port int `yaml:"PORT"`
11+
MonPort int `yaml:"MON_PORT"`
12+
GRPCPort int `yaml:"GRPC_PORT"`
13+
EnablePprof bool `yaml:"ENABLE_PPROF"`
14+
MaxRequestDuration string `yaml:"MAX_REQUEST_DURATION"`
15+
TokenExchangeJWTKeySetURL string `yaml:"TOKEN_EXCHANGE_JWK_KEY_SET_URL"`
16+
TokenExchangeIssuer string `yaml:"TOKEN_EXCHANGE_ISSUER_URL"`
1717
CloudEventBucket string `yaml:"CLOUDEVENT_BUCKET"`
1818
EphemeralBucket string `yaml:"EPHEMERAL_BUCKET"`
1919
ParquetBucket string `yaml:"PARQUET_BUCKET"`

internal/graph/base.resolvers.go

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

internal/graph/cloud_event.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,9 +7,10 @@ import (
77
)
88

99
// CloudEventWrapper holds a pointer to a RawEvent so resolvers can expose
10-
// header, data, and dataBase64 without copying the underlying event.
10+
// header, data, dataBase64, and dataUrl without copying the underlying event.
1111
type CloudEventWrapper struct {
12-
Raw *cloudevent.RawEvent
12+
Raw *cloudevent.RawEvent
13+
DataURL string // non-empty when the payload is a blob served via presigned URL
1314
}
1415

1516
// RawJSON is the raw bytes of a JSON value. It implements graphql.Marshaler by

internal/graph/generated.go

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

pkg/eventrepo/event_repo_test.go

Lines changed: 9 additions & 9 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, "")
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, "")
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, "")
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, "")
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) {
@@ -424,7 +424,7 @@ func TestGetEventWithAllHeaderFields(t *testing.T) {
424424
eventDataEnvelope := []byte(`{"data":` + string(eventData) + `}`)
425425

426426
// Create service
427-
indexService := eventrepo.New(conn, mockS3Client, "")
427+
indexService := eventrepo.New(conn, mockS3Client, nil, "")
428428

429429
// Test retrieving the event
430430
t.Run("retrieve event with full headers", func(t *testing.T) {
@@ -576,7 +576,7 @@ func TestGetCloudEventFromIndex_ParquetRef(t *testing.T) {
576576
ctrl := gomock.NewController(t)
577577
mockS3 := mockS3ParquetReader(t, ctrl, parquetBytes)
578578

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

581581
// Build the index object as GetCloudEventFromIndex expects
582582
index := cloudevent.CloudEvent[eventrepo.ObjectInfo]{
@@ -655,7 +655,7 @@ func TestListCloudEventsFromIndexes_ParquetCaching(t *testing.T) {
655655
ctrl := gomock.NewController(t)
656656
mockS3 := mockS3ParquetReader(t, ctrl, parquetBytes)
657657

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

660660
indexes := []cloudevent.CloudEvent[eventrepo.ObjectInfo]{
661661
{CloudEventHeader: hdr0, Data: eventrepo.ObjectInfo{Key: indexKeys[0]}},
@@ -747,7 +747,7 @@ func TestListIndexesAdvanced(t *testing.T) {
747747
keyTypeStatusSource1Producer3 := insertTestData(t, ctx, conn, eventIdx3)
748748
keyTypeStatusSource3Producer4 := insertTestData(t, ctx, conn, eventIdx4)
749749

750-
indexService := eventrepo.New(conn, nil, "")
750+
indexService := eventrepo.New(conn, nil, nil, "")
751751

752752
tests := []struct {
753753
name string
@@ -1032,7 +1032,7 @@ func TestGetCloudEventTypeSummaries(t *testing.T) {
10321032
insertTestData(t, ctx, conn, status3)
10331033
insertTestData(t, ctx, conn, fp1)
10341034

1035-
indexService := eventrepo.New(conn, nil, "")
1035+
indexService := eventrepo.New(conn, nil, nil, "")
10361036

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

pkg/eventrepo/eventrepo.go

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,8 @@ import (
1515
chindexer "github.com/DIMO-Network/cloudevent/clickhouse"
1616
"github.com/DIMO-Network/cloudevent/parquet"
1717
"github.com/DIMO-Network/fetch-api/pkg/grpc"
18+
"github.com/aws/aws-sdk-go-v2/aws"
19+
v4 "github.com/aws/aws-sdk-go-v2/aws/signer/v4"
1820
"github.com/aws/aws-sdk-go-v2/service/s3"
1921
"github.com/volatiletech/sqlboiler/v4/drivers"
2022
"github.com/volatiletech/sqlboiler/v4/queries"
@@ -30,6 +32,7 @@ const tagsColumn = "JSONExtract(extras, 'tags', 'Array(String)')"
3032
// Service manages and retrieves data messages from indexed objects in S3.
3133
type Service struct {
3234
objGetter ObjectGetter
35+
presigner Presigner
3336
chConn clickhouse.Conn
3437
// parquetBucket is the object storage bucket for Iceberg Parquet files.
3538
parquetBucket string
@@ -46,15 +49,43 @@ type ObjectGetter interface {
4649
PutObject(ctx context.Context, params *s3.PutObjectInput, optFns ...func(*s3.Options)) (*s3.PutObjectOutput, error)
4750
}
4851

52+
// Presigner generates presigned S3 GET URLs.
53+
type Presigner interface {
54+
PresignGetObject(ctx context.Context, params *s3.GetObjectInput, optFns ...func(*s3.PresignOptions)) (*v4.PresignedHTTPRequest, error)
55+
}
56+
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+
61+
// presignTTL is the lifetime of generated presigned S3 URLs.
62+
const presignTTL = 15 * time.Minute
63+
4964
// New creates a new instance of Service.
50-
func New(chConn clickhouse.Conn, objGetter ObjectGetter, parquetBucket string) *Service {
65+
func New(chConn clickhouse.Conn, objGetter ObjectGetter, presigner Presigner, parquetBucket string) *Service {
5166
return &Service{
5267
objGetter: objGetter,
68+
presigner: presigner,
5369
chConn: chConn,
5470
parquetBucket: parquetBucket,
5571
}
5672
}
5773

74+
// PresignBlobURL returns a short-lived presigned GET URL for the given S3 key and bucket.
75+
func (s *Service) PresignBlobURL(ctx context.Context, key, bucket string) (string, error) {
76+
if s.presigner == nil {
77+
return "", fmt.Errorf("presigner not configured")
78+
}
79+
req, err := s.presigner.PresignGetObject(ctx, &s3.GetObjectInput{
80+
Bucket: aws.String(bucket),
81+
Key: aws.String(key),
82+
}, s3.WithPresignExpires(presignTTL))
83+
if err != nil {
84+
return "", fmt.Errorf("presign %s/%s: %w", bucket, key, err)
85+
}
86+
return req.URL, nil
87+
}
88+
5889
// GetLatestIndex returns the latest cloud event index that matches the given options.
5990
func (s *Service) GetLatestIndex(ctx context.Context, opts *grpc.SearchOptions) (cloudevent.CloudEvent[ObjectInfo], error) {
6091
advancedOpts := convertSearchOptionsToAdvanced(opts)

0 commit comments

Comments
 (0)