GNL
Dokümanlar menüsü
Core · Ücretsiz@gnldev/queue

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.

kuyruğa ekle ve koştur
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.

bir cron tetikleyicisi
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.

yay ve tüket
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.

karantina ve serbest bırakma
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#

fnenqueue

WorkStore'a bir iş yazar. QueueDepthExceededError sınırsız birikmeye karşı korur.

fncreateWorker

İşleri kapar ve koşturur. Seçenekler: owner, ttlMs (bayat kilidi geri alma), pollMs, maxAttempts (dead-letter eşiği), backoff, maxPollMs.

fnlistJobs

İş durumu listesi — dead-letter olanlar dahil.

fnretryJob

Dead-letter olmuş bir işi diriltir.

fnscheduleWorkflow

Bir tetikleyici kaydeder: { name, input?, at | every | cron, maxAttempts?, misfire? }.

fnpollScheduler

Vakti 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.

fnlistTriggers

Kayıtlı tetikleyiciler ve durumları — son hata dahil.

fncreateScheduler

pollScheduler etrafında başlat/durdur döngüsü; aynı lockTtlMs'i alır.

fnemit

Log'a bir olay ekler. EventDepthExceededError birikmeyi sınırlar.

fncreateConsumer

Kendi 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).

fnlistDeadEvents

Bir (topic, consumer) çifti için karantinadaki teslimatlar; deneme sayısı ve son işleyici hatasıyla birlikte.

fnretryDeadEvent

Karantinadaki 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.

İşleyiciler idempotent olmalı — ve ucuza olabilir
Kuyruk geri alması da olay tekrar teslimi de işleyicinizi yeniden koşturur. İşleyici dayanıklıysa bu yapısal olarak güvenlidir: iş ya da olay kimliğinden türetilmiş sabit bir 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.
Birikmeler bilerek sınırlıdır
Hem 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.