diff --git a/docs/pipeline.md b/docs/pipeline.md new file mode 100644 index 0000000..bb62b37 --- /dev/null +++ b/docs/pipeline.md @@ -0,0 +1,193 @@ +## `tcli` pipeline (YAML) schema v1 + +An alternative to shell-pipe scenario tests. A pipeline file describes named +steps, how data flows between them, and pipeline-wide defaults/variables. +Runnable as `tcli pipeline run ` and importable as a Go library +from `pkg/pipeline`. + +### Why not just bash pipes? +- Bash pipelines are linear no fan-out (reuse a step's result in two + downstream branches) and no fan-in (converge two branches into one). +- No named steps, so logs and reruns are hard to correlate. +- No shared defaults/variables `-format` strings and endpoints get + duplicated and misquoted. +- Not portable to Windows shells. + +### File shape + +```yaml +name: # pipeline name (required) +description: # human summary (optional) + +variables: # scalar reusable values + : # referenced as ${{ variables. }} + +concurrency: # inter-step parallelism controls (optional) + maxParallel: # cap simultaneous steps; unlimited if unset + onFailure: cancel|drain # cancel (default): abort in-flight siblings on + # first failure; drain: let them finish + +defaults: # applied to every step; step-level overrides win + verbose: + ignoreErrors: + retryCount: + parallelism: # see the Parallelism section below + statusCode: + basePath: # endpoint overrides; usually per-module but + scheme: # settable pipeline-wide when a whole run + server: # targets a non-default host + jwt: # bearer token; prefer ${{ env. }} + +steps: # ordered list; execution order is a DAG derived + # from inputFrom + ${{ steps.* }} references + + # explicit dependsOn + - name: # unique step id (required) + command: # " [] " (required) + params: {...} # tcli flags/parameters for the operation + body: # sugar for params.body; object is auto-serialized + format: # jq expression applied to this step's output + # (equivalent to tcli's -format flag) + + inputFrom: # feed 's stdout as this step's stdin + # exactly like a shell pipe + + outputs: # capture named scalars from this step's response + : # later steps read ${{ steps..outputs. }} + + dependsOn: [] # explicit ordering; usually inferred from + # inputFrom / ${{ steps.* }} references + + condition: # skip this step if expression is false/null + # applied against pipeline state (see below) + continueOnError: # do not abort the pipeline if this step fails + + # any tcli flag can be set as a top-level step field for convenience + count: + parallelism: # unset or 1 = sequential; N>1 = N workers; + # 0 = auto (GOMAXPROCS, matches CLI -parallel) + retryCount: + ignoreErrors: + statusCode: + verbose: +``` + +### Wiring: `inputFrom` vs `${{ steps.*.outputs.* }}` + +Two orthogonal mechanisms. Pick per use case. + +| Mechanism | Use when | Semantics | +|---|---|---| +| `inputFrom: ` (+ optional `format`) | Bulk streaming; N records flow through | Equivalent to a shell pipe `A \| B`. Runner wires `A`'s stdout to `B`'s stdin. If both steps declare `parallel: true`, it streams; otherwise `B` starts after `A` completes. | +| `${{ steps..outputs. }}` inside `params:` (with an `outputs:` block on the source) | A later step needs one scalar value from an earlier response | Buffers the source step's response; extracts values with jq. Automatically creates a `dependsOn` edge. | + +Both mechanisms can coexist on the same step: a step can stream data in via +`inputFrom` *and* reference `${{ steps.other.outputs.token }}` in its params. + +**`outputs:` on multi-record steps.** The `outputs:` jq expression is +evaluated against each response the step produces. Recommend using +`outputs:` only on single-response steps (the common create-then-reference +pattern). If a downstream step references `${{ steps.X.outputs.var }}` +where `X` produced N responses, the value is the last one captured. +For bulk data flow between multi-record steps, use `inputFrom` instead. + +### Reserved keys + +The following tcli control flags MUST be set via top-level `camelCase` step +fields (or under `defaults:`), NOT inside `params:`. The loader rejects them +under `params:` to prevent two spellings for the same setting: + +`format`, `count`, `parallelism`, `verbose`, `retryCount`, `ignoreErrors`, +`statusCode`, `basePath`, `scheme`, `server`, `jwt`. + +Everything else under `params:` maps directly to a swagger operation +parameter, using the parameter name as declared in the spec. + +### Expression syntax + +`${{ }}` interpolation is supported everywhere a value is expected +(string params, format strings, condition expressions). Recognized paths: + +- `${{ variables. }}` from top-level `variables:` +- `${{ steps..outputs. }}` from a step's `outputs:` block +- `${{ steps..status }}` HTTP status code of a completed step +- `${{ steps..ok }}` boolean; true if step succeeded +- `${{ env. }}` from process env + +### Execution model + +1. Parse the file. Validate step names are unique and all references resolve. +2. Build a DAG from `inputFrom`, `${{ steps.* }}` references, and explicit + `dependsOn`. Reject cycles. +3. Run steps in topological order. Independent branches can run concurrently, + capped by `concurrency.maxParallel` (unlimited if unset). +4. A step fails if the underlying command returns non-zero, unless the step + has `continueOnError: true` or `ignoreErrors: true`. On first failure, + `concurrency.onFailure` decides whether in-flight sibling steps are + cancelled (default) or drained. +5. Pipeline exit code is 0 iff every non-`continueOnError` step succeeded. + +### Parallelism + +Two orthogonal axes: + +| Axis | Field | What it parallelizes | +|---|---|---| +| Inter-step | `concurrency.maxParallel` (pipeline) | Independent branches of the DAG run at the same time | +| Intra-step | `parallelism: ` (step) | One step processes its N input records concurrently. `parallelism: 0` matches the existing `-parallel` tcli flag (auto = GOMAXPROCS). | + +The two compose freely: a step can be single-threaded internally +(`parallelism` unset or `1`) but still run at the same wall-clock time +as other independent steps. + +For the rest of this section, a step is called **concurrent** when +`parallelism` is unset-or-`1` (sequential) vs. anything else (workers). + +#### `inputFrom` parallelism semantics + +| Producer | Consumer | Behavior | +|---|---|---| +| sequential | sequential | Producer emits all records, then consumer processes them one at a time. Deterministic order. | +| sequential | concurrent | Producer runs to completion; consumer fans records out to workers. Output order nondeterministic. | +| concurrent | sequential | Producer streams; consumer processes serially in arrival order. Matches `A -parallel \| B` today. | +| concurrent | concurrent | **Streaming**: producer and consumer run concurrently, records handed off as they're produced. Output order nondeterministic. | + +#### Fan-out over a stream + +Two steps can both `inputFrom: A`. The runner buffers `A`'s output and +replays it to each consumer a stream can't be teed without buffering, +so plan for the memory cost if `A` produces a large volume. + +#### Fanning out over a parameter grid + +There is no dedicated `matrix:` field. Express a grid as records emitted +by an upstream step; the consumer's `parallelism` handles the concurrent +invocations, one per record: + +```yaml +- name: grid + command: utils echo + format: '[["us-east-1","us-west-2","eu-west-1"][], [10,50][]] + | {region:.[0], batchSize:.[1]}' + # emits 6 records: {region:"us-east-1",batchSize:10}, + +- name: teardown + command: petstore pet deletePet + inputFrom: grid + parallelism: 6 # -> 6 concurrent invocations +``` + +### Minimal example (equivalent to `examples/example.sh`) + +See [`examples/pipeline/petstore_crud.yaml`](/examples/pipeline/petstore_crud.yaml). + +### Open questions (v1 v2) + +- **Includes / templates** pull a step group from another file so a + create/verify/delete triple can be reused across pipelines. +- **Secrets** pipeline-level `secrets:` mapped to env for HTTP auth. + +### References + +- [Bash example explanation](/docs/example_explanation.md) +- [How to add a module](/docs/modules.md) +- [jq language](https://github.com/jqlang/jq/wiki/jq-Language-Description) diff --git a/examples/README.md b/examples/README.md index c796c49..ab5cca1 100644 --- a/examples/README.md +++ b/examples/README.md @@ -14,6 +14,7 @@ See [extending support for custom commands](/examples/_pubsub/README.md) ### References - [README.md](/README.md) +- [Pipeline example explanation](/docs/pipelines.md) - [how to build and run](/docs/build_and_run.md) - Format - [jq lang](https://github.com/jqlang/jq/wiki/jq-Language-Description) diff --git a/examples/pipeline/petstore_crud.yaml b/examples/pipeline/petstore_crud.yaml new file mode 100644 index 0000000..ea223e5 --- /dev/null +++ b/examples/pipeline/petstore_crud.yaml @@ -0,0 +1,43 @@ +name: petstore-crud +description: | + 1:1 equivalent of examples/example.sh, expressed as a tcli pipeline. + Seed N pets, create them, read them back, delete them, and verify + that a follow-up read returns 404. + +variables: + petCount: 1 + status_notfound: 404 + +defaults: + verbose: false + retryCount: 2 + +steps: + # utils echo emits N seed records on stdout, wrapped in {body: } + # so the next step (addPet) picks up its "body" parameter directly. +- name: seed + command: utils echo + format: 'range(1;${{ variables.petCount }}+1) | {body: {id:., name:.|tostring, photourls:[.|tostring]}}' + +- name: create + command: petstore pet addPet + inputFrom: seed + format: '{petId:.id}' + +- name: read + command: petstore pet getPetById + inputFrom: create + format: '{petId:.id, api_key:.id}' + +- name: delete + command: petstore pet deletePet + inputFrom: read + format: '{petId:.message}' + + # A follow-up read must 404 for every deleted id; step fails otherwise. +- name: verify_gone + command: petstore pet getPetById + inputFrom: delete + params: + doc: shell + statusCode: ${{ variables.status_notfound }} diff --git a/examples/pipeline/petstore_dag.yaml b/examples/pipeline/petstore_dag.yaml new file mode 100644 index 0000000..d62f0b5 --- /dev/null +++ b/examples/pipeline/petstore_dag.yaml @@ -0,0 +1,54 @@ +name: petstore-dag +description: | + Demonstrates flows a bash pipeline cannot express: + - fan-out: one create feeds two independent verifications + - scalar reference: a single field from a response is injected as a + param on a later step via ${{ steps..outputs. }} + - conditional teardown: delete only runs if the create step succeeded + +variables: + petId: 4242 + status_ok: 200 + status_notfound: 404 + +steps: +- name: create_one + command: petstore pet addPet + body: + id: ${{ variables.petId }} + name: fido + photoUrls: [http://img.example/fido.png] + outputs: + newId: .id # capture the created id off the response + newName: .name + + # -------- fan-out: two branches read from create_one in parallel -------- + +- name: read_by_id + command: petstore pet getPetById + params: + petId: ${{ steps.create_one.outputs.newId }} + statusCode: ${{ variables.status_ok }} + +- name: search_by_status + command: petstore pet findPetsByStatus + params: + status: available + dependsOn: [create_one] # explicit; no data reference + + # -------- fan-in: teardown depends on both branches finishing -------- + +- name: delete_one + command: petstore pet deletePet + params: + petId: ${{ steps.create_one.outputs.newId }} + dependsOn: [read_by_id, search_by_status] + condition: ${{ steps.create_one.ok }} # skip teardown if create failed + continueOnError: true # never block the pipeline on cleanup + +- name: verify_gone + command: petstore pet getPetById + params: + petId: ${{ steps.create_one.outputs.newId }} + dependsOn: [delete_one] + statusCode: ${{ variables.status_notfound }} diff --git a/examples/pipeline/petstore_parallel.yaml b/examples/pipeline/petstore_parallel.yaml new file mode 100644 index 0000000..c1010b5 --- /dev/null +++ b/examples/pipeline/petstore_parallel.yaml @@ -0,0 +1,62 @@ +name: petstore-parallel +description: | + Exercises both parallelism axes: + - inter-step: two verification branches run at the same time (DAG fan-out) + - intra-step: create/verify/teardown each process 100 records concurrently + Also shows the "fan out over a parameter grid" pattern via stdin+jq, + which replaces the need for a dedicated matrix: field. + +concurrency: + maxParallel: 8 # cap simultaneous DAG branches + onFailure: cancel # abort in-flight siblings when any step fails + +variables: + petCount: 100 + status_ok: 200 + status_notfound: 404 + +defaults: + retryCount: 5 + +steps: +- name: seed + command: utils echo + format: 'range(1;${{ variables.petCount }}+1) | {body: {id:., name:.|tostring, photourls:[.|tostring]}}' + + # Intra-step parallelism: 100 addPet calls fan out across 16 workers. + # Streams from `seed` because the consumer is concurrent. +- name: create + command: petstore pet addPet + inputFrom: seed + parallelism: 16 + format: '{petId:.id}' + + # ---- inter-step fan-out: these two branches run concurrently ---- + +- name: verify_read + command: petstore pet getPetById + inputFrom: create # runner buffers `create`'s output (two consumers) + parallelism: 8 + statusCode: ${{ variables.status_ok }} + +- name: verify_search + command: petstore pet findPetsByStatus + params: + status: available + dependsOn: [create] # non-data dependency; runs alongside verify_read + + # ---- fan-in: teardown waits for both branches ---- + +- name: teardown + command: petstore pet deletePet + inputFrom: create + parallelism: 8 + dependsOn: [verify_read, verify_search] + format: '{petId:.message}' + continueOnError: true + +- name: verify_gone + command: petstore pet getPetById + inputFrom: teardown + parallelism: 8 + statusCode: ${{ variables.status_notfound }}