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():
- Copies
input/<id>todoing/<id>and logsstation_started. - Runs your
body(wp)with the path of thedoingcopy. Add files there. - If
bodyreturns: logsstation_complete, movesdoing/<id>tooutput/<id>, and moves the untouchedinput/<id>todone/<id>. Returns{ status: 'ok', value }. - If
bodythrows: logsstation_failed, moves the untouched input tofailed/<id>, writeserror.jsonanderror.txtthere, and removes thedoingcopy. 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.