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
This commit is contained in:
Jacob Mathison
2026-05-08 21:42:46 +00:00
parent 7b820b18c2
commit a080ca59d8
35 changed files with 2587 additions and 108 deletions
+137 -41
View File
@@ -14,18 +14,30 @@ import (
"queryorchestration/internal/database/repository"
"queryorchestration/internal/document/batch"
"queryorchestration/internal/document/batch/foldercleanup"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/google/uuid"
)
type batchWorkerService interface {
AddFailedFilename(ctx context.Context, batchID uuid.UUID, filename string) error
CheckDuplicate(ctx context.Context, clientID string, hash string) (*uuid.UUID, error)
MarkCompleted(ctx context.Context, batchID uuid.UUID) error
MarkFailed(ctx context.Context, batchID uuid.UUID) error
RecordOutcome(ctx context.Context, batchID uuid.UUID, filename string, outcome repository.BatchOutcomeStatus, errorDetail *string, documentID *uuid.UUID) error
SetTotalDocuments(ctx context.Context, batchID uuid.UUID, total int32) error
UpdateProgress(ctx context.Context, clientID string, batchID uuid.UUID, processed int32, failed int32, invalid int32) error
}
// Config key constants for batch worker
const (
ConfigKeyBatchService = "batchService"
ConfigKeyS3Client = "s3Client"
ConfigKeyBucket = "bucket"
ConfigKeyUploadHandler = "uploadHandler"
ConfigKeyFolderCleanup = "folderCleanupService"
)
// acceptedExtensions defines the file extensions allowed in batch uploads.
@@ -115,6 +127,15 @@ func processBatchWork(ctx context.Context, logger *slog.Logger, config map[strin
return fmt.Errorf("upload handler not found in config")
}
cleanupService, _ := config[ConfigKeyFolderCleanup].(*foldercleanup.Service)
if cleanupService != nil {
if finalized, err := cleanupService.FinalizeReadyBatches(ctx, foldercleanup.DefaultReadyBatchLimit); err != nil {
logger.Error("Failed to finalize ready batch folders", "error", err)
} else if finalized > 0 {
logger.Info("Finalized ready batch folders", "count", finalized)
}
}
// Query unprocessed batches
batches, err := batchService.ListUnprocessed(ctx)
if err != nil {
@@ -136,7 +157,7 @@ func processBatchWork(ctx context.Context, logger *slog.Logger, config map[strin
continue
}
if err := processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails); err != nil {
if err := processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, cleanupService); err != nil {
logger.Error("Failed to process batch",
"batch_id", batchDetails.ID,
"client_id", batchDetails.ClientID,
@@ -146,6 +167,14 @@ func processBatchWork(ctx context.Context, logger *slog.Logger, config map[strin
}
}
if cleanupService != nil {
if deleted, err := cleanupService.DeleteSafeAutoCreatedFolders(ctx, foldercleanup.DefaultFolderDeleteLimit); err != nil {
logger.Error("Failed to delete safe auto-created folders", "error", err)
} else if deleted > 0 {
logger.Info("Deleted safe auto-created folders", "count", deleted)
}
}
return nil
}
@@ -153,11 +182,12 @@ func processBatchWork(ctx context.Context, logger *slog.Logger, config map[strin
func processSingleBatch(
ctx context.Context,
logger *slog.Logger,
batchService *batch.Service,
batchService batchWorkerService,
s3Client *s3.Client,
bucket string,
uploadHandler func(ctx context.Context, clientID string, docData io.Reader, filename string, batchID *uuid.UUID) error,
batchInfo *batch.BatchUploadDetails,
cleanupService *foldercleanup.Service,
) error {
logger.Info("Processing batch",
"batch_id", batchInfo.ID,
@@ -171,6 +201,7 @@ func processSingleBatch(
if err != nil {
// Mark batch as failed if we can't download the zip
_ = batchService.MarkFailed(ctx, batchInfo.ID)
tryFinalizeBatchFolders(ctx, logger, cleanupService, batchInfo.ID)
return fmt.Errorf("failed to download zip from S3: %w", err)
}
@@ -188,17 +219,34 @@ func processSingleBatch(
zipReader, err := zip.NewReader(bytes.NewReader(zipData), int64(len(zipData)))
if err != nil {
_ = batchService.MarkFailed(ctx, batchInfo.ID)
tryFinalizeBatchFolders(ctx, logger, cleanupService, batchInfo.ID)
return fmt.Errorf("failed to open zip archive: %w", err)
}
// Process each file in the zip
for _, file := range zipReader.File {
processed, failed, invalid := processZipFile(ctx, logger, batchService, s3Client, bucket, extractPath, uploadHandler, batchInfo, file)
processed, failed, invalid, err := processZipFile(ctx, logger, batchService, s3Client, bucket, extractPath, uploadHandler, batchInfo, file)
if err != nil {
_ = batchService.MarkFailed(ctx, batchInfo.ID)
tryFinalizeBatchFolders(ctx, logger, cleanupService, batchInfo.ID)
return fmt.Errorf("failed to process zip file %q: %w", file.Name, err)
}
processedCount += processed
failedCount += failed
invalidCount += invalid
}
totalProcessed := processedCount + failedCount + invalidCount
if totalProcessed == 0 {
if err := recordNoProcessableFilesOutcome(ctx, batchService, batchInfo); err != nil {
_ = batchService.MarkFailed(ctx, batchInfo.ID)
tryFinalizeBatchFolders(ctx, logger, cleanupService, batchInfo.ID)
return fmt.Errorf("record no-processable-files outcome: %w", err)
}
failedCount = 1
totalProcessed = 1
}
// Update batch progress - convert safely to int32
const maxInt32 = 2147483647
var processedInt32, failedInt32, invalidInt32 int32
@@ -229,8 +277,6 @@ func processSingleBatch(
logger.Error("Failed to update batch progress", "batch_id", batchInfo.ID, "error", err)
}
// Mark batch as completed if any files were processed successfully
totalProcessed := processedCount + failedCount + invalidCount
if totalProcessed > 0 {
if processedCount > 0 {
// Mark as completed if at least one document was processed successfully
@@ -243,6 +289,7 @@ func processSingleBatch(
logger.Error("Failed to mark batch as failed", "batch_id", batchInfo.ID, "error", err)
}
}
tryFinalizeBatchFolders(ctx, logger, cleanupService, batchInfo.ID)
}
logger.Info("Batch processing completed",
@@ -254,6 +301,15 @@ func processSingleBatch(
return nil
}
func tryFinalizeBatchFolders(ctx context.Context, logger *slog.Logger, cleanupService *foldercleanup.Service, batchID uuid.UUID) {
if cleanupService == nil {
return
}
if _, err := cleanupService.TryFinalizeBatch(ctx, batchID); err != nil {
logger.Error("Failed to finalize batch folder cleanup", "batch_id", batchID, "error", err)
}
}
// downloadZipFromS3 downloads a zip file from S3
func downloadZipFromS3(ctx context.Context, s3Client *s3.Client, bucket, key string) ([]byte, error) {
result, err := s3Client.GetObject(ctx, &s3.GetObjectInput{
@@ -291,33 +347,23 @@ func uploadToS3(ctx context.Context, s3Client *s3.Client, bucket, key string, da
func processZipFile(
ctx context.Context,
logger *slog.Logger,
batchService *batch.Service,
batchService batchWorkerService,
s3Client *s3.Client,
bucket, extractPath string,
uploadHandler func(ctx context.Context, clientID string, docData io.Reader, filename string, batchID *uuid.UUID) error,
batchInfo *batch.BatchUploadDetails,
file *zip.File,
) (int, int, int) {
) (int, int, int, error) {
if file.FileInfo().IsDir() {
return 0, 0, 0
return 0, 0, 0, nil
}
// Skip system files and hidden files (but allow relative path prefixes like ./file.pdf)
filename := file.Name
if strings.HasPrefix(filename, "__MACOSX/") {
return 0, 0, 0
}
// Skip hidden files (files that start with . but are not relative paths like ./file)
if strings.HasPrefix(filename, ".") && !strings.HasPrefix(filename, "./") {
return 0, 0, 0
}
// Security: Clean filename to prevent path traversal attacks while preserving internal structure
filename = sanitizeZipPath(filename)
if filename == "" {
logger.Warn("Skipping file with invalid path", "original", file.Name)
return 0, 0, 0
filename, skip := batchZipFilename(file.Name)
if skip {
if filename == "" {
logger.Warn("Skipping file with invalid path", "original", file.Name)
}
return 0, 0, 0, nil
}
// Extract file content
@@ -326,8 +372,10 @@ func processZipFile(
logger.Error("Failed to open file in zip", "filename", filename, "error", err)
_ = batchService.AddFailedFilename(ctx, batchInfo.ID, filename)
errMsg := err.Error()
_ = batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedOpen, &errMsg, nil)
return 0, 1, 0
if recordErr := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedOpen, &errMsg, nil); recordErr != nil {
return 0, 0, 0, fmt.Errorf("record failed_open outcome: %w", recordErr)
}
return 0, 1, 0, nil
}
fileData, err := io.ReadAll(fileReader)
@@ -336,8 +384,10 @@ func processZipFile(
logger.Error("Failed to read file from zip", "filename", filename, "error", err)
_ = batchService.AddFailedFilename(ctx, batchInfo.ID, filename)
errMsg := err.Error()
_ = batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedRead, &errMsg, nil)
return 0, 1, 0
if recordErr := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedRead, &errMsg, nil); recordErr != nil {
return 0, 0, 0, fmt.Errorf("record failed_read outcome: %w", recordErr)
}
return 0, 1, 0, nil
}
// Store extracted file in S3 temporary directory (filename now includes sanitized path)
@@ -346,16 +396,20 @@ func processZipFile(
logger.Error("Failed to upload extracted file to S3", "filename", filename, "error", err)
_ = batchService.AddFailedFilename(ctx, batchInfo.ID, filename)
errMsg := err.Error()
_ = batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedS3Upload, &errMsg, nil)
return 0, 1, 0
if recordErr := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedS3Upload, &errMsg, nil); recordErr != nil {
return 0, 0, 0, fmt.Errorf("record failed_s3_upload outcome: %w", recordErr)
}
return 0, 1, 0, nil
}
// Check if file type is supported
if !isAcceptedExtension(filename) {
logger.Warn("Invalid file type", "filename", filename)
_ = batchService.AddFailedFilename(ctx, batchInfo.ID, filename)
_ = batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusInvalidType, nil, nil)
return 0, 0, 1
if recordErr := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusInvalidType, nil, nil); recordErr != nil {
return 0, 0, 0, fmt.Errorf("record invalid_type outcome: %w", recordErr)
}
return 0, 0, 1, nil
}
// Call the document upload handler with batch ID
@@ -363,20 +417,62 @@ func processZipFile(
logger.Error("Failed to upload document", "filename", filename, "error", err)
_ = batchService.AddFailedFilename(ctx, batchInfo.ID, filename)
errMsg := err.Error()
_ = batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedUpload, &errMsg, nil)
return 0, 1, 0
if recordErr := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedUpload, &errMsg, nil); recordErr != nil {
return 0, 0, 0, fmt.Errorf("record failed_upload outcome: %w", recordErr)
}
return 0, 1, 0, nil
}
// Check for duplicate by computing file hash and looking up existing document
if err := recordSuccessfulBatchOutcome(ctx, batchService, batchInfo, filename, fileData); err != nil {
return 0, 0, 0, err
}
logger.Info("Successfully processed file", "filename", filename, "batch_id", batchInfo.ID)
return 1, 0, 0, nil
}
func recordNoProcessableFilesOutcome(ctx context.Context, batchService batchWorkerService, batchInfo *batch.BatchUploadDetails) error {
filename := batchInfo.OriginalFilename
if filename == "" {
filename = batchInfo.ID.String() + ".zip"
}
errMsg := "archive contained no processable files"
_ = batchService.AddFailedFilename(ctx, batchInfo.ID, filename)
if err := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusFailedUpload, &errMsg, nil); err != nil {
return fmt.Errorf("record failed_upload outcome: %w", err)
}
return nil
}
func batchZipFilename(filename string) (string, bool) {
if strings.HasPrefix(filename, "__MACOSX/") {
return filename, true
}
if strings.HasPrefix(filename, ".") && !strings.HasPrefix(filename, "./") {
return filename, true
}
cleanFilename := sanitizeZipPath(filename)
return cleanFilename, cleanFilename == ""
}
func recordSuccessfulBatchOutcome(
ctx context.Context,
batchService batchWorkerService,
batchInfo *batch.BatchUploadDetails,
filename string,
fileData []byte,
) error {
hash := sha256.Sum256(fileData)
hashStr := hex.EncodeToString(hash[:])
existingID, _ := batchService.CheckDuplicate(ctx, batchInfo.ClientID, hashStr)
if existingID != nil {
_ = batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusDuplicate, nil, existingID)
} else {
_ = batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusSubmitted, nil, nil)
if err := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusDuplicate, nil, existingID); err != nil {
return fmt.Errorf("record duplicate outcome: %w", err)
}
return nil
}
logger.Info("Successfully processed file", "filename", filename, "batch_id", batchInfo.ID)
return 1, 0, 0
if err := batchService.RecordOutcome(ctx, batchInfo.ID, filename, repository.BatchOutcomeStatusSubmitted, nil, nil); err != nil {
return fmt.Errorf("record submitted outcome: %w", err)
}
return nil
}
+430 -7
View File
@@ -6,14 +6,17 @@ import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"io"
"log/slog"
"os"
"strings"
"testing"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/document/batch"
"queryorchestration/internal/document/batch/foldercleanup"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/serviceconfig/objectstore"
"queryorchestration/internal/test"
@@ -21,10 +24,20 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
type outcomeFailingBatchService struct {
*batch.Service
}
func (s *outcomeFailingBatchService) RecordOutcome(ctx context.Context, batchID uuid.UUID, filename string,
outcome repository.BatchOutcomeStatus, errorDetail *string, documentID *uuid.UUID) error {
return errors.New("forced outcome write failure")
}
// TestProcessBatchWork tests the main batch processing work function
func TestProcessBatchWork(t *testing.T) {
if testing.Short() {
@@ -110,6 +123,33 @@ func TestProcessBatchWork(t *testing.T) {
}
}
func TestProcessBatchWorkRunsFolderCleanupWithoutBatches(t *testing.T) {
ctx := t.Context()
cfg := &BaseConfig{}
test.CreateDB(t, cfg)
workerConfig := map[string]any{
ConfigKeyBatchService: batch.New(cfg),
ConfigKeyS3Client: &s3.Client{},
ConfigKeyBucket: "bucket",
ConfigKeyUploadHandler: func(context.Context, string, io.Reader, string, *uuid.UUID) error { return nil },
ConfigKeyFolderCleanup: foldercleanup.New(cfg),
}
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
require.NoError(t, processBatchWork(ctx, logger, workerConfig))
}
func TestTryFinalizeBatchFoldersHandlesNilAndErrors(t *testing.T) {
ctx := t.Context()
cfg := &BaseConfig{}
test.CreateDB(t, cfg)
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
tryFinalizeBatchFolders(ctx, logger, nil, uuid.New())
tryFinalizeBatchFolders(ctx, logger, foldercleanup.New(cfg), uuid.New())
}
// TestProcessSingleBatch tests processing a single batch
func TestProcessSingleBatch(t *testing.T) {
if testing.Short() {
@@ -169,7 +209,7 @@ func TestProcessSingleBatch(t *testing.T) {
// Process the batch
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails)
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, nil)
require.NoError(t, err)
// All three file types (PDF, TXT, DOCX) are now accepted
@@ -253,7 +293,7 @@ func TestProcessSingleBatchWithFailure(t *testing.T) {
// Process the batch
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails)
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, nil)
require.NoError(t, err)
// Verify results
@@ -357,6 +397,61 @@ func createTestZIPForWorker(t *testing.T, numFiles int) []byte {
return buf.Bytes()
}
func firstZipFile(t *testing.T, files map[string][]byte) *zip.File {
t.Helper()
zipFiles := zipFiles(t, files)
require.NotEmpty(t, zipFiles)
return zipFiles[0]
}
func zipFiles(t *testing.T, files map[string][]byte) []*zip.File {
t.Helper()
var buf bytes.Buffer
zipWriter := zip.NewWriter(&buf)
for name, content := range files {
if content == nil && strings.HasSuffix(name, "/") {
header := &zip.FileHeader{Name: name}
header.SetMode(os.ModeDir | 0o755)
_, err := zipWriter.CreateHeader(header)
require.NoError(t, err)
continue
}
fileWriter, err := zipWriter.Create(name)
require.NoError(t, err)
_, err = fileWriter.Write(content)
require.NoError(t, err)
}
require.NoError(t, zipWriter.Close())
reader, err := zip.NewReader(bytes.NewReader(buf.Bytes()), int64(buf.Len()))
require.NoError(t, err)
return reader.File
}
func zipBytes(t *testing.T, files map[string][]byte) []byte {
t.Helper()
var buf bytes.Buffer
zipWriter := zip.NewWriter(&buf)
for name, content := range files {
if content == nil && strings.HasSuffix(name, "/") {
header := &zip.FileHeader{Name: name}
header.SetMode(os.ModeDir | 0o755)
_, err := zipWriter.CreateHeader(header)
require.NoError(t, err)
continue
}
fileWriter, err := zipWriter.Create(name)
require.NoError(t, err)
_, err = fileWriter.Write(content)
require.NoError(t, err)
}
require.NoError(t, zipWriter.Close())
return buf.Bytes()
}
// Helper function to create a ZIP with mixed content types
func createMixedContentZIP(t *testing.T) []byte {
var buf bytes.Buffer
@@ -539,7 +634,7 @@ func TestProcessBatchWithNestedFolders(t *testing.T) {
// Process the batch
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails)
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, nil)
require.NoError(t, err)
// Verify batch status
@@ -677,7 +772,7 @@ func TestProcessBatchWithPathTraversal(t *testing.T) {
// Process the batch
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails)
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, nil)
require.NoError(t, err)
// Verify that path traversal attempts were blocked
@@ -757,6 +852,334 @@ func TestProcessZipFile_RecordsOutcomes(t *testing.T) {
}
}
func TestProcessZipFile_RecordsFailedS3UploadOutcome(t *testing.T) {
if testing.Short() {
t.SkipNow()
}
ctx := t.Context()
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
test.CreateDB(t, cfg)
test.CreateAWSResources(t, cfg)
clientID := "test_zip_failed_s3"
test.CreateTestClient(t, cfg, clientID, "Test ZIP Failed S3 Client")
batchService := batch.New(cfg)
s3ClientInterface := cfg.GetStoreClient()
s3Client, ok := s3ClientInterface.(*s3.Client)
require.True(t, ok)
batchID, err := batchService.CreateWithStorage(ctx, clientID, "failed-s3.zip", 1, cfg.GetBucket(), "archive.zip", 1)
require.NoError(t, err)
file := firstZipFile(t, map[string][]byte{
"document.pdf": []byte("%%PDF-1.4\n%%Test PDF\n%%%%EOF"),
})
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
processed, failed, invalid, err := processZipFile(
ctx,
logger,
batchService,
s3Client,
"missing-bucket",
"extract/path",
func(context.Context, string, io.Reader, string, *uuid.UUID) error {
t.Fatal("upload handler should not be called when extracted S3 upload fails")
return nil
},
&batch.BatchUploadDetails{
BatchUploadSummary: batch.BatchUploadSummary{ID: batchID, ClientID: clientID},
},
file,
)
require.NoError(t, err)
assert.Equal(t, 0, processed)
assert.Equal(t, 1, failed)
assert.Equal(t, 0, invalid)
outcomes, err := batchService.ListOutcomes(ctx, batchID)
require.NoError(t, err)
require.Len(t, outcomes, 1)
assert.Equal(t, string(repository.BatchOutcomeStatusFailedS3Upload), outcomes[0].Outcome)
assert.Equal(t, "document.pdf", outcomes[0].Filename)
require.NotNil(t, outcomes[0].ErrorDetail)
assert.Contains(t, *outcomes[0].ErrorDetail, "failed to put object to S3")
}
func TestProcessZipFile_RecordsInvalidTypeOutcome(t *testing.T) {
if testing.Short() {
t.SkipNow()
}
ctx := t.Context()
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
test.CreateDB(t, cfg)
test.CreateAWSResources(t, cfg)
clientID := "test_zip_invalid_type"
test.CreateTestClient(t, cfg, clientID, "Test ZIP Invalid Type Client")
batchService := batch.New(cfg)
s3ClientInterface := cfg.GetStoreClient()
s3Client, ok := s3ClientInterface.(*s3.Client)
require.True(t, ok)
batchID, err := batchService.CreateWithStorage(ctx, clientID, "invalid-type.zip", 1, cfg.GetBucket(), "archive.zip", 1)
require.NoError(t, err)
file := firstZipFile(t, map[string][]byte{
"unsupported.exe": []byte("not a supported document"),
})
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
processed, failed, invalid, err := processZipFile(
ctx,
logger,
batchService,
s3Client,
cfg.GetBucket(),
"extract/path",
func(context.Context, string, io.Reader, string, *uuid.UUID) error {
t.Fatal("upload handler should not be called for invalid file types")
return nil
},
&batch.BatchUploadDetails{
BatchUploadSummary: batch.BatchUploadSummary{ID: batchID, ClientID: clientID},
},
file,
)
require.NoError(t, err)
assert.Equal(t, 0, processed)
assert.Equal(t, 0, failed)
assert.Equal(t, 1, invalid)
outcomes, err := batchService.ListOutcomes(ctx, batchID)
require.NoError(t, err)
require.Len(t, outcomes, 1)
assert.Equal(t, string(repository.BatchOutcomeStatusInvalidType), outcomes[0].Outcome)
assert.Equal(t, "unsupported.exe", outcomes[0].Filename)
}
func TestProcessZipFile_SkipsDirectoriesAndHiddenEntries(t *testing.T) {
ctx := t.Context()
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
files := zipFiles(t, map[string][]byte{
"folder/": nil,
"__MACOSX/document.pdf": []byte("metadata"),
".DS_Store": []byte("metadata"),
"../escape.pdf": []byte("blocked"),
})
for _, file := range files {
processed, failed, invalid, err := processZipFile(
ctx,
logger,
nil,
nil,
"",
"",
func(context.Context, string, io.Reader, string, *uuid.UUID) error {
t.Fatal("upload handler should not be called for skipped entries")
return nil
},
&batch.BatchUploadDetails{},
file,
)
require.NoError(t, err)
assert.Equal(t, 0, processed)
assert.Equal(t, 0, failed)
assert.Equal(t, 0, invalid)
}
}
func TestOutcomeWriteFailureDoesNotCompleteBatch(t *testing.T) {
if testing.Short() {
t.SkipNow()
}
ctx := t.Context()
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
test.CreateDB(t, cfg)
test.CreateAWSResources(t, cfg)
clientID := "test_outcome_write_failure"
test.CreateTestClient(t, cfg, clientID, "Test Outcome Write Failure Client")
batchService := batch.New(cfg)
failingBatchService := &outcomeFailingBatchService{Service: batchService}
s3ClientInterface := cfg.GetStoreClient()
s3Client, ok := s3ClientInterface.(*s3.Client)
require.True(t, ok)
bucket := cfg.GetBucket()
zipContent := createTestZIPForWorker(t, 1)
archiveKey := fmt.Sprintf("test/%s/outcome-failure.zip", clientID)
_, err := s3Client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(bucket),
Key: aws.String(archiveKey),
Body: bytes.NewReader(zipContent),
})
require.NoError(t, err)
batchID, err := batchService.CreateWithStorage(ctx, clientID, "outcome-failure.zip", 0, bucket, archiveKey, int64(len(zipContent)))
require.NoError(t, err)
folder, err := cfg.GetDBQueries().CreateFolder(ctx, &repository.CreateFolderParams{
Path: "outcome-failure-folder",
Parentid: nil,
Clientid: clientID,
Createdby: "test",
})
require.NoError(t, err)
require.NoError(t, cfg.GetDBQueries().RecordBatchFolderCandidate(ctx, &repository.RecordBatchFolderCandidateParams{
BatchID: batchID,
FolderID: folder.ID,
Path: folder.Path,
ExistedBefore: false,
}))
uploadHandler := func(ctx context.Context, clientID string, docData io.Reader, filename string, batchID *uuid.UUID) error {
_, err := io.ReadAll(docData)
return err
}
batchDetails := &batch.BatchUploadDetails{
BatchUploadSummary: batch.BatchUploadSummary{
ID: batchID,
ClientID: clientID,
OriginalFilename: "outcome-failure.zip",
Status: "processing",
},
ArchiveKey: archiveKey,
ArchiveBucket: bucket,
}
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
cleanupService := foldercleanup.New(cfg)
err = processSingleBatch(ctx, logger, failingBatchService, s3Client, bucket, uploadHandler, batchDetails, cleanupService)
require.Error(t, err)
batchInfo, err := cfg.GetDBQueries().GetBatchUploadWithStorage(ctx, &repository.GetBatchUploadWithStorageParams{
ID: batchID,
ClientID: clientID,
})
require.NoError(t, err)
assert.Equal(t, repository.BatchStatusFailed, batchInfo.Status)
assert.Zero(t, batchInfo.TotalDocuments, "failed outcome writes must not count toward batch totals")
outcomeCount, err := cfg.GetDBQueries().CountBatchOutcomes(ctx, batchID)
require.NoError(t, err)
assert.Zero(t, outcomeCount)
cleanupBatch, err := cfg.GetDBQueries().LockBatchUploadForFolderCleanup(ctx, batchID)
require.NoError(t, err)
assert.False(t, cleanupBatch.FolderCleanupCompletedAt.Valid)
_, err = cfg.GetDBQueries().GetFolderByID(ctx, folder.ID)
require.NoError(t, err)
}
func TestAllSkippedZipRecordsFailedOutcomeAndFinalizes(t *testing.T) {
if testing.Short() {
t.SkipNow()
}
ctx := t.Context()
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
test.CreateDB(t, cfg)
test.CreateAWSResources(t, cfg)
clientID := "test_all_skipped_zip"
test.CreateTestClient(t, cfg, clientID, "Test All Skipped ZIP Client")
batchService := batch.New(cfg)
s3ClientInterface := cfg.GetStoreClient()
s3Client, ok := s3ClientInterface.(*s3.Client)
require.True(t, ok)
bucket := cfg.GetBucket()
zipContent := zipBytes(t, map[string][]byte{
"folder/": nil,
"__MACOSX/document.pdf": []byte("metadata"),
".DS_Store": []byte("metadata"),
"../escape.pdf": []byte("blocked"),
})
archiveKey := fmt.Sprintf("test/%s/all-skipped.zip", clientID)
_, err := s3Client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(bucket),
Key: aws.String(archiveKey),
Body: bytes.NewReader(zipContent),
})
require.NoError(t, err)
batchID, err := batchService.CreateWithStorage(ctx, clientID, "all-skipped.zip", 0, bucket, archiveKey, int64(len(zipContent)))
require.NoError(t, err)
folder, err := cfg.GetDBQueries().CreateFolder(ctx, &repository.CreateFolderParams{
Path: "/all-skipped",
Parentid: nil,
Clientid: clientID,
Createdby: "test",
})
require.NoError(t, err)
require.NoError(t, cfg.GetDBQueries().RecordBatchFolderCandidate(ctx, &repository.RecordBatchFolderCandidateParams{
BatchID: batchID,
FolderID: folder.ID,
Path: folder.Path,
ExistedBefore: false,
}))
batchDetails := &batch.BatchUploadDetails{
BatchUploadSummary: batch.BatchUploadSummary{
ID: batchID,
ClientID: clientID,
OriginalFilename: "all-skipped.zip",
Status: "processing",
},
ArchiveKey: archiveKey,
ArchiveBucket: bucket,
}
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket,
func(context.Context, string, io.Reader, string, *uuid.UUID) error {
t.Fatal("upload handler should not be called for skipped-only archives")
return nil
},
batchDetails,
foldercleanup.New(cfg),
)
require.NoError(t, err)
batchInfo, err := cfg.GetDBQueries().GetBatchUploadWithStorage(ctx, &repository.GetBatchUploadWithStorageParams{
ID: batchID,
ClientID: clientID,
})
require.NoError(t, err)
assert.Equal(t, repository.BatchStatusFailed, batchInfo.Status)
assert.Equal(t, int32(1), batchInfo.TotalDocuments)
assert.Equal(t, int32(1), batchInfo.FailedDocuments)
outcomes, err := batchService.ListOutcomes(ctx, batchID)
require.NoError(t, err)
require.Len(t, outcomes, 1)
assert.Equal(t, "all-skipped.zip", outcomes[0].Filename)
assert.Equal(t, string(repository.BatchOutcomeStatusFailedUpload), outcomes[0].Outcome)
require.NotNil(t, outcomes[0].ErrorDetail)
assert.Contains(t, *outcomes[0].ErrorDetail, "archive contained no processable files")
unprocessed, err := cfg.GetDBQueries().GetUnprocessedBatches(ctx)
require.NoError(t, err)
for _, batch := range unprocessed {
assert.NotEqual(t, batchID, batch.ID)
}
cleanupBatch, err := cfg.GetDBQueries().LockBatchUploadForFolderCleanup(ctx, batchID)
require.NoError(t, err)
assert.True(t, cleanupBatch.FolderCleanupCompletedAt.Valid)
_, err = cfg.GetDBQueries().GetFolderByID(ctx, folder.ID)
assert.True(t, errors.Is(err, pgx.ErrNoRows))
}
// TestProcessZipFile_DetectsDuplicates verifies that when a file in the batch
// has the same content hash as an existing document in the system, it is
// recorded as a duplicate with the existing document's ID.
@@ -999,7 +1422,7 @@ func TestProcessBatchWithDuplicateFilenames(t *testing.T) {
// Process the batch
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails)
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, nil)
require.NoError(t, err)
// Verify all files were uploaded with correct paths
@@ -1066,7 +1489,7 @@ func TestProcessBatchWithAllSupportedFileTypes(t *testing.T) {
}
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails)
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, nil)
require.NoError(t, err)
// Verify batch completed successfully
@@ -1166,7 +1589,7 @@ func TestProcessBatchWithMixedValidInvalidFileTypes(t *testing.T) {
}
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails)
err = processSingleBatch(ctx, logger, batchService, s3Client, bucket, uploadHandler, batchDetails, nil)
require.NoError(t, err)
// Batch should still complete (not fail) because some files succeeded
+10
View File
@@ -16,6 +16,7 @@ import (
"queryorchestration/internal/backgroundtask"
"queryorchestration/internal/cognitoauth"
"queryorchestration/internal/document/batch"
"queryorchestration/internal/document/batch/foldercleanup"
documentupload "queryorchestration/internal/document/upload"
awsc "queryorchestration/internal/serviceconfig/aws"
"queryorchestration/internal/serviceconfig/build"
@@ -224,8 +225,16 @@ func New(ctx context.Context, cfg Config) (*Server, error) {
return nil, err
}
if err := ensureLocalDevClient(ctx, cfg); err != nil {
return nil, err
}
if err := ensureLocalDevEula(ctx, cfg); err != nil {
return nil, err
}
// Initialize background task runner for batch processing
batchService := batch.New(cfg)
cleanupService := foldercleanup.New(cfg)
// Create upload handler function that calls the document upload service
uploadHandler := createDocumentUploadHandler(cfg)
@@ -236,6 +245,7 @@ func New(ctx context.Context, cfg Config) (*Server, error) {
ConfigKeyS3Client: cfg.GetStoreClient(),
ConfigKeyBucket: cfg.GetBucket(),
ConfigKeyUploadHandler: uploadHandler,
ConfigKeyFolderCleanup: cleanupService,
}
runner, err := backgroundtask.Initialize(
+127
View File
@@ -0,0 +1,127 @@
package api
import (
"context"
"errors"
"log/slog"
"os"
"strings"
"queryorchestration/internal/client"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgconn"
)
const (
localDevClientIDEnv = "LOCAL_DEV_CLIENT_ID"
localDevClientNameEnv = "LOCAL_DEV_CLIENT_NAME"
localDevEulaVersionEnv = "LOCAL_DEV_EULA_VERSION"
localDevEulaTitleEnv = "LOCAL_DEV_EULA_TITLE"
localDevEulaAutoAcceptSubjectEnv = "LOCAL_DEV_EULA_AUTO_ACCEPT_SUBJECT"
localDevEulaAutoAcceptEmailEnv = "LOCAL_DEV_EULA_AUTO_ACCEPT_EMAIL"
)
const localDevEulaContent = `<h3>Acceptance of Terms</h3><p>By accessing and using AArete DoczyAI in this local development environment, you agree to use the application only for authorized testing and evaluation.</p><h3>Authorized Use</h3><p>You are responsible for ensuring uploaded documents and generated outputs are appropriate for local development use and comply with applicable internal policies.</p><h3>Data Handling</h3><p>This local instance is intended for development workflows. Do not upload production, confidential, or regulated data unless the environment has been approved for that data.</p><h3>No Warranty</h3><p>This local development instance is provided as-is for testing. Features, data, and outputs may change during development.</p>`
type localDevClientConfig interface {
serviceconfig.ConfigProvider
GetLogger() *slog.Logger
}
func ensureLocalDevClient(ctx context.Context, cfg localDevClientConfig) error {
if strings.ToLower(os.Getenv("DISABLE_AUTH")) != "true" {
return nil
}
clientID := strings.TrimSpace(os.Getenv(localDevClientIDEnv))
if clientID == "" {
return nil
}
clientName := strings.TrimSpace(os.Getenv(localDevClientNameEnv))
if clientName == "" {
clientName = clientID
}
_, err := cfg.GetDBQueries().GetClient(ctx, clientID)
if err == nil {
return nil
}
if !errors.Is(err, pgx.ErrNoRows) {
var pgErr *pgconn.PgError
if errors.As(err, &pgErr) && pgErr.Code == "42P01" {
localDevBootstrapLogger(cfg).Warn("Skipping local dev client bootstrap because client tables are unavailable")
return nil
}
return err
}
_, err = client.New(cfg).Create(ctx, client.CreateParams{
ID: clientID,
Name: clientName,
})
if err != nil {
return err
}
localDevBootstrapLogger(cfg).Info("Created local dev client", slog.String("client_id", clientID))
return nil
}
func ensureLocalDevEula(ctx context.Context, cfg localDevClientConfig) error {
if strings.ToLower(os.Getenv("DISABLE_AUTH")) != "true" {
return nil
}
version := strings.TrimSpace(os.Getenv(localDevEulaVersionEnv))
if version == "" {
return nil
}
title := strings.TrimSpace(os.Getenv(localDevEulaTitleEnv))
if title == "" {
title = "AArete DoczyAI Terms of Use"
}
versionID, err := cfg.GetDBQueries().ActivateLocalDevEulaVersion(ctx, &repository.ActivateLocalDevEulaVersionParams{
Version: version,
Title: title,
Content: localDevEulaContent,
})
if err != nil {
var pgErr *pgconn.PgError
if errors.As(err, &pgErr) && pgErr.Code == "42P01" {
localDevBootstrapLogger(cfg).Warn("Skipping local dev EULA bootstrap because EULA tables are unavailable")
return nil
}
return err
}
subject := strings.TrimSpace(os.Getenv(localDevEulaAutoAcceptSubjectEnv))
email := strings.TrimSpace(os.Getenv(localDevEulaAutoAcceptEmailEnv))
if subject != "" {
if email == "" {
email = subject
}
if err := cfg.GetDBQueries().CreateLocalDevEulaAgreement(ctx, &repository.CreateLocalDevEulaAgreementParams{
CognitoSubjectID: subject,
UserEmail: email,
EulaVersionID: versionID,
}); err != nil {
return err
}
}
localDevBootstrapLogger(cfg).Info("Activated local dev EULA", slog.String("version", version))
return nil
}
func localDevBootstrapLogger(cfg localDevClientConfig) *slog.Logger {
if logger := cfg.GetLogger(); logger != nil {
return logger
}
return slog.Default()
}
@@ -0,0 +1,120 @@
package api
import (
"testing"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/serviceconfig"
"queryorchestration/internal/test"
"github.com/stretchr/testify/require"
)
func TestEnsureLocalDevClientCreatesConfiguredClient(t *testing.T) {
if testing.Short() {
t.SkipNow()
}
t.Setenv("DISABLE_AUTH", "true")
t.Setenv(localDevClientIDEnv, "mathisonprojects-local")
t.Setenv(localDevClientNameEnv, "Mathison Projects Local")
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
test.CreateDB(t, cfg)
err := ensureLocalDevClient(t.Context(), cfg)
require.NoError(t, err)
dbClient, err := cfg.GetDBQueries().GetClient(t.Context(), "mathisonprojects-local")
require.NoError(t, err)
require.Equal(t, "Mathison Projects Local", dbClient.Name)
rootFolder, err := cfg.GetDBQueries().GetFolderByPath(t.Context(), &repository.GetFolderByPathParams{
Clientid: "mathisonprojects-local",
Path: "/",
})
require.NoError(t, err)
require.Equal(t, "/", rootFolder.Path)
require.NoError(t, ensureLocalDevClient(t.Context(), cfg))
}
func TestEnsureLocalDevClientSkipsWhenUnconfigured(t *testing.T) {
t.Setenv("DISABLE_AUTH", "true")
t.Setenv(localDevClientIDEnv, "")
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
require.NoError(t, ensureLocalDevClient(t.Context(), cfg))
}
func TestEnsureLocalDevEula(t *testing.T) {
if testing.Short() {
t.SkipNow()
}
t.Setenv("DISABLE_AUTH", "true")
t.Setenv(localDevEulaVersionEnv, "local-2026-05-07")
t.Setenv(localDevEulaTitleEnv, "AArete DoczyAI Terms of Use")
t.Setenv(localDevEulaAutoAcceptSubjectEnv, "jacob+2@mathisonprojects.com")
t.Setenv(localDevEulaAutoAcceptEmailEnv, "jacob+2@mathisonprojects.com")
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
test.CreateDB(t, cfg)
require.NoError(t, ensureLocalDevEula(t.Context(), cfg))
current, err := cfg.GetDBQueries().GetCurrentEulaVersion(t.Context())
require.NoError(t, err)
require.Equal(t, "local-2026-05-07", current.Version)
require.Equal(t, "AArete DoczyAI Terms of Use", current.Title)
agreement, err := cfg.GetDBQueries().GetUserCurrentEulaAgreement(t.Context(), "jacob+2@mathisonprojects.com")
require.NoError(t, err)
require.Equal(t, current.ID, agreement.Eulaversionid)
require.NoError(t, ensureLocalDevEula(t.Context(), cfg))
}
func TestEnsureLocalDevEulaUsesDefaults(t *testing.T) {
if testing.Short() {
t.SkipNow()
}
t.Setenv("DISABLE_AUTH", "true")
t.Setenv(localDevEulaVersionEnv, "local-defaults")
t.Setenv(localDevEulaTitleEnv, "")
t.Setenv(localDevEulaAutoAcceptSubjectEnv, "local-subject")
t.Setenv(localDevEulaAutoAcceptEmailEnv, "")
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
test.CreateDB(t, cfg)
require.NoError(t, ensureLocalDevEula(t.Context(), cfg))
current, err := cfg.GetDBQueries().GetCurrentEulaVersion(t.Context())
require.NoError(t, err)
require.Equal(t, "local-defaults", current.Version)
require.Equal(t, "AArete DoczyAI Terms of Use", current.Title)
agreement, err := cfg.GetDBQueries().GetUserCurrentEulaAgreement(t.Context(), "local-subject")
require.NoError(t, err)
require.Equal(t, "local-subject", agreement.Useremail)
}
func TestEnsureLocalDevEulaSkipsWhenDisabledOrUnconfigured(t *testing.T) {
cfg := &BaseConfig{}
_ = serviceconfig.InitializeConfig(cfg)
t.Setenv("DISABLE_AUTH", "false")
t.Setenv(localDevEulaVersionEnv, "local-disabled")
require.NoError(t, ensureLocalDevEula(t.Context(), cfg))
t.Setenv("DISABLE_AUTH", "true")
t.Setenv(localDevEulaVersionEnv, "")
require.NoError(t, ensureLocalDevEula(t.Context(), cfg))
}