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).
This server cannot be deployed
Maintenance
ActivityMaintained
ResponsivenessNo issues