# 10 - Runtime, vykonna cast a kapacita Navrh, ne popis stavu. Runtime neexistuje, dnes se ulozeny strom nevykonava. Souvisejici navrh datovych modelu je v [09-navrh-rozsireni.md](09-navrh-rozsireni.md). Tenhle soubor odpovida na tri veci: jak se vyhodnocuji kroky, jak se prijimaji udalosti a co to znamena pri 150 klientech. ## Nejdriv cisla, pak architektura Zadani: 150 klientu, kazdy asi 5 systemu, z nich chodi radove desitky udalosti. To je 750 napojeni. "Desitky udalosti" ma dve cteni a **odpoved se mezi nimi podstatne lisi**, takze obe: | Scenar | Desitky udalosti za | Udalosti/den | Kroku/den | Prumer | Spicka | | ------ | ------------------- | ------------ | ---------- | ------- | -------- | | A | den a system | 37 tisic | 375 tisic | 4 kr/s | 30-60/s | | B | hodinu a system | 450 tisic | 4,5 mil | 52 kr/s | 200-400/s| Pocitano s 10 kroky na udalost, coz je stredni automatizace. Spicka vychazi z toho, ze provoz je v osmihodinovem okne a uvnitr nerovnomerny, tedy radove osmkrat nad prumerem. **Planovat se musi na kroky, ne na udalosti.** Jedna udalost s trisetkrokovym stromem stoji tristakrat vic nez udalost s jednim krokem. Az bude runtime bezet, je metrika kroku za sekundu ta jedina, podle ktere se da neco rict. ### Co je a co neni uzke misto | Vec | Scenar A | Scenar B | | ----------------------- | --------------- | ------------------------------ | | Fronta v Postgresu | par procent | zvladne, ale s davkovym odberem| | Soubezne HTTP volani | 15 soubezne | 90 soubezne, Node se nezapoti | | **Zapis `run_step`** | 22 GB/mesic | **270 GB/mesic, nutne zkratit**| | Limity cizich API | uzke misto | uzke misto | Fronta nad Postgresem s `FOR UPDATE SKIP LOCKED` uklidne obslouzi radove 200 az 500 uloh za sekundu na jednom uzlu, kdyz se odebira davkove. Scenar A je tedy nezajimavy a scenar B je v pohodlnem pasmu. **Skutecne uzke misto je objem zapisu.** Radek `run_step` s vstupem a vystupem v JSONB ma realne 1 az 3 kB. Pri scenari B je to 9 GB denne, coz za pul roku nikdo neuklidi. Reseni je v sekci o retenci a je to jedina vec z celeho navrhu, kterou **nelze odlozit na potom**. Druhe uzke misto je za nasimi hranicemi. iDoklad, WhatsApp ani ekonomicky system nesnesou desitky pozadavku za sekundu na jeden ucet. Skalovani workeru bez limitu za napojeni znamena jen rychleji dojit k odpovedi 429. ### Co pri teto velikosti nepotrebujeme Rict to nahlas, aby se to nestavelo: **zadna Kafka, zadny Redis, zadny Kubernetes, zadne sharding.** 150 klientu je pro jeden Postgres a par procesu Node maly provoz. Usili patri do idempotence, spravedlnosti mezi klienty a pozorovatelnosti, ne do infrastruktury. Odhad velikosti: | Scenar | Postgres | Workeri | Kde to bezi | | ------ | ------------------ | --------------------------- | ----------- | | A | 4 vCPU, 16 GB | 2 procesy, 50 soubezne | jeden stroj | | B | 8-16 vCPU, 32 GB, NVMe | 4-6 procesu, 100 soubezne | dva stroje | ## Moznosti u databaze | Varianta | Verdikt pri 150 klientech | | -------------------------------- | ------------------------------------------------ | | Jedno DB, `tenant_id` ve sloupci | **ano, tohle** | | Schema na klienta | ne: 150 x 20 tabulek je 3000 tabulek, migrace se stanou nespolehlivymi | | Databaze na klienta | ne, ale nechat si dvere otevrene | | Partitionovani podle klienta | ne, oddily by byly velikostne nesouvisle | | Partitionovani podle casu | **ano, u pripisovacich tabulek** | ### Dvere k oddelene databazi za par korun Jednou prijde klient, ktery bude chtit vlastni databazi, nebo bude delat tricet procent provozu. Aby to pak nebyla prestavba, staci **jedna vec od zacatku**: pristup k poolu jen pres `dbFor(tenantId)`, i kdyz zpocatku vraci porad tentyz pool. K tomu **nikdy nespojovat dotazem dva klienty**, coz uz vynucuje povinny argument `tenantIds` v ulozistich. Splneni tehle dvou podminek znamena, ze presun jednoho klienta do vlastni databaze je konfigurace, ne prepisovani dotazu. ### Pool a jedno pravidlo, na kterem to stoji nebo pada **Worker nesmi drzet spojeni do databaze po dobu volani ciziho API.** Je to nejcastejsi zpusob, jak takovy system umre. Volani do iDokladu trva 300 ms. Kdyz drzi spojeni, znamena 90 soubeznych kroku 90 obsazenych spojeni a pool skonci. Spravne poradi: ``` transakce: odeber ulohu (2 ms) - spojeni drzim uvolni spojeni volani ciziho API (300 ms) - spojeni nedrzim transakce: zapis vysledek (2 ms) - spojeni drzim znovu ``` Pri tomhle poradi staci na 90 soubeznych kroku 3 az 5 spojeni. Bez nej 90. Az bude workeru vic, prijde PgBouncer v transakcnim rezimu. **Pozor: v transakcnim rezimu nefunguje `LISTEN/NOTIFY`**, a prave na nem ma podle [09-navrh-rozsireni.md](09-navrh-rozsireni.md) stat sbernice udalosti pro SSE. Ta potrebuje prime spojeni mimo PgBouncer. Zjistit to az pri nasazeni znamena rozbity zivy dashboard. ## Cesta udalosti: tri oddelene faze ``` POST /webhook/:token -> event_inbox 202, rychle a hloupe | dispatcher udalost na N behu | worker krok po kroku ``` Rozdeleni na tri faze neni akademicke. Kazda ma jinou vlastnost: prijem musi byt rychly, rozeslani musi byt idempotentni, vykonavani musi byt prerusitelne. ### 1. Prijem: rychle a hloupe ```sql event_inbox(id, tenant_id, token_id, source_kind, source_id, idempotency_key, payload jsonb, received_at, status, dispatched_at, cause_run_id, depth) unique index on (token_id, idempotency_key) ``` Webhook **nesmi vyhodnocovat strom**. Overi token, overi velikost, zapise jeden radek, vrati 202. Cil je p99 pod 50 ms. Duvod je praktickeho razu: odesilatel pri timeoutu opakuje. Kdyz webhook ceka na iDoklad, pomaly iDoklad zpusobi, ze tataz objednavka prijde tri krat. - **Idempotence pri prijmu.** Hlavicka `Idempotency-Key`, nebo hash tela, kdyz ji odesilatel neposila. Unikatni index nad `(token_id, idempotency_key)` v okne 24 hodin. Duplikat vrati 202 a stejne ID udalosti, ne chybu - pro odesilatele to je uspech, protoze jeho udalost je prijata. - **Limit za token.** Jeden rozbity klient ve smycce nesmi zaplnit inbox. - **Strop velikosti tela**, radove 256 kB, vetsi 413. - **Backpressure.** Kdyz hloubka fronty prekroci hranici, vracet 429 tokenum, ktere nejsou oznacene jako kriticke. Odesilatele 429 umi, na rozdil od tiche latence rostouci do minut. ### 2. Vlastni udalosti nechodi pres HTTP Podle zadani budou udalosti vznikat volanim na webhook z jinych automatizaci. U cizich odesilatelu ano. **U nasich vlastnich ne.** Volat vlastni HTTP endpoint na sebe pridava latenci, obsazuje spojeni a zaklada poruchu, ktera nemusi existovat. Vnitrni udalost zapise radek do `event_inbox` primo, prochazi tim samym dispatcherem a je v tom samem prehledu. Webhook zustava pro to, co prichazi zvenci. ### 3. Ochrana proti smycce musi projit skrz udalost Tohle je nejvaznejsi dusledek toho, ze automatizace vyrabeji udalosti pro jine automatizace. `depth` na behu chrani jen **uvnitr jednoho behu**. Kdyz automatizace A vyrobi udalost, ktera spusti B, a B vyrobi udalost, ktera spusti A, tak kazdy jednotlivy beh ma hloubku 1 a kontrola nikdy nezasahne. Smycka pobezi, dokud ji nekdo nevsimne na uctu za cizi API. Proto `event_inbox` nese `cause_run_id` a `depth`, a plati: ``` depth nove udalosti = depth behu, ktery ji vyrobil, + 1 depth > 5 -> udalost se odmitne, zapise se do logu a upozorni se ``` Bez tohohle jednoho sloupce je navrh z bodu 9 generator nekonecnych smycek. ### 4. Rozeslani Dispatcher najde automatizace, ktere na dvojici klient a spoustec sedi, a zalozi beh pro kazdou. - **Filtr na spousteci se vyhodnoti tady**, jeste pred zalozenim behu. Je to odpoved na otevrene rozhodnuti z konce [06-tickety.md](06-tickety.md) a pri tomhle objemu to neni kosmetika: beh, ktery hned skonci, stejne zaplati zapis do `run`, `run_step` i `ticket_trace`. - **Jedna transakce**: oznac udalost jako rozeslanou, zaloz behy, zaloz prvni ulohy. Bud vse, nebo nic. - **Unikatni index nad `(event_id, automation_id)`.** Kdyz dispatcher padne uprostred, opakovane rozeslani nezalozi druhy beh. ## Vykonna cast: co ridi prechod mezi kroky Dve veci, oddelene. **1. Cista funkce.** `next(tree, path, outputs): FlowPath | null` rozhodne, ktery krok je dalsi. Zadny stav, zadne IO, testovatelne. Podminka se vyhodnoti nad `outputs` a vybere vetev. Chuze po strome uz z poloviny existuje ve `web/src/lib/flow.ts` a `src/data/flowScope.ts`. **2. Radek v tabulce.** Fakticky prubeh behu drzi zaznam ulohy v databazi, ne pamet procesu. **Nikdy `setTimeout`, nikdy dlouhy retez promisu, nikdy rekurze drzici cely beh.** Restart containeru je bezna vec a beh ho musi prezit. ```sql run(id, tenant_id, source_kind, source_id, source_version, event_id, trigger_type, status, depth, queue_key, started_at, finished_at, error) run_step(id, run_id, path, step_id, attempt, status, input jsonb, output jsonb, error, started_at, finished_at) job(id, run_id, tenant_id, next_path, run_after, attempts, queue_key, locked_by, locked_until, priority) ``` `source_kind` a `source_version` rikaji, jestli beh patri automatizaci nebo akci a podle ktere verze jeji definice se ma dokoncit. To je to, co dela z cekaciho kroku fungujici vec: beh cekajici tyden dobehne podle stromu, ktery platil pri jeho spusteni. `run_after` v tabulce `job` je zaroven **cele cekani z bodu 8**. Krok `wait` neni v runtimu vyjimka, je to obycejny krok, ktery misto volani sluzby nastavi `run_after` a skonci. Proto ten bod skoro nic nestoji. ### Smycka workeru ``` 1. transakce: odeber DAVKU uloh SELECT ... FROM job WHERE run_after <= now() AND (locked_until IS NULL OR locked_until < now()) AND tenant_id <> ALL (:klienti_na_stropu) ORDER BY priority, run_after FOR UPDATE SKIP LOCKED LIMIT 25 UPDATE job SET locked_by = :worker, locked_until = now() + interval '2 min' 2. uvolni spojeni, vykonej ulohy soubezne, kazda PRAVE JEDEN krok 3. za kazdou ulohu transakce: zapis vysledek kroku smaz hotovou ulohu zarad dalsi ulohu podle next() ``` **Krok 3 v jedne transakci je cely trik.** Bud se zapise vysledek i dalsi uloha, nebo nic. Nikdy nevznikne beh, ktery ma hotovy krok a nema pokracovani, ani dvakrat zarazeny stejny krok. **Davkovy odber je nejvetsi pacidlo na propustnost.** Po jedne uloze znamena pri 300 krocich za sekundu 300 odberovych transakci za sekundu. Po dvaceti peti je jich dvanact. Je to jedna zmena `LIMIT` a nekolikanasobne mensi zatez. `SKIP LOCKED` znamena, ze workeru muze byt libovolne mnoho a nepotrebuji koordinatora. Za zvazeni stoji `graphile-worker` - je to tentyz princip nad Postgresem, overeny provozem. Rucne az kdyz bude potreba spravedlnost podle `queue_key`, kterou hotova knihovna neresi. ### Spravedlnost mezi klienty Ciste FIFO znamena, ze jeden vecerni import u jednoho klienta zastavi ostatnich 149. Pri 150 klientech to neni hypoteza, je to otazka casu. Prakticky pouzitelna verze je dvojice: - **Semafor za klienta ve workeru**: nejvyse N soubeznych kroku na klienta. - **Odberovy dotaz preskoci klienty na stropu** (`tenant_id <> ALL (...)`). Seznam si worker drzi sam a je aktualni na jednu davku. Presna spravedlnost cistym SQL je slozita a nevyplati se. Tohle je odhadem o dva rady jednodussi a rozdil nikdo nepozna. K tomu dve dalsi hranice: - **Limit a rychlostni strop za napojeni, ne za konektor.** Kvota je na uctu klienta v iDokladu, ne na tom, ze iDoklad existuje. Zetonovy kosik jednim `UPDATE ... RETURNING` nad radkem napojeni. - **Serializace nad jednim ticketem.** `queue_key = ticket:` a jen jedna bezici uloha na klic. Bez toho dve automatizace prepisuji stav teze veci a poradi neni dane. ### Pady - **Lease.** Worker padne uprostred kroku, `locked_until` vyprsi, ulohu si vezme jiny. Nic se neztrati. - **Kroky musi byt idempotentni.** Kazdy krok dostane `idempotencyKey = runId + ':' + path`, **stabilni pres vsechny pokusy**. Konektory, ktere umi `Idempotency-Key`, ho dostanou a druhy pokus nevystavi druhou fakturu. Klic za pokus by byl k nicemu, o tom to cele je. - **Retry.** Exponencialni backoff s jitterem, radove 5 pokusu. Rozlisit opakovatelne (timeout, spojeni, 429, 5xx) od koncovych (400, 401, 403, validace). Koncovou chybu neopakovat, jen se tim vypali kvota. - **Po vycerpani pokusu** dostane beh stav `failed`, zapise se do logu ticketu a vznikne udalost. Beh zustane a **lze ho pokracovat od padleho kroku**, protoze stav je per krok, ne per beh. - **Limit kroku na beh** a globalni timeout behu. - **Exactly-once neexistuje.** Cil je at-least-once plus idempotence. Kdo slibi exactly-once, jen jeste nenasel pripad, kdy to nedrzi. ## Kde bezi skripty Skripty z bodu 9 jsou cisty prevod dat, ale i tak maji vlastni provozni pravidlo, a je dulezite: **Skript nesmi bezet v hlavnim vlakne workeru.** Skript, ktery pocita pet set milisekund, zablokuje smycku udalosti a s ni **vsechny ostatni soubezne kroky toho workeru**. Jeden nepovedeny cyklus u jednoho klienta tim zastavi provoz vsech ostatnich, a v logu to vypada jako pomala cizi API. Navrh: - Bazen `worker_threads`, v kazdem `isolated-vm`. Radove tolik vlaken, kolik je jader. - Tvrdy timeout 50 az 200 ms a **zabiti vlakna** pri prekroceni, ne zdvorile preruseni. Prerusit smycku `while (true)` jinak nejde. - Zkompilovany skript se drzi v cache za verzi, kontext se po N spustenich zahodi kvuli unikum pameti. - Rezie kontextu je radove milisekunda, takze tisic skriptovych kroku za sekundu neni problem. Kdyz nekdy bude potreba skript, ktery neco vola nebo dlouho pocita, nedostane vic pravomoci. Stane se **vlastni sluzbou v AppFactory** a v katalogu konektorem, ktery ji vola. Tim pro nej zacne platit retry, rate limit i audit jako pro kazdy jiny krok. Skripty se nikdy nespousti v procesu API. Portal nesmi zpomalit kvuli tomu, ze nekdo ulozil spatny cyklus. ## Retence a objem, hned pri navrhu schematu Tohle je jedina vec, ktera pri scenari B rozhoduje o tom, jestli to za pul roku jde provozovat. **Zkracovani obsahu.** Plny vstup a vystup kroku se uklada jen u kroku, ktere selhaly, plus u male vzorku uspesnych. U ostatnich se uklada velikost, hash a prvnich radove 512 bajtu. Snizi to objem radove desetkrat a neztrati to nic, co by nekdo cetl - do uspesneho kroku se nikdo nechodi divat. **Partitionovani po mesicich** u pripisovacich tabulek: `run_step`, `ticket_trace`, `event_inbox`, `audit`. Mazani stareho oddilu je pak `DROP TABLE`, ne `DELETE` bezici pres noc. **Retence** podle toho, kdo to cte: | Data | Jak dlouho | | --------------------- | ----------------- | | Vstupy a vystupy kroku| 30 dni | | Souhrn behu | 12 mesicu | | Log ticketu | 90 dni | | Audit | dele, dane pravni potrebou | Cisla patri do nastaveni za klienta, protoze delsi retence je dobry duvod pro drazsi tarif. ## Co se monitoruje Ne CPU. Ctyri veci, a kazda odpovida na jinou otazku: | Metrika | Odpovida na | | ----------------------------------- | --------------------------------- | | Hloubka fronty | stiha se to | | **Vek nejstarsi pripravene ulohy** | je to zahlcene, nebo zaseknute | | Kroku za sekundu, p95 za konektor | kde to drhne | | Padle behy za hodinu, podil opakovani| co je rozbite | | Podil kroku za klienta | kdo je hlucny soused | | Zpozdeni inboxu (prijato az rozeslano) | stiha dispatcher | Bez veku nejstarsi ulohy se neda odlisit "je hodne prace" od "nic se nedeje", a to jsou dva uplne jine problemy se stejnou hloubkou fronty. ## Kdy zmenit architekturu Aby se to nemuselo rozhodovat dopredu. Do te doby plati navrh vyse. | Signal | Co udelat | | ---------------------------------------- | ------------------------------------- | | Fronta zere nad 30 % CPU databaze | vetsi davky, pak fronta v Redisu (BullMQ) | | Zapisy `run_step` prevalcuji IO | zkratit obsah, vzorkovat, velka tela do objektoveho uloziste | | Jeden klient dela nad 30 % provozu | vlastni bazen workeru, pak vlastni databaze | | Fronta roste kazdy den ve spicce | pridat workery, jsou bezstavove | | Prevazuji chyby 429 z cizich API | limity za napojeni, pak vyjednat kvoty| | Cekajici behy jdou do stovek tisic | oddelena fronta pro dlouha cekani, aby nezdrzovala bezny odber | ## Jeden container, dve role AppFactory nasazuje jednu aplikaci, takze worker nebude zvlastni sluzba. Rozdelit ho **procesne uvnitr image** pres `APP_ROLE=api|worker|both` s vychozim `both`. Az bude spicka takova, ze behy zpomaluji portal, nasadi se druha instance s `APP_ROLE=worker`. Zmena je jedna environment variable, zadny zasah do infrastruktury AppFactory. Migrace pri startu potrebuji poradovy zamek (`pg_advisory_lock`), aby je pri soubeznem nasazeni nespustilo vic instanci najednou. ## Poradi, v jakem to stavet 1. `event_inbox`, webhook, idempotence pri prijmu, `depth` skrz udalost. Bez toho zbytek nema co zpracovavat a smycky jsou otevrene. 2. `run`, `run_step`, `job`, `executeStep`, smycka po jedne uloze. Nejmensi verze, ktera vykona strom. 3. Retry, lease, idempotency key ke konektorum. 4. Dispatcher s filtrem na spousteci. 5. Krok `wait` a prehled cekajicich behu. 6. Davkovy odber, semafor za klienta, limity za napojeni. 7. Zkracovani obsahu, partitionovani, retence. 8. Bazen vlaken pro skripty. Body 1 az 3 jsou nutne, aby vubec neco bezelo. Body 6 a 7 jsou to, co odlisuje scenar A od scenare B, a **daji se dodelat pozdeji bez prestavby** - vyzaduji ale, aby uz od zacatku existovaly sloupce `tenant_id` a `queue_key` v tabulce `job` a partitionovani u `run_step`. Pridat oddily do nejvetsi tabulky v systemu az potom je ta jedina cast, ktera by opravdu bolela.