package storeeventrunner import ( "context" "log/slog" "strings" documentstore "queryorchestration/internal/document/store" ) const Name = "storeEventRunner" type Services struct { Store *documentstore.Service } type Runner struct { svc *Services } func New(svc *Services) Runner { return Runner{ svc: svc, } } type S3Bucket struct { Name string `json:"name"` Arn string `json:"arn"` } type S3Object struct { Key string `json:"key"` Size int64 `json:"size"` ETag string `json:"eTag"` VersionId string `json:"versionId"` } type S3EventRecordDetails struct { Bucket S3Bucket `json:"bucket"` Object S3Object `json:"object"` } type S3EventRecord struct { EventVersion string `json:"eventVersion"` EventSource string `json:"eventSource"` AwsRegion Region `json:"awsRegion"` EventTime string `json:"eventTime"` EventName EventS3 `json:"eventName"` S3 S3EventRecordDetails `json:"s3"` } type S3EventNotification struct { Records []S3EventRecord `json:"Records"` } type Region string type EventS3 string const ( EventS3ObjectCreatedPut = "ObjectCreated:Put" EventS3ObjectCreatedCompleteMultipartUpload = "ObjectCreated:CompleteMultipartUpload" ) func (s Runner) Process(ctx context.Context, body S3EventNotification) bool { for _, record := range body.Records { // Skip batch processing keys - these are handled separately if strings.HasPrefix(record.S3.Object.Key, "batches/") { slog.Debug("discarding batch message", "key", record.S3.Object.Key, "bucket", record.S3.Bucket.Name, "event", record.EventName) continue } err := s.svc.Store.Process(ctx, documentstore.Params{ Bucket: record.S3.Bucket.Name, Key: record.S3.Object.Key, Hash: record.S3.Object.ETag, Event: documentstore.EventS3(record.EventName), }) if err != nil { slog.Error("unable to process", "error", err) return false } } return true }