Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion agent/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"
3 changes: 3 additions & 0 deletions agent/config.example.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
Expand Down
27 changes: 27 additions & 0 deletions agent/core/clock.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
30 changes: 30 additions & 0 deletions agent/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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(),
Expand Down
171 changes: 171 additions & 0 deletions agent/core/flush.py
Original file line number Diff line number Diff line change
@@ -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,
}
106 changes: 106 additions & 0 deletions agent/core/gpu.py
Original file line number Diff line number Diff line change
@@ -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
Loading
Loading