Actor — saf event-driven worker.
Yalnızca istemcilerden event geldiğinde çalışır; kendi başına hiçbir şey yapmaz.
Bildirim sistemi, komut işleyici, durum yöneticisi için idealdir.
use async_trait::async_trait;
use common::types::SessionId;
use serde_json::{Value, json};
use worker_runtime::{actor::Actor, context::WorkerContext};
pub struct NotificationWorker {
sent_count: u64,
}
impl NotificationWorker {
pub fn new() -> Box<dyn Actor> { Box::new(Self { sent_count: 0 }) }
}
#[async_trait]
impl Actor for NotificationWorker {
fn name(&self) -> &'static str { "notification" }
// Istemci abone oldugunda bir kez calisir
async fn on_session_connected(
&mut self, ctx: &WorkerContext, sid: SessionId,
) -> anyhow::Result<()> {
ctx.emit_to(sid, "welcome", json!({ "msg": "Baglandiniz" })).await?;
Ok(())
}
// Istemciden event geldiginde calisir
async fn on_event(
&mut self, ctx: &WorkerContext,
sid: SessionId, event: &str, payload: Value,
correlation_id: Option<String>,
) -> anyhow::Result<()> {
match event {
// Tek kullaniciya bildirim
"notify" => {
self.sent_count += 1;
ctx.emit_to(sid, "notification", json!({
"title": payload["title"],
"body": payload["body"],
"id": self.sent_count,
"corr_id": correlation_id,
})).await?;
}
// Bu worker'a abone TUM istemcilere yayin
"broadcast" => {
ctx.broadcast("announcement", json!({
"message": payload["message"],
"severity": payload.get("severity").unwrap_or(&json!("info")),
})).await?;
}
// Istatistik sorgula
"stats" => {
ctx.emit_to(sid, "stats", json!({ "sent": self.sent_count })).await?;
}
_ => {}
}
Ok(())
}
async fn on_session_disconnected(&mut self, _ctx: &WorkerContext, _sid: SessionId)
-> anyhow::Result<()> { Ok(()) }
async fn on_tick(&mut self, _ctx: &WorkerContext) -> anyhow::Result<()> { Ok(()) }
async fn on_stop(&mut self, _ctx: &WorkerContext) -> anyhow::Result<()> { Ok(()) }
}
Background — otonom döngü worker'ı.
WebSocket istemcisi olmadan da çalışır; sadece periyodik tick'lerle iş yapar.
Zamanlayıcı, veri senkronizasyonu, sistem izleme için idealdir.
on_event çağrılmaz.
use async_trait::async_trait;
use common::types::SessionId;
use serde_json::{Value, json};
use worker_runtime::{actor::Actor, context::WorkerContext};
pub struct HealthCheckWorker {
tick: u64,
last_ok: bool,
}
impl HealthCheckWorker {
pub fn new() -> Box<dyn Actor> {
Box::new(Self { tick: 0, last_ok: true })
}
}
#[async_trait]
impl Actor for HealthCheckWorker {
fn name(&self) -> &'static str { "health_check" }
// Worker baslarken bir kez calisir
async fn on_start(&mut self, _ctx: &WorkerContext) -> anyhow::Result<()> {
tracing::info!("health_check worker basladi");
Ok(())
}
// ~100ms'de bir supervisor tarafindan cagirilir (WorkerKind::Background)
async fn on_tick(&mut self, ctx: &WorkerContext) -> anyhow::Result<()> {
self.tick += 1;
// Her 100 tick'te bir (yaklasik 10 sn) kontrol et
if self.tick % 100 == 0 {
let db_host = ctx.config("DB_HOST").unwrap_or("localhost");
let ok = check_db_reachable(db_host).await;
if ok != self.last_ok {
self.last_ok = ok;
// Durum degisince abone istemcilere bildir
ctx.broadcast("health_changed", json!({
"service": db_host,
"healthy": ok,
"tick": self.tick,
})).await?;
}
}
Ok(())
}
async fn on_stop(&mut self, _ctx: &WorkerContext) -> anyhow::Result<()> {
tracing::info!("health_check worker durdu");
Ok(())
}
// Background worker'da bu metodlar cagrilmaz; bos birakilabilir:
async fn on_event(&mut self, _ctx: &WorkerContext, _sid: SessionId,
_e: &str, _p: Value, _c: Option<String>) -> anyhow::Result<()> { Ok(()) }
async fn on_session_connected(&mut self, _ctx: &WorkerContext, _sid: SessionId)
-> anyhow::Result<()> { Ok(()) }
async fn on_session_disconnected(&mut self, _ctx: &WorkerContext, _sid: SessionId)
-> anyhow::Result<()> { Ok(()) }
}
async fn check_db_reachable(_host: &str) -> bool {
// Gercek uygulamada: TCP baglantisi veya HTTP ping
true
}
Hybrid — hem event-driven hem otonom.
Istemciden komut alır VE periyodik tick ile kendi işini de yapar.
Mail kuyruğu, cache yöneticisi, toplu iş işleyici için idealdir.
use async_trait::async_trait;
use common::types::SessionId;
use serde_json::{Value, json};
use worker_runtime::{actor::Actor, context::WorkerContext};
pub struct MailQueueWorker {
queue: Vec<MailItem>,
flush_tick: u64,
sent_total: u64,
}
struct MailItem { to: String, subject: String, body: String }
impl MailQueueWorker {
pub fn new() -> Box<dyn Actor> {
Box::new(Self { queue: Vec::new(), flush_tick: 0, sent_total: 0 })
}
}
#[async_trait]
impl Actor for MailQueueWorker {
fn name(&self) -> &'static str { "mail_queue" }
// Istemciden komut: kuyruga ekle veya sorgu yap
async fn on_event(
&mut self, ctx: &WorkerContext,
sid: SessionId, event: &str, payload: Value,
_corr: Option<String>,
) -> anyhow::Result<()> {
match event {
"enqueue" => {
self.queue.push(MailItem {
to: payload["to"].as_str().unwrap_or("").into(),
subject: payload["subject"].as_str().unwrap_or("").into(),
body: payload["body"].as_str().unwrap_or("").into(),
});
ctx.emit_to(sid, "enqueued", json!({
"queue_depth": self.queue.len()
})).await?;
}
"flush" => {
let sent = self.do_flush().await;
ctx.emit_to(sid, "flushed", json!({ "sent": sent })).await?;
}
"status" => {
ctx.emit_to(sid, "status", json!({
"queued": self.queue.len(),
"sent_total": self.sent_total,
})).await?;
}
_ => {}
}
Ok(())
}
// Her ~3 saniyede bir otomatik flush (30 tick * 100ms = 3000ms)
async fn on_tick(&mut self, _ctx: &WorkerContext) -> anyhow::Result<()> {
self.flush_tick += 1;
if self.flush_tick % 30 == 0 && !self.queue.is_empty() {
let sent = self.do_flush().await;
tracing::info!(sent, queued = self.queue.len(), "otomatik flush");
}
Ok(())
}
// Kapanmadan once kuyrugu bosalt
async fn on_stop(&mut self, _ctx: &WorkerContext) -> anyhow::Result<()> {
self.do_flush().await;
Ok(())
}
async fn on_session_connected(&mut self, _ctx: &WorkerContext, _sid: SessionId)
-> anyhow::Result<()> { Ok(()) }
async fn on_session_disconnected(&mut self, _ctx: &WorkerContext, _sid: SessionId)
-> anyhow::Result<()> { Ok(()) }
}
impl MailQueueWorker {
async fn do_flush(&mut self) -> usize {
let items: Vec<_> = self.queue.drain(..).collect();
let n = items.len();
for item in items {
// Gercek uygulamada: SMTP veya mail API cagrisi
tracing::info!(to = %item.to, subject = %item.subject, "mail gonderildi");
self.sent_total += 1;
}
n
}
}
Helper — proje altinda, worker'larin uzerinde calisan
ortak/uzun omurlu kaynak. Proje yasadigi surece yasar (DB baglantisi, cache, API client).
Bir helper bir kez kurulur, o projedeki izinli TUM worker'lar paylasir.
1. Helper kodu (panelde yazilir, dylib derlenir)
use std::collections::HashMap;
use std::sync::{Mutex, OnceLock};
use async_trait::async_trait;
use serde_json::{json, Value};
use worker_runtime::helper::Helper;
pub struct DbHelper {
// Proje config'i — init'te yakalanir, her yerde self.cfg("KEY") ile okunur
// (worker'daki ctx.config() ile ayni).
config: OnceLock<HashMap<String, String>>,
conn: Mutex<Option<String>>, // gercekte: rusqlite::Connection / pool
}
impl DbHelper {
pub fn new() -> Box<dyn Helper> {
Box::new(Self { config: OnceLock::new(), conn: Mutex::new(None) })
}
// Proje sabitlerini her yerden oku — ctx.config() esdegeri.
fn cfg(&self, key: &str) -> Option<String> {
self.config.get().and_then(|c| c.get(key).cloned())
}
}
#[async_trait]
impl Helper for DbHelper {
fn name(&self) -> &'static str { "db" }
// Bir kez, yuklenirken cagrilir. config = proje config degiskenleri.
async fn init(&self, config: &HashMap<String, String>) -> anyhow::Result<()> {
let _ = self.config.set(config.clone()); // sakla → self.cfg(...) ile eris
let url = self.cfg("DB_URL").unwrap_or_default();
*self.conn.lock().unwrap() = Some(url); // baglantiyi burada kur, sakla
Ok(())
}
// Worker'lardan gelen cagrilar. &self → durum Mutex ardinda.
async fn call(&self, method: &str, args: Value) -> anyhow::Result<Value> {
match method {
"query" => {
let sql = args["sql"].as_str().unwrap_or_default();
let url = self.cfg("DB_URL").unwrap_or_default(); // config burada da erisilir
// self.conn ... uzerinden sorgu calistir
Ok(json!({ "rows": [], "sql": sql, "db": url }))
}
other => Err(anyhow::anyhow!("bilinmeyen metod: {other}")),
}
}
}
#[no_mangle]
pub unsafe extern "C" fn _create_helper() -> *mut Box<dyn Helper> {
Box::into_raw(Box::new(DbHelper::new()))
}
2. Worker kodunda kullanim
// Worker olustururken "db" helper'ina izin verilmis olmali.
async fn on_event(&mut self, ctx: &WorkerContext, sid: SessionId,
event: &str, payload: Value, corr: Option<String>) -> anyhow::Result<()> {
if let Some(db) = ctx.helper("db") {
let result = db.call("query",
json!({ "sql": "SELECT * FROM users" })).await?;
ctx.emit_to(sid, "users", result, corr).await;
}
Ok(())
}
Adimlar
- Proje sayfasinda Helpers → Ekle ile helper olustur (orn.
db).
- Helper'a tikla, Starter Template yukle, kodu duzenle.
- Kaydet → Build & Deploy. Yesil nokta = yuklendi.
- Worker olustururken izin verilen helper listesinden sec.
- Worker kodunda
ctx.helper("db") ile eris.
Onemli: call takes &self —
ayni anda bircok worker cagirir. Degisken state mutlaka Mutex/RwLock ardinda olmali.
Config: Proje sabitleri init aninda yakalanir (snapshot). Worker'lar gibi —
proje config'i degisirse helper'i tekrar Build edip yeni degerleri al.
Servis yeniden baslatilinca helper'lar tekrar Build edilmeli.
WorkerContext — worker'in dis dunyayla iletisim arayuzu.
Her event callback'ine ctx olarak gecilir.
// Proje Config okuma (panel "Proje Config" bolumunden eklenen KEY=VALUE)
// Degerler worker spawn edildiginde anlik goruntu alinir.
let db_host = ctx.config("DB_HOST").unwrap_or("localhost");
let api_key = ctx.config("API_KEY").unwrap_or_default();
// Tekil emit — sadece belirtilen session'a gonder
ctx.emit_to(session_id, "event_adi", serde_json::json!({
"key": "value",
"sayi": 42,
})).await?;
// Yayin — bu worker'a abone TUM aktif istemcilere gonder
ctx.broadcast("event_adi", serde_json::json!({
"timestamp": chrono::Utc::now().timestamp(),
"data": "...",
})).await?;
// Worker kimligi
let project = &ctx.id.project; // "crm"
let worker = &ctx.id.worker; // "notification"
let version = &ctx.id.version; // "v1"
// Loglama (JSON formatiyla loglanir)
tracing::info!(worker = %ctx.id, "Islem tamamlandi");
tracing::warn!(err = %e, "Hata olustu");
tracing::error!("Kritik hata!");
• emit_to ve broadcast async — her zaman .await? ile cagirin.
• Config panelden guncellendikten sonra worker yeniden baslatilmalidir.
• broadcast o anda bagli ve abone olan oturumlara gonderir.
WebSocket Protokolu — Gateway port 9000'de JSON frame'leri alir/gonderir.
Her mesaj type alaniyla tanimlanir.
Istemci → Sunucu
// 1. Baglanti
const ws = new WebSocket("ws://localhost:9000/?token=dev-secret");
// 2. Worker'a abone ol
ws.send(JSON.stringify({
type: "subscribe", project: "crm", worker: "notification", version: "v1"
}));
// 3. Event gonder
ws.send(JSON.stringify({
type: "event",
project: "crm",
worker: "notification",
event: "notify",
payload: { title: "Merhaba", body: "Test mesaji", level: "info" },
correlation_id: "req-abc123" // istege bagli; yanit frame'inde ayni deger doner
}));
// 4. Aboneligi iptal et
ws.send(JSON.stringify({ type: "unsubscribe", project: "crm", worker: "notification" }));
// 5. Canlilik kontrolu
ws.send(JSON.stringify({ type: "ping" })); // {"type":"pong"} doner
Sunucu → Istemci
ws.onmessage = (e) => {
const frame = JSON.parse(e.data);
switch (frame.type) {
case "subscribed": // Abonelik onaylandi
console.log("Abone:", frame.project, frame.worker, frame.version);
break;
case "event": // Worker'dan gelen event
console.log(frame.event, frame.payload);
// frame.correlation_id ile istek-yanit eslestirilebilir
break;
case "worker_state_changed": // Worker durumu degisti
console.log(frame.worker, "->", frame.status);
// status: "running" | "paused" | "stopped" | "crashed(...)"
break;
case "error": // Protokol veya is mantigi hatasi
console.error(frame.code, frame.message);
break;
case "pong": break;
}
};
Yeni Worker Olusturma
- Proje sayfasinda Ekle butonuna tikla — worker adi, versiyon, tur sec.
- Listede worker belirir. Uzerine tikla — editorde acilir.
- Starter Template Olustur butonuyla temel kod yukle.
- Kodu duzenle. Tab = 4 bosluk, Cmd/Ctrl+S = kaydet.
- Kaydet — yeni versiyon adi gir — onayla.
- Basla butonuyla worker'i calistir.
Versiyonlama
Kaydet her zaman yeni versiyon olusturur — mevcut dosya dokunulmaz.
Yeni versiyon, ayni adi tasiyan worker'in factory'sini devralir ve hemen baslatilabilir.
Eski versiyon calismaya devam eder; hazir oldugunda Durdur ile kapat.
Dosyalar: workers/{proje}/{ad}/{versiyon}/src/lib.rs
Lifecycle Butonlari
Basla — durdurulmus worker'i yeniden baslatir. Zaten calisiyorsa no-op.
Duraklat — event islemeyi dondurur; kuyruga yeni eventler dusmez.
Devam — duraklatilmis worker'i devam ettirir.
Durdur — Shutdown sinyali gonderir; on_stop callback'i calistirilir.
Gateway: ws://localhost:9000 | Panel: http://localhost:9001
DEV_TOKEN: dev-secret | ADMIN_TOKEN: admin-secret
Actor = sadece event | Background = sadece tick | Hybrid = ikisi birden