From 305e127e16e5d88fe81c2a2589652b2aed80d1e2 Mon Sep 17 00:00:00 2001 From: Mahesh Date: Wed, 23 Sep 2026 16:50:02 +0530 Subject: [PATCH 1/3] feat(file): read matched files concurrently Keep a file source on one pipeline worker so a glob is not parsed twice, and use task_concurrency to bound parallel fetches of those files. --- go.mod | 2 +- internal/pkg/pipeline/task/file/README.md | 3 +- internal/pkg/pipeline/task/file/file.go | 52 +++++++++++----- internal/pkg/pipeline/task/file/file_test.go | 59 +++++++++++++++++++ .../region-0/batch-0/file-0.txt | 1 + .../region-0/batch-0/file-1.txt | 1 + .../region-0/batch-0/file-2.txt | 1 + .../region-0/batch-0/file-3.txt | 1 + .../region-0/batch-0/file-4.txt | 1 + .../region-0/batch-1/file-0.txt | 1 + .../region-0/batch-1/file-1.txt | 1 + .../region-0/batch-1/file-2.txt | 1 + .../region-0/batch-1/file-3.txt | 1 + .../region-0/batch-1/file-4.txt | 1 + .../region-0/batch-2/file-0.txt | 1 + .../region-0/batch-2/file-1.txt | 1 + .../region-0/batch-2/file-2.txt | 1 + .../region-0/batch-2/file-3.txt | 1 + .../region-0/batch-2/file-4.txt | 1 + .../region-0/batch-3/file-0.txt | 1 + .../region-0/batch-3/file-1.txt | 1 + .../region-0/batch-3/file-2.txt | 1 + .../region-0/batch-3/file-3.txt | 1 + .../region-0/batch-3/file-4.txt | 1 + .../region-1/batch-0/file-0.txt | 1 + .../region-1/batch-0/file-1.txt | 1 + .../region-1/batch-0/file-2.txt | 1 + .../region-1/batch-0/file-3.txt | 1 + .../region-1/batch-0/file-4.txt | 1 + .../region-1/batch-1/file-0.txt | 1 + .../region-1/batch-1/file-1.txt | 1 + .../region-1/batch-1/file-2.txt | 1 + .../region-1/batch-1/file-3.txt | 1 + .../region-1/batch-1/file-4.txt | 1 + .../region-1/batch-2/file-0.txt | 1 + .../region-1/batch-2/file-1.txt | 1 + .../region-1/batch-2/file-2.txt | 1 + .../region-1/batch-2/file-3.txt | 1 + .../region-1/batch-2/file-4.txt | 1 + .../region-1/batch-3/file-0.txt | 1 + .../region-1/batch-3/file-1.txt | 1 + .../region-1/batch-3/file-2.txt | 1 + .../region-1/batch-3/file-3.txt | 1 + .../region-1/batch-3/file-4.txt | 1 + .../region-2/batch-0/file-0.txt | 1 + .../region-2/batch-0/file-1.txt | 1 + .../region-2/batch-0/file-2.txt | 1 + .../region-2/batch-0/file-3.txt | 1 + .../region-2/batch-0/file-4.txt | 1 + .../region-2/batch-1/file-0.txt | 1 + .../region-2/batch-1/file-1.txt | 1 + .../region-2/batch-1/file-2.txt | 1 + .../region-2/batch-1/file-3.txt | 1 + .../region-2/batch-1/file-4.txt | 1 + .../region-2/batch-2/file-0.txt | 1 + .../region-2/batch-2/file-1.txt | 1 + .../region-2/batch-2/file-2.txt | 1 + .../region-2/batch-2/file-3.txt | 1 + .../region-2/batch-2/file-4.txt | 1 + .../region-2/batch-3/file-0.txt | 1 + .../region-2/batch-3/file-1.txt | 1 + .../region-2/batch-3/file-2.txt | 1 + .../region-2/batch-3/file-3.txt | 1 + .../region-2/batch-3/file-4.txt | 1 + .../region-3/batch-0/file-0.txt | 1 + .../region-3/batch-0/file-1.txt | 1 + .../region-3/batch-0/file-2.txt | 1 + .../region-3/batch-0/file-3.txt | 1 + .../region-3/batch-0/file-4.txt | 1 + .../region-3/batch-1/file-0.txt | 1 + .../region-3/batch-1/file-1.txt | 1 + .../region-3/batch-1/file-2.txt | 1 + .../region-3/batch-1/file-3.txt | 1 + .../region-3/batch-1/file-4.txt | 1 + .../region-3/batch-2/file-0.txt | 1 + .../region-3/batch-2/file-1.txt | 1 + .../region-3/batch-2/file-2.txt | 1 + .../region-3/batch-2/file-3.txt | 1 + .../region-3/batch-2/file-4.txt | 1 + .../region-3/batch-3/file-0.txt | 1 + .../region-3/batch-3/file-1.txt | 1 + .../region-3/batch-3/file-2.txt | 1 + .../region-3/batch-3/file-3.txt | 1 + .../region-3/batch-3/file-4.txt | 1 + .../region-4/batch-0/file-0.txt | 1 + .../region-4/batch-0/file-1.txt | 1 + .../region-4/batch-0/file-2.txt | 1 + .../region-4/batch-0/file-3.txt | 1 + .../region-4/batch-0/file-4.txt | 1 + .../region-4/batch-1/file-0.txt | 1 + .../region-4/batch-1/file-1.txt | 1 + .../region-4/batch-1/file-2.txt | 1 + .../region-4/batch-1/file-3.txt | 1 + .../region-4/batch-1/file-4.txt | 1 + .../region-4/batch-2/file-0.txt | 1 + .../region-4/batch-2/file-1.txt | 1 + .../region-4/batch-2/file-2.txt | 1 + .../region-4/batch-2/file-3.txt | 1 + .../region-4/batch-2/file-4.txt | 1 + .../region-4/batch-3/file-0.txt | 1 + .../region-4/batch-3/file-1.txt | 1 + .../region-4/batch-3/file-2.txt | 1 + .../region-4/batch-3/file-3.txt | 1 + .../region-4/batch-3/file-4.txt | 1 + test/pipelines/file_concurrency_test.yaml | 14 ++--- test/pipelines/file_read_failure.yaml | 11 ++++ .../file_read_failure/not-a-file.txt/.gitkeep | 0 test/pipelines/file_read_failure/ok-0.txt | 1 + test/pipelines/file_read_failure/ok-1.txt | 1 + 109 files changed, 217 insertions(+), 26 deletions(-) create mode 100644 internal/pkg/pipeline/task/file/file_test.go create mode 100644 test/pipelines/file_concurrency/region-0/batch-0/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-0/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-0/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-0/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-0/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-1/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-1/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-1/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-1/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-1/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-2/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-2/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-2/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-2/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-2/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-3/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-3/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-3/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-3/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-0/batch-3/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-0/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-0/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-0/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-0/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-0/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-1/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-1/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-1/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-1/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-1/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-2/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-2/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-2/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-2/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-2/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-3/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-3/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-3/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-3/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-1/batch-3/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-0/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-0/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-0/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-0/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-0/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-1/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-1/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-1/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-1/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-1/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-2/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-2/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-2/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-2/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-2/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-3/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-3/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-3/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-3/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-2/batch-3/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-0/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-0/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-0/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-0/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-0/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-1/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-1/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-1/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-1/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-1/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-2/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-2/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-2/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-2/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-2/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-3/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-3/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-3/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-3/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-3/batch-3/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-0/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-0/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-0/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-0/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-0/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-1/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-1/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-1/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-1/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-1/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-2/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-2/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-2/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-2/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-2/file-4.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-3/file-0.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-3/file-1.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-3/file-2.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-3/file-3.txt create mode 100644 test/pipelines/file_concurrency/region-4/batch-3/file-4.txt create mode 100644 test/pipelines/file_read_failure.yaml create mode 100644 test/pipelines/file_read_failure/not-a-file.txt/.gitkeep create mode 100644 test/pipelines/file_read_failure/ok-0.txt create mode 100644 test/pipelines/file_read_failure/ok-1.txt diff --git a/go.mod b/go.mod index 92ecfc0..73e4548 100644 --- a/go.mod +++ b/go.mod @@ -26,6 +26,7 @@ require ( github.com/yamitzky/xlrd-go v0.1.0 golang.org/x/crypto v0.55.0 golang.org/x/net v0.58.0 + golang.org/x/sync v0.22.0 google.golang.org/protobuf v1.36.12 gopkg.in/yaml.v3 v3.0.1 ) @@ -90,7 +91,6 @@ require ( github.com/xuri/efp v0.0.1 // indirect github.com/xuri/nfp v0.0.2-0.20250530014748-2ddeb826f9a9 // indirect golang.org/x/oauth2 v0.36.0 // indirect - golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.41.0 // indirect ) diff --git a/internal/pkg/pipeline/task/file/README.md b/internal/pkg/pipeline/task/file/README.md index d300512..4c21b3c 100644 --- a/internal/pkg/pipeline/task/file/README.md +++ b/internal/pkg/pipeline/task/file/README.md @@ -37,7 +37,7 @@ In read mode, two values are stored in each record's context: | `tags` | map[string]string | - | S3 **write** only: object tags applied on `PutObject`. Ignored for local paths. Values support macros and context templates. See [S3 object tags](#s3-object-tags). | | `success_file` | bool | `false` | Whether to create a success file after writing | | `success_file_name` | string | `_SUCCESS` | Name of the success file | -| `task_concurrency` | int | `1` | Number of competing-consumer workers for this task | +| `task_concurrency` | int | `1` | Write: competing-consumer workers. Read: parallel file fetches (the pipeline still runs a single source worker so files are not duplicated) | | `context` | map | - | JQ expressions whose results are stored on each record for downstream tasks | | `fail_on_error` | bool | `false` | Whether to stop the pipeline if this task encounters an error | @@ -176,6 +176,7 @@ With the inputs above, the writes land at: - `test/pipelines/file.yaml` - Basic file operations - `test/pipelines/context_test.yaml` - File task with context variables +- `test/pipelines/file_concurrency_test.yaml` - Concurrent read of 100 nested files (`task_concurrency: 3`) ## Use Cases diff --git a/internal/pkg/pipeline/task/file/file.go b/internal/pkg/pipeline/task/file/file.go index 19bc5d0..e6e0a0d 100644 --- a/internal/pkg/pipeline/task/file/file.go +++ b/internal/pkg/pipeline/task/file/file.go @@ -15,6 +15,7 @@ import ( "github.com/patterninc/caterpillar/internal/pkg/pipeline/record" "github.com/patterninc/caterpillar/internal/pkg/pipeline/task" "github.com/patterninc/caterpillar/internal/pkg/textutil" + "golang.org/x/sync/errgroup" ) const ( @@ -131,30 +132,49 @@ func (f *file) readFile(output chan<- *record.Record) error { return err } - for _, path := range paths { + // if no paths are found, return nil + if len(paths) == 0 { + return nil + } - readerCloser, err := reader.read(path) - if err != nil { - return err - } - defer readerCloser.Close() + // set the error group limit to the minimum of the task concurrency and the number of paths + g, gctx := errgroup.WithContext(context.Background()) + g.SetLimit(min(f.GetTaskConcurrency(), len(paths))) - content, err := io.ReadAll(readerCloser) - if err != nil { - return err + // iterate over the paths and emit the file + for _, path := range paths { + if gctx.Err() != nil { + break } + g.Go(func() error { + return f.emitFile(reader, path, output) + }) + } - // Create a default record with context - fileName := textutil.SlugifyFileName(filepath.Base(path)) - rc := &record.Record{Context: ctx} - rc.SetContextValue(string(task.CtxKeyFileNameWrite), fileName) - rc.SetContextValue(string(task.CtxKeyFilePathWrite), textutil.SlugifyFilePath(path)) + return g.Wait() - // let's write content to output channel - f.SendData(rc.Context, content, output) +} + +func (f *file) emitFile(r reader, path string, output chan<- *record.Record) error { + + readerCloser, err := r.read(path) + if err != nil { + return err + } + defer readerCloser.Close() + content, err := io.ReadAll(readerCloser) + if err != nil { + return err } + fileName := textutil.SlugifyFileName(filepath.Base(path)) + rc := &record.Record{Context: ctx} + rc.SetContextValue(string(task.CtxKeyFileNameWrite), fileName) + rc.SetContextValue(string(task.CtxKeyFilePathWrite), textutil.SlugifyFilePath(path)) + + f.SendData(rc.Context, content, output) + return nil } diff --git a/internal/pkg/pipeline/task/file/file_test.go b/internal/pkg/pipeline/task/file/file_test.go new file mode 100644 index 0000000..78885d1 --- /dev/null +++ b/internal/pkg/pipeline/task/file/file_test.go @@ -0,0 +1,59 @@ +package file + +import ( + "fmt" + "os" + "path/filepath" + "testing" + + "github.com/patterninc/caterpillar/internal/pkg/config" + "github.com/patterninc/caterpillar/internal/pkg/pipeline/record" + "github.com/patterninc/caterpillar/internal/pkg/pipeline/task" +) + +func TestReadFileConcurrent(t *testing.T) { + + dir := t.TempDir() + want := make(map[string]struct{}, 10) + for i := 0; i < 10; i++ { + content := fmt.Sprintf("content-%d", i) + if err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("f%d.txt", i)), []byte(content), 0o644); err != nil { + t.Fatal(err) + } + want[content] = struct{}{} + } + + f := &file{ + Base: task.Base{ + Name: "read", + Type: "file", + TaskConcurrency: 4, + }, + Path: config.String(filepath.Join(dir, "*.txt")), + } + + out := make(chan *record.Record, 20) + errCh := make(chan error, 1) + go func() { + errCh <- f.readFile(out) + close(out) + }() + + got := make(map[string]struct{}) + for r := range out { + got[string(r.Data)] = struct{}{} + } + if err := <-errCh; err != nil { + t.Fatalf("readFile: %v", err) + } + + if len(got) != len(want) { + t.Fatalf("got %d records, want %d", len(got), len(want)) + } + for content := range want { + if _, ok := got[content]; !ok { + t.Errorf("missing content %q", content) + } + } + +} diff --git a/test/pipelines/file_concurrency/region-0/batch-0/file-0.txt b/test/pipelines/file_concurrency/region-0/batch-0/file-0.txt new file mode 100644 index 0000000..47b81db --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-0/file-0.txt @@ -0,0 +1 @@ +region-0/batch-0/file-0 diff --git a/test/pipelines/file_concurrency/region-0/batch-0/file-1.txt b/test/pipelines/file_concurrency/region-0/batch-0/file-1.txt new file mode 100644 index 0000000..8891260 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-0/file-1.txt @@ -0,0 +1 @@ +region-0/batch-0/file-1 diff --git a/test/pipelines/file_concurrency/region-0/batch-0/file-2.txt b/test/pipelines/file_concurrency/region-0/batch-0/file-2.txt new file mode 100644 index 0000000..4e604db --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-0/file-2.txt @@ -0,0 +1 @@ +region-0/batch-0/file-2 diff --git a/test/pipelines/file_concurrency/region-0/batch-0/file-3.txt b/test/pipelines/file_concurrency/region-0/batch-0/file-3.txt new file mode 100644 index 0000000..c9e3921 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-0/file-3.txt @@ -0,0 +1 @@ +region-0/batch-0/file-3 diff --git a/test/pipelines/file_concurrency/region-0/batch-0/file-4.txt b/test/pipelines/file_concurrency/region-0/batch-0/file-4.txt new file mode 100644 index 0000000..3fd0d91 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-0/file-4.txt @@ -0,0 +1 @@ +region-0/batch-0/file-4 diff --git a/test/pipelines/file_concurrency/region-0/batch-1/file-0.txt b/test/pipelines/file_concurrency/region-0/batch-1/file-0.txt new file mode 100644 index 0000000..2c28845 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-1/file-0.txt @@ -0,0 +1 @@ +region-0/batch-1/file-0 diff --git a/test/pipelines/file_concurrency/region-0/batch-1/file-1.txt b/test/pipelines/file_concurrency/region-0/batch-1/file-1.txt new file mode 100644 index 0000000..f710d31 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-1/file-1.txt @@ -0,0 +1 @@ +region-0/batch-1/file-1 diff --git a/test/pipelines/file_concurrency/region-0/batch-1/file-2.txt b/test/pipelines/file_concurrency/region-0/batch-1/file-2.txt new file mode 100644 index 0000000..8bc5802 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-1/file-2.txt @@ -0,0 +1 @@ +region-0/batch-1/file-2 diff --git a/test/pipelines/file_concurrency/region-0/batch-1/file-3.txt b/test/pipelines/file_concurrency/region-0/batch-1/file-3.txt new file mode 100644 index 0000000..83b4063 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-1/file-3.txt @@ -0,0 +1 @@ +region-0/batch-1/file-3 diff --git a/test/pipelines/file_concurrency/region-0/batch-1/file-4.txt b/test/pipelines/file_concurrency/region-0/batch-1/file-4.txt new file mode 100644 index 0000000..ce6ee01 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-1/file-4.txt @@ -0,0 +1 @@ +region-0/batch-1/file-4 diff --git a/test/pipelines/file_concurrency/region-0/batch-2/file-0.txt b/test/pipelines/file_concurrency/region-0/batch-2/file-0.txt new file mode 100644 index 0000000..2bc5e50 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-2/file-0.txt @@ -0,0 +1 @@ +region-0/batch-2/file-0 diff --git a/test/pipelines/file_concurrency/region-0/batch-2/file-1.txt b/test/pipelines/file_concurrency/region-0/batch-2/file-1.txt new file mode 100644 index 0000000..05bdf86 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-2/file-1.txt @@ -0,0 +1 @@ +region-0/batch-2/file-1 diff --git a/test/pipelines/file_concurrency/region-0/batch-2/file-2.txt b/test/pipelines/file_concurrency/region-0/batch-2/file-2.txt new file mode 100644 index 0000000..fea1b1e --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-2/file-2.txt @@ -0,0 +1 @@ +region-0/batch-2/file-2 diff --git a/test/pipelines/file_concurrency/region-0/batch-2/file-3.txt b/test/pipelines/file_concurrency/region-0/batch-2/file-3.txt new file mode 100644 index 0000000..c0804e9 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-2/file-3.txt @@ -0,0 +1 @@ +region-0/batch-2/file-3 diff --git a/test/pipelines/file_concurrency/region-0/batch-2/file-4.txt b/test/pipelines/file_concurrency/region-0/batch-2/file-4.txt new file mode 100644 index 0000000..0082793 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-2/file-4.txt @@ -0,0 +1 @@ +region-0/batch-2/file-4 diff --git a/test/pipelines/file_concurrency/region-0/batch-3/file-0.txt b/test/pipelines/file_concurrency/region-0/batch-3/file-0.txt new file mode 100644 index 0000000..b698621 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-3/file-0.txt @@ -0,0 +1 @@ +region-0/batch-3/file-0 diff --git a/test/pipelines/file_concurrency/region-0/batch-3/file-1.txt b/test/pipelines/file_concurrency/region-0/batch-3/file-1.txt new file mode 100644 index 0000000..80a3e53 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-3/file-1.txt @@ -0,0 +1 @@ +region-0/batch-3/file-1 diff --git a/test/pipelines/file_concurrency/region-0/batch-3/file-2.txt b/test/pipelines/file_concurrency/region-0/batch-3/file-2.txt new file mode 100644 index 0000000..0ced86f --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-3/file-2.txt @@ -0,0 +1 @@ +region-0/batch-3/file-2 diff --git a/test/pipelines/file_concurrency/region-0/batch-3/file-3.txt b/test/pipelines/file_concurrency/region-0/batch-3/file-3.txt new file mode 100644 index 0000000..b044938 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-3/file-3.txt @@ -0,0 +1 @@ +region-0/batch-3/file-3 diff --git a/test/pipelines/file_concurrency/region-0/batch-3/file-4.txt b/test/pipelines/file_concurrency/region-0/batch-3/file-4.txt new file mode 100644 index 0000000..1063f63 --- /dev/null +++ b/test/pipelines/file_concurrency/region-0/batch-3/file-4.txt @@ -0,0 +1 @@ +region-0/batch-3/file-4 diff --git a/test/pipelines/file_concurrency/region-1/batch-0/file-0.txt b/test/pipelines/file_concurrency/region-1/batch-0/file-0.txt new file mode 100644 index 0000000..89ea31c --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-0/file-0.txt @@ -0,0 +1 @@ +region-1/batch-0/file-0 diff --git a/test/pipelines/file_concurrency/region-1/batch-0/file-1.txt b/test/pipelines/file_concurrency/region-1/batch-0/file-1.txt new file mode 100644 index 0000000..aec9522 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-0/file-1.txt @@ -0,0 +1 @@ +region-1/batch-0/file-1 diff --git a/test/pipelines/file_concurrency/region-1/batch-0/file-2.txt b/test/pipelines/file_concurrency/region-1/batch-0/file-2.txt new file mode 100644 index 0000000..757f972 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-0/file-2.txt @@ -0,0 +1 @@ +region-1/batch-0/file-2 diff --git a/test/pipelines/file_concurrency/region-1/batch-0/file-3.txt b/test/pipelines/file_concurrency/region-1/batch-0/file-3.txt new file mode 100644 index 0000000..f055813 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-0/file-3.txt @@ -0,0 +1 @@ +region-1/batch-0/file-3 diff --git a/test/pipelines/file_concurrency/region-1/batch-0/file-4.txt b/test/pipelines/file_concurrency/region-1/batch-0/file-4.txt new file mode 100644 index 0000000..c8cf6e7 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-0/file-4.txt @@ -0,0 +1 @@ +region-1/batch-0/file-4 diff --git a/test/pipelines/file_concurrency/region-1/batch-1/file-0.txt b/test/pipelines/file_concurrency/region-1/batch-1/file-0.txt new file mode 100644 index 0000000..1bd5dd4 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-1/file-0.txt @@ -0,0 +1 @@ +region-1/batch-1/file-0 diff --git a/test/pipelines/file_concurrency/region-1/batch-1/file-1.txt b/test/pipelines/file_concurrency/region-1/batch-1/file-1.txt new file mode 100644 index 0000000..24215a6 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-1/file-1.txt @@ -0,0 +1 @@ +region-1/batch-1/file-1 diff --git a/test/pipelines/file_concurrency/region-1/batch-1/file-2.txt b/test/pipelines/file_concurrency/region-1/batch-1/file-2.txt new file mode 100644 index 0000000..f4dde6c --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-1/file-2.txt @@ -0,0 +1 @@ +region-1/batch-1/file-2 diff --git a/test/pipelines/file_concurrency/region-1/batch-1/file-3.txt b/test/pipelines/file_concurrency/region-1/batch-1/file-3.txt new file mode 100644 index 0000000..7628eb0 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-1/file-3.txt @@ -0,0 +1 @@ +region-1/batch-1/file-3 diff --git a/test/pipelines/file_concurrency/region-1/batch-1/file-4.txt b/test/pipelines/file_concurrency/region-1/batch-1/file-4.txt new file mode 100644 index 0000000..9877696 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-1/file-4.txt @@ -0,0 +1 @@ +region-1/batch-1/file-4 diff --git a/test/pipelines/file_concurrency/region-1/batch-2/file-0.txt b/test/pipelines/file_concurrency/region-1/batch-2/file-0.txt new file mode 100644 index 0000000..428980a --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-2/file-0.txt @@ -0,0 +1 @@ +region-1/batch-2/file-0 diff --git a/test/pipelines/file_concurrency/region-1/batch-2/file-1.txt b/test/pipelines/file_concurrency/region-1/batch-2/file-1.txt new file mode 100644 index 0000000..ee5e81b --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-2/file-1.txt @@ -0,0 +1 @@ +region-1/batch-2/file-1 diff --git a/test/pipelines/file_concurrency/region-1/batch-2/file-2.txt b/test/pipelines/file_concurrency/region-1/batch-2/file-2.txt new file mode 100644 index 0000000..3936c41 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-2/file-2.txt @@ -0,0 +1 @@ +region-1/batch-2/file-2 diff --git a/test/pipelines/file_concurrency/region-1/batch-2/file-3.txt b/test/pipelines/file_concurrency/region-1/batch-2/file-3.txt new file mode 100644 index 0000000..e690961 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-2/file-3.txt @@ -0,0 +1 @@ +region-1/batch-2/file-3 diff --git a/test/pipelines/file_concurrency/region-1/batch-2/file-4.txt b/test/pipelines/file_concurrency/region-1/batch-2/file-4.txt new file mode 100644 index 0000000..f3bd131 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-2/file-4.txt @@ -0,0 +1 @@ +region-1/batch-2/file-4 diff --git a/test/pipelines/file_concurrency/region-1/batch-3/file-0.txt b/test/pipelines/file_concurrency/region-1/batch-3/file-0.txt new file mode 100644 index 0000000..1d7ab9d --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-3/file-0.txt @@ -0,0 +1 @@ +region-1/batch-3/file-0 diff --git a/test/pipelines/file_concurrency/region-1/batch-3/file-1.txt b/test/pipelines/file_concurrency/region-1/batch-3/file-1.txt new file mode 100644 index 0000000..819d067 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-3/file-1.txt @@ -0,0 +1 @@ +region-1/batch-3/file-1 diff --git a/test/pipelines/file_concurrency/region-1/batch-3/file-2.txt b/test/pipelines/file_concurrency/region-1/batch-3/file-2.txt new file mode 100644 index 0000000..2fd157f --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-3/file-2.txt @@ -0,0 +1 @@ +region-1/batch-3/file-2 diff --git a/test/pipelines/file_concurrency/region-1/batch-3/file-3.txt b/test/pipelines/file_concurrency/region-1/batch-3/file-3.txt new file mode 100644 index 0000000..eb10005 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-3/file-3.txt @@ -0,0 +1 @@ +region-1/batch-3/file-3 diff --git a/test/pipelines/file_concurrency/region-1/batch-3/file-4.txt b/test/pipelines/file_concurrency/region-1/batch-3/file-4.txt new file mode 100644 index 0000000..6ee64f0 --- /dev/null +++ b/test/pipelines/file_concurrency/region-1/batch-3/file-4.txt @@ -0,0 +1 @@ +region-1/batch-3/file-4 diff --git a/test/pipelines/file_concurrency/region-2/batch-0/file-0.txt b/test/pipelines/file_concurrency/region-2/batch-0/file-0.txt new file mode 100644 index 0000000..578f28e --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-0/file-0.txt @@ -0,0 +1 @@ +region-2/batch-0/file-0 diff --git a/test/pipelines/file_concurrency/region-2/batch-0/file-1.txt b/test/pipelines/file_concurrency/region-2/batch-0/file-1.txt new file mode 100644 index 0000000..3917253 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-0/file-1.txt @@ -0,0 +1 @@ +region-2/batch-0/file-1 diff --git a/test/pipelines/file_concurrency/region-2/batch-0/file-2.txt b/test/pipelines/file_concurrency/region-2/batch-0/file-2.txt new file mode 100644 index 0000000..3fe9f2d --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-0/file-2.txt @@ -0,0 +1 @@ +region-2/batch-0/file-2 diff --git a/test/pipelines/file_concurrency/region-2/batch-0/file-3.txt b/test/pipelines/file_concurrency/region-2/batch-0/file-3.txt new file mode 100644 index 0000000..2d9cf4b --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-0/file-3.txt @@ -0,0 +1 @@ +region-2/batch-0/file-3 diff --git a/test/pipelines/file_concurrency/region-2/batch-0/file-4.txt b/test/pipelines/file_concurrency/region-2/batch-0/file-4.txt new file mode 100644 index 0000000..7c03720 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-0/file-4.txt @@ -0,0 +1 @@ +region-2/batch-0/file-4 diff --git a/test/pipelines/file_concurrency/region-2/batch-1/file-0.txt b/test/pipelines/file_concurrency/region-2/batch-1/file-0.txt new file mode 100644 index 0000000..6efa518 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-1/file-0.txt @@ -0,0 +1 @@ +region-2/batch-1/file-0 diff --git a/test/pipelines/file_concurrency/region-2/batch-1/file-1.txt b/test/pipelines/file_concurrency/region-2/batch-1/file-1.txt new file mode 100644 index 0000000..c516d32 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-1/file-1.txt @@ -0,0 +1 @@ +region-2/batch-1/file-1 diff --git a/test/pipelines/file_concurrency/region-2/batch-1/file-2.txt b/test/pipelines/file_concurrency/region-2/batch-1/file-2.txt new file mode 100644 index 0000000..4a4d906 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-1/file-2.txt @@ -0,0 +1 @@ +region-2/batch-1/file-2 diff --git a/test/pipelines/file_concurrency/region-2/batch-1/file-3.txt b/test/pipelines/file_concurrency/region-2/batch-1/file-3.txt new file mode 100644 index 0000000..a43bec0 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-1/file-3.txt @@ -0,0 +1 @@ +region-2/batch-1/file-3 diff --git a/test/pipelines/file_concurrency/region-2/batch-1/file-4.txt b/test/pipelines/file_concurrency/region-2/batch-1/file-4.txt new file mode 100644 index 0000000..b9fc44a --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-1/file-4.txt @@ -0,0 +1 @@ +region-2/batch-1/file-4 diff --git a/test/pipelines/file_concurrency/region-2/batch-2/file-0.txt b/test/pipelines/file_concurrency/region-2/batch-2/file-0.txt new file mode 100644 index 0000000..57a28d4 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-2/file-0.txt @@ -0,0 +1 @@ +region-2/batch-2/file-0 diff --git a/test/pipelines/file_concurrency/region-2/batch-2/file-1.txt b/test/pipelines/file_concurrency/region-2/batch-2/file-1.txt new file mode 100644 index 0000000..7f1e58d --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-2/file-1.txt @@ -0,0 +1 @@ +region-2/batch-2/file-1 diff --git a/test/pipelines/file_concurrency/region-2/batch-2/file-2.txt b/test/pipelines/file_concurrency/region-2/batch-2/file-2.txt new file mode 100644 index 0000000..7e95314 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-2/file-2.txt @@ -0,0 +1 @@ +region-2/batch-2/file-2 diff --git a/test/pipelines/file_concurrency/region-2/batch-2/file-3.txt b/test/pipelines/file_concurrency/region-2/batch-2/file-3.txt new file mode 100644 index 0000000..259be4f --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-2/file-3.txt @@ -0,0 +1 @@ +region-2/batch-2/file-3 diff --git a/test/pipelines/file_concurrency/region-2/batch-2/file-4.txt b/test/pipelines/file_concurrency/region-2/batch-2/file-4.txt new file mode 100644 index 0000000..d21182f --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-2/file-4.txt @@ -0,0 +1 @@ +region-2/batch-2/file-4 diff --git a/test/pipelines/file_concurrency/region-2/batch-3/file-0.txt b/test/pipelines/file_concurrency/region-2/batch-3/file-0.txt new file mode 100644 index 0000000..d29f21e --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-3/file-0.txt @@ -0,0 +1 @@ +region-2/batch-3/file-0 diff --git a/test/pipelines/file_concurrency/region-2/batch-3/file-1.txt b/test/pipelines/file_concurrency/region-2/batch-3/file-1.txt new file mode 100644 index 0000000..728bc1c --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-3/file-1.txt @@ -0,0 +1 @@ +region-2/batch-3/file-1 diff --git a/test/pipelines/file_concurrency/region-2/batch-3/file-2.txt b/test/pipelines/file_concurrency/region-2/batch-3/file-2.txt new file mode 100644 index 0000000..d5565a9 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-3/file-2.txt @@ -0,0 +1 @@ +region-2/batch-3/file-2 diff --git a/test/pipelines/file_concurrency/region-2/batch-3/file-3.txt b/test/pipelines/file_concurrency/region-2/batch-3/file-3.txt new file mode 100644 index 0000000..876c63f --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-3/file-3.txt @@ -0,0 +1 @@ +region-2/batch-3/file-3 diff --git a/test/pipelines/file_concurrency/region-2/batch-3/file-4.txt b/test/pipelines/file_concurrency/region-2/batch-3/file-4.txt new file mode 100644 index 0000000..ea3f625 --- /dev/null +++ b/test/pipelines/file_concurrency/region-2/batch-3/file-4.txt @@ -0,0 +1 @@ +region-2/batch-3/file-4 diff --git a/test/pipelines/file_concurrency/region-3/batch-0/file-0.txt b/test/pipelines/file_concurrency/region-3/batch-0/file-0.txt new file mode 100644 index 0000000..c62351e --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-0/file-0.txt @@ -0,0 +1 @@ +region-3/batch-0/file-0 diff --git a/test/pipelines/file_concurrency/region-3/batch-0/file-1.txt b/test/pipelines/file_concurrency/region-3/batch-0/file-1.txt new file mode 100644 index 0000000..9a59add --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-0/file-1.txt @@ -0,0 +1 @@ +region-3/batch-0/file-1 diff --git a/test/pipelines/file_concurrency/region-3/batch-0/file-2.txt b/test/pipelines/file_concurrency/region-3/batch-0/file-2.txt new file mode 100644 index 0000000..7a06294 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-0/file-2.txt @@ -0,0 +1 @@ +region-3/batch-0/file-2 diff --git a/test/pipelines/file_concurrency/region-3/batch-0/file-3.txt b/test/pipelines/file_concurrency/region-3/batch-0/file-3.txt new file mode 100644 index 0000000..f45fc22 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-0/file-3.txt @@ -0,0 +1 @@ +region-3/batch-0/file-3 diff --git a/test/pipelines/file_concurrency/region-3/batch-0/file-4.txt b/test/pipelines/file_concurrency/region-3/batch-0/file-4.txt new file mode 100644 index 0000000..c29dd75 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-0/file-4.txt @@ -0,0 +1 @@ +region-3/batch-0/file-4 diff --git a/test/pipelines/file_concurrency/region-3/batch-1/file-0.txt b/test/pipelines/file_concurrency/region-3/batch-1/file-0.txt new file mode 100644 index 0000000..e1a7f50 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-1/file-0.txt @@ -0,0 +1 @@ +region-3/batch-1/file-0 diff --git a/test/pipelines/file_concurrency/region-3/batch-1/file-1.txt b/test/pipelines/file_concurrency/region-3/batch-1/file-1.txt new file mode 100644 index 0000000..8319692 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-1/file-1.txt @@ -0,0 +1 @@ +region-3/batch-1/file-1 diff --git a/test/pipelines/file_concurrency/region-3/batch-1/file-2.txt b/test/pipelines/file_concurrency/region-3/batch-1/file-2.txt new file mode 100644 index 0000000..2043aff --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-1/file-2.txt @@ -0,0 +1 @@ +region-3/batch-1/file-2 diff --git a/test/pipelines/file_concurrency/region-3/batch-1/file-3.txt b/test/pipelines/file_concurrency/region-3/batch-1/file-3.txt new file mode 100644 index 0000000..5e07401 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-1/file-3.txt @@ -0,0 +1 @@ +region-3/batch-1/file-3 diff --git a/test/pipelines/file_concurrency/region-3/batch-1/file-4.txt b/test/pipelines/file_concurrency/region-3/batch-1/file-4.txt new file mode 100644 index 0000000..d72db21 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-1/file-4.txt @@ -0,0 +1 @@ +region-3/batch-1/file-4 diff --git a/test/pipelines/file_concurrency/region-3/batch-2/file-0.txt b/test/pipelines/file_concurrency/region-3/batch-2/file-0.txt new file mode 100644 index 0000000..c6066ac --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-2/file-0.txt @@ -0,0 +1 @@ +region-3/batch-2/file-0 diff --git a/test/pipelines/file_concurrency/region-3/batch-2/file-1.txt b/test/pipelines/file_concurrency/region-3/batch-2/file-1.txt new file mode 100644 index 0000000..9f2426a --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-2/file-1.txt @@ -0,0 +1 @@ +region-3/batch-2/file-1 diff --git a/test/pipelines/file_concurrency/region-3/batch-2/file-2.txt b/test/pipelines/file_concurrency/region-3/batch-2/file-2.txt new file mode 100644 index 0000000..d22fc27 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-2/file-2.txt @@ -0,0 +1 @@ +region-3/batch-2/file-2 diff --git a/test/pipelines/file_concurrency/region-3/batch-2/file-3.txt b/test/pipelines/file_concurrency/region-3/batch-2/file-3.txt new file mode 100644 index 0000000..995a25d --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-2/file-3.txt @@ -0,0 +1 @@ +region-3/batch-2/file-3 diff --git a/test/pipelines/file_concurrency/region-3/batch-2/file-4.txt b/test/pipelines/file_concurrency/region-3/batch-2/file-4.txt new file mode 100644 index 0000000..271e476 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-2/file-4.txt @@ -0,0 +1 @@ +region-3/batch-2/file-4 diff --git a/test/pipelines/file_concurrency/region-3/batch-3/file-0.txt b/test/pipelines/file_concurrency/region-3/batch-3/file-0.txt new file mode 100644 index 0000000..f01d53c --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-3/file-0.txt @@ -0,0 +1 @@ +region-3/batch-3/file-0 diff --git a/test/pipelines/file_concurrency/region-3/batch-3/file-1.txt b/test/pipelines/file_concurrency/region-3/batch-3/file-1.txt new file mode 100644 index 0000000..f3bba7a --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-3/file-1.txt @@ -0,0 +1 @@ +region-3/batch-3/file-1 diff --git a/test/pipelines/file_concurrency/region-3/batch-3/file-2.txt b/test/pipelines/file_concurrency/region-3/batch-3/file-2.txt new file mode 100644 index 0000000..415a10c --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-3/file-2.txt @@ -0,0 +1 @@ +region-3/batch-3/file-2 diff --git a/test/pipelines/file_concurrency/region-3/batch-3/file-3.txt b/test/pipelines/file_concurrency/region-3/batch-3/file-3.txt new file mode 100644 index 0000000..1df66b6 --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-3/file-3.txt @@ -0,0 +1 @@ +region-3/batch-3/file-3 diff --git a/test/pipelines/file_concurrency/region-3/batch-3/file-4.txt b/test/pipelines/file_concurrency/region-3/batch-3/file-4.txt new file mode 100644 index 0000000..ea10fee --- /dev/null +++ b/test/pipelines/file_concurrency/region-3/batch-3/file-4.txt @@ -0,0 +1 @@ +region-3/batch-3/file-4 diff --git a/test/pipelines/file_concurrency/region-4/batch-0/file-0.txt b/test/pipelines/file_concurrency/region-4/batch-0/file-0.txt new file mode 100644 index 0000000..d649906 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-0/file-0.txt @@ -0,0 +1 @@ +region-4/batch-0/file-0 diff --git a/test/pipelines/file_concurrency/region-4/batch-0/file-1.txt b/test/pipelines/file_concurrency/region-4/batch-0/file-1.txt new file mode 100644 index 0000000..f1c51ee --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-0/file-1.txt @@ -0,0 +1 @@ +region-4/batch-0/file-1 diff --git a/test/pipelines/file_concurrency/region-4/batch-0/file-2.txt b/test/pipelines/file_concurrency/region-4/batch-0/file-2.txt new file mode 100644 index 0000000..99a812e --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-0/file-2.txt @@ -0,0 +1 @@ +region-4/batch-0/file-2 diff --git a/test/pipelines/file_concurrency/region-4/batch-0/file-3.txt b/test/pipelines/file_concurrency/region-4/batch-0/file-3.txt new file mode 100644 index 0000000..64563f6 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-0/file-3.txt @@ -0,0 +1 @@ +region-4/batch-0/file-3 diff --git a/test/pipelines/file_concurrency/region-4/batch-0/file-4.txt b/test/pipelines/file_concurrency/region-4/batch-0/file-4.txt new file mode 100644 index 0000000..2e5554f --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-0/file-4.txt @@ -0,0 +1 @@ +region-4/batch-0/file-4 diff --git a/test/pipelines/file_concurrency/region-4/batch-1/file-0.txt b/test/pipelines/file_concurrency/region-4/batch-1/file-0.txt new file mode 100644 index 0000000..c1c6c86 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-1/file-0.txt @@ -0,0 +1 @@ +region-4/batch-1/file-0 diff --git a/test/pipelines/file_concurrency/region-4/batch-1/file-1.txt b/test/pipelines/file_concurrency/region-4/batch-1/file-1.txt new file mode 100644 index 0000000..34ab405 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-1/file-1.txt @@ -0,0 +1 @@ +region-4/batch-1/file-1 diff --git a/test/pipelines/file_concurrency/region-4/batch-1/file-2.txt b/test/pipelines/file_concurrency/region-4/batch-1/file-2.txt new file mode 100644 index 0000000..22ae037 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-1/file-2.txt @@ -0,0 +1 @@ +region-4/batch-1/file-2 diff --git a/test/pipelines/file_concurrency/region-4/batch-1/file-3.txt b/test/pipelines/file_concurrency/region-4/batch-1/file-3.txt new file mode 100644 index 0000000..4459a4d --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-1/file-3.txt @@ -0,0 +1 @@ +region-4/batch-1/file-3 diff --git a/test/pipelines/file_concurrency/region-4/batch-1/file-4.txt b/test/pipelines/file_concurrency/region-4/batch-1/file-4.txt new file mode 100644 index 0000000..55b908a --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-1/file-4.txt @@ -0,0 +1 @@ +region-4/batch-1/file-4 diff --git a/test/pipelines/file_concurrency/region-4/batch-2/file-0.txt b/test/pipelines/file_concurrency/region-4/batch-2/file-0.txt new file mode 100644 index 0000000..3237176 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-2/file-0.txt @@ -0,0 +1 @@ +region-4/batch-2/file-0 diff --git a/test/pipelines/file_concurrency/region-4/batch-2/file-1.txt b/test/pipelines/file_concurrency/region-4/batch-2/file-1.txt new file mode 100644 index 0000000..ee61d44 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-2/file-1.txt @@ -0,0 +1 @@ +region-4/batch-2/file-1 diff --git a/test/pipelines/file_concurrency/region-4/batch-2/file-2.txt b/test/pipelines/file_concurrency/region-4/batch-2/file-2.txt new file mode 100644 index 0000000..2454cb2 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-2/file-2.txt @@ -0,0 +1 @@ +region-4/batch-2/file-2 diff --git a/test/pipelines/file_concurrency/region-4/batch-2/file-3.txt b/test/pipelines/file_concurrency/region-4/batch-2/file-3.txt new file mode 100644 index 0000000..951d48b --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-2/file-3.txt @@ -0,0 +1 @@ +region-4/batch-2/file-3 diff --git a/test/pipelines/file_concurrency/region-4/batch-2/file-4.txt b/test/pipelines/file_concurrency/region-4/batch-2/file-4.txt new file mode 100644 index 0000000..f233da9 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-2/file-4.txt @@ -0,0 +1 @@ +region-4/batch-2/file-4 diff --git a/test/pipelines/file_concurrency/region-4/batch-3/file-0.txt b/test/pipelines/file_concurrency/region-4/batch-3/file-0.txt new file mode 100644 index 0000000..47e9b4b --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-3/file-0.txt @@ -0,0 +1 @@ +region-4/batch-3/file-0 diff --git a/test/pipelines/file_concurrency/region-4/batch-3/file-1.txt b/test/pipelines/file_concurrency/region-4/batch-3/file-1.txt new file mode 100644 index 0000000..a7bc08b --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-3/file-1.txt @@ -0,0 +1 @@ +region-4/batch-3/file-1 diff --git a/test/pipelines/file_concurrency/region-4/batch-3/file-2.txt b/test/pipelines/file_concurrency/region-4/batch-3/file-2.txt new file mode 100644 index 0000000..45cef44 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-3/file-2.txt @@ -0,0 +1 @@ +region-4/batch-3/file-2 diff --git a/test/pipelines/file_concurrency/region-4/batch-3/file-3.txt b/test/pipelines/file_concurrency/region-4/batch-3/file-3.txt new file mode 100644 index 0000000..2c7b888 --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-3/file-3.txt @@ -0,0 +1 @@ +region-4/batch-3/file-3 diff --git a/test/pipelines/file_concurrency/region-4/batch-3/file-4.txt b/test/pipelines/file_concurrency/region-4/batch-3/file-4.txt new file mode 100644 index 0000000..5a1741e --- /dev/null +++ b/test/pipelines/file_concurrency/region-4/batch-3/file-4.txt @@ -0,0 +1 @@ +region-4/batch-3/file-4 diff --git a/test/pipelines/file_concurrency_test.yaml b/test/pipelines/file_concurrency_test.yaml index b61e3d8..3381622 100644 --- a/test/pipelines/file_concurrency_test.yaml +++ b/test/pipelines/file_concurrency_test.yaml @@ -1,11 +1,9 @@ tasks: - - name: read_glob_concurrently + - name: read_nested_concurrently type: file - path: test/pipelines/*.txt - task_concurrency: 4 - - name: split_to_lines - type: split - - name: write_concurrently + path: test/pipelines/file_concurrency/**/*.txt + task_concurrency: 3 + fail_on_error: true + - name: write_to_local type: file - path: test/pipelines/output/{{ macro "uuid" }}.txt - task_concurrency: 8 + path: ./output/{{ context "CATERPILLAR_FILE_NAME_WRITE" }} diff --git a/test/pipelines/file_read_failure.yaml b/test/pipelines/file_read_failure.yaml new file mode 100644 index 0000000..1748b56 --- /dev/null +++ b/test/pipelines/file_read_failure.yaml @@ -0,0 +1,11 @@ +# Glob matches two readable files and one directory named like a file. +# Opening that directory fails the read, and fail_on_error fails the run. +tasks: + - name: read_with_failure + type: file + path: test/pipelines/file_read_failure/* + task_concurrency: 2 + fail_on_error: true + - name: echo + type: echo + only_data: true diff --git a/test/pipelines/file_read_failure/not-a-file.txt/.gitkeep b/test/pipelines/file_read_failure/not-a-file.txt/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/test/pipelines/file_read_failure/ok-0.txt b/test/pipelines/file_read_failure/ok-0.txt new file mode 100644 index 0000000..6c8be44 --- /dev/null +++ b/test/pipelines/file_read_failure/ok-0.txt @@ -0,0 +1 @@ +ok-0 diff --git a/test/pipelines/file_read_failure/ok-1.txt b/test/pipelines/file_read_failure/ok-1.txt new file mode 100644 index 0000000..a4b08e7 --- /dev/null +++ b/test/pipelines/file_read_failure/ok-1.txt @@ -0,0 +1 @@ +ok-1 From 38878743d83ce696ccf6144f8323331e2ec395f9 Mon Sep 17 00:00:00 2001 From: Mahesh Date: Wed, 23 Sep 2026 16:55:26 +0530 Subject: [PATCH 2/3] test(file): use uuid macro for output path in concurrency test --- test/pipelines/file_concurrency_test.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/pipelines/file_concurrency_test.yaml b/test/pipelines/file_concurrency_test.yaml index 3381622..884ab79 100644 --- a/test/pipelines/file_concurrency_test.yaml +++ b/test/pipelines/file_concurrency_test.yaml @@ -6,4 +6,4 @@ tasks: fail_on_error: true - name: write_to_local type: file - path: ./output/{{ context "CATERPILLAR_FILE_NAME_WRITE" }} + path: ./output/{{ macro "uuid" }}.txt From c9cb49055a0c9a6ee82e84ae3e4322ae90fcc1d7 Mon Sep 17 00:00:00 2001 From: Mahesh Date: Wed, 30 Sep 2026 14:20:21 +0530 Subject: [PATCH 3/3] docs(file): trim task_concurrency note and drop test file Address PR #128 review: shorten the read-mode task_concurrency description and remove file_test.go per the repo's convention of keeping test pipelines rather than Go test files. Co-Authored-By: Claude Opus 4.8 (1M context) --- internal/pkg/pipeline/task/file/README.md | 2 +- internal/pkg/pipeline/task/file/file_test.go | 59 -------------------- 2 files changed, 1 insertion(+), 60 deletions(-) delete mode 100644 internal/pkg/pipeline/task/file/file_test.go diff --git a/internal/pkg/pipeline/task/file/README.md b/internal/pkg/pipeline/task/file/README.md index 4c21b3c..364d469 100644 --- a/internal/pkg/pipeline/task/file/README.md +++ b/internal/pkg/pipeline/task/file/README.md @@ -37,7 +37,7 @@ In read mode, two values are stored in each record's context: | `tags` | map[string]string | - | S3 **write** only: object tags applied on `PutObject`. Ignored for local paths. Values support macros and context templates. See [S3 object tags](#s3-object-tags). | | `success_file` | bool | `false` | Whether to create a success file after writing | | `success_file_name` | string | `_SUCCESS` | Name of the success file | -| `task_concurrency` | int | `1` | Write: competing-consumer workers. Read: parallel file fetches (the pipeline still runs a single source worker so files are not duplicated) | +| `task_concurrency` | int | `1` | Write: competing-consumer workers. Read: parallel file fetches | | `context` | map | - | JQ expressions whose results are stored on each record for downstream tasks | | `fail_on_error` | bool | `false` | Whether to stop the pipeline if this task encounters an error | diff --git a/internal/pkg/pipeline/task/file/file_test.go b/internal/pkg/pipeline/task/file/file_test.go deleted file mode 100644 index 78885d1..0000000 --- a/internal/pkg/pipeline/task/file/file_test.go +++ /dev/null @@ -1,59 +0,0 @@ -package file - -import ( - "fmt" - "os" - "path/filepath" - "testing" - - "github.com/patterninc/caterpillar/internal/pkg/config" - "github.com/patterninc/caterpillar/internal/pkg/pipeline/record" - "github.com/patterninc/caterpillar/internal/pkg/pipeline/task" -) - -func TestReadFileConcurrent(t *testing.T) { - - dir := t.TempDir() - want := make(map[string]struct{}, 10) - for i := 0; i < 10; i++ { - content := fmt.Sprintf("content-%d", i) - if err := os.WriteFile(filepath.Join(dir, fmt.Sprintf("f%d.txt", i)), []byte(content), 0o644); err != nil { - t.Fatal(err) - } - want[content] = struct{}{} - } - - f := &file{ - Base: task.Base{ - Name: "read", - Type: "file", - TaskConcurrency: 4, - }, - Path: config.String(filepath.Join(dir, "*.txt")), - } - - out := make(chan *record.Record, 20) - errCh := make(chan error, 1) - go func() { - errCh <- f.readFile(out) - close(out) - }() - - got := make(map[string]struct{}) - for r := range out { - got[string(r.Data)] = struct{}{} - } - if err := <-errCh; err != nil { - t.Fatalf("readFile: %v", err) - } - - if len(got) != len(want) { - t.Fatalf("got %d records, want %d", len(got), len(want)) - } - for content := range want { - if _, ok := got[content]; !ok { - t.Errorf("missing content %q", content) - } - } - -}