@aylonmuramatsu/flow-pipes
v1.1.0
Published
Motor de pipelines tipado com contexto compartilhado, Issues e execução paralela
Maintainers
Readme
Flow Pipes
Languages: English · Português (Brasil)
Por que criar services e fluxos complexos não pode ser divertido e simples — como uma pipeline?
Cansou de service que só “conversa” na base do throw? A gente também.
Flow Pipes é um motor de pipelines tipado pra Node.js / TypeScript. A ideia é quase ofensiva de tão simples: monta uma esteira de etapas, compartilha estado sem drama, trata erro de negócio como Issue (sem derrubar o fluxo por birra) e você decide quando a festa acaba.
Menos try/catch em cascata. Mais código que dá pra ler sem café intravenoso.
Instalação (a parte chata, rápida)
yarn add @aylonmuramatsu/flow-pipesSó pede Node.js >= 18. Sem ritual.
Novidades 1.1.0
| Feature | Uso |
| ------- | --- |
| ctx.issue({ ... }) | objeto → a lib cria a classe (sem declarar Issue no app) |
| .step([a, b, c]) | vários steps em sequência sem encadear |
| .when(pred, step\|steps) / { when } | step condicional (skipped na trail) |
| .each(..., { fatalScope }) | "flow" (default) para o lote; "item" só o item |
| ctx.stop({ reason, code }) | freio com stopInfo (≠ Fatal) |
| issue.code + metadata tipada | Issues genéricas + code first-class |
| FieldRequiredIssue | message + metadata livre; default Error |
| traceId + durationMs | correlação do run + duração na trail |
| retry no step | só em throw; fixed / exponential |
| run(input, { signal, timeoutMs, traceId }) | abort / timeout → ABORTED |
Por que isso existe (spoiler: dor de verdade)
Exception não é contrato de negócio
Quando um service só fala com o outro via exception, qualquer caminho “silencioso” vira festival de try/catch. Você queria processar pedido; terminou de bombeiro de stack trace às 18h47.
Aqui o domínio anota Issues no contexto. Você processa. No final: commit, rollback, resposta parcial ou “pode parar, pessoal”.
Excel com 800 linhas e a regra é “não pode morrer na primeira”
Já viveu isso? Regras por campo, precisa passar por tudo, fazer staging — e só no fim cancelar os inserts e devolver quais erros atrapalharam a festa.
“Estoura na linha 2” não é estratégia. É preguiça com stack. .each + Issues por linha + decisão no fim: aí sim.
if perdido no meio do service? A gente sente
Validação espalhada é esconde-esconde de bug. Cada preocupação vira um step. Issue soft? Segue o baile. FatalIssue / ctx.stop()? Agora encerra de propósito, não por acidente.
Uma “ação” que você reaproveita de verdade (não aquele ctrl+c emocionado)
Agrupa operações num Flow. Outro service chama a mesma ação. Integração deixa de ser fanfic copiada e vira composição. Seu eu do futuro vai te mandar um obrigado (ou pelo menos não vai te xingar no code review).
Shared: adeus variável global disfarçada de “tô só passando a ref”
Sem gambiarra de parâmetro em parâmetro. O shared é o estado do fluxo. Steps leem e escrevem ali. O motor cuida de ctx.result (trail, falhas). Domínio de um lado, fofoca do motor do outro.
Paralelo quando o relógio aperta
.parallel pra fan-out de I/O. .each({ concurrency }) pra lote com workers. Flow.runAll pra vários fluxos independentes. Otimiza sem inventar worker pool toda sprint como se fosse hobby.
Mental model invertido (e bem melhor)
- Monta a esteira
- Processa
- Olha Issues /
stopped/result - Aí commit, rollback ou responde o cliente
Primeiro o trabalho. Depois o veredito. Tipo série boa: sem spoiler no episódio 1.
Conceitos essenciais (versão café)
| Peça | Em uma frase |
| --------------------------- | ----------------------------------------------------- |
| Flow | Sua esteira — a “ação” completa |
| step | Uma etapa, uma responsabilidade. Sem novela mexicana. |
| ctx.input | O que entrou nesse step (no .each, é o item da vez) |
| ctx.shared | Onde o domínio mora durante o rolê |
| ctx.issue(...) | “Anota aí: deu ruim de negócio” |
| Issue | Problema de domínio. Não é exception de fantasia. |
| FatalIssue / ctx.stop({ reason?, code? }) | “Pode parar, pessoal.” (+ stopInfo) |
| ctx.traceId / trail durationMs | Correlação do run + quanto cada step demorou |
| ctx.result | O que o motor viu depois do run |
| ctx.stopped / stopInfo | Alguém puxou o freio (e por quê) |
| .each / fatalScope | Um filho por item; fatal no item ou no lote |
| .when / retry | Condicional e re-tentativa em throw |
| .parallel / concurrency | Mais rápido, menos fila no caixa |
Steps se falam pelo shared. O motor monta o result no final. Cada um no seu quadrado.
Import (o kit básico do herói)
import {
Flow,
FlowContext,
FatalIssue,
BusinessRuleIssue,
FieldRequiredIssue,
Severity,
Issue,
buildFailureReport,
collectAllIssues,
collectTrail,
} from "@aylonmuramatsu/flow-pipes";Issues, validations e rules (o kit completo)
Aqui mora o coração da lib: não jogue exception pra tudo. Modele o problema, registre, continue (ou pare) com intenção.
O caminho feliz: objeto no ctx.issue (sem classe no seu app)
A classe é obrigação da lib. Você passa o descritor:
ctx.issue({
type: "FieldRequired", // ou omitir (= BusinessRule); custom: "MinAge", "InsufficientStock"
message: "Email não preenchido.",
code: "FIELD_REQUIRED",
metadata: { path: "email" },
});
ctx.issue({
type: "Fatal",
message: "Lock negado.",
code: "JOB_LOCKED",
metadata: { jobId },
});
// filtrar sem instanceof:
ctx.hasCode("FIELD_REQUIRED");
ctx.findByCode("MIN_AGE");
ctx.findByType("MinAge");
// built-ins ainda batem com classe (se precisar):
ctx.has(FieldRequiredIssue);new FieldRequiredIssue(...) continua válido — opcional, avançado.
Anatomia de uma Issue
Toda Issue:
- estende a classe abstrata
Issue - declara um
severity(info|warning|error|fatal) - passa
message+metadatanosuper - ganha de brinde:
toJSON(),createdAt,isFatal(),code?,name= nome da classe
import { Issue, Severity } from "@aylonmuramatsu/flow-pipes";
export class MinAgeIssue extends Issue {
public readonly severity = Severity.Error;
constructor(age: number, line: number) {
super("Cliente menor de idade.", { code: "MIN_AGE", age, line, minAge: 18 });
}
}metadata é o seu JSON de bolso pro front, log e suporte. Coloque path, field, sku, code — o que for útil depois.
Se metadata.code for string, vira também issue.code (e aparece no toJSON()).
Severidade: o que cada uma quer dizer
| Severity | Clima | Flow continua? |
| --------- | ----------------- | ------------------------------------------------------ |
| info | “Só avisando” | Sim |
| warning | “Hmm, olha isso” | Sim |
| error | “Regra quebrada” | Sim (até você dar stop()) |
| fatal | “Sem recuperação” | Não — ctx.issue(fatal) já chama stop() sozinho |
Duas alavancas, dois papéis:
severity→ relatório, filtro, HTTP mappingctx.stop({ reason, code })→ “não rode o próximo step” +stopInfotipável
FatalIssue / severity === fatal puxa o freio automático (com code da Issue). Nas outras, você decide se só anota ou se encerra.
Issues padrão (tempero de cozinha)
| Issue | Quando usar |
| -------------------- | --------------------------------------------------------------------------- |
| FieldRequiredIssue | Campo obrigatório / não preenchido — message + metadata livre (tipável) |
| BusinessRuleIssue | Regra genérica (“menor de idade”, “CPF duplicado”) — também tipável |
| FatalIssue | Lock negado, estado impossível, “não tem como seguir” — também tipável |
new FieldRequiredIssue("Email não preenchido.", {
code: "FIELD_REQUIRED",
path: "email",
});
// Excel ainda funciona — você escolhe o shape:
new FieldRequiredIssue("O campo 'cpf' é obrigatório.", {
code: "FIELD_REQUIRED",
line: 3,
field: "cpf",
});
// Soft em import parcial:
new FieldRequiredIssue("Coluna opcional vazia.", { field: "notes" }, Severity.Warning);Ótimos pra começar. Ruins pra virar deus de mil nomes — quando a regra tem identidade (estoque, gateway, KYC), crie a classe dela.
Criando Issues custom (receitas do mundo real)
import { Issue, Severity } from "@aylonmuramatsu/flow-pipes";
/** Validação de formato / domínio tipado */
export class InvalidCpfIssue extends Issue {
public readonly severity = Severity.Error;
constructor(cpf: string, line: number) {
super("CPF inválido.", { cpf, line, code: "INVALID_CPF" });
}
}
/** Regra de negócio com contexto rico */
export class InsufficientStockIssue extends Issue {
public readonly severity = Severity.Error;
constructor(sku: string, requested: number, available: number) {
super(`Estoque insuficiente para ${sku}.`, {
sku,
requested,
available,
code: "INSUFFICIENT_STOCK",
});
}
}
/** Aviso que não bloqueia (relatório / compliance) */
export class SuspiciousAmountIssue extends Issue {
public readonly severity = Severity.Warning;
constructor(amount: number) {
super("Valor fora do padrão histórico.", {
amount,
code: "SUSPICIOUS_AMOUNT",
});
}
}
/** Integração externa sem volta */
export class PaymentGatewayIssue extends Issue {
public readonly severity = Severity.Fatal;
constructor(gatewayCode: string, raw?: unknown) {
super(`Gateway recusou (${gatewayCode}).`, {
gatewayCode,
raw,
code: "PAYMENT_GATEWAY",
});
}
}
/** Política / autorização */
export class ForbiddenOperationIssue extends Issue {
public readonly severity = Severity.Fatal;
constructor(action: string, userId: string) {
super(`Operação não permitida: ${action}.`, {
action,
userId,
code: "FORBIDDEN",
});
}
}
export class DuplicateCpfIssue extends Issue {
public readonly severity = Severity.Error;
constructor(cpf: string, line: number) {
super(`CPF duplicado: ${cpf}`, { cpf, line, code: "DUPLICATE_CPF" });
}
}Dica de ouro: um campo code estável no metadata faz o front/i18n/monitoramento felizes. A message pode ser humana; o code é contrato.
Padrões de validation (steps que só julgam)
Validação boa é chata e previsível: lê input / shared, emite Issue, às vezes marca metadata, às vezes dá stop(). Não grava no banco. Não chama gateway. Só julga, igual a vizinha na janela! brincadeirinha! haha.
1) Acumular tudo (ótimo pra formulário / Excel)
Não para no primeiro erro — junta a lista completa pra devolver de uma vez.
async function validateCustomer(ctx: FlowContext<CustomerInput, Shared>) {
const { name, email, document } = ctx.input;
if (!name)
ctx.issue(
new FieldRequiredIssue("Nome não preenchido.", {
code: "FIELD_REQUIRED",
path: "name",
}),
);
if (!email)
ctx.issue(
new FieldRequiredIssue("Email não preenchido.", {
code: "FIELD_REQUIRED",
path: "email",
}),
);
if (!document)
ctx.issue(
new FieldRequiredIssue("Documento não preenchido.", {
code: "FIELD_REQUIRED",
path: "document",
}),
);
if (document && !isValidCpf(document)) {
ctx.issue(new InvalidCpfIssue(document, 0));
}
if (ctx.issues.length > 0) {
ctx.metadata.set("invalid", true);
ctx.stop({
code: "VALIDATION_STOP",
reason: "validation_failed",
});
}
}2) Soft no .each (linha zoada não mata o lote)
async function validateRow(ctx: FlowContext<ExcelRow, ImportShared>) {
const row = ctx.input;
if (!row.name) {
ctx.issue(new FieldRequiredIssue("O campo 'name' é obrigatório.", { code: "FIELD_REQUIRED", line: row.line, field: "name" }));
}
if (row.age != null && row.age < 18) {
ctx.issue(new MinAgeIssue(row.age, row.line));
}
if (row.cpf && ctx.shared.seenCpfs.has(row.cpf)) {
ctx.issue(new DuplicateCpfIssue(row.cpf, row.line));
} else if (row.cpf) {
ctx.shared.seenCpfs.add(row.cpf);
}
if (ctx.issues.length > 0) {
ctx.metadata.set("invalid", true);
// sem stop no filho → próximas linhas seguem
}
}
async function stageIfValid(ctx: FlowContext<ExcelRow, ImportShared>) {
if (ctx.metadata.get("invalid")) return;
// staging / insert “provisório”
}3) Fail-fast quando o resto não faz sentido
async function ensureJobLock(ctx) {
if (!(await tryLock(ctx.input.jobId))) {
ctx.issue(new FatalIssue("job already locked", { jobId: ctx.input.jobId }));
// FatalIssue já dá stop() — próximos steps nem acordam
}
}4) Pipeline de validation (vários steps, uma responsabilidade cada)
Flow.create<Shared>()
.shared({ ... })
.step(validateRequiredFields) // presence
.step(validateFormats) // CPF, email, UUID
.step(validateBusinessRules) // idade, estoque, status
.step(validatePermissions) // authz
.step(persist)
.run(input);Cada step pequeno = testável, reutilizável, sem novela de 200 linhas.
Padrões de rules (regras de negócio)
Rule ≠ “if solto”. Rule é decisão de domínio com Issue nomeada (e, se precisar, efeito no shared).
async function applyDiscountRules(ctx: FlowContext<CartInput, CheckoutShared>) {
const { coupon, total } = ctx.shared;
if (coupon && coupon.expired) {
ctx.issue(
new BusinessRuleIssue("Cupom expirado.", { code: "COUPON_EXPIRED" }),
);
ctx.shared.coupon = null; // regra também limpa estado
return;
}
if (coupon && total < coupon.minAmount) {
ctx.issue(
new BusinessRuleIssue("Cupom exige valor mínimo.", {
code: "COUPON_MIN_AMOUNT",
minAmount: coupon.minAmount,
total,
}),
);
return;
}
if (total > 10_000) {
ctx.issue(new SuspiciousAmountIssue(total)); // warning: segue, mas fica no radar
}
}
async function reserveStock(ctx: FlowContext<OrderInput, OrderShared>) {
for (const item of ctx.shared.items) {
const available = await stock.of(item.sku);
if (available < item.qty) {
ctx.issue(new InsufficientStockIssue(item.sku, item.qty, available));
}
}
if (ctx.has(InsufficientStockIssue)) {
ctx.stop(); // não cobra, não emite NF
}
}Consultando Issues depois do run (o buffet)
const ctx = await flow.run(input);
ctx.has(FieldRequiredIssue); // boolean
ctx.find(FieldRequiredIssue); // lista tipada
ctx.first(BusinessRuleIssue); // a primeira
ctx.count(InvalidCpfIssue); // quantas
ctx.issues; // todas deste contexto
collectAllIssues(ctx); // pai + children (`.each`)
// HTTP / API
if (ctx.has(FieldRequiredIssue) || ctx.metadata.get("invalid")) {
return {
status: 400,
body: { issues: collectAllIssues(ctx).map((i) => i.toJSON()) },
};
}
if (ctx.stopped) {
return {
status: 422,
body: {
failure: ctx.result?.failure ?? buildFailureReport(ctx),
issues: collectAllIssues(ctx).map((i) => i.toJSON()),
},
};
}toJSON() já serializa type, severity, message, metadata, createdAt — pronto pra log e response.
Monte o Flow como “casos de uso”
Três jeitos que funcionam bem em time real:
// 1) Factory da ação (reuso entre HTTP, fila, cron)
export function importCustomersFlow(db: Db) {
return Flow.create<ImportShared>()
.shared({ db, rows: [], seenCpfs: new Set(), staged: [] })
.step(readFile)
.each((c) => c.shared.rows, [validateRow, applyRowRules, stageRow])
.step(decideCommitOrStop);
}
// 2) Validations puras + rules + side-effects separados
Flow.create()
.step(validateInput) // presence + format
.step(applyDomainRules) // estoque, cupom, KYC
.step(chargePayment) // I/O
.step(persistOrder);
// 3) Soft no lote + hard no final
Flow.create()
.each((c) => c.shared.rows, [validateRow, stageRow])
.step(async (ctx) => {
if (collectAllIssues(ctx).some((i) => i.severity === "error")) {
ctx.shared.errors = collectAllIssues(ctx).map((i) => i.toJSON());
ctx.stop(); // chamador: rollback
}
});Mini-mapa: quando usar o quê
| Situação | O que fazer |
| --------------------------------- | -------------------------------------------------------- |
| Campo vazio | FieldRequiredIssue (ou custom de campo) |
| Formato inválido | Issue tipada (InvalidCpfIssue, …) |
| Regra de domínio | BusinessRuleIssue ou Issue dedicada |
| Dá pra seguir, mas quero rastrear | Severity.Warning + sem stop |
| Não faz sentido continuar | FatalIssue ou Issue + ctx.stop() |
| Lote: uma linha ruim | Issue no filho, metadata.invalid, sem stop no each |
| Lote: no fim cancela tudo | collectAllIssues + stop no step final |
| Lib externa estourou | deixa throwar — runner vira Fatal + result.failure |
| Resposta HTTP 400 rica | find / collectAllIssues + toJSON() |
Exemplos (cole, adapte, seja feliz)
Isso aqui é referência, não um reality show com botão “Run”. Copia pro seu projeto e faz a mágica.
1. Flow básico — “oi, mundo”, mas com atitude
interface Shared {
messages: string[];
}
async function greet(ctx: FlowContext<{ name: string }, Shared>) {
if (!ctx.input.name) {
ctx.issue(new FieldRequiredIssue("Nome não preenchido.", { code: "FIELD_REQUIRED", path: "name" }));
ctx.stop();
return;
}
ctx.shared.messages.push(`Olá, ${ctx.input.name}`);
}
const ctx = await Flow.create<Shared>()
.shared({ messages: [] })
.step(greet)
.run({ name: "Ana" });
console.log(ctx.shared.messages);
console.log(ctx.result);2. Service / API — porque o HTTP também merece esteira
Lib externa estourou? O processo Node não cai junto. O runner vira FatalIssue e te devolve o contexto pra você responder com classe (ou com 503, sem julgamentos).
interface CreateOrderInput {
customerId: string;
amount: number;
}
interface OrderShared {
total: number;
orderId?: string;
}
async function validate(ctx: FlowContext<CreateOrderInput, OrderShared>) {
if (!ctx.input.customerId) {
ctx.issue(new FieldRequiredIssue("customerId obrigatório.", { code: "FIELD_REQUIRED", path: "customerId" }));
ctx.stop();
}
}
async function charge(ctx: FlowContext<CreateOrderInput, OrderShared>) {
if (ctx.input.amount < 0) {
throw new Error("Gateway indisponível");
}
ctx.shared.total = ctx.input.amount;
}
async function persist(ctx: FlowContext<CreateOrderInput, OrderShared>) {
ctx.shared.orderId = `ord_${Date.now()}`;
}
export class CreateOrderService {
async execute(input: CreateOrderInput) {
const ctx = await Flow.create<OrderShared>()
.shared({ total: 0 })
.step(validate)
.step(charge)
.step(persist)
.run(input);
if (ctx.has(FieldRequiredIssue)) {
return {
status: 400,
body: { issues: ctx.find(FieldRequiredIssue).map((i) => i.toJSON()) },
};
}
if (ctx.stopped || ctx.has(FatalIssue)) {
return {
status: 503,
body: {
error: ctx.result?.failure?.message ?? ctx.first(FatalIssue)?.message,
},
};
}
return {
status: 201,
body: { data: ctx.shared },
};
}
}3. Importação com relatório — o clássico “processa tudo, chora depois”
Valida o lote, faz staging, e só no fim: rollback + “olha a lista do que deu errado”.
interface Row {
line: number;
name: string;
cpf: string;
}
interface Shared {
rows: Row[];
staged: Row[];
errors: unknown[];
}
async function loadRows(ctx) {
ctx.shared.rows = ctx.input.rows;
}
async function validateRow(ctx) {
const row = ctx.input as Row;
if (!row.name) {
ctx.issue(new FieldRequiredIssue("O campo 'name' é obrigatório.", { code: "FIELD_REQUIRED", line: row.line, field: "name" }));
}
if (!row.cpf) {
ctx.issue(new FieldRequiredIssue("O campo 'cpf' é obrigatório.", { code: "FIELD_REQUIRED", line: row.line, field: "cpf" }));
}
}
async function stageIfValid(ctx) {
if (ctx.issues.length) return;
ctx.shared.staged.push(ctx.input as Row);
}
async function decide(ctx) {
const issues = collectAllIssues(ctx);
if (issues.length > 0) {
ctx.shared.errors = issues.map((i) => i.toJSON());
ctx.stop();
}
}
const ctx = await Flow.create<Shared>()
.shared({ rows: [], staged: [], errors: [] })
.step(loadRows)
.each((c) => c.shared.rows, [validateRow, stageIfValid])
.step(decide)
.run({
rows: [
/* ... */
],
});
if (ctx.stopped) {
// rollback + manda ctx.shared.errors pro cliente
} else {
// commit de ctx.shared.staged e vai tomar café
}4. Lote com .each — um por um, sem drama coletivo
Cada item ganha um contexto filho (ctx.children) dividindo o mesmo shared. Uma linha zoada não precisa cancelar o passeio inteiro.
interface Row {
id: string;
name: string;
}
interface Shared {
rows: Row[];
saved: string[];
}
async function saveRow(ctx: FlowContext<Row, Shared>) {
if (!ctx.input.name) {
ctx.issue(new FieldRequiredIssue("Nome não preenchido.", { code: "FIELD_REQUIRED", path: "name" }));
return; // só essa linha fica de fora
}
ctx.shared.saved.push(ctx.input.id);
}
const ctx = await Flow.create<Shared>()
.shared({
rows: [
{ id: "1", name: "Ana" },
{ id: "2", name: "" },
{ id: "3", name: "Bruno" },
],
saved: [],
})
.each((c) => c.shared.rows, [saveRow], {
name: "persistRows",
concurrency: 4, // opcional; default 1 (modo zen)
})
.run();
console.log(ctx.shared.saved);
console.log(ctx.children.map((c) => c.issues));5. .parallel — porque esperar em fila é coisa do século passado
Dica de ouro: chaves distintas no shared (profile, wallet, score). Senão vira briga de vizinho no mesmo endereço.
interface Shared {
profile?: { name: string };
wallet?: { balance: number };
score?: { value: number };
}
async function fetchProfile(ctx: FlowContext<{ userId: string }, Shared>) {
ctx.shared.profile = { name: `User ${ctx.input.userId}` };
}
async function fetchWallet(ctx: FlowContext<{ userId: string }, Shared>) {
ctx.shared.wallet = { balance: 1500 };
}
async function fetchScore(ctx: FlowContext<{ userId: string }, Shared>) {
ctx.shared.score = { value: 820 };
}
const ctx = await Flow.create<Shared>()
.shared({})
.parallel([fetchProfile, fetchWallet, fetchScore], { name: "enrich" })
.run({ userId: "u-42" });
console.log(ctx.shared);6. Flow.runAll — vários flows, um só “vai”
const userFlow = Flow.create().step(async (ctx) => {
/* ... */
});
const billingFlow = Flow.create().step(async (ctx) => {
/* ... */
});
const results = await Flow.runAll([
{ name: "UserService", run: () => userFlow.run(input) },
{ name: "BillingService", run: () => billingFlow.run(input) },
]);
// default "settle" — um tropeço não cancela o resto do show
// modo dramático (fail-fast):
await Flow.runAll(tasks, { mode: "all" });7. Ação reutilizável — escreve uma vez, usa em todo canto
export function createCustomerFlow() {
return Flow.create<CustomerShared>()
.shared({ user: null, wallet: null })
.step(validatePayload)
.step(insertUser)
.step(createWallet)
.step(linkAddress);
}
// HTTP
const ctx = await createCustomerFlow().run(req.body);
if (ctx.stopped) return res.status(422).json(ctx.result);
// fila / outro contexto — mesma ação, zero copy-paste emocionado
await createCustomerFlow().run(payloadFromQueue);8. Soft vs fatal — você no volante, não o destino
async function lockJob(ctx) {
if (!(await tryLock(ctx.input.jobId))) {
ctx.issue(new FatalIssue("job already locked"));
ctx.stop();
}
}
async function softValidate(ctx) {
if (ctx.input.warning) {
ctx.issue(new BusinessRuleIssue("linha suspeita")); // segue andando, de olho aberto
}
}9. Quando o inesperado aparece (e você quer saber onde)
const ctx = await flow.run(input);
if (ctx.stopped) {
const report = ctx.result?.failure ?? buildFailureReport(ctx);
// report.location → step / eachIndex
// report.error.stack → o vilão
// report.trail → o que já tinha dado certo antes do plot twist
console.error(report);
}Controle do runner (1.1.0)
Flow.create<Shared>()
.shared({ rows: [], invalid: false })
.step(loadRows)
.each(
(c) => c.shared.rows,
[
validateRow,
{ run: stageRow, when: (c) => !c.metadata.get("invalid") },
],
{ fatalScope: "item" }, // FatalIssue/throw só mata o item
)
.when((c) => !c.stopped, decideCommit)
.step(callGateway, {
name: "callGateway",
retry: { attempts: 3, delayMs: 100, backoff: "exponential" },
})
.run({ file: "clientes.xlsx" }, { timeoutMs: 30_000, traceId: "req-42" });
// ctx.traceId, ctx.stopInfo, ctx.result.trail[].durationMsfatalScope: "flow" (default) propaga Fatal do item pro lote. "item" continua as próximas linhas.
API resumida (cola na parede do time)
Flow.create<TShared>()
.shared({ ... })
.step(fn, { name?, when?, retry? })
.step([a, b, { run, name?, when?, retry? }])
.when(pred, fn | [steps], { name?, retry? })
.each(resolver, steps, { name?, concurrency?, fatalScope?: "flow" | "item" })
.parallel(steps, { name? })
.run(input?, { signal?, timeoutMs?, traceId? })
Flow.runAll([{ name, run }], { mode?: "settle" | "all" })
ctx.stop({ reason?, code? }) // → stopInfo
ctx.traceId
ctx.signal // AbortSignal do run
ctx.result.trail[].durationMs
issue.code // espelho de metadata.coderetry: só em throw (não em Issue). Ex.: { attempts: 3, delayMs: 50, backoff: "exponential" }.
Nome do step: options.name → fn.name → arquivo → "anonymous".
"anonymous" é o vilão silencioso do log. Dê nome às suas funções — elas merecem.
Ideias de onde plugar isso
| Cenário | Por que encaixa (além do feeling) |
| ------------------------ | ------------------------------------------------------- |
| Import Excel / CSV / XML | Todas as linhas, Issues, rollback + relatório bonitinho |
| Cadastro multi-etapa | Uma ação; falhou no meio → stop + rollback no chamador |
| Jobs / filas | Lock, soft fail vs fatal, status “quase” |
| Enriquecer entidade | .parallel em APIs externas sem drama |
| Sync em lote | .each + concurrency sem inventar roda |
| Orquestrar services | O Flow vira o caso de uso tipado |
| API com erro rico | collectAllIssues / buildFailureReport |
| Validação em pipeline | Steps pequenos no lugar de ifs nômades |
O que isso não é (combinado?)
- Não aposenta
try/catchde infraestrutura. Rede caiu? Bug esquisito? Exception ainda existe — o runner localiza e joga emctx.result. - Não é BPMN nem aquele workflow eterno que vive num banco esquecido. É esteira em processo: leve, tipada, no seu service/job/script.
- Não esconde o domínio: negócio no
shared, fofoca do motor noresult.
Licença
MIT — use, abuse (com carinho) e construa fluxos que você ainda entenda daqui a seis meses.
Divirta-se na esteira. 🚂
