Workflows

Trabajo en segundo plano que sobrevive a los reinicios. Un workflow es una función cuyos pasos se guardan a medida que se completan: una ejecución que falla se reintenta con espera creciente, una interrumpida por un despliegue continúa desde su último paso terminado, y todas aparecen con su entrada, salida y error en Workflows, en el panel de PocketBase.

Los workflows se ejecutan en el motor OpenWorkflow integrado en el PocketBase de Vela, usando directamente la biblioteca openworkflow. Todo proyecto creado con Vela los incluye; un proyecto creado antes de que vinieran de serie los obtiene con vela enable workflows.

Sintaxis

$ vela workflows <command>
  • list - Lista las ejecuciones recientes y su estado
  • run - Inicia una ejecución desde el terminal
  • cancel - Cancela una ejecución pendiente o en curso

Cómo encaja

src/lib/server/workflows.ts contiene el cliente y el worker. Exporta ow, el cliente de OpenWorkflow con el que se definen los workflows, getAdmin(), un cliente superusuario de PocketBase para usar dentro de los pasos, y startWorker(), al que el hook init de src/hooks.server.ts llama una vez al arrancar el servidor.

El worker corre dentro del servidor web, tanto en desarrollo como en producción, así que no hay un segundo proceso que ejecutar ni desplegar. Reclama ejecuciones de PocketBase y las procesa en el mismo proceso; varios servidores sobre una misma base de datos se reparten el trabajo sin conflicto, porque cada ejecución se presta a un solo worker a la vez. Un proyecto sin archivos de workflow nunca arranca un worker.

Definir un workflow

Los workflows viven en src/lib/workflows/, un archivo por workflow, y vela generate workflow escribe uno:

$ vela generate workflow send-welcome-email
import { z } from 'zod';
import { ow } from '$lib/server/workflows';

export const sendWelcomeEmail = ow.defineWorkflow(
	{
		name: 'send-welcome-email',
		schema: z.object({ userId: z.string() }),
		retryPolicy: { maximumAttempts: 3 }
	},
	async ({ input, step }) => {
		const user = await step.run({ name: 'load-user' }, async () => {
			const admin = await getAdmin();
			return admin.collection('users').getOne(input.userId);
		});

		await step.run({ name: 'send' }, async () => {
			// ...
		});
	}
);

schema es lo que acepta run(), comprobado antes de encolar la ejecución. Cada step.run guarda su valor de retorno, de modo que un reintento o un reinicio continúa tras el último paso terminado en lugar de empezar de cero. Eso convierte a los pasos en la unidad de reintento: haz que cada uno sea seguro de repetir. El valor que devuelve un paso tiene que ser JSON.

Un workflow tiene un solo intento salvo que retryPolicy diga otra cosa. step.sleep, step.waitForSignal y los workflows hijos mediante step.runWorkflow son la API de la propia biblioteca y están documentados en openworkflow.dev.

Iniciar una ejecución

Importa el workflow en cualquier punto del servidor — una acción de formulario, una ruta de API, un hook, otro workflow — y llama a run():

import { sendWelcomeEmail } from '$lib/workflows/send-welcome-email';

await sendWelcomeEmail.run({ userId: user.id }, { idempotencyKey: user.id });

run() devuelve en cuanto la ejecución queda encolada. Su handle tiene result() para esperar la salida y cancel(); esperar dentro de un manejador de petición rara vez es lo que quieres, porque la idea es sacar el trabajo del camino de la petición.

  • idempotencyKey - Llamar a run() de nuevo con la misma clave en un plazo de 24 horas devuelve la ejecución existente en lugar de iniciar otra. Usa como clave aquello de lo que trata la ejecución — un id de usuario, un id de evento de Stripe — y una petición reintentada o un webhook reenviado no podrán hacer el trabajo dos veces.
  • availableAt - Un Date o una duración como '10m'; la ejecución espera hasta entonces.
  • deadlineAt - Un Date tras el cual la ejecución se marca como fallida en vez de reintentarse.

Workflows recurrentes

Un archivo que exporta una expresión cron ejecuta todos los workflows que define con ese horario:

$ vela generate workflow sync-prices --cron '*/5 * * * *'
export const cron = '*/5 * * * *';

El horario se dispara como mucho una vez por minuto, y una sola vez entre todos los servidores: cada disparo se identifica por el nombre del workflow y el minuto, así que por muchos servidores que haya en marcha, se inicia una única ejecución. Un minuto que pasa sin nada en marcha se salta, no se recupera. En las pruebas el horario está desactivado; tickCron() de $lib/server/workflows inicia a demanda todos los workflows recurrentes.

Pruebas

El generador escribe un <name>.server.test.ts junto a cada workflow, que ejecuta vela test:server. El proceso de pruebas arranca su propio worker, así que una ejecución iniciada en una prueba se completa sin pasar por el servidor de desarrollo:

import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import { startWorker, stopWorker } from '$lib/server/workflows';
import { sendWelcomeEmail } from './send-welcome-email';

describe('send-welcome-email', () => {
	beforeAll(() => startWorker());
	afterAll(() => stopWorker());

	it('completes', async (context) => {
		const handle = await sendWelcomeEmail.run({ userId: context.user.id });
		await expect(handle.result({ timeoutMs: 15_000 })).resolves.toBeNull();
	});
});

Ajustes

Dos variables de entorno regulan el worker; ninguna es obligatoria:

WORKFLOWS_CONCURRENCY=5
WORKFLOWS_ENABLED=true

WORKFLOWS_CONCURRENCY limita cuántas ejecuciones procesa un servidor a la vez. WORKFLOWS_ENABLED=false impide que un servidor procese ejecuciones sin dejar de encolarlas — el interruptor para mover el trabajo a un proceso dedicado más adelante.

En producción

El worker vive dentro del servicio vela-web, así que un despliegue o un reinicio lo detiene junto con el servidor web: termina las ejecuciones que tiene en curso, y una que no llegue a terminar a tiempo se retoma cuando expira su préstamo, desde el último paso completado. Las ejecuciones y sus pasos se guardan en openworkflow.db, junto a la base de datos de la aplicación, y forman parte de toda copia de seguridad.

Listar

Lista las ejecuciones recientes y su estado.

Ejecutar

Inicia una ejecución desde el terminal.

Cancelar

Cancela una ejecución pendiente o en curso.