Skip to main content
Glama
halilogia

swarm-orchestrator

by halilogia
README.md
# ⚡ SwarmOrchestrator

**Birden fazla AI modelini eşzamanlı çalıştıran, hata toleranslı çoklu ajan orkestrasyon motoru.**

SwarmOrchestrator, yerel OpenAI uyumlu router'ınız (`9Router` / `LiteLLM` / `OpenRouter`, varsayılan `http://localhost:20128/v1`) veya kalıcı **Antigravity CLI (`agy`)** alt süreci üzerinden bir görev listesini sınırlı eşzamanlılıkla koşturur. Sonuçları tipli modellerle döner, isteğe bağlı olarak diske yazar ve her batch'i JSON olarak loglar. Antigravity ve diğer MCP istemcileri için **6 MCP aracı** sunar.

> **Durum:** Erken sürüm. Motor, sağlayıcı fabrikası, retry/backoff, iptal, maliyet/gecikme logu, sanitizasyon, CLI ve MCP araçları kod olarak mevcuttur; offline testler bu mantığı doğrular. **Canlı 9Router koşusu bu repoda doğrulanmamıştır** — router erişilebilirliği ortama bağlıdır. Gerçek ağ davranışı yalnızca router açıkken `python tests/test_agy_provider.py` ile görülebilir.

---

## 🏛️ Mimari Şema

```text
                                  ┌────────────────────────┐
                                  │  Antigravity / User    │
                                  └───────────┬────────────┘
                                              │
                                     (Tasks / Batches)
                                              ▼
                                 ┌─────────────────────────┐
                                 │    SwarmEngine (Core)   │
                                 │  - Asyncio Semaphore    │
                                 │  - Retry + Backoff      │
                                 │  - Run Logger           │
                                 └────────────┬────────────┘
                                              │
                       ┌──────────────────────┼──────────────────────┐
                       ▼                      ▼                      ▼
               ┌───────────────┐      ┌───────────────┐      ┌───────────────┐
               │   Worker #1   │      │   Worker #2   │      │   Worker #N   │
               └───────┬───────┘      └───────┬───────┘      └───────┬───────┘
                       │                      │                      │
                       └──────────────────────┼──────────────────────┘
                                              ▼
                                 ┌─────────────────────────┐
                                 │  Provider (fabrika)     │
                                 │  AsyncRouterClient      │
                                 │  veya AGYSession        │
                                 └─────────────────────────┘
```

Eşzamanlılık `asyncio.Semaphore` ile sınırlanır (`MAX_CONCURRENT_WORKERS`, varsayılan `8`). Her worker bağımsız çalışır; bir görevin hatası diğerlerini durdurmaz.

---

## 🚀 Hızlı Başlangıç

### 1. Ortam Değişkenleri (`.env`)

Proje kökündeki `.env.example` dosyasını `.env` olarak kopyalayıp doldurun:

```powershell
Copy-Item .env.example .env
```

```env
LOCAL_ROUTER_BASE_URL=http://localhost:20128/v1
LOCAL_ROUTER_API_KEY=your_local_key_here

DEFAULT_MODEL=all
FALLBACK_MODEL=ag/gemini-3.7-flash-medium

MAX_CONCURRENT_WORKERS=8
REQUEST_TIMEOUT_SECONDS=150
MAX_RETRIES=3
RETRY_DELAY_SECONDS=2

TASK_TIMEOUT_SECONDS=0
BATCH_TIMEOUT_SECONDS=0

SAVE_RUN_LOGS=true
RUNS_DIR=./runs
METRICS_LOG=./runs/metrics.jsonl
```

`.env` **asla commit edilmez** (bkz. `.gitignore`).

### 2. Bağımlılıklar

```powershell
pip install -r requirements.txt
```

`asyncio` listedelenmez: stdlib modülüdür (3.4+) ve PyPI'daki `asyncio` paketi
Python 3.4'ten eski sürümler içindir — kurulursa stdlib'i gölgeler ve projedeki
tüm `import asyncio` ifadelerini kırar.

### 3. Kalite kapısı (offline)

```powershell
python -m ruff check .     # statik kontrol
python -m pytest           # 134 offline test, ağ veya router gerekmez
```

`tests/` altındaki testlerin tamamı sahte (fake) istemci kullanır ve
`conftest.py` herhangi bir bağlantı denemesini hata olarak işaretler. Canlı
altyapı gerektiren tek dosya `tests/test_agy_provider.py`'dır; bu dosya pytest
tarafından toplanmaz, doğrudan çalıştırılır (bkz. `AGENTS.md`).

---

## 💻 Kullanım Yöntemleri

### Yöntem A: MCP Sunucusu (önerilen)

`mcp/server.py` bir FastMCP sunucusudur. Antigravity veya MCP istemcinizin yapılandırmasına ekleyin:

```json
{
  "mcpServers": {
    "swarm-orchestrator": {
      "command": "python",
      "args": ["C:/Users/<you>/Desktop/GitHub/Public/SwarmOrchestrator/mcp/server.py"],
      "env": {
        "PYTHONPATH": "C:/Users/<you>/Desktop/GitHub/Public/SwarmOrchestrator"
      }
    }
  }
}
```

`args` ve `PYTHONPATH` bu deponun köküne işaret etmelidir. Sunucu, çalışma dizininden bağımsız olarak proje kökünü `sys.path`'e kendisi ekler.

### Yöntem B: Python API'si

```python
import asyncio
from core.swarm_engine import SwarmEngine
from core.task_models import SwarmTask

async def main():
    tasks = [
        SwarmTask(
            id="t1",
            title="Örnek",
            system_prompt="You are a precise assistant.",
            user_prompt="Reply with exactly: OK",
        )
    ]
    summary = await SwarmEngine(max_workers=4).run_batch(tasks=tasks)
    print(summary.completed_tasks, "/", summary.total_tasks)

asyncio.run(main())
```

### Yöntem C: Hazır iş akışları

```python
from pathlib import Path
from workflows.translation_workflow import TranslationWorkflow

workflow = TranslationWorkflow()          # use_cache=True varsayılan
summary = await workflow.execute(files=[Path("ornek.tsx")])
```

`workflows/` altında iki hazır akış vardır:

| İş akışı | Sınıf | Rol |
| --- | --- | --- |
| `translation_workflow.py` | `TranslationWorkflow` | Dosyaları teknik İngilizceye çevirme; hash tabanlı cache ile atlama |
| `generic_workflow.py` | `GenericSwarmWorkflow` | Serbest görev tanımlarından batch üretme |

`base_workflow.py` içindeki `BaseWorkflow`, `prepare_tasks()` uygulanmasını bekleyen soyut şablondur.

> **Yöntem D: CLI.** Terminalden çalıştırmak için `cli.py` kullanılır (aşağıya bakın).

---

## 🖥️ CLI

`cli.py` hiçbir yol sabitlemez; tüm dizinler argümandır.

| Komut | İşlev |
| --- | --- |
| `python cli.py translate -i <kaynak> -o <hedef>` | Bir dizindeki dosyaları çevirir, `--dry-run` yalnızca planı listeler |
| `python cli.py run -t tasks.json -o sonuclar/` | JSON görev listesini paralel koşar, görev başına dosya yazar |
| `python cli.py models` | Router'daki canlı model kimliklerini ve eşleşen fiyatı listeler |
| `python cli.py info` | Etkin yapılandırmayı yazar (API anahtarı gösterilmez) |
| `python cli.py serve` | MCP sunucusunu bu depodan başlatır |

```powershell
# Klasör çevirisi (önce --dry-run ile planı görün)
python cli.py translate -i .\src\pages -o .\locales\en\pages --dry-run
python cli.py translate -i .\src\pages -o .\locales\en\pages -c 8

# JSON görev listesi
python cli.py run -t .\tasks.json -o .\out -c 4 --provider agy

# Belirli görevleri hiç çalıştırmadan atla (per-task cancel)
python cli.py run -t .\tasks.json --cancel task_3 --cancel task_7
```

Çeviri komutu mevcut çıktıları otomatik atlar (`--force` ile hepsini yeniden
çevirir), içerik hash'ine dayalı cache'i kullanır ve başarılı dosyalardan sonra
cache'i günceller. `run_page_translations.py` artık bu komutun ince bir sarmalayıcısıdır:

```powershell
python run_page_translations.py --input <kaynak> --output <hedef>
```

### Maliyet ve gecikme raporu

Her batch sonrası terminalde özet, ayrıca `runs/metrics.jsonl` içine batch başına
tek satır yazılır (gecikme yüzdilleri, token sayıları, görev başına maliyet):

```text
batch translate_1 @ 2026-09-27T09:04:44
  tasks      : 2 total | 2 completed | 0 failed | 0 cancelled
  wall clock : 0.02s
  provider   : mean 0.01s | p50 0.01s | p95 0.01s | max 0.01s
  per task   : mean 0.02s | p50 0.02s | p95 0.02s
  tokens     : 300 (in 200 / out 100)
  cost       : $0.0000 (incomplete: no rate configured for all)
```

> **Maliyet dürüstlüğü:** Bu proje fiyat listesi taşımaz. `MODEL_PRICING_JSON`
> tanımlanmadan `cost` **$0.0000 (incomplete)** olarak raporlanır — ücretsiz
> olduğu anlamına gelmez. Fiyat girdisi:
> ```env
> MODEL_PRICING_JSON='{"paid/":{"input":3.0,"output":15.0}}'
> ```
> Token sayıları router `usage` alanı döndürürse dolar (bkz. `REQUEST_STREAM_USAGE`).

### İptal (cancel) davranışı

| Tetik | Etki |
| --- | --- |
| `Ctrl+C` (CLI) | Uçuşta olan tüm görevler iptal edilir; batch özeti yine yazdırılır |
| `--cancel <id>` | Belirtilen görev hiç sağlayıcıya ulaşmaz |
| `--task-timeout` | Tek görev bütçesi; aşan görev **başarısız** olur, diğerleri sürer |
| `--batch-timeout` | Batch bütçesi; bitmeyen görevler **iptal** edilir |
| `engine.drain()` | Yeni deneme başlatılmaz; uçuşta olan deneme tamamlanır |

İptal edilen görev `cancelled` durumuna geçer ve tek başına batch'i düşürmez.

---

## 🧰 MCP Araçları

`mcp/server.py` aşağıdaki araçları sunar:

| Araç | İşlev |
| --- | --- |
| `get_available_models` | Router'dan canlı model kimliklerini ve varsayılan komboyu döner |
| `get_live_models` | Modelleri yeteneklere göre inceler (tools / reasoning / vision / context) ve filtreler |
| `check_model_health_and_quota` | Router bağlantısını ve gecikmeyi ölçer; sağlayıcı dağılımını raporlar |
| `get_router_logs` | `runs/` altındaki son batch loglarının özetini döner |
| `run_smart_task` | Tek görevi görev tipine göre model seçerek koşar; başarısızlıkta fallback |
| `run_swarm_batch` | Görev listesini eşzamanlı koşturur; `target_file` verilirse sonucu diske yazar |

Örnek: MCP istemcisinden *"Şu 40 dosyayı 8 eşzamanlı worker ile çevir"* dediğinizde ajan `run_swarm_batch` aracını çağırır.

---

## 🔀 Sağlayıcılar

`core/client.py` içindeki `get_ai_client(provider=...)` bir fabrikadır:

| `provider` değeri | Dönen istemci | Davranış |
| --- | --- | --- |
| `default` / `9router` | `AsyncRouterClient` | HTTP; streaming yanıtları toleranslı biçimde ayrıştırır |
| `agy` / `antigravity` / `cli` | `AGYSession` | Kalıcı `agy` alt süreci; NDJSON (`stream-json`); **tekil oturum** (eşzamanlılık = 1) |

`AGYSession` bir singleton'dır: `asyncio.Lock` ile istekleri sıraya alır, izole bir geçici dizinde (`tempfile.mkdtemp`) çalışır ve `--dangerously-skip-permissions` **kullanmaz**.

> **Eşzamanlılık uyarısı:** `run_swarm_batch` içinde `provider="agy"` seçerseniz, `concurrency` parametresi ne olursa olsun gerçek paralellik **1**'dir. Yüksek paralellik için router sağlayıcısını kullanın.

---

## 📂 Proje Yapısı

```text
SwarmOrchestrator/
├── cli.py                       # Terminal arayüzü: translate / run / models / info / serve
├── config/
│   └── settings.py              # Pydantic tabanlı merkezi ayarlar
├── core/
│   ├── client.py                # Async HTTP istemcisi + sağlayıcı fabrikası
│   ├── agy_session.py           # Kalıcı Antigravity CLI (agy) oturumu
│   ├── worker.py                # Görev yaşam döngüsü, adaptif timeout, retry, iptal
│   ├── swarm_engine.py          # Semafor, iptal API'si, batch motoru + run logger
│   ├── metrics.py               # Token/gecikme istatistikleri, JSONL maliyet logu
│   ├── pricing.py               # Model prefix'i ile fiyat eşleştirme
│   ├── task_models.py           # Tipli veri modelleri (Pydantic)
│   ├── sanitizer.py             # Markdown/konuşma temizliği
│   └── cache.py                 # SHA-256 tabanlı görev cache'i
├── workflows/
│   ├── base_workflow.py         # Soyut iş akışı şablonu
│   ├── translation_workflow.py  # Teknik çeviri iş akışı
│   └── generic_workflow.py      # Genel görev iş akışı
├── mcp/
│   └── server.py                # FastMCP sunucusu (6 araç)
├── tests/
│   ├── conftest.py              # Offline koruması + sahte istemci
│   ├── test_router_parser.py    # SSE/JSON ayrıştırma (offline)
│   ├── test_worker_retry.py     # Retry/backoff/iptal/maliyet (offline)
│   ├── test_engine.py           # Batch sayımı, eşzamanlılık, iptal (offline)
│   ├── test_metrics.py          # Yüzdelik, fiyat, JSONL log (offline)
│   ├── test_sanitizer.py        # Fence/önsöz temizliği (offline)
│   ├── test_cli.py              # Argüman/yol/çıkış kodu (offline)
│   └── test_agy_provider.py     # CANLI sağlayıcı testi (pytest dışı, gerçek model çağırır)
├── run_page_translations.py     # Geriye uyumlu sarmalayıcı -> cli.py translate
├── scripts/
│   └── sync_to_gemini.ps1       # Kanonik depoyu .gemini çalışma kopyasına yansıtır
├── archives/                    # Eski/arşivlenmiş kod
├── docs/
│   └── KNOWLEDGE.md             # Kalıcı teknik bilgi
├── runs/                        # Batch JSON + metrics.jsonl (git-ignored)
├── pyproject.toml               # ruff + pytest yapılandırması
├── AGENTS.md                    # Ajan kılavuzu
├── ARCHITECTURE.md              # Katmanlar ve veri akışı
├── CHANGELOG.md                 # Sürüm tarihçesi
├── ROADMAP.md                   # Hedef fazlar
├── .env.example
├── requirements.txt
└── LICENSE
```

---

## 🛡️ Dayanıklılık Davranışları

- **Adaptif timeout:** `SwarmWorker._compute_adaptive_timeout()` istem uzunluğuna göre, tabanı `REQUEST_TIMEOUT_SECONDS` olan ve 240 sn ile sınırlanan bir değer hesaplar.
- **Retry + üstel backoff:** `MAX_RETRIES` kadar deneme; her denemede `RETRY_DELAY_SECONDS * 2^(n-1)` bekleme. Backoff sırasında iptal gelirse bekleme kesilir.
- **Boş yanıt başarısızlıktır:** Model boş içerik döndürürse görev `completed` sayılmaz, yeniden denenir; diske boş dosya yazılmaz.
- **İptal:** `SwarmEngine.cancel(id)`, `drain()` ve `cancel_all()`; ayrıca görev ve batch zaman bütçeleri. Tek bir görevin iptali/hatası batch'i düşürmez.
- **Toleranslı ayrıştırma:** `AsyncRouterClient._parse_router_response()` düz JSON, `raw_decode` ve satır satır SSE toplama olmak üzere üç aşamalı ayrıştırır; SSE `usage` alanlarını toplar.
- **Maliyet/ gecikme logu:** Her batch `runs/metrics.jsonl` içine tek satır yazar; fiyat tablosu yoksa maliyet "unpriced" olarak işaretlenir.
- **Sanitizasyon:** Her sonuç `CodeSanitizer` üzerinden geçer (konuşma önsözü ve markdown fence temizliği).

---

## ⚠️ Bilinen Sınırlamalar

- **Canlı router bağımlılığı:** Motorun router yolu (`AsyncRouterClient`) bu repoda uçtan uca doğrulanmamıştır. `tests/` altındaki testler sahte istemci kullanır ve ağ davranışını kanıtlamaz.
- **`agy` paralelliği 1'dir:** Kalıcı CLI oturumu tekil olduğundan ölçeklenmez.
- **Fiyat listesi yok:** `MODEL_PRICING_JSON` tanımlanmadıkça maliyet "unpriced" raporlanır; token sayıları da ancak router `usage` dönerse dolar (`REQUEST_STREAM_USAGE`).
- **Sanitizasyon TS'e özeldir:** `CodeSanitizer.clean_typescript_code()` adı ve desenleri TypeScript odaklıdır; TS dışı çıktı için özelleştirme gerekir.
- **`runs/` git-ignored:** Koşu logları ve `metrics.jsonl` yerelde kalır; repoya girmez.
- **UTF-8 yanıt varsayımı:** `AsyncRouterClient` yanıt gövdesini UTF-8 olarak çözer.
- **Ayrıştırıcı kusuru (bilinen):** `data:` öneki olmayan ardışık iki JSON nesnesi birleştirilmez; `raw_decode` aşaması ilk nesneyi olduğu gibi döndürür. Yalnızca SSE gövdeleri toplama aşamasına ulaşır.

---

## 🔁 Kaynak Depo ve Çalışma Kopyası

Bu depo **tek doğruluk kaynağıdır** (source of truth). Ancak Antigravity, MCP sunucusunu `.gemini` altındaki bir kopyadan yükler. İki kopya zamanla ayrışmaması için değişiklikleri senkronize edin:

```powershell
# Önizleme (hiçbir şey kopyalanmaz)
powershell -ExecutionPolicy Bypass -File .\scripts\sync_to_gemini.ps1 -WhatIf

# Gerçek senkronizasyon
powershell -ExecutionPolicy Bypass -File .\scripts\sync_to_gemini.ps1
```

Betik `robocopy /MIR` kullanır; bu depoda sildiğiniz bir dosya hedefte de silinir (istenen davranış — ayrışmayı önler).

**Kasıtlı olarak korunanlar** (hedefte asla dokunulmaz veya silinmez):

| Öğe | Neden |
| --- | --- |
| `.env` | Yerel router anahtarı; hedefe özgü ve git-ignored |
| `runs/` | Yerel koşu logları |
| `.venv/`, `venv/` | Yerel sanal ortam |

**Akış:** kodu bu depoda düzenle → commit et → `sync_to_gemini.ps1` çalıştır → Antigravity'yi yeniden başlat.

> **Dikkat:** Yalnızca bu depoyu düzenleyin. `.gemini` kopyasında yapılan doğrudan değişiklikler bir sonraki senkronizasyonda **kaybolur**.

---

## 📄 Lisans

GNU General Public License v3.0 — bkz. [LICENSE](LICENSE).