diff --git a/documentation/01-prehled-a-stav.md b/documentation/01-prehled-a-stav.md index 2571663..6099b04 100644 --- a/documentation/01-prehled-a-stav.md +++ b/documentation/01-prehled-a-stav.md @@ -47,7 +47,11 @@ React aplikaci ze slozky `dist/public`. | Telo akce jako strom | hotovo | tentyz editor jako automatizace | | Audit a prepnuti na jiny ucet | hotovo | prepnuti je vychozi jen pro cteni, vse v auditu | | Bugs a wishes | chybi | vyvojarska agenda, samostatna evidence vedle ticketu | -| Beh automatizaci | castecne | strom se vykona, ale synchronne a bez fronty | +| Beh automatizaci | hotovo | fronta, worker, opakovani, ochrana proti smycce | +| Prijem udalosti do fronty | hotovo | webhook odpovi 202, praci dela worker | +| Pravidelne dotazovani sluzeb | hotovo | planovac pro postu a zpravy, perioda u spoustece | +| Upozorneni na pridelenou praci | hotovo | cislo u zalozky a hlaska v portalu | +| Incident z chyby | hotovo | popis pro klienta, podrobnosti pro admina | | Uloziste konektoru | hotovo | Postgres, nebo JSON soubor. Udaje vzdy sifrovane | | Uloziste pro zbytek | hotovo | tickety, automatizace, incidenty, rozlozeni, entity | | Monetizace a cena za krok | navrh | popis v 16-monetizace.md, neni naprogramovane | @@ -127,4 +131,5 @@ pro frontu, beh kroku a rozpocet na 150 klientu, a | [17-nastaveni-a-prava.md](17-nastaveni-a-prava.md) | prava, typy, akce, widgety, prepnuti uctu | | [18-ticketovaci-system.md](18-ticketovaci-system.md) | udalosti, externi ID, statistiky, pohledy | | [19-kapacita-200-firem.md](19-kapacita-200-firem.md) | zmereno, co zvladne soucasny stav | +| [20-fronta-a-runtime.md](20-fronta-a-runtime.md) | fronta, worker, spoustece, ochrana proti smycce | | [99-zmeny.md](99-zmeny.md) | zaznam zmen, nejnovejsi nahore | diff --git a/documentation/04-api.md b/documentation/04-api.md index 2600d76..011c5a9 100644 --- a/documentation/04-api.md +++ b/documentation/04-api.md @@ -36,6 +36,9 @@ Vyzaduji `Authorization: Bearer `: | GET | `/api/dashboard/intake` | | POST | `/api/dashboard/intake/regenerate` | | GET | `/api/dashboard/settings/actions/:id/scope` | +| GET | `/api/dashboard/notifications` | +| POST | `/api/dashboard/notifications/read` | +| GET | `/api/dashboard/runs` | | GET | `/api/dashboard/tickets` | | GET | `/api/dashboard/tickets/workload` | | GET | `/api/dashboard/tickets/:id` | @@ -227,6 +230,21 @@ tentyz klic. Neznamy `typeId` se zahodi a zaloguje, ticket vznikne bez typu. Odmitnout celou udalost kvuli jednomu poli by znamenalo ztratu dat. +## Fronta behu + +Popis je v [20-fronta-a-runtime.md](20-fronta-a-runtime.md). + +`POST /webhook/:token` vraci **202**, ne 200: data jsme prevzali a strom se +vykona na pozadi. Vysledek se hleda v `GET /api/dashboard/runs` nebo v logu +ticketu. Cekat na cizi sluzbu v requestu nejde - za jeji rychlost nerucime +a odesilateli by vyprsel timeout. + +`GET /webhook/:token` vraci **kontrakt**: co se v tele ceka, na jakych cestach +a ukazku. Bez toho by musel ten, kdo webhook zapojuje, hadat. + +`GET /api/dashboard/runs` ma u kazdeho behu cele chybove hlaseni, pocet pokusu +a kdy se to zkusi znovu. + ## Prava a navigace `GET /api/dashboard/access` vraci `permissions` (efektivni prava po slouceni diff --git a/documentation/15-rejstrik-funkci.md b/documentation/15-rejstrik-funkci.md index 895d52b..117a645 100644 --- a/documentation/15-rejstrik-funkci.md +++ b/documentation/15-rejstrik-funkci.md @@ -37,6 +37,13 @@ Volající nikdy nezjišťuje, jestli běží Postgres, soubor, nebo pamět. | `hasPermission(...)` | `src/data/permissions.ts` | Jedna kontrola. Používá ji `crudRouter` i ruční handlery. | | `navFor(...)` | `src/data/tenantFeatures.ts` | Průnik toho, co firma má, a toho, na co má člověk právo. Navigace chodí ze serveru. | | `recordAudit(input)` | `src/data/audit.ts` | Zápis do auditu. Nevrací chybu a nečeká se - rozbitý audit nesmí rozbít aplikaci. | +| `enqueue(input)` | `src/runtime/queue.ts` | Zařadí běh. Klíč proti dvojímu zařazení drží jeden běh na jednu událost. | +| `claimBatch(limit)` | `src/runtime/queue.ts` | Vezme další práci, spravedlivě po firmách. Místo, kde nad Postgresem musí být SKIP LOCKED. | +| `onTicketEvent(kind, ticket)` | `src/runtime/triggers.ts` | Změna ticketu zařadí navázané automatizace, včetně ochrany proti smyčce. | +| `withRun(marker, work)` | `src/runtime/context.ts` | Označí, který běh práci způsobil. Bez toho automatizace spouští sama sebe. | +| `findBuiltinStep(...)` | `src/runtime/builtinSteps.ts` | Kroky, které sahají do našeho úložiště, ne ven přes HTTP. | +| `findPersonByExternalId(...)` | `src/data/people.ts` | Řešitel podle ID z cizí aplikace, například voicebotId. | +| `notify(input)` | `src/data/notifications.ts` | Upozorní člověka. Nečeká se a nevyhazuje chyby, stejně jako audit. | | `runFlow(steps, context, options)` | `src/runtime/executor.ts` | Vykoná strom kroků. Nikdy nevyhodí výjimku, chyba je výsledek. Používá to akce na ticketu i webhook, aby se strom choval všude stejně. | | `widgetCatalog(tenantIds, userId)` | `src/data/widgets.ts` | Jediná definice toho, co jde položit na dashboard. Používá ji nabídka i kontrola ukládaného rozložení. | | `intakeEvent(input)` | `src/data/ticketStore.ts` | Přijme událost zvenku: podle externího ID buď založí ticket, nebo ji navěsí na existující. Jediná cesta, kterou se událost stává ticketem. | diff --git a/documentation/20-fronta-a-runtime.md b/documentation/20-fronta-a-runtime.md new file mode 100644 index 0000000..d0554d7 --- /dev/null +++ b/documentation/20-fronta-a-runtime.md @@ -0,0 +1,210 @@ +# Fronta, worker a spouštěče + +Jak se událost dostane od webhooku k vykonanému stromu. Rozbor kapacity je +v [19-kapacita-200-firem.md](19-kapacita-200-firem.md), model ticketu +v [18-ticketovaci-system.md](18-ticketovaci-system.md). + +## Webhook odpoví hned, práci udělá worker + +``` +POST /webhook/ + -> kontrola těla podle kontraktu + -> zápis do fronty + -> 202 Accepted (do jednotek milisekund) + +worker (na pozadí) + -> vezme z fronty + -> vykoná strom + -> zapíše výsledek do logu ticketu + -> při chybě naplánuje další pokus, nebo založí incident +``` + +Odesílatel **nikdy nečeká** na cizí službu. Důvody: + +- Za konektory neručíme. iDoklad může odpovídat pět sekund nebo být hodinu + mimo. Kdyby se čekalo v requestu, odesílateli vyprší timeout a událost je + pryč, přestože jsme ji dostali. +- Pád procesu nesmí ztratit práci. Záznam ve frontě restart přežije, + rozdělaný běh v paměti ne. +- Bez fronty není kam si poznamenat, že se to má za minutu zkusit znovu. + +Odpověď 202 znamená "převzali jsme to", ne "hotovo". Výsledek se hledá +v `GET /api/dashboard/runs` nebo v logu ticketu. + +## Co frontu plní + +| Druh | Kdo to spustí | Příklad | +| --- | --- | --- | +| Push | cizí služba zavolá nás | e-shop pošle novou objednávku | +| Vnitřní událost | něco se stalo u nás | vznikl nebo se změnil ticket | +| Pull | ptáme se sami | e-mail, zprávy z Messengeru | + +### Pull, tedy pravidelné dotazování + +Většina služeb webhooky nemá. U e-mailu a schránek zpráv se **musíme ptát**. +Dělá to plánovač: každých 30 sekund projde automatizace, jejichž spouštěč je +"musí se obvolávat", a u těch, kterým uplynula perioda, zařadí běh. + +Plánovač **sám nic nevolá**. Jen řekne "je čas" a samotný dotaz je první krok +stromu. Díky tomu se dotazování chová stejně jako cokoliv jiného: má záznam +v logu, opakuje se při chybě a jde ho změnit bez zásahu do kódu. + +Výchozí periody: pošta 60 s, zprávy 30 s, plánovač 60 s. Automatizace si to +může přepsat polem `intervalSec` u spouštěče, minimum je 10 sekund - kratší už +není dotazování, ale útok na cizí službu. + +Jeden čekající dotaz na automatizaci: když předchozí ještě běží, další se +nezařadí. Jinak by se fronta zaplnila dotazy na službu, která stejně nestíhá. + +## Kontrakt webhooku + +Odesílatelé posílají různé tvary. Jeden `{"a":"aaa"}`, druhý celý model +s vnořenými objekty a poli. Proto má každý parametr spouštěče **cestu**: + +```json +{ + "document": { "id": "D-99" }, + "errors": [{ "code": "OCR_FAIL", "message": "Nepodařilo se přečíst částku" }] +} +``` + +| Parametr | Cesta | Typ | +| --- | --- | --- | +| `docId` | `document.id` | string | +| `errorMessage` | `errors.0.message` | string | + +Ve stromu se pak píše `{{docId}}` bez ohledu na to, jak hluboko to odesílatel +schoval. Celé tělo je navíc pod `_body`, takže se nic neztratí. + +Chybějící povinný parametr vrací 400 s tím, který to je a kde se hledal. +Přebytek v těle nevadí - odesílatel často posílá víc, než potřebujeme, a +odmítnout ho kvůli tomu by znamenalo, že webhook nejde zapojit. + +`GET` na tutéž adresu vrátí nápovědu: co se čeká, na jakých cestách a ukázku +těla. V portálu je u adresy vidět totéž včetně metody, a kopíruje se **celá +adresa včetně domény**. + +## Opakování a vzdání se + +| Pokus | Kdy | +| --- | --- | +| 1. | hned | +| 2. | za 30 s | +| 3. | za 2 min | +| 4. | za 10 min | +| 5. | za hodinu | + +Pak běh skončí jako `failed` a zůstane k nahlédnutí. Nemaže se: bez záznamu +by nikdo nezjistil, že se něco nestalo. + +**Marná chyba se neopakuje vůbec.** Chybějící skript nebo neexistující skupina +za minutu existovat nezačne, takže se běh rovnou vzdá. Opakuje se jen to, co +může pominout: nedostupná služba, timeout. + +## Incident z chyby + +Každá chyba, kterou už nemá smysl zkoušet, založí incident se **dvěma +úrovněmi**: + +- `title` a `impact` čte **klient**. Bez názvů kroků a ID běhů: "Automatizace + u ticketu TK-4822 nedoběhla do konce, data jsme neztratili." +- `detail` čte **admin**. Je v něm všechno: která automatizace, který běh, + kolik pokusů, který krok selhal, celé hlášení, výpis všech kroků a data, + která přišla na vstupu. + +`detail` se vrací **jen správci platformy**. Je to naše diagnostika, ne +informace pro zákazníka. + +Stejná příčina nezakládá druhý incident, dokud je první otevřený. Jinak by +deset stejných chyb znamenalo deset incidentů a nikdo by se v tom nevyznal. + +## Ochrana proti smyčce + +Automatizace navázaná na změnu ticketu ticket změní, čímž se spustí znovu. +Bez ochrany to server položí, což se při vývoji stalo. + +Dvě pojistky: + +1. **Označení běhu.** Změna, kterou udělala automatizace X, nespustí + automatizaci X. Používá se na to `AsyncLocalStorage`, protože běží čtyři + běhy naráz a obyčejná proměnná by patřila všem. +2. **Strop na ticket.** Jedna automatizace smí nad jedním ticketem běžet + nejvýš pětkrát za minutu. Chytí to i smyčku mezi dvěma automatizacemi, + kterou první pojistka nepozná. Překročení se zaloguje, aby to šlo spravit. + +## Spravedlivost mezi firmami + +Z každé firmy se na jedno kolo vezme nejvýš jeden běh. Jedna firma s tisícem +událostí tak nezablokuje ostatní - bez toho stačí jeden rozbitý e-shop +a servicedesk stojí všem. + +## Kroky, které děláme my + +Založit ticket nebo přehodit ho na člověka není volání cizí služby, takže to +nejde přes skript - sahá to do našeho úložiště. Pro uživatele je to v katalogu +operace jako každá jiná. + +| Krok | Co dělá | +| --- | --- | +| `ticket/upsert` | podle externího ID založí ticket, nebo na existující navěsí událost | +| `ticket/assign-least-busy` | předá nejvolnějšímu ze skupiny, při shodě rozhoduje podíl ke kapacitě | +| `ticket/set-type` | nastaví typ, za kterým stojí vlastní pole | +| `ticket/set-stage` | posune do další fáze workflow daného typu | +| `ticket/add-tags` | přidá štítky, existující nechá | +| `ticket/set-status` | změní stav v životním cyklu | +| `incident/create` | založí incident | +| `flow/pause`, `flow/log` | pauza a zápis do logu | + +## Tři osy na ticketu + +| Osa | Kdo ji určuje | K čemu | +| --- | --- | --- | +| `status` | pevná čtveřice (nový, v řešení, čeká, vyřešeno) | životní cyklus, počítají se z něj statistiky a fronta | +| `stage` | firma u typu ticketu (`TicketType.statuses`) | postup uvnitř typu: čeká na zabalení, předáno dopravci | +| `tags` | kdokoliv, volně | označení, která spolu nemusí souviset | + +Fáze může být **jen jedna**, proto se na ni dá spolehnout v podmínce. Přes +štítky by to fungovalo taky, ale ticket by mohl mít "čeká na zabalení" +i "expedováno" naráz a nikdo by nepoznal, co platí. + +Fáze mimo workflow typu se odmítne. Překlep by jinak tiše vyřadil podmínku, +která na fázi stojí. + +## Živý dashboard + +Dlaždice nad našimi daty se překreslí na událost ze streamu, tedy hned. + +Data z konektorů ne: server je drží v mezipaměti podle `ttlSec` u widgetu. +Přehled se po uplynutí té doby zeptá znovu a dokud je mezipaměť čerstvá, +dostane ji zpátky bez volání cizí služby. U dlaždice je vidět stáří dat +a tlačítko, které vynutí načtení znovu. + +Bez mezipaměti by otevření přehledu znamenalo volání cizího API za každou +dlaždici, a to má limity a někdy se za to platí. + +## Ověřený scénář + +Scénář jedné firmy: sklad (3 lidi), expedice (2), IT (3), typy ticketu +objednávka a chyba, dvě automatizace. + +1. Web pošle `{"kind":"order.created","data":{"order":{"id":"5001"}}}`. +2. Webhook odpoví **202 za 12 ms**, nic nečeká. +3. Worker založí ticket, dá mu typ objednávka a štítek čeká na zabalení. +4. Změna ticketu spustí druhou automatizaci, ta podle typu a štítku předá + práci **nejvolnějšímu ze skladu**. +5. Druhá objednávka jde **jinému člověku**, protože první už jednu má. +6. Chyba z převodníku dokladů přijde s vnořenou cestou `errors.0.message`, + vznikne ticket typu chyba, dostane ho IT a **založí se incident**. + +Ověřeno 19 kontrolami proti běžícímu serveru. + +## Co zbývá + +- **Víc instancí.** Výběr z fronty je v paměti jednoho procesu. Nad Postgresem + to musí být `SELECT ... FOR UPDATE SKIP LOCKED`, jinak si dva workery + vezmou tentýž běh. Místo je označené v `runtime/queue.ts`. +- **Strop souběžných volání na dvojici firma a služba** a vypnutí služby po + sérii chyb. Timeout a rozlišení "zkusit znovu / marné" už ve + `scripts/http.ts` je. +- **Dlouhé čekání** (pošli e-mail za tři dny) přes běh naplánovaný na později. + Krok `flow/pause` umí nejvýš minutu, protože blokuje běh. diff --git a/documentation/99-zmeny.md b/documentation/99-zmeny.md index a3f7630..014eb41 100644 --- a/documentation/99-zmeny.md +++ b/documentation/99-zmeny.md @@ -2,6 +2,73 @@ Nejnovejsi nahore. +## 2026-08-13 - fronta, worker a spoustece + +Popis v [20-fronta-a-runtime.md](20-fronta-a-runtime.md). + +### Zmeneno zasadne + +- **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. + +### Pridano + +- `runtime/queue.ts`: fronta behu v ulozisti. Opakovani s rostouci prodlevou + (30 s, 2 min, 10 min, hodina), spravedlive poradi po firmach, navrat + zaseknutych behu po restartu, uklid hotovych. +- `runtime/worker.ts`: bere praci z fronty, ctyri behy naraz. +- `runtime/triggers.ts`: tri druhy spoustecu. Push (webhook), vnitrni udalost + (vznik a zmena ticketu) a **pull, tedy pravidelne dotazovani** u sluzeb, + ktere webhooky nemaji - posta, zpravy. Planovac jen rekne "je cas", samotny + dotaz je prvni krok stromu. +- **Kontrakt tela webhooku.** Kazdy parametr ma cestu (`data.order.id`, + `errors.0.message`), takze jde napojit i odesilatel s vnorenym modelem. + U adresy je videt metoda, ukazka tela podle parametru a kopiruje se cela + adresa vcetne domeny. +- Vnitrni kroky: `ticket/upsert` (zaloz nebo doplň podle externiho ID), + `assign-least-busy`, `assign-by-external`, `set-type`, `set-stage`, + `add-tags`, `set-status`, `incident/create`, `flow/pause`, `flow/log`. +- **Faze ticketu** (`stage`) 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** (`externalIds`). Voicebot posle + `voicebotId` a ticket skonci u toho, komu patri. Vazba je na jednom miste, + ne v kazde automatizaci. +- **Upozorneni**: komu prijde ticket, ten to vidi hned, vcetne cisla u zalozky. +- **Incident z kazde chyby** 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` se vraci jen spravci platformy. +- **Ochrana proti smycce.** Automatizace navazana na zmenu ticketu ticket + meni, cimz se spousti znovu - pri vyvoji to server polozilo. Resi to + oznaceni behu (`AsyncLocalStorage`) a strop peti behu na ticket za minutu. +- Zivy dashboard: dlazdice nad nasimi daty se prekresli na udalost, data + z konektoru drzi server podle `ttlSec` a jde vynutit nacteni znovu. + +### Opraveno + +- `path` a `intervalSec` u spoustece se pri ulozeni zahazovaly, takze kontrakt + webhooku nefungoval. +- Nad seznamem neslo pouzit `contains`, takze na stitky neslo postavit + podminku. Prave na tom stoji prideleni prace. +- Novejsi vystup kroku ted prekryje starsi se stejnym jmenem. Driv to builder + hlasil jako konflikt i tam, kde zadny nebyl. +- Marna chyba (chybejici skript, neexistujici skupina) se uz neopakuje petkrat. + +### Overeno + +Dva scenare proti bezicimu serveru, 34 kontrol celkem: + +1. **Firma se skladem, expedici a IT.** Webhook odpovedel za 12 ms, worker + zalozil ticket, dal mu typ a stitek, druha automatizace ho podle typu + a stitku predala nejvolnejsimu ze skladu. Druha objednavka sla jinemu + cloveku. Chyba z prevodniku dokladu prisla vnorenou cestou, skoncila u IT + a zalozila incident. +2. **Hovory z voicebota.** Telo `{callSid, status, voicebotId}`: callSid do + externiho ID, status do faze, prirazeni podle voicebotId. Tri zpravy + o tomtez hovoru daly **jeden ticket** se tremi udalostmi. Neznamy voicebot + neskoncil tise - je videt ve fronte i jako incident. + ## 2026-08-13 - runtime, prokliky z widgetu a kapacitni rozbor ### Pridano diff --git a/src/config.ts b/src/config.ts index 72cec93..7387053 100644 --- a/src/config.ts +++ b/src/config.ts @@ -152,6 +152,21 @@ export const config = { * Jen pro lokalni vyvoj, v nasazeni musi zustat vypnute. */ allowPrivateTargets: process.env.ALLOW_PRIVATE_TARGETS === 'true', + /** + * Zpracovava tenhle proces frontu behu? + * + * Vychozi ano, takze jedna instance umi obojí. Az bude potreba oddelit + * vykon od API, spusti se tentyz obraz podruhe s `WORKER=1` a u API + * se nastavi `WORKER=0`. + */ + workerEnabled: process.env.WORKER !== '0', + /** + * Jak casto plánovač hleda, co je na case (v sekundach). + * + * Casovane spoustece nemaji sekundovou presnost a nepotrebuji ji. Kratsi + * interval znamena jen vic dotazu do fronty. + */ + schedulerIntervalSec: positiveNumber(process.env.SCHEDULER_INTERVAL_SEC, 30), }; /** Zaklad verejne adresy aplikace vcetne prefixu proxy. */ diff --git a/src/data/automationStore.ts b/src/data/automationStore.ts index 698e686..55fb99a 100644 --- a/src/data/automationStore.ts +++ b/src/data/automationStore.ts @@ -24,6 +24,16 @@ export interface TriggerField { name: string; type: FieldType; required: boolean; + /** + * Kde ta hodnota v tele je, kdyz to neni primo `name`. + * + * Odesilatele posilaji ruzne tvary: jeden `{"a":"aaa"}`, druhy cely model + * s vnorenymi objekty a poli. Bez cesty by slo napojit jen ploche telo + * a slozitejsi odesilatel by se musel prizpusobovat nam, coz nejde. + * + * Priklady: `customer.id`, `errors.0.message`, `data.items`. + */ + path?: string; } export interface FlowTrigger { @@ -31,6 +41,11 @@ export interface FlowTrigger { operationId: string; /** Deklarovane vstupni parametry. Podminky se odkazuji na jejich `id`. */ fields: TriggerField[]; + /** + * Jak casto se ma sluzba obvolavat, kdyz nam sama nezavola (v sekundach). + * Plati jen u spoustecu, ktere se musi ptat. Minimum je 10 s. + */ + intervalSec?: number; /** * Neodhadnutelny token v adrese webhooku. Generuje VZDY server, * klient ho nesmi urcovat ani menit. diff --git a/src/data/bootstrap.ts b/src/data/bootstrap.ts index e051fce..568fe1c 100644 --- a/src/data/bootstrap.ts +++ b/src/data/bootstrap.ts @@ -12,6 +12,10 @@ import { refreshActions, actionPermissions, actionStore, seedActions } from './ticketActions.js'; import { refreshCustomWidgets, customWidgetStore, seedCustomWidgets } from './customWidgets.js'; import { auditStore } from './audit.js'; +import { notificationStore, refreshNotifications } from './notifications.js'; +import { onTicket } from './ticketHooks.js'; +import { onTicketEvent } from '../runtime/triggers.js'; +import { initQueue } from '../runtime/queue.js'; import { groupStore, personStore, @@ -56,6 +60,7 @@ const entities: Array<{ { store: customWidgetStore as EntityStore, seed: seedCustomWidgets }, // Audit vychozi sadu nema, zaznamy vznikaji az provozem. { store: auditStore as EntityStore }, + { store: notificationStore as EntityStore }, ]; /** @@ -69,6 +74,7 @@ const runtimeData: Array<{ name: string; init: () => Promise }> = [ { name: 'automatizace', init: initAutomations }, { name: 'incidenty', init: initIncidents }, { name: 'rozlozeni dashboardu', init: initLayouts }, + { name: 'fronta behu', init: initQueue }, ]; /** @@ -86,6 +92,7 @@ export async function refreshCaches(): Promise { refreshTicketTypes(), refreshActions(), refreshCustomWidgets(), + refreshNotifications(), ]); // Prava k akcim vznikaji z definic akci, takze se registruji az po jejich nacteni. @@ -128,5 +135,11 @@ export async function bootstrapData(options: { databaseReady: boolean }): Promis } console.info(`[data] nactena provozni data: ${runtimeData.map((i) => i.name).join(', ')}`); + /* + * Zmena ticketu zaradi navazane automatizace. Registruje se az tady, aby + * uloziste ticketu nemuselo vedet o runtime - jinak by vznikl kruh. + */ + onTicket((change, ticket) => onTicketEvent(change, ticket)); + return mode; } diff --git a/src/data/conditions.ts b/src/data/conditions.ts index deb40dc..c21eb01 100644 --- a/src/data/conditions.ts +++ b/src/data/conditions.ts @@ -47,7 +47,12 @@ export const operatorsByType: Record = { // Nad strukturou ma smysl jen to, jestli vubec neco prisla. Porovnavat dva // objekty by znamenalo urcit, co je "stejny", a to zalezi na pripadu. object: ['isEmpty', 'isNotEmpty'], - list: ['isEmpty', 'isNotEmpty'], + /* + * U seznamu ma `contains` smysl: stitky ticketu jsou seznam a bezna otazka + * zni "je mezi nimi 'ceka na zabaleni'?". Bez toho by na stitky neslo + * postavit podminku a cele prideleni prace by se muselo delat jinak. + */ + list: ['contains', 'isEmpty', 'isNotEmpty'], }; /** Operatory, ktere nepotrebuji hodnotu k porovnani. */ diff --git a/src/data/flowScope.ts b/src/data/flowScope.ts index 4d7ba47..42f9259 100644 --- a/src/data/flowScope.ts +++ b/src/data/flowScope.ts @@ -62,7 +62,15 @@ export function collectScopes(flow: AutomationFlow): FlowScopes { if (outputs.length === 0) continue; for (const output of outputs) all.set(output.id, output); - available = [...available, ...outputs]; + /* + * Novejsi vystup **prekryje** starsi se stejnym jmenem. Presne to dela + * i runtime pri dosazovani do sablon, takze `{{assigneeId}}` znamena + * vzdycky ten posledni. Kdyby se tady jen pricitalo, hlasil by builder + * konflikt tam, kde zadny neni - napriklad kdyz krok vraci totez, co uz + * dal spoustec. + */ + const names = new Set(outputs.map((output) => output.name)); + available = [...available.filter((field) => !names.has(field.name)), ...outputs]; } }; diff --git a/src/data/incidentStore.ts b/src/data/incidentStore.ts index 51b702d..8610acd 100644 --- a/src/data/incidentStore.ts +++ b/src/data/incidentStore.ts @@ -16,10 +16,33 @@ export type IncidentStatus = 'investigating' | 'identified' | 'monitoring' | 're export interface Incident { id: string; + /** + * Firma, ktere se to tyka. `null` = nas vlastni, platformni incident + * (vypadek sluzby, ktera je spolecna vsem). + */ + tenantId: string | null; + /** Srozumitelne pro klienta. Zadne stack trace, zadne ID kroku. */ title: string; service: string; severity: IncidentSeverity; status: IncidentStatus; + /** + * Co s tim ma klient delat, nebo co to pro nej znamena. + * + * "Objednávka 5001 se nepřenesla do účetnictví, zkusíme to znovu" je + * pouzitelne. "TypeError: undefined" neni. + */ + impact: string; + /** + * Vsechno, co potrebuje admin: cele hlaseni, ktery beh, ktery krok, co + * prislo na vstupu. Zamerne oddelene od `title` - klient tohle videt nema + * a admin bez toho nema z ceho vychazet. + */ + detail: string | null; + /** Odkud incident vznikl, napr. `run_ab12` nebo `rucne`. */ + source: string | null; + /** Ticket, u ktereho to prasklo. */ + ticketId: string | null; startedAt: string; resolvedAt: string | null; } @@ -31,28 +54,43 @@ function minutesAgo(minutes: number): string { const incidents: Incident[] = [ { id: 'INC-231', + tenantId: null, title: 'Zvýšená latence hlasové brány (region EU-West)', service: 'Voicebot Gateway', severity: 'sev2', status: 'monitoring', + impact: 'Hovory se spojují pomaleji, žádný se ale neztrácí.', + detail: null, + source: null, + ticketId: null, startedAt: minutesAgo(88), resolvedAt: null, }, { id: 'INC-230', + tenantId: null, title: 'Timeouty při zápisu do fakturačního API', service: 'Integrace / iDoklad', severity: 'sev3', status: 'identified', + impact: 'Část faktur se vystaví se zpožděním. Nic se neztrácí, opakuje se to samo.', + detail: null, + source: null, + ticketId: null, startedAt: minutesAgo(240), resolvedAt: null, }, { id: 'INC-228', + tenantId: null, title: 'Neúspěšné doručení webhooků z e-shopu', service: 'Webhook Router', severity: 'sev3', status: 'resolved', + impact: 'Objednávky z e-shopu chvíli nepadaly do systému. Doplnily se zpětně.', + detail: null, + source: null, + ticketId: null, startedAt: minutesAgo(2_900), resolvedAt: minutesAgo(2_700), }, @@ -61,16 +99,13 @@ const incidents: Incident[] = [ let counter = 231; /** Tvar v ulozisti. Cas vzniku je `startedAt`, `createdAt` je kvuli rozhrani. */ -interface StoredIncident extends Incident, TenantEntity { - tenantId: null; -} +interface StoredIncident extends Incident, TenantEntity {} const mirror = withMirror(defineStore('incident')); function toStored(incident: Incident): StoredIncident { return { ...incident, - tenantId: null, createdAt: incident.startedAt, updatedAt: nowIso(), }; @@ -87,8 +122,8 @@ export async function initIncidents(): Promise { incidents.length = 0; for (const row of rows) { - // Sloupce uloziste zpatky nepatri, incident je nezna. - const { tenantId: _tenantId, createdAt: _createdAt, updatedAt: _updatedAt, ...incident } = row; + // Casy uloziste zpatky nepatri, incident ma vlastni `startedAt`. + const { createdAt: _createdAt, updatedAt: _updatedAt, ...incident } = row; incidents.push(incident); } @@ -98,26 +133,51 @@ export async function initIncidents(): Promise { } } -export function listIncidents(): Incident[] { - return [...incidents].sort((a, b) => { +/** + * Incidenty firmy plus platformni. + * + * Bez omezeni na firmy vraci prazdno, ne vse. Zapomenuty filtr nesmi znamenat + * "ukaz cizi incidenty" - je to stejne pravidlo jako u ticketu. + */ +export function listIncidents(tenantIds?: string[]): Incident[] { + const visible = tenantIds + ? incidents.filter((item) => item.tenantId === null || tenantIds.includes(item.tenantId)) + : incidents; + + return [...visible].sort((a, b) => { if (a.status === 'resolved' && b.status !== 'resolved') return 1; if (b.status === 'resolved' && a.status !== 'resolved') return -1; return b.startedAt.localeCompare(a.startedAt); }); } +/** Uz je stejny problem otevreny? Aby deset stejnych chyb nedelalo deset incidentu. */ +export function findOpenIncident(source: string): Incident | undefined { + return incidents.find((item) => item.source === source && item.status !== 'resolved'); +} + export function createIncident(input: { title: string; service: string; severity: IncidentSeverity; + tenantId?: string | null; + impact?: string; + detail?: string | null; + source?: string | null; + ticketId?: string | null; }): Incident { counter += 1; const incident: Incident = { id: `INC-${counter}`, + tenantId: input.tenantId ?? null, title: input.title, service: input.service, severity: input.severity, status: 'investigating', + impact: input.impact ?? '', + detail: input.detail ?? null, + source: input.source ?? null, + ticketId: input.ticketId ?? null, startedAt: new Date().toISOString(), resolvedAt: null, }; diff --git a/src/data/notifications.ts b/src/data/notifications.ts new file mode 100644 index 0000000..f25536a --- /dev/null +++ b/src/data/notifications.ts @@ -0,0 +1,145 @@ +/** + * Upozorneni pro cloveka. + * + * Kdyz nekomu prijde ticket, musi se to dozvedet. Bez toho je prirazeni jen + * zmena udaje v databazi, o ktere nikdo neni. + * + * Cesta je: prirazeni -> upozorneni -> udalost na sbernici -> cislo u zalozky + * a hlaska v portalu. Kdo ma portal zavreny, uvidi to pri prihlaseni, protoze + * upozorneni jsou ulozena, ne jen poslana. + * + * Spojka mezi resitelem a uctem je **e-mail**: resitel nemusi mit ucet + * (viz data/people.ts), a kdyz ho nema, upozorneni se zahodi. + */ + +import { randomUUID } from 'node:crypto'; +import { publish } from '../events/bus.js'; +import { findPerson } from './people.js'; +import { defineStore, nowIso, type TenantEntity } from './store/index.js'; +import { withCache } from './store/cached.js'; +import { findUserByEmail } from './users.js'; + +export type NotificationKind = 'ticket.assigned' | 'ticket.mentioned' | 'automation.failed'; + +export interface Notification extends TenantEntity { + tenantId: string; + /** Komu to patri. ID uzivatelskeho uctu, ne resitele. */ + userId: string; + kind: NotificationKind; + title: string; + /** Kam se proklikne. */ + href: string | null; + ticketId: string | null; + readAt: string | null; +} + +export const notificationStore = defineStore('notification'); +const cache = withCache(notificationStore); + +/** Kolik se jich drzi. Starsi se odmazavaji, neni to archiv. */ +const MAX_PER_USER = 200; + +export async function refreshNotifications(): Promise { + await cache.refresh(); +} + +export interface NotifyInput { + tenantId: string; + /** Resitel, kteremu to patri. Ucet se dohleda pres e-mail. */ + personId: string; + kind: NotificationKind; + title: string; + href?: string | null; + ticketId?: string | null; +} + +/** + * Zalozi upozorneni. + * + * **Necekana se a nevyhazuje chyby.** Rozbite upozorneni nesmi shodit + * prirazeni ticketu - stejna dohoda jako u auditu. + */ +export function notify(input: NotifyInput): void { + const person = findPerson(input.personId); + if (!person) return; + + const user = findUserByEmail(person.email); + if (!user) { + // Resitel bez uctu je bezna vec. Neni komu to ukazat, ale neni to chyba. + return; + } + + const timestamp = nowIso(); + const notification: Notification = { + id: `ntf_${randomUUID().slice(0, 12)}`, + tenantId: input.tenantId, + userId: user.id, + kind: input.kind, + title: input.title, + href: input.href ?? null, + ticketId: input.ticketId ?? null, + readAt: null, + createdAt: timestamp, + updatedAt: timestamp, + }; + + void notificationStore + .create(notification) + .then(async () => { + await cache.refresh(); + // Portal ma cislo u zalozky prekreslit hned, ne az pri obnoveni stranky. + publish('notification.created', input.title, { + userId: user.id, + ticketId: notification.ticketId, + }); + await trim(user.id); + }) + .catch((err: unknown) => { + console.error('[upozorneni] zapis selhal:', err); + }); +} + +async function trim(userId: string): Promise { + const mine = cache + .all() + .filter((item) => item.userId === userId) + .sort((a, b) => a.createdAt.localeCompare(b.createdAt)); + + const excess = mine.slice(0, Math.max(0, mine.length - MAX_PER_USER)); + for (const item of excess) { + await notificationStore.remove(item.id, { tenantIds: [item.tenantId], includeGlobal: true }); + } + if (excess.length > 0) await cache.refresh(); +} + +/** Upozorneni uzivatele, nejnovejsi nahore. */ +export function listNotifications(userId: string, limit = 50): Notification[] { + return cache + .all() + .filter((item) => item.userId === userId) + .sort((a, b) => b.createdAt.localeCompare(a.createdAt)) + .slice(0, limit); +} + +export function unreadCount(userId: string): number { + return cache.all().filter((item) => item.userId === userId && item.readAt === null).length; +} + +/** Oznaci prectene. Bez `ids` vsechny. */ +export async function markRead(userId: string, ids?: string[]): Promise { + const timestamp = nowIso(); + const mine = cache + .all() + .filter((item) => item.userId === userId && item.readAt === null) + .filter((item) => !ids || ids.includes(item.id)); + + for (const item of mine) { + await notificationStore.update( + item.id, + { readAt: timestamp }, + { tenantIds: [item.tenantId], includeGlobal: true }, + ); + } + if (mine.length > 0) await cache.refresh(); + return mine.length; +} diff --git a/src/data/people.ts b/src/data/people.ts index efa7f12..6e684f4 100644 --- a/src/data/people.ts +++ b/src/data/people.ts @@ -24,6 +24,17 @@ export interface Person extends TenantEntity { capacity: number; /** Vypnuty resitel se nenabizi k prirazeni, ale stare tickety nespadnou. */ enabled: boolean; + /** + * ID, pod kterymi cloveka znaji cizi aplikace. + * + * Voicebot posle `voicebotId`, telefonni ustredna klapku, chat svoje ID - + * a my z toho musime poznat, komu ticket patri. Bez toho by se to muselo + * mapovat v kazde automatizaci zvlast a pri zmene cloveka opravovat na + * peti mistech. + * + * Jeden clovek jich muze mit vic, protoze aplikaci je vic. + */ + externalIds: string[]; } export const personStore = defineStore('person'); @@ -31,7 +42,7 @@ const cache = withCache(personStore); export function seedPeople(): Person[] { const timestamp = nowIso(); - const base = { enabled: true, createdAt: timestamp, updatedAt: timestamp }; + const base = { enabled: true, externalIds: [], createdAt: timestamp, updatedAt: timestamp }; return [ { ...base, id: 'ppl_vomacka', tenantId: 'tnt_automia', name: 'Karel Vomáčka', email: 'karel.vomacka@automia.cz', role: 'Servicedesk', capacity: 8 }, @@ -66,6 +77,26 @@ export function listAllPeople(tenantIds: string[]): Person[] { .sort((a, b) => a.name.localeCompare(b.name, 'cs')); } +/** + * Resitel podle ID z cizi aplikace. + * + * Hleda se **jen ve vybranych firmach**: dve firmy mohou mit voicebota se + * stejnym ID a ticket nesmi skoncit u cizi firmy. + */ +export function findPersonByExternalId(value: string, tenantIds: string[]): Person | undefined { + const needle = value.trim(); + if (needle === '') return undefined; + + return cache + .all() + .find( + (person) => + tenantIds.includes(person.tenantId) && + person.enabled !== false && + (person.externalIds ?? []).some((id) => id.trim() === needle), + ); +} + export function findPerson(id: string): Person | undefined { return cache.byId(id); } diff --git a/src/data/services.ts b/src/data/services.ts index f713bf8..9efaada 100644 --- a/src/data/services.ts +++ b/src/data/services.ts @@ -334,6 +334,53 @@ export const services: Service[] = [ * kanaly do nich ustuji (WhatsApp, e-mail, hlas) a zalozeny ticket * je zase spoustecem navazne automatizace - typicky "mame zakaznika?". */ + { + id: 'incident', + name: 'Incidenty', + category: 'obecne', + description: + 'Výpadek nebo porucha, která se týká víc lidí najednou. Na rozdíl od ticketu ' + + 'neřeší jednoho zákazníka, ale stav služby.', + icon: 'AlarmClock', + status: 'available', + general: true, + appId: null, + visibility: { mode: 'everyone', tenantIds: [], userIds: [] }, + credentials: [], + triggers: [], + actions: [ + { + id: 'create', + name: 'Založit incident', + description: + 'Když se chyba netýká jednoho ticketu, ale celé služby. Typicky navazuje ' + + 'na ticket typu chyba.', + implementation: 'script', + inputs: [ + { id: 'title', label: 'Název', kind: 'text', required: true }, + { + id: 'service', + label: 'Čeho se týká', + kind: 'text', + required: false, + hint: 'Název služby nebo aplikace, například Web nebo Voicebot.', + }, + { + id: 'severity', + label: 'Závažnost', + kind: 'choice', + required: false, + options: [ + { value: 'sev1', label: 'SEV1, kritická' }, + { value: 'sev2', label: 'SEV2, vážná' }, + { value: 'sev3', label: 'SEV3, menší' }, + ], + }, + ], + outputFields: [{ id: 'incidentId', name: 'incidentId', type: 'string', required: true }], + }, + ], + }, { id: 'ticket', name: 'Tickety', @@ -392,6 +439,25 @@ export const services: Service[] = [ { id: 'ticket.priority', name: 'priority', type: 'string', required: true }, ], }, + { + id: 'changed', + name: 'Ticket vznikl nebo se změnil', + description: + 'Spustí se při každé změně ticketu, včetně vzniku. Na tomhle stojí ' + + 'automatické přidělování práce: podle typu a štítku se rozhodne, kdo to dostane.', + providedFields: [ + { id: 'ticket.id', name: 'ticketId', type: 'string', required: true }, + { id: 'ticket.externalId', name: 'externalId', type: 'string', required: false }, + { id: 'ticket.subject', name: 'subject', type: 'string', required: true }, + { id: 'ticket.status', name: 'status', type: 'string', required: true }, + { id: 'ticket.priority', name: 'priority', type: 'string', required: true }, + { id: 'ticket.typeId', name: 'typeId', type: 'string', required: false }, + { id: 'ticket.stage', name: 'stage', type: 'string', required: false }, + { id: 'ticket.tags', name: 'tags', type: 'list', required: false }, + { id: 'ticket.assigneeId', name: 'assigneeId', type: 'string', required: false }, + { id: 'ticket.company', name: 'company', type: 'string', required: false }, + ], + }, { id: 'status-changed', name: 'Změna stavu ticketu', @@ -406,6 +472,184 @@ export const services: Service[] = [ }, ], actions: [ + { + id: 'upsert', + name: 'Založit nebo doplnit ticket', + description: + 'Podle externího ID buď založí nový ticket, nebo na existující navěsí událost. ' + + 'Externí ID je unikátní v rámci firmy, takže druhá zpráva o téže objednávce ' + + 'skončí na jednom místě.', + implementation: 'script', + inputs: [ + { + id: 'externalId', + label: 'Externí ID', + kind: 'text', + required: false, + hint: 'ID u odesílatele, typicky číslo objednávky. Bez něj vznikne vždy nový ticket.', + }, + { id: 'subject', label: 'Předmět', kind: 'text', required: false }, + { id: 'body', label: 'Obsah', kind: 'longtext', required: false }, + { + id: 'typeId', + label: 'Typ ticketu', + kind: 'text', + required: false, + hint: 'ID typu. Za typem stojí vlastní pole, takže z nich pak akce mohou čerpat.', + }, + { + id: 'tags', + label: 'Štítky', + kind: 'text', + required: false, + hint: 'Oddělené čárkou. Přidají se, existující se nemažou.', + }, + { + id: 'priority', + label: 'Priorita', + kind: 'choice', + required: false, + options: [ + { value: 'low', label: 'Nízká' }, + { value: 'normal', label: 'Běžná' }, + { value: 'high', label: 'Vysoká' }, + { value: 'critical', label: 'Kritická' }, + ], + }, + { + id: 'stage', + label: 'Fáze', + kind: 'text', + required: false, + hint: 'Musí být z workflow daného typu. Například stav hovoru od voicebota.', + }, + { id: 'event', label: 'Typ události', kind: 'text', required: false }, + { id: 'label', label: 'Popisek do časové osy', kind: 'text', required: false }, + ], + outputFields: [ + { id: 'ticketId', name: 'ticketId', type: 'string', required: true }, + { id: 'created', name: 'created', type: 'boolean', required: true }, + ], + }, + { + id: 'assign-least-busy', + name: 'Předat nejvolnějšímu ze skupiny', + description: + 'Najde ve skupině toho, kdo má nejmíň nevyřízených ticketů, a předá mu to. ' + + 'Při shodě rozhoduje podíl ke kapacitě. Vypnutí lidé se přeskočí.', + implementation: 'script', + inputs: [ + { + id: 'groupId', + label: 'Skupina', + kind: 'choice', + required: true, + hint: 'Například sklad nebo IT. Skupiny se spravují v Nastavení.', + }, + { + id: 'ticketId', + label: 'Ticket', + kind: 'text', + required: false, + hint: 'Prázdné = ticket, kvůli kterému běh vznikl.', + }, + ], + outputFields: [ + { id: 'assigneeId', name: 'assigneeId', type: 'string', required: true }, + { id: 'assigneeName', name: 'assigneeName', type: 'string', required: true }, + ], + }, + { + id: 'assign-by-external', + name: 'Předat podle ID z cizí aplikace', + description: + 'Najde řešitele, který má u sebe uvedené externí ID, a předá mu ticket. ' + + 'Typicky voicebotId nebo klapka. Vazba se nastavuje u řešitele, takže ' + + 'při změně člověka se opravuje na jednom místě.', + implementation: 'script', + inputs: [ + { + id: 'value', + label: 'Hodnota', + kind: 'text', + required: true, + hint: 'Například {{voicebotId}}.', + }, + { + id: 'fallbackGroupId', + label: 'Náhradní skupina', + kind: 'text', + required: false, + hint: 'Když se nikdo nenajde, předá se nejvolnějšímu z této skupiny.', + }, + { id: 'ticketId', label: 'Ticket', kind: 'text', required: false }, + ], + outputFields: [ + { id: 'assigneeId', name: 'assigneeId', type: 'string', required: true }, + { id: 'assigneeName', name: 'assigneeName', type: 'string', required: true }, + ], + }, + { + id: 'set-type', + name: 'Nastavit typ ticketu', + description: 'Za typem stojí vlastní pole a podle typu se ukazují akce.', + implementation: 'script', + inputs: [ + { id: 'typeId', label: 'Typ ticketu', kind: 'text', required: true }, + { id: 'ticketId', label: 'Ticket', kind: 'text', required: false }, + ], + }, + { + id: 'add-tags', + name: 'Přidat štítky', + description: 'Existující štítky zůstanou, jinak by se dva kroky přebíjely.', + implementation: 'script', + inputs: [ + { id: 'tags', label: 'Štítky', kind: 'text', required: true, hint: 'Oddělené čárkou.' }, + { id: 'ticketId', label: 'Ticket', kind: 'text', required: false }, + ], + }, + { + id: 'set-stage', + name: 'Posunout do další fáze', + description: + 'Fáze je vlastní workflow typu ticketu, například čeká na zabalení nebo ' + + 'předáno dopravci. Na rozdíl od štítků může být jen jedna, takže se na ni ' + + 'dá spolehnout v podmínce.', + implementation: 'script', + inputs: [ + { + id: 'stage', + label: 'Fáze', + kind: 'text', + required: true, + hint: 'Musí být z workflow daného typu, jinak krok selže.', + }, + { id: 'ticketId', label: 'Ticket', kind: 'text', required: false }, + ], + outputFields: [{ id: 'stage', name: 'stage', type: 'string', required: true }], + }, + { + id: 'set-status', + name: 'Změnit stav ticketu', + description: 'Například vyřešeno, když automatizace dokončila, co měla.', + implementation: 'script', + inputs: [ + { + id: 'status', + label: 'Stav', + kind: 'choice', + required: true, + options: [ + { value: 'new', label: 'Nový' }, + { value: 'open', label: 'V řešení' }, + { value: 'waiting', label: 'Čeká na klienta' }, + { value: 'resolved', label: 'Vyřešeno' }, + ], + }, + { id: 'ticketId', label: 'Ticket', kind: 'text', required: false }, + ], + }, { id: 'create', name: 'Založit ticket', diff --git a/src/data/ticketHooks.ts b/src/data/ticketHooks.ts new file mode 100644 index 0000000..63c8a76 --- /dev/null +++ b/src/data/ticketHooks.ts @@ -0,0 +1,38 @@ +/** + * Co se ma stat po zmene ticketu. + * + * Uloziste ticketu **nesmi vedet o runtime**. Kdyby si ho zavolalo primo, + * vznikl by kruh: uloziste vola frontu, fronta vola automatizace, automatizace + * sahaji do uloziste. Takhle uloziste jen rekne "stalo se tohle" a kdo chce, + * si to poslechne. + * + * Registruje se pri startu, viz `data/bootstrap.ts`. + */ + +import type { Ticket } from './ticketStore.js'; + +export type TicketChange = 'ticket.created' | 'ticket.updated'; + +type Listener = (change: TicketChange, ticket: Ticket) => void; + +const listeners: Listener[] = []; + +export function onTicket(listener: Listener): void { + listeners.push(listener); +} + +/** + * Oznami zmenu. + * + * **Nikdy nevyhodi vyjimku a necekana se.** Rozbity posluchac nesmi shodit + * ulozeni ticketu - to je horsi nez nespustena automatizace. + */ +export function onTicketChanged(change: TicketChange, ticket: Ticket): void { + for (const listener of listeners) { + try { + listener(change, ticket); + } catch (err) { + console.error(`[tickets] posluchac zmeny selhal (${change}):`, err); + } + } +} diff --git a/src/data/ticketStore.ts b/src/data/ticketStore.ts index 824e65b..1c73786 100644 --- a/src/data/ticketStore.ts +++ b/src/data/ticketStore.ts @@ -17,6 +17,8 @@ */ import { publish } from '../events/bus.js'; +import { onTicketChanged } from './ticketHooks.js'; +import { notify } from './notifications.js'; import { findPerson, type Person } from './people.js'; import { defineStore } from './store/index.js'; import { withMirror } from './store/mirror.js'; @@ -137,6 +139,20 @@ export interface Ticket { typeId: string | null; /** Hodnoty vlastnich poli typu. Klic je `TicketTypeField.key`. */ fields: Record; + /** + * Faze ve workflow **daneho typu**, napr. "ceka na zabaleni". + * + * Tri osy zamerne, kazda ma jinou praci: + * - `status` je zivotni cyklus (novy, v reseni, ceka, vyreseno). Podle nej + * se pocita fronta i statistiky, takze musi zustat pevny. + * - `stage` je **postup uvnitr typu** a definuje si ho firma u typu ticketu + * (`TicketType.statuses`). Objednavka ma jine faze nez reklamace. + * - `tags` jsou volne stitky, ktere spolu nemusi souviset. + * + * Delat faze pres stitky by fungovalo, ale nesla by na nich postavit + * kontrola: ticket by mohl mit "ceka na zabaleni" i "expedovano" naraz. + */ + stage: string | null; /** * Volne oznaceni. Na rozdil od typu jich muze byt vic a nestoji za nimi * zadna pole - proto se hodi na filtry a widgety, ne na akce, ktere @@ -185,6 +201,7 @@ interface StoredTicket | 'typeId' | 'fields' | 'tags' + | 'stage' | 'assigneeGroupId' | 'externalId' | 'externalSource' @@ -198,6 +215,7 @@ interface StoredTicket typeId?: string | null; fields?: Record; tags?: string[]; + stage?: string | null; externalId?: string | null; externalSource?: string | null; firstResponseAt?: string | null; @@ -259,6 +277,9 @@ function persist(ticket: StoredTicket): void { function touch(ticket: StoredTicket): void { ticket.updatedAt = new Date().toISOString(); persist(ticket); + // Automatizace navazane na zmenu ticketu. Necekana se a chyby nevyhazuje, + // jinak by rozbita fronta rozbila ukladani ticketu. + onTicketChanged('ticket.updated', toTicket(ticket)); } /** @@ -783,6 +804,7 @@ function toTicket(stored: StoredTicket): Ticket { typeId: stored.typeId ?? null, fields: stored.fields ?? {}, tags: stored.tags ?? [], + stage: stored.stage ?? null, externalId: stored.externalId ?? null, externalSource: stored.externalSource ?? null, firstResponseAt: stored.firstResponseAt ?? null, @@ -816,6 +838,8 @@ export interface TicketFilter { channel?: TicketChannel; /** Typ ticketu. `none` = tickety bez typu. */ typeId?: string; + /** Faze ve workflow typu. `none` = tickety bez faze. */ + stage?: string; /** Jeden tag. `none` = tickety bez tagu. */ tag?: string; /** Skupina resitelu. `none` = bez skupiny. */ @@ -829,6 +853,8 @@ export function listTickets(filter: TicketFilter): Ticket[] { if (filter.channel && ticket.channel !== filter.channel) return false; if (filter.typeId === 'none' ? ticket.typeId : filter.typeId && ticket.typeId !== filter.typeId) return false; + if (filter.stage === 'none' ? ticket.stage : filter.stage && ticket.stage !== filter.stage) + return false; if (filter.tag === 'none') { if ((ticket.tags ?? []).length > 0) return false; } else if (filter.tag && !(ticket.tags ?? []).includes(filter.tag)) { @@ -943,6 +969,13 @@ function markResponded(ticket: StoredTicket): void { if (!ticket.firstResponseAt) ticket.firstResponseAt = new Date().toISOString(); } +/** Ticket podle naseho ID, jen z povolenych firem. */ +export function findTicket(id: string, tenantIds: string[]): Ticket | undefined { + const stored = tickets.find((ticket) => ticket.id === id); + if (!stored || !tenantIds.includes(stored.tenantId)) return undefined; + return toTicket(stored); +} + /** Ticket firmy podle externiho ID. Klic je dvojice firma a ID, ne ID samotne. */ export function findByExternalId(tenantId: string, externalId: string): Ticket | undefined { const stored = tickets.find( @@ -1197,6 +1230,7 @@ export interface CreateTicketInput { typeId?: string | null; fields?: Record; tags?: string[]; + stage?: string | null; automationId?: string | null; /** Log toho, jak ticket vznikl. Bez nej je ticket nedohledatelny. */ trace?: TraceInput[]; @@ -1236,6 +1270,7 @@ export function createTicket(input: CreateTicketInput): Ticket { typeId: input.typeId ?? null, fields: input.fields ?? {}, tags: input.tags ?? [], + stage: input.stage ?? null, automationId: input.automationId ?? null, createdAt: now, updatedAt: now, @@ -1245,6 +1280,8 @@ export function createTicket(input: CreateTicketInput): Ticket { events.set(stored.id, []); persist(stored); + onTicketChanged('ticket.created', toTicket(stored)); + publish('ticket.created', `Nový ticket ${stored.id}: ${stored.subject}`, { ticketId: stored.id, channel: stored.channel, @@ -1339,10 +1376,24 @@ export function assignTicket( return undefined; } + const previousAssignee = ticket.assigneeId; ticket.assigneeId = person?.id ?? null; if (person) markResponded(ticket); touch(ticket); + // Komu ticket prisel, ten se to musi dozvedet. Znovu prirazeni tomu samemu + // cloveku upozorneni negeneruje, jinak by mu chodilo pri kazde drobnosti. + if (person && person.id !== previousAssignee) { + notify({ + tenantId: ticket.tenantId, + personId: person.id, + kind: 'ticket.assigned', + title: `Máte nový ticket ${ticket.id}: ${ticket.subject}`, + href: `/dashboard/tickety/${ticket.id}`, + ticketId: ticket.id, + }); + } + appendTrace(id, [ { kind: 'note', @@ -1389,6 +1440,41 @@ export function setTicketType( } /** Tagy se prepisuji cele. Prirustkova zmena by u vic lidi naraz kolidovala. */ +/** + * Nastavi fazi. + * + * Faze musi byt z workflow daneho typu. Cizi hodnota se odmitne - jinak by + * se do dat dostal preklep a podminka nad fazi by tise prestala platit. + */ +export function setTicketStage( + id: string, + stage: string | null, + allowed: string[], + tenantIds: string[], +): Ticket | undefined { + const ticket = findWritable(id, tenantIds); + if (!ticket) return undefined; + + if (stage !== null && allowed.length > 0 && !allowed.includes(stage)) { + console.warn(`[tickets] ${id}: faze ${stage} neni ve workflow typu`); + return undefined; + } + + const previous = ticket.stage ?? null; + ticket.stage = stage; + touch(ticket); + + appendTrace(id, [ + { + kind: 'note', + status: 'info', + label: `Fáze: ${previous ?? 'bez fáze'} -> ${stage ?? 'bez fáze'}`, + }, + ]); + publish('ticket.updated', `Ticket ${ticket.id} má novou fázi`, { ticketId: ticket.id }); + return toTicket(ticket); +} + export function setTicketTags(id: string, tags: string[], tenantIds: string[]): Ticket | undefined { const ticket = findWritable(id, tenantIds); if (!ticket) return undefined; diff --git a/src/events/bus.ts b/src/events/bus.ts index 5bc1b21..eb87679 100644 --- a/src/events/bus.ts +++ b/src/events/bus.ts @@ -21,7 +21,9 @@ export type DashboardEventType = | 'automation.updated' | 'automation.deleted' | 'automation.run' - | 'webhook.received'; + | 'webhook.received' + /** Nekomu prislo upozorneni. Portal podle toho prekresli cislo u zalozky. */ + | 'notification.created'; export interface DashboardEvent { id: string; diff --git a/src/index.ts b/src/index.ts index 36bd4a7..a770a24 100644 --- a/src/index.ts +++ b/src/index.ts @@ -22,6 +22,8 @@ import { adminRouter } from './routes/admin.js'; import { runMigrations } from './db/migrate.js'; import { closeDatabase, databaseHealth, isDatabaseEnabled } from './db/pool.js'; import { ensureLoaded, scriptsDir } from './scripts/registry.js'; +import { startWorker, stopWorker } from './runtime/worker.js'; +import { startScheduler, stopScheduler } from './runtime/triggers.js'; const here = path.dirname(fileURLToPath(import.meta.url)); /** Zbuildovana SPA. Vite ji zapisuje do dist/public, viz vite.config.ts. */ @@ -228,6 +230,20 @@ await initConnectorStore({ databaseReady }); */ await bootstrapData({ databaseReady }); +/* + * Worker a planovac. Bezi ve stejnem procesu jako API, ale mimo request: + * webhook zapise udalost do fronty a hned odpovi, praci udela worker. + * + * `WORKER=0` je vypne. To se hodi, az bude vykon oddeleny od API - tentyz + * obraz se pak spusti podruhe jen jako worker. + */ +if (config.workerEnabled) { + startWorker(); + startScheduler(); +} else { + console.info('[start] worker vypnuty (WORKER=0), fronta se v tomhle procesu nezpracovava'); +} + // Poslouchat na vsech rozhranich containeru, ne jen na localhost (AGENTS.md). const server = app.listen(config.port, '0.0.0.0', () => { console.info(`[start] csbot-prototype bezi na portu ${config.port}`); @@ -245,7 +261,9 @@ for (const signal of ['SIGTERM', 'SIGINT'] as const) { server.close(() => { // Rozepsany zapis do souboru se musi dokoncit, jinak se posledni zmena // ztrati - debounce je kratky, ale nenulovy. - void Promise.all([flushConnectorStore(), flushStores()]) + stopScheduler(); + void stopWorker() + .then(() => Promise.all([flushConnectorStore(), flushStores()])) .catch((err: unknown) => console.error('[stop] zapis dat selhal:', err)) .then(() => closeDatabase()) .finally(() => process.exit(0)); diff --git a/src/openapi.ts b/src/openapi.ts index 11c0ae3..dbcb21c 100644 --- a/src/openapi.ts +++ b/src/openapi.ts @@ -1968,6 +1968,178 @@ export function buildOpenApiDocument() { }, }, }, + '/api/dashboard/notifications': { + get: { + tags: ['Dashboard'], + summary: 'Upozorneni prihlaseneho', + description: + 'Cislo u zalozky Tickety a hlasky o pridelene praci. Upozorneni jsou ulozena, ' + + 'takze je najde i ten, kdo mel portal zavreny.', + security: [{ bearerAuth: [] }], + responses: { + '200': { + description: 'Upozorneni', + content: { + 'application/json': { + schema: { + type: 'object', + properties: { + items: { type: 'array', items: { type: 'object' } }, + unread: { type: 'integer' }, + mine: { type: 'integer', description: 'Kolik ticketu ma volajici u sebe.' }, + }, + }, + }, + }, + }, + }, + }, + }, + '/api/dashboard/notifications/read': { + post: { + tags: ['Dashboard'], + summary: 'Oznacit upozorneni jako prectena', + security: [{ bearerAuth: [] }], + requestBody: { + required: false, + content: { + 'application/json': { + schema: { + type: 'object', + properties: { + ids: { + type: 'array', + items: { type: 'string' }, + description: 'Bez seznamu se oznaci vsechna.', + }, + }, + }, + }, + }, + }, + responses: { '200': { description: 'Oznaceno' } }, + }, + }, + '/api/dashboard/runs': { + get: { + tags: ['Automatizace'], + summary: 'Stav fronty behu', + description: + 'Kdyz neco nefunguje, tohle je prvni misto, kam se clovek podiva: ceka fronta, ' + + 'nebo uz to nekolikrat selhalo? U kazdeho behu je cele chybove hlaseni.', + security: [{ bearerAuth: [] }], + responses: { + '200': { + description: 'Fronta a posledni behy', + content: { + 'application/json': { + schema: { + type: 'object', + properties: { + stats: { + type: 'object', + properties: { + pending: { type: 'integer' }, + running: { type: 'integer' }, + done: { type: 'integer' }, + failed: { type: 'integer' }, + oldestPendingAt: { type: 'string', nullable: true }, + }, + }, + items: { type: 'array', items: { type: 'object' } }, + }, + }, + }, + }, + }, + }, + }, + }, + '/webhook/ticket/{token}': { + post: { + tags: ['Webhook'], + summary: 'Prijem udalosti do ticketu', + description: + 'VEREJNY endpoint, autorizuje token firmy v adrese. Se stejnym externalId se ' + + 'udalost navesi na existujici ticket, jinak vznikne novy. externalId je ' + + 'unikatni v ramci firmy.', + parameters: [{ name: 'token', in: 'path', required: true, schema: { type: 'string' } }], + requestBody: { + required: true, + content: { + 'application/json': { + schema: { + type: 'object', + properties: { + externalId: { oneOf: [{ type: 'string' }, { type: 'number' }] }, + source: { type: 'string', example: 'eshop' }, + event: { type: 'string', example: 'order.created' }, + subject: { type: 'string' }, + typeId: { type: 'string' }, + tags: { type: 'array', items: { type: 'string' } }, + fields: { type: 'object', additionalProperties: true }, + }, + }, + }, + }, + }, + responses: { + '201': { description: 'Ticket vznikl' }, + '200': { description: 'Udalost se navesila na existujici ticket' }, + '400': { description: 'Neplatna data' }, + '404': { description: 'Neznamy token' }, + }, + }, + }, + '/webhook/{token}': { + post: { + tags: ['Webhook'], + summary: 'Prijem dat do automatizace', + description: + 'VEREJNY endpoint. **Odpovi hned** (202) a strom vykona worker na pozadi - ' + + 'cizi sluzba muze odpovidat pomalu a odesilateli by vyprsel timeout. ' + + 'Telo se kontroluje proti kontraktu spoustece, vcetne vnorenych cest.', + parameters: [{ name: 'token', in: 'path', required: true, schema: { type: 'string' } }], + requestBody: { + required: true, + content: { + 'application/json': { + schema: { type: 'object', additionalProperties: true }, + }, + }, + }, + responses: { + '202': { + description: 'Prijato, zpracuje se na pozadi', + content: { + 'application/json': { + schema: { + type: 'object', + properties: { + accepted: { type: 'boolean' }, + automationId: { type: 'string' }, + runId: { type: 'string', nullable: true }, + }, + }, + }, + }, + }, + '400': { description: 'Telo neodpovida kontraktu spoustece' }, + '404': { description: 'Neznamy token' }, + '409': { description: 'Automatizace je pozastavena' }, + }, + }, + get: { + tags: ['Webhook'], + summary: 'Napoveda: co se v tele ceka', + description: 'Vraci metodu, seznam parametru vcetne cest a ukazku tela.', + parameters: [{ name: 'token', in: 'path', required: true, schema: { type: 'string' } }], + responses: { + '200': { description: 'Kontrakt' }, + '404': { description: 'Neznamy token' }, + }, + }, + }, '/api/contact': { post: { tags: ['Kontakt'], diff --git a/src/routes/dashboard.ts b/src/routes/dashboard.ts index c3e397e..1e86f22 100644 --- a/src/routes/dashboard.ts +++ b/src/routes/dashboard.ts @@ -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(), }) diff --git a/src/routes/settings.ts b/src/routes/settings.ts index 531dc05..402b98d 100644 --- a/src/routes/settings.ts +++ b/src/routes/settings.ts @@ -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()) diff --git a/src/routes/webhook.ts b/src/routes/webhook.ts index 4d5d1fe..2ca3daa 100644 --- a/src/routes/webhook.ts +++ b/src/routes/webhook.ts @@ -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; - const problems = validatePayload(automation.flow.trigger.fields, payload); + const body = (req.body ?? {}) as Record; + 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[] { + body: Record, +): { values: Record; problems: string[] } { + const values: Record = {}; 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 { + const example: Record = {}; + + 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; + } + 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'; } } diff --git a/src/routes/widgetData.ts b/src/routes/widgetData.ts index f7be23b..dfbbc74 100644 --- a/src/routes/widgetData.ts +++ b/src/routes/widgetData.ts @@ -53,6 +53,13 @@ const statusLabels: Record = { 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, scope: ResolvedScope, + force = false, ): Promise { 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, + 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) { diff --git a/src/runtime/builtinSteps.ts b/src/runtime/builtinSteps.ts new file mode 100644 index 0000000..61b662a --- /dev/null +++ b/src/runtime/builtinSteps.ts @@ -0,0 +1,349 @@ +/** + * Kroky, ktere delame my, ne cizi sluzba. + * + * Zalozit ticket, prehodit ho na nejvolnejsiho cloveka ze skupiny nebo zalozit + * incident nejde pres `runScript`: skript vola HTTP ven, tohle sahá do naseho + * uloziste. Proto vlastni registr, ktery runtime zkusi driv nez skripty. + * + * Pro uzivatele je to k nerozeznani - v katalogu je to operace jako kazda jina + * a ve stromu se nastavuje stejne. + */ + +import { createIncident } from '../data/incidentStore.js'; +import { findGroup, findPersonByExternalId, listPeople } from '../data/people.js'; +import { findTicketType } from '../data/ticketTypes.js'; +import { + assignTicket, + assignTicketGroup, + findByExternalId, + findTicket, + getWorkload, + intakeEvent, + setTicketStage, + setTicketTags, + setTicketType, + updateTicketStatus, + type Ticket, + type TicketPriority, + type TicketStatus, +} from '../data/ticketStore.js'; + +export interface StepContext { + tenantId: string; + /** Ticket, ke kteremu beh patri. Nekdy vznikne az behem nej. */ + ticketId: string | null; +} + +export interface StepOutcome { + ok: boolean; + summary: string; + detail?: string | null; + outputs: Record; + /** Vyplnene, kdyz krok zalozil nebo nasel ticket. Dalsi kroky ho pak maji. */ + ticketId?: string; +} + +type Handler = ( + inputs: Record, + context: StepContext, +) => Promise | StepOutcome; + +/** Neco chybelo. Vraci se jako vysledek, ne jako vyjimka. */ +function missing(what: string): StepOutcome { + return { ok: false, summary: `chybí ${what}`, detail: null, outputs: {} }; +} + +const statuses: TicketStatus[] = ['new', 'open', 'waiting', 'resolved']; +const priorities: TicketPriority[] = ['low', 'normal', 'high', 'critical']; + +/** + * Registr kroku. Klic je `serviceId/operationId` z katalogu. + */ +const handlers: Record = { + /** + * Zalozi ticket, nebo doplni existujici podle externiho ID. + * + * Tohle je krok, kterym vetsina automatizaci zacina: prislo neco zvenku + * a ma z toho byt ticket. Kdyz uz ticket se stejnym externim ID ve firme + * je, **udalost se na nej navesi** misto zalozeni druheho. + */ + 'ticket/upsert': (inputs, context) => { + const subject = inputs.subject?.trim(); + if (!subject && !inputs.externalId) return missing('předmět nebo externí ID'); + + const result = intakeEvent({ + tenantId: context.tenantId, + externalId: inputs.externalId?.trim() || null, + externalSource: inputs.source?.trim() || 'automatizace', + type: inputs.event?.trim() || 'automation', + label: inputs.label?.trim() || subject || 'Událost', + payload: inputs.payload ? safeJson(inputs.payload) : {}, + create: { + subject: subject || inputs.externalId || 'Bez předmětu', + body: inputs.body ?? '', + typeId: inputs.typeId?.trim() || null, + priority: priorities.includes(inputs.priority as TicketPriority) + ? (inputs.priority as TicketPriority) + : 'normal', + }, + addTags: inputs.tags + ? inputs.tags.split(',').map((tag) => tag.trim()).filter(Boolean) + : undefined, + }); + + /* + * Faze se nastavi az po zalozeni: u noveho ticketu jeste neni typ, takze + * by se nemela proti cemu overit. Prazdna hodnota nic nemeni. + */ + if (inputs.stage?.trim()) { + const type = result.ticket.typeId ? findTicketType(result.ticket.typeId) : undefined; + setTicketStage(result.ticket.id, inputs.stage.trim(), type?.statuses ?? [], [ + context.tenantId, + ]); + } + + return { + ok: true, + summary: result.created ? `založen ${result.ticket.id}` : `doplněn ${result.ticket.id}`, + detail: null, + ticketId: result.ticket.id, + outputs: { + ticketId: result.ticket.id, + created: result.created, + externalId: result.ticket.externalId, + }, + }; + }, + + /** + * Preda ticket **nejvolnejsimu** cloveku ze skupiny. + * + * Nejvolnejsi znamena nejmene nevyrizenych ticketu, pri shode nejmensi podil + * ke kapacite. Kdo je vypnuty, nedostane nic. Kdyz nikdo neni, krok selze - + * tise nechat ticket lezet ve fronte je horsi nez to rict nahlas. + */ + 'ticket/assign-least-busy': (inputs, context) => { + const ticketId = inputs.ticketId?.trim() || context.ticketId; + if (!ticketId) return missing('ticket'); + + const groupId = inputs.groupId?.trim(); + if (!groupId) return missing('skupina'); + + const group = findGroup(groupId); + if (!group || group.tenantId !== context.tenantId) { + return { ok: false, summary: `skupina ${groupId} neexistuje`, outputs: {} }; + } + + const people = listPeople([context.tenantId]).filter( + (person) => group.personIds.includes(person.id) && person.enabled !== false, + ); + if (people.length === 0) { + return { ok: false, summary: `skupina ${group.name} nemá koho`, outputs: {} }; + } + + const workload = getWorkload(people, [context.tenantId]); + const best = [...workload.rows].sort((a, b) => { + if (a.open !== b.open) return a.open - b.open; + // Pri shode rozhoduje podil ke kapacite: dva tickety u cloveka s kapacitou + // tri je vic prace nez dva u cloveka s kapacitou deset. + const aShare = a.person.capacity > 0 ? a.open / a.person.capacity : Infinity; + const bShare = b.person.capacity > 0 ? b.open / b.person.capacity : Infinity; + return aShare - bShare; + })[0]; + + if (!best) return { ok: false, summary: 'nikdo volný', outputs: {} }; + + // Skupina zustava na ticketu jako informace, kdo to ma resit. + assignTicketGroup(ticketId, groupId, [context.tenantId]); + const updated = assignTicket(ticketId, best.person.id, [context.tenantId]); + if (!updated) return { ok: false, summary: 'přiřazení se nepodařilo', outputs: {} }; + + return { + ok: true, + summary: `${best.person.name} (${best.open} u sebe)`, + detail: null, + outputs: { assigneeId: best.person.id, assigneeName: best.person.name }, + }; + }, + + /** + * Preda ticket cloveku podle ID z cizi aplikace. + * + * Typicky pripad: voicebot posle `voicebotId` a ticket ma skoncit u toho, + * komu ten voicebot patri. Vazba je u resitele (`externalIds`), takze pri + * zmene cloveka se meni na jednom miste, ne v kazde automatizaci. + * + * Kdyz se nikdo nenajde, muze se pouzit nahradni skupina - jinak by ticket + * tise zustal lezet ve fronte. + */ + 'ticket/assign-by-external': (inputs, context) => { + const ticketId = inputs.ticketId?.trim() || context.ticketId; + if (!ticketId) return missing('ticket'); + + const value = inputs.value?.trim(); + if (!value) return missing('hodnotu, podle které se hledá'); + + const person = findPersonByExternalId(value, [context.tenantId]); + if (!person) { + const fallback = inputs.fallbackGroupId?.trim(); + if (!fallback) { + return { + ok: false, + summary: `nikdo nemá externí ID ${value}`, + detail: + 'Doplňte to ID u řešitele v Nastavení, nebo u kroku vyplňte náhradní skupinu.', + outputs: {}, + }; + } + + const handler = handlers['ticket/assign-least-busy']; + return handler({ ticketId, groupId: fallback }, context); + } + + const updated = assignTicket(ticketId, person.id, [context.tenantId]); + if (!updated) return { ok: false, summary: 'přiřazení se nepodařilo', outputs: {} }; + + return { + ok: true, + summary: `${person.name} (podle ${value})`, + detail: null, + outputs: { assigneeId: person.id, assigneeName: person.name }, + }; + }, + + /** Nastavi typ ticketu. Za typem stoji vlastni pole, proto je to krok. */ + 'ticket/set-type': (inputs, context) => { + const ticketId = inputs.ticketId?.trim() || context.ticketId; + if (!ticketId) return missing('ticket'); + if (!inputs.typeId?.trim()) return missing('typ'); + + const updated = setTicketType(ticketId, inputs.typeId.trim(), undefined, [context.tenantId]); + if (!updated) return { ok: false, summary: 'typ se nepodařilo nastavit', outputs: {} }; + return { ok: true, summary: inputs.typeId, detail: null, outputs: { typeId: inputs.typeId } }; + }, + + /** Prida tagy. Existujici se nemazou, jinak by dva kroky bojovaly. */ + 'ticket/add-tags': (inputs, context) => { + const ticketId = inputs.ticketId?.trim() || context.ticketId; + if (!ticketId) return missing('ticket'); + + const wanted = (inputs.tags ?? '').split(',').map((tag) => tag.trim()).filter(Boolean); + if (wanted.length === 0) return missing('tagy'); + + const current = currentTicket(ticketId, context.tenantId); + const merged = [...new Set([...(current?.tags ?? []), ...wanted])]; + const updated = setTicketTags(ticketId, merged, [context.tenantId]); + if (!updated) return { ok: false, summary: 'tagy se nepodařilo uložit', outputs: {} }; + + return { ok: true, summary: merged.join(', '), detail: null, outputs: { tags: merged } }; + }, + + /** + * Posune ticket do dalsi faze. + * + * Faze je z workflow typu, takze `objednavka` ma jine nez `reklamace`. + * Cizi hodnota se odmitne - preklep by tise vyradil podminku nad fazi. + */ + 'ticket/set-stage': (inputs, context) => { + const ticketId = inputs.ticketId?.trim() || context.ticketId; + if (!ticketId) return missing('ticket'); + + const stage = inputs.stage?.trim(); + if (!stage) return missing('fáze'); + + const ticket = findTicket(ticketId, [context.tenantId]); + const type = ticket?.typeId ? findTicketType(ticket.typeId) : undefined; + const allowed = type?.statuses ?? []; + + const updated = setTicketStage(ticketId, stage, allowed, [context.tenantId]); + if (!updated) { + return { + ok: false, + summary: `fáze ${stage} není ve workflow typu`, + detail: allowed.length > 0 ? `Typ dovoluje: ${allowed.join(', ')}.` : null, + outputs: {}, + }; + } + + return { ok: true, summary: stage, detail: null, outputs: { stage } }; + }, + + /** Zmeni stav ticketu. */ + 'ticket/set-status': (inputs, context) => { + const ticketId = inputs.ticketId?.trim() || context.ticketId; + if (!ticketId) return missing('ticket'); + + const status = inputs.status?.trim() as TicketStatus; + if (!statuses.includes(status)) return missing('platný stav'); + + const updated = updateTicketStatus(ticketId, status, [context.tenantId]); + if (!updated) return { ok: false, summary: 'stav se nepodařilo změnit', outputs: {} }; + return { ok: true, summary: status, detail: null, outputs: { status } }; + }, + + /** + * Zalozi incident. + * + * Incident je neco jineho nez ticket: ticket je pozadavek jednoho zakaznika, + * incident je "nefunguje to a tyka se to vic lidi". Proto vlastni krok. + */ + 'incident/create': (inputs) => { + const title = inputs.title?.trim(); + if (!title) return missing('název'); + + const incident = createIncident({ + title, + service: inputs.service?.trim() || 'Neurčeno', + severity: (['sev1', 'sev2', 'sev3'] as const).includes(inputs.severity as 'sev1') + ? (inputs.severity as 'sev1' | 'sev2' | 'sev3') + : 'sev2', + }); + + return { + ok: true, + summary: incident.id, + detail: null, + outputs: { incidentId: incident.id }, + }; + }, + + /** + * Pauza mezi kroky. + * + * Zamerne kratka a **blokujici v ramci behu**. Dlouhe cekani (poslat e-mail + * za tri dny) takhle delat nejde - to patri do fronty jako beh naplanovany + * na pozdeji, viz `nextAttemptAt`. Strop je proto minuta. + */ + 'flow/pause': async (inputs) => { + const seconds = Math.min(Math.max(Number(inputs.seconds ?? 1) || 1, 0), 60); + await new Promise((resolve) => setTimeout(resolve, seconds * 1_000)); + return { ok: true, summary: `${seconds} s`, detail: null, outputs: {} }; + }, + + /** Zapis do logu ticketu. Na overeni, ze strom dosel, kam mel. */ + 'flow/log': (inputs) => ({ + ok: true, + summary: inputs.message?.trim() || 'zápis', + detail: null, + outputs: {}, + }), +}; + +function currentTicket(ticketId: string, tenantId: string): Ticket | undefined { + return findTicket(ticketId, [tenantId]) ?? findByExternalId(tenantId, ticketId); +} + +function safeJson(text: string): Record { + try { + const parsed: unknown = JSON.parse(text); + return parsed && typeof parsed === 'object' ? (parsed as Record) : { value: parsed }; + } catch { + // Nerozparsovany text neni duvod krok shodit, ulozi se jak je. + return { value: text }; + } +} + +/** Ma tenhle krok vlastni obsluhu? */ +export function findBuiltinStep(serviceId: string, operationId: string): Handler | undefined { + return handlers[`${serviceId}/${operationId}`]; +} diff --git a/src/runtime/context.ts b/src/runtime/context.ts new file mode 100644 index 0000000..12d7aca --- /dev/null +++ b/src/runtime/context.ts @@ -0,0 +1,31 @@ +/** + * Kdo prave bezi. + * + * Bez tohohle vznikne smycka: automatizace navazana na zmenu ticketu ticket + * zmeni, cimz se spusti znovu, a tak porad dokola. Prvni pokus o to skoncil + * tim, ze server prestal odpovidat. + * + * `AsyncLocalStorage` je na to spravny nastroj: bezi ctyri behy naraz a obycejna + * promenna by patrila vsem. Takhle ma kazdy beh svuj kontext, ktery se veze + * pres vsechna `await`. + */ + +import { AsyncLocalStorage } from 'node:async_hooks'; + +export interface RunMarker { + runId: string; + automationId: string; + tenantId: string; +} + +const storage = new AsyncLocalStorage(); + +/** Spusti praci s oznacenim, ze patri tomuhle behu. */ +export function withRun(marker: RunMarker, work: () => Promise): Promise { + return storage.run(marker, work); +} + +/** Ktery beh zpusobil to, co se prave deje. `undefined` = zmena od cloveka. */ +export function currentRun(): RunMarker | undefined { + return storage.getStore(); +} diff --git a/src/runtime/executor.ts b/src/runtime/executor.ts index 96fe355..f2248c9 100644 --- a/src/runtime/executor.ts +++ b/src/runtime/executor.ts @@ -28,6 +28,7 @@ import { renderTemplate } from '../data/templates.js'; import { appendTrace, type TraceInput } from '../data/ticketStore.js'; import { scriptIdFor } from '../scripts/lookup.js'; import { runScript } from '../scripts/runner.js'; +import { findBuiltinStep } from './builtinSteps.js'; /** * Hodnoty, na ktere jde ve strome odkazovat. @@ -53,6 +54,13 @@ export interface StepResult { kind: 'action' | 'condition'; label: string; ok: boolean; + /** + * false = opakovat nema smysl. + * + * Chybejici skript nebo neznama skupina za minutu existovat nezacne. + * Opakovat se ma jen to, co muze pominout: nedostupna sluzba, timeout. + */ + retryable?: boolean; /** U podminky, kterou vetvi se slo. */ branch?: 'yes' | 'no'; summary: string; @@ -63,6 +71,8 @@ export interface StepResult { export interface RunResult { ok: boolean; + /** false = beh nema smysl opakovat, chyba sama nepomine. */ + retryable: boolean; /** Kolik kroku se opravdu vykonalo, vcetne podminek. */ steps: StepResult[]; /** Kontext po behu, tedy i vystupy kroku. */ @@ -136,8 +146,13 @@ export async function runFlow( appendTrace(options.ticketId, results.map(toTrace)); } + const failed = results.find((step) => !step.ok); + return { ok: error === null, + // Kdyz krok nerekl jinak, opakovat se smi: vypadek cizi sluzby je + // nejcastejsi duvod chyby a ten pomine. + retryable: failed?.retryable !== false, steps: results, context: working, error, @@ -154,6 +169,60 @@ async function runAction( const startedAt = Date.now(); const label = `${step.serviceId}/${step.operationId}`; + const inputsForStep = fillTemplates(step.inputs ?? {}, context); + + /* + * Nejdriv nase vlastni kroky. Zalozit ticket nebo prehodit ho na cloveka + * neni volani cizi sluzby, takze to nejde pres skript - sahá to do naseho + * uloziste. Pro uzivatele je to v katalogu operace jako kazda jina. + */ + const builtin = findBuiltinStep(step.serviceId, step.operationId); + if (builtin) { + try { + const outcome = await builtin(inputsForStep, { + tenantId: options.tenantId, + ticketId: options.ticketId, + }); + + if (outcome.ticketId && !options.ticketId) { + // Krok zalozil ticket. Dalsi kroky uz vedi, ke kteremu patri. + options.ticketId = outcome.ticketId; + } + + context[step.id] = outcome.outputs; + for (const [key, value] of Object.entries(outcome.outputs)) { + context[`${step.id}.${key}`] = value; + // Vystup je i pod holym jmenem, aby se v sablone dalo psat + // `{{ticketId}}` misto `{{krok1.ticketId}}`. Pozdejsi krok prepise + // starsi, coz je to, co clovek ceka. + context[key] = value; + } + + return { + stepId: step.id, + kind: 'action', + label, + ok: outcome.ok, + summary: outcome.summary, + detail: outcome.detail ?? null, + // Vnitrni krok selhava na spatnem nastaveni, ne na vypadku. Opakovani + // by jen pettkrat zopakovalo tutéz chybu. + retryable: false, + durationMs: Date.now() - startedAt, + }; + } catch (err) { + return { + stepId: step.id, + kind: 'action', + label, + ok: false, + summary: 'krok skončil chybou', + detail: err instanceof Error ? err.message : String(err), + durationMs: Date.now() - startedAt, + }; + } + } + const scriptId = scriptIdFor(step.serviceId, step.operationId); if (!scriptId) { return { @@ -165,6 +234,8 @@ async function runAction( detail: `Služba ${step.serviceId} nemá pro operaci ${step.operationId} skript, ` + 'takže ji nelze vykonat. Zbytek stromu se nespustil.', + // Skript nezacne existovat sam od sebe, opakovat to nema smysl. + retryable: false, durationMs: Date.now() - startedAt, }; } @@ -187,8 +258,7 @@ async function runAction( }; } - const inputs = fillTemplates(step.inputs ?? {}, context); - const result = await runScript(scriptId, inputs, { + const result = await runScript(scriptId, inputsForStep, { connector: connector ?? null, // Klic je stabilni na dvojici beh a krok, takze opakovane odeslani // nevystavi druhou fakturu. diff --git a/src/runtime/queue.ts b/src/runtime/queue.ts new file mode 100644 index 0000000..354dc6c --- /dev/null +++ b/src/runtime/queue.ts @@ -0,0 +1,305 @@ +/** + * Fronta behu. + * + * Udalost se **nezpracuje hned**. Zapise se do fronty a odesilateli se odpovi, + * ze je prijata. Zpracuje ji worker. Duvody jsou tri a kazdy sam o sobe staci: + * + * - **Za konektory nerucime.** Cizi sluzba muze odpovidat pet sekund nebo + * byt hodinu mimo. Kdyby se cekalo v requestu, odesilateli vyprsi timeout + * a udalost je pryc, i kdyz jsme ji dostali. + * - **Pad procesu nesmi ztratit praci.** Zaznam ve fronte prezije restart, + * rozdelany beh v pameti ne. + * - **Opakovani.** Chyba cizi sluzby casto pomine sama. Bez fronty neni kam + * si poznamenat, ze se to ma za minutu zkusit znovu. + * + * Vyber dalsi prace je **spravedlivy po firmach**: jedna firma s tisicem + * udalosti nesmi zablokovat ostatni. Bez toho staci jeden rozbity e-shop + * a servisdesk stoji vsem. + * + * Uloziste je obecne (`defineStore`), takze to same funguje nad Postgresem + * i nad JSON souborem. Pri vic instancich je potreba `SELECT ... FOR UPDATE + * SKIP LOCKED`, viz `claimBatch` nize a documentation/10-runtime-a-kapacita.md. + */ + +import { randomUUID } from 'node:crypto'; +import { defineStore, nowIso, type TenantEntity } from '../data/store/index.js'; + +export type QueueStatus = 'pending' | 'running' | 'done' | 'failed'; + +/** Co beh spustilo. Podle toho se pozna, co je v `payload`. */ +export type TriggerKind = + /** Prislo na adresu webhooku automatizace. */ + | 'webhook' + /** Vznikl ticket. */ + | 'ticket.created' + /** Ticket se zmenil. */ + | 'ticket.updated' + /** Rucni spusteni z portalu. */ + | 'manual'; + +export interface QueueItem extends TenantEntity { + tenantId: string; + automationId: string; + trigger: TriggerKind; + /** Data spoustece. U webhooku telo requestu, u ticketu jeho udaje. */ + payload: Record; + /** Ticket, ke kteremu beh patri. Do jeho logu se zapisuje prubeh. */ + ticketId: string | null; + status: QueueStatus; + attempts: number; + /** Driv se to nezkusi. Prazdne = hned. */ + nextAttemptAt: string | null; + /** Cele hlaseni posledni chyby. Nezkracene, je to jedina stopa. */ + lastError: string | null; + /** Kdy si to worker vzal. Podle toho se pozna zaseknuty beh. */ + claimedAt: string | null; + finishedAt: string | null; + /** + * Klic proti dvojimu zarazeni. Dve stejne udalosti behem chvile znamenaji + * jeden beh, ne dva - odesilatele casto posilaji opakovane. + */ + dedupeKey: string | null; +} + +export const queueStore = defineStore('runQueue'); + +/** + * Kolikrat se beh zkusi, nez se vzda. + * + * Prodlevy rostou: minuta staci na kratky vypadek, hodina na delsi. Zkouset + * to pordad by u trvale rozbite sluzby znamenalo nekonecnou zatez. + */ +const BACKOFF_MS = [30_000, 2 * 60_000, 10 * 60_000, 60 * 60_000]; +export const MAX_ATTEMPTS = BACKOFF_MS.length + 1; + +/** Po jake dobe se bezici beh povazuje za zaseknuty a vrati se do fronty. */ +const STUCK_AFTER_MS = 10 * 60_000; + +/** Kopie fronty v pameti. Cte se pri kazdem kole workeru, meni se zridka. */ +let items: QueueItem[] = []; + +export async function initQueue(): Promise { + await queueStore.init(); + items = await queueStore.listAll(); + + // Beh, ktery byl rozdelany pri padu procesu, se vrati do fronty. Bez toho + // by zustal navzdy ve stavu `running` a nikdo by ho uz nevzal. + const stuck = items.filter((item) => item.status === 'running'); + for (const item of stuck) { + item.status = 'pending'; + item.claimedAt = null; + await save(item); + } + if (stuck.length > 0) { + console.warn(`[fronta] ${stuck.length} rozdelanych behu vraceno do fronty po restartu`); + } + + const pending = items.filter((item) => item.status === 'pending').length; + console.info(`[fronta] nactena, ceka ${pending} behu`); +} + +async function save(item: QueueItem): Promise { + item.updatedAt = nowIso(); + await queueStore.put(item); +} + +export interface EnqueueInput { + tenantId: string; + automationId: string; + trigger: TriggerKind; + payload?: Record; + ticketId?: string | null; + dedupeKey?: string | null; +} + +/** + * Zaradi beh. + * + * Vraci zaznam, nebo `null`, kdyz se stejny beh uz ve fronte ceka - to neni + * chyba, je to smysl klice proti dvojimu zarazeni. + */ +export async function enqueue(input: EnqueueInput): Promise { + if (input.dedupeKey) { + const waiting = items.find( + (item) => + item.dedupeKey === input.dedupeKey && + (item.status === 'pending' || item.status === 'running'), + ); + if (waiting) return null; + } + + const timestamp = nowIso(); + const item: QueueItem = { + id: `run_${randomUUID().slice(0, 12)}`, + tenantId: input.tenantId, + automationId: input.automationId, + trigger: input.trigger, + payload: input.payload ?? {}, + ticketId: input.ticketId ?? null, + status: 'pending', + attempts: 0, + nextAttemptAt: null, + lastError: null, + claimedAt: null, + finishedAt: null, + dedupeKey: input.dedupeKey ?? null, + createdAt: timestamp, + updatedAt: timestamp, + }; + + items.push(item); + await queueStore.put(item); + return item; +} + +/** + * Vezme dalsi praci. + * + * Spravedlive po firmach: z kazde firmy nejvys jeden beh na kolo, takze + * jedna firma s tisicem udalosti nezablokuje ostatni. V ramci firmy se bere + * nejstarsi. + * + * Pri vic instancich musi vyber probehnout v databazi (`FOR UPDATE SKIP + * LOCKED`), jinak si dva workeri vezmou tentyz beh. Tady je to v pameti + * jednoho procesu, coz pro jednu instanci staci a je to videt na jednom miste. + */ +export async function claimBatch(limit: number): Promise { + const now = Date.now(); + + // Zaseknuty beh se vraci do fronty, aby se neztratil kvuli spadlemu workeru. + for (const item of items) { + if (item.status !== 'running' || !item.claimedAt) continue; + if (now - new Date(item.claimedAt).getTime() < STUCK_AFTER_MS) continue; + console.warn(`[fronta] ${item.id} visel v behu moc dlouho, vraci se do fronty`); + item.status = 'pending'; + item.claimedAt = null; + await save(item); + } + + const ready = items + .filter((item) => item.status === 'pending') + .filter((item) => !item.nextAttemptAt || new Date(item.nextAttemptAt).getTime() <= now) + .sort((a, b) => a.createdAt.localeCompare(b.createdAt)); + + const claimed: QueueItem[] = []; + const perTenant = new Set(); + + for (const item of ready) { + if (claimed.length >= limit) break; + if (perTenant.has(item.tenantId)) continue; + perTenant.add(item.tenantId); + + item.status = 'running'; + item.claimedAt = nowIso(); + item.attempts += 1; + await save(item); + claimed.push(item); + } + + return claimed; +} + +export async function markDone(item: QueueItem): Promise { + item.status = 'done'; + item.finishedAt = nowIso(); + item.lastError = null; + await save(item); +} + +/** + * Beh selhal. + * + * Do vycerpani pokusu se vrati do fronty s rostouci prodlevou, pak skonci + * jako `failed` a zustane k nahlednuti. Nemazat: bez zaznamu by nikdo + * nezjistil, ze se neco nestalo. + */ +export async function markFailed( + item: QueueItem, + error: string, + retryable = true, +): Promise { + item.lastError = error; + + if (!retryable) { + // Chyba, ktera sama nepomine: spatne nastaveny krok, chybejici skupina. + // Zkouset to petkrat by jen petkrat zopakovalo tutéz hlasku. + item.status = 'failed'; + item.finishedAt = nowIso(); + console.error(`[fronta] ${item.id} skoncil, opakovat nema smysl: ${error}`); + await save(item); + return; + } + + if (item.attempts >= MAX_ATTEMPTS) { + item.status = 'failed'; + item.finishedAt = nowIso(); + console.error(`[fronta] ${item.id} se vzdal po ${item.attempts} pokusech: ${error}`); + } else { + const wait = BACKOFF_MS[Math.min(item.attempts - 1, BACKOFF_MS.length - 1)]; + item.status = 'pending'; + item.claimedAt = null; + item.nextAttemptAt = new Date(Date.now() + wait).toISOString(); + console.warn( + `[fronta] ${item.id} selhal (${item.attempts}. pokus), zkusi se za ${Math.round(wait / 1000)} s: ${error}`, + ); + } + + await save(item); +} + +export interface QueueStats { + pending: number; + running: number; + failed: number; + done: number; + /** Nejstarsi cekajici beh. Podle toho se pozna, ze fronta nestiha. */ + oldestPendingAt: string | null; +} + +export function queueStats(tenantIds?: string[]): QueueStats { + const visible = tenantIds ? items.filter((item) => tenantIds.includes(item.tenantId)) : items; + const pending = visible.filter((item) => item.status === 'pending'); + + return { + pending: pending.length, + running: visible.filter((item) => item.status === 'running').length, + failed: visible.filter((item) => item.status === 'failed').length, + done: visible.filter((item) => item.status === 'done').length, + oldestPendingAt: pending.reduce( + (acc, item) => (acc === null || item.createdAt < acc ? item.createdAt : acc), + null, + ), + }; +} + +/** Posledni behy, nejnovejsi nahore. Pro prehled v portalu. */ +export function recentRuns(tenantIds: string[], limit = 50): QueueItem[] { + return items + .filter((item) => tenantIds.includes(item.tenantId)) + .sort((a, b) => b.createdAt.localeCompare(a.createdAt)) + .slice(0, limit); +} + +/** + * Uklid hotovych behu. + * + * Hotove se drzi jen chvili, neuspesne dele - ty se jeste resi. Bez uklidu + * by fronta rostla do nekonecna, viz rozpocet v dokumentaci kapacity. + */ +export async function trimQueue(keepDone = 1_000, keepFailed = 500): Promise { + const done = items.filter((item) => item.status === 'done').sort((a, b) => a.createdAt.localeCompare(b.createdAt)); + const failed = items.filter((item) => item.status === 'failed').sort((a, b) => a.createdAt.localeCompare(b.createdAt)); + + const remove = [ + ...done.slice(0, Math.max(0, done.length - keepDone)), + ...failed.slice(0, Math.max(0, failed.length - keepFailed)), + ]; + if (remove.length === 0) return; + + const ids = new Set(remove.map((item) => item.id)); + items = items.filter((item) => !ids.has(item.id)); + for (const item of remove) { + await queueStore.remove(item.id, { tenantIds: [item.tenantId], includeGlobal: true }); + } + console.info(`[fronta] uklizeno ${remove.length} starych behu`); +} diff --git a/src/runtime/triggers.ts b/src/runtime/triggers.ts new file mode 100644 index 0000000..0ee8e8b --- /dev/null +++ b/src/runtime/triggers.ts @@ -0,0 +1,250 @@ +/** + * Co plni frontu. + * + * Tri druhy spoustecu a kazdy se chova jinak: + * + * - **Push**: cizi sluzba nam sama zavola webhook. Nejlevnejsi a nejrychlejsi, + * ale umi to jen sluzba, ktera webhooky ma. + * - **Vnitrni udalost**: neco se stalo u nas, typicky vznikl nebo se zmenil + * ticket. Automatizace na to muze navazat, aniz by kdokoliv volal zvenku. + * - **Pull**: sluzba webhooky nema, takze se ji musime **pravidelne ptat**. + * Tak funguje e-mail i vetsina schranek zprav. Dela to planovac nize. + * + * Rozhodovani, ktera automatizace se spusti, je **na jednom miste**. Kdyby + * bylo v kazde route zvlast, jedna by se casem chovala jinak nez druha. + */ + +import { config } from '../config.js'; +import { listAutomations, getAutomation, type AutomationDetail } from '../data/automationStore.js'; +import { listActiveTenants } from '../data/tenants.js'; +import type { Ticket } from '../data/ticketStore.js'; +import { currentRun } from './context.js'; +import { enqueue } from './queue.js'; + +/** + * Kolikrat smi jedna automatizace bezet nad jednim ticketem za minutu. + * + * Pojistka proti smycce, kterou nezachyti oznaceni behu: dve automatizace, + * kde kazda meni ticket a spousti tu druhou. Prekroceni se zaloguje, aby to + * nekdo poznal a spravil, misto aby se to tise zahazovalo. + */ +const MAX_RUNS_PER_TICKET = 5; +const WINDOW_MS = 60_000; + +const recentRuns = new Map(); + +function tooOften(automationId: string, ticketId: string): boolean { + const key = `${automationId}:${ticketId}`; + const now = Date.now(); + const times = (recentRuns.get(key) ?? []).filter((time) => now - time < WINDOW_MS); + + if (times.length >= MAX_RUNS_PER_TICKET) { + recentRuns.set(key, times); + return true; + } + + times.push(now); + recentRuns.set(key, times); + + // Mapa nesmi rust bez konce. Pri kazdem stem zapisu se stare klice zahodi. + if (recentRuns.size > 1_000) { + for (const [existing, stamps] of recentRuns) { + if (stamps.every((time) => now - time >= WINDOW_MS)) recentRuns.delete(existing); + } + } + + return false; +} + +/** + * Spoustece, ktere se musi obvolavat, protoze nam sluzba sama nezavola. + * + * Klic je `serviceId/operationId` z katalogu. Hodnota je vychozi perioda + * v sekundach - jak casto se ma sluzby zeptat, jestli je neco noveho. + * + * Perioda je kompromis: kratsi znamena rychlejsi reakci a vic volani cizi + * sluzby, ktera ma casto limit. Minuta u posty je bezne, u zprav se cekat + * nechce, u reportu staci hodina. + */ +const POLLED: Record = { + 'email/received': 60, + 'messenger/message': 30, + 'whatsapp/message': 30, + 'scheduler/every': 60, +}; + +/** Je tenhle spoustec typu "musime se ptat sami"? */ +export function isPolled(serviceId: string, operationId: string): boolean { + return `${serviceId}/${operationId}` in POLLED; +} + +function periodOf(automation: AutomationDetail): number { + const trigger = automation.flow.trigger; + if (!trigger) return 0; + + const key = `${trigger.serviceId}/${trigger.operationId}`; + const base = POLLED[key] ?? 0; + + /* + * Uzivatel si periodu muze zkratit nebo prodlouzit polem `intervalSec` + * u spoustece. Pod deset sekund se nejde: to uz neni dotazovani, to je + * utok na cizi sluzbu. + */ + const own = Number((trigger as { intervalSec?: unknown }).intervalSec); + if (Number.isFinite(own) && own >= 10) return own; + return base; +} + +/** Kdy se naposledy ptalo. V pameti procesu, po restartu se zepta hned. */ +const lastPolledAt = new Map(); + +/** + * Zaradi behy automatizaci, ktere reaguji na vnitrni udalost. + * + * Vola se po vzniku a po zmene ticketu. Necekana se: kdyby zarazeni do fronty + * shodilo ulozeni ticketu, rozbita fronta by rozbila ticketovaci system. + */ +export function onTicketEvent(kind: 'ticket.created' | 'ticket.updated', ticket: Ticket): void { + const matching = listAutomations([ticket.tenantId]).filter((automation) => { + if (!automation.enabled) return false; + const detail = getAutomation(automation.id, [ticket.tenantId]); + const trigger = detail?.flow.trigger; + if (!trigger || trigger.serviceId !== 'ticket') return false; + + // `changed` reaguje na obojí, aby se nemusely delat dve skoro stejne. + if (trigger.operationId === 'changed') return true; + return trigger.operationId === (kind === 'ticket.created' ? 'created' : 'updated'); + }); + + const run = currentRun(); + + for (const automation of matching) { + /* + * Zmena, kterou udelala tatáz automatizace, ji nesmi spustit znovu. + * Prirazeni resitele meni ticket, takze bez tohohle by se automatizace + * volala do nekonecna. + */ + if (run?.automationId === automation.id) continue; + + if (tooOften(automation.id, ticket.id)) { + console.warn( + `[spoustec] ${automation.id} uz bezela nad ${ticket.id} ${MAX_RUNS_PER_TICKET}x za minutu, ` + + 'preskakuji - vypada to na smycku mezi automatizacemi', + ); + continue; + } + + void enqueue({ + tenantId: ticket.tenantId, + automationId: automation.id, + trigger: kind, + ticketId: ticket.id, + payload: ticketPayload(ticket), + /* + * Klic drzi jeden beh na jednu zmenu ticketu. Bez nej by pet zmen behem + * sekundy znamenalo pet behu, ktere delaji totez - a u prirazeni resitele + * by se prace rozdelila petkrat. + */ + dedupeKey: `${automation.id}:${ticket.id}:${ticket.updatedAt}`, + }).catch((err: unknown) => { + console.error(`[spoustec] ${automation.id} se nepodarilo zaradit:`, err); + }); + } +} + +/** + * Udaje ticketu tak, jak je vidi strom. + * + * Jmena bez teckove cesty, aby se v sablonach psalo `{{subject}}`, ne + * `{{ticket.subject}}`. Vlastni pole typu jsou na stejne urovni, protoze + * pro uzivatele je to tentyz druh udaje. + */ +export function ticketPayload(ticket: Ticket): Record { + return { + ticketId: ticket.id, + externalId: ticket.externalId, + subject: ticket.subject, + body: ticket.body, + status: ticket.status, + priority: ticket.priority, + channel: ticket.channel, + typeId: ticket.typeId, + stage: ticket.stage, + tags: ticket.tags, + assigneeId: ticket.assignee?.id ?? null, + assigneeGroupId: ticket.assigneeGroupId, + company: ticket.customer.company, + contact: ticket.customer.contact, + reply: ticket.customer.reply, + ...(ticket.fields ?? {}), + }; +} + +/** + * Planovac: zaradi behy automatizaci, ktere se musi ptat samy. + * + * Nedela ten dotaz sam. Jen rekne "je cas", zaradi beh a **dotaz je prvni krok + * stromu**. Diky tomu se dotazovani chova stejne jako cokoliv jineho: ma svuj + * zaznam v logu, opakuje se pri chybe a jde ho zmenit bez zasahu do kodu. + */ +export async function planPolled(): Promise { + const now = Date.now(); + let planned = 0; + + for (const tenant of listActiveTenants()) { + for (const summary of listAutomations([tenant.id])) { + if (!summary.enabled) continue; + + const automation = getAutomation(summary.id, [tenant.id]); + const trigger = automation?.flow.trigger; + if (!automation || !trigger) continue; + if (!isPolled(trigger.serviceId, trigger.operationId)) continue; + + const period = periodOf(automation) * 1_000; + if (period <= 0) continue; + + const last = lastPolledAt.get(automation.id) ?? 0; + if (now - last < period) continue; + lastPolledAt.set(automation.id, now); + + const item = await enqueue({ + tenantId: tenant.id, + automationId: automation.id, + trigger: 'manual', + payload: { plannedAt: new Date(now).toISOString() }, + /* + * Klic drzi jeden cekajici dotaz na automatizaci. Kdyz predchozi jeste + * bezi, protoze sluzba odpovida pomalu, dalsi se nezaradi - jinak by + * se fronta zaplnila dotazy na sluzbu, ktera stejne nestiha. + */ + dedupeKey: `poll:${automation.id}`, + }); + if (item) planned += 1; + } + } + + return planned; +} + +let timer: NodeJS.Timeout | null = null; + +/** Spusti planovac. Vola se pri startu, kdyz je proces zaroven workerem. */ +export function startScheduler(): void { + if (timer) return; + + const every = Math.max(config.schedulerIntervalSec, 5) * 1_000; + timer = setInterval(() => { + void planPolled().catch((err: unknown) => { + // Rozbity planovac nesmi shodit proces. Za chvili to zkusi znovu. + console.error('[planovac] selhal:', err); + }); + }, every); + timer.unref?.(); + + console.info(`[planovac] spusten, kontrola kazdych ${Math.round(every / 1000)} s`); +} + +export function stopScheduler(): void { + if (timer) clearInterval(timer); + timer = null; +} diff --git a/src/runtime/worker.ts b/src/runtime/worker.ts new file mode 100644 index 0000000..c9833e6 --- /dev/null +++ b/src/runtime/worker.ts @@ -0,0 +1,239 @@ +/** + * Worker: bere praci z fronty a vykonava ji. + * + * Bezi ve stejnem procesu jako API, ale **mimo request**. Webhook zapise + * udalost do fronty a hned odpovi; co se s ni stane dal, uz odesilatele + * nezajima a hlavne na to necekal. + * + * Az bude potreba vic vykonu, spusti se tentyz obraz s `WORKER_ONLY=1` jako + * druhy proces a API se od vykonu oddeli. Do te doby je to jeden proces, coz + * je pro jednu instanci spravne - dva procesy kvuli deseti behum za minutu + * jsou jen dve veci, ktere mohou spadnout. + */ + +import { config } from '../config.js'; +import { getAutomation, recordRun } from '../data/automationStore.js'; +import { createIncident, findOpenIncident } from '../data/incidentStore.js'; +import { publish } from '../events/bus.js'; +import { withRun } from './context.js'; +import { runFlow, type RunResult } from './executor.js'; +import { + claimBatch, + markDone, + markFailed, + trimQueue, + MAX_ATTEMPTS, + type QueueItem, +} from './queue.js'; + +/** Kolik behu naraz. Kroky cekaji na cizi sluzby, procesor se skoro nepouzije. */ +const CONCURRENCY = 4; + +/** Jak casto se fronta kontroluje, kdyz je prazdna. */ +const IDLE_MS = 1_000; + +/** Jak casto se uklizi hotove behy. */ +const CLEANUP_EVERY = 200; + +let running = false; +let stopping = false; +let rounds = 0; + +/** + * Spusti workera. + * + * Neceka se na nej - bezi na pozadi po celou dobu behu procesu. Chyba jednoho + * behu nesmi workera zabit, proto je cela smycka v try/catch. + */ +export function startWorker(): void { + if (running) return; + running = true; + stopping = false; + console.info(`[worker] spusten, ${CONCURRENCY} behu naraz`); + void loop(); +} + +export async function stopWorker(): Promise { + stopping = true; +} + +async function loop(): Promise { + while (!stopping) { + let claimed: QueueItem[] = []; + try { + claimed = await claimBatch(CONCURRENCY); + } catch (err) { + // Nedostupne uloziste nesmi workera zabit. Za chvili to zkusi znovu. + console.error('[worker] frontu se nepodarilo precist:', err); + } + + if (claimed.length === 0) { + await sleep(IDLE_MS); + continue; + } + + // Behy jdou soubezne, cekaji stejne na cizi sluzby. + await Promise.all(claimed.map((item) => runOne(item))); + + rounds += 1; + if (rounds % CLEANUP_EVERY === 0) { + await trimQueue().catch((err: unknown) => console.error('[worker] uklid selhal:', err)); + } + } + + running = false; + console.info('[worker] zastaven'); +} + +async function runOne(item: QueueItem): Promise { + const automation = getAutomation(item.automationId, [item.tenantId]); + + if (!automation) { + // Automatizace mezitim zmizela. Neni co delat a opakovat to nema smysl. + await markFailed(item, `Automatizace ${item.automationId} neexistuje.`); + return; + } + + if (!automation.enabled) { + await markFailed(item, `Automatizace ${automation.name} je pozastavená.`); + return; + } + + try { + /* + * Beh se oznaci, aby zmeny, ktere udela, nespustily tutéz automatizaci + * znovu. Bez toho vznikne smycka: automatizace zmeni ticket, zmena + * ticketu ji spusti, a tak porad. + */ + const result = await withRun( + { runId: item.id, automationId: item.automationId, tenantId: item.tenantId }, + () => + runFlow(automation.flow.steps, item.payload, { + tenantId: item.tenantId, + ticketId: item.ticketId, + // Klic je stabilni na beh, takze opakovany pokus nevystavi druhou fakturu. + idempotencyKey: `run:${item.id}`, + }), + ); + + recordRun(automation.id, result.ok); + + if (result.ok) { + await markDone(item); + publish('automation.run', `Automatizace ${automation.name} proběhla`, { + automationId: automation.id, + ok: true, + }); + return; + } + + const failed = result.steps.find((step) => !step.ok); + // Cele hlaseni, vcetne toho, co sluzba vratila. Zkratit ho tady by + // znamenalo, ze se pricina uz nikde nedozvime. + const message = [result.error, failed?.detail].filter(Boolean).join('\n'); + // Marnou chybu (chybejici skript, spatne nastaveny krok) nema smysl + // opakovat - petkrat by se zopakovala tatáz hlaska. + await markFailed(item, message || 'Běh selhal bez hlášení.', result.retryable); + + /* + * Kazda chyba, kterou uz nema smysl zkouset, zaklada incident. + * + * Dve urovne zamerne: `title` a `impact` cte klient a musi z toho poznat, + * co to pro nej znamena. `detail` cte admin a je v nem vsechno - cely beh, + * ktery krok, co prislo na vstupu. + * + * Zaklada se az kdyz uz se to nebude opakovat. Kdyby vznikl hned, mel by + * klient incident u kazdeho vypadku site, ktery se sam za minutu spravi. + */ + if (!result.retryable || item.attempts >= MAX_ATTEMPTS) { + reportIncident(item, automation.name, result, failed, message); + } + + publish('automation.run', `Automatizace ${automation.name} skončila chybou`, { + automationId: automation.id, + ok: false, + }); + } catch (err) { + // `runFlow` chyby nevyhazuje, tohle je posledni pojistka. + console.error(`[worker] neocekavana chyba behu ${item.id}:`, err); + await markFailed(item, err instanceof Error ? err.message : String(err)); + } +} + +/** + * Zalozi incident z neuspesneho behu. + * + * Stejna pricina nezaklada druhy incident, dokud je ten prvni otevreny - + * jinak by deset stejnych chyb znamenalo deset incidentu a nikdo by se v tom + * nevyznal. + */ +function reportIncident( + item: QueueItem, + automationName: string, + result: RunResult, + failed: RunResult['steps'][number] | undefined, + message: string, +): void { + const source = `automation:${item.automationId}:${failed?.stepId ?? 'flow'}`; + if (findOpenIncident(source)) return; + + const detail = [ + `Automatizace: ${automationName} (${item.automationId})`, + `Běh: ${item.id}, spouštěč ${item.trigger}, pokusů ${item.attempts}`, + item.ticketId ? `Ticket: ${item.ticketId}` : null, + failed ? `Krok, který selhal: ${failed.label} (${failed.stepId})` : null, + '', + 'Hlášení:', + message || '(bez hlášení)', + '', + 'Kroky běhu:', + ...result.steps.map( + (step) => ` ${step.ok ? 'ok ' : 'CHYBA'} ${step.label}: ${step.summary}`, + ), + '', + 'Data na vstupu:', + JSON.stringify(item.payload, null, 2), + ] + .filter((line) => line !== null) + .join('\n'); + + createIncident({ + tenantId: item.tenantId, + // Klient nema cist nazev kroku ani ID behu. Ma poznat, co nefunguje. + title: `Automatizace ${automationName} neproběhla`, + service: failed ? failed.label.split('/')[0] : 'Automatizace', + severity: 'sev3', + impact: clientMessage(item, failed), + detail, + source, + ticketId: item.ticketId, + }); +} + +/** + * Vysvetleni pro klienta. + * + * Zamerne bez technickych podrobnosti: klient potrebuje vedet, co se nestalo + * a jestli s tim ma neco delat. Cele hlaseni je v `detail` pro admina. + */ +function clientMessage(item: QueueItem, failed: RunResult['steps'][number] | undefined): string { + const what = item.ticketId ? `u ticketu ${item.ticketId}` : 'na pozadí'; + const where = failed ? ` Zastavilo se to na kroku ${failed.label}.` : ''; + return ( + `Automatizace ${what} nedoběhla do konce, takže část kroků se neprovedla.${where} ` + + 'Data jsme neztratili, po opravě je možné běh zopakovat.' + ); +} + +function sleep(ms: number): Promise { + return new Promise((resolve) => { + const timer = setTimeout(resolve, ms); + // Cekani nesmi drzet proces nazivu pri ukonceni. + timer.unref?.(); + }); +} + +/** Je worker v tomhle procesu zapnuty? Pro diagnostiku v `/health/ready`. */ +export function workerEnabled(): boolean { + return config.workerEnabled; +} diff --git a/web/src/components/dashboard/DashboardLayout.tsx b/web/src/components/dashboard/DashboardLayout.tsx index d2579c3..e3c6efb 100644 --- a/web/src/components/dashboard/DashboardLayout.tsx +++ b/web/src/components/dashboard/DashboardLayout.tsx @@ -74,6 +74,15 @@ function accessLabel(user: AuthUser | null, access: Access | null): string { function DashboardShell() { const { user, logout } = useAuth(); const access = useApiQuery('/api/dashboard/access'); + + /* + * Cislo u zalozky Tickety a upozorneni. Obnovuje se na udalost, ne casovacem - + * kdyz nekomu prijde ticket, ma to videt hned, ne za minutu. + */ + const notifications = useApiQuery<{ unread: number; mine: number }>( + '/api/dashboard/notifications', + { refetchOn: ['notification.created', 'ticket.assigned', 'ticket.resolved'] }, + ); const navigate = useNavigate(); const location = useLocation(); const [sidebarOpen, setSidebarOpen] = useState(false); @@ -128,7 +137,19 @@ function DashboardShell() { } > - {entry.label} + {entry.label} + {/* + Kolik ticketu ma prihlaseny u sebe. Nula se nekresli - prazdny + odznak je jen sum. + */} + {entry.key === 'tickets' && (notifications.data?.mine ?? 0) > 0 && ( + + {notifications.data?.mine} + + )} ); })} diff --git a/web/src/components/dashboard/EventToasts.tsx b/web/src/components/dashboard/EventToasts.tsx index ec6acd1..e27d10c 100644 --- a/web/src/components/dashboard/EventToasts.tsx +++ b/web/src/components/dashboard/EventToasts.tsx @@ -1,4 +1,4 @@ -import { AlarmClock, CheckCircle2, LifeBuoy, UserCheck, Webhook, Workflow, X } from 'lucide-react'; +import { AlarmClock, Bell, CheckCircle2, LifeBuoy, UserCheck, Webhook, Workflow, X } from 'lucide-react'; import { useEffect, useState } from 'react'; import type { LucideIcon } from 'lucide-react'; import { useEventStream } from '@/components/dashboard/EventStreamProvider'; @@ -22,6 +22,7 @@ const style: Record = { 'automation.deleted': { icon: Workflow, tone: 'text-white/50 bg-white/5' }, 'automation.run': { icon: Workflow, tone: 'text-ok-400 bg-ok-500/12' }, 'webhook.received': { icon: Webhook, tone: 'text-accent-300 bg-accent-500/12' }, + 'notification.created': { icon: Bell, tone: 'text-brand-300 bg-brand-500/12' }, }; /** Bubliny o tom, co se prave stalo. Bez nich by zive zmeny nebyly poznat. */ diff --git a/web/src/components/dashboard/flow/TriggerConfig.tsx b/web/src/components/dashboard/flow/TriggerConfig.tsx index 599ac9e..d8079fd 100644 --- a/web/src/components/dashboard/flow/TriggerConfig.tsx +++ b/web/src/components/dashboard/flow/TriggerConfig.tsx @@ -42,6 +42,7 @@ export function TriggerConfig({ @@ -151,6 +152,21 @@ function CustomFields({ ))} + {/* + Cesta v tele. Prazdna = hodnota lezi primo pod nazvem. + Bez toho by slo napojit jen ploche telo a odesilatel + s vnorenym modelem by se musel prizpusobovat nam. + */} + + updateField(field.id, { path: event.target.value.trim() || undefined }) + } + placeholder="cesta v těle, např. data.order.id" + title="Kde hodnota v přijatém JSONu je. Prázdné = přímo pod názvem." + className="min-w-40 flex-1 rounded-lg border border-ink-600/70 bg-ink-900/70 px-2.5 py-1.5 font-mono text-xs text-white placeholder:text-white/25 focus:border-brand-400/70 focus:outline-none" + /> +