npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

@aylonmuramatsu/flow-pipes

v1.1.0

Published

Motor de pipelines tipado com contexto compartilhado, Issues e execução paralela

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-pipes

Só 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)

  1. Monta a esteira
  2. Processa
  3. Olha Issues / stopped / result
  4. 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:

  1. estende a classe abstrata Issue
  2. declara um severity (info | warning | error | fatal)
  3. passa message + metadata no super
  4. 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 mapping
  • ctx.stop({ reason, code }) → “não rode o próximo step” + stopInfo tipá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[].durationMs

fatalScope: "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.code

retry: 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/catch de infraestrutura. Rede caiu? Bug esquisito? Exception ainda existe — o runner localiza e joga em ctx.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 no result.

Licença

MIT — use, abuse (com carinho) e construa fluxos que você ainda entenda daqui a seis meses.

Divirta-se na esteira. 🚂