diff --git a/agent/__init__.py b/agent/__init__.py index eabc0e2..bc2776c 100644 --- a/agent/__init__.py +++ b/agent/__init__.py @@ -31,4 +31,8 @@ # 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" +# +# 0.3.0: acil gönderim (flush) devrede — eşik aşıldığında veri 30 saniyelik +# turu beklemiyor. Bir cihazın verisi "neden tam çöküş anında var / yok" +# sorusunun cevabı doğrudan bu sürüm sınırıdır. +__version__ = "0.3.0" diff --git a/agent/config.example.toml b/agent/config.example.toml index a4cfa74..78194ac 100644 --- a/agent/config.example.toml +++ b/agent/config.example.toml @@ -59,6 +59,9 @@ spool_max_size_mb = 200 # "load_avg" yük ortalaması (Linux) -> metrics.load_avg_1/5/15 # "gpu" GPU kullanımı + VRAM -> metrics.gpu_usage_percent, gpu_vram_used_mb # "external_ip" dış IP (statik) -> devices.external_ip +# Değeri agent GÖNDERMEZ: isteği alan taraf (collector) +# bağlantının kaynak IP'sinden yazar. Buradaki tercih +# yalnızca "yazılsın mı" sorusunu cevaplar. # "crash_processes" flush anında en çok kaynak yiyen 5 süreç -> crash_snapshots # # Örnek: enabled_addons = ["swap", "load_avg"] diff --git a/agent/core/clock.py b/agent/core/clock.py index e47145e..4606bac 100644 --- a/agent/core/clock.py +++ b/agent/core/clock.py @@ -31,3 +31,30 @@ def epoch_to_utc_iso(epoch_seconds: float) -> str: hesaba katılmaz, değer doğrudan UTC olarak yorumlanır. """ return datetime.fromtimestamp(epoch_seconds, tz=timezone.utc).isoformat(timespec=_TIMESPEC) + + +def seconds_since_iso(timestamp: str | None) -> float | None: + """Verilen ISO 8601 damgasından bu yana geçen saniyeyi döndürür. + + İki durumda None döner ve çağıran bunu "ölçülemedi" olarak yorumlar: + * timestamp None ya da boş — karşılaştırılacak bir an yok, + * metin ISO 8601 olarak çözülemiyor — state.json elle düzenlenmiş olabilir. + + Sonuç NEGATİF de çıkabilir: damga gelecekte kalmışsa (sistem saati geri + alınmış) fark eksi olur. Çağıran bu durumu kendi kuralına göre yorumlar; + burada gizlenmez, çünkü "geçen süre" sorusunun dürüst cevabı budur. + """ + if not timestamp: + return None + + try: + moment = datetime.fromisoformat(timestamp) + except ValueError: + return None + + # Saat dilimi taşımayan bir damga UTC kabul edilir: bu modülün ürettiği her + # damga zaten UTC'dir, dilimsiz bir değer ancak dışarıdan gelmiş olabilir. + if moment.tzinfo is None: + moment = moment.replace(tzinfo=timezone.utc) + + return (datetime.now(timezone.utc) - moment).total_seconds() diff --git a/agent/core/config.py b/agent/core/config.py index bd6b1bb..f2a703f 100644 --- a/agent/core/config.py +++ b/agent/core/config.py @@ -29,6 +29,25 @@ stat.S_IRGRP | stat.S_IWGRP | stat.S_IXGRP | stat.S_IROTH | stat.S_IWOTH | stat.S_IXOTH ) +# Seçilebilir eklentilerin adları. Tek tanım yeri BURASI: eklentiyi toplayan +# modüller (metrics, inventory, flush) bu sabitleri import eder, böylece +# config.toml'daki metinle kodun beklediği metin ayrışamaz. +ADDON_TEMPERATURE = "temperature" +ADDON_SWAP = "swap" +ADDON_LOAD_AVG = "load_avg" +ADDON_GPU = "gpu" +ADDON_EXTERNAL_IP = "external_ip" +ADDON_CRASH_PROCESSES = "crash_processes" + +KNOWN_ADDONS = ( + ADDON_TEMPERATURE, + ADDON_SWAP, + ADDON_LOAD_AVG, + ADDON_GPU, + ADDON_EXTERNAL_IP, + ADDON_CRASH_PROCESSES, +) + # Config'de bulunması ZORUNLU alanlar. Eksikse agent açılışta durur; varsayılan # uydurmak, yanlış adrese veri göndermeye çalışan bir agent üretirdi. REQUIRED_KEYS = ("collector_url", "device_key") @@ -141,6 +160,17 @@ def _parse(raw: dict, *, warn) -> Config: if not isinstance(addons, list) or not all(isinstance(a, str) for a in addons): raise ConfigError("'enabled_addons' string listesi olmalı") + # Tanınmayan ad HATA DEĞİL, uyarıdır: yazım hatası yüzünden agent'ı + # başlatmamak, bir eklentinin toplanmamasından daha ağır bir sonuç olurdu. + # Ama sessiz de kalınmaz — "temprature" yazan kullanıcı, sıcaklık sütunu + # neden hep null diye günlerce bakabilirdi. + unknown = [name for name in addons if name not in KNOWN_ADDONS] + if unknown: + warn( + f"enabled_addons içinde tanınmayan ad: {', '.join(unknown)} — " + f"yok sayılıyor. Geçerli değerler: {', '.join(KNOWN_ADDONS)}" + ) + return Config( collector_url=raw["collector_url"].strip().rstrip("/"), device_key=raw["device_key"].strip(), diff --git a/agent/core/flush.py b/agent/core/flush.py new file mode 100644 index 0000000..8c6e4f1 --- /dev/null +++ b/agent/core/flush.py @@ -0,0 +1,171 @@ +""" +Acil gönderim — eşik aşıldığında 30 saniyelik gönderim turu beklenmez. + +Modül iki soruyu cevaplar ve başka hiçbir şey yapmaz: + * "şu an bir eşik aşıldı mı, aşıldıysa hangisi?" — evaluate() + * "aşıldı ama çok yakın zamanda flush ettik mi?" — cooldown_active() + +Gönderimin kendisi (spool'u boşaltmak, shipper'ı çağırmak, last_flush_at'i +yazmak) döngünün işidir; burada karar üretilir, yan etki üretilmez. Tek istisna +build_crash_snapshot(): süreç listesini okumak için psutil'e dokunur, ama o da +yalnızca okur. + +CLAUDE.md §7 — eşikler cpu>90 / ram>90 / disk>95 ve error|critical seviyeli log. +""" + +from __future__ import annotations + +import time +import uuid + +import psutil + +from agent.core.clock import seconds_since_iso, utc_now_iso +from agent.core.config import ADDON_CRASH_PROCESSES, Config +from agent.core.metrics import BYTES_PER_MB, MetricSample + +# crash_snapshots.trigger_reason'ın alabileceği dört değer. Şemadaki check +# kısıtı ve collector'daki Literal ile birebir aynı olmak zorunda. +REASON_LOG = "log" +REASON_RAM = "ram" +REASON_CPU = "cpu" +REASON_DISK = "disk" + +# Aynı anda birden fazla eşik aşılabilir ama sütun tek değer alır. Sıralama +# yukarıdan aşağıya denenir ve ilk tutan yazılır: log en üstte, çünkü diğer üçü +# "yük yüksek" derken log "bir şey bozuldu" der. +REASON_ORDER = (REASON_LOG, REASON_RAM, REASON_CPU, REASON_DISK) + +# Snapshot'a kaç süreç girer. +TOP_PROCESS_COUNT = 5 + +# Süreç başına CPU yüzdesi iki okuma arasındaki farktan hesaplanır; ilk okuma +# her süreç için 0.0 döner. Aradaki bu kısa bekleme olmadan liste tamamen +# sıfırlardan oluşur ve sıralama anlamsızlaşır. +PROCESS_SAMPLE_SECONDS = 0.1 + + +def _above(value: float | None, threshold: int) -> bool: + """Ölçüm eşiği aştı mı. Ölçülemeyen alan (None) eşiği aşmış sayılmaz.""" + return value is not None and value > threshold + + +def evaluate( + *, + sample: MetricSample, + ram_percent: float | None, + urgent_log_count: int, + config: Config, +) -> str | None: + """Eşik aşıldıysa trigger_reason'ı, aşılmadıysa None döndürür. + + ram_percent ölçümün yanında AYRICA taşınır: MetricSample doğrudan wire + gövdesi olarak gidiyor ve collector sözleşme dışı alanı 422 ile reddediyor, + yani yüzde o nesneye eklenemez. Yine de aynı ölçüm anına aittir — eşik + kararı ile kaydedilen satır arasında zaman farkı olmaz. + + Sıra REASON_ORDER'dır; ilk tutan kazanır. + """ + if urgent_log_count > 0: + return REASON_LOG + if _above(ram_percent, config.flush_ram_threshold): + return REASON_RAM + if _above(sample.cpu_percent, config.flush_cpu_threshold): + return REASON_CPU + if _above(sample.disk_percent, config.flush_disk_threshold): + return REASON_DISK + return None + + +def cooldown_active(last_flush_at: str | None, cooldown_seconds: int) -> bool: + """Son flush'ın üzerinden cooldown süresi geçmediyse True. + + Ölçüm duvar saatiyle yapılır (monotonic ile değil): monotonic agent her + yeniden başladığında sıfırlanır, o anda cooldown da sıfırlanırdı. + + İki durumda False döner, yani flush'a izin verilir: + * damga yok ya da çözülemedi — daha önce hiç flush edilmemiş kabul edilir, + * geçen süre NEGATİF — damga gelecekte kalmış, yani sistem saati geri + alınmış. Bu durumda cooldown'ı açık saymak flush'ı süresiz kilitlerdi; + izin verildiğinde damga yeniden yazılır ve hesap kendiliğinden düzelir. + """ + elapsed = seconds_since_iso(last_flush_at) + if elapsed is None or elapsed < 0: + return False + return elapsed < cooldown_seconds + + +def _top_processes(reason: str, limit: int) -> list[dict]: + """En çok kaynak tüketen süreçleri {name, cpu, ram_mb} sözlükleri olarak verir. + + İki turlu okuma: ilk tur her sürecin CPU sayacına taban değeri koyar, + PROCESS_SAMPLE_SECONDS kadar beklenir, ikinci tur o tabana göre gerçek + yüzdeyi verir. + + Sıralama ölçütü tetikleyiciye göre değişir: RAM eşiği aşıldıysa belleğe, + diğer hallerde CPU'ya bakılır — "kaynak-yiyen" ifadesinin karşılığı, o an + tükenen kaynaktır. İkinci alan eşitlik bozucudur. + + Okuma sırasında ölen ya da izin vermeyen süreçler sessizce atlanır: snapshot + tam olmasa da alınır, çünkü alındığı an bir daha gelmez. + """ + for process in psutil.process_iter(): + try: + process.cpu_percent() + except (psutil.NoSuchProcess, psutil.AccessDenied): + continue + + time.sleep(PROCESS_SAMPLE_SECONDS) + + rows: list[dict] = [] + for process in psutil.process_iter(["name", "memory_info"]): + try: + cpu = process.cpu_percent() + info = process.info + except (psutil.NoSuchProcess, psutil.AccessDenied): + continue + + memory = info.get("memory_info") + rows.append( + { + # name boş dönebilir (çekirdek thread'leri); sütun metin + # bekliyor, boş metin yerine görünür bir işaret konur. + "name": info.get("name") or "?", + "cpu": round(cpu, 1), + "ram_mb": memory.rss // BYTES_PER_MB if memory else 0, + } + ) + + if reason == REASON_RAM: + rows.sort(key=lambda row: (row["ram_mb"], row["cpu"]), reverse=True) + else: + rows.sort(key=lambda row: (row["cpu"], row["ram_mb"]), reverse=True) + + return rows[:limit] + + +def build_crash_snapshot(reason: str, config: Config) -> dict: + """POST /ingest gövdesindeki crash_snapshots satırını üretir (CLAUDE.md §4.2). + + Satır HER flush'ta yazılır. crash_processes eklentisi kapalıysa processes + boş kalır ama trigger_reason ile measured_at yine kaydedilir; metrikler + "CPU %95'ti" der, bu satır "flush gerçekten attı" der. + + Süreçler okunamazsa (psutil beklenmedik bir hata verirse) snapshot boş + süreç listesiyle döner: eksik bir kayıt, hiç kayıt olmamasından iyidir. + """ + processes: list[dict] = [] + # Süreç listesi yalnızca eklenti açıkken doldurulur; kapalıyken satır + # yine yazılır ama processes boş kalır. + if ADDON_CRASH_PROCESSES in config.enabled_addons: + try: + processes = _top_processes(reason, TOP_PROCESS_COUNT) + except psutil.Error: + processes = [] + + return { + "uuid": str(uuid.uuid4()), + "measured_at": utc_now_iso(), + "trigger_reason": reason, + "processes": processes, + } diff --git a/agent/core/gpu.py b/agent/core/gpu.py new file mode 100644 index 0000000..89f72c8 --- /dev/null +++ b/agent/core/gpu.py @@ -0,0 +1,106 @@ +""" +GPU okuma — `nvidia-smi` çıktısının tek muhatabı. + +NEDEN AYRI MODÜL: GPU iki farklı yere veri verir — model adı envantere +(statik), kullanım ve VRAM metriklere (her ölçümde). İkisi de aynı harici +programa dayanır. Ayrı bir modül olmasaydı `nvidia-smi`nin komut satırı, +çıktı biçimi ve hata halleri iki dosyada birden tekrarlanırdı. + +NEDEN SUBPROCESS: pynvml gibi bir kütüphane daha temiz görünür ama kurulum +maliyeti getirir; `nvidia-smi` sürücüyle birlikte zaten gelir. Aynı gerekçe +journald okuması için de kullanıldı (systemd-python yerine `journalctl` +çağrısı) — bağımlılık eklemek yerine sistemde HAZIR olanı çağırmak. + +KAPSAM: yalnızca NVIDIA. AMD/Intel GPU'lar için değer null kalır; eklenti +açıksa bile yanlış bir sayı üretilmez. +""" + +from __future__ import annotations + +import subprocess + +# Sürücü kuruluysa PATH'te bulunur. Bulunamaması hata değildir: eklenti açık +# ama makinede NVIDIA GPU yok demektir. +NVIDIA_SMI = "nvidia-smi" + +# Ölçüm aralığı saniyelerle ifade ediliyor; yanıt vermeyen bir süreç döngüyü +# bekletmemeli. journald'ın 30 saniyelik payına karşılık burada süre kısa +# tutuldu, çünkü GPU okuması her ölçümde tekrarlanır. +QUERY_TIMEOUT_SECONDS = 2.0 + +# Çıktı biçimi: başlıksız CSV. `nounits` sayıların yanına birim yazılmasını +# engeller ("34 %" yerine "34"), yani ayrıştırma tek bir float() çağrısıdır. +_CSV_FLAGS = "--format=csv,noheader,nounits" + + +class GpuReader: + """nvidia-smi çağrılarını yapan okuyucu. + + Durum tutar çünkü tek bir şeyi hatırlaması gerekir: program sistemde HİÇ + yoksa bir daha denemeye gerek yok. Ölçüm aralığı saniyelerle ifade edildiği + için, olmayan bir programı her turda başlatmaya çalışmak sürekli ve boş bir + süreç yaratma maliyetidir. + + Diğer hatalar (zaman aşımı, sıfırdan farklı çıkış kodu) kalıcı sayılmaz: + sürücü geçici olarak meşgul olabilir, bir sonraki tur yeniden denenir. + """ + + def __init__(self) -> None: + self._missing = False + + def read_usage(self) -> tuple[float | None, int | None]: + """(kullanım yüzdesi, kullanılan VRAM MB) — okunamazsa (None, None). + + Çoklu GPU'da yalnızca İLK kart okunur: şemada tek bir gpu_usage_percent + sütunu var. Kart başına satır tutmak zaman serisinin şeklini + değiştirirdi; bugünkü soru "GPU yüklü müydü", "hangi kart" değil. + """ + line = self._query("utilization.gpu,memory.used") + if line is None: + return None, None + + parts = [part.strip() for part in line.split(",")] + if len(parts) != 2: + return None, None + + try: + return round(float(parts[0]), 1), int(float(parts[1])) + except ValueError: + # [N/A] gibi sayı olmayan değerler: sürücü o alanı raporlamıyor. + return None, None + + def read_model(self) -> str | None: + """GPU model adı (envanter için) — okunamazsa None.""" + line = self._query("name") + return line or None + + def _query(self, fields: str) -> str | None: + """Sorguyu çalıştırır ve çıktının İLK satırını döndürür. + + Her hata None'a indirgenir: GPU bir EKLENTİdir, onun yokluğu ölçüm + turunu düşürmemeli. Ayrım yalnızca "program yok" halinde yapılır, + çünkü tek kalıcı olan odur. + """ + if self._missing: + return None + + try: + result = subprocess.run( + [NVIDIA_SMI, f"--query-gpu={fields}", _CSV_FLAGS], + capture_output=True, + text=True, + errors="replace", + timeout=QUERY_TIMEOUT_SECONDS, + check=False, + ) + except FileNotFoundError: + self._missing = True + return None + except (subprocess.TimeoutExpired, OSError): + return None + + if result.returncode != 0: + return None + + first_line = result.stdout.strip().splitlines() + return first_line[0].strip() if first_line else None diff --git a/agent/core/inventory.py b/agent/core/inventory.py index 2777118..5162777 100644 --- a/agent/core/inventory.py +++ b/agent/core/inventory.py @@ -5,9 +5,14 @@ toplam RAM, disk, işletim sistemi sürümü. Bu yüzden ayrı bir uç noktaya gider (POST /inventory) ve devices satırının üzerine yazılır — zaman serisi değildir. -Statik eklentiler (gpu_model, external_ip) M7'de bu modüle eklenecek. - -M2 KAPSAMI: okuma ve karşılaştırma gerçek, gönderim yok. +STATİK EKLENTİLER — ikisi de "nadiren değişir" tanımına uyar ama farklı +yerlerden gelir: + * gpu_model — makinede okunur, eklenti açıksa doldurulur (aşağıda). + * external_ip — agent GÖNDERMEZ. Cihazın kendi dış IP'sini bildirmesi, + cihazın kendi kimliği hakkında sunucuya bilgi vermesi demekti; agent + yanlış bir IP yazabilirdi. Değeri, isteği gerçekten alan tarafta + (collector, proxy başlığından) yazılır. Kullanıcının açık/kapalı tercihi + yine buradan gider: enabled_addons listesi payload'da taşınıyor. """ from __future__ import annotations @@ -20,7 +25,8 @@ from agent import __version__ from agent.core.clock import epoch_to_utc_iso -from agent.core.config import Config +from agent.core.config import ADDON_GPU, Config +from agent.core.gpu import GpuReader from agent.core.metrics import BYTES_PER_MB, DISK_MOUNT_POINT # İşlemci modelinin okunduğu yer. platform.processor() Linux'ta çoğunlukla @@ -46,6 +52,9 @@ class Inventory: agent_version: str enabled_addons: list[str] + # Statik eklenti: gpu eklentisi kapalıyken None kalır. + gpu_model: str | None = None + def as_dict(self) -> dict: """Karşılaştırmaya ve gönderime uygun sözlük hali. @@ -77,6 +86,9 @@ def collect_inventory(config: Config) -> Inventory: last_boot=epoch_to_utc_iso(psutil.boot_time()), agent_version=__version__, enabled_addons=list(config.enabled_addons), + # Envanter açılışta bir kez okunur; GpuReader'ın burada durum + # taşımasına gerek yok, tek çağrılık bir nesne yeter. + gpu_model=GpuReader().read_model() if ADDON_GPU in config.enabled_addons else None, ) diff --git a/agent/core/loop.py b/agent/core/loop.py index fe00a0b..61d8486 100644 --- a/agent/core/loop.py +++ b/agent/core/loop.py @@ -7,9 +7,9 @@ Tick sabit 1 saniyedir ve config'den okunmaz; collect/send/poll aralıkları birbirinden bağımsız sayaçlardır. -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. +M7 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, eşik aşıldığında +gönderim turu beklenmeden acil flush yapılır. """ from __future__ import annotations @@ -23,16 +23,23 @@ from agent import __version__ from agent.core import commands as commands_module +from agent.core import flush as flush_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 +from agent.core.metrics import MetricReading, MetricSample, MetricsCollector from agent.core.shipper import Shipper -from agent.core.spool import RECORD_LOG, RECORD_METRIC, Spool +from agent.core.spool import RECORD_CRASH, RECORD_LOG, RECORD_METRIC, Spool from agent.core.state import State, StateStore -from agent.logsources.base import LEVELS, LogRecord, LogSource, LogSourceError +from agent.logsources.base import ( + LEVELS, + URGENT_LEVELS, + LogRecord, + LogSource, + LogSourceError, +) # Döngünün nabzı. Config'e AÇILMAZ: ölçüm sıklığı (insan sınırı) ile döngü ritmi # (sistem sabiti) ayrı kavramlardır. @@ -113,15 +120,19 @@ def _startup_inventory(config: Config, state: State) -> Inventory | None: return current -def _collect(collector: MetricsCollector, spool: Spool) -> None: - """Ölçüm alıp spool'a yazar. +def _collect(collector: MetricsCollector, spool: Spool, config: Config) -> MetricReading: + """Ölçüm alıp spool'a yazar ve okumayı geri döndürür. Pause'da da çalışır: pause yalnızca buluta göndermeyi durdurur, yerel kaydı değil. + + Spool'a yalnızca reading.sample yazılır; yanındaki ram_percent kaydedilmez, + eşik karşılaştırmasını yapacak olan çağırana verilir. """ - sample = collector.collect() - spool.add(RECORD_METRIC, asdict(sample)) - _log(f"[collect] {_format_sample(sample)}") + reading = collector.collect(config) + spool.add(RECORD_METRIC, asdict(reading.sample)) + _log(f"[collect] {_format_sample(reading.sample)}") + return reading def _log_payload(record: LogRecord) -> dict: @@ -151,9 +162,14 @@ def _level_summary(records: list[LogRecord]) -> str: return " ".join(f"{level}={counts[level]}" for level in LEVELS if counts[level]) -def _collect_logs(source: LogSource, spool: Spool, state: State, store: StateStore) -> None: +def _collect_logs(source: LogSource, spool: Spool, state: State, store: StateStore) -> int: """Cursor'dan beri biriken logları okuyup spool'a yazar. + Dönen değer, bu turda okunan error|critical kayıtların sayısıdır — acil + gönderim kararını döngü buna bakarak verir (CLAUDE.md §7). Sayı BURADA + üretilir çünkü kayıtlar yalnızca burada elde tutulur; spool'a yazıldıktan + sonra hangisinin bu turda geldiğini ayırt etmenin ucuz bir yolu kalmaz. + Pause'da da çalışır: metrik toplama gibi, log toplama da yerel kayıttır. SIRA ÖNEMLİDİR — önce kayıtlar spool'a, sonra cursor state'e yazılır. Ters @@ -167,7 +183,7 @@ def _collect_logs(source: LogSource, spool: Spool, state: State, store: StateSto # Log kaynağı erişilemez diye metrik toplama ve gönderim durmaz; # tur log'suz sürer, sorun bir sonraki turda yeniden denenir. _log(f"[logs] okunamadı: {error}") - return + return 0 for record in records: spool.add(RECORD_LOG, _log_payload(record)) @@ -179,6 +195,8 @@ def _collect_logs(source: LogSource, spool: Spool, state: State, store: StateSto if records: _log(f"[logs] {len(records)} kayıt ({_level_summary(records)})") + return sum(1 for record in records if record.level in URGENT_LEVELS) + def _prune_acked(state: State, store: StateStore, acked: list[str]) -> None: """Onaylanan komut id'lerini state'ten düşer. @@ -291,6 +309,70 @@ def _send_spool( _log(f"[send] {result.sent} kayıt gönderildi (spool: {spool.count()}).") +def _maybe_flush( + reading: MetricReading, + urgent_log_count: int, + config: Config, + state: State, + store: StateStore, + spool: Spool, + shipper: Shipper, +) -> bool: + """Eşik aşıldıysa acil gönderim yapar. Gönderim yapıldıysa True döner. + + Pause kapısı BURADADIR, çağıranda değil: acil gönderim de "buluta + yükleme"dir ve pause bunu durdurur (CLAUDE.md §7). Kural, yüklemeyi yapan + fonksiyonun içinde durursa yeni bir çağıran eklendiğinde unutulamaz. + Komut poll'ünün aksine burada istisna yoktur — flush telemetridir, + teardown kontrol mesajı değil. + + Pause'da eşiğe hiç bakılmaz: ölçüm ve loglar spool'a yazılmaya devam eder, + resume anında hepsi çıkar. crash_snapshots satırı da yazılmaz, çünkü o + satırın anlamı "flush attı"dır — atmadığı bir anda yazılsa yalan söylerdi. + """ + if not state.logging_enabled: + return False + + reason = flush_module.evaluate( + sample=reading.sample, + ram_percent=reading.ram_percent, + urgent_log_count=urgent_log_count, + config=config, + ) + if reason is None: + return False + + if flush_module.cooldown_active(state.last_flush_at, config.flush_cooldown_seconds): + # Veri kaybolmaz: eşiği aşan ölçüm de, tetikleyen log da spool'da + # duruyor ve normal gönderim turunda çıkacak. Bastırılan tek şey + # ACELE etmek — cooldown'ın amacı zaten flush selini önlemek. + _log(f"[flush] {reason} eşiği aşıldı, cooldown sürüyor — atlandı.") + return False + + # SIRA ÖNEMLİDİR: snapshot önce spool'a yazılır, sonra gönderim yapılır. + # Ters sırada snapshot bir sonraki tura kalır ve kendisini tetikleyen + # ölçümden ayrı bir istekte giderdi. + spool.add(RECORD_CRASH, flush_module.build_crash_snapshot(reason, config)) + + # Damga gönderimden ÖNCE yazılır: gönderim başarısız olsa bile cooldown + # başlamış sayılır. Aksi halde collector erişilemezken eşik her turda + # yeniden tutar, her tur yeni bir snapshot üretilir ve spool kesintinin + # sürdüğü süre boyunca boş yere şişerdi. + state.last_flush_at = utc_now_iso() + store.save(state) + + if not shipper.ready(): + _log( + f"[flush] {reason} eşiği aşıldı — backoff sürüyor, " + f"{shipper.backoff_seconds:.0f} sn sonra gönderilecek." + ) + return False + + _log(f"[flush] {reason} eşiği aşıldı — acil gönderim.") + _send_spool(shipper, config, state, store, spool) + return True + + def run(loader: ConfigLoader, store: StateStore, log_source: LogSource) -> None: """Agent'ı açılıştan kapanışa kadar çalıştırır. @@ -323,6 +405,10 @@ def run(loader: ConfigLoader, store: StateStore, log_source: LogSource) -> None: f"poll={config.command_poll_seconds}s (tick={TICK_SECONDS}s)" ) _log(f"[start] logging_enabled={state.logging_enabled}") + _log( + "[start] eklentiler: " + + (", ".join(config.enabled_addons) if config.enabled_addons else "yok (yalnızca çekirdek)") + ) _log( "[start] journal cursor: " + ("kayıtlı — kaldığı yerden" if state.journal_cursor else "yok — şimdiden başlanacak") @@ -348,10 +434,19 @@ def run(loader: ConfigLoader, store: StateStore, log_source: LogSource) -> None: config = loader.load() if now >= next_collect: - _collect(collector, spool) - _collect_logs(log_source, spool, state, store) + reading = _collect(collector, spool, config) + urgent_logs = _collect_logs(log_source, spool, state, store) next_collect = now + config.collect_interval_seconds + # Eşik ölçümün HEMEN ardından değerlendirilir: tetikleyen örnek + # ve loglar spool'a yeni yazıldı, yani acil gönderim onları da + # götürür. Pause denetimi _maybe_flush'ın içindedir. + # Flush gerçekten gönderdiyse normal gönderim sayacı ileri + # alınır — az önce boşalan spool'u saniyeler sonra bir kez daha + # yoklamanın anlamı yok. + if _maybe_flush(reading, urgent_logs, config, state, store, spool, shipper): + next_send = now + config.send_interval_seconds + if now >= next_poll: if _poll_commands(poller, config, state, store, spool, shipper): # delete uygulandı: cihaz kaydı sunucudan silindi, yerel diff --git a/agent/core/metrics.py b/agent/core/metrics.py index 49f74d5..dd951a1 100644 --- a/agent/core/metrics.py +++ b/agent/core/metrics.py @@ -1,11 +1,14 @@ """ Ölçüm toplama — psutil ile CPU, RAM, disk ve ağ. -Çekirdek metrikler: cpu_percent, ram_used_mb, disk_percent, +Çekirdek metrikler HER ZAMAN toplanır: cpu_percent, ram_used_mb, disk_percent, net_sent_mb, net_recv_mb. Eklenti alanları (sıcaklık, swap, load average, GPU) -M7'de bu modüle eklenecek; şu an toplanmıyor. +yalnızca config'in enabled_addons listesinde adı geçiyorsa okunur; kapalıyken +alan null kalır. -M2 KAPSAMI: ölçüm gerçek, gönderim yok. Örnekler yalnızca ekrana basılır. +Eklentilerin varsayılan olarak KAPALI olmasının sebebi ölçüm maliyeti değil, +anlam maliyetidir: her makinede olmayan bir sütun (sıcaklık sensörü, +NVIDIA GPU, Linux'a özgü load average) açıkça istenmeden doldurulmaz. """ from __future__ import annotations @@ -17,12 +20,31 @@ import psutil from agent.core.clock import utc_now_iso +from agent.core.config import ( + ADDON_GPU, + ADDON_LOAD_AVG, + ADDON_SWAP, + ADDON_TEMPERATURE, + Config, +) +from agent.core.gpu import GpuReader # Bayt -> MB çevrimi ikili tabanda (MiB). Şema sütunları "mb" adını taşır ama # RAM ve disk değerleri işletim sisteminin raporladığı ikili birimdir; tek bir # çevrim sabiti kullanmak RAM ve ağ sayılarının aynı ölçekte kalmasını sağlar. BYTES_PER_MB = 1024 * 1024 +# Sıcaklık okunurken denenecek sensör adları, sırayla. İlk bulunan kullanılır. +# +# Sıra keyfi değil, DARALAN güvenilirlikte: coretemp (Intel) ve k10temp (AMD) +# doğrudan CPU çekirdeğini ölçer; cpu_thermal ARM kartlarının (Raspberry Pi) +# karşılığıdır; acpitz ise anakart sensörüdür ve CPU'ya yalnızca yakındır. +# +# Şemada sensör ADInı tutan bir sütun YOK — yani temperature_c'nin neyi +# ölçtüğü satırdan okunamaz. Bu yüzden "bulduğun ilk sensörü al" yaklaşımı +# kullanılmaz: aynı grafikteki iki nokta iki farklı şeyi anlatabilirdi. +CPU_SENSOR_NAMES = ("coretemp", "k10temp", "cpu_thermal", "acpitz") + # Doluluk oranının okunduğu bağlama noktası. Envanterdeki disk_total_mb de # buradan gelir, böylece yüzde ile toplam aynı diski anlatır. DISK_MOUNT_POINT = "/" @@ -44,6 +66,36 @@ class MetricSample: net_sent_mb: float | None net_recv_mb: float | None + # --- eklentiler: kapalıyken None, yani sütun null kalır --- + # Varsayılan değerleri var çünkü çekirdek alanların aksine bunların + # OLMAMASI normaldir; her çağrının hepsini vermesi gerekmez. + temperature_c: float | None = None + swap_used_mb: int | None = None + load_avg_1: float | None = None + load_avg_5: float | None = None + load_avg_15: float | None = None + gpu_usage_percent: float | None = None + gpu_vram_used_mb: int | None = None + + +@dataclass(frozen=True) +class MetricReading: + """Tek bir toplama turunun sonucu. + + İki parça taşır çünkü ikisinin gideceği yer farklıdır: + * sample — spool'a yazılıp collector'a gönderilir, + * ram_percent — yalnızca flush eşiği karşılaştırmasında kullanılır, + hiçbir yere kaydedilmez. + + Yüzde MetricSample'ın İÇİNE konamaz: o nesne asdict() ile doğrudan wire + gövdesine dönüşüyor ve collector sözleşme dışı alanı 422 ile reddediyor. + Ayrı taşınması, eşiğin bakacağı değerin kaydedilen satırla aynı ölçüm anına + ait olmasını sağlar. + """ + + sample: MetricSample + ram_percent: float | None + class MetricsCollector: """Ardışık ölçümler arasında fark gerektiren alanların durumunu tutar. @@ -61,20 +113,35 @@ def __init__(self, disk_mount_point: str = DISK_MOUNT_POINT) -> None: self._previous_net: tuple[int, int, float] | None = None # psutil.cpu_percent'in taban değeri alındı mı. self._cpu_primed = False + # GPU eklentisi kapalıyken hiç kullanılmaz; nesne durum tuttuğu için + # (nvidia-smi var mı) toplayıcıyla aynı ömrü paylaşır. + self._gpu = GpuReader() - def collect(self) -> MetricSample: + def collect(self, config: Config) -> MetricReading: """Tek bir ölçüm alır. + config her çağrıda YENİDEN verilir: kullanıcı config.toml'da bir + eklentiyi açtığında değişiklik servis yeniden başlatılmadan geçerli + olsun diye (döngü her tick'te dosyayı yeniden okuyor). + Fark gerektiren alanlar (cpu_percent, net_*) hesaplanamadığında None döner; çağıran bunu doğrudan null olarak kaydeder. 0.0 yazmak "yük yoktu" / "trafik yoktu" anlamına gelirdi ve ölçülemeyen bir anı sıfırla karıştırırdı. + + Dönen MetricReading, kaydedilecek örneğin yanında RAM yüzdesini de + taşır; ikisi de AYNI psutil okumasından çıkar. """ memory = psutil.virtual_memory() disk = psutil.disk_usage(self._disk_mount_point) net_sent_mb, net_recv_mb = self._network_rates() - return MetricSample( + addons = config.enabled_addons + gpu_usage_percent, gpu_vram_used_mb = ( + self._gpu.read_usage() if ADDON_GPU in addons else (None, None) + ) + + sample = MetricSample( uuid=str(uuid.uuid4()), measured_at=utc_now_iso(), cpu_percent=self._cpu_percent(), @@ -87,8 +154,18 @@ def collect(self) -> MetricSample: disk_percent=round(disk.percent, 1), net_sent_mb=net_sent_mb, net_recv_mb=net_recv_mb, + temperature_c=_cpu_temperature() if ADDON_TEMPERATURE in addons else None, + swap_used_mb=_swap_used_mb() if ADDON_SWAP in addons else None, + **_load_average(enabled=ADDON_LOAD_AVG in addons), + gpu_usage_percent=gpu_usage_percent, + gpu_vram_used_mb=gpu_vram_used_mb, ) + # memory.percent, ram_used_mb ile aynı tanımı kullanır: + # (total - available) / total. Yani yüzde ile mutlak değer aynı şeyi + # iki ölçekte anlatır, biri diğeriyle çelişmez. + return MetricReading(sample=sample, ram_percent=round(memory.percent, 1)) + def _cpu_percent(self) -> float | None: """CPU kullanım yüzdesi — son çağrıdan bu yana geçen süre üzerinden. @@ -133,3 +210,67 @@ def _network_rates(self) -> tuple[float | None, float | None]: sent_rate = (counters.bytes_sent - previous_sent) / BYTES_PER_MB / elapsed recv_rate = (counters.bytes_recv - previous_recv) / BYTES_PER_MB / elapsed return round(sent_rate, 3), round(recv_rate, 3) + + +def _cpu_temperature() -> float | None: + """CPU sıcaklığı (°C) — bilinen sensörlerden ilk bulunan. + + Tanınan sensör yoksa None döner. "Ne bulursan onu al" DEĞİL, çünkü şemada + sensör adı için sütun yok: aynı sütuna bir makinede CPU, başka bir + makinede NVMe diskinin sıcaklığı yazılsaydı sayı karşılaştırılamaz olurdu. + + sensors_temperatures Linux dışında hiç tanımlı değildir; getattr ile + sorulur, böylece taşınabilirlik tek satırda çözülür. + """ + reader = getattr(psutil, "sensors_temperatures", None) + if reader is None: + return None + + try: + sensors = reader() + except (OSError, AttributeError): + return None + + for name in CPU_SENSOR_NAMES: + readings = sensors.get(name) + if readings and readings[0].current is not None: + return round(readings[0].current, 1) + + return None + + +def _swap_used_mb() -> int | None: + """Kullanılan swap (MB). Swap tanımlı değilse 0 döner — bu bir ölçümdür. + + Burada None YALNIZCA okuma başarısız olduğunda döner. Swap'ı olmayan bir + makinede doğru cevap "ölçülemedi" değil, "sıfır kullanılıyor"dur. + """ + try: + return psutil.swap_memory().used // BYTES_PER_MB + except (OSError, RuntimeError): + return None + + +def _load_average(*, enabled: bool) -> dict[str, float | None]: + """1/5/15 dakikalık yük ortalaması — Linux'a özgü. + + Üç alan tek fonksiyondan çıkar çünkü tek bir çağrının üç parçasıdır; + birini alıp diğerini alamamak mümkün değil. + + Windows'ta psutil bu değeri taklit eder ama ilk çağrıdan sonra 5 saniye + boyunca anlamsız değer verir; MVP Linux olduğu için sorun bugün yok, + yine de hata hali sessizce null'a düşer. + """ + if not enabled: + return {"load_avg_1": None, "load_avg_5": None, "load_avg_15": None} + + try: + one, five, fifteen = psutil.getloadavg() + except (OSError, AttributeError): + return {"load_avg_1": None, "load_avg_5": None, "load_avg_15": None} + + return { + "load_avg_1": round(one, 2), + "load_avg_5": round(five, 2), + "load_avg_15": round(fifteen, 2), + } diff --git a/collector/endpoints_ingest.py b/collector/endpoints_ingest.py index 9a6f654..9d1c5c3 100644 --- a/collector/endpoints_ingest.py +++ b/collector/endpoints_ingest.py @@ -3,14 +3,20 @@ Üçü de cihaz anahtarı ile korunur. Payload'da `device_id` YOKTUR; satırlara `device_id` ve `account_id` doğrulanmış anahtardan eklenir. + +Aynı ilkenin ikinci uygulaması `external_ip`tir: cihaz onu da GÖNDERMEZ. +Kimlik gibi, adres de cihazın kendi beyanı olamaz — değeri isteği gerçekten +alan taraf yazar (aşağıda `_external_ip`). """ from __future__ import annotations +import logging +from ipaddress import ip_address from typing import Annotated, Any, Literal from uuid import UUID -from fastapi import APIRouter +from fastapi import APIRouter, Header from pydantic import AwareDatetime, BaseModel, ConfigDict, Field from auth import AuthenticatedDevice, DeviceIdentity @@ -19,6 +25,8 @@ from version import COLLECTOR_VERSION from supabase_client import get_client +logger = logging.getLogger("tracebox.ingest") + router = APIRouter() # Tek istekte kabul edilen azami satır sayısı (tablo başına). Agent spool'u @@ -30,6 +38,24 @@ UUID_FIELD = "uuid" ID_COLUMN = "id" +# Dış IP'nin okunduğu başlık. Fly'ın proxy'si bunu KENDİSİ yazar ve istemcinin +# gönderdiği değeri ezer; `X-Forwarded-For` ise istemcinin önüne kendi +# uydurduğu adresleri ekleyebildiği bir listedir, yani kaynak olarak +# kullanılamaz. +# +# Başlık yoksa (yerel çalıştırma, başka bir barındırıcı) değer null kalır. +# Soketin karşı ucuna düşmek bir seçenek DEĞİL: proxy arkasında o adres +# proxy'nin kendisidir ve cihazın adresi diye kaydedilmesi, boş bırakmaktan +# daha kötüdür. +CLIENT_IP_HEADER = "Fly-Client-IP" + +# Eklentinin adı burada TEKRAR tanımlanır; agent'tan import EDİLMEZ. İki taraf +# ayrı deploy edilen ayrı programlar (collector imajında agent kodu yok); +# paylaştıkları şey Python nesnesi değil, wire sözleşmesidir. Adın iki tarafta +# aynı kaldığını tests/test_ingest_external_ip.py sınar — aynı yöntem log +# seviyeleri ve flush sebepleri için de kullanılıyor. +EXTERNAL_IP_ADDON = "external_ip" + class _Payload(BaseModel): """Ortak model ayarları. @@ -56,7 +82,8 @@ class InventoryIn(_Payload): last_boot: AwareDatetime | None = None agent_version: str | None = None gpu_model: str | None = None - external_ip: str | None = None + # external_ip BURADA YOK. `extra="forbid"` sayesinde alanın yokluğu pasif + # bir eksiklik değil aktif bir REDDİR: göndermeye çalışan agent 422 alır. enabled_addons: list[str] = Field(default_factory=list) @@ -127,12 +154,24 @@ class IngestIn(_Payload): @router.post("/inventory") -async def post_inventory(payload: InventoryIn, device: AuthenticatedDevice) -> dict: +async def post_inventory( + payload: InventoryIn, + device: AuthenticatedDevice, + fly_client_ip: Annotated[str | None, Header()] = None, +) -> dict: """Envanteri cihaz satırına yazar. Envanter zaman serisi değildir: her gönderim bir öncekinin üzerine yazar. + + `external_ip` gövdeden değil BAŞLIKTAN türetilir ve yalnızca burada yazılır. + Adresin tazeliği envanterin tazeliği kadardır (açılışta + değiştiğinde); + daha sık güncellemek `/ingest`e de aynı iki satırı koymak demek olurdu ama + o zaman rızayı okumak için cihaz satırındaki `enabled_addons`a bakmak + gerekirdi. Alan "statik eklenti" olarak tanımlandığı için bu değiş tokuş + bugün yapılmadı. """ fields = payload.model_dump(mode="json") + fields["external_ip"] = _external_ip(fly_client_ip, payload.enabled_addons) fields["last_seen"] = server_now() await call_or_503(lambda: get_client().update_device(device.id, fields)) @@ -193,6 +232,32 @@ async def get_verify(device: AuthenticatedDevice) -> dict: } +def _external_ip(header_value: str | None, enabled_addons: list[str]) -> str | None: + """Proxy'nin bildirdiği istemci adresi — eklenti kapalıysa None. + + Rıza her gönderimde YENİDEN sorulur ve kapalıyken açıkça None yazılır: + kullanıcı eklentiyi kapattığında `enabled_addons` değiştiği için envanter + zaten yeniden gönderilir, o gönderim de daha önce kaydedilmiş adresi siler. + "Yazmamak" yetmezdi — eski değer satırda kalırdı. + + Değer ayrıca IP olarak ÇÖZÜLEBİLDİĞİ doğrulanır. Beklenmedik bir şey gelmesi + proxy zincirinin varsayıldığı gibi olmadığını gösterir; onu olduğu gibi + kaydetmek metin sütununa güvenilmeyen girdi yazmak olurdu. + """ + if EXTERNAL_IP_ADDON not in enabled_addons: + return None + + if not header_value: + return None + + try: + return str(ip_address(header_value.strip())) + except ValueError: + # Değerin kendisi loglanmaz: doğrulanmamış, dışarıdan gelen bir metin. + logger.warning("%s başlığı IP adresi olarak çözülemedi", CLIENT_IP_HEADER) + return None + + def _row(item: BaseModel, device: DeviceIdentity) -> dict[str, Any]: """Payload kaydını tablo satırına çevirir: uuid → id, kimlik ve varış zamanı eklenir. diff --git a/collector/supabase_client.py b/collector/supabase_client.py index 902b888..7b35e6f 100644 --- a/collector/supabase_client.py +++ b/collector/supabase_client.py @@ -69,7 +69,7 @@ { # komut ack'inin türettiği durum "logging_enabled", - # agent'ın envanterden bildirdikleri (InventoryIn ile aynı 14 alan) + # agent'ın envanterden bildirdikleri (InventoryIn ile aynı 13 alan) "cpu_model", "cpu_cores_physical", "cpu_cores_logical", @@ -82,10 +82,12 @@ "last_boot", "agent_version", "gpu_model", - "external_ip", "enabled_addons", - # sunucunun damgaladığı + # sunucunun kendi bildiklerinden yazdıkları: varış anı ve bağlantının + # geldiği adres. external_ip listede ama InventoryIn'de DEĞİL — cihaz + # onu gönderemez, collector proxy başlığından türetir (M7). "last_seen", + "external_ip", } ) diff --git a/collector/version.py b/collector/version.py index d5be0d8..1bad26c 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.4.0" +COLLECTOR_VERSION = "0.5.0" diff --git a/tests/test_addons.py b/tests/test_addons.py new file mode 100644 index 0000000..9169dc5 --- /dev/null +++ b/tests/test_addons.py @@ -0,0 +1,436 @@ +""" +Seçilebilir eklentiler — enabled_addons filtresi, sensör seçimi ve GPU okuma. + +Eklentilerin ortak riski şudur: **kapalıyken sessizce açık olmak ya da açıkken +sessizce yanlış şeyi ölçmek.** İkisi de hata üretmez; biri kullanıcının +istemediği veriyi toplar, diğeri sütuna makineden makineye farklı anlam taşıyan +bir sayı yazar. + +Gerçek donanım burada kullanılamaz (bu makinede sensör olabilir, CI'da olmaz; +NVIDIA kartı olabilir, olmayabilir). Bu yüzden psutil ve nvidia-smi taklit +edilir: test edilen şey donanım değil, KARAR — hangi sensör seçiliyor, hangi +alan ne zaman dolduruluyor, hata nasıl yutuluyor. +""" + +from __future__ import annotations + +import subprocess +from dataclasses import asdict, replace +from pathlib import Path +from types import SimpleNamespace + +import psutil +import pytest + +import agent +from agent.core import metrics as metrics_module +from agent.core.config import ( + ADDON_EXTERNAL_IP, + ADDON_GPU, + ADDON_SWAP, + ADDON_TEMPERATURE, + KNOWN_ADDONS, + Config, + ConfigLoader, +) +from agent.core.gpu import GpuReader +from agent.core.inventory import collect_inventory +from agent.core.metrics import MetricsCollector + +BARE = Config(collector_url="https://collector.test", device_key="tbx_live_test") +ALL_ADDONS = replace(BARE, enabled_addons=KNOWN_ADDONS) + +# MetricSample'daki eklenti alanları — kapalıyken hepsi null kalmalı. +ADDON_FIELDS = ( + "temperature_c", + "swap_used_mb", + "load_avg_1", + "load_avg_5", + "load_avg_15", + "gpu_usage_percent", + "gpu_vram_used_mb", +) + + +class Warnings(list): + """`warn` yerine geçer; basılan uyarıları toplar.""" + + def __call__(self, message: str) -> None: + self.append(message) + + +def sensor(current: float | None, label: str = "") -> SimpleNamespace: + """psutil.sensors_temperatures çıktısındaki tek okuma.""" + return SimpleNamespace(label=label, current=current, high=None, critical=None) + + +class FakeRun: + """subprocess.run yerine geçer; çağrıları sayar, sonucu testten alır.""" + + def __init__(self, stdout: str = "", returncode: int = 0, error: Exception | None = None): + self._stdout = stdout + self._returncode = returncode + self._error = error + self.calls: list[list[str]] = [] + + def __call__(self, command, **kwargs): + self.calls.append(list(command)) + if self._error is not None: + raise self._error + return subprocess.CompletedProcess(command, self._returncode, self._stdout, "") + + +@pytest.fixture +def collector(): + """CPU yüzdesinin tabanı alınmış bir toplayıcı — ikinci ölçüm gerçek olur.""" + instance = MetricsCollector() + instance.collect(BARE) + return instance + + +# --- enabled_addons filtresi ------------------------------------------------ + + +def test_disabled_addons_leave_every_column_null(collector): + """Varsayılan config'te eklenti alanlarının HEPSİ null olmalı. + + Bu testin koruduğu şey bir gizlilik/rıza kuralıdır: kullanıcı istemediği + hiçbir ölçümü göndermemiş olmalı. Sızıntı sessizdir — sütun dolar, kimse + fark etmez. + """ + sample = asdict(collector.collect(BARE).sample) + + assert [field for field in ADDON_FIELDS if sample[field] is not None] == [] + + +def test_enabling_an_addon_fills_only_its_own_columns(collector, monkeypatch): + """Bir eklentiyi açmak diğerlerini açmaz. + + Tek bir `if` yanlış yazılırsa ("herhangi biri açıksa hepsini oku") kullanıcı + swap isterken GPU'yu da göndermeye başlar. + """ + monkeypatch.setattr(psutil, "swap_memory", lambda: SimpleNamespace(used=512 * 1024 * 1024)) + + sample = asdict(collector.collect(replace(BARE, enabled_addons=(ADDON_SWAP,))).sample) + + assert sample["swap_used_mb"] == 512 + assert [field for field in ADDON_FIELDS if field != "swap_used_mb" and sample[field]] == [] + + +def test_addon_list_is_read_from_the_config_on_every_collect(collector, monkeypatch): + """Eklenti açmak servisi yeniden başlatmayı gerektirmez. + + Döngü config'i her tick okuyor; toplayıcı listeyi __init__'te + saklasaydı kullanıcının config.toml'da yaptığı değişiklik ancak yeniden + başlatmada geçerli olurdu. + """ + monkeypatch.setattr(psutil, "swap_memory", lambda: SimpleNamespace(used=0)) + + before = collector.collect(BARE).sample.swap_used_mb + after = collector.collect(replace(BARE, enabled_addons=(ADDON_SWAP,))).sample.swap_used_mb + + assert before is None + assert after == 0 + + +# --- sıcaklık: hangi sensör ------------------------------------------------- + + +def test_known_cpu_sensors_are_preferred_in_order(monkeypatch): + """Sıra: coretemp -> k10temp -> cpu_thermal -> acpitz. + + Aynı makinede birden fazla sensör bulunur. acpitz (anakart) CPU'ya yalnızca + yakındır; coretemp çekirdeği doğrudan ölçer. Yanlış sıra, sütuna sistematik + olarak birkaç derece kaymış bir değer yazar — hata değil, sessiz yanlışlık. + """ + monkeypatch.setattr( + psutil, + "sensors_temperatures", + lambda: {"acpitz": [sensor(38.0)], "coretemp": [sensor(61.0)]}, + raising=False, + ) + + assert metrics_module._cpu_temperature() == 61.0 + + +def test_unknown_sensors_are_never_used(monkeypatch): + """Tanınmayan sensör varsa değer null kalır — rastgele bir sensör seçilmez. + + Şemada sensör adı için sütun YOK. "Ne bulursan yaz" kuralı, aynı sütuna bir + makinede CPU, diğerinde NVMe diskinin sıcaklığını yazardı; iki satır + karşılaştırılamaz hale gelirdi. + """ + monkeypatch.setattr( + psutil, + "sensors_temperatures", + lambda: {"nvme": [sensor(52.0)], "iwlwifi_1": [sensor(44.0)]}, + raising=False, + ) + + assert metrics_module._cpu_temperature() is None + + +def test_sensor_without_a_reading_is_skipped(monkeypatch): + """Sensör listesi boş ya da değeri None ise sütun null kalır.""" + monkeypatch.setattr( + psutil, + "sensors_temperatures", + lambda: {"coretemp": [], "k10temp": [sensor(None)]}, + raising=False, + ) + + assert metrics_module._cpu_temperature() is None + + +def test_platform_without_temperature_support_returns_null(monkeypatch): + """psutil.sensors_temperatures Linux dışında hiç TANIMLI DEĞİLDİR. + + Yokluğunu getattr ile sormak yerine doğrudan çağırmak, Windows'ta + AttributeError ile ölçüm turunu düşürürdü. + """ + monkeypatch.delattr(psutil, "sensors_temperatures", raising=False) + + assert metrics_module._cpu_temperature() is None + + +def test_sensor_read_failure_does_not_break_the_sample(monkeypatch): + """Sensör okunamazsa yalnızca o alan null olur, ölçüm turu sürer.""" + + def explode(): + raise OSError("sensör okunamadı") + + monkeypatch.setattr(psutil, "sensors_temperatures", explode, raising=False) + + assert metrics_module._cpu_temperature() is None + + +# --- swap ve load average --------------------------------------------------- + + +def test_zero_swap_is_a_measurement_not_a_missing_value(monkeypatch): + """Swap'ı olmayan makinede doğru cevap 0'dır, null değil. + + Fark grafikte görünür: null "ölçemedim" der, 0 "kullanılmıyor" der. + """ + monkeypatch.setattr(psutil, "swap_memory", lambda: SimpleNamespace(used=0)) + + assert metrics_module._swap_used_mb() == 0 + + +def test_swap_read_failure_falls_back_to_null(monkeypatch): + def explode(): + raise OSError("swap okunamadı") + + monkeypatch.setattr(psutil, "swap_memory", explode) + + assert metrics_module._swap_used_mb() is None + + +def test_load_average_fills_all_three_or_none_of_them(): + """Üç alan tek çağrıdan gelir; biri dolup diğeri boş kalamaz.""" + enabled = metrics_module._load_average(enabled=True) + disabled = metrics_module._load_average(enabled=False) + + assert all(value is not None for value in enabled.values()) + assert all(value is None for value in disabled.values()) + + +def test_load_average_failure_nulls_all_three(monkeypatch): + """Platform desteklemiyorsa üçü birden null olur, ölçüm düşmez.""" + + def explode(): + raise OSError("desteklenmiyor") + + monkeypatch.setattr(psutil, "getloadavg", explode, raising=False) + + assert metrics_module._load_average(enabled=True) == { + "load_avg_1": None, + "load_avg_5": None, + "load_avg_15": None, + } + + +# --- GPU: nvidia-smi -------------------------------------------------------- + + +def test_gpu_usage_is_parsed_from_the_csv_output(monkeypatch): + """`--format=csv,noheader,nounits` sayıyı birimsiz verir.""" + run = FakeRun(stdout="34, 2100\n") + monkeypatch.setattr(subprocess, "run", run) + + assert GpuReader().read_usage() == (34.0, 2100) + assert "--query-gpu=utilization.gpu,memory.used" in run.calls[0] + + +def test_only_the_first_gpu_is_read(monkeypatch): + """Çoklu kartta ilk satır alınır — şemada tek sütun var.""" + monkeypatch.setattr(subprocess, "run", FakeRun(stdout="34, 2100\n90, 8000\n")) + + assert GpuReader().read_usage() == (34.0, 2100) + + +def test_non_numeric_gpu_values_become_null(monkeypatch): + """Sürücü bir alanı raporlamıyorsa [N/A] yazar; sayı uydurulmaz.""" + monkeypatch.setattr(subprocess, "run", FakeRun(stdout="[N/A], [N/A]\n")) + + assert GpuReader().read_usage() == (None, None) + + +def test_failed_gpu_query_becomes_null(monkeypatch): + """Sıfırdan farklı çıkış kodu: sürücü hatası — çıktı ayrıştırılmaz. + + Taklit çıktı BİLEREK ayrıştırılabilir bırakıldı: çıkış kodu denetimi + kaldırılsaydı, boş bir stdout ile bu test yine yeşil kalırdı. Yarım kalmış + bir çıktının sayıya çevrilmesi, sürücü hatasını "GPU %34 yüklüydü" diye + kaydetmek demektir. + """ + monkeypatch.setattr(subprocess, "run", FakeRun(stdout="34, 2100\n", returncode=9)) + + assert GpuReader().read_usage() == (None, None) + + +def test_gpu_timeout_becomes_null(monkeypatch): + """Yanıt vermeyen nvidia-smi ölçüm döngüsünü bekletmez.""" + monkeypatch.setattr( + subprocess, "run", FakeRun(error=subprocess.TimeoutExpired("nvidia-smi", 2.0)) + ) + + assert GpuReader().read_usage() == (None, None) + + +def test_missing_nvidia_smi_is_remembered(monkeypatch): + """Program yoksa BİR kez denenir, bir daha denenmez. + + Ölçüm aralığı saniyelerle ifade ediliyor: olmayan bir programı her turda + başlatmaya çalışmak, GPU'su olmayan her makinede sürekli ve tamamen boş bir + süreç yaratma maliyetidir. + """ + run = FakeRun(error=FileNotFoundError("nvidia-smi yok")) + monkeypatch.setattr(subprocess, "run", run) + reader = GpuReader() + + assert reader.read_usage() == (None, None) + assert reader.read_usage() == (None, None) + assert reader.read_model() is None + assert len(run.calls) == 1, "olmayan program tekrar tekrar çağrıldı" + + +def test_a_timeout_is_not_remembered(monkeypatch): + """Geçici hata kalıcı sayılmaz — sürücü meşgulse sonraki tur yeniden dener.""" + run = FakeRun(error=subprocess.TimeoutExpired("nvidia-smi", 2.0)) + monkeypatch.setattr(subprocess, "run", run) + reader = GpuReader() + + reader.read_usage() + reader.read_usage() + + assert len(run.calls) == 2 + + +def test_gpu_model_is_read_for_the_inventory(monkeypatch): + run = FakeRun(stdout="NVIDIA GeForce RTX 4050 Laptop GPU\n") + monkeypatch.setattr(subprocess, "run", run) + + assert GpuReader().read_model() == "NVIDIA GeForce RTX 4050 Laptop GPU" + assert "--query-gpu=name" in run.calls[0] + + +def test_gpu_is_not_queried_while_the_addon_is_off(collector, monkeypatch): + """Eklenti kapalıyken nvidia-smi HİÇ çalıştırılmaz. + + Kapalı bir eklentinin bedeli sıfır olmalı; süreç başlatıp sonucu atmak + "kapalı" demek değildir. + """ + run = FakeRun(stdout="34, 2100\n") + monkeypatch.setattr(subprocess, "run", run) + + collector.collect(BARE) + + assert run.calls == [] + + +def test_gpu_model_enters_the_inventory_only_when_enabled(monkeypatch): + monkeypatch.setattr(subprocess, "run", FakeRun(stdout="NVIDIA Test Card\n")) + + assert collect_inventory(BARE).gpu_model is None + assert collect_inventory(replace(BARE, enabled_addons=(ADDON_GPU,))).gpu_model == ( + "NVIDIA Test Card" + ) + + +# --- external_ip: agent'ın göndermediği alan -------------------------------- + + +def test_the_agent_never_reports_its_own_external_ip(): + """external_ip envanterde ALAN OLARAK YOK. + + Cihazın kendi dış IP'sini bildirmesi, doğruluğu cihazın insafına bırakırdı: + bir agent istediği IP'yi yazabilirdi. Değeri, isteği gerçekten alan taraf + (collector) bağlantının kaynağından yazar. + + Kullanıcının tercihi yine de agent'tan gider — enabled_addons listesiyle. + """ + inventory = collect_inventory(ALL_ADDONS) + + assert "external_ip" not in asdict(inventory) + # Tercih yine de gidiyor: collector'ın "yazayım mı" sorusunun cevabı burada. + assert ADDON_EXTERNAL_IP in inventory.enabled_addons + + +# --- config: tanınmayan eklenti adı ----------------------------------------- + + +def test_unknown_addon_name_warns_without_stopping_the_agent(tmp_path): + """"temprature" yazan kullanıcı ne agent'ı kaybeder ne de sessiz kalır. + + Hata yükseltmek, tek harflik yazım hatası yüzünden izlemeyi tamamen + durdururdu — izleme aracının yapabileceği en kötü şey. + """ + path = tmp_path / "config.toml" + path.write_text( + 'collector_url = "https://collector.test"\n' + 'device_key = "tbx_live_test"\n' + 'enabled_addons = ["temprature", "swap"]\n' + ) + path.chmod(0o600) + warnings = Warnings() + + config = ConfigLoader(path, warn=warnings).load() + + assert config.enabled_addons == ("temprature", "swap") + assert len(warnings) == 1 + assert "temprature" in warnings[0] + assert ADDON_TEMPERATURE in warnings[0], "uyarı doğru yazımı göstermiyor" + + +def test_every_known_addon_name_passes_without_a_warning(tmp_path): + """Tanınan adların tamamı sessizce kabul edilmeli. + + Uyarı listesi bir adı yakalarsa, ya sabit listede ya da doğrulamada + tutarsızlık var demektir. + """ + path = tmp_path / "config.toml" + names = ", ".join(f'"{name}"' for name in KNOWN_ADDONS) + path.write_text( + 'collector_url = "https://collector.test"\n' + 'device_key = "tbx_live_test"\n' + f"enabled_addons = [{names}]\n" + ) + path.chmod(0o600) + warnings = Warnings() + + ConfigLoader(path, warn=warnings).load() + + assert warnings == [] + + +def test_the_example_config_lists_every_known_addon(): + """KNOWN_ADDONS ile config.example.toml aynı adları anmalı. + + Koda eklenip örnek dosyada anılmayan bir eklenti, kullanıcının varlığını + hiç öğrenemeyeceği bir eklentidir; tersi (örnekte olup kodda olmayan) ise + kurulumda uyarı basar. + """ + body = (Path(agent.__file__).parent / "config.example.toml").read_text(encoding="utf-8") + + assert [name for name in KNOWN_ADDONS if f'"{name}"' not in body] == [] diff --git a/tests/test_flush.py b/tests/test_flush.py new file mode 100644 index 0000000..4dd6468 --- /dev/null +++ b/tests/test_flush.py @@ -0,0 +1,532 @@ +""" +agent/core/flush.py + loop._maybe_flush — acil gönderim. + +Buradaki soruların ortak yanı **sessiz** olmalarıdır: eşik yanlış bağlanırsa, +cooldown yanlış tarafa kayarsa ya da pause'da flush atarsa hiçbir hata mesajı +çıkmaz. Agent çalışmaya devam eder; yalnızca çöküş anındaki veri ya hiç gelmez +ya da gereksiz yere sel olur. + +Testler üç katmanı ayırır: + * evaluate() — eşik kararı, + * cooldown_active() — sel koruması, + * build_crash_snapshot()— wire satırı, + * loop._maybe_flush() — bunların döngüdeki sırası ve yan etkileri. +""" + +from __future__ import annotations + +import uuid +from dataclasses import replace +from datetime import datetime, timedelta, timezone +from types import SimpleNamespace +from typing import Literal, get_args, get_origin + +import psutil +import pytest + +from agent.core import flush, loop +from agent.core.clock import utc_now_iso +from agent.core.config import ADDON_CRASH_PROCESSES, Config +from agent.core.metrics import MetricReading, MetricSample +from agent.core.shipper import SendResult +from agent.core.spool import RECORD_CRASH, Spool +from agent.core.state import StateStore +from agent.logsources.base import LogRecord, LogSourceError + +CONFIG = Config(collector_url="https://collector.test", device_key="tbx_live_test") + + +def make_reading( + *, + cpu: float | None = 1.0, + ram_percent: float | None = 10.0, + disk: float | None = 1.0, + ram_used_mb: int = 100, +) -> MetricReading: + """Eşiklerin ÇOK ALTINDA bir ölçüm; test yalnızca ilgilendiği alanı yükseltir.""" + sample = MetricSample( + uuid=str(uuid.uuid4()), + measured_at=utc_now_iso(), + cpu_percent=cpu, + ram_used_mb=ram_used_mb, + disk_percent=disk, + net_sent_mb=0.0, + net_recv_mb=0.0, + ) + return MetricReading(sample=sample, ram_percent=ram_percent) + + +def evaluate(reading: MetricReading, *, urgent: int = 0, config: Config = CONFIG) -> str | None: + return flush.evaluate( + sample=reading.sample, + ram_percent=reading.ram_percent, + urgent_log_count=urgent, + config=config, + ) + + +def iso_seconds_ago(seconds: float) -> str: + """`seconds` saniye önceki anın ISO damgası (negatif değer geleceği verir).""" + moment = datetime.now(timezone.utc) - timedelta(seconds=seconds) + return moment.isoformat(timespec="seconds") + + +class FakeShipper: + """`send_pending` çağrılarını sayar; hazır olma ve sonuç testten verilir.""" + + def __init__(self, *, ready: bool = True, ok: bool = True) -> None: + self._ready = ready + self._ok = ok + self.backoff_seconds = 0.0 if ready else 30.0 + self.calls: list[list[str]] = [] + + def ready(self) -> bool: + return self._ready + + def send_pending(self, config, applied_command_ids: list[str]) -> SendResult: + self.calls.append(list(applied_command_ids)) + return SendResult(ok=self._ok, sent=1, detail="" if self._ok else "HTTP 500") + + +class FakeLogSource: + """Verilen kayıtları bir kez döndürür; ikinci okumada boş gelir.""" + + def __init__(self, *records: LogRecord, error: Exception | None = None) -> None: + self._records = list(records) + self._error = error + + def read_since(self, cursor): + if self._error is not None: + raise self._error + records, self._records = self._records, [] + return records, "cursor-2" + + +def log_record(level: str) -> LogRecord: + return LogRecord(timestamp=utc_now_iso(), level=level, message=f"{level} mesajı") + + +class FakeProcess: + """psutil.Process'in _top_processes'in dokunduğu yüzeyi kadarı.""" + + def __init__(self, name: str, cpu: float, ram_mb: int, error: Exception | None = None) -> None: + self._name = name + self._cpu = cpu + self._ram_mb = ram_mb + self._error = error + + def cpu_percent(self) -> float: + if self._error is not None: + raise self._error + return self._cpu + + @property + def info(self) -> dict: + return {"name": self._name, "memory_info": SimpleNamespace(rss=self._ram_mb * 1024 * 1024)} + + +@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 fake_processes(monkeypatch): + """psutil.process_iter'ı verilen listeyle değiştirir ve beklemeyi sıfırlar.""" + + def install(processes: list[FakeProcess]) -> None: + monkeypatch.setattr(psutil, "process_iter", lambda attrs=None: iter(processes)) + # Gerçek örnekleme aralığı testte yalnızca beklemeye yol açar. + monkeypatch.setattr(flush, "PROCESS_SAMPLE_SECONDS", 0) + + return install + + +# --- evaluate: eşik kararı ------------------------------------------------- + + +def test_quiet_machine_does_not_flush(): + """Hiçbir eşik tutmuyorsa acil gönderim yok — normal tur yeterli.""" + assert evaluate(make_reading()) is None + + +def test_threshold_must_be_exceeded_not_merely_reached(): + """Eşik `>` ile karşılaştırılır, `>=` ile değil. + + Fark tek bir kıyaslama işaretidir ama sonucu büyüktür: disk eşiği 95'te + sabit duran bir makine `>=` ile her ölçümde flush ederdi ve cooldown + dolduğu her an yeniden — "acil" kavramı sürekli hale gelirdi. + """ + assert evaluate(make_reading(cpu=90.0)) is None + assert evaluate(make_reading(cpu=90.1)) == flush.REASON_CPU + + +def test_unmeasurable_field_never_crosses_a_threshold(): + """None "ölçülemedi" demektir; eşiği aşmış sayılmaz. + + İlk ölçümde cpu_percent bilerek None döner (metrics.py). None'ı eşikle + karşılaştırmaya çalışan bir kod TypeError ile döngüyü düşürürdü; sessizce + "aşıldı" sayan bir kod ise her agent açılışında bir flush üretirdi. + """ + assert evaluate(make_reading(cpu=None, ram_percent=None, disk=None)) is None + + +def test_ram_threshold_reads_the_percentage_not_the_megabytes(): + """Eşik yüzdedir; kayda giren ram_used_mb ile karıştırılamaz. + + Bu yüzden collect() ikisini birden döndürüyor. Yanlış alana bağlanırsa + 16 GB RAM'li bir makinede ram_used_mb neredeyse her zaman 90'ın üstünde + olur ve agent durmadan flush eder. + """ + assert evaluate(make_reading(ram_percent=10.0, ram_used_mb=15000)) is None + assert evaluate(make_reading(ram_percent=95.0, ram_used_mb=100)) == flush.REASON_RAM + + +def test_urgent_log_alone_triggers_a_flush(): + """error|critical log, yük normalken bile acil gönderim sebebidir. + + Projenin vaadi bu: kritik log 30 saniyelik turu beklemez. + """ + assert evaluate(make_reading(), urgent=1) == flush.REASON_LOG + + +def test_reason_priority_is_log_then_ram_then_cpu_then_disk(): + """Hepsi aynı anda tutabilir ama sütun tek değer alır. + + Sıra keyfi değil: log "bir şey bozuldu" der, diğer üçü yalnızca "yük + yüksek" der. Yanlış sıra, çöküş sonrası bakılan satırda daha az bilgi + taşıyan sebebi gösterir. + """ + everything = make_reading(cpu=99.0, ram_percent=99.0, disk=99.0) + + assert evaluate(everything, urgent=1) == flush.REASON_LOG + assert evaluate(everything) == flush.REASON_RAM + assert evaluate(make_reading(cpu=99.0, disk=99.0)) == flush.REASON_CPU + assert evaluate(make_reading(disk=99.0)) == flush.REASON_DISK + + +def test_thresholds_come_from_the_config_not_from_constants(): + """Eşikler config'ten okunur — kullanıcı kendi sınırını koyabilir.""" + strict = replace(CONFIG, flush_cpu_threshold=10) + + assert evaluate(make_reading(cpu=50.0)) is None + assert evaluate(make_reading(cpu=50.0), config=strict) == flush.REASON_CPU + + +# --- cooldown: sel koruması ------------------------------------------------ + + +def test_first_flush_is_never_blocked(): + """Damga yoksa daha önce hiç flush edilmemiştir; cooldown kapalıdır.""" + assert flush.cooldown_active(None, 20) is False + + +def test_unreadable_stamp_does_not_block_the_flush(): + """state.json elle bozulmuşsa flush kilitlenmez. + + Okunamayan bir damga "çok yakın zamanda flush ettik" anlamına gelemez; + aksi yönde yorumlanırsa tek bir bozuk satır acil gönderimi kalıcı olarak + kapatırdı. + """ + assert flush.cooldown_active("dün", 20) is False + + +def test_recent_flush_blocks_the_next_one(): + assert flush.cooldown_active(iso_seconds_ago(5), 20) is True + + +def test_cooldown_expires(): + assert flush.cooldown_active(iso_seconds_ago(25), 20) is False + + +def test_a_stamp_from_the_future_does_not_lock_the_flush_forever(): + """Sistem saati geri alındıysa damga gelecekte kalır. + + Geçen süre negatif çıkar. "Süre dolmadı" sayılırsa cooldown saat farkı + kadar — saatlerce, günlerce — açık kalır ve acil gönderim tamamen durur. + O yüzden negatif fark, cooldown'ın kapalı olması demektir. + """ + assert flush.cooldown_active(iso_seconds_ago(-3600), 20) is False + + +# --- build_crash_snapshot: wire satırı ------------------------------------- + + +def test_snapshot_is_written_even_when_the_addon_is_off(): + """crash_processes kapalıyken (varsayılan) satır yine yazılır, süreçler boş. + + Metrikler "CPU %95'ti" der; bu satır "flush GERÇEKTEN attı" der. İkisi + farklı sorulardır — satır atlanırsa acil gönderimin çalıştığına dair tek + kanıt kaybolur. + """ + snapshot = flush.build_crash_snapshot(flush.REASON_CPU, CONFIG) + + assert snapshot["processes"] == [] + assert snapshot["trigger_reason"] == flush.REASON_CPU + assert snapshot["measured_at"] + assert uuid.UUID(snapshot["uuid"]) + + +def test_snapshot_carries_processes_when_the_addon_is_on(fake_processes): + """Eklenti açıkken en çok kaynak yiyen süreçler listelenir.""" + fake_processes([FakeProcess(f"p{index}", cpu=float(index), ram_mb=index) for index in range(9)]) + config = replace(CONFIG, enabled_addons=(ADDON_CRASH_PROCESSES,)) + + snapshot = flush.build_crash_snapshot(flush.REASON_CPU, config) + + assert len(snapshot["processes"]) == flush.TOP_PROCESS_COUNT + assert [row["name"] for row in snapshot["processes"]] == ["p8", "p7", "p6", "p5", "p4"] + assert snapshot["processes"][0] == {"name": "p8", "cpu": 8.0, "ram_mb": 8} + + +def test_ram_triggered_snapshot_ranks_by_memory(fake_processes): + """Sıralama ölçütü, o an TÜKENEN kaynaktır. + + RAM eşiği aştığında CPU'ya göre sıralanmış bir liste yanlış beş süreci + gösterir — belleği bitiren süreç listede hiç görünmeyebilir. + """ + fake_processes( + [ + FakeProcess("cpu-yiyen", cpu=99.0, ram_mb=1), + FakeProcess("ram-yiyen", cpu=0.1, ram_mb=8000), + ] + ) + config = replace(CONFIG, enabled_addons=(ADDON_CRASH_PROCESSES,)) + + by_ram = flush.build_crash_snapshot(flush.REASON_RAM, config) + by_cpu = flush.build_crash_snapshot(flush.REASON_CPU, config) + + assert [row["name"] for row in by_ram["processes"]] == ["ram-yiyen", "cpu-yiyen"] + assert [row["name"] for row in by_cpu["processes"]] == ["cpu-yiyen", "ram-yiyen"] + + +def test_processes_that_vanish_mid_read_are_skipped(fake_processes): + """Okuma sırasında ölen ya da izin vermeyen süreç snapshot'ı düşürmez. + + Eksik bir snapshot, hiç snapshot olmamasından iyidir: o an bir daha gelmez. + """ + fake_processes( + [ + FakeProcess("ölen", cpu=99.0, ram_mb=10, error=psutil.NoSuchProcess(pid=1)), + FakeProcess("kapalı", cpu=98.0, ram_mb=10, error=psutil.AccessDenied()), + FakeProcess("sağlam", cpu=5.0, ram_mb=10), + ] + ) + config = replace(CONFIG, enabled_addons=(ADDON_CRASH_PROCESSES,)) + + snapshot = flush.build_crash_snapshot(flush.REASON_CPU, config) + + assert [row["name"] for row in snapshot["processes"]] == ["sağlam"] + + +def test_psutil_failure_leaves_the_snapshot_without_processes(monkeypatch): + """psutil beklenmedik bir hata verirse snapshot yine üretilir. + + Süreç listesi bir EKLENTİdir; onun hatası acil gönderimin kaydını + engellememeli. + """ + monkeypatch.setattr(flush, "_top_processes", _raise_psutil_error) + config = replace(CONFIG, enabled_addons=(ADDON_CRASH_PROCESSES,)) + + snapshot = flush.build_crash_snapshot(flush.REASON_RAM, config) + + assert snapshot["processes"] == [] + assert snapshot["trigger_reason"] == flush.REASON_RAM + + +def _raise_psutil_error(reason, limit): + raise psutil.Error() + + +def test_reasons_match_the_collectors_contract(): + """flush.py'nin ürettiği her sebep collector'ın kabul ettiği kümede olmalı. + + İki dosya birbirini import etmiyor; sözleşme yalnızca CLAUDE.md §4.2'de + yazılı. Buraya yeni bir sebep eklenip collector'daki Literal'a + eklenmezse agent 422 alır ve o snapshot HİÇ kaydedilmez — üstelik + yalnızca gerçek bir çöküş anında, yani test edilmesi en zor anda. + """ + from collector.endpoints_ingest import CrashSnapshotIn + + annotation = CrashSnapshotIn.model_fields["trigger_reason"].annotation + literal = next(arg for arg in get_args(annotation) if get_origin(arg) is Literal) + + assert set(flush.REASON_ORDER) == set(get_args(literal)) + + +# --- loop._maybe_flush: döngüdeki sıra ve yan etkiler ---------------------- + + +def test_collect_logs_counts_only_the_urgent_levels(store, spool): + """_collect_logs, flush'ın girdisini üretir: bu turda kaç acil log geldi. + + Sayı yalnızca okuma anında bilinebilir — kayıtlar spool'a karıştıktan sonra + hangisinin bu turda geldiğini ayırt etmenin ucuz yolu yok. Sayım info ve + warning'i de katarsa agent her sıradan log satırında flush eder; hiç + saymazsa "kritik log 30 saniye beklemez" vaadi sessizce çöker. + """ + source = FakeLogSource( + log_record("info"), + log_record("warning"), + log_record("error"), + log_record("critical"), + ) + + urgent = loop._collect_logs(source, spool, store.load(), store) + + assert urgent == 2 + assert spool.count() == 4, "acil olmayan loglar da spool'a yazılmalı" + + +def test_unreadable_log_source_reports_no_urgent_records(store, spool): + """Log okunamadıysa acil log sayısı sıfırdır — flush uydurulmaz.""" + source = FakeLogSource(error=LogSourceError("journalctl yok")) + + assert loop._collect_logs(source, spool, store.load(), store) == 0 + + +def maybe_flush(reading, store, spool, shipper, *, urgent: int = 0, config: Config = CONFIG): + state = store.load() + sent = loop._maybe_flush(reading, urgent, config, state, store, spool, shipper) + return sent, state + + +def crash_records(spool: Spool) -> list[dict]: + return [record.payload for record in spool.take(100) if record.type == RECORD_CRASH] + + +def test_crossing_a_threshold_sends_without_waiting_for_the_send_interval(store, spool): + """Projenin ana vaadinin kod karşılığı: veri 30 saniyeyi beklemez.""" + shipper = FakeShipper() + + sent, state = maybe_flush(make_reading(cpu=99.0), store, spool, shipper) + + assert sent is True + assert len(shipper.calls) == 1, "acil gönderim yapılmadı" + assert crash_records(spool)[0]["trigger_reason"] == flush.REASON_CPU + assert state.last_flush_at, "cooldown damgası yazılmadı" + # Damganın DİSKE yazıldığı burada kanıtlanamaz: başarılı gönderim zaten + # last_send için state'i kaydediyor, yani iddia yanlışlıkla tatmin olurdu. + # Kalıcılığın tek gerçek kanıtı gönderimin yapılmadığı yollardadır — + # test_backoff_records_the_snapshot_but_does_not_send ve + # test_failed_send_still_starts_the_cooldown. + + +def test_quiet_tick_touches_nothing(store, spool): + """Eşik tutmuyorsa ne snapshot ne gönderim ne damga.""" + shipper = FakeShipper() + + sent, state = maybe_flush(make_reading(), store, spool, shipper) + + assert sent is False + assert shipper.calls == [] + assert spool.count() == 0 + assert state.last_flush_at is None + + +def test_urgent_log_reaches_the_loop(store, spool): + """_collect_logs'un saydığı acil log, döngüde gerçekten flush'a dönüşür.""" + shipper = FakeShipper() + + sent, _ = maybe_flush(make_reading(), store, spool, shipper, urgent=1) + + assert sent is True + assert crash_records(spool)[0]["trigger_reason"] == flush.REASON_LOG + + +def test_cooldown_suppresses_both_the_snapshot_and_the_send(store, spool): + """Cooldown içindeyken hiçbir yan etki olmaz. + + Veri kaybolmuyor: eşiği aşan ölçüm ve loglar spool'da duruyor, normal + gönderim turunda çıkacak. Bastırılan tek şey ACELE etmek. + """ + state = store.load() + state.last_flush_at = iso_seconds_ago(1) + store.save(state) + shipper = FakeShipper() + + sent, _ = maybe_flush(make_reading(cpu=99.0), store, spool, shipper) + + assert sent is False + assert shipper.calls == [] + assert crash_records(spool) == [] + assert store.load().last_flush_at == state.last_flush_at, "damga tazelendi" + + +def test_paused_agent_never_flushes(store, spool): + """Pause = buluta yükleme yok. Acil gönderim de bir yüklemedir. + + Komut poll'ünün aksine bunun istisnası yok: flush telemetridir, teardown + kontrol mesajı değil. Snapshot da yazılmaz — o satırın anlamı "flush + attı"dır, atmadığı bir anda yazılsa yalan söylerdi. + """ + state = store.load() + state.logging_enabled = False + store.save(state) + shipper = FakeShipper() + + sent, _ = maybe_flush(make_reading(cpu=99.0), store, spool, shipper) + + assert sent is False + assert shipper.calls == [] + assert crash_records(spool) == [] + assert store.load().last_flush_at is None + + +def test_backoff_records_the_snapshot_but_does_not_send(store, spool): + """Collector erişilemezken snapshot yine alınır, gönderim ertelenir. + + Damga da yazılır: yazılmasaydı kesinti boyunca eşik her turda yeniden + tutar, her tur yeni bir snapshot üretilir ve spool boş yere şişerdi. + """ + shipper = FakeShipper(ready=False) + + sent, _ = maybe_flush(make_reading(cpu=99.0), store, spool, shipper) + + assert sent is False, "backoff sürerken gönderim denenmemeli" + assert shipper.calls == [] + assert len(crash_records(spool)) == 1 + assert store.load().last_flush_at is not None + + +def test_failed_send_still_starts_the_cooldown(store, spool): + """Gönderim 500 alsa bile cooldown başlar. + + Aksi halde sunucu hata verdiği sürece her ölçüm turu yeni bir acil + gönderim denemesi üretir — hatanın üstüne yük binerdi. + """ + shipper = FakeShipper(ok=False) + + sent, _ = maybe_flush(make_reading(cpu=99.0), store, spool, shipper) + + assert len(shipper.calls) == 1 + assert sent is True, "gönderim denendi — sayaç ileri alınmalı" + assert store.load().last_flush_at is not None + + +def test_snapshot_enters_the_spool_before_the_send(store, spool): + """Snapshot, kendisini tetikleyen veriyle AYNI istekte gitmelidir. + + Ters sırada yazılsaydı bir sonraki tura kalır ve çöküş anının kaydı, + çöküşü anlatan metriklerden ayrı düşerdi. + """ + seen: list[int] = [] + shipper = FakeShipper() + original = shipper.send_pending + + def record_then_send(config, applied_command_ids): + seen.append(len(crash_records(spool))) + return original(config, applied_command_ids) + + shipper.send_pending = record_then_send + + maybe_flush(make_reading(disk=99.0), store, spool, shipper) + + assert seen == [1], "gönderim anında snapshot henüz spool'da değildi" diff --git a/tests/test_ingest_external_ip.py b/tests/test_ingest_external_ip.py new file mode 100644 index 0000000..74bb31d --- /dev/null +++ b/tests/test_ingest_external_ip.py @@ -0,0 +1,228 @@ +""" +POST /inventory — dış IP'yi cihaz DEĞİL, isteği alan taraf yazar. + +Bu, `device_id`de verilen kararın ikinci uygulaması: cihazın kendisi hakkında +söylediği hiçbir şey doğruluk kaynağı değildir. `device_id` anahtardan +türetiliyor; `external_ip` de bağlantının kendisinden türetilir. Agent +gönderseydi, istediği adresi yazabilirdi — dashboard'da "bu makine nereden +bağlanıyor" sorusunun cevabı cihazın beyanı olurdu. + +İkinci mesele RIZA: alan bir eklentidir. Kullanıcı kapattığında yalnızca yeni +yazma durmaz, daha önce kaydedilmiş adres de silinir. + +Supabase taklit ediliyor, cihaz doğrulaması bağımlılık override'ıyla +sabitleniyor — anahtarın kendisi test_collector_security.py'de sınanıyor. +""" + +from __future__ import annotations + +from dataclasses import fields as dataclass_fields + +import pytest +from fastapi.testclient import TestClient + +import auth +import endpoints_commands +import endpoints_ingest +from agent.core.config import ADDON_EXTERNAL_IP +from agent.core.inventory import Inventory +from endpoints_ingest import EXTERNAL_IP_ADDON, InventoryIn +from main import app +from supabase_client import DEVICE_WRITABLE_COLUMNS + +DEVICE_ID = "33333333-3333-3333-3333-333333333333" +ACCOUNT_ID = "11111111-1111-1111-1111-111111111111" + +CLIENT_IP = "203.0.113.7" +HEADER = {"Fly-Client-IP": CLIENT_IP} + + +class FakeSupabase: + """Yalnızca cihaz satırına yazmayı kaydeden sahte istemci.""" + + def __init__(self) -> None: + self.updates: list[dict] = [] + self.inserted: list[tuple[str, list[dict]]] = [] + + async def update_device(self, device_id: str, fields: dict) -> None: + self.updates.append(fields) + + async def insert_rows(self, table: str, rows: list[dict]) -> None: + self.inserted.append((table, rows)) + + async def list_pending_commands(self, device_id: str) -> list[dict]: + return [] + + +@pytest.fixture +def fake_supabase(monkeypatch): + client = FakeSupabase() + monkeypatch.setattr(endpoints_ingest, "get_client", lambda: client) + monkeypatch.setattr(endpoints_commands, "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 send(client, *, addons=(ADDON_EXTERNAL_IP,), headers=HEADER, **body): + """POST /inventory — gövde varsayılanı yalnızca eklenti listesini taşır.""" + return client.post( + "/inventory", + json={"enabled_addons": list(addons), **body}, + headers=headers, + ) + + +def written(fake_supabase) -> dict: + """Cihaz satırına yazılan alanlar.""" + assert fake_supabase.updates, "update_device hiç çağrılmadı" + return fake_supabase.updates[-1] + + +# --- adresin kaynağı -------------------------------------------------------- + + +def test_the_proxy_header_is_written_to_the_device_row(client, fake_supabase): + """Eklenti açıkken adres başlıktan okunup satıra yazılır.""" + assert send(client).status_code == 200 + + assert written(fake_supabase)["external_ip"] == CLIENT_IP + + +def test_the_body_may_not_carry_an_external_ip(client, fake_supabase): + """Agent adresi göndermeye çalışırsa istek REDDEDİLİR, yok sayılmaz. + + 422, sözleşmenin sessizce esnemediğini gösterir: alan yok sayılsaydı + gönderen taraf gönderdiğini sanmaya devam ederdi. + """ + response = send(client, external_ip="198.51.100.9") + + assert response.status_code == 422 + assert fake_supabase.updates == [], "reddedilen istek yine de satıra dokundu" + + +def test_the_forwarded_for_header_is_not_a_source(client, fake_supabase): + """X-Forwarded-For adresin kaynağı değildir. + + O başlık bir LİSTEDİR ve istemci listenin başına istediğini ekleyebilir; + kaynak olarak kullanmak, reddedilen "gövdede gönder" yolunu başka bir adla + geri açmak olurdu. + """ + send(client, headers={"X-Forwarded-For": "198.51.100.9"}) + + assert written(fake_supabase)["external_ip"] is None + + +def test_a_missing_header_leaves_the_column_null(client, fake_supabase): + """Proxy başlığı yoksa (yerel çalıştırma) sütun null kalır. + + Soketin karşı ucuna düşülmez: proxy arkasında o adres proxy'nin kendisidir + ve cihazın adresi diye kaydedilmesi, boş bırakmaktan daha kötüdür. + """ + send(client, headers={}) + + assert written(fake_supabase)["external_ip"] is None + + +def test_a_malformed_header_is_not_stored(client, fake_supabase): + """IP olarak çözülemeyen değer sütuna OLDUĞU GİBİ yazılmaz.""" + send(client, headers={"Fly-Client-IP": "