> ## Documentation Index
> Fetch the complete documentation index at: https://docs.bijection.com/llms.txt
> Use this file to discover all available pages before exploring further.

# Pipelines

> Capture input tables, transform them and publish the result

A pipeline turns input tables into [published tables](/publications/published-tables).
You declare it with `defineEtl` from the `@bijection/pipelines` component: the
inputs to capture, the tables to publish, and a `transform` function. Bijection
reruns the pipeline whenever an input changes and publishes every output table
together.

This page builds the `orderSummary` pipeline, which publishes the
`open_orders` and `customer_totals` tables declared on the
[previous page](/publications/published-tables).

## Setting up a pipeline

<Steps>
  <Step title="Install from `npm`">
    Install the pipelines component, and the datasets package whose access
    types the policy below uses.

    ```sh theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
    npm i @bijection/pipelines @bijection/datasets
    ```
  </Step>

  <Step title="Add the pipelines component">
    Add the component to your app by calling `app.use` in the
    `bijection.config.ts` file in your app's `bijection/` folder. The next
    `bijection dev` generates `components.pipelines`, which the pipeline
    definition below refers to. See [using components](/components/using).

    ```ts bijection/bijection.config.ts theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
    import { defineApp } from "bijection/server";
    import pipelines from "@bijection/pipelines/bijection.config.js";

    const app = defineApp();
    app.use(pipelines);
    export default app;
    ```
  </Step>

  <Step title="Decide who may run it">
    A pipeline checks every operation against an access policy: an internal
    query that receives the caller's `actor` (their `tokenIdentifier`), the
    requested `operation` and the pipeline as `resource`. It returns
    `{ scope }` to grant the operation, or `null` to refuse it.

    ```ts bijection/access.ts theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
    import { internalQuery } from "./_generated/server";
    import { accessRequest, accessGrant } from "@bijection/datasets/access";

    export const pipelineAccess = internalQuery({
      args: accessRequest,
      returns: accessGrant,
      handler: async (ctx, { actor }) => {
        // `operators` is an ordinary table of users allowed to run pipelines.
        const operator = await ctx.db
          .query("operators")
          .withIndex("by_token", (q) => q.eq("tokenIdentifier", actor))
          .unique();
        return operator ? { scope: "orders" } : null;
      },
    });
    ```

    Return the same `scope` for the same actor every time. Background work runs
    as the user who configured the pipeline, and the policy is asked again for
    that user as the pipeline captures and publishes. If it stops granting
    them, the pipeline's published tables become unreadable.
  </Step>

  <Step title="Define the pipeline">
    Call `defineEtl` and export the functions it generates.

    ```ts bijection/summaries.ts theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
    import { defineEtl, Etl } from "@bijection/pipelines/etl";
    import {
      makeFunctionReference,
      type FunctionReference,
    } from "bijection/server";
    import { components, internal } from "./_generated/api";
    import { mutation, query } from "./_generated/server";
    import { order, openOrders, customerTotals } from "./schema";

    // Named by string: the definition refers to its own exports.
    const ref = (name: string) =>
      makeFunctionReference<"query">(name) as unknown as FunctionReference<
        "query",
        "internal"
      >;

    const plan = defineEtl({
      key: "orderSummary",
      component: components.pipelines,
      access: internal.access.pipelineAccess,
      inputs: {
        orders: { table: "orders", row: order },
      },
      outputs: {
        openOrders: { table: "open_orders", key: "order_number", row: openOrders },
        customerTotals: { table: "customer_totals", key: "customer", row: customerTotals },
      },
      publisher: {
        manifest: ref("summaries:manifest"),
        steps: ref("summaries:steps"),
        page: ref("summaries:page"),
        rebuild: "summaries:rebuild",
      },
      transform({ orders }) {
        const open = orders.filter((o) => o.status === "open");
        const totals = new Map<string, { open_count: bigint; open_amount: bigint }>();
        for (const o of open) {
          const t = totals.get(o.customer) ?? { open_count: 0n, open_amount: 0n };
          t.open_count += 1n;
          t.open_amount += o.amount;
          totals.set(o.customer, t);
        }
        return {
          output: {
            openOrders: open.map(({ order_number, customer, amount }) => ({
              order_number,
              customer,
              amount,
            })),
            customerTotals: [...totals].map(([customer, t]) => ({ customer, ...t })),
          },
          next: null,
        };
      },
    });

    export const capture = plan.capture;
    export const transform = plan.transform;
    export const manifest = plan.manifestQuery;
    export const steps = plan.stepsQuery;
    export const page = plan.pageQuery;
    export const publisher = plan.publisher;

    const orderSummary: Etl = new Etl(components.pipelines, {
      ...plan,
      capture: internal.summaries.capture,
      transform: internal.summaries.transform,
    });
    export const rebuild = orderSummary.deploymentMutation();
    ```

    Deploy with `bijection dev` or `bijection deploy`. Deployment checks that
    `manifest`, `steps`, `page` and `rebuild` name functions your app exports.
  </Step>

  <Step title="Configure it">
    Call `configure` from a mutation run by an authenticated user the policy
    grants. It returns the pipeline's group ID.

    ```ts bijection/summaries.ts theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
    export const configure = mutation({
      args: {},
      handler: (ctx) => orderSummary.configure(ctx),
    });
    ```

    The first run starts on its own. From then on, every committed change to an
    input table starts a new run. Changes that arrive during a run are
    collected, and one successor run picks them up. Calling `configure` again
    with an unchanged definition returns the same group; after you change the
    definition, call it again to apply the change.
  </Step>
</Steps>

## The pipeline definition

`defineEtl` takes these options:

| Option | Description |
| - | - |
| `key` | The pipeline's name. It must equal the `producer` of every output table. |
| `component` | The installed pipelines component, `components.pipelines`. |
| `access` | The access policy query. It cannot be replaced after the first `configure`. |
| `inputs` | 1 to 4 inputs, each `{ table, row }`, where `row` is the validator of that table's documents. |
| `outputs` | 1 to 4 outputs. Each name must equal its table's `member`; each gives the `table`, its `definePublished` declaration as `row`, and the `key` field. |
| `publisher` | The exported names of the generated `manifest`, `steps` and `page` queries, and of the `rebuild` mutation. |
| `transform` | The function that computes the outputs. |
| `stream` | Optional. The input to read in pages. See [streaming large inputs](#streaming-large-inputs). |
| `state` | Optional. A validator for the state carried between pages. |
| `capturePageRows` | Optional. A smaller page size, from 1 to 64 rows, for a streamed pipeline. |
| `validate` | Optional. Expectations over the complete result. See [checking the complete result](#checking-the-complete-result). |
| `freshness` | Optional. `{ maxAgeMs }`, from one second to 30 days. When the serving result is older, the pipeline raises a freshness alert. |

Pipeline, input and output names, table names and key fields must start with a
letter and contain only letters, digits and `_`, up to 64 characters.

## The publisher declaration

Bijection does not let your code write published rows. Instead, the engine
reads a result that is ready to publish through three internal queries, and
verifies every row it is handed before writing any of it. A publisher
declaration, created with `definePublisher` from `bijection/server`, names those
queries:

| Option | Description |
| - | - |
| `producer` | The producer whose tables this publisher fills. |
| `manifest` | An internal query answering which results are ready to publish. |
| `steps` | An internal query paging one ready result's sealed step records. |
| `page` | An internal query paging one output's staged rows. |
| `rebuild` | An internal mutation that deployment schedules when it reserves a new owner for the group. |
| `transition` | Optional. See [changing a published group](#changing-a-published-group). |

Each function is given as a reference or as a `"module:function"` string, and
`manifest`, `steps`, `page` and `rebuild` must be four different functions. The
declaration must be exported from a module of your app, where deployment
discovers it.

`defineEtl` builds all of this for you. `plan.publisher` is the
`definePublisher` declaration, `plan.manifestQuery`, `plan.stepsQuery` and
`plan.pageQuery` are the three queries, and `deploymentMutation()` is the
`rebuild` mutation. The names in the `publisher` option are where you exported
them.

Only your app's root may declare a publisher. A component cannot publish into
its own tables this way.

## Writing the transform

`transform` receives an object with one array per input, the carried `state`
(`null` on the first page) and `{ step, final }`. Input rows include `_id` and
`_creationTime`. It returns `{ output, next }`: one array of rows per output,
and the state for the next page, or `null`.

Input and output rows are limited to these types:

* `v.string()`, `v.int64()`, `v.number()` (finite values only), `v.boolean()`
  and `v.null()`
* nested `v.object(...)`, up to eight levels deep
* optional fields with `v.optional(...)`

Arrays, unions, literals, IDs, bytes and records are not supported in pipeline
rows. An input row must match its declared validator exactly; a document with
an undeclared field fails the capture.

Every output row must match its published table's validator. Its `key` field
must be a required string, and each key may appear only once in an output over
the whole run. A duplicate key fails the run. An empty output is valid, and
publishes an empty table.

The transform runs in a restricted action:

* It has no network access. A `fetch` fails the step.
* It has about two seconds of execution time per step.
* It cannot change its input; it works on a copy.

If `transform` throws, the run fails with that error and nothing is published.
The previous publication keeps serving.

<Warning>
  Without `stream`, a run captures every input in full into a single page of
  at most 64 rows and 256 KiB in total. A larger input fails the run. Stream
  any input that can grow beyond that.
</Warning>

## Streaming large inputs

Set `stream` to one input to read it in pages. The transform then runs once per
page and can carry `state` from one page to the next. Every page comes from the
same snapshot of the input, taken when the run started.

```ts bijection/invoices.ts {4-5,10-21} theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
const plan = defineEtl({
  key: "invoiceTotals",
  inputs: { invoices: { table: "invoices", row: invoice } },
  stream: "invoices",
  state: v.object({ invoice_count: v.int64(), amount: v.int64() }),
  outputs: {
    totals: { table: "invoice_totals", key: "scope", row: invoiceTotals },
  },
  // component, access and publisher as above
  transform({ invoices }, state, meta) {
    const next = {
      invoice_count: (state?.invoice_count ?? 0n) + BigInt(invoices.length),
      amount:
        (state?.amount ?? 0n) +
        invoices.reduce((sum, invoice) => sum + invoice.amount, 0n),
    };
    // Emit the total once, after the last page.
    return meta.final
      ? { output: { totals: [{ scope: "all", ...next }] }, next: null }
      : { output: { totals: [] }, next };
  },
});
```

A step may emit rows on any page; the published tables hold the rows of every
step. `meta.final` is `true` on the last page of the input. Each page holds at
most 64 rows of the streamed input, fewer when other inputs are present: every
input that is not streamed is captured in full beside each page and counts
against the same 64-row budget.

The state is a validated value of at most 32 KiB, built from the same types as
rows. It cannot hold arrays or records, so an aggregate over an open-ended set
of keys belongs in the output rows, not in the state.

A streamed input holds at most 8,192 rows, captured within 240 seconds. If the
snapshot expires while pages are being read, the capture starts over from a
fresh snapshot, for at most three attempts in total. A run whose input exceeds 8,192 rows fails
with `ETL input exceeds 8192 rows`; correcting the input starts a new run.

## Checking the complete result

`validate` declares expectations over the complete result of a run, evaluated
after its last step and before it can be published. A run with a failed
blocking expectation publishes nothing, and the previous publication keeps
serving.

```ts bijection/summaries.ts theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
validate: {
  state: v.object({ orders: v.int64(), totals: v.int64() }),
  reduce({ openOrders, customerTotals }, prior) {
    return {
      orders:
        (prior?.orders ?? 0n) + openOrders.reduce((n, r) => n + r.amount, 0n),
      totals:
        (prior?.totals ?? 0n) +
        customerTotals.reduce((n, r) => n + r.open_amount, 0n),
    };
  },
  check(state, candidate) {
    return [
      {
        expectation: "totals match orders",
        severity: "blocking",
        passed: state.orders === state.totals,
      },
      {
        expectation: "open orders did not collapse",
        severity: "advisory",
        passed:
          !candidate.previous ||
          candidate.rows.openOrders * 2 >= candidate.previous.rows.openOrders,
      },
    ];
  },
  references: [{ from: "openOrders", field: "customer", to: "customerTotals" }],
},
```

* `reduce(outputs, state, meta)` folds the output of each step, in step order,
  into a validated `state` of at most 32 KiB. The first call receives `null`.
* `check(state, candidate)` runs once, after the last fold, and returns a list
  of verdicts `{ expectation, severity, passed, detail? }`. `candidate.rows`
  holds each output's row count. `candidate.previous` describes the serving
  publication it would replace, or is `null`.
* `references`, at most two, require every `from` row's string `field` to be a
  key of the `to` output in the same result.
* `rejects: { member, max }` names an output that holds rejected input records.
  Its table needs a required string `reason` field. More than `max` rejected
  rows blocks the result; any rejected row is an advisory warning. The rejected
  rows publish with the other outputs, readable under that table's own access
  rule. `reduce` never receives them.

A `blocking` verdict that fails holds the result back with the error
`Expectation <name> failed: <detail>`. An `advisory` verdict that fails is kept
as a warning and does not prevent publication. A throw in `reduce` or `check`
counts as a failed blocking verdict named `validation`. You can return at most
32 verdicts, with unique names of 1 to 128 characters that do not contain `:`,
and details of at most 1 KiB.

## Managing runs

The `Etl` object provides these operations. Call them from your own
mutations and queries; each one requires an authenticated user and checks the
access policy for them.

| Method | What it does |
| - | - |
| `configure(ctx)` | Create the pipeline, or apply a changed definition. |
| `publication(ctx)` | Report which result is selected and which is serving. See [reading publications](/publications/reading#knowing-what-is-serving). |
| `refresh(ctx, group)` | Start a run now. |
| `pause(ctx, group)`, `resume(ctx, group)` | Stop and restart automatic runs. The serving publication stays readable while paused. |
| `retry(ctx, run)` | Rerun a failed run on exactly its retained input. |
| `rebuild(ctx, group)` | Capture current inputs and run the current code, even when the input is unchanged. |
| `select(ctx, run, { reason, refresh })` | Publish a retained successful run, for example to roll back. `refresh` is `"pause"` or `"continue"`. |
| `withdraw(ctx, run, reason, evidence, defect?)` | Mark a run defective, by its input (`"basis"`, the default) or its code (`"code"`), and make it unreadable. |

Mutations also accept a final `{ request }` option. Repeating an operation with
the same request key returns the recorded outcome instead of acting again,
which makes a retry after a lost reply safe.

`select` returns `{ status: "requested", run, selection }`. Publication happens
afterwards, in its own commit; use `publication(ctx)` to see when it is
serving. A withdrawn publication stops serving at once, and its tables stay
unavailable until another run publishes. Data a client read earlier cannot be
recalled.

## Changing a published group

<Warning>Ownership transitions are in beta.</Warning>

Changing a pipeline's transform code needs no special step: the next run uses
the new code. Changing which tables the group publishes does. Renaming the
pipeline's `key`, moving its `publisher` export to another module or name,
adding, removing or renaming an output, changing an output table's row type,
or retiring the pipeline all require an explicit `transition` in the publisher
declaration. A deployment that makes such a change without one is refused and
changes nothing.

```ts bijection/summaries.ts {6-14} theme={"theme":{"light":"github-light-default","dark":"github-dark-default"}}
publisher: {
  manifest: ref("summaries:manifest"),
  steps: ref("summaries:steps"),
  page: ref("summaries:page"),
  rebuild: "summaries:rebuild",
  transition: {
    from: {
      publisher: "published/summaries.js:publisher",
      producer: "orderSummary",
      tables: ["open_orders", "customer_totals"],
      expectedEpoch: "<current epoch>",
    },
    toProducer: "orderSummary",
  },
},
```

`from` names the installed owner exactly: its publisher, producer, complete
list of output tables, and the current owner epoch. Read these from the
deployment's publisher ownership inspection, a `POST` to
`/api/get_config_hashes` on the deployment URL with `publisherOwnership: true`
and the deployment's admin key. The CLI does not fill in the epoch or show the
plan before deploying. A transition whose epoch is no longer current is refused
with `PublisherOwnership`. Set `toProducer: null` to retire the group; remove
the declaration on a later deployment.

A transition makes the whole group unavailable until the new owner publishes a
complete run. Reads fail with `PublishedUnavailable` in between; unchanged rows
are recomputed, not carried over. Deployment schedules the `rebuild` mutation,
which starts that run without any further call. A retired table stays reserved,
unreadable and unwritable, until you delete it.

## Limits

| Limit | Value |
| - | - |
| Inputs per pipeline | 1 to 4 |
| Outputs per pipeline | 1 to 4 |
| Capture page, all inputs together | 64 rows, 256 KiB |
| Streamed input per run | 8,192 rows (128 pages), captured within 240 s |
| Transform steps per run | 256 |
| Output of one step, all outputs together | 512 rows, 64 KiB |
| Output of one run | 65,536 rows, 8 MiB |
| Carried state, validation state | 32 KiB each |
| Transform execution time per step | about 2 s |
| Pipelines per component instance | 32 |
| Pipelines per producer | 1 |
| Tables one pipeline may claim over its lifetime | 8 |
| Concurrent captures per deployment | 32; further captures wait for a free slot |
| Pipeline definition | 16 KiB |

Every limit that a definition can exceed is refused when you call `configure`,
before it affects running pipelines. The limits on input and output size are
checked as the run proceeds, and a run that exceeds one fails without
publishing.

<Note>
  These limits bound what a run may do; they are not measured capacity.
  Throughput at the upper limits has not been measured.
</Note>
