package docinitrunner import ( "context" "database/sql" "errors" "log/slog" "queryorchestration/internal/database/repository" "queryorchestration/internal/document/batch/outcome" documentinit "queryorchestration/internal/document/init" "queryorchestration/internal/serviceconfig/objectstore" "github.com/google/uuid" ) const Name = "docInitRunner" // Services holds the dependencies for the docInitRunner. // // Fields: // - Document: service for creating/deduplicating documents // - Queries: SQLC-generated database queries // - Outcome: reporter for writing batch pipeline outcomes type Services struct { Document *documentinit.Service Queries *repository.Queries Outcome *outcome.Reporter } // Runner processes document init messages from the SQS queue. type Runner struct { svc *Services } // New creates a new docInitRunner. // // Parameters: // - svc: the runner's service dependencies // // Returns: // - Runner: the runner instance func New(svc *Services) Runner { return Runner{ svc: svc, } } // Body is the SQS message body for docInitRunner. type Body struct { Bucket string `json:"bucket" validate:"required"` Key string `json:"key" validate:"required"` Hash string `json:"hash" validate:"required"` } // Process handles a single document init message. Returns true to acknowledge // the message (success), or false to requeue (infrastructure error). func (s Runner) Process(ctx context.Context, body Body) bool { key, err := objectstore.ParseBucketKey(body.Key) if err != nil { slog.Error("unable to parse key", "key", body.Key, "error", err) return false } // Get filename and folder_id from documentUploads table var filename *string var folderID *uuid.UUID upload, err := s.svc.Queries.GetDocumentUploadByKey(ctx, &repository.GetDocumentUploadByKeyParams{ Bucket: body.Bucket, Key: body.Key, }) if err != nil && !errors.Is(err, sql.ErrNoRows) { slog.Error("unable to get document upload", "bucket", body.Bucket, "key", body.Key, "error", err) return false } if err == nil { if upload.Filename != nil { filename = upload.Filename } folderID = upload.FolderID } result, err := s.svc.Document.Create(ctx, &documentinit.Create{ Key: key, Bucket: body.Bucket, Hash: body.Hash, Filename: filename, OriginalPath: filename, // For batch uploads, filename contains the full original path FolderID: folderID, }) if err != nil { return false } slog.Debug("created document", "id", result.ID, "filename", filename, "folder_id", folderID) // Report outcome for batch documents outcomeName := repository.BatchOutcomeStatusInitComplete if result.IsDuplicate { outcomeName = repository.BatchOutcomeStatusInitDuplicate } var resolveFilename string if result.Filename != nil { resolveFilename = *result.Filename } if resolveErr := s.svc.Outcome.Resolve(ctx, result.BatchID, resolveFilename, result.ID, outcomeName, nil); resolveErr != nil { slog.Error("failed to resolve batch outcome", "document_id", result.ID, "outcome", outcomeName, "error", resolveErr) // Non-fatal: document was created successfully } return true }