Fronta a worker: webhook odpovi hned, praci udelaji workeri

Webhook uz nic nevykonava v requestu. Zapise udalost do fronty a odpovi 202
do jednotek milisekund; strom vykona worker na pozadi. Za konektory nerucime,
takze cekat na cizi sluzbu v requestu znamena ztracet udalosti pri timeoutu.

Fronta ma opakovani s rostouci prodlevou (30 s, 2 min, 10 min, hodina),
spravedlive poradi po firmach (jedna firma s tisicem udalosti nezablokuje
ostatni), navrat zaseknutych behu po restartu a uklid hotovych. Marna chyba
se neopakuje - chybejici skript za minutu existovat nezacne.

Tri druhy spoustecu: push (webhook), vnitrni udalost (vznik a zmena ticketu)
a pull, tedy pravidelne dotazovani u sluzeb bez webhooku (posta, zpravy).
Planovac jen rekne "je cas", samotny dotaz je prvni krok stromu, takze ma
zaznam v logu a opakuje se pri chybe jako cokoliv jineho.

Kontrakt tela webhooku: kazdy parametr ma cestu (data.order.id,
errors.0.message), takze jde napojit i odesilatel s vnorenym modelem.
U adresy je metoda, ukazka tela a kopiruje se cela adresa vcetne domeny.

Vnitrni kroky, ktere sahaji do naseho uloziste: ticket/upsert (zaloz nebo
dopln podle externiho ID), assign-least-busy, assign-by-external, set-type,
set-stage, add-tags, set-status, incident/create, flow/pause a flow/log.

Faze ticketu jako treti osa vedle stavu a stitku. Stav je zivotni cyklus
a pocitaji se z nej statistiky, faze je workflow daneho typu a muze byt jen
jedna, takze se na ni da spolehnout v podmince.

ID z cizich aplikaci u resitele: voicebot posle voicebotId a ticket skonci
u toho, komu patri. Vazba je na jednom miste, ne v kazde automatizaci.

Kazda chyba zaklada incident se dvema urovnemi: impact cte klient a je
srozumitelny, detail cte admin a je v nem cely beh, ktery krok selhal, cele
hlaseni a data na vstupu. Detail vidi jen spravce platformy.

Ochrana proti smycce: automatizace navazana na zmenu ticketu ticket meni,
cimz se spousti znovu - pri vyvoji to server polozilo. Resi to oznaceni behu
pres AsyncLocalStorage a strop peti behu na jeden ticket za minutu.

Upozorneni pri prideleni prace vcetne cisla u zalozky Tickety. Zivy dashboard:
dlazdice nad nasimi daty na udalost, data z konektoru podle ttlSec s moznosti
vynutit nacteni znovu.

Opraveno: path a intervalSec u spoustece se pri ulozeni zahazovaly; nad
seznamem neslo pouzit contains, takze na stitky neslo postavit podminku;
novejsi vystup kroku ted prekryje starsi misto hlaseni konfliktu.

Overeno dvema scenari proti bezicimu serveru, 34 kontrol: firma se skladem,
expedici a IT, a hovory z voicebota (callSid do externiho ID, status do faze,
prirazeni podle voicebotId, tri zpravy = jeden ticket se tremi udalostmi).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
JiriUhlir
2026-08-13 16:41:02 +02:00
co-authored by Claude Opus 5
parent 5d186dcd2e
commit a57eca123e
38 changed files with 2942 additions and 112 deletions
+101 -2
View File
@@ -44,6 +44,8 @@ import { listIncidents } from '../data/incidentStore.js';
import { getSummary } from '../data/mock.js';
import { findPersonByEmail, listGroups, listPeople } from '../data/people.js';
import { hasPermission } from '../data/permissions.js';
import { listNotifications, markRead, unreadCount } from '../data/notifications.js';
import { queueStats, recentRuns } from '../runtime/queue.js';
import { recordAudit } from '../data/audit.js';
import { findTenant, generateIntakeToken, refreshTenants, tenantStore } from '../data/tenants.js';
import { listTicketTypes } from '../data/ticketTypes.js';
@@ -111,8 +113,24 @@ dashboardRouter.get('/storage', (_req, res) => {
res.json(storageStatus());
});
dashboardRouter.get('/incidents', (_req, res) => {
res.json({ items: listIncidents() });
/**
* Incidenty firmy plus platformni.
*
* Klient vidi `title` a `impact`, tedy co to pro nej znamena. `detail` s celym
* hlasenim, ID behu a daty na vstupu vidi **jen spravce platformy** - je to
* nase diagnostika, ne informace pro zakaznika.
*/
dashboardRouter.get('/incidents', (req, res) => {
const scope = scopeOrDeny(req, res);
if (!scope) return;
const forAdmin = req.user!.platformAdmin;
return res.json({
items: listIncidents(scope.tenantIds).map((incident) => ({
...incident,
detail: forAdmin ? incident.detail : null,
})),
});
});
// ------------------------------------------------------- rozlozeni dashboardu
@@ -339,6 +357,80 @@ dashboardRouter.post('/intake/regenerate', (req, res) => {
});
});
// ------------------------------------------------------------- upozorneni
/**
* Upozorneni prihlaseneho.
*
* Cislo u zalozky "Moje tickety" a hlaska pri prirazeni. Nechodi to pres
* stream jako jedina cesta - kdo mel portal zavreny, musi to najit i po
* prihlaseni, proto jsou upozorneni ulozena.
*/
dashboardRouter.get('/notifications', (req, res) => {
const items = listNotifications(req.user!.id);
return res.json({
items,
unread: unreadCount(req.user!.id),
/** Kolik ticketu ma prihlaseny u sebe. To je to cislo u zalozky. */
mine: myOpenTickets(req),
});
});
/** Oznaci prectene. Bez seznamu vsechny. */
dashboardRouter.post('/notifications/read', (req, res) => {
const ids = Array.isArray(req.body?.ids)
? (req.body.ids as unknown[]).filter((id): id is string => typeof id === 'string')
: undefined;
return void markRead(req.user!.id, ids)
.then((count) => res.json({ marked: count, unread: unreadCount(req.user!.id) }))
.catch((err: unknown) => {
console.error('[upozorneni] oznaceni selhalo:', err);
return res.status(500).json({ error: 'internal_error', message: 'Nepodařilo se uložit.' });
});
});
/** Kolik nevyrizenych ma prihlaseny u sebe. */
function myOpenTickets(req: Request): number {
const access = accessFor(req.user!);
if (!access.personId) return 0;
return listTickets({
tenantIds: access.tenants.map((tenant) => tenant.id),
assignee: access.personId,
}).filter((ticket) => ticket.status !== 'resolved').length;
}
// ------------------------------------------------------------------- fronta
/**
* Stav fronty behu.
*
* Kdyz neco nefunguje, tohle je prvni misto, kam se clovek podiva: ceka
* fronta, nebo uz to nekolikrat selhalo a vzdalo se?
*/
dashboardRouter.get('/runs', (req, res) => {
const scope = scopeOrDeny(req, res);
if (!scope) return;
return res.json({
stats: queueStats(scope.tenantIds),
items: recentRuns(scope.tenantIds).map((run) => ({
id: run.id,
automationId: run.automationId,
trigger: run.trigger,
status: run.status,
attempts: run.attempts,
ticketId: run.ticketId,
// Cele hlaseni. Zkratit ho tady znamena, ze se pricina uz nedozvime.
lastError: run.lastError,
nextAttemptAt: run.nextAttemptAt,
createdAt: run.createdAt,
finishedAt: run.finishedAt,
})),
});
});
// ------------------------------------------------------------------- tickety
/**
@@ -636,6 +728,11 @@ const fieldSchema = z.object({
// dosadit nejde - predava se jen jako celek dalsimu kroku.
type: z.enum(['string', 'number', 'boolean', 'date', 'object', 'list']),
required: z.boolean(),
/**
* Kde ta hodnota v prichozim tele je, kdyz to neni primo `name`.
* Bez toho by slo napojit jen ploche telo, viz `TriggerField.path`.
*/
path: z.string().trim().max(200).optional(),
});
const flowSchema = z.object({
@@ -644,6 +741,8 @@ const flowSchema = z.object({
serviceId: z.string().min(1),
operationId: z.string().min(1),
fields: z.array(fieldSchema),
/** Jak casto se ma sluzba obvolavat. Plati jen u spoustecu, co se ptaji. */
intervalSec: z.number().int().min(10).max(86_400).optional(),
// Token generuje server. Cokoliv od klienta se ignoruje.
webhookToken: z.string().optional(),
})
+4
View File
@@ -278,6 +278,8 @@ const personCreate = z.object({
email: z.string().trim().email(),
role: z.string().trim().max(60).optional(),
capacity: z.number().int().min(1).max(200).optional(),
/** ID, pod kterymi cloveka znaji cizi aplikace, napr. voicebotId. */
externalIds: z.array(z.string().trim().min(1).max(120)).max(20).optional(),
});
settingsRouter.use(
@@ -292,6 +294,7 @@ settingsRouter.use(
role: z.string().trim().max(60).optional(),
capacity: z.number().int().min(1).max(200).optional(),
enabled: z.boolean().optional(),
externalIds: z.array(z.string().trim().min(1).max(120)).max(20).optional(),
}),
writePermission: 'people.manage',
build: (input) => ({
@@ -300,6 +303,7 @@ settingsRouter.use(
role: input.role ?? '',
capacity: input.capacity ?? 8,
enabled: true,
externalIds: input.externalIds ?? [],
}),
validate: (person, all) =>
all.some((other) => other.email.toLowerCase() === person.email.toLowerCase())
+128 -69
View File
@@ -1,12 +1,13 @@
import { Router } from 'express';
import { z } from 'zod';
import { findByWebhookToken, recordRun, type TriggerField } from '../data/automationStore.js';
import { findByWebhookToken, type TriggerField } from '../data/automationStore.js';
import type { FieldType } from '../data/conditions.js';
import { findByIntakeToken } from '../data/tenants.js';
import { runFlow } from '../runtime/executor.js';
import { intakeEvent } from '../data/ticketStore.js';
import { findTicketType, listTicketTypes } from '../data/ticketTypes.js';
import { publish } from '../events/bus.js';
import { getPath } from '../scripts/mapping.js';
import { enqueue } from '../runtime/queue.js';
export const webhookRouter = Router();
@@ -153,9 +154,13 @@ webhookRouter.get('/ticket/:token', (req, res) => {
});
/**
* VEREJNY endpoint - zamerne BEZ prihlaseni. Autorizaci resi neodhadnutelny
* token v adrese (32 znaku, base64url), presne tak, jak to dela vetsina
* webhookovych sluzeb.
* VEREJNY endpoint automatizace - zamerne BEZ prihlaseni. Autorizaci resi
* neodhadnutelny token v adrese (32 znaku, base64url), presne tak, jak to
* dela vetsina webhookovych sluzeb.
*
* **Odpovi hned.** Strom se nevykonava v requestu: cizi sluzba muze odpovidat
* pomalu nebo byt mimo a odesilateli by mezitim vyprsel timeout. Udalost se
* zapise do fronty a zpracuje ji worker, viz `runtime/queue.ts`.
*
* Neznamy token vraci 404 a nikdy neprozradi, ze nejaka automatizace existuje.
*/
@@ -176,8 +181,8 @@ webhookRouter.post('/:token', (req, res) => {
});
}
const payload = (req.body ?? {}) as Record<string, unknown>;
const problems = validatePayload(automation.flow.trigger.fields, payload);
const body = (req.body ?? {}) as Record<string, unknown>;
const { values, problems } = readPayload(automation.flow.trigger.fields, body);
if (problems.length > 0) {
console.warn(`[webhook] ${automation.id}: neplatna data - ${problems.join(' ')}`);
@@ -190,106 +195,160 @@ webhookRouter.post('/:token', (req, res) => {
publish('webhook.received', `Webhook přijal data pro ${automation.id}`, {
automationId: automation.id,
fields: Object.keys(payload),
fields: Object.keys(values),
});
console.info(
`[webhook] ${automation.id}: prijato (${Object.keys(payload).join(', ') || 'bez dat'})`,
);
if (automation.flow.steps.length === 0) {
recordRun(automation.id);
return res.status(202).json({
accepted: true,
automationId: automation.id,
message: 'Přijato. Automatizace nemá žádné kroky.',
});
}
/*
* Strom se vykona hned v requestu. Odesilatel tim dostane skutecny vysledek
* misto "prijato" - u desitek kroku je to v poradku, u tisicu udalosti za
* minutu to patri do fronty, viz documentation/10-runtime-a-kapacita.md.
*/
return void runFlow(automation.flow.steps, payload, {
return void enqueue({
tenantId: automation.tenantId,
ticketId: null,
idempotencyKey: `webhook:${automation.id}:${token.slice(0, 8)}`,
automationId: automation.id,
trigger: 'webhook',
/*
* Do stromu jdou **pojmenovane hodnoty** podle kontraktu, plus cele telo
* pod `_body`. Diky tomu se v sablonach pise `{{orderId}}` bez ohledu na
* to, jak hluboko to odesilatel schoval, a zaroven se nic neztrati.
*/
payload: { ...values, _body: body },
})
.then((run) => {
recordRun(automation.id, run.ok);
const failed = run.steps.find((step) => !step.ok);
if (!run.ok) {
console.warn(`[webhook] ${automation.id}: beh selhal - ${run.error}`);
}
// Chyba stromu neni chyba prijmu. Data jsme prevzali, jen se nepovedlo
// je zpracovat - proto 200 s `ok: false`, ne 500.
return res.status(200).json({
.then((item) => {
console.info(
`[webhook] ${automation.id}: prijato, zarazeno jako ${item?.id ?? '(duplicita)'}`,
);
// 202: prevzato, zpracuje se. Ne 200, ktera by rikala "hotovo".
return res.status(202).json({
accepted: true,
automationId: automation.id,
ok: run.ok,
steps: run.steps.length,
durationMs: run.durationMs,
error: run.error,
// Cele hlaseni toho kroku, ktery to zastavil.
detail: failed?.detail ?? null,
runId: item?.id ?? null,
message: 'Přijato, zpracuje se na pozadí.',
});
})
.catch((err: unknown) => {
// `runFlow` chyby nevyhazuje, tohle je posledni pojistka.
console.error(`[webhook] ${automation.id}: neocekavana chyba behu:`, err);
recordRun(automation.id, false);
return res.status(500).json({
error: 'internal_error',
message: err instanceof Error ? err.message : 'Běh selhal.',
console.error(`[webhook] ${automation.id}: zarazeni selhalo:`, err);
return res.status(503).json({
error: 'queue_unavailable',
message: 'Frontu se nepodařilo zapsat, zkuste to znovu.',
});
});
});
/** Overi prichozi data proti deklarovanym parametrum spoustece. */
function validatePayload(
/**
* Napoveda ke konkretni adrese webhooku.
*
* Rekne, co presne se v tele ceka. Bez toho by ten, kdo webhook zapojuje,
* musel hadat nebo se ptat - a hadat se bude spatne.
*/
webhookRouter.get('/:token', (req, res) => {
const automation = findByWebhookToken(req.params.token);
if (!automation?.flow.trigger) {
return res.status(404).json({ error: 'not_found', message: 'Webhook neexistuje.' });
}
const fields = automation.flow.trigger.fields;
return res.json({
method: 'POST',
contentType: 'application/json',
automation: automation.name,
enabled: automation.enabled,
/** Co se v tele ceka. `path` je misto v tele, kdyz to neni primo `name`. */
expects: fields.map((field) => ({
name: field.name,
path: field.path ?? field.name,
type: field.type,
required: field.required,
})),
example: exampleBody(fields),
note:
'Odpověď 202 znamená, že jsme data převzali. Strom se vykoná na pozadí, ' +
'takže tady výsledek nenajdete.',
});
});
/**
* Precte telo podle kontraktu.
*
* Kazde deklarovane pole se hleda na sve ceste. Chybejici povinne pole je
* chyba, chybejici nepovinne se preskoci. Prebytek v tele nikomu nevadi -
* odesilatel casto posila vic, nez potrebujeme, a odmitnout ho kvuli tomu
* by znamenalo, ze webhook nejde zapojit.
*/
function readPayload(
fields: TriggerField[],
payload: Record<string, unknown>,
): string[] {
body: Record<string, unknown>,
): { values: Record<string, unknown>; problems: string[] } {
const values: Record<string, unknown> = {};
const problems: string[] = [];
for (const field of fields) {
const value = payload[field.name];
const value = getPath(body, field.path ?? field.name);
if (value === undefined || value === null) {
if (field.required) problems.push(`Chybí povinný parametr „${field.name}".`);
if (field.required) {
const where = field.path && field.path !== field.name ? ` (cesta ${field.path})` : '';
problems.push(`Chybí povinný parametr „${field.name}"${where}.`);
}
continue;
}
if (!matchesType(value, field.type)) {
problems.push(`Parametr „${field.name}" má mít typ ${field.type}.`);
continue;
}
values[field.name] = value;
}
// Neznama pole jen zalogujeme - odmitat je by rozbilo odesilatele,
// kteri posilaji navic i sva vlastni data.
const declared = new Set(fields.map((f) => f.name));
const extra = Object.keys(payload).filter((key) => !declared.has(key));
if (extra.length > 0) {
console.info(`[webhook] nedeklarovane parametry navic: ${extra.join(', ')}`);
return { values, problems };
}
/** Ukazkove telo podle kontraktu, at je videt, jak to ma vypadat. */
function exampleBody(fields: TriggerField[]): Record<string, unknown> {
const example: Record<string, unknown> = {};
for (const field of fields) {
const value = exampleValue(field.type);
const path = (field.path ?? field.name).split('.');
let target = example;
for (let i = 0; i < path.length - 1; i += 1) {
const key = path[i];
if (typeof target[key] !== 'object' || target[key] === null) target[key] = {};
target = target[key] as Record<string, unknown>;
}
target[path[path.length - 1]] = value;
}
return problems;
return example;
}
function exampleValue(type: FieldType): unknown {
switch (type) {
case 'number':
return 42;
case 'boolean':
return true;
case 'date':
return '2026-01-31';
case 'object':
return { klic: 'hodnota' };
case 'list':
return ['první', 'druhá'];
default:
return 'text';
}
}
function matchesType(value: unknown, type: FieldType): boolean {
switch (type) {
case 'string':
return typeof value === 'string';
case 'number':
return typeof value === 'number' && Number.isFinite(value);
case 'boolean':
return typeof value === 'boolean';
case 'date':
return typeof value === 'string' && !Number.isNaN(new Date(value).getTime());
return typeof value === 'string' && !Number.isNaN(Date.parse(value));
case 'object':
return typeof value === 'object' && value !== null && !Array.isArray(value);
case 'list':
return Array.isArray(value);
default:
console.warn(`[webhook] neznamy typ parametru: ${type}`);
return false;
return typeof value === 'string';
}
}
+36 -8
View File
@@ -53,6 +53,13 @@ const statusLabels: Record<string, string> = {
const requestSchema = z.object({
widgetIds: z.array(z.string().min(1)).max(24),
/**
* true = natáhnout data z konektoru znovu, i kdyz jsou v mezipameti cerstva.
*
* Pouziva to tlacitko "obnovit" u dlazdice. Pravidelnou obnovu to obchazet
* nema - proto to neni vychozi chovani.
*/
force: z.boolean().default(false),
});
/**
@@ -88,6 +95,13 @@ type WidgetValue =
rows: Array<{ key: string; label: string; value: number }>;
at: string;
stale: boolean;
/**
* Kdy se data zase natáhnou. Klient podle toho pozna, jestli ma smysl
* cekat, nebo si vyzadat obnovu rucne.
*/
nextAt: string;
/** Jak dlouho se drzi, v sekundach. Nastavuje se u widgetu. */
ttlSec: number;
};
interface WidgetResult {
@@ -229,12 +243,19 @@ async function computeConnector(
widgetId: string,
source: Extract<WidgetSource, { kind: 'connector' }>,
scope: ResolvedScope,
force = false,
): Promise<WidgetValue> {
const key = cacheKey(widgetId, scope.tenantIds);
const ttl = Math.max(source.ttlSec ?? 300, 30) * 1_000;
const cached = externalCache.get(key);
if (cached && Date.now() - cached.at < ttl) {
return cached.value.kind === 'external' ? { ...cached.value, stale: true } : cached.value;
if (cached && Date.now() - cached.at < ttl && !force) {
return cached.value.kind === 'external'
? {
...cached.value,
stale: true,
nextAt: new Date(cached.at + ttl).toISOString(),
}
: cached.value;
}
const scriptId = scriptIdFor(source.serviceId, source.operationId);
@@ -259,7 +280,7 @@ async function computeConnector(
}
const picked = source.path ? getPath(result.outputs, source.path) : result.outputs;
const value = toExternal(picked, source);
const value = toExternal(picked, source, ttl);
rememberExternal(key, value);
return value;
}
@@ -268,11 +289,14 @@ async function computeConnector(
function toExternal(
picked: unknown,
source: Extract<WidgetSource, { kind: 'connector' }>,
ttl: number,
): WidgetValue {
const at = new Date().toISOString();
const now = Date.now();
const at = new Date(now).toISOString();
const nextAt = new Date(now + ttl).toISOString();
if (typeof picked === 'number') {
return { kind: 'external', value: picked, text: null, rows: [], at, stale: false };
return { kind: 'external', value: picked, text: null, rows: [], at, nextAt, ttlSec: ttl / 1_000, stale: false };
}
if (Array.isArray(picked)) {
@@ -286,11 +310,11 @@ function toExternal(
};
});
// Delka pole je casto to jedine cislo, ktere dava smysl.
return { kind: 'external', value: picked.length, text: null, rows, at, stale: false };
return { kind: 'external', value: picked.length, text: null, rows, at, nextAt, ttlSec: ttl / 1_000, stale: false };
}
if (picked === null || picked === undefined) {
return { kind: 'external', value: null, text: null, rows: [], at, stale: false };
return { kind: 'external', value: null, text: null, rows: [], at, nextAt, ttlSec: ttl / 1_000, stale: false };
}
if (typeof picked === 'object') {
@@ -304,6 +328,8 @@ function toExternal(
text: rows.length > 0 ? null : JSON.stringify(picked).slice(0, 200),
rows,
at,
nextAt,
ttlSec: ttl / 1_000,
stale: false,
};
}
@@ -316,6 +342,8 @@ function toExternal(
text,
rows: [],
at,
nextAt,
ttlSec: ttl / 1_000,
stale: false,
};
}
@@ -499,7 +527,7 @@ widgetDataRouter.post('/', async (req, res) => {
try {
const value =
widget.source.kind === 'connector'
? await computeConnector(widget.id, widget.source, scope)
? await computeConnector(widget.id, widget.source, scope, parsed.data.force)
: computeWidget(widget, scope);
return { id, ok: true, value };
} catch (err) {