Fob
Browse the docs

Workpieces and bins

On a line, each station works through the workpieces in its input bin and leaves results in output for the next station. lib-worker does the bookkeeping: moving folders between bins, logging events, and recording errors. Read Concepts first for the five bins.

Bins live under temp/stations/ in the worker repo. bin('IN1', 'input') returns the absolute path of a bin and creates it if needed.

At the head of a line: create workpieces#

The first station creates one folder per unit of work directly in its output bin, writes pointer.json, and logs workpiece_created:

import fs from 'node:fs';
import path from 'node:path';
import { defineStep, bin, logEvent } from '@fob/lib-worker';
import { z } from 'zod';

export default defineStep({
  slug: 'IN0_01_seed',
  name: 'Seed invoices',
  description: 'Creates one workpiece per invoice',
  inputSchema: z.object({
    invoices: z.array(z.object({ number: z.string(), vendor: z.string(), amount: z.number() })),
  }),
  outputSchema: z.object({ seeded: z.number() }),
  execute: async (config, { work_record }) => {
    for (const invoice of config.invoices) {
      const wp = path.join(bin('IN0', 'output'), invoice.number);
      fs.mkdirSync(wp, { recursive: true });
      fs.writeFileSync(path.join(wp, 'pointer.json'), JSON.stringify({ workpiece_id: invoice.number, vendor: invoice.vendor }, null, 2));
      fs.writeFileSync(path.join(wp, 'invoice.json'), JSON.stringify(invoice, null, 2));
      logEvent(wp, 'IN0', work_record.id, 'workpiece_created');
    }
    return { seeded: config.invoices.length };
  },
});

Use a stable id from the source (an invoice number, a message id, a file hash) as the folder name, so the same item doesn't become two workpieces. pointer.json is written once and not changed by later stations.

Downstream: processWorkpiece()#

A later station's processing step loops over its input bin and hands each workpiece to processWorkpiece():

import fs from 'node:fs';
import path from 'node:path';
import { defineStep, bin, listWorkpieces, processWorkpiece } from '@fob/lib-worker';
import { z } from 'zod';

export default defineStep({
  slug: 'IN1_01_check',
  name: 'Check invoices',
  description: 'Fails invoices above the approval limit',
  inputSchema: z.object({ approval_limit: z.number().default(5000) }),
  outputSchema: z.object({ ok: z.number(), failed: z.number() }),
  execute: async (config, { work_record }) => {
    const counts = { ok: 0, failed: 0 };
    for (const workpiece_id of listWorkpieces(bin('IN1', 'input'))) {
      const result = await processWorkpiece({
        station: 'IN1',
        workpiece_id,
        work_record_id: work_record.id,
        body: async (wp) => {
          const invoice = JSON.parse(fs.readFileSync(path.join(wp, 'invoice.json'), 'utf8'));
          if (invoice.amount > config.approval_limit) {
            throw new Error(`Amount ${invoice.amount} is above the approval limit of ${config.approval_limit}`);
          }
          fs.writeFileSync(path.join(wp, 'check.md'), `# ${invoice.number}\n\nWithin the approval limit.\n`);
        },
      });
      counts[result.status === 'ok' ? 'ok' : 'failed'] += 1;
    }
    return counts;
  },
});

For each workpiece, processWorkpiece():

  1. Copies input/<id> to doing/<id> and logs station_started.
  2. Runs your body(wp) with the path of the doing copy. Add files there.
  3. If body returns: logs station_complete, moves doing/<id> to output/<id>, and moves the untouched input/<id> to done/<id>. Returns { status: 'ok', value }.
  4. If body throws: logs station_failed, moves the untouched input to failed/<id>, writes error.json and error.txt there, and removes the doing copy. Returns { status: 'failed', error }. It doesn't throw, so one bad workpiece doesn't fail the whole run.

Extra properties on the error you throw (such as status or body from an HTTP response) are saved in error.json.

The log#

Every workpiece's log.jsonl has one event per line:

{"ts":"2026-10-03T14:35:02.855Z","station":"IN0","wr":"local-wr-1791038102853","event":"workpiece_created"}
{"ts":"2026-10-03T14:35:03.032Z","station":"IN1","wr":"local-wr-1791038103027","event":"station_started"}
{"ts":"2026-10-03T14:35:03.032Z","station":"IN1","wr":"local-wr-1791038103027","event":"station_failed"}

wr is the work record of the run that did it, so you can find the run in the Orchestrator. fob-worker workpieces show <id> prints this as a journey.

Failures and retries#

fob-worker lines status IN                 # ⚠ 1 at IN1/failed
fob-worker workpieces show INV-1002        # where it is and what happened
cat temp/stations/IN1/failed/INV-1002/error.txt

To retry after fixing the cause, move it back to input:

mv temp/stations/IN1/failed/INV-1002 temp/stations/IN1/input/

The next run of the station picks it up. Locally, run the step again with fob-worker steps run IN1_01_check --station IN1.

Testing a line locally#

fob-worker steps run runs your own steps but not the built-in move_files conveyor. Between stations, move the workpieces yourself:

fob-worker steps run IN0_01_seed --station IN0
mkdir -p temp/stations/IN1/input && mv temp/stations/IN0/output/* temp/stations/IN1/input/
fob-worker steps run IN1_01_check --station IN1
fob-worker lines status IN

To start over, empty the bins: fob-worker lines empty-bins IN --all-bins. It shows what it will delete and asks first.

Interrupted runs#

If the worker stops mid-run, a workpiece can be left in doing. The worker clears stale doing folders when it starts; the workpiece is still in input and runs again.