diff --git a/README.md b/README.md index 32a587d..453b700 100644 --- a/README.md +++ b/README.md @@ -161,10 +161,12 @@ Her veri parçasının sahibi tam olarak tek bir bileşendir: `state.json`'ı ya Agent her kaydı önce disk'teki spool'a yazar ve ancak `200` cevabını aldıktan sonra siler. Bu yüzden retry beklenen bir durumdur — dolayısıyla her kayıt agent'ın ürettiği bir `UUID` taşır ve server `ON CONFLICT DO NOTHING` ile insert eder. Sonuç: aynı kayıt iki kez gönderilse bile duplicate oluşmaz, veri kaybı içinse disk'in kendisinin arızalanması gerekir. **Pause, kaydı durdurmaz.** -Bir cihazı pause etmek *upload*'ı durdurur, *toplamayı* değil. Veri yerelde birikmeye devam eder ve resume'da sırayla akar. Pause sırasında komut yoklaması da devam eder — aksi hâlde `resume` komutu cihaza hiçbir zaman ulaşamazdı. +Bir cihazı pause etmek *upload*'ı durdurur, *toplamayı* değil. Veri yerelde birikmeye devam eder ve resume'da sırayla akar. Pause sırasında komut yoklaması da devam eder — aksi hâlde `resume` komutu cihaza hiçbir zaman ulaşamazdı. Komut ack'i de aynı sebeple durmaz: o bir telemetri değil kontrol mesajıdır ve tek bir ölçüm satırı taşımaz. Durdurulsaydı server komutun uygulandığını hiç öğrenemez, aynı `pause`u sonsuza kadar yeniden gönderir ve dashboard cihazı hâlâ "çalışıyor" gösterirdi. **Silme işleminin bir sırası vardır.** -Bir cihazı kaldırmak satırı hemen silmez. Önce kuyruğa bir `delete` komutu girer; agent kendini yerelde temizler (config, state, spool ve key dahil), ack'ler ve collector satırı **ancak ondan sonra** düşürür. Ters sırada yapılsaydı key anında geçersizleşir, agent kendisini uninstall etmesi gerektiğini hiç öğrenemez ve makinede sonsuza dek çalışmaya devam ederdi. +Bir cihazı kaldırmak satırı hemen silmez. Önce kuyruğa bir `delete` komutu girer; agent komutu poll'da alır, **önce** ack'ler ve `200` cevabını gördükten sonra yerelini temizler. Collector satırı ancak o ack ile düşürür. Sıra her iki yönde de kritiktir: satır erken silinseydi key anında geçersizleşir ve agent kendisini kaldırması gerektiğini hiç öğrenemezdi; yerel temizlik ack'ten önce yapılsaydı key ile birlikte ack'i gönderme imkânı da giderdi ve satır sunucuda ölümsüz kalırdı. + +Temizliğin ikinci yarısı ise agent'ın yetkisi **dışındadır**: servis yetkisiz bir kullanıcıyla, `NoNewPrivileges=yes` ve `ProtectSystem=strict` altında çalışır — kendi kurulumunu kaldıramaz, systemd'ye dokunamaz. Bu yüzden agent yalnızca yazabildiği tek yere, kendi state dizinine bir işaret dosyası bırakır; root tarafında bekleyen bir systemd `path` unit'i onu görür ve `uninstall.sh`'i çalıştırır. Böylece delete uçtan uca tamamlanır ama agent'ın yetkisi bir gram artmaz. **Agent'ın sınırları vardır.** Disk'teki spool hem yaş (10 gün) hem boyut (200 MB) ile sınırlanmış bir ring buffer'dır; sınır aşılınca en eski kayıt düşer. İzlediği makinenin disk'ini dolduran bir monitoring aracı, açıklaması beklenen outage'a kendisi sebep olmuş olur. diff --git a/agent/__init__.py b/agent/__init__.py index f9b7604..eabc0e2 100644 --- a/agent/__init__.py +++ b/agent/__init__.py @@ -27,4 +27,8 @@ # # 0.x = MVP, sözleşmeler henüz değişebilir. Wire payload'ları dondurulduğunda # 1.0.0'a çıkılacak. -__version__ = "0.1.0" +# +# 0.2.0: agent artık uzaktan komut alıyor (pause/resume/delete). Davranış farkı +# teşhis için önemli — dashboard'dan verilen bir pause'un neden işlemediğinin +# ilk cevabı "cihazda hangi sürüm var?" sorusudur. +__version__ = "0.2.0" diff --git a/agent/__main__.py b/agent/__main__.py index 627fc21..884f3a3 100644 --- a/agent/__main__.py +++ b/agent/__main__.py @@ -56,6 +56,19 @@ def main(argv: list[str] | None = None) -> int: return _run_verify(config) store = StateStore() + + # Bu cihaza `delete` komutu uygulanmışsa agent bir daha açılmaz. İşaret + # dosyası dururken açılsaydı, kaydı silinmiş cihaz her poll'da 401 alan bir + # süreç olarak sonsuza kadar dönerdi. Çıkış kodu 0: systemd + # Restart=on-failure ile çalışıyor, yani yeniden başlatmaz. + if store.is_deleted(): + print( + f"Bu cihaz silindi ({store.deleted_marker_path}); agent başlatılmadı.\n" + "Kurulumu tamamen kaldırmak için: sudo /opt/tracebox/uninstall.sh --yes", + flush=True, + ) + return EXIT_OK + # Tek platform seçimi buradadır: Windows desteği geldiğinde değişecek satır # bu, döngü değil. log_source = JournaldSource() diff --git a/agent/core/commands.py b/agent/core/commands.py new file mode 100644 index 0000000..6b1f1a0 --- /dev/null +++ b/agent/core/commands.py @@ -0,0 +1,235 @@ +""" +Komutlar — dashboard'un verdiği talimatları alır ve uygular. + +Kuyruk tek yönlüdür: dashboard `commands` tablosuna bir satır ekler, agent +`GET /commands` ile onu alır, uygular ve id'sini ack olarak geri yollar. Sunucu +ack'i görene kadar aynı komutu vermeye devam eder — yani teslim garantisi +burada da **at-least-once**'tır ve uygulama bu yüzden **idempotent** olmalıdır: +zaten duraklatılmış bir agent'a gelen ikinci `pause` hiçbir şey değiştirmez. + +Poll'ün kendi HTTP istemcisi var, Shipper'ınki kullanılmaz: gönderim durmuşken +(pause) ya da backoff sırasında bile komut sorulmaya devam eder — durursa +`resume` cihaza hiç ulaşamaz. + +`state.logging_enabled`'ın TEK YAZARI burasıdır. Sunucudaki +`devices.logging_enabled` sütunu bunun kopyasıdır ve ancak ack ulaştığında +güncellenir. +""" + +from __future__ import annotations + +from dataclasses import dataclass, field + +import httpx + +COMMANDS_PATH = "/commands" + +# Poll, döngünün turunu bloklar; collector yanıt vermiyorsa tur uzamasın. +REQUEST_TIMEOUT_SECONDS = 10.0 + +PAUSE = "pause" +RESUME = "resume" +DELETE = "delete" + +# pause/resume'un state karşılığı. `delete` burada yok: o bir ayar değişikliği +# değil, agent'ın kendini kapatması. +LOGGING_STATE_BY_TYPE = {PAUSE: False, RESUME: True} + + +class CommandError(Exception): + """Komutlar sorulamadı — ağ, yetki ya da beklenmeyen yanıt.""" + + +@dataclass(frozen=True) +class Command: + """Sunucudan gelen tek bir komut. Sözleşme iki alandan ibaret (§4.2).""" + + id: str + type: str + + +@dataclass(frozen=True) +class CommandResult: + """Bir poll turunun sonucu. + + `applied_ids`: ack'lenmeyi bekleyen YENİ id'ler. + `state_changed`: state.json diske yazılmalı mı. + `deleted`: cihaz kaydı silindi, agent duracak. + """ + + applied_ids: list[str] = field(default_factory=list) + state_changed: bool = False + deleted: bool = False + + +class CommandPoller: + """`GET /commands` — kendi istemcisiyle, backoff'suz.""" + + def __init__(self) -> None: + self._client = httpx.Client(timeout=REQUEST_TIMEOUT_SECONDS) + + def fetch(self, config) -> list[Command]: + """Bekleyen komutları getirir; sorun varsa CommandError fırlatır. + + Backoff yok: poll zaten seyrektir (varsayılan 10 sn) ve kapalı bir + collector'a atılan istek başarısız olup geçer. Gönderimdeki backoff'un + sebebi birikmiş veriyi tekrar tekrar yüklemeye çalışmaktı; burada + gövde boş. + """ + url = f"{config.collector_url.rstrip('/')}{COMMANDS_PATH}" + headers = {"Authorization": f"Bearer {config.device_key}"} + + try: + response = self._client.get(url, headers=headers) + except httpx.HTTPError as error: + raise CommandError(f"bağlanılamadı ({error.__class__.__name__})") from error + + if response.status_code == 401: + raise CommandError("cihaz anahtarı reddedildi (401)") + + if response.status_code != 200: + raise CommandError(f"HTTP {response.status_code}") + + try: + body = response.json() + except ValueError as error: + raise CommandError("yanıt JSON değil") from error + + return _parse(body) + + def close(self) -> None: + self._client.close() + + +def _parse(body) -> list[Command]: + """Yanıt gövdesini Command listesine çevirir. + + Biçimi bozuk tek bir kayıt turu düşürmez, yalnızca kendisi atlanır: bir + komut yüzünden diğerlerini (özellikle `resume`u) kaybetmek daha kötüdür. + """ + if not isinstance(body, dict) or not isinstance(body.get("commands"), list): + raise CommandError("yanıt beklenen şekilde değil") + + commands = [] + for item in body["commands"]: + if not isinstance(item, dict): + continue + identifier, type_ = item.get("id"), item.get("type") + if isinstance(identifier, str) and isinstance(type_, str): + commands.append(Command(id=identifier, type=type_)) + + return commands + + +def apply_commands(commands, *, config, state, store, spool, shipper, log) -> CommandResult: + """Komutları sırayla uygular ve turun sonucunu döndürür. + + Sunucu satırları `created_at` artan sırada verir; liste bu sırayla işlenir, + yani aynı turda pause + resume geldiyse son verilen kazanır. + """ + applied: list[str] = [] + state_changed = False + + for command in commands: + if command.type == DELETE: + if _delete(command, config=config, store=store, spool=spool, shipper=shipper, log=log): + return CommandResult(applied_ids=applied, state_changed=state_changed, deleted=True) + # Ack gitmediyse silme yarım kalmasın diye tur burada biter: komut + # sunucuda `pending` durur ve bir sonraki poll'da tekrar gelir. + # Arkasındaki komutları uygulamanın da anlamı yok, cihaz gidiyor. + break + + if command.type in LOGGING_STATE_BY_TYPE: + wanted = LOGGING_STATE_BY_TYPE[command.type] + if state.logging_enabled != wanted: + state.logging_enabled = wanted + state_changed = True + log(f"[cmd] {command.type} uygulandı — logging_enabled={wanted}") + else: + # Ack henüz ulaşmadığı için tekrar gönderilmiş komut. Uygulama + # idempotent: durum zaten istenen değerde. + log(f"[cmd] {command.type} zaten uygulanmış — ack tekrar denenecek") + + # Zaten ack listesinde olsa bile buraya yazılır: komutun tekrar + # gelmesi ack'in ulaşmadığı anlamına gelir, yani tekrar denenmeli. + if command.id not in applied: + applied.append(command.id) + continue + + # Bilinmeyen tür ACK EDİLMEZ. Ack etseydik sunucu komutu `applied` + # sayar ve bir daha vermezdi; yani agent'ın anlamadığı bir talimat + # sessizce uygulanmış görünürdü. Ack edilmeyince komut kuyrukta kalır + # ve agent güncellendiğinde uygulanır. + log(f"[cmd] bilinmeyen komut türü '{command.type}' — yok sayıldı (ack edilmedi)") + + return CommandResult(applied_ids=applied, state_changed=state_changed) + + +def ack_now(applied_ids: list[str], *, config, shipper, log) -> list[str]: + """Uygulanan komutları hemen ack'ler; onaylanan id'leri döndürür. + + Normalde ack'ler `POST /ingest` gövdesine biner (§4.2) — ama pause hâlinde + gönderim durur ve o gövde hiç yola çıkmaz. Ack de gitmezse sunucu komutu + her poll'da yeniden verir ve `devices.logging_enabled` kopyası hiç + güncellenmez, yani dashboard cihazı hâlâ "çalışıyor" gösterir. + + Bu yüzden ack telemetriye bağlanmaz: uygulanan bir komut, gönderim açık da + kapalı da olsa kendi küçük isteğiyle bildirilir. Bu bir KONTROL mesajıdır + (spec'in `delete` ack'i için tanıdığı istisnanın aynı gerekçesi), ölçüm + değil — gövdesinde tek bir ölçüm satırı bile yok. + + Başarısız olursa id'ler state'te kalır ve bir sonraki gönderime piggyback + olur; kayıp yok, yalnızca gecikme. + """ + if not applied_ids: + return [] + + result = shipper.send_acks(config, applied_ids) + if not result.ok: + log(f"[cmd] ack gönderilemedi: {result.detail} — sonraki gönderime bırakıldı") + return [] + + log(f"[cmd] {len(applied_ids)} komut ack'lendi") + return list(applied_ids) + + +def _delete(command, *, config, store, spool, shipper, log) -> bool: + """Silme akışı. SIRA KRİTİKTİR (CLAUDE.md §11 Boşluk E). + + 1. ack gönderilir — pause'da bile, çünkü bu teardown kontrol mesajıdır, + 2. 200 alınır: collector `devices` satırını sildi (CASCADE ile veri de), + 3. yerel veri silinir (spool + state), + 4. kaldırma işareti bırakılır: gerisi root'un işi (aşağıda). + + Ters sıra iki türlü de bozardı: önce yerel silme yapılsaydı anahtar giderdi + ve ack hiç atılamazdı; sunucu satırı erken silinseydi agent 401 alır, + komutu hiç göremezdi. + """ + log("[cmd] delete alındı — önce ack gönderiliyor.") + + result = shipper.send_acks(config, [command.id]) + if not result.ok: + log(f"[cmd] delete ack gönderilemedi: {result.detail} — silme ertelendi.") + return False + + log("[cmd] ack onaylandı: cihaz kaydı sunucudan silindi. Yerel temizlik başlıyor.") + + spool.wipe() + log(f"[cmd] spool silindi: {spool.path}") + + store.wipe() + log(f"[cmd] state silindi: {store.path}") + + # Kalanı (systemd servisi, /opt, /etc ve anahtarın kendisi) agent SİLEMEZ: + # yetkisiz `tracebox` kullanıcısıyla, NoNewPrivileges=yes ve + # ProtectSystem=strict altında çalışır — yazabildiği tek yer kendi state + # dizinidir. Bu bilinçli bir kilit; delete uğruna gevşetilmez. + # + # Bunun yerine oraya bir işaret dosyası bırakılır: root tarafında çalışan + # tracebox-uninstall.path onu görür ve uninstall.sh'i çalıştırır. Böylece + # kaldırma, agent'ın yetkisini artırmadan tamamlanır. + marker = store.mark_deleted() + log(f"[cmd] kaldırma işareti bırakıldı: {marker}") + log("[cmd] kaldırma tamamlanmazsa elle: sudo /opt/tracebox/uninstall.sh --yes") + + return True diff --git a/agent/core/loop.py b/agent/core/loop.py index 6ba6f0a..fe00a0b 100644 --- a/agent/core/loop.py +++ b/agent/core/loop.py @@ -7,8 +7,9 @@ Tick sabit 1 saniyedir ve config'den okunmaz; collect/send/poll aralıkları birbirinden bağımsız sayaçlardır. -M4 KAPSAMI: ölçümler ve sistem logları spool'a yazılıp collector'a gönderilir. -Komut poll'u M6'da, acil flush M7'de bu döngüye bağlanacak. +M6 KAPSAMI: ölçümler ve sistem logları spool'a yazılıp collector'a gönderilir, +komutlar (pause/resume/delete) poll ile alınıp uygulanır. Acil flush M7'de +bu döngüye bağlanacak. """ from __future__ import annotations @@ -21,8 +22,10 @@ from dataclasses import asdict from agent import __version__ +from agent.core import commands as commands_module from agent.core import inventory as inventory_module from agent.core.clock import utc_now_iso +from agent.core.commands import CommandError, CommandPoller from agent.core.config import Config, ConfigLoader from agent.core.inventory import Inventory from agent.core.metrics import MetricSample, MetricsCollector @@ -177,13 +180,76 @@ def _collect_logs(source: LogSource, spool: Spool, state: State, store: StateSto _log(f"[logs] {len(records)} kayıt ({_level_summary(records)})") -def _poll_commands(config: Config) -> None: - """Komut kuyruğunu yoklama adımı — M6'da GET /commands buraya bağlanır. +def _prune_acked(state: State, store: StateStore, acked: list[str]) -> None: + """Onaylanan komut id'lerini state'ten düşer. + + "Listeyi boşalt" DEĞİL, "onaylananları çıkar": aradaki fark bugün görünmez + (tek döngü var, gönderim sırasında yeni komut uygulanamaz) ama kural + bugünden doğru yazılırsa ileride bozulmaz. + """ + if not acked: + return + + remaining = [value for value in state.applied_command_ids if value not in set(acked)] + if remaining != state.applied_command_ids: + state.applied_command_ids = remaining + store.save(state) + + +def _poll_commands( + poller: CommandPoller, + config: Config, + state: State, + store: StateStore, + spool: Spool, + shipper: Shipper, +) -> bool: + """Komutları sorar, uygular ve ack'ler. Cihaz silindiyse True döner. Pause'da da çalışır; durmasaydı resume ve delete komutları cihaza hiç - ulaşamazdı. + ulaşamazdı — duraklatılmış bir agent'ı geri açmanın başka yolu kalmazdı. """ - _log(f"[poll] komut sorulacak (her {config.command_poll_seconds} sn) — M6") + try: + commands = poller.fetch(config) + except CommandError as error: + # Komut alınamaması toplamayı ve gönderimi durdurmaz; tur komutsuz + # geçer, sorun bir sonraki poll'da yeniden denenir. + _log(f"[poll] komutlar alınamadı: {error}") + return False + + if not commands: + return False + + _log(f"[poll] {len(commands)} komut: {', '.join(command.type for command in commands)}") + + result = commands_module.apply_commands( + commands, + config=config, + state=state, + store=store, + spool=spool, + shipper=shipper, + log=_log, + ) + + if result.deleted: + return True + + # Zaten listede olan id yeniden eklenmez; ack denemesi yine de yapılır, + # çünkü sunucu ack'i görene kadar aynı komutu vermeye devam eder. + new_ids = [value for value in result.applied_ids if value not in state.applied_command_ids] + if new_ids: + state.applied_command_ids.extend(new_ids) + + if result.state_changed or new_ids: + store.save(state) + + _prune_acked( + state, + store, + commands_module.ack_now(result.applied_ids, config=config, shipper=shipper, log=_log), + ) + return False def _send_inventory( @@ -206,6 +272,11 @@ def _send_spool( ) -> None: """Spool'u gönderir; başarılıysa last_send'i günceller.""" result = shipper.send_pending(config, state.applied_command_ids) + + # Onaylanan ack'ler gönderim yarıda kalsa bile düşer: ilk istek 200 almış + # olabilir, o id'ler artık sunucuda `applied`. + _prune_acked(state, store, result.acked) + if not result.ok: _log( f"[send] gönderilemedi: {result.detail} — {spool.count()} kayıt bekliyor, " @@ -238,6 +309,7 @@ def run(loader: ConfigLoader, store: StateStore, log_source: LogSource) -> None: max_size_mb=config.spool_max_size_mb, ) shipper = Shipper(spool) + poller = CommandPoller() _log(f"[start] TraceBox agent {__version__}") _log(f"[start] config: {loader.path}") @@ -266,6 +338,7 @@ def run(loader: ConfigLoader, store: StateStore, log_source: LogSource) -> None: next_send = now + config.send_interval_seconds # --- KALP ATIŞI --- + deleted = False try: while not stop.is_set(): now = time.monotonic() @@ -280,7 +353,11 @@ def run(loader: ConfigLoader, store: StateStore, log_source: LogSource) -> None: next_collect = now + config.collect_interval_seconds if now >= next_poll: - _poll_commands(config) + if _poll_commands(poller, config, state, store, spool, shipper): + # delete uygulandı: cihaz kaydı sunucudan silindi, yerel + # veri temizlendi. Toplamaya devam etmenin anlamı yok. + deleted = True + break next_poll = now + config.command_poll_seconds # Aşağısı yalnızca gönderim açıkken çalışır. Pause sırasında @@ -303,6 +380,11 @@ def run(loader: ConfigLoader, store: StateStore, log_source: LogSource) -> None: # --- KAPANIŞ --- # Durum diskte zaten günceldir (her save anında yazıldı); burada yalnızca # açık dosya ve bağlantılar kapatılır. + poller.close() shipper.close() spool.close() - _log("[stop] döngü durdu.") + + if deleted: + _log("[stop] cihaz silindi — agent duruyor, systemd yeniden başlatmayacak.") + else: + _log("[stop] döngü durdu.") diff --git a/agent/core/shipper.py b/agent/core/shipper.py index 8ad5768..0b748ea 100644 --- a/agent/core/shipper.py +++ b/agent/core/shipper.py @@ -12,7 +12,7 @@ from __future__ import annotations import time -from dataclasses import dataclass +from dataclasses import dataclass, field import httpx @@ -46,11 +46,19 @@ @dataclass(frozen=True) class SendResult: - """Bir gönderim turunun sonucu.""" + """Bir gönderim turunun sonucu. + + `acked`: 200 ile onaylanmış komut id'leri. Çağıran bunları state'ten düşer. + "Listeyi boşalt" değil "gönderdiklerini düş" kuralı bilerek böyle yazıldı: + gönderim sırasında yeni bir komut uygulanmış olsaydı, boşaltmak o ack'i + sessizce yutardı ([[decisions]] → "`applied_command_ids` 200 alınınca + küçültülür"). + """ ok: bool sent: int = 0 detail: str = "" + acked: list[str] = field(default_factory=list) class Shipper: @@ -84,6 +92,7 @@ def send_pending(self, config, applied_command_ids: list[str]) -> SendResult: """ sent = 0 acks = list(applied_command_ids) + confirmed: list[str] = [] for _ in range(MAX_BATCHES_PER_CYCLE): records = self._spool.take(BATCH_ROWS) @@ -93,17 +102,32 @@ def send_pending(self, config, applied_command_ids: list[str]) -> SendResult: payload = _build_payload(records, acks) ok, detail = self._post(config, INGEST_PATH, payload) if not ok: - return SendResult(ok=False, sent=sent, detail=detail) + return SendResult(ok=False, sent=sent, detail=detail, acked=confirmed) # 200 alındı: kayıtlar artık sunucuda, spool'dan düşebilirler. self._spool.ack([record.uuid for record in records]) sent += len(records) + confirmed.extend(acks) acks = [] if len(records) < BATCH_ROWS: break - return SendResult(ok=True, sent=sent) + return SendResult(ok=True, sent=sent, acked=confirmed) + + def send_acks(self, config, command_ids: list[str]) -> SendResult: + """Yalnızca ack taşıyan küçük bir `POST /ingest`. + + Gövdesinde tek bir ölçüm satırı yoktur; bu bir KONTROL mesajıdır. Bu + yüzden gönderim kapalıyken (pause) de atılabilir: pause'un durdurduğu + şey telemetridir, agent'ın "komutu uyguladım" demesi değil. + """ + if not command_ids: + return SendResult(ok=True) + + payload = _build_payload([], list(command_ids)) + ok, detail = self._post(config, INGEST_PATH, payload) + return SendResult(ok=ok, detail=detail, acked=list(command_ids) if ok else []) def close(self) -> None: self._client.close() diff --git a/agent/core/spool.py b/agent/core/spool.py index 482b7ca..03d6bd3 100644 --- a/agent/core/spool.py +++ b/agent/core/spool.py @@ -128,6 +128,17 @@ def size_bytes(self) -> int: def close(self) -> None: self._connection.close() + def wipe(self) -> None: + """Bekleyen her şeyi ve dosyanın kendisini siler (`delete` komutu). + + Tablo boşaltmak yetmez: silinen satırların izi WAL dosyasında kalabilir + ve bu veri artık cihazda BULUNMAMALIDIR — kullanıcı cihazı sildi. + Bağlantı da kapatılır; açık bir tanıtıcı silinen dosyayı canlı tutar. + """ + self.close() + for suffix in ("", "-wal", "-shm"): + Path(f"{self._path}{suffix}").unlink(missing_ok=True) + def _configure(self) -> None: """Bağlantı ayarları ve şema. diff --git a/agent/core/state.py b/agent/core/state.py index f9a4bcd..1bd682b 100644 --- a/agent/core/state.py +++ b/agent/core/state.py @@ -15,6 +15,8 @@ from pathlib import Path from typing import Any +from agent.core.clock import utc_now_iso + # Üretimdeki çalışma dizini. TRACEBOX_STATE_DIR tanımlıysa onun değeri kullanılır # (config.py'deki override ile aynı mantık). DEFAULT_STATE_DIR = Path("/var/lib/tracebox") @@ -23,6 +25,11 @@ STATE_FILENAME = "state.json" LOCK_FILENAME = "agent.lock" +# `delete` komutu uygulandığında bırakılan işaret. Agent'ın yazabildiği tek +# dizinde durur ve iki iş görür: root tarafındaki tracebox-uninstall.path onu +# görüp kaldırmayı başlatır, agent da yeniden açılırsa buradan durur. +DELETED_FILENAME = "deleted" + @dataclass class State: @@ -134,6 +141,36 @@ def save(self, state: State) -> None: finally: os.close(dir_fd) + def wipe(self) -> None: + """state.json'u siler — `delete` komutunun yerel temizliğinin parçası. + + Dizin bırakılır: kaldırma işareti oraya yazılacak ve dizini silmek + zaten agent'ın yetkisi dışında (uninstall.sh'in işi). + """ + self._path.unlink(missing_ok=True) + self._path.with_name(f"{STATE_FILENAME}.tmp").unlink(missing_ok=True) + + @property + def deleted_marker_path(self) -> Path: + return self._dir / DELETED_FILENAME + + def mark_deleted(self) -> Path: + """Kaldırma işaretini bırakır ve yolunu döndürür. + + İçeriği okunmaz; önemli olan dosyanın VAR olmasıdır (systemd path + unit'i `PathExists` ile bakar). Yine de zaman damgası yazılır: teşhis + sırasında "bu makine ne zaman silindi?" sorusunun tek cevabı budur. + """ + self._dir.mkdir(parents=True, exist_ok=True) + self.deleted_marker_path.write_text( + f"{utc_now_iso()} delete komutu uygulandı\n", encoding="utf-8" + ) + return self.deleted_marker_path + + def is_deleted(self) -> bool: + """Bu cihaz `delete` komutuyla silinmiş mi.""" + return self.deleted_marker_path.exists() + def _quarantine(self, exc: Exception) -> None: """Okunamayan state dosyasını kenara alır.""" quarantine_path = self._path.with_name(f"{STATE_FILENAME}.corrupt") diff --git a/agent/install.sh b/agent/install.sh index cbef2b3..49f78ae 100755 --- a/agent/install.sh +++ b/agent/install.sh @@ -31,6 +31,16 @@ STATE_DIR="/var/lib/tracebox" SERVICE_USER="tracebox" SERVICE_NAME="tracebox-agent.service" UNIT_PATH="/etc/systemd/system/${SERVICE_NAME}" + +# `delete` komutunun root tarafı. Agent yetkisiz çalıştığı için kendi kurulumunu +# kaldıramaz; yapabildiği tek şey state dizinine bir işaret dosyası bırakmaktır. +# Aşağıdaki path unit'i o dosyayı görür ve kaldırma servisini root olarak +# çalıştırır. +UNINSTALL_SERVICE="tracebox-uninstall.service" +UNINSTALL_PATH_UNIT="tracebox-uninstall.path" +UNINSTALL_SERVICE_PATH="/etc/systemd/system/${UNINSTALL_SERVICE}" +UNINSTALL_PATH_UNIT_PATH="/etc/systemd/system/${UNINSTALL_PATH_UNIT}" +DELETED_MARKER="${STATE_DIR}/deleted" TTY_DEVICE="/dev/tty" # --- Çıktı yardımcıları ---------------------------------------------------- @@ -182,7 +192,10 @@ tar -xzf "${TARBALL}" -C "${SOURCE_ROOT}" --strip-components=1 \ || fail "Kaynak arşivi açılamadı." SOURCE_AGENT="${SOURCE_ROOT}/agent" -for required in "__main__.py" "requirements.txt" "tracebox-agent.service"; do +for required in \ + "__main__.py" "requirements.txt" \ + "tracebox-agent.service" "tracebox-uninstall.service" "tracebox-uninstall.path" +do [[ -f "${SOURCE_AGENT}/${required}" ]] \ || fail "İndirilen arşiv eksik: agent/${required} yok." done @@ -222,6 +235,14 @@ for dir in "${INSTALL_DIR}" "${CONFIG_DIR}" "${STATE_DIR}"; do fi done +# Önceki kurulum `delete` ile bitmişse geride kaldırma işareti kalmış olabilir. +# Silinmeden path unit etkinleştirilirse, yeni kurulum daha ayağa kalkmadan +# kendini kaldırır. +if [[ -e "${DELETED_MARKER}" ]]; then + rm -f "${DELETED_MARKER}" + say "önceki kurulumdan kalan kaldırma işareti silindi" +fi + # Kod her kurulumda sıfırdan kopyalanır: eski sürümden kalan bir dosya # yenisinin yanında durmasın. rm -rf "${INSTALL_DIR}/agent" @@ -335,12 +356,21 @@ ROLLBACK_ENABLED=0 step "6/7 systemd servisi" install -m 644 "${INSTALL_DIR}/agent/tracebox-agent.service" "${UNIT_PATH}" +install -m 644 "${INSTALL_DIR}/agent/${UNINSTALL_SERVICE}" "${UNINSTALL_SERVICE_PATH}" +install -m 644 "${INSTALL_DIR}/agent/${UNINSTALL_PATH_UNIT}" "${UNINSTALL_PATH_UNIT_PATH}" systemctl daemon-reload systemctl enable "${SERVICE_NAME}" >/dev/null 2>&1 systemctl restart "${SERVICE_NAME}" say "kuruldu ve başlatıldı: ${SERVICE_NAME}" +# İzleyici şimdi başlar ve kaldırma işaretini beklemeye koyulur. `enable --now` +# olmadan yalnızca bir sonraki açılışta devreye girerdi: bu makinede verilen +# ilk `delete` komutu, makine yeniden başlatılana kadar tamamlanmazdı. +systemctl enable --now "${UNINSTALL_PATH_UNIT}" >/dev/null 2>&1 \ + || warn "${UNINSTALL_PATH_UNIT} etkinleştirilemedi; delete komutu geldiğinde kaldırma elle yapılmalı" +say "kaldırma izleyicisi etkin: ${UNINSTALL_PATH_UNIT} (delete komutunun root tarafı)" + # Servisin ilk saniyede düşüp düşmediğini görmek için kısa bir bekleme. sleep 2 if systemctl is-active --quiet "${SERVICE_NAME}"; then diff --git a/agent/tracebox-uninstall.path b/agent/tracebox-uninstall.path new file mode 100644 index 0000000..d4843e9 --- /dev/null +++ b/agent/tracebox-uninstall.path @@ -0,0 +1,20 @@ +[Unit] +Description=TraceBox — silinen cihazda kaldırmayı tetikler +Documentation=https://github.com/denisergocmen/tracebox + +[Path] +# Agent kendi kurulumunu KALDIRAMAZ: yetkisiz `tracebox` kullanıcısıyla, +# NoNewPrivileges=yes ve ProtectSystem=strict altında çalışır — yazabildiği tek +# yer /var/lib/tracebox'tır. `delete` komutunu uygulayınca oraya bu işaret +# dosyasını bırakır; root tarafında bekleyen bu birim onu görüp kaldırmayı +# başlatır. Böylece delete tamamlanır ama agent'ın yetkisi bir gram artmaz. +# +# PathExists (PathChanged değil): dosya kaldırma sırasında ya da makine +# kapalıyken oluşmuş olabilir. PathExists zaten var olan dosyayı da yakalar, +# yani agent silme işaretini bırakıp makine hemen kapansa bile kaldırma bir +# sonraki açılışta yapılır. +PathExists=/var/lib/tracebox/deleted +Unit=tracebox-uninstall.service + +[Install] +WantedBy=multi-user.target diff --git a/agent/tracebox-uninstall.service b/agent/tracebox-uninstall.service new file mode 100644 index 0000000..96d519e --- /dev/null +++ b/agent/tracebox-uninstall.service @@ -0,0 +1,17 @@ +[Unit] +Description=TraceBox Agent'ı kaldırır (delete komutunun root tarafı) +Documentation=https://github.com/denisergocmen/tracebox + +# Kaldırma betiği yoksa çalışacak bir şey de yok; birim `failed` yerine sessizce +# atlanır. +ConditionPathExists=/opt/tracebox/uninstall.sh + +[Service] +Type=oneshot + +# User= YOK — bu birim BİLEREK root çalışır: systemctl, /etc, /opt ve kullanıcı +# silme yetkisi ister. Kaldırmanın ayrı bir birim olması da şart: agent'ın +# servisi içinden çalıştırılsaydı, betiğin ilk işi olan `systemctl disable --now +# tracebox-agent.service` betiği kendi cgroup'uyla birlikte öldürür ve kaldırma +# daha ilk adımda yarım kalırdı. +ExecStart=/opt/tracebox/uninstall.sh --yes diff --git a/agent/uninstall.sh b/agent/uninstall.sh index 4b4bb6a..accf866 100755 --- a/agent/uninstall.sh +++ b/agent/uninstall.sh @@ -5,13 +5,23 @@ # Servisi durdurur, dosyaları ve yetkisiz kullanıcıyı siler. Cihaz anahtarı # config.toml içinde durduğu için bu betik anahtarı da yok eder. # -# İki yerden çağrılır: kullanıcı elle (sudo ./uninstall.sh) ve M6'daki `delete` -# komutu (self-uninstall). O yüzden --yes ile soru sormadan da çalışabilir. +# İki yerden çağrılır: kullanıcı elle (sudo ./uninstall.sh) ve `delete` komutu. +# İkinci yol dolaylıdır — agent yetkisiz çalıştığı için bu betiği kendisi +# çağıramaz; state dizinine bir işaret dosyası bırakır, root tarafındaki +# tracebox-uninstall.path onu görüp betiği çalıştırır. O çağrıda soru soracak +# kimse yoktur, bu yüzden --yes ile onaysız da çalışır. set -euo pipefail SERVICE_NAME="tracebox-agent.service" UNIT_PATH="/etc/systemd/system/${SERVICE_NAME}" + +# Kaldırmayı tetikleyen çift. Bu betik ÇOĞU ZAMAN tracebox-uninstall.service +# tarafından çalıştırılır — yani kendi birimini durdurmaya çalışmamalıdır. +UNINSTALL_SERVICE="tracebox-uninstall.service" +UNINSTALL_PATH_UNIT="tracebox-uninstall.path" +UNINSTALL_SERVICE_PATH="/etc/systemd/system/${UNINSTALL_SERVICE}" +UNINSTALL_PATH_UNIT_PATH="/etc/systemd/system/${UNINSTALL_PATH_UNIT}" INSTALL_DIR="/opt/tracebox" CONFIG_DIR="/etc/tracebox" STATE_DIR="/var/lib/tracebox" @@ -29,6 +39,7 @@ fail() { printf '\n✗ %s\n' "$*" >&2; exit 1; } if [[ "${1:-}" != "--yes" ]]; then printf 'TraceBox Agent kaldırılacak:\n' printf ' - %s durdurulup devre dışı bırakılacak\n' "${SERVICE_NAME}" + printf ' - %s izleyicisi kaldırılacak\n' "${UNINSTALL_PATH_UNIT}" printf ' - %s, %s, %s silinecek\n' "${INSTALL_DIR}" "${CONFIG_DIR}" "${STATE_DIR}" printf ' - %s kullanıcısı silinecek\n' "${SERVICE_USER}" printf '\nCihaz anahtarı da silinir; cihazı tekrar eklemek için yeni anahtar gerekir.\n' @@ -44,16 +55,30 @@ step "Servis durduruluyor" if command -v systemctl >/dev/null 2>&1; then systemctl disable --now "${SERVICE_NAME}" >/dev/null 2>&1 || true say "durduruldu ve devre dışı bırakıldı" + + # İzleyici de kapatılır; işaret dosyası birazdan silinecek dizinde duruyor. + systemctl disable --now "${UNINSTALL_PATH_UNIT}" >/dev/null 2>&1 || true + say "izleyici kapatıldı: ${UNINSTALL_PATH_UNIT}" + + # DİKKAT: ${UNINSTALL_SERVICE} durdurulmaz. Bu betiği şu anda o birim + # çalıştırıyor olabilir; `stop` demek kendi süreç ağacını öldürmek, yani + # kaldırmayı tam ortasında yarıda kesmek olurdu. Birimin [Install] bölümü de + # yok (yalnızca path unit tetikler), dolayısıyla disable edilecek bir bağ + # zaten yok — unit dosyasını silmek yeterli. else say "systemctl yok — atlandı" fi -if [[ -f "${UNIT_PATH}" ]]; then - rm -f "${UNIT_PATH}" - say "unit dosyası silindi: ${UNIT_PATH}" -fi +for unit_file in "${UNIT_PATH}" "${UNINSTALL_PATH_UNIT_PATH}" "${UNINSTALL_SERVICE_PATH}"; do + if [[ -f "${unit_file}" ]]; then + rm -f "${unit_file}" + say "unit dosyası silindi: ${unit_file}" + fi +done if command -v systemctl >/dev/null 2>&1; then + # daemon-reload çalışan birimi durdurmaz: systemd onu bellekte tutar, yani + # bu betik kendi unit dosyası silindikten sonra da sorunsuz devam eder. systemctl daemon-reload || true # Durdurulmuş servisin `failed` kaydı systemctl listesinde kalmasın. systemctl reset-failed "${SERVICE_NAME}" >/dev/null 2>&1 || true diff --git a/collector/auth.py b/collector/auth.py index ce5fac8..59f4e74 100644 --- a/collector/auth.py +++ b/collector/auth.py @@ -72,7 +72,6 @@ class DeviceIdentity: id: str account_id: str device_name: str - pending_delete: bool @dataclass(frozen=True) @@ -234,7 +233,6 @@ async def require_device( id=row["id"], account_id=row["account_id"], device_name=row["device_name"], - pending_delete=row["pending_delete"], ) diff --git a/collector/db_access.py b/collector/db_access.py new file mode 100644 index 0000000..4850708 --- /dev/null +++ b/collector/db_access.py @@ -0,0 +1,40 @@ +""" +Cihaz uçlarının paylaştığı iki küçük yardımcı: sunucu saati ve 503 sarmalayıcı. + +`endpoints_ingest` ile `endpoints_commands` aynı iki davranışa ihtiyaç duyar; +ikisini de tek bir uç noktanın modülünde tutmak ya kopyalamayı ya da bir uç +noktanın diğerinin iç adlarını içe aktarmasını gerektirirdi. +""" + +from __future__ import annotations + +from datetime import datetime, timezone + +from fastapi import HTTPException, status + +from supabase_client import SupabaseError + + +def server_now() -> str: + """`last_seen` ve `applied_at` için sunucu saati (ISO 8601, UTC). + + Agent'ın damgası kullanılmaz: saati kaymış bir cihaz aksi halde offline + tespitini yanıltırdı. + """ + return datetime.now(timezone.utc).isoformat() + + +async def call_or_503(operation): + """Supabase çağrısını çalıştırır, sonucunu döndürür; hatayı 503'e çevirir. + + 503, agent'a "veriyi tut, sonra tekrar dene" demektir; spool kaydı ancak 200 + sonrası silinir. Aynı kural komut ack'i için de geçerlidir: ack yazılamazsa + agent id'yi state'inde tutar ve bir sonraki gönderimde tekrar yollar. + """ + try: + return await operation() + except SupabaseError as error: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail="Kayıt şu an yazılamıyor.", + ) from error diff --git a/collector/endpoints_commands.py b/collector/endpoints_commands.py new file mode 100644 index 0000000..c6d7a22 --- /dev/null +++ b/collector/endpoints_commands.py @@ -0,0 +1,101 @@ +""" +Komut uçları: GET /commands ve `POST /ingest` gövdesindeki ack'in işlenmesi. + +Kuyruk tek yönlü akar: dashboard `commands` tablosuna `pending` bir satır ekler, +agent onu poll ile alır, uygular ve id'sini bir sonraki gönderimde geri yollar. +Bu modül o döngünün sunucu tarafındaki iki ucunu tutar — komutu VERMEK ve +uygulandığını KAYDETMEK. + +Ack'in işlenmesi burada durur, `endpoints_ingest`'te değil: komutun durumunu +değiştiren tek yer burasıdır, ingest yalnızca çağırır. +""" + +from __future__ import annotations + +from dataclasses import dataclass +from uuid import UUID + +from fastapi import APIRouter + +from auth import AuthenticatedDevice, DeviceIdentity +from db_access import call_or_503, server_now +from supabase_client import get_client + +router = APIRouter() + +# Komut türünün `logging_enabled` sunucu kopyasına karşılığı. `delete` burada +# yok: o komut satırı hiç bırakmaz, cihaz kaydının tamamı silinir. +LOGGING_STATE_BY_TYPE = {"pause": False, "resume": True} + +DELETE_COMMAND = "delete" + + +@dataclass(frozen=True) +class AckResult: + """Ack işlendikten sonra çağıranın bilmesi gerekenler. + + `device_deleted`: cihaz satırı silindi — geriye yazılacak bir satır yok. + `logging_enabled`: pause/resume uygulandıysa yeni durum, yoksa None. + """ + + device_deleted: bool = False + logging_enabled: bool | None = None + + +@router.get("/commands") +async def get_commands(device: AuthenticatedDevice) -> dict: + """Cihazın bekleyen komutlarını döndürür. + + `last_seen` burada da tazelenir. Duraklatılmış (pause) bir agent veri + göndermeyi durdurur ama komut poll'ünü sürdürür — yoksa `resume` ona hiç + ulaşmazdı. Yalnızca ingest `last_seen` yazsaydı, o cihaz susduğu için + dashboard'da "offline" görünürdü; oysa erişilebilir durumda. Sütunun anlamı + bu yüzden "veri geldi" değil, "cihazdan haber alındı". + """ + client = get_client() + + commands = await call_or_503(lambda: client.list_pending_commands(device.id)) + await call_or_503( + lambda: client.update_device(device.id, {"last_seen": server_now()}) + ) + + return {"commands": [{"id": row["id"], "type": row["type"]} for row in commands]} + + +async def process_acks( + client, device: DeviceIdentity, command_ids: list[UUID] +) -> AckResult: + """Ack edilen komutları `applied` yapar ve doğurduğu durumu uygular. + + Sıra önemlidir: + 1. komutlar `applied` işaretlenir ve güncellenen satırlar geri okunur, + 2. dönen satırlarda `delete` varsa cihaz satırı silinir (CASCADE), + 3. yoksa pause/resume'un getirdiği `logging_enabled` çağırana bildirilir. + + Uygulanacak durum, agent'ın gövdesinden değil veritabanından dönen `type` + alanından türetilir; agent'ın gönderdiği tek şey komut id'sidir. + + Aynı turda birden fazla pause/resume gelirse EN SON verilen komut kazanır: + satırlar `created_at` artan sırada döner ve döngü sonuncuyu yazar. + """ + if not command_ids: + return AckResult() + + applied = await call_or_503( + lambda: client.mark_commands_applied( + device.id, [str(cid) for cid in command_ids], server_now() + ) + ) + + logging_enabled: bool | None = None + for row in applied: + if row["type"] == DELETE_COMMAND: + # Silme her şeyin önüne geçer: satır gidince pause/resume'un + # yazılacağı yer de kalmaz. + await call_or_503(lambda: client.delete_device(device.id)) + return AckResult(device_deleted=True) + + if row["type"] in LOGGING_STATE_BY_TYPE: + logging_enabled = LOGGING_STATE_BY_TYPE[row["type"]] + + return AckResult(logging_enabled=logging_enabled) diff --git a/collector/endpoints_ingest.py b/collector/endpoints_ingest.py index 9052d56..9a6f654 100644 --- a/collector/endpoints_ingest.py +++ b/collector/endpoints_ingest.py @@ -7,16 +7,17 @@ from __future__ import annotations -from datetime import datetime, timezone from typing import Annotated, Any, Literal from uuid import UUID -from fastapi import APIRouter, HTTPException, status +from fastapi import APIRouter from pydantic import AwareDatetime, BaseModel, ConfigDict, Field from auth import AuthenticatedDevice, DeviceIdentity +from db_access import call_or_503, server_now +from endpoints_commands import process_acks from version import COLLECTOR_VERSION -from supabase_client import SupabaseError, get_client +from supabase_client import get_client router = APIRouter() @@ -117,9 +118,12 @@ class IngestIn(_Payload): crash_snapshots: Annotated[ list[CrashSnapshotIn], Field(max_length=MAX_ROWS_PER_TABLE) ] = Field(default_factory=list) - # M6'da işlenecek: ack edilen komutlar 'applied' yapılacak, delete komutu - # cihaz satırının silinmesini tetikleyecek. - applied_command_ids: list[UUID] = Field(default_factory=list) + # Uygulanmış komutların id'leri (ack). Ayrı bir uç yerine bu gövdeye + # binerler: agent zaten düzenli olarak buraya istek atıyor, ack için ikinci + # bir tur açmak boşuna trafik olurdu. + applied_command_ids: Annotated[ + list[UUID], Field(max_length=MAX_ROWS_PER_TABLE) + ] = Field(default_factory=list) @router.post("/inventory") @@ -129,9 +133,9 @@ async def post_inventory(payload: InventoryIn, device: AuthenticatedDevice) -> d Envanter zaman serisi değildir: her gönderim bir öncekinin üzerine yazar. """ fields = payload.model_dump(mode="json") - fields["last_seen"] = _server_now() + fields["last_seen"] = server_now() - await _write(lambda: get_client().update_device(device.id, fields)) + await call_or_503(lambda: get_client().update_device(device.id, fields)) return {"status": "ok", "device_name": device.device_name} @@ -151,9 +155,21 @@ async def post_ingest(payload: IngestIn, device: AuthenticatedDevice) -> dict: for table, items in tables: rows = [_row(item, device) for item in items] - await _write(lambda table=table, rows=rows: client.insert_rows(table, rows)) - - await _write(lambda: client.update_device(device.id, {"last_seen": _server_now()})) + await call_or_503(lambda table=table, rows=rows: client.insert_rows(table, rows)) + + # Ack, veri yazıldıktan SONRA işlenir: `delete` ack'i cihaz satırını siler ve + # o satıra bağlı her şey CASCADE ile gider. Ters sırada, aynı gövdede gelen + # ölçümler silinmiş bir cihaza yazılmaya çalışılırdı (foreign key hatası). + ack = await process_acks(client, device, payload.applied_command_ids) + + if not ack.device_deleted: + fields: dict[str, Any] = {"last_seen": server_now()} + if ack.logging_enabled is not None: + # pause/resume'un sunucu kopyası. Doğruluk kaynağı agent'ın + # state.json'ı; bu sütun dashboard rozeti için tutulur ve bu yüzden + # komut VERİLDİĞİNDE değil, agent UYGULADIĞINI bildirdiğinde yazılır. + fields["logging_enabled"] = ack.logging_enabled + await call_or_503(lambda: client.update_device(device.id, fields)) return { "status": "ok", @@ -190,29 +206,5 @@ def _row(item: BaseModel, device: DeviceIdentity) -> dict[str, Any]: row[ID_COLUMN] = row.pop(UUID_FIELD) row["device_id"] = device.id row["account_id"] = device.account_id - row["received_at"] = _server_now() + row["received_at"] = server_now() return row - - -def _server_now() -> str: - """`last_seen` için sunucu saati. - - Agent'ın damgası kullanılmaz: saati kaymış bir cihaz aksi halde offline - tespitini yanıltırdı. - """ - return datetime.now(timezone.utc).isoformat() - - -async def _write(operation) -> None: - """Supabase çağrısını çalıştırır; hatayı 503'e çevirir. - - 503, agent'a "veriyi tut, sonra tekrar dene" demektir; spool kaydı ancak 200 - sonrası silinir. - """ - try: - await operation() - except SupabaseError as error: - raise HTTPException( - status_code=status.HTTP_503_SERVICE_UNAVAILABLE, - detail="Kayıt şu an yazılamıyor.", - ) from error diff --git a/collector/main.py b/collector/main.py index 63a0610..847ecc4 100644 --- a/collector/main.py +++ b/collector/main.py @@ -8,9 +8,7 @@ Bağlı router'lar: endpoints_device.py POST /devices (user JWT) endpoints_ingest.py POST /inventory, POST /ingest, GET /verify (device key) - -Sonraki milestone'da eklenecek: - M6 -> endpoints_commands.py GET /commands (device key) + endpoints_commands.py GET /commands (device key) """ from __future__ import annotations @@ -21,6 +19,7 @@ from fastapi import FastAPI import supabase_client +from endpoints_commands import router as commands_router from endpoints_device import router as device_router from endpoints_ingest import router as ingest_router from version import COLLECTOR_VERSION @@ -58,6 +57,7 @@ async def lifespan(app: FastAPI): app.include_router(device_router) app.include_router(ingest_router) +app.include_router(commands_router) @app.get("/") diff --git a/collector/supabase_client.py b/collector/supabase_client.py index 6fb09d6..902b888 100644 --- a/collector/supabase_client.py +++ b/collector/supabase_client.py @@ -58,13 +58,17 @@ # hesabın dashboard'una veri enjekte edebilirdi. # key_hash — kimlik kanıtının kendisi; cihaz kendi anahtarını seçemez. # device_name — dashboard'un alanı (db/rls.sql: grant update (device_name)). -# logging_enabled — pause/resume durumu; sunucu kopyasını kimin yazacağı M6'da -# pending_delete karara bağlanacak. O karar verilene kadar collector bu -# sütunlara dokunmaz. +# +# logging_enabled M6'da listeye EKLENDİ (bkz. md/memory/decisions.md → "Komutlar +# (M6)"): pause/resume durumunun sunucu kopyasını, agent komutu ack'leyince +# collector yazar. Değer istek gövdesinden gelmez — `commands` satırındaki +# `type` alanından türetilir; agent'ın gönderdiği tek şey komut id'sidir. # # Buraya sütun eklemek bilinçli bir güvenlik kararıdır. DEVICE_WRITABLE_COLUMNS = frozenset( { + # komut ack'inin türettiği durum + "logging_enabled", # agent'ın envanterden bildirdikleri (InventoryIn ile aynı 14 alan) "cpu_model", "cpu_cores_physical", @@ -124,7 +128,7 @@ async def find_device_by_key_hash(self, key_hash: str) -> dict[str, Any] | None: "/devices", params={ "key_hash": f"eq.{key_hash}", - "select": "id,account_id,device_name,key_hash,logging_enabled,pending_delete", + "select": "id,account_id,device_name,key_hash,logging_enabled", "limit": "1", }, ) @@ -184,6 +188,77 @@ async def insert_device(self, row: dict[str, Any]) -> dict[str, Any]: return rows[0] + async def delete_device(self, device_id: str) -> None: + """Cihaz satırını siler. + + Yalnızca `delete` komutunun ack'i bu yola girer. Satırla birlikte + metrics / logs / crash_snapshots / commands satırları da gider — şemadaki + foreign key'ler `on delete cascade` taşır. + """ + await self._request( + "DELETE", + "/devices", + params={"id": f"eq.{device_id}"}, + headers={"Prefer": PREFER_MINIMAL}, + ) + + async def list_pending_commands(self, device_id: str) -> list[dict[str, Any]]: + """Cihazın bekleyen komutlarını eskiden yeniye döndürür. + + Filtre `device_id` ile sınırlıdır: service key RLS'i bypass ettiği için + satır sahipliğini bu sorgu kurar. Sıra `created_at` artan — agent birden + fazla komutu tek turda alırsa verildikleri sırayla uygular. + """ + response = await self._request( + "GET", + "/commands", + params={ + "device_id": f"eq.{device_id}", + "status": "eq.pending", + "select": "id,type", + "order": "created_at.asc", + }, + ) + return response.json() + + async def mark_commands_applied( + self, device_id: str, command_ids: list[str], applied_at: str + ) -> list[dict[str, Any]]: + """Verilen komutları `applied` yapar ve GERÇEKTEN güncellenen satırları döndürür. + + İki filtre birden uygulanır: `id` listede olacak VE satır bu cihaza ait + olacak. İkincisi olmasaydı bir cihaz, başka bir cihazın komutunu + ack'leyip onu uygulanmış gösterebilirdi — kurbanın agent'ı komutu hiç + görmezdi (agent yalnızca `pending` olanları çeker). + + `status` filtresi BİLEREK yok: zaten `applied` olan bir komut yeniden + yazılır. Ack tekrarı normaldir (200 yolda kaybolursa agent aynı id'yi + bir daha yollar) ve dönen satırlar delete akışının tetikleyicisidir — + `applied` olanları eleseydik, ilk turda satır silme adımı yarım kalan + bir delete bir daha asla tamamlanamazdı. Bedeli: tekrar eden ack + `applied_at`'i tazeler. + + Dönen satırlar `type` alanını taşır; çağıran hangi durumun uygulanacağını + (pause/resume/delete) buradan öğrenir — agent'ın gövdesinden değil. + """ + if not command_ids: + return [] + + id_list = ",".join(command_ids) + response = await self._request( + "PATCH", + "/commands", + params={ + "id": f"in.({id_list})", + "device_id": f"eq.{device_id}", + "select": "id,type,created_at", + "order": "created_at.asc", + }, + json={"status": "applied", "applied_at": applied_at}, + headers={"Prefer": PREFER_REPRESENTATION}, + ) + return response.json() + async def insert_rows(self, table: str, rows: list[dict[str, Any]]) -> None: """Satırları ekler; `id` çakışanları sessizce atlar.""" if not rows: diff --git a/collector/version.py b/collector/version.py index b20ec6c..d5be0d8 100644 --- a/collector/version.py +++ b/collector/version.py @@ -13,4 +13,4 @@ # Agent'ın bildirdiği agent_version'dan bağımsız, collector'ın kendi sürümü. # Ayakta olan sürüm GET /verify üzerinden doğrulanır — kimliksiz uçlar (/ ve # /health) sürüm döndürmez. -COLLECTOR_VERSION = "0.3.0" +COLLECTOR_VERSION = "0.4.0" diff --git a/db/migrations/0003_drop_pending_delete.sql b/db/migrations/0003_drop_pending_delete.sql new file mode 100644 index 0000000..0520d4d --- /dev/null +++ b/db/migrations/0003_drop_pending_delete.sql @@ -0,0 +1,47 @@ +-- ============================================================================= +-- TraceBox — db/migrations/0003_drop_pending_delete.sql +-- +-- NE: `devices` tablosundan `pending_delete` sütunu kaldırılır. +-- +-- NEDEN: Sütun, silme akışının ara durumunu ("emir verildi, cihaz henüz +-- duymadı") işaretlemek için tasarlanmıştı. Ama yazacak kimse yoktu: +-- +-- * dashboard yazamaz — db/rls.sql `devices` üzerindeki UPDATE yetkisini +-- tek sütuna daraltıyor: grant update (device_name). +-- * collector yazmıyor — supabase_client.py'deki DEVICE_WRITABLE_COLUMNS +-- listesinin dışında; yanındaki not "kimin yazacağı M6'da karara +-- bağlanacak" diyordu. +-- +-- Yazan olmayınca sütun her satırda `false` kaldı: kimsenin dolduramadığı +-- bir bayrak, hiçbir şey işaretlemiyor. Bilgi zaten `commands` tablosunda +-- duruyor — bekleyen silme = (type='delete' and status='pending'). +-- Kopyayı yaşatmak için ya kilit gevşetilecekti ya yeni bir mekanizma +-- (trigger / ayrı endpoint) eklenecekti; kaldırmak üçünün de bedelini +-- sıfırlıyor. +-- +-- VERİ KAYBI: yok. Sütun hiç yazılmadı, tüm satırlarda varsayılan `false`. +-- +-- SIRA (ÖNEMLİ): bu migration, `pending_delete`'i artık SELECT etmeyen +-- collector sürümü Fly'a deploy EDİLDİKTEN SONRA çalıştırılır. Ters sırada +-- çalıştırılırsa canlı collector olmayan bir sütunu istemeye devam eder ve +-- cihaz kimliği doğrulanamaz (device key ile gelen her istek hata alır). +-- +-- TARİH: 2026-08-26 — M6 (Komutlar). Karar: md/memory/decisions.md. +-- ============================================================================= + +alter table public.devices drop column if exists pending_delete; + +-- ============================================================================= +-- DOĞRULAMA — SIFIR satır dönmeli. Bir satır dönerse sütun hâlâ duruyordur. +-- ============================================================================= +select column_name + from information_schema.columns + where table_schema = 'public' + and table_name = 'devices' + and column_name = 'pending_delete'; + +-- ----------------------------------------------------------------------------- +-- GERİ ALMA (çalıştırılmaz — sadece kayıt): +-- alter table public.devices +-- add column pending_delete boolean not null default false; +-- ----------------------------------------------------------------------------- diff --git a/db/rls.sql b/db/rls.sql index f905721..0aa2e98 100644 --- a/db/rls.sql +++ b/db/rls.sql @@ -89,8 +89,8 @@ create policy del_devices on public.devices -- ----------------------------------------------------------------------------- -- SORUN: yukarıdaki upd_devices politikası SATIR düzeyinde çalışır; hangi -- SÜTUNLARIN yazılabileceğini söylemez. Tek başına bırakılırsa kullanıcı kendi --- cihaz satırının key_hash / last_seen / logging_enabled / pending_delete --- alanlarını da yazabilir. Bu: +-- cihaz satırının key_hash / last_seen / logging_enabled alanlarını da +-- yazabilir. Bu: -- * "single writer" ilkesini bozar (bu alanların yazarı collector'dır), -- * last_seen'i elle ileri atarak offline tespitini yanıltmayı, -- * key_hash'i değiştirip kendi agent'ını kilitlemeyi mümkün kılar. diff --git a/db/schema.sql b/db/schema.sql index 48c288f..35c45c2 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -112,10 +112,10 @@ create table devices ( -- (ör. now() - last_seen > 2 dk ise offline sayılır). last_seen timestamptz, - -- Delete akışının ara durumu. Dashboard "sil" dediğinde satır - -- hemen silinmez: kuyruğa bir 'delete' komutu girer ve bu bayrak true olur. - -- Satırı, agent komutu uygulayıp ack'ledikten sonra collector siler. - pending_delete boolean not null default false, + -- Silme akışının ara durumu ("emir verildi, cihaz henüz duymadı") burada + -- bir bayrakla tutulmuyor: bilgi `commands` tablosunda zaten var — + -- type='delete' ve status='pending' olan satır. Satırı, agent komutu + -- uygulayıp ack'ledikten sonra collector siler (migration 0003). created_at timestamptz not null default now() ); diff --git a/tests/test_commands.py b/tests/test_commands.py new file mode 100644 index 0000000..87528f2 --- /dev/null +++ b/tests/test_commands.py @@ -0,0 +1,471 @@ +"""agent/core/commands.py — komut teslimi, uygulanması ve `delete` sırası. + +Üç ayrı soru sınanır: + +1. **Poll dayanıklı mı?** Komut alınamaması agent'ı durdurmamalı; bozuk bir + kayıt diğerlerini (özellikle `resume`u) düşürmemeli. +2. **Uygulama idempotent mi?** Sunucu ack'i görene kadar aynı komutu vermeye + devam eder; ikinci `pause` hiçbir şeyi değiştirmemeli ama ack'i yeniden + denenmeli — denenmezse duraklatılmış bir agent sonsuza kadar öyle kalır. +3. **`delete` sırası doğru mu?** Önce ack, sonra yerel silme (CLAUDE.md §11 + Boşluk E). Ters sırada anahtar giderdi ve ack hiç atılamazdı: sunucudaki + cihaz kaydı ölümsüz kalırdı. + +Ağa çıkılmaz. Poll için httpx'in MockTransport'u, ack için sahte bir shipper +kullanılır; gerçek olan state ve spool'dur — silmenin gerçekten olup olmadığı +ancak diskteki dosyalara bakarak ölçülebilir. +""" + +import httpx +import pytest + +from agent.core import commands as commands_module +from agent.core.commands import Command, CommandError, CommandPoller +from agent.core.config import Config +from agent.core.shipper import SendResult +from agent.core.spool import RECORD_METRIC, Spool +from agent.core.state import StateStore + +CONFIG = Config(collector_url="https://collector.test", device_key="tbx_live_test") + + +class FakeShipper: + """`send_acks` çağrılarını kaydeden sahte shipper. + + `on_send`, ack atıldığı ANDA çalışır. Silme sırasını ölçmenin başka yolu + yok: sıranın doğru olduğunu ancak "ack sırasında yerel veri hâlâ duruyor + muydu?" sorusuna bakarak anlarız. + """ + + def __init__(self, *outcomes: bool, on_send=None) -> None: + self._outcomes = list(outcomes) + self._on_send = on_send + self.calls: list[list[str]] = [] + + def send_acks(self, config, command_ids: list[str]) -> SendResult: + self.calls.append(list(command_ids)) + if self._on_send is not None: + self._on_send() + + ok = self._outcomes.pop(0) if self._outcomes else True + return SendResult( + ok=ok, + detail="" if ok else "HTTP 500", + acked=list(command_ids) if ok else [], + ) + + +@pytest.fixture +def spool(tmp_path): + instance = Spool(tmp_path, max_age_days=10, max_size_mb=200) + yield instance + instance.close() + + +@pytest.fixture +def store(tmp_path): + return StateStore(tmp_path) + + +@pytest.fixture +def messages() -> list[str]: + """Toplanan konsol satırları — `log` yerine geçer.""" + return [] + + +def make_poller(handler) -> CommandPoller: + """Poller'ı kurar ve HTTP istemcisini sahte collector'a bağlar.""" + poller = CommandPoller() + poller.close() + poller._client = httpx.Client(transport=httpx.MockTransport(handler)) + return poller + + +def respond(body, status: int = 200): + """Sabit bir yanıt döndüren MockTransport işleyicisi.""" + + def handler(request: httpx.Request) -> httpx.Response: + if isinstance(body, str): + return httpx.Response(status, text=body) + return httpx.Response(status, json=body) + + return handler + + +def apply(commands, *, state, store, spool, shipper, messages, config=CONFIG): + return commands_module.apply_commands( + commands, + config=config, + state=state, + store=store, + spool=spool, + shipper=shipper, + log=messages.append, + ) + + +# --- Poll ------------------------------------------------------------------ + + +def test_poll_asks_the_commands_endpoint_with_the_device_key(): + """Cihaz kimliği yalnızca anahtardan türetilir (CLAUDE.md §11 Boşluk A). + + URL'de device_id taşınsaydı, herkes başkasının cihaz id'sini yazarak onun + komutlarını okumayı deneyebilirdi. + """ + seen: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen.append(request) + return httpx.Response(200, json={"commands": []}) + + make_poller(handler).fetch(CONFIG) + + assert str(seen[0].url) == "https://collector.test/commands" + assert seen[0].headers["Authorization"] == "Bearer tbx_live_test" + + +def test_commands_arrive_in_the_order_the_server_gave_them(): + """Sıra korunur: aynı turda pause + resume geldiyse sonuncusu kazanmalı.""" + poller = make_poller( + respond({"commands": [{"id": "a", "type": "pause"}, {"id": "b", "type": "resume"}]}) + ) + + assert poller.fetch(CONFIG) == [Command(id="a", type="pause"), Command(id="b", type="resume")] + + +def test_empty_queue_is_not_an_error(): + """Beklenen durum bu: çoğu poll boş döner.""" + assert make_poller(respond({"commands": []})).fetch(CONFIG) == [] + + +@pytest.mark.parametrize( + "status, expected", + [(401, "401"), (500, "500"), (404, "404")], + ids=["anahtar reddedildi", "sunucu hatası", "yol yok"], +) +def test_failed_poll_is_reported_not_swallowed(status, expected): + """Hata sessizce "komut yok"a dönüşmemeli. + + Dönüşseydi, anahtarı reddedilen bir agent hiçbir uyarı vermeden sonsuza + kadar boş kuyruk görürdü. + """ + poller = make_poller(respond({}, status)) + + with pytest.raises(CommandError) as error: + poller.fetch(CONFIG) + + assert expected in str(error.value) + + +def test_connection_failure_is_reported(): + """Collector kapalıyken poll patlamaz, anlaşılır bir hata verir.""" + + def handler(request: httpx.Request) -> httpx.Response: + raise httpx.ConnectError("bağlantı yok") + + with pytest.raises(CommandError): + make_poller(handler).fetch(CONFIG) + + +def test_response_that_is_not_json_is_reported(): + """Araya giren bir vekil (proxy) HTML hata sayfası döndürebilir.""" + with pytest.raises(CommandError): + make_poller(respond("502")).fetch(CONFIG) + + +def test_response_with_an_unexpected_shape_is_reported(): + """Sözleşme `{"commands": [...]}`; başka bir şekil komut kaybı demektir.""" + with pytest.raises(CommandError): + make_poller(respond({"items": []})).fetch(CONFIG) + + +def test_a_malformed_command_does_not_take_the_others_down(): + """Bozuk tek kayıt atlanır, turun geri kalanı uygulanır. + + Tüm turu düşürmek `resume`u kaybetmek olurdu: duraklatılmış agent, bozuk + bir komut yüzünden bir daha hiç açılmazdı. + """ + poller = make_poller( + respond( + { + "commands": [ + {"id": 7, "type": "pause"}, + "çöp", + {"type": "pause"}, + {"id": "b", "type": "resume"}, + ] + } + ) + ) + + assert poller.fetch(CONFIG) == [Command(id="b", type="resume")] + + +# --- pause / resume -------------------------------------------------------- + + +def test_pause_stops_sending_and_is_acked(store, spool, messages): + state = store.load() + + result = apply( + [Command(id="k1", type="pause")], + state=state, + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert state.logging_enabled is False + assert result.state_changed is True + assert result.applied_ids == ["k1"] + + +def test_resume_turns_sending_back_on(store, spool, messages): + state = store.load() + state.logging_enabled = False + + apply( + [Command(id="k2", type="resume")], + state=state, + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert state.logging_enabled is True + + +def test_a_repeated_command_changes_nothing_but_is_acked_again(store, spool, messages): + """Bu dosyadaki en kritik testlerden biri. + + Komutun tekrar gelmesi, ack'in ULAŞMADIĞI anlamına gelir. Uygulama + idempotenttir (durum zaten istenen değerde, değiştirilecek bir şey yok) ama + ack yeniden denenmelidir. Denenmezse sunucu komutu her poll'da yeniden + verir, agent her seferinde "zaten uygulanmış" deyip susar ve + `devices.logging_enabled` kopyası hiç güncellenmez: dashboard duraklatılmış + cihazı sonsuza kadar "çalışıyor" gösterir. + """ + state = store.load() + state.logging_enabled = False + + result = apply( + [Command(id="k1", type="pause")], + state=state, + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert state.logging_enabled is False + assert result.state_changed is False, "değişmeyen durum için diske yazılıyor" + assert result.applied_ids == ["k1"], "tekrar gelen komut ack edilmiyor" + + +def test_the_last_command_of_a_round_wins(store, spool, messages): + """Sunucu satırları created_at sırasıyla verir; son verilen geçerlidir.""" + state = store.load() + + result = apply( + [Command(id="k1", type="pause"), Command(id="k2", type="resume")], + state=state, + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert state.logging_enabled is True + assert result.applied_ids == ["k1", "k2"], "ara komut ack edilmemiş" + + +def test_an_unknown_command_type_is_not_acked(store, spool, messages): + """Anlaşılmayan komut ack EDİLMEZ. + + Ack edilseydi sunucu onu `applied` sayar ve bir daha vermezdi: agent'ın + hiç uygulamadığı bir talimat, dashboard'da uygulanmış görünürdü. Ack + edilmeyince komut kuyrukta bekler ve agent güncellendiğinde çalışır. + """ + state = store.load() + + result = apply( + [Command(id="k9", type="reboot")], + state=state, + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert result.applied_ids == [] + assert result.state_changed is False + + +# --- ack ------------------------------------------------------------------ + + +def test_applied_commands_are_acked_right_away(messages): + """Ack telemetriye bağlanmaz — pause hâlinde gönderim durur, ack durmaz.""" + shipper = FakeShipper() + + confirmed = commands_module.ack_now( + ["k1", "k2"], config=CONFIG, shipper=shipper, log=messages.append + ) + + assert shipper.calls == [["k1", "k2"]] + assert confirmed == ["k1", "k2"] + + +def test_nothing_is_acked_when_there_is_nothing_to_ack(messages): + shipper = FakeShipper() + + assert commands_module.ack_now([], config=CONFIG, shipper=shipper, log=messages.append) == [] + assert shipper.calls == [] + + +def test_a_failed_ack_confirms_nothing(messages): + """Başarısız ack'te id'ler state'te kalır ve sonraki gönderime piggyback olur.""" + confirmed = commands_module.ack_now( + ["k1"], config=CONFIG, shipper=FakeShipper(False), log=messages.append + ) + + assert confirmed == [] + + +# --- delete --------------------------------------------------------------- + + +def fill(spool: Spool, store: StateStore) -> None: + """Silinecek gerçek veri: birkaç spool kaydı ve diskte bir state.json.""" + spool.add(RECORD_METRIC, {"uuid": "kayit-1", "measured_at": "2026-08-26T10:00:00Z"}) + store.save(store.load()) + + +def test_delete_acks_before_wiping_anything(store, spool, messages): + """Bu dosyanın en kritik testi (CLAUDE.md §11 Boşluk E). + + Sıra tersine dönerse anahtar (ve state) ack atılmadan silinir; ack hiç + gitmez, collector cihaz satırını hiç silmez ve kayıt sunucuda ölümsüz + kalır. Kod yine çalışır, hiçbir hata görünmez — bozukluğun tek bekçisi bu + testtir. + """ + fill(spool, store) + alive_at_ack: list[bool] = [] + + shipper = FakeShipper(on_send=lambda: alive_at_ack.append(spool.path.exists())) + + apply( + [Command(id="sil", type="delete")], + state=store.load(), + store=store, + spool=spool, + shipper=shipper, + messages=messages, + ) + + assert shipper.calls == [["sil"]] + assert alive_at_ack == [True], "ack atılmadan önce yerel veri silinmiş" + + +def test_delete_removes_the_local_data(store, spool, messages): + """Cihaz silindi: ölçümler de state de makinede kalmamalı.""" + fill(spool, store) + + result = apply( + [Command(id="sil", type="delete")], + state=store.load(), + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert result.deleted is True + assert not spool.path.exists() + assert not store.path.exists() + + +def test_delete_leaves_the_marker_that_triggers_the_uninstall(store, spool, messages): + """Kaldırmanın root tarafını başlatan tek şey bu dosyadır. + + Agent yetkisiz çalışır ve /opt, /etc ile systemd'ye dokunamaz; bırakabildiği + tek iz, kendi state dizinindeki bu işaret. Düşerse cihaz sunucudan silinir + ama servis makinede çalışmaya devam eder — her poll'da 401 alan bir hayalet. + """ + fill(spool, store) + + apply( + [Command(id="sil", type="delete")], + state=store.load(), + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert store.is_deleted() is True + assert store.deleted_marker_path.exists() + + +def test_delete_is_postponed_when_the_ack_does_not_get_through(store, spool, messages): + """Ack ulaşmadıysa hiçbir şey silinmez. + + Silinseydi anahtar giderdi: komut sunucuda `pending` kalır, agent bir daha + ack atamaz ve cihaz kaydı elle silinene kadar orada durur. + """ + fill(spool, store) + + result = apply( + [Command(id="sil", type="delete")], + state=store.load(), + store=store, + spool=spool, + shipper=FakeShipper(False), + messages=messages, + ) + + assert result.deleted is False + assert result.applied_ids == [] + assert spool.path.exists() + assert store.path.exists() + assert store.is_deleted() is False + + +def test_delete_ends_the_round(store, spool, messages): + """Silmeden sonraki komutlar uygulanmaz — uygulanacak bir cihaz kalmadı.""" + state = store.load() + state.logging_enabled = False + fill(spool, store) + + apply( + [Command(id="sil", type="delete"), Command(id="k2", type="resume")], + state=state, + store=store, + spool=spool, + shipper=FakeShipper(), + messages=messages, + ) + + assert state.logging_enabled is False, "silinen cihazda sonraki komut uygulanmış" + + +def test_a_postponed_delete_also_ends_the_round(store, spool, messages): + """Ack başarısızsa da tur biter: cihaz gitmek üzere, arkasını uygulamak anlamsız.""" + state = store.load() + state.logging_enabled = False + fill(spool, store) + + result = apply( + [Command(id="sil", type="delete"), Command(id="k2", type="resume")], + state=state, + store=store, + spool=spool, + shipper=FakeShipper(False), + messages=messages, + ) + + assert state.logging_enabled is False + assert result.applied_ids == [] diff --git a/tests/test_endpoints_commands.py b/tests/test_endpoints_commands.py new file mode 100644 index 0000000..285a123 --- /dev/null +++ b/tests/test_endpoints_commands.py @@ -0,0 +1,387 @@ +""" +collector/endpoints_commands.py — komut kuyruğunun sunucu tarafı. + +İki uç sınanır: komutun VERİLMESİ (`GET /commands`) ve uygulandığının +KAYDEDİLMESİ (`POST /ingest` gövdesindeki `applied_command_ids`). İkisi ayrı +dosyada durmaz çünkü aynı döngünün iki yarısıdır; ack ingest'e binmiş olsa da +mantığı bu modülde yaşar. + +Buradaki hataların hepsi sessizdir: komut teslim edilmezse agent hiç +duraklamaz, ack işlenmezse aynı komut sonsuza kadar tekrar gelir, `delete` +yanlış sırada işlenirse cihaz kaydı ortada kalır. Hiçbiri istisna fırlatmaz. + +Supabase taklit ediliyor (`get_client` sahte bir istemciyle değiştiriliyor) ve +cihaz doğrulaması bağımlılık override'ıyla sabitleniyor — anahtarın kendisi +zaten test_collector_security.py'de sınanıyor, burada sınanan uçların MANTIĞI. +""" + +from __future__ import annotations + +import pytest +from fastapi.testclient import TestClient + +import auth +import endpoints_commands +import endpoints_ingest +from main import app +from supabase_client import DEVICE_WRITABLE_COLUMNS, SupabaseError + +DEVICE_ID = "33333333-3333-3333-3333-333333333333" +ACCOUNT_ID = "11111111-1111-1111-1111-111111111111" + +# Kurbanın cihazı — sahte istemcinin "yanlış cihaza dokunuldu mu?" kontrolü için. +OTHER_DEVICE_ID = "44444444-4444-4444-4444-444444444444" + +PAUSE_ID = "aaaaaaaa-0000-0000-0000-000000000001" +RESUME_ID = "aaaaaaaa-0000-0000-0000-000000000002" +DELETE_ID = "aaaaaaaa-0000-0000-0000-000000000003" + + +def command_row(command_id: str, command_type: str, created_at: str = "2026-08-26T10:00:00Z"): + """`commands` tablosundan dönen satır.""" + return {"id": command_id, "type": command_type, "created_at": created_at} + + +class FakeSupabase: + """Çağrıları sırasıyla kaydeden, ağa çıkmayan sahte istemci. + + `calls` listesi hem NE yapıldığını hem HANGİ SIRADA yapıldığını tutar; + ack testlerinin bir kısmı yalnızca sıraya bakar. + """ + + def __init__(self) -> None: + self.calls: list[tuple] = [] + self.pending: list[dict] = [] + self.applied: list[dict] = [] + self.list_error: SupabaseError | None = None + self.mark_error: SupabaseError | None = None + self.delete_error: SupabaseError | None = None + + @property + def call_names(self) -> list[str]: + return [call[0] for call in self.calls] + + def last(self, name: str) -> tuple: + """Adı verilen son çağrı — yoksa test AssertionError ile durur.""" + matches = [call for call in self.calls if call[0] == name] + assert matches, f"{name} hiç çağrılmadı" + return matches[-1] + + async def list_pending_commands(self, device_id: str) -> list[dict]: + self.calls.append(("list_pending_commands", device_id)) + if self.list_error is not None: + raise self.list_error + return self.pending + + async def mark_commands_applied( + self, device_id: str, command_ids: list[str], applied_at: str + ) -> list[dict]: + self.calls.append(("mark_commands_applied", device_id, command_ids, applied_at)) + if self.mark_error is not None: + raise self.mark_error + return self.applied + + async def delete_device(self, device_id: str) -> None: + self.calls.append(("delete_device", device_id)) + if self.delete_error is not None: + raise self.delete_error + + async def update_device(self, device_id: str, fields: dict) -> None: + self.calls.append(("update_device", device_id, fields)) + + async def insert_rows(self, table: str, rows: list[dict]) -> None: + self.calls.append(("insert_rows", table, rows)) + + +@pytest.fixture +def fake_supabase(monkeypatch): + """İki uç noktayı da aynı sahte veritabanına bağlar.""" + client = FakeSupabase() + monkeypatch.setattr(endpoints_commands, "get_client", lambda: client) + monkeypatch.setattr(endpoints_ingest, "get_client", lambda: client) + app.dependency_overrides[auth.require_device] = lambda: auth.DeviceIdentity( + id=DEVICE_ID, account_id=ACCOUNT_ID, device_name="dizustu" + ) + yield client + app.dependency_overrides.clear() + + +@pytest.fixture +def client(fake_supabase): + return TestClient(app) + + +def poll(client): + """GET /commands — agent'ın komut sorması.""" + return client.get("/commands") + + +def ack(client, *command_ids, **body): + """POST /ingest — ack piggyback. Gövde varsayılanı boş bir gönderimdir.""" + return client.post( + "/ingest", json={"applied_command_ids": list(command_ids), **body} + ) + + +# --- GET /commands: komutun teslimi ---------------------------------------- + + +def test_pending_commands_are_delivered(client, fake_supabase): + """Bekleyen komut agent'a ulaşmalı — yoksa pause/resume/delete hiç çalışmaz.""" + fake_supabase.pending = [command_row(PAUSE_ID, "pause")] + + body = poll(client).json() + + assert body["commands"] == [{"id": PAUSE_ID, "type": "pause"}] + + +def test_empty_queue_returns_an_empty_list(client): + """Komut yokken de aynı şekil dönmeli; agent tek bir kod yolu izler.""" + assert poll(client).json() == {"commands": []} + + +def test_only_the_authenticated_device_is_queried(client, fake_supabase): + """Sorgu doğrulanmış anahtarın cihazına kilitli. + + İstekte device_id YOK (Boşluk A) — olsaydı bir cihaz başkasının komutlarını + çekip onun agent'ından gizleyebilirdi. + """ + poll(client) + assert fake_supabase.last("list_pending_commands")[1] == DEVICE_ID + + +def test_only_id_and_type_reach_the_agent(client, fake_supabase): + """Satırın diğer sütunları yanıta sızmamalı — sözleşme iki alandan ibaret.""" + fake_supabase.pending = [command_row(PAUSE_ID, "pause")] + + command = poll(client).json()["commands"][0] + + assert set(command) == {"id", "type"} + + +def test_poll_refreshes_last_seen(client, fake_supabase): + """Bu dosyadaki en kolay kaçırılan davranış. + + Duraklatılmış agent veri göndermez ama komut sormayı sürdürür. `last_seen` + yalnızca ingest'te yazılsaydı o cihaz dashboard'da OFFLINE görünürdü — oysa + ulaşılabilir durumda ve `resume` komutunu bekliyor. Sütunun anlamı "veri + geldi" değil, "cihazdan haber alındı". + """ + poll(client) + + name, device_id, fields = fake_supabase.last("update_device") + assert device_id == DEVICE_ID + assert "last_seen" in fields + + +def test_database_failure_becomes_503(client, fake_supabase): + """Komut listesi okunamazsa agent'a "sonra tekrar dene" denir.""" + fake_supabase.list_error = SupabaseError("kesinti") + + assert poll(client).status_code == 503 + + +# --- Ack: pause / resume --------------------------------------------------- + + +def test_acked_commands_are_marked_applied(client, fake_supabase): + """Ack işlenmezse aynı komut her poll'da tekrar gelir — sonsuz döngü.""" + fake_supabase.applied = [command_row(PAUSE_ID, "pause")] + + ack(client, PAUSE_ID) + + name, device_id, ids, applied_at = fake_supabase.last("mark_commands_applied") + assert ids == [PAUSE_ID] + assert applied_at + + +def test_ack_is_scoped_to_the_authenticated_device(client, fake_supabase): + """Bir cihaz başkasının komutunu ack'leyemez. + + Ack'lenen komut `applied` olur ve bir daha teslim edilmez; kurbanın agent'ı + o komutu hiç görmezdi. + """ + fake_supabase.applied = [command_row(PAUSE_ID, "pause")] + + ack(client, PAUSE_ID) + + assert fake_supabase.last("mark_commands_applied")[1] == DEVICE_ID + + +def test_pause_ack_writes_the_server_copy(client, fake_supabase): + """pause uygulandığında `devices.logging_enabled` false olur.""" + fake_supabase.applied = [command_row(PAUSE_ID, "pause")] + + ack(client, PAUSE_ID) + + assert fake_supabase.last("update_device")[2]["logging_enabled"] is False + + +def test_resume_ack_writes_the_server_copy(client, fake_supabase): + """resume uygulandığında true olur.""" + fake_supabase.applied = [command_row(RESUME_ID, "resume")] + + ack(client, RESUME_ID) + + assert fake_supabase.last("update_device")[2]["logging_enabled"] is True + + +def test_plain_ingest_does_not_touch_logging_enabled(client, fake_supabase): + """Ack yoksa sütuna dokunulmaz. + + Her gönderimde yazılsaydı sunucu kopyası, agent'ın gerçek durumunu değil + son isteğin varsayılanını yansıtırdı. + """ + ack(client) + + assert "logging_enabled" not in fake_supabase.last("update_device")[2] + + +def test_no_ack_skips_the_commands_table(client, fake_supabase): + """Boş ack listesi için veritabanına hiç gidilmez — her 30 saniyede bir + boşa yazma isteği anlamına gelirdi.""" + ack(client) + + assert "mark_commands_applied" not in fake_supabase.call_names + + +def test_state_comes_from_the_database_not_the_request(client, fake_supabase): + """Agent yalnızca id gönderir; ne yapılacağını veritabanındaki `type` söyler. + + Tür gövdeden okunsaydı, cihaz kendi sunucu kopyasını istediği gibi + yazabilirdi — hiç verilmemiş bir `resume`u bildirmek gibi. + """ + fake_supabase.applied = [command_row(PAUSE_ID, "pause")] + + ack(client, PAUSE_ID) + + assert fake_supabase.last("update_device")[2]["logging_enabled"] is False + + +def test_the_last_command_wins(client, fake_supabase): + """Aynı turda pause ve resume ack'lenirse SON verilen komut kazanır. + + Satırlar `created_at` artan sırada döner; döngü sonuncuyu yazar. + """ + fake_supabase.applied = [ + command_row(PAUSE_ID, "pause", "2026-08-26T10:00:00Z"), + command_row(RESUME_ID, "resume", "2026-08-26T10:05:00Z"), + ] + + ack(client, PAUSE_ID, RESUME_ID) + + assert fake_supabase.last("update_device")[2]["logging_enabled"] is True + + +def test_written_fields_are_allowed_by_the_client(client, fake_supabase): + """Uç noktanın yazdığı her sütun `DEVICE_WRITABLE_COLUMNS` içinde olmalı. + + Sahte istemci allowlist'i çalıştırmaz; gerçek istemci çalıştırır ve + listede olmayan bir sütun canlıda ValueError olur — testte değil. + """ + fake_supabase.applied = [command_row(PAUSE_ID, "pause")] + + ack(client, PAUSE_ID) + + assert set(fake_supabase.last("update_device")[2]) <= DEVICE_WRITABLE_COLUMNS + + +def test_failed_ack_becomes_503(client, fake_supabase): + """Ack yazılamazsa 200 dönmemeli. + + 200 dönseydi agent id'yi state'inden silerdi ve komut sonsuza kadar + `pending` kalırdı: agent uyguladığını sanır, sunucu beklemeye devam eder. + """ + fake_supabase.mark_error = SupabaseError("kesinti") + + assert ack(client, PAUSE_ID).status_code == 503 + + +# --- Ack: delete ----------------------------------------------------------- + + +def test_delete_ack_removes_the_device_row(client, fake_supabase): + """Agent kendini sildiğini bildirince cihaz kaydı da gider (CASCADE).""" + fake_supabase.applied = [command_row(DELETE_ID, "delete")] + + ack(client, DELETE_ID) + + assert fake_supabase.last("delete_device")[1] == DEVICE_ID + + +def test_delete_ack_deletes_only_the_authenticated_device(client, fake_supabase): + """Silinen satır her zaman anahtarın sahibi — gövdeden gelen bir id değil.""" + fake_supabase.applied = [command_row(DELETE_ID, "delete")] + + ack(client, DELETE_ID) + + assert OTHER_DEVICE_ID not in fake_supabase.last("delete_device") + + +def test_delete_ack_does_not_write_the_row_afterwards(client, fake_supabase): + """Satır silindikten sonra `last_seen` yazılmaya çalışılmamalı. + + PostgREST silinmiş satıra yapılan güncellemeyi hata saymaz — sessizce sıfır + satır günceller. Yani bu yanlış canlıda hiç görünmez; yalnızca boşa bir + istek olarak kalır ve delete akışının sırası bozulduğunda ilk kanıt budur. + """ + fake_supabase.applied = [command_row(DELETE_ID, "delete")] + + ack(client, DELETE_ID) + + order = fake_supabase.call_names + assert "update_device" not in order[order.index("delete_device"):] + + +def test_delete_wins_over_pause_in_the_same_batch(client, fake_supabase): + """Aynı turda delete varsa cihaz gider; pause'un yazılacağı satır kalmaz.""" + fake_supabase.applied = [ + command_row(PAUSE_ID, "pause", "2026-08-26T10:00:00Z"), + command_row(DELETE_ID, "delete", "2026-08-26T10:05:00Z"), + ] + + ack(client, PAUSE_ID, DELETE_ID) + + assert "delete_device" in fake_supabase.call_names + assert "update_device" not in fake_supabase.call_names + + +def test_failed_delete_becomes_503(client, fake_supabase): + """Cihaz satırı silinemezse agent 200 almamalı. + + 200 alsaydı yerel wipe'ı yapar ve anahtarını silerdi; kayıt sunucuda öksüz + kalır, bir daha kimse onu silmeye zorlayamazdı (force remove hariç). + """ + fake_supabase.applied = [command_row(DELETE_ID, "delete")] + fake_supabase.delete_error = SupabaseError("kesinti") + + assert ack(client, DELETE_ID).status_code == 503 + + +# --- Ack ile veri yazımının sırası ----------------------------------------- + + +def test_rows_are_written_before_the_delete_ack(client, fake_supabase): + """Aynı gövdedeki son ölçümler, cihaz satırı silinmeden önce yazılmalı. + + Ters sırada foreign key ihlali olurdu: `metrics.device_id` artık var olmayan + bir satırı gösterirdi. Ve tam olarak bu veri en değerlisidir — çöküşten + hemen önceki son kayıtlar. + """ + fake_supabase.applied = [command_row(DELETE_ID, "delete")] + + ack( + client, + DELETE_ID, + logs=[ + { + "uuid": "bbbbbbbb-0000-0000-0000-000000000001", + "measured_at": "2026-08-26T10:00:00Z", + "level": "critical", + "message": "son soz", + } + ], + ) + + order = fake_supabase.call_names + assert order.index("insert_rows") < order.index("delete_device") diff --git a/tests/test_entrypoint.py b/tests/test_entrypoint.py new file mode 100644 index 0000000..09908d7 --- /dev/null +++ b/tests/test_entrypoint.py @@ -0,0 +1,78 @@ +""" +agent/__main__.py — açılış kararları. + +Giriş noktası iş yapmaz, yalnızca "çalışılsın mı, çalışılmasın mı"ya karar +verir. Buradaki tek soru silinmiş cihazla ilgili: `delete` komutu uygulanmış +bir makinede agent bir daha AÇILMAMALIDIR. + +Açılırsa görünürde bir hata olmaz — ve sorun tam olarak budur: cihaz kaydı +sunucudan silindiği için anahtar artık geçersizdir, süreç her poll'da 401 alan +bir hayalete dönüşür ve journald'ı sonsuza kadar hata satırlarıyla doldurur. +""" + +from __future__ import annotations + +import agent.__main__ as entrypoint +from agent.core.state import StateStore + +CONFIG_BODY = """ +collector_url = "https://collector.example" +device_key = "tbx_live_test" +""" + + +def prepare(tmp_path, monkeypatch): + """Geçerli bir config ve boş bir state dizini kurar.""" + config_path = tmp_path / "config.toml" + config_path.write_text(CONFIG_BODY) + config_path.chmod(0o600) + + state_dir = tmp_path / "state" + state_dir.mkdir() + + monkeypatch.setenv("TRACEBOX_CONFIG", str(config_path)) + monkeypatch.setenv("TRACEBOX_STATE_DIR", str(state_dir)) + return state_dir + + +def test_a_deleted_device_does_not_start_the_loop(tmp_path, monkeypatch, capsys): + """İşaret dosyası duruyorsa döngüye hiç girilmez.""" + state_dir = prepare(tmp_path, monkeypatch) + StateStore(state_dir).mark_deleted() + + started = [] + monkeypatch.setattr(entrypoint.loop, "run", lambda *args: started.append(args)) + + exit_code = entrypoint.main([]) + + assert started == [], "silinmiş cihazda döngü başlatıldı" + assert "silindi" in capsys.readouterr().out + + +def test_a_deleted_device_exits_without_asking_for_a_restart(tmp_path, monkeypatch, capsys): + """Çıkış kodu SIFIR olmalı. + + Servis `Restart=on-failure` ile çalışır: sıfırdan farklı her kod systemd'yi + yeniden başlatmaya çağırır. Silinmiş cihazda bu, saniyede bir açılıp kapanan + bir döngü demektir — ta ki çökme freni servisi `failed` bırakana kadar. + """ + state_dir = prepare(tmp_path, monkeypatch) + StateStore(state_dir).mark_deleted() + monkeypatch.setattr(entrypoint.loop, "run", lambda *args: None) + + assert entrypoint.main([]) == entrypoint.EXIT_OK + + +def test_a_normal_device_still_starts(tmp_path, monkeypatch, capsys): + """İşaret yoksa açılış olağan yolundan devam eder. + + Kontrolün fazla geniş yazılması (örneğin state dizininin varlığına bakmak) + her agent'ı kalıcı olarak durdururdu; bu test o kaymayı tutar. + """ + prepare(tmp_path, monkeypatch) + + started = [] + monkeypatch.setattr(entrypoint.loop, "run", lambda *args: started.append(args)) + + assert entrypoint.main([]) == entrypoint.EXIT_OK + assert len(started) == 1 diff --git a/tests/test_install_scripts.py b/tests/test_install_scripts.py index 50233c3..69a1893 100644 --- a/tests/test_install_scripts.py +++ b/tests/test_install_scripts.py @@ -1,7 +1,11 @@ """ -agent/install.sh · agent/uninstall.sh · agent/tracebox-agent.service — sözleşme testleri. +Kurulum betikleri ve systemd birimleri — sözleşme testleri. -Bu üç dosya bir birim testinde ÇALIŞTIRILAMAZ: kullanıcı oluşturur, sistem +Kapsanan dosyalar: agent/install.sh · agent/uninstall.sh · +agent/tracebox-agent.service · agent/tracebox-uninstall.path · +agent/tracebox-uninstall.service. + +Bu dosyalar bir birim testinde ÇALIŞTIRILAMAZ: kullanıcı oluşturur, sistem dizinlerine yazar, systemd'ye dokunur. Doğru çalıştıkları tek seferlik olarak atılabilir bir Docker konteynerinde elle doğrulandı. @@ -20,12 +24,21 @@ import pytest +from agent.core.state import DEFAULT_STATE_DIR, DELETED_FILENAME + AGENT_DIR = Path(__file__).resolve().parent.parent / "agent" INSTALL = AGENT_DIR / "install.sh" UNINSTALL = AGENT_DIR / "uninstall.sh" UNIT = AGENT_DIR / "tracebox-agent.service" +# `delete` komutunun root tarafı: agent yetkisiz çalıştığı için kendi kurulumunu +# kaldıramaz, yalnızca state dizinine bir işaret bırakır. Path unit onu görür, +# service unit kaldırma betiğini root olarak çalıştırır. +PATH_UNIT = AGENT_DIR / "tracebox-uninstall.path" +UNINSTALL_UNIT = AGENT_DIR / "tracebox-uninstall.service" + +UNITS = [UNIT, PATH_UNIT, UNINSTALL_UNIT] SCRIPTS = [INSTALL, UNINSTALL] # Yalnızca geliştirme için var olan override'lar. Üretim yolunu ezerler; kurulum @@ -62,10 +75,20 @@ def unit() -> str: return read(UNIT) +@pytest.fixture(scope="module") +def path_unit() -> str: + return read(PATH_UNIT) + + +@pytest.fixture(scope="module") +def uninstall_unit() -> str: + return read(UNINSTALL_UNIT) + + # --- Dosyaların kendisi ---------------------------------------------------- -@pytest.mark.parametrize("path", SCRIPTS + [UNIT], ids=lambda p: p.name) +@pytest.mark.parametrize("path", SCRIPTS + UNITS, ids=lambda p: p.name) def test_file_exists(path): """install.sh kurulum akışının TEK giriş noktası; eksikse akış kopar.""" assert path.is_file() @@ -109,7 +132,7 @@ def test_script_stops_on_the_first_error(path): # --- Kilitli kural: geliştirme override'ları sızmamalı --------------------- -@pytest.mark.parametrize("path", SCRIPTS + [UNIT], ids=lambda p: p.name) +@pytest.mark.parametrize("path", SCRIPTS + UNITS, ids=lambda p: p.name) @pytest.mark.parametrize("variable", DEVELOPMENT_OVERRIDES) def test_development_override_does_not_leak_into_production(path, variable): """Bu dosyadaki en kritik test. @@ -286,3 +309,147 @@ def test_uninstall_stops_the_service_before_deleting_files(uninstall_sh): def test_uninstall_removes_the_service_user(uninstall_sh): """Geride yetim bir sistem kullanıcısı bırakılmamalı.""" assert "userdel" in uninstall_sh + + +# --- delete komutunun root tarafı ------------------------------------------ +# Agent yetkisiz çalışır: /opt, /etc ve systemd ona kapalı. Bu yüzden `delete` +# komutunu uygularken yalnızca state dizinine bir işaret dosyası bırakır; +# kaldırmayı root tarafında bekleyen path unit üstlenir. Zincir üç halkadan +# oluşur (agent → işaret → path unit → uninstall.sh) ve her halka ayrı bir +# dosyada durduğu için sessizce kopabilir. + + +def test_the_watcher_looks_at_the_path_the_agent_actually_writes(path_unit): + """Bu dosyadaki en kritik testlerden biri. + + İzlenen yol ile agent'ın yazdığı yol AYNI olmalı. Biri değişip diğeri + kalırsa hiçbir hata görünmez: agent işaretini bırakır, path unit başka bir + dosyayı bekler ve kaldırma hiç başlamaz. Cihaz sunucudan silinmiş olur ama + servis makinede çalışmaya devam eder. + """ + expected = DEFAULT_STATE_DIR / DELETED_FILENAME + + assert re.search(rf"^PathExists={re.escape(str(expected))}$", path_unit, re.MULTILINE), ( + f"path unit {expected} yolunu izlemiyor" + ) + + +def test_the_watcher_catches_a_marker_that_is_already_there(path_unit): + """`PathExists`, `PathChanged` değil. + + PathChanged yalnızca dosya OLUŞTUĞU anda tetiklenir. Agent işaretini + bıraktıktan hemen sonra makine kapanırsa (ya da o an path unit çalışmıyorsa) + tetikleme kaçar ve kaldırma bir daha hiç yapılmaz. PathExists zaten duran + dosyayı da görür: kaldırma en geç bir sonraki açılışta tamamlanır. + """ + assert re.search(r"^PathExists=", path_unit, re.MULTILINE) + assert not re.search(r"^PathChanged=", path_unit, re.MULTILINE) + + +def test_the_watcher_triggers_the_uninstall_unit(path_unit): + """Zincirin ikinci halkası: path unit hangi birimi çalıştıracak?""" + assert re.search(rf"^Unit={re.escape(UNINSTALL_UNIT.name)}$", path_unit, re.MULTILINE) + + +def test_the_watcher_starts_at_boot(path_unit): + """[Install] bölümü olmayan bir path unit `enable` EDİLEMEZ. + + Edilemezse yeniden başlatmadan sonra izleyici hiç ayağa kalkmaz ve o + aradaki bir `delete` komutu sonsuza kadar yarım kalır. + """ + assert re.search(r"^WantedBy=multi-user\.target$", path_unit, re.MULTILINE) + + +def test_the_uninstall_unit_runs_the_script_without_asking(uninstall_unit): + """Kaldırma servisi soru soramaz: karşısında terminal yok. + + `--yes` düşerse betik onay bekler, oneshot birim yanıtsız takılır ve + kaldırma tamamlanmaz. + """ + exec_start = re.search(r"^ExecStart=(.+)$", uninstall_unit, re.MULTILINE) + + assert exec_start, "ExecStart yok" + assert exec_start.group(1) == "/opt/tracebox/uninstall.sh --yes" + + +def test_the_uninstall_unit_runs_as_root(uninstall_unit): + """Bu birim BİLEREK root çalışır — istisna burada, agent'ta değil. + + `User=tracebox` eklenirse kaldırma sessizce başarısız olur: systemctl, + /etc ve /opt o kullanıcıya kapalıdır. Testin işi, istisnanın yalnızca bu + dosyada kaldığını ve agent'ın kendi biriminin yetkisiz kalmaya devam + ettiğini (test_service_does_not_run_as_root) ayrı ayrı tutmaktır. + """ + assert not re.search(r"^User=", uninstall_unit, re.MULTILINE) + + +def test_the_uninstall_unit_is_separate_from_the_agent_unit(unit): + """Kaldırma agent'ın kendi biriminden çalıştırılamaz. + + Çalıştırılsaydı betiğin ilk işi olan `systemctl disable --now + tracebox-agent.service` betiği kendi cgroup'uyla birlikte öldürür ve + kaldırma daha ilk adımda yarıda kalırdı. + """ + assert "uninstall.sh" not in code(unit) + + +def test_install_puts_both_uninstall_units_in_place(install_sh): + """Kurulum bu iki dosyayı /etc/systemd/system'e koymazsa zincir hiç kurulmaz.""" + body = code(install_sh) + + for variable, unit_file in ( + ("UNINSTALL_SERVICE", UNINSTALL_UNIT), + ("UNINSTALL_PATH_UNIT", PATH_UNIT), + ): + # Betiğin adlandırdığı dosya, repodaki dosyanın kendisi olmalı: birim + # yeniden adlandırılıp betik güncellenmezse kurulum var olmayan bir + # dosyayı kopyalamaya çalışır. + assert f'{variable}="{unit_file.name}"' in body, f"{unit_file.name} adı betikte yok" + assert f'{variable}_PATH="/etc/systemd/system/${{{variable}}}"' in body + assert re.search( + rf'install -m 644 "\$\{{INSTALL_DIR\}}/agent/\$\{{{variable}\}}" ' + rf'"\$\{{{variable}_PATH\}}"', + body, + ), f"{unit_file.name} kurulmuyor" + + assert 'enable --now "${UNINSTALL_PATH_UNIT}"' in body, "path unit etkinleştirilmiyor" + + +def test_install_clears_a_leftover_marker_before_starting_the_watcher(install_sh): + """İkinci en kritik test. + + Önceki kurulum `delete` ile bittiyse geride kaldırma işareti kalmış olabilir. + İzleyici o dosya dururken etkinleştirilirse yeni kurulum, daha ilk saniyede + kendini kaldırır — kullanıcının anahtarı da beraberinde gider. + """ + body = code(install_sh) + remove_index = body.index('rm -f "${DELETED_MARKER}"') + enable_index = body.index('enable --now "${UNINSTALL_PATH_UNIT}"') + + assert remove_index < enable_index, "izleyici, eski işaret silinmeden başlatılıyor" + + +def test_uninstall_removes_both_uninstall_units(uninstall_sh): + """Geride yetim unit dosyası kalmamalı; sonraki kurulum onları yeniler.""" + body = code(uninstall_sh) + + assert "${UNINSTALL_PATH_UNIT_PATH}" in body + assert "${UNINSTALL_SERVICE_PATH}" in body + + +def test_uninstall_does_not_stop_the_unit_that_is_running_it(uninstall_sh): + """En sinsi hata burada olurdu. + + Betik çoğu zaman tracebox-uninstall.service tarafından çalıştırılır. O birimi + durdurmak (`stop` ya da `disable --now`) kendi süreç ağacını öldürmek + demektir: kaldırma tam ortasında kesilir, makinede yarım silinmiş bir + kurulum kalır. İzleyici (path unit) durdurulabilir — betiği o çalıştırmıyor. + """ + body = code(uninstall_sh) + + assert not re.search(r"systemctl (stop|disable --now) .*UNINSTALL_SERVICE", body), ( + "betik kendini çalıştıran birimi durduruyor" + ) + assert re.search(r'systemctl disable --now "\$\{UNINSTALL_PATH_UNIT\}"', body), ( + "izleyici kapatılmıyor" + ) diff --git a/tests/test_loop_commands.py b/tests/test_loop_commands.py new file mode 100644 index 0000000..6e5af68 --- /dev/null +++ b/tests/test_loop_commands.py @@ -0,0 +1,255 @@ +""" +agent/core/loop.py — komut turunun döngüye bağlanışı. + +Uygulama mantığı commands.py'de sınandı; buradaki soru farklı: **ack edilen +id'ler state'te ne zaman durur, ne zaman düşer?** + +Bu, gözle görülmeyen bir muhasebedir. Yanlış tarafa kayarsa iki sessiz bozukluk +üretir: id erken düşerse komut sunucuda sonsuza kadar `pending` kalır (dashboard +"duraklatıldı" demez); geç düşerse agent aynı ack'i her gönderimde tekrar +tekrar yollar. İkisinde de hata mesajı yok. +""" + +from __future__ import annotations + +import pytest + +from agent.core import loop +from agent.core.commands import Command, CommandError +from agent.core.config import Config +from agent.core.shipper import SendResult +from agent.core.spool import Spool +from agent.core.state import StateStore + +CONFIG = Config(collector_url="https://collector.test", device_key="tbx_live_test") + + +class FakePoller: + """Sırayla tüketilen poll sonuçları; eleman bir liste ya da fırlatılacak hata.""" + + def __init__(self, *rounds) -> None: + self._rounds = list(rounds) + + def fetch(self, config): + outcome = self._rounds.pop(0) if self._rounds else [] + if isinstance(outcome, Exception): + raise outcome + return outcome + + +class FakeShipper: + """`send_acks` çağrılarını kaydeder; `ok` sonucu testten verilir.""" + + def __init__(self, ok: bool = True) -> None: + self._ok = ok + self.calls: list[list[str]] = [] + + def send_acks(self, config, command_ids: list[str]) -> SendResult: + self.calls.append(list(command_ids)) + return SendResult(ok=self._ok, detail="" if self._ok else "HTTP 500") + + +@pytest.fixture +def spool(tmp_path): + instance = Spool(tmp_path, max_age_days=10, max_size_mb=200) + yield instance + instance.close() + + +@pytest.fixture +def store(tmp_path): + return StateStore(tmp_path) + + +def poll(poller, state, store, spool, shipper) -> bool: + return loop._poll_commands(poller, CONFIG, state, store, spool, shipper) + + +# --- state muhasebesi ------------------------------------------------------ + + +def test_only_the_confirmed_ids_leave_the_state(store): + """"Listeyi boşalt" değil, "onaylananları çıkar". + + Bugün fark görünmez (tek döngü var), ama kural yanlış yazılırsa ileride + gönderim sırasında uygulanan bir komutun ack'i sessizce yutulur. + """ + state = store.load() + state.applied_command_ids = ["k1", "k2", "k3"] + + loop._prune_acked(state, store, ["k2"]) + + assert state.applied_command_ids == ["k1", "k3"] + assert store.load().applied_command_ids == ["k1", "k3"], "değişiklik diske yazılmadı" + + +def test_nothing_is_written_when_nothing_was_confirmed(store): + """Onay yoksa diske yazma da yok — her turda gereksiz fsync yapılmaz.""" + state = store.load() + state.applied_command_ids = ["k1"] + + loop._prune_acked(state, store, []) + + assert state.applied_command_ids == ["k1"] + assert not store.path.exists() + + +def test_an_acked_command_leaves_no_debt_behind(store, spool): + """Olağan akış: komut uygulanır, hemen ack'lenir ve state temiz kalır.""" + state = store.load() + shipper = FakeShipper() + + poll(FakePoller([Command(id="k1", type="pause")]), state, store, spool, shipper) + + assert state.logging_enabled is False + assert shipper.calls == [["k1"]] + assert state.applied_command_ids == [] + assert store.load().applied_command_ids == [] + + +def test_an_unconfirmed_ack_stays_in_the_state_for_the_next_send(store, spool): + """Ack ulaşmadıysa id state'te KALIR ve sonraki gönderime piggyback olur. + + Kalmazsa komut hiç bildirilmemiş olur: sunucu onu her poll'da yeniden verir, + `devices.logging_enabled` kopyası hiç güncellenmez. + """ + state = store.load() + + poll(FakePoller([Command(id="k1", type="pause")]), state, store, spool, FakeShipper(ok=False)) + + assert state.applied_command_ids == ["k1"] + assert store.load().applied_command_ids == ["k1"], "borç diske yazılmadı" + + +def test_a_redelivered_command_is_acked_again_without_being_listed_twice(store, spool): + """Bu dosyanın en kritik testi. + + Komutun tekrar gelmesi ack'in ulaşmadığı anlamına gelir; ack yeniden + denenmelidir. Ama id state'e İKİNCİ kez yazılmamalıdır — yazılsaydı liste + her turda büyür ve aynı id gönderimde defalarca tekrarlanırdı. + """ + state = store.load() + state.logging_enabled = False + state.applied_command_ids = ["k1"] + shipper = FakeShipper(ok=False) + + poll(FakePoller([Command(id="k1", type="pause")]), state, store, spool, shipper) + + assert shipper.calls == [["k1"]], "tekrar gelen komut yeniden ack edilmedi" + assert state.applied_command_ids == ["k1"], "aynı id listeye iki kez yazıldı" + + +def test_an_unknown_command_creates_no_debt(store, spool): + """Uygulanmayan komut ack edilmez, dolayısıyla state'e de girmez.""" + state = store.load() + shipper = FakeShipper() + + poll(FakePoller([Command(id="k9", type="reboot")]), state, store, spool, shipper) + + assert shipper.calls == [] + assert state.applied_command_ids == [] + + +# --- turun döngüye etkisi -------------------------------------------------- + + +def test_a_failed_poll_does_not_stop_the_agent(store, spool, capsys): + """Collector'a ulaşılamaması toplamayı ve gönderimi durdurmaz. + + Durdursaydı geçici bir ağ kesintisi agent'ı tamamen susturur, çöküş anına + ait veri hiç toplanmazdı. + """ + state = store.load() + + stopped = poll(FakePoller(CommandError("bağlanılamadı")), state, store, spool, FakeShipper()) + + assert stopped is False + assert "komutlar alınamadı" in capsys.readouterr().out + + +def test_an_empty_queue_costs_nothing(store, spool): + """Çoğu poll boş döner: ne ack atılır ne diske yazılır.""" + state = store.load() + shipper = FakeShipper() + + assert poll(FakePoller([]), state, store, spool, shipper) is False + assert shipper.calls == [] + assert not store.path.exists() + + +def test_delete_tells_the_loop_to_stop(store, spool): + """`delete` uygulandıktan sonra döngü devam etmemeli. + + Etseydi agent, kaydı sunucudan silinmiş bir cihaz için ölçüm toplamaya ve + 401 alan istekler atmaya devam ederdi. + """ + state = store.load() + + assert poll(FakePoller([Command(id="sil", type="delete")]), state, store, spool, FakeShipper()) + + +def test_a_postponed_delete_lets_the_loop_continue(store, spool): + """Ack gitmediyse silme olmamıştır; agent çalışmaya devam eder ve tekrar dener.""" + state = store.load() + + stopped = poll( + FakePoller([Command(id="sil", type="delete")]), state, store, spool, FakeShipper(ok=False) + ) + + assert stopped is False + assert store.is_deleted() is False + + +# --- normal gönderime binen ack'ler ---------------------------------------- + + +class FakeSender: + """`send_pending` sonucunu testin verdiği şekilde döndüren sahte shipper.""" + + def __init__(self, result: SendResult) -> None: + self._result = result + self.backoff_seconds = 0.0 + self.seen: list[list[str]] = [] + + def send_pending(self, config, applied_command_ids: list[str]) -> SendResult: + self.seen.append(list(applied_command_ids)) + return self._result + + +def test_acks_that_ride_along_with_a_normal_send_leave_the_state(store, spool): + """Ack'in ikinci yolu: `POST /ingest` gövdesine binmek (CLAUDE.md §4.2). + + Ack'ler oraya da bindiği için düşme kuralı iki yerde birden geçerlidir; + yalnızca poll tarafında uygulanırsa id'ler sonsuza kadar her gövdede + tekrarlanır. + """ + state = store.load() + state.applied_command_ids = ["k1", "k2"] + + loop._send_spool( + FakeSender(SendResult(ok=True, sent=2, acked=["k1"])), CONFIG, state, store, spool + ) + + assert state.applied_command_ids == ["k2"] + + +def test_acks_confirmed_before_a_failure_still_leave_the_state(store, spool): + """Tur yarıda kalsa bile ilk isteğin ack'i onaylanmıştır. + + Ack'ler yalnızca ilk isteğe biner; o istek 200 aldıysa komutlar sunucuda + `applied` olmuştur. Sonraki batch'in hatası bunu geri almaz — id'ler yine + düşer, yoksa bir daha hiç bildirilmeyecek komutlar için sonsuza kadar + taşınırlar. + """ + state = store.load() + state.applied_command_ids = ["k1"] + + loop._send_spool( + FakeSender(SendResult(ok=False, sent=2, detail="HTTP 500", acked=["k1"])), + CONFIG, + state, + store, + spool, + ) + + assert state.applied_command_ids == [] diff --git a/tests/test_shipper.py b/tests/test_shipper.py index cf59eec..4f3c1a8 100644 --- a/tests/test_shipper.py +++ b/tests/test_shipper.py @@ -232,3 +232,95 @@ def test_inventory_is_confirmed_only_by_a_200(spool): assert shipper.send_inventory(CONFIG, {"os_name": "Ubuntu"}).ok is False assert shipper.send_inventory(CONFIG, {"os_name": "Ubuntu"}).ok is True + + +# --- Komut ack'leri -------------------------------------------------------- + + +def test_confirmed_acks_are_reported_back(spool): + """200 alınan ack'ler sonuçta bildirilir; çağıran onları state'ten düşer. + + Bildirilmezlerse id'ler state'te kalır ve her gönderimde tekrar tekrar + gönderilir — sunucu onları çoktan `applied` yapmışken. + """ + fill(spool, 1) + shipper = make_shipper(spool, FakeCollector(200)) + + assert shipper.send_pending(CONFIG, ["komut-1"]).acked == ["komut-1"] + + +def test_acks_are_not_reported_when_the_request_fails(spool): + """İstek başarısızsa ack onaylanmış sayılmaz. + + Sayılsaydı id state'ten düşer, sunucuya hiç ulaşmaz ve komut sonsuza kadar + `pending` kalırdı: her poll'da yeniden gelir, agent'ın uyguladığı hiç + bilinmezdi. + """ + fill(spool, 1) + shipper = make_shipper(spool, FakeCollector(500)) + + result = shipper.send_pending(CONFIG, ["komut-1"]) + + assert result.ok is False + assert result.acked == [] + + +def test_acks_confirmed_by_the_first_batch_survive_a_later_failure(spool, monkeypatch): + """Tur yarıda kalsa bile ilk isteğin ack'i onaylanmıştır. + + Ack'ler yalnızca ilk isteğe biner; o istek 200 aldıysa komutlar sunucuda + `applied` olmuştur. Sonraki batch'in hatası bu gerçeği geri almaz. + """ + monkeypatch.setattr("agent.core.shipper.BATCH_ROWS", 2) + fill(spool, 5) + shipper = make_shipper(spool, FakeCollector(200, 500)) + + result = shipper.send_pending(CONFIG, ["komut-1"]) + + assert result.ok is False + assert result.acked == ["komut-1"] + + +def test_ack_only_request_carries_no_measurements(spool): + """send_acks bir KONTROL mesajıdır: gövdesinde tek ölçüm satırı yoktur. + + Spool dolu olsa bile ona dokunulmaz — bu istek pause sırasında da atılır ve + telemetri taşısaydı pause'un anlamını çiğnerdi. + """ + fill(spool, 3) + collector = FakeCollector(200) + shipper = make_shipper(spool, collector) + + result = shipper.send_acks(CONFIG, ["komut-1", "komut-2"]) + + assert result.ok is True + assert result.acked == ["komut-1", "komut-2"] + assert collector.bodies == [ + { + "metrics": [], + "logs": [], + "crash_snapshots": [], + "applied_command_ids": ["komut-1", "komut-2"], + } + ] + # Spool'a dokunulmadı: kayıtlar normal gönderimi bekliyor. + assert spool.count() == 3 + + +def test_ack_only_request_is_skipped_when_there_is_nothing_to_ack(spool): + """Ack yoksa istek de yok — boş gövdenin karşılığı yok.""" + collector = FakeCollector() + shipper = make_shipper(spool, collector) + + assert shipper.send_acks(CONFIG, []).ok is True + assert collector.requests == [] + + +def test_failed_ack_only_request_confirms_nothing(spool): + """Ack isteği başarısızsa hiçbir id onaylanmaz; sonraki gönderime kalırlar.""" + shipper = make_shipper(spool, FakeCollector(500)) + + result = shipper.send_acks(CONFIG, ["komut-1"]) + + assert result.ok is False + assert result.acked == [] diff --git a/tests/test_spool.py b/tests/test_spool.py index 85ab4d2..712a066 100644 --- a/tests/test_spool.py +++ b/tests/test_spool.py @@ -16,6 +16,8 @@ sayılmaz ve halka tampon (ring buffer) sessizce çalışmayı bırakır. """ +import sqlite3 + import pytest from agent.core import spool as spool_module @@ -146,3 +148,47 @@ def test_size_limit_drops_the_oldest_records(tmp_path): assert spool.size_bytes() <= 1 * 1024 * 1024 finally: spool.close() + + +def test_wipe_leaves_no_trace_of_the_data(tmp_path): + """`delete` komutunun yerel temizliği: dosya da yan dosyaları da gitmeli. + + Tabloyu boşaltmak yetmez ve dosyayı silmek de tek başına yetmez: WAL + modunda yazılan satırlar `-wal` yan dosyasında durur. Temiz bir kapanışta + SQLite onu kendi toplar — ama son bağlantı kapanmadıysa (ikinci bir tanıtıcı, + çökmüş bir süreç) yan dosya olduğu gibi kalır ve içinde ölçümler vardır. + Kullanıcı cihazı sildiğinde o veri makinede BULUNMAMALIDIR; bu yüzden silme + SQLite'ın nezaketine bırakılmaz. + + Test o durumu bilerek kurar: ikinci bir bağlantı açık tutulur, böylece + kapanış WAL'ı temizleyemez ve geriye yalnızca `wipe`ın kendi silmesi kalır. + """ + spool = Spool(tmp_path, max_age_days=10, max_size_mb=200) + for index in range(5): + spool.add(RECORD_METRIC, make_record(f"kayit-{index}", message=PAYLOAD_FILLER)) + + sidecar = sqlite3.connect(spool.path) + try: + # Bağlantı gerçekten AÇILMALI: sqlite3.connect tembeldir, dosyaya ilk + # sorguda dokunur. Sorgu atılmazsa spool'unki son bağlantı sayılır, + # kapanışta WAL'ı toplar ve testin kurduğu durum hiç oluşmaz. + sidecar.execute("select count(*) from pending").fetchone() + + wal = spool.path.with_name(spool.path.name + "-wal") + assert wal.exists(), "WAL yan dosyası hiç oluşmamış — test kurgusu geçersiz" + + spool.wipe() + + for suffix in ("", "-wal", "-shm"): + leftover = spool.path.with_name(spool.path.name + suffix) + assert not leftover.exists(), f"geride kaldı: {leftover}" + + # Silmeden sonra spool ÖLÜDÜR. Bağlantı açık bırakılsaydı sıradaki + # yazma dosyayı yeniden yaratır ve silinmiş cihazda yeni ölçümler + # birikmeye başlardı — hem de kimse fark etmeden. + with pytest.raises(sqlite3.ProgrammingError): + spool.add(RECORD_METRIC, make_record("silme-sonrasi")) + + assert not spool.path.exists() + finally: + sidecar.close()