Files
query-orchestration/api/docInitRunner/runner.go
T
Jacob Mathison a080ca59d8 Merged in jmathison/v2-upload-batch (pull request #224)
Implement v2 upload batch cleanup

* Implement v2 upload batch cleanup

* Merge remote-tracking branch 'origin/main' into jmathison/v2-upload-batch

* Address upload batch review feedback

* Raise batch worker coverage

* Fix batch cleanup review issues


Approved-by: Jay Brown
2026-05-08 21:42:46 +00:00

123 lines
3.5 KiB
Go

package docinitrunner
import (
"context"
"database/sql"
"errors"
"log/slog"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/document/batch/foldercleanup"
"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
Cleanup *foldercleanup.Service
}
// 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
}
if outcomeName == repository.BatchOutcomeStatusInitDuplicate {
if s.svc.Cleanup != nil && result.BatchID != nil {
if _, err := s.svc.Cleanup.TryFinalizeBatch(ctx, *result.BatchID); err != nil {
slog.Error("failed to finalize batch folder cleanup", "batch_id", *result.BatchID, "error", err)
}
}
}
return true
}