Dokümanlar menüsü
Arka plan işleri: kuyruk, zamanlayıcı ve olaylar
Ajan işini istek dışında koşturun — dayanıklı işler, cron tetikleyicileri ve bir olay yolu; her biri kendi sınırında exactly-once, hepsi aynı journal'ın üstünde.
Ne işe yarar#
Her ajan çalıştırması bir HTTP isteğiyle başlamaz. Birinin biriken işi işlemesi, 03:00'te tetiklenmesi ya da bir olaya tepki vermesi gerekir — ve bunların her biri, bir çökmenin yan etkiyi ikileyebileceği bir yerdir. Bu üç paket, aynı garantiyi bu üç sınıra koyar.
Hiçbiri dayanıklılığı yeniden yazmaz. Bir iş işleyicisi tipik olarak runDurable çağırır, bir tetikleyici dayanıklı bir workflow koşar, bir olay tüketicisinden de aynısı beklenir. Her paketin eklediği şey, journal'ın tek başına veremeyeceği parçadır: iki worker'ın aynı işi almaması için bir kilit, bir cron aralığının iki kez koşmaması için bir tetikleme kaydı, ve bir olayın bitmeden bitti işaretlenmemesi için bir ack işareti.
Kuyruk — worker'dan sağ çıkan işler#
enqueue işi yazar, createWorker onu kapıp koşturur. Bir worker TTL'li bir kilit tutar; çökerse kilit bayatlar, başka bir worker işi geri alır, işleyici tekrar koşar — ve işleyici aynı runId'yi devam ettirdiği için yan etki yine bir kez olur. Başarısızlık maxAttempts ile sınırlıdır; sonrasında iş dead-letter olur ve retryJob ile diriltilebilir.
Boş-yoklama geri çekilmesi varsayılan olarak açıktır: bir yoklama hiçbir şey kapamazsa aralık bir tavana kadar ikiye katlanır ve iş kapılır kapılmaz sıfırlanır. Bu varsayılanın sebebi bir denetim bulgusudur — boş bir kuyruğu sabit aralıkla yoklayan bin tüketici saniyede on binlerce sorgu üretir.
import { enqueue, createWorker } from '@gnldev/queue';
await enqueue(storage.work, 'reindex', { docId: 'd-9' });
// one handler per job type; ctx carries the job's own durable runId and journal
const worker = createWorker(
storage,
{
reindex: async (payload, ctx) => {
// a reclaim resumes this runId, it does not restart it
await runDurable({ runId: ctx.runId, journal: ctx.journal, model, tools, prompt: payload.docId });
},
},
{ ttlMs: 30_000, maxAttempts: 3 },
);
worker.start();Zamanlayıcı — çift tetiklemesiz cron#
scheduleWorkflow bir tetikleyici kaydeder: tek seferlik at, periyodik every ya da 5 alanlı cron (UTC, dakika çözünürlüğü). Tetikleyici tanımı journal'da değişmezdir, yalnız durumu hareket eder. Her tetikleme bir çalıştırma kilidi alır ve runId'sini tetikleme sayacından türetir; böylece aynı aralık iki kez başlayamaz.
Zaman veri olarak ele alınır: sonraki çalışma zamanı çözülür ve sonra journal'a dondurulur — bir replay'in duvar saatiyle kaymasını engelleyen şey budur.
misfire, kaçırılmış bir pencerenin ne anlama geldiğine karar verir. 'skip' (varsayılan) planlanan ızgaradaki bir sonraki aralığa atlar, böylece kayma birikmez; 'catchup' kaçırılan her tetiklemeyi sırayla, yoklama başına bir tane koşar. Gecelik bir rapor genelde skip ister; bir faturalama döngüsü genelde catchup.
import { scheduleWorkflow, createScheduler } from '@gnldev/scheduler';
await scheduleWorkflow(journal, {
name: 'nightly-report',
cron: '0 3 * * *', // 03:00 UTC
misfire: 'skip', // or 'catchup' to replay every missed slot
});
createScheduler(journal, runner).start();Olaylar — dürüst garantili fan-out#
emit bir log'a ekler; her tüketici olayları kendi ack işareti altında işler, yani N tüketicinin her biri her olayı görür. İşaret, işleyici başarılı olduktan sonra bir CAS ile yazılır. Bütün tasarım bu sıralamadadır: işleyici hata fırlatırsa ya da süreç ölürse olay kaybolmaz, sonraki yoklama tekrar dener.
Bedeli gizlenmez, yazılır: bu exactly-once işaretleme, at-least-once teslimtir. Başarılı bir işleyici ile ack'i arasındaki bir çökme — ya da eşzamanlı bir yoklama yarışı — tekrar teslime yol açabilir. İşleyiciyi idempotent yazın ya da içine runDurable / claim koyun; o zaman tekrar teslimin bir maliyeti kalmaz.
import { emit, createConsumer } from '@gnldev/events';
await emit(storage.work, 'order.paid', { orderId: 'o-7' });
const consumer = createConsumer(
storage.work,
'order.paid',
async (payload, meta) => {
// the ack marker is written only after this returns — so make it idempotent
await runDurable({ runId: `evt:${meta.id}`, journal, model, tools, prompt: payload.orderId });
},
{ name: 'fulfilment' }, // ack markers are separated by this name (fan-out)
);
consumer.start();Tek bir hatalı işleyici artık konuyu kilitlemiyor#
Bir tüketicinin imleci, ancak işleyicisinin kabul ettiği olayın ötesine geçer. Sürekli hata fırlatan bir işleyici bu yüzden imleci sonsuza kadar park ediyordu: hatalı olayın arkasındaki her olay, asla başarılı olmayacak bir teslimatı bekliyordu. Kaybolmuş mesaj değil — durmuş bir konu.
Artık bir teslimat maxAttempts başarısızlıktan sonra (varsayılan 8) karantinaya alınıyor ve imleç ilerliyor. Denemeler retryDelayMs ile aralanıyor — 60 saniyeden başlayıp ikiye katlanan, bir saatte sınırlanan bir üstel dizi — böylece bütçe yoklama aralığında bir saniyede tükenmek yerine yaklaşık iki saate yayılıyor. Bir yeniden deneme bütçesi, ancak var olduğu kesintiden uzun yaşarsa bütçedir.
Vakti gelmemiş bir olay atlanır ama imleci yine dondurur: hâlâ yeniden denenebilir olduğu için arkasındakiler geçilmiş sayılamaz. Beklemek, vazgeçmek değildir. Hiç karantinaya almak istemiyorsanız maxAttempts: Infinity bu şekli kalıcı yapar — diğer olaylar yine akar, ama bu tüketici park kalır ve her yoklamada oradan itibaren yeniden tarar.
import { createConsumer, listDeadEvents, retryDeadEvent } from '@gnldev/events';
const consumer = createConsumer(storage.work, 'order.paid', handler, {
name: 'fulfilment',
maxAttempts: 8, // quarantine after this many FAILED attempts (the default)
retryDelayMs: 60_000, // fixed spacing; omit for exponential 60s → 1h (the default)
});
// what is parked, and why
const dead = await listDeadEvents(storage.work, 'order.paid', 'fulfilment');
// [{ id: 'e-3', attempts: 8, error: 'ECONNREFUSED', at: 173... }]
// hand one back once the downstream is healthy again
await retryDeadEvent(storage.work, 'order.paid', 'fulfilment', 'e-3');listDeadEvents(work, topic, consumer) parkta ne olduğunu gösterir, retryDeadEvent(work, topic, consumer, eventId) birini geri verir. Studio ikisini de sunar, operatörün betik yazmasına gerek yok. Adresleme üçlü ile yapılır: (topic, consumer, id) — bir konu dallanır, yani tek bir olayın tüketici başına bir karantina kaydı olur ve tek başına id bunların birkaçını birden gösterir.
API#
enqueueWorkStore'a bir iş yazar. QueueDepthExceededError sınırsız birikmeye karşı korur.
createWorkerİşleri kapar ve koşturur. Seçenekler: owner, ttlMs (bayat kilidi geri alma), pollMs, maxAttempts (dead-letter eşiği), backoff, maxPollMs.
listJobsİş durumu listesi — dead-letter olanlar dahil.
retryJobDead-letter olmuş bir işi diriltir.
scheduleWorkflowBir tetikleyici kaydeder: { name, input?, at | every | cron, maxAttempts?, misfire? }.
pollSchedulerVakti gelen tetikleyicileri ateşler; her ateşleme bir koşu kilidi ve ateşleme-sayacı runId ile tam-bir-kezdir. lockTtlMs (varsayılan 60sn) çöken bir yoklayıcının tetikleyiciyi ne kadar bloke edeceğini sınırlar — kilidi tutan çalışırken yeniliyor, yani kurtarma süresi en yavaş koşu değil TTL kadardır.
listTriggersKayıtlı tetikleyiciler ve durumları — son hata dahil.
createSchedulerpollScheduler etrafında başlat/durdur döngüsü; aynı lockTtlMs'i alır.
emitLog'a bir olay ekler. EventDepthExceededError birikmeyi sınırlar.
createConsumerKendi ack işaretleri olan adlandırılmış tüketici; işaret ancak işleyici başarılı olduktan sonra yazılır. Seçenekler: pollMs, backoff, maxPollMs, maxAttempts (karantina eşiği, varsayılan 8), retryDelayMs (denemeler arası aralık, varsayılan 60 saniyeden üstel).
listDeadEventsBir (topic, consumer) çifti için karantinadaki teslimatlar; deneme sayısı ve son işleyici hatasıyla birlikte.
retryDeadEventKarantinadaki bir teslimatı tüketiciye geri verir. Etkisiz-tekrarlanabilir: zaten serbest bırakılmış bir olayı yeniden bırakmak belgelenmiş bir işlemdir, hata değil.
runId ile runDurable çağırın; içindeki tamamlanmış tool'lar bir daha çalışmaz. Bu, HTTP yolunun aldığı garantinin aynısıdır, başka bir kapıdan girilmiş hâli.enqueue hem emit, birikme çok derinleştiğinde reddedebilir (QueueDepthExceededError, EventDepthExceededError). Sınırsız bir kuyruk, yavaş bir tüketiciyi depolama dolduğunda fark edilen bir kesintiye çevirir; reddedilen bir yazım ise hemen fark edilir.