From 3480766c3e74bafbea62d0812a67de76e30c870d Mon Sep 17 00:00:00 2001 From: Nenad Djukic Date: Sat, 25 Jul 2026 17:40:27 +0200 Subject: [PATCH] Prijemnik dijagnostike sa TERMINIA uredjaja (logger/) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - UDP prijemnik (podrazumevano 46465, NE 46464 gde uredjaji slusaju discovery): dnevnik rada sa zaglavljem #TERMINIA, heartbeat JSON, sve ostalo se odbacuje - paket bez zaglavlja (FW < 0.13.4) se vodi kao "nepoznat uredjaj" po izvornoj IP adresi i tako je i oznacen u prikazu - dnevnik u fajlove po uredjaju i danu, rotacija po velicini + tri kvote (starost, po uredjaju, ukupno); otvoren fajl se nikad ne brise - heartbeat u memoriji; interval se meri iz razmaka javljanja, "tih" = 3x interval - web stranica (podrazumevano 127.0.0.1:8081): spisak uredjaja, crveno za tihe, poslednjih N linija sa pretragom, /dnevnik.txt i /api/uredjaji - bez baze i bez spoljnih zavisnosti; ne dira licencni servis, njegovu bazu, rute ni nginx - sanitizacija svega iz paketa (duzine, kontrolni znaci, ime foldera) — bez izlaska iz foldera dnevnika - systemd unit + primer env fajla + README sa uputstvom za pustanje na server - testovi: parsiranje (zaglavlje/JSON/bez zaglavlja/smece), registar, rotacija i kvote, web rute Nije pusteno na server — puštanje radi vlasnik po logger/README.md. --- .gitignore | 1 + .../02-uredjaji-heartbeat-discovery.md | 7 +- logger/README.md | 152 ++++++ logger/deploy/terminia-logger.env.example | 35 ++ logger/deploy/terminia-logger.service | 42 ++ logger/packet.go | 268 ++++++++++ logger/packet_test.go | 213 ++++++++ logger/registry.go | 194 +++++++ logger/registry_test.go | 113 +++++ logger/service.go | 259 ++++++++++ logger/service_test.go | 200 ++++++++ logger/store.go | 476 ++++++++++++++++++ logger/store_test.go | 245 +++++++++ logger/web.go | 365 ++++++++++++++ 14 files changed, 2569 insertions(+), 1 deletion(-) create mode 100644 logger/README.md create mode 100644 logger/deploy/terminia-logger.env.example create mode 100644 logger/deploy/terminia-logger.service create mode 100644 logger/packet.go create mode 100644 logger/packet_test.go create mode 100644 logger/registry.go create mode 100644 logger/registry_test.go create mode 100644 logger/service.go create mode 100644 logger/service_test.go create mode 100644 logger/store.go create mode 100644 logger/store_test.go create mode 100644 logger/web.go diff --git a/.gitignore b/.gitignore index 724362b..a504acf 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,7 @@ .env *.exe licence-server +terminia-logger crypto/private.pem log/*.log data/ diff --git a/docs/loggerservice/02-uredjaji-heartbeat-discovery.md b/docs/loggerservice/02-uredjaji-heartbeat-discovery.md index bc8d928..48bfe4a 100644 --- a/docs/loggerservice/02-uredjaji-heartbeat-discovery.md +++ b/docs/loggerservice/02-uredjaji-heartbeat-discovery.md @@ -1,7 +1,12 @@ # Uređaji: šta terminal VEĆ šalje (formati) i šta server treba da primi > Firmware strana je GOTOVA (FW 0.10.0+, modul `firmware/10_terminia/svc.cpp`). -> Server strana (prijemnik) NE postoji — to je zadatak u 04. +> ~~Server strana (prijemnik) NE postoji — to je zadatak u 04.~~ +> **UDP prijemnik je NAPISAN 2026-07-25:** [`logger/`](../../logger/README.md) +> (Go, bez baze, zaseban proces i systemd unit; dnevnik u fajlove sa rotacijom, +> heartbeat u memoriji, mala web stranica). Nije još pušten na server — +> puštanje radi vlasnik po uputstvu u `logger/README.md`. Upis u `devices` +> tabelu (tačka 2 ispod) ostaje za RBAC fazu. ## 1. LAN discovery (za servisera na licu mesta / PC alat) diff --git a/logger/README.md b/logger/README.md new file mode 100644 index 0000000..75ef500 --- /dev/null +++ b/logger/README.md @@ -0,0 +1,152 @@ +# Prijemnik dijagnostike sa TERMINIA uređaja + +Mali samostalan servis: sluša UDP pakete koje terminal šalje (dnevnik rada i +heartbeat), piše ih u fajlove po uređaju i po danu, i prikazuje na jednoj +web stranici. + +- **Nema bazu**, nema spoljnih Go zavisnosti (samo standardna biblioteka). +- **Ne dira licencni servis** — zaseban proces, zaseban port, zaseban + systemd unit. Deli samo Go modul (isti repo). +- Ugovor o formatu paketa: ASP repo, `docs/24-servisni-kanal.md`. + +## Šta prima + +| Vrsta | Kako se prepoznaje | Šta radi sa njim | +|---|---|---| +| Dnevnik rada (dbg) | čist tekst; prvi red `#TERMINIA ` | linije se upisuju u fajl uređaja | +| Heartbeat | paket počinje sa `{` (JSON sa `"hb":1`) | pamti se u memoriji (naziv, MAC, IP, FW, vreme) | +| Sve ostalo | — | tiho se odbacuje, brojač „odbačeno" na stranici | + +**Zaglavlje dnevnika postoji tek od FW 0.13.4.** Stariji uređaj šalje goli +tekst — takav se vodi kao **nepoznat uređaj**, razvrstava po izvornoj IP +adresi i tako je i **označen na stranici** (ne pravimo se da znamo ko je). + +Prijemnik **ne odgovara** na pakete. Prijem naloga na uređaju nije +implementiran i traži potpis (docs/24 §3.2), pa slanja naloga ovde nema. + +## Web stranica + +| Ruta | Šta daje | +|---|---| +| `/` | spisak uređaja: naziv, MAC, IP, FW, „poslednji put pre X"; **crveno** = nije se javio duže od 3× svog intervala | +| `/uredjaj?k=&q=&n=` | poslednjih N linija dnevnika, sa pretragom po tekstu | +| `/dnevnik.txt?k=&n=&q=` | isto, kao čist tekst (za `curl`) | +| `/api/uredjaji` | JSON spisak uređaja (za alate) | + +Interval javljanja se **meri** iz razmaka između dva heartbeat-a; dok ga nema, +koristi se `LOGGER_HB_DEFAULT_SEC`. + +## Pokretanje lokalno + +```powershell +cd D:\Terminia\TERMINIA-SOFTWARE\licence-server +go build -o terminia-logger.exe .\logger\cmd\terminia-logger + +$env:LOGGER_UDP_PORT="46465" +$env:LOGGER_HTTP_PORT="8081" +$env:LOGGER_DIR="D:\Terminia\TERMINIA-SOFTWARE\licence-server\data\logger" +.\terminia-logger.exe +``` + +Zatim otvoriti `http://127.0.0.1:8081`. + +Proba bez uređaja (PowerShell, u drugom prozoru): + +```powershell +$u=New-Object System.Net.Sockets.UdpClient; $u.Connect("127.0.0.1",46465) +function Send($s){ $b=[System.Text.Encoding]::UTF8.GetBytes($s); [void]$u.Send($b,$b.Length) } +Send '{"app":"TERMINIA","name":"TERM-001","mac":"1C:DB:D4:A6:A4:60","ip":"192.168.0.20","fw":"0.13.4","hb":1}' +Send "#TERMINIA TERM-001 1C:DB:D4:A6:A4:60 0.13.4`n[stampac] javio se na 192.168.123.100`n" +$u.Close() +``` + +Testovi: + +```powershell +go test ./logger/... +``` + +## Env promenljive + +| Promenljiva | Podrazumevano | Značenje | +|---|---|---| +| `LOGGER_UDP_PORT` | `46465` | port na koji uređaji šalju. **Ne 46464** — tamo uređaji slušaju discovery | +| `LOGGER_UDP_HOST` | (prazno) | adresa za slušanje; prazno = sve | +| `LOGGER_HTTP_HOST` | `127.0.0.1` | web sluša samo lokalno jer dnevnik sadrži **brojeve karata** | +| `LOGGER_HTTP_PORT` | `8081` | port web pregleda | +| `LOGGER_TOKEN` | (prazno) | ako je postavljen, web traži `?t=` (obavezno ako web nije samo lokalni) | +| `LOGGER_DIR` | `./data/logger` | folder sa dnevnicima | +| `LOGGER_KEEP_DAYS` | `14` | koliko dana se dnevnici čuvaju | +| `LOGGER_MAX_FILE_MB` | `10` | preko toga se dnevni fajl rotira (`.log.001`, `.002`…) | +| `LOGGER_MAX_DEVICE_MB` | `50` | kvota po uređaju; preko toga se briše najstarije | +| `LOGGER_MAX_TOTAL_MB` | `1024` | ukupna kvota za sve uređaje | +| `LOGGER_MAX_DEVICES` | `500` | koliko uređaja se drži u memoriji | +| `LOGGER_HB_DEFAULT_SEC` | `60` | pretpostavljeni interval dok se ne izmeri | + +Neispravna vrednost = servis se **ne pokreće** i kaže koja promenljiva ne valja +(bolje nego da tiho radi sa pogrešnim podešavanjem). + +## Šta se upisuje u uređaj + +U SPA uređaja: **Podešavanja → Servis** + +| Polje | Vrednost | +|---|---| +| `svcLogHost` | `157.180.119.151:46465` | +| „Šalji dnevnik rada na server" | **uključeno** | +| heartbeat (`svcHbSeconds`) | npr. `60` | + +Uređaj šalje **samo brojčani IP:port** — domen se ne može upisati (docs/24 §1). + +> ⚠️ Ovo ide preko interneta **u čistom tekstu**, a u dnevniku ima brojeva +> karata. Za trajno praćenje flote predviđen je HTTPS kanal (docs/24 §3); +> ovaj prijemnik je za dijagnostiku dok se problem lovi. + +## Puštanje na server (radi vlasnik) + +```bash +# 1) build na serveru +cd /root/projects/licence-server +git pull +go build -o terminia-logger ./logger/cmd/terminia-logger + +# 2) podešavanja +cp logger/deploy/terminia-logger.env.example logger.env +$EDITOR logger.env + +# 3) systemd +cp logger/deploy/terminia-logger.service /etc/systemd/system/ +systemctl daemon-reload +systemctl enable --now terminia-logger +systemctl status terminia-logger +journalctl -u terminia-logger -f + +# 4) firewall — SAMO taj port, UDP +ufw allow 46465/udp + +# 5) provera sa strane (ne sme da vrati odgovor, samo da ne pukne) +echo '#TERMINIA PROBA AA:BB:CC:DD:EE:FF 0.13.4 +proba veze' | nc -u -w1 157.180.119.151 46465 + +# 6) pregled (web je vezan za 127.0.0.1) — SSH tunel sa svog racunara: +ssh -L 8081:127.0.0.1:8081 root@157.180.119.151 +# pa u pregledacu: http://127.0.0.1:8081 +``` + +Ako web treba da bude dostupan bez tunela: postaviti `LOGGER_TOKEN`, staviti +nginx ispred sa TLS-om i lozinkom, i tek onda promeniti `LOGGER_HTTP_HOST`. +**Ne otvarati port 8081 direktno na internet.** + +Šta ovaj servis **ne dira**: bazu `licence_db`, proces `licence-server`, nginx +konfiguraciju, portove 8090/8091/8092. + +## Granice i odbrane + +- Paket može poslati bilo ko: sve iz paketa je ograničeno po dužini, kontrolni + znaci se izbacuju, a ime foldera se pravi isključivo od belih znakova + (`[A-Za-z0-9_-]`) — sadržaj paketa ne može da odluči gde se piše. +- Broj uređaja u memoriji je ograničen; kad se granica dostigne, izbacuje se + onaj koji se najduže nije javio. +- Dnevnici se brišu po tri pravila redom: starost → kvota po uređaju → ukupna + kvota. Fajl koji je trenutno otvoren za pisanje se nikad ne briše. +- Greška pri upisu ne obara prijem — vidi se kao brojač „greške upisa". diff --git a/logger/deploy/terminia-logger.env.example b/logger/deploy/terminia-logger.env.example new file mode 100644 index 0000000..39b2ff3 --- /dev/null +++ b/logger/deploy/terminia-logger.env.example @@ -0,0 +1,35 @@ +# Podesavanja prijemnika dijagnostike (systemd EnvironmentFile). +# Kopirati u /root/projects/licence-server/logger.env i izmeniti po potrebi. +# Sve promenljive su opcione - bez njih vaze podrazumevane vrednosti iz koda. + +# --- mreza --------------------------------------------------------------- +# UDP port na koji uredjaji salju dnevnik i heartbeat (svcLogHost = IP:PORT). +# NE koristiti 46464 - na tom portu uredjaji slusaju discovery. +LOGGER_UDP_PORT=46465 +# Prazno = slusaj na svim adresama. +LOGGER_UDP_HOST= + +# Web pregled. Podrazumevano SAMO lokalno, jer dnevnik sadrzi brojeve karata. +# Napolje ide kroz SSH tunel ili nginx sa lozinkom. +LOGGER_HTTP_HOST=127.0.0.1 +LOGGER_HTTP_PORT=8081 + +# Ako se web ipak otvori ka mrezi, obavezno postaviti token (pristup: ?t=). +# LOGGER_TOKEN= + +# --- dnevnici ------------------------------------------------------------ +# Folder sa dnevnicima (systemd StateDirectory pravi /var/lib/terminia-logger). +LOGGER_DIR=/var/lib/terminia-logger + +# Granice da dnevnici ne napune disk (najstarije se brise prvo). +LOGGER_KEEP_DAYS=14 +LOGGER_MAX_FILE_MB=10 +LOGGER_MAX_DEVICE_MB=50 +LOGGER_MAX_TOTAL_MB=1024 + +# Najveci broj uredjaja koji se drzi u memoriji (paket moze poslati bilo ko). +LOGGER_MAX_DEVICES=500 + +# Pretpostavljeni interval javljanja kad ga uredjaj jos nije pokazao. +# Uredjaj je "tih" kad se ne javi duze od 3x ovog intervala. +LOGGER_HB_DEFAULT_SEC=60 diff --git a/logger/deploy/terminia-logger.service b/logger/deploy/terminia-logger.service new file mode 100644 index 0000000..9a3a30c --- /dev/null +++ b/logger/deploy/terminia-logger.service @@ -0,0 +1,42 @@ +# Prijemnik servisne dijagnostike sa TERMINIA uredjaja. +# Obrazac: licence-server.service (vidi docs/SETUP.md), ali BEZ MySQL zavisnosti - +# ovaj servis nema bazu i ne dodiruje licencni servis. +# +# Instalacija: +# cp logger/deploy/terminia-logger.service /etc/systemd/system/ +# cp logger/deploy/terminia-logger.env.example /root/projects/licence-server/logger.env +# systemctl daemon-reload && systemctl enable --now terminia-logger + +[Unit] +Description=TERMINIA prijemnik dijagnostike (UDP dnevnik + heartbeat) +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +User=root +WorkingDirectory=/root/projects/licence-server +ExecStart=/root/projects/licence-server/terminia-logger +EnvironmentFile=-/root/projects/licence-server/logger.env +Restart=always +RestartSec=5 +StandardOutput=journal +StandardError=journal +SyslogIdentifier=terminia-logger +KillMode=mixed +KillSignal=SIGTERM +TimeoutStopSec=15 + +# Dnevnici sadrze BROJEVE KARATA - folder je citljiv samo vlasniku procesa. +StateDirectory=terminia-logger +StateDirectoryMode=0750 + +# Osnovno ogradjivanje: servis samo slusa UDP i pise u svoj folder. +NoNewPrivileges=true +PrivateTmp=true +ProtectSystem=full +# ProtectHome se NE ukljucuje: binar i EnvironmentFile stoje pod /root/projects +# (isto kao licence-server), a ProtectHome bi taj folder sakrio od servisa. + +[Install] +WantedBy=multi-user.target diff --git a/logger/packet.go b/logger/packet.go new file mode 100644 index 0000000..5f8e121 --- /dev/null +++ b/logger/packet.go @@ -0,0 +1,268 @@ +// Package logger je prijemnik servisne dijagnostike sa TERMINIA uredjaja. +// +// Uredjaj salje dve vrste UDP paketa na isto odrediste (cfg.svcLogHost): +// - dnevnik rada (dbg): cist tekst, prvi red je zaglavlje "#TERMINIA ..." +// - heartbeat: JSON identitet sa "hb":1 +// +// Ugovor je opisan u ASP repou, docs/24-servisni-kanal.md. +package logger + +import ( + "encoding/json" + "net" + "strings" + "unicode" + "unicode/utf8" +) + +// Vrsta paketa. +type Kind int + +const ( + KindGarbage Kind = iota // ne prepoznajemo ga - odbacuje se + KindHeartbeat // JSON identitet + KindLog // linije dnevnika +) + +// Tvrde granice. UDP paket moze doci od bilo koga, pa se sve sto ulazi u +// memoriju ili u fajl mora ograniciti pre nego sto se dodirne. +const ( + MaxPacketBytes = 4096 // uredjaj salje do ~1100 B; ostalo je visak + MaxLinesPerPacket = 64 + MaxLineLen = 400 + MaxNameLen = 40 + MaxFWLen = 24 + MaxSerialLen = 40 + MaxKeyLen = 48 +) + +const headerPrefix = "#TERMINIA" + +// Packet je rezultat parsiranja jednog UDP paketa. Nikad ne sadrzi sirove +// bajtove iz mreze - sva polja su prosla kroz sanitizaciju. +type Packet struct { + Kind Kind + Name string + MAC string // kanonski "AA:BB:CC:DD:EE:FF", prazan ako ga nema + IP string + FW string + Serial string + Identified bool // paket sam kaze ko je (zaglavlje ili heartbeat) + Lines []string + SourceIP string +} + +// heartbeatJSON prati format iz firmvera (pure::discoveryReply + "hb"). +// Polje "id" postoji zbog starijih/drugacijih klijenata koji serijski broj +// salju pod tim imenom - ne kostaje nista, a spasava kompatibilnost. +type heartbeatJSON struct { + App string `json:"app"` + Name string `json:"name"` + MAC string `json:"mac"` + IP string `json:"ip"` + FW string `json:"fw"` + Serial string `json:"serial"` + ID string `json:"id"` + HB json.RawMessage `json:"hb"` +} + +// ParsePacket razvrstava jedan primljen paket. Nikad ne panici i nikad ne +// vraca gresku - los paket je samo KindGarbage. +func ParsePacket(data []byte, sourceIP string) Packet { + p := Packet{Kind: KindGarbage, SourceIP: sanitizeText(sourceIP, 45)} + if len(data) == 0 { + return p + } + if len(data) > MaxPacketBytes { + data = data[:MaxPacketBytes] + } + body := strings.TrimLeft(string(data), " \t\r\n\x00") + if body == "" { + return p + } + if body[0] == '{' { + return parseHeartbeat(body, p) + } + return parseLog(body, p) +} + +func parseHeartbeat(body string, p Packet) Packet { + var hb heartbeatJSON + if err := json.Unmarshal([]byte(body), &hb); err != nil { + return p // krnj JSON - odbacujemo, servis se ne ruzi zbog toga + } + mac, okMac := NormalizeMAC(hb.MAC) + name := sanitizeText(hb.Name, MaxNameLen) + serial := hb.Serial + if serial == "" { + serial = hb.ID + } + // Bez ijednog upotrebljivog polja identiteta paket nema smisla: to je + // neki tudji JSON koji je slucajno stigao na nas port. + if !okMac && name == "" { + return p + } + p.Kind = KindHeartbeat + p.Identified = true + p.Name = name + p.MAC = mac + p.IP = sanitizeIP(hb.IP) + p.FW = sanitizeText(hb.FW, MaxFWLen) + p.Serial = sanitizeText(serial, MaxSerialLen) + return p +} + +func parseLog(body string, p Packet) Packet { + p.Kind = KindLog + raw := strings.Split(body, "\n") + if len(raw) > 0 && strings.HasPrefix(raw[0], headerPrefix) { + parseHeader(raw[0], &p) + raw = raw[1:] + } + for _, ln := range raw { + if len(p.Lines) >= MaxLinesPerPacket { + break + } + s := sanitizeText(ln, MaxLineLen) + if s == "" { + continue + } + p.Lines = append(p.Lines, s) + } + if len(p.Lines) == 0 && !p.Identified { + p.Kind = KindGarbage + } + return p +} + +// parseHeader cita "#TERMINIA ". +// +// Naziv uredjaja SME da sadrzi razmake (korisnik ga slobodno kuca), pa se ne +// oslanjamo na redni broj polja: trazi se prvo polje koje je validan MAC, sve +// pre njega je naziv, prvo posle njega je verzija firmvera. +func parseHeader(line string, p *Packet) { + fields := strings.Fields(strings.TrimPrefix(line, headerPrefix)) + macAt := -1 + for i, f := range fields { + if _, ok := NormalizeMAC(f); ok { + macAt = i + break + } + } + if macAt < 0 { + // Zaglavlje bez MAC-a nije upotrebljivo za razvrstavanje; ipak + // uzimamo naziv ako ga ima, ali uredjaj ostaje "neidentifikovan". + if len(fields) > 0 { + p.Name = sanitizeText(strings.Join(fields, " "), MaxNameLen) + } + return + } + p.MAC, _ = NormalizeMAC(fields[macAt]) + p.Name = sanitizeText(strings.Join(fields[:macAt], " "), MaxNameLen) + if macAt+1 < len(fields) { + p.FW = sanitizeText(fields[macAt+1], MaxFWLen) + } + p.Identified = true +} + +// NormalizeMAC prima "1C:DB:D4:A6:A4:60", "1c-db-d4-a6-a4-60" ili +// "1CDBD4A6A460" i vraca kanonski oblik sa dvotackama. +func NormalizeMAC(s string) (string, bool) { + if len(s) > 24 { + return "", false + } + var hex strings.Builder + for _, r := range s { + switch { + case r >= '0' && r <= '9', r >= 'A' && r <= 'F': + hex.WriteRune(r) + case r >= 'a' && r <= 'f': + hex.WriteRune(unicode.ToUpper(r)) + case r == ':' || r == '-' || r == '.': + // razdvajaci se preskacu + default: + return "", false + } + } + h := hex.String() + if len(h) != 12 { + return "", false + } + var out strings.Builder + for i := 0; i < 12; i += 2 { + if i > 0 { + out.WriteByte(':') + } + out.WriteString(h[i : i+2]) + } + return out.String(), true +} + +// DeviceKey je kljuc pod kojim se uredjaj vodi u registru I u imenu foldera. +// Kad MAC nedostaje (FW stariji od 0.13.4 salje dnevnik bez zaglavlja), jedino +// sto imamo je izvorna IP adresa - takav uredjaj se JASNO oznacava kao +// nepoznat, da se ne pravimo da znamo ko je. +func DeviceKey(mac, sourceIP string) string { + if m, ok := NormalizeMAC(mac); ok { + return SanitizeSegment(strings.ReplaceAll(m, ":", "")) + } + ip := sourceIP + if ip == "" { + ip = "bez-adrese" + } + return SanitizeSegment("nepoznat-" + ip) +} + +// SanitizeSegment pravi bezbedan deo putanje. Sve sto nije [A-Za-z0-9_-] +// postaje '_', pa ".." i "/" ne mogu da prezive - podatak iz paketa nikad ne +// sme da odluci gde se pise. +func SanitizeSegment(s string) string { + var b strings.Builder + for _, r := range s { + switch { + case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9', r == '-', r == '_': + b.WriteRune(r) + default: + b.WriteByte('_') + } + if b.Len() >= MaxKeyLen { + break + } + } + out := strings.Trim(b.String(), "_-") + if out == "" { + return "nepoznat" + } + return out +} + +// sanitizeText izbacuje kontrolne znake (ukljucujuci CR i ESC sekvence koje bi +// unakazile ispis u terminalu) i sece string na zadatu duzinu. +func sanitizeText(s string, max int) string { + var b strings.Builder + n := 0 + for _, r := range s { + if r == '\t' { + r = ' ' + } + // utf8.RuneError se javlja na neispravnim bajtovima - to je binarno + // smece koje je stiglo na port, ne tekst uredjaja. + if r < 0x20 || r == 0x7f || r == utf8.RuneError || (!unicode.IsPrint(r) && r != ' ') { + continue + } + b.WriteRune(r) + n++ + if n >= max { + break + } + } + return strings.TrimSpace(b.String()) +} + +func sanitizeIP(s string) string { + s = strings.TrimSpace(s) + if ip := net.ParseIP(s); ip != nil { + return ip.String() + } + return "" +} diff --git a/logger/packet_test.go b/logger/packet_test.go new file mode 100644 index 0000000..fd5a9ff --- /dev/null +++ b/logger/packet_test.go @@ -0,0 +1,213 @@ +package logger + +import ( + "strings" + "testing" +) + +func TestParseLogSaZaglavljem(t *testing.T) { + raw := "#TERMINIA TERM-001 1C:DB:D4:A6:A4:60 0.13.4\n" + + "[stampac] opet se javlja na 192.168.123.100 (ping)\n" + + "[izvestaj] dan 20659: 14 karata, 2100 RSD\n" + p := ParsePacket([]byte(raw), "192.168.0.20") + + if p.Kind != KindLog { + t.Fatalf("vrsta = %v, ocekivano KindLog", p.Kind) + } + if !p.Identified { + t.Error("paket sa zaglavljem mora biti identifikovan") + } + if p.Name != "TERM-001" || p.MAC != "1C:DB:D4:A6:A4:60" || p.FW != "0.13.4" { + t.Errorf("identitet = %q / %q / %q", p.Name, p.MAC, p.FW) + } + if len(p.Lines) != 2 { + t.Fatalf("linija = %d, ocekivano 2: %q", len(p.Lines), p.Lines) + } + if !strings.HasPrefix(p.Lines[0], "[stampac]") { + t.Errorf("prva linija = %q", p.Lines[0]) + } + if got := DeviceKey(p.MAC, p.SourceIP); got != "1CDBD4A6A460" { + t.Errorf("kljuc = %q", got) + } +} + +func TestParseZaglavljeSaRazmakomUNazivu(t *testing.T) { + // Naziv uredjaja kuca korisnik i sme da ima razmake. + p := ParsePacket([]byte("#TERMINIA Bazen Ulaz A 1c-db-d4-a6-a4-60 0.13.4\nzdravo\n"), "10.0.0.5") + if p.Name != "Bazen Ulaz A" { + t.Errorf("naziv = %q", p.Name) + } + if p.MAC != "1C:DB:D4:A6:A4:60" { + t.Errorf("mac = %q", p.MAC) + } + if p.FW != "0.13.4" { + t.Errorf("fw = %q", p.FW) + } +} + +func TestParseHeartbeat(t *testing.T) { + raw := `{"app":"TERMINIA","name":"TERM-001","mac":"1C:DB:D4:A6:A4:60",` + + `"ip":"192.168.0.20","fw":"0.13.4","serial":"1CDBD4A6A460","hb":1}` + p := ParsePacket([]byte(raw), "192.168.0.20") + + if p.Kind != KindHeartbeat { + t.Fatalf("vrsta = %v, ocekivano KindHeartbeat", p.Kind) + } + if !p.Identified || p.MAC != "1C:DB:D4:A6:A4:60" || p.IP != "192.168.0.20" || + p.FW != "0.13.4" || p.Serial != "1CDBD4A6A460" || p.Name != "TERM-001" { + t.Errorf("identitet = %+v", p) + } + if len(p.Lines) != 0 { + t.Errorf("heartbeat ne sme imati linije: %q", p.Lines) + } +} + +func TestParseHeartbeatSaIdPoljem(t *testing.T) { + // Stariji/drugaciji klijent salje serijski pod imenom "id". + p := ParsePacket([]byte(`{"name":"T","mac":"AABBCCDDEEFF","id":"XYZ","hb":1}`), "10.0.0.1") + if p.Kind != KindHeartbeat || p.Serial != "XYZ" { + t.Errorf("rezultat = %+v", p) + } +} + +func TestParseBezZaglavlja(t *testing.T) { + // FW stariji od 0.13.4: goli tekst, znamo samo izvornu adresu. + p := ParsePacket([]byte("nesto se desilo\ndrugi red\n"), "192.168.0.77") + if p.Kind != KindLog { + t.Fatalf("vrsta = %v", p.Kind) + } + if p.Identified { + t.Error("paket bez zaglavlja NE SME da vazi za identifikovan") + } + if p.MAC != "" { + t.Errorf("mac = %q, ocekivano prazno", p.MAC) + } + key := DeviceKey(p.MAC, p.SourceIP) + if key != "nepoznat-192_168_0_77" { + t.Errorf("kljuc = %q", key) + } +} + +func TestParseSmeceNeObaraServis(t *testing.T) { + slucajevi := [][]byte{ + nil, + {}, + []byte("\n\n\n"), + []byte("{"), // krnj JSON + []byte(`{"app":"TERMINIA"`), // presecen paket + []byte(`{"nesto":"tudje"}`), // JSON bez identiteta + {0x00, 0x01, 0x02, 0xff, 0xfe}, // binarno smece + []byte("#TERMINIA\n"), // zaglavlje bez ijednog polja + []byte(strings.Repeat("A", MaxPacketBytes*3)), // predugacak paket + } + for i, c := range slucajevi { + func() { + defer func() { + if r := recover(); r != nil { + t.Fatalf("slucaj %d: panika: %v", i, r) + } + }() + p := ParsePacket(c, "1.2.3.4") + // Bilo koji ishod je prihvatljiv osim pada; ako je nesto proslo, + // mora biti u granicama. + if len(p.Lines) > MaxLinesPerPacket { + t.Errorf("slucaj %d: previse linija %d", i, len(p.Lines)) + } + for _, ln := range p.Lines { + if len([]rune(ln)) > MaxLineLen { + t.Errorf("slucaj %d: linija duza od granice", i) + } + } + }() + } +} + +func TestParseSeceIKontrolneZnake(t *testing.T) { + dugacko := strings.Repeat("x", 1000) + raw := "#TERMINIA " + strings.Repeat("N", 200) + " 1C:DB:D4:A6:A4:60 " + strings.Repeat("9", 100) + "\n" + + "linija\x1b[31m sa ESC \x07 i \x00 nulom\n" + dugacko + "\n" + p := ParsePacket([]byte(raw), "1.2.3.4") + + if len([]rune(p.Name)) > MaxNameLen { + t.Errorf("naziv nije skracen: %d", len([]rune(p.Name))) + } + if len([]rune(p.FW)) > MaxFWLen { + t.Errorf("fw nije skracen: %d", len([]rune(p.FW))) + } + if strings.ContainsAny(p.Lines[0], "\x1b\x07\x00") { + t.Errorf("kontrolni znaci nisu uklonjeni: %q", p.Lines[0]) + } + if len([]rune(p.Lines[1])) > MaxLineLen { + t.Errorf("linija nije skracena: %d", len([]rune(p.Lines[1]))) + } +} + +func TestParsePreviseLinija(t *testing.T) { + var b strings.Builder + b.WriteString("#TERMINIA T 1C:DB:D4:A6:A4:60 0.13.4\n") + for i := 0; i < 500; i++ { + b.WriteString("linija\n") + } + p := ParsePacket([]byte(b.String()), "1.2.3.4") + if len(p.Lines) != MaxLinesPerPacket { + t.Errorf("linija = %d, ocekivano %d", len(p.Lines), MaxLinesPerPacket) + } +} + +func TestNormalizeMAC(t *testing.T) { + dobri := map[string]string{ + "1C:DB:D4:A6:A4:60": "1C:DB:D4:A6:A4:60", + "1c:db:d4:a6:a4:60": "1C:DB:D4:A6:A4:60", + "1cdbd4a6a460": "1C:DB:D4:A6:A4:60", + "1C-DB-D4-A6-A4-60": "1C:DB:D4:A6:A4:60", + } + for in, want := range dobri { + got, ok := NormalizeMAC(in) + if !ok || got != want { + t.Errorf("NormalizeMAC(%q) = %q,%v", in, got, ok) + } + } + losi := []string{"", "1C:DB:D4:A6:A4", "1C:DB:D4:A6:A4:6G", "../../etc/passwd", + strings.Repeat("A", 200), "1C:DB:D4:A6:A4:60:70"} + for _, in := range losi { + if _, ok := NormalizeMAC(in); ok { + t.Errorf("NormalizeMAC(%q) je prosao, a ne bi smeo", in) + } + } +} + +func TestSanitizeSegmentNemaIzlaskaIzFoldera(t *testing.T) { + slucajevi := map[string]string{ + "../../etc/passwd": "etc_passwd", + "/apsolutno/putanja": "apsolutno_putanja", + "C:\\Windows\\win": "C__Windows_win", + "": "nepoznat", + "...": "nepoznat", + "1CDBD4A6A460": "1CDBD4A6A460", + } + for in, want := range slucajevi { + got := SanitizeSegment(in) + if got != want { + t.Errorf("SanitizeSegment(%q) = %q, ocekivano %q", in, got, want) + } + if strings.ContainsAny(got, `/\`) || got == ".." { + t.Errorf("SanitizeSegment(%q) = %q — pusta izlazak iz foldera", in, got) + } + } + if n := len(SanitizeSegment(strings.Repeat("a", 500))); n > MaxKeyLen { + t.Errorf("kljuc nije ogranicen: %d", n) + } +} + +func TestDeviceKeyPoIPKadNemaMAC(t *testing.T) { + if got := DeviceKey("", "10.20.30.40"); got != "nepoznat-10_20_30_40" { + t.Errorf("kljuc = %q", got) + } + if got := DeviceKey("", ""); got != "nepoznat-bez-adrese" { + t.Errorf("kljuc bez adrese = %q", got) + } + // Lazna izvorna adresa ne sme da napravi putanju. + if got := DeviceKey("", "../../root"); strings.ContainsAny(got, `/\`) { + t.Errorf("kljuc = %q", got) + } +} diff --git a/logger/registry.go b/logger/registry.go new file mode 100644 index 0000000..24b4c77 --- /dev/null +++ b/logger/registry.go @@ -0,0 +1,194 @@ +package logger + +import ( + "sort" + "sync" + "time" +) + +// Device je sve sto pamtimo o jednom uredjaju. Sve je u memoriji i sve je +// ograniceno po velicini - dnevnik ide u fajlove, ne ovde. +type Device struct { + Key string + Name string + MAC string + IP string // adresa koju uredjaj sam prijavi (heartbeat) + SourceIP string // adresa sa koje je paket stvarno stigao + FW string + Serial string + Identified bool // znamo li ko je (zaglavlje/heartbeat) ili samo nagadjamo po IP-u + + FirstSeen time.Time + LastSeen time.Time + LastHeartbeat time.Time + + Heartbeats uint64 + LogPackets uint64 + LogLines uint64 + + IntervalSec int // izmeren razmak izmedju dva heartbeata (0 = jos ne znamo) + LastLine string // poslednja linija dnevnika, samo za pregled +} + +// Registry cuva uredjaje u memoriji. Broj uredjaja je ogranicen jer paket moze +// poslati bilo ko - bez granice bi jedan izvor sa laznim MAC-ovima pojeo RAM. +type Registry struct { + mu sync.RWMutex + devices map[string]*Device + maxDevices int + defaultHb int // pretpostavljeni interval kad uredjaj nikad nije poslao heartbeat +} + +func NewRegistry(maxDevices, defaultHbSec int) *Registry { + if maxDevices <= 0 { + maxDevices = 500 + } + if defaultHbSec <= 0 { + defaultHbSec = 60 + } + return &Registry{ + devices: make(map[string]*Device), + maxDevices: maxDevices, + defaultHb: defaultHbSec, + } +} + +// Update ubacuje paket u registar i vraca kljuc uredjaja (i da li je primljen). +// Paket koji nije prepoznat se ignorise. +func (r *Registry) Update(p Packet, now time.Time) (string, bool) { + if p.Kind == KindGarbage { + return "", false + } + key := DeviceKey(p.MAC, p.SourceIP) + + r.mu.Lock() + defer r.mu.Unlock() + + d, ok := r.devices[key] + if !ok { + if len(r.devices) >= r.maxDevices { + r.evictOldestLocked() + } + d = &Device{Key: key, FirstSeen: now} + r.devices[key] = d + } + + d.LastSeen = now + d.SourceIP = p.SourceIP + if p.Name != "" { + d.Name = p.Name + } + if p.MAC != "" { + d.MAC = p.MAC + } + if p.FW != "" { + d.FW = p.FW + } + if p.Serial != "" { + d.Serial = p.Serial + } + if p.Identified { + d.Identified = true + } + + switch p.Kind { + case KindHeartbeat: + if p.IP != "" { + d.IP = p.IP + } + // Interval ne dolazi u paketu - meri se izmedju dva javljanja. + // Uzima se poslednji razmak, uz clamp da jedan izgubljen paket ili + // restart uredjaja ne razvuku prag na sate. + if !d.LastHeartbeat.IsZero() { + if sec := int(now.Sub(d.LastHeartbeat).Seconds()); sec >= 5 && sec <= 3600 { + d.IntervalSec = sec + } + } + d.LastHeartbeat = now + d.Heartbeats++ + case KindLog: + d.LogPackets++ + d.LogLines += uint64(len(p.Lines)) + if n := len(p.Lines); n > 0 { + d.LastLine = p.Lines[n-1] + } + } + return key, true +} + +// evictOldestLocked izbacuje uredjaj koji se najduze nije javio. Poziva se samo +// kad je registar pun. +func (r *Registry) evictOldestLocked() { + var oldestKey string + var oldest time.Time + for k, d := range r.devices { + if oldestKey == "" || d.LastSeen.Before(oldest) { + oldestKey, oldest = k, d.LastSeen + } + } + if oldestKey != "" { + delete(r.devices, oldestKey) + } +} + +// Interval vraca interval po kom se procenjuje da li je uredjaj tih. +func (r *Registry) Interval(d *Device) int { + if d.IntervalSec > 0 { + return d.IntervalSec + } + return r.defaultHb +} + +// Stale je tacno kad se uredjaj nije javio duze od 3x svog intervala. +func (r *Registry) Stale(d *Device, now time.Time) bool { + limit := time.Duration(3*r.Interval(d)) * time.Second + return now.Sub(d.LastSeen) > limit +} + +// Snapshot vraca kopije uredjaja, sortirane: prvo tihi (da bode oci), pa po +// nazivu/kljucu. +func (r *Registry) Snapshot(now time.Time) []Device { + r.mu.RLock() + out := make([]Device, 0, len(r.devices)) + for _, d := range r.devices { + out = append(out, *d) + } + r.mu.RUnlock() + + sort.Slice(out, func(i, j int) bool { + si, sj := r.Stale(&out[i], now), r.Stale(&out[j], now) + if si != sj { + return si + } + ni, nj := out[i].Name, out[j].Name + if ni == "" { + ni = out[i].Key + } + if nj == "" { + nj = out[j].Key + } + if ni != nj { + return ni < nj + } + return out[i].Key < out[j].Key + }) + return out +} + +// Get vraca kopiju jednog uredjaja. +func (r *Registry) Get(key string) (Device, bool) { + r.mu.RLock() + defer r.mu.RUnlock() + d, ok := r.devices[key] + if !ok { + return Device{}, false + } + return *d, true +} + +// Count vraca broj uredjaja u registru. +func (r *Registry) Count() int { + r.mu.RLock() + defer r.mu.RUnlock() + return len(r.devices) +} diff --git a/logger/registry_test.go b/logger/registry_test.go new file mode 100644 index 0000000..21adca9 --- /dev/null +++ b/logger/registry_test.go @@ -0,0 +1,113 @@ +package logger + +import ( + "fmt" + "testing" + "time" +) + +func hbPaket(mac string) Packet { + return ParsePacket([]byte(fmt.Sprintf( + `{"app":"TERMINIA","name":"TERM","mac":"%s","ip":"192.168.0.20","fw":"0.13.4","hb":1}`, mac)), + "192.168.0.20") +} + +func TestRegistryMeriIntervalHeartbeata(t *testing.T) { + r := NewRegistry(10, 60) + t0 := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + + r.Update(hbPaket("1C:DB:D4:A6:A4:60"), t0) + key, _ := r.Update(hbPaket("1C:DB:D4:A6:A4:60"), t0.Add(30*time.Second)) + + d, ok := r.Get(key) + if !ok { + t.Fatal("uredjaj nije u registru") + } + if d.IntervalSec != 30 { + t.Errorf("interval = %d, ocekivano 30", d.IntervalSec) + } + if d.Heartbeats != 2 { + t.Errorf("heartbeata = %d", d.Heartbeats) + } + // 3x30 s = 90 s praga (poslednji paket je stigao u t0+30 s) + if r.Stale(&d, t0.Add(150*time.Second)) == false { + t.Error("posle 90 s tisine uredjaj mora biti oznacen kao tih") + } + if r.Stale(&d, t0.Add(40*time.Second)) { + t.Error("uredjaj koji se javio pre 10 s ne sme biti tih") + } +} + +func TestRegistryBezHeartbeataKoristiPodrazumevaniInterval(t *testing.T) { + r := NewRegistry(10, 60) + t0 := time.Now() + key, _ := r.Update(ParsePacket([]byte("#TERMINIA T 1C:DB:D4:A6:A4:60 0.13.4\nred\n"), "1.2.3.4"), t0) + d, _ := r.Get(key) + if got := r.Interval(&d); got != 60 { + t.Errorf("interval = %d, ocekivano 60 (podrazumevani)", got) + } + if !r.Stale(&d, t0.Add(200*time.Second)) { + t.Error("posle 180 s mora biti tih") + } +} + +func TestRegistryOgranicavaBrojUredjaja(t *testing.T) { + r := NewRegistry(3, 60) + t0 := time.Now() + for i := 0; i < 50; i++ { + p := ParsePacket([]byte("neka linija\n"), fmt.Sprintf("10.0.0.%d", i)) + r.Update(p, t0.Add(time.Duration(i)*time.Second)) + } + if n := r.Count(); n != 3 { + t.Errorf("uredjaja u registru = %d, ocekivano 3 (granica)", n) + } + // Najstariji mora ispasti, najnoviji ostati. + if _, ok := r.Get("nepoznat-10_0_0_49"); !ok { + t.Error("poslednji uredjaj mora ostati u registru") + } + if _, ok := r.Get("nepoznat-10_0_0_0"); ok { + t.Error("najstariji uredjaj je morao biti izbacen") + } +} + +func TestRegistryZadrzavaIdentitetIzZaglavlja(t *testing.T) { + r := NewRegistry(10, 60) + t0 := time.Now() + key, _ := r.Update(ParsePacket([]byte("#TERMINIA TERM-001 1C:DB:D4:A6:A4:60 0.13.4\nprvi red\n"), "1.2.3.4"), t0) + // Kasniji paket bez zaglavlja sa istim MAC-om ne postoji, ali heartbeat + // dopunjuje IP i serijski. + r.Update(hbPaket("1C:DB:D4:A6:A4:60"), t0.Add(time.Second)) + + d, _ := r.Get(key) + if !d.Identified || d.MAC != "1C:DB:D4:A6:A4:60" || d.IP != "192.168.0.20" { + t.Errorf("uredjaj = %+v", d) + } + if d.LogLines != 1 || d.LastLine != "prvi red" { + t.Errorf("linije = %d / %q", d.LogLines, d.LastLine) + } +} + +func TestRegistrySnapshotStavljaTiheNaVrh(t *testing.T) { + r := NewRegistry(10, 60) + t0 := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + r.Update(hbPaket("AA:AA:AA:AA:AA:AA"), t0.Add(-time.Hour)) + r.Update(hbPaket("BB:BB:BB:BB:BB:BB"), t0) + + snap := r.Snapshot(t0) + if len(snap) != 2 { + t.Fatalf("snapshot = %d", len(snap)) + } + if snap[0].MAC != "AA:AA:AA:AA:AA:AA" { + t.Errorf("tih uredjaj mora biti prvi, dobijeno %s", snap[0].MAC) + } +} + +func TestRegistryIgnorisSmece(t *testing.T) { + r := NewRegistry(10, 60) + if _, ok := r.Update(ParsePacket([]byte("{"), "1.2.3.4"), time.Now()); ok { + t.Error("smece ne sme uci u registar") + } + if r.Count() != 0 { + t.Errorf("registar nije prazan: %d", r.Count()) + } +} diff --git a/logger/service.go b/logger/service.go new file mode 100644 index 0000000..76abb2b --- /dev/null +++ b/logger/service.go @@ -0,0 +1,259 @@ +package logger + +import ( + "context" + "errors" + "fmt" + "log" + "net" + "net/http" + "os" + "strconv" + "sync/atomic" + "time" +) + +// Config je citav podesivi deo prijemnika. Sve dolazi iz env promenljivih +// (systemd EnvironmentFile), nista nije zakucano u kodu. +type Config struct { + UDPHost string // prazno = sve adrese + UDPPort int + HTTPHost string + HTTPPort int + Dir string + Token string // ako je postavljen, web trazi ?t= + KeepDays int + MaxFileMB int + MaxDeviceMB int + MaxTotalMB int + MaxDevices int + HbDefaultSec int +} + +// DefaultConfig su podrazumevane vrednosti. +// +// UDP 46465 namerno NIJE 46464: na tom portu uredjaji slusaju discovery, pa bi +// prijemnik na istoj mrezi ulazio u sukob sa njima. +// +// Web sluza samo na 127.0.0.1 jer dnevnik sadrzi BROJEVE KARATA - napolje se +// pusta iskljucivo kroz SSH tunel ili nginx sa autentikacijom. +func DefaultConfig() Config { + return Config{ + UDPHost: "", + UDPPort: 46465, + HTTPHost: "127.0.0.1", + HTTPPort: 8081, + Dir: "./data/logger", + KeepDays: 14, + MaxFileMB: 10, + MaxDeviceMB: 50, + MaxTotalMB: 1024, + MaxDevices: 500, + HbDefaultSec: 60, + } +} + +// ConfigFromEnv cita podesavanja preko podrazumevanih vrednosti. +func ConfigFromEnv() (Config, error) { + c := DefaultConfig() + var err error + str := func(key string, dst *string) { + if v, ok := os.LookupEnv(key); ok { + *dst = v + } + } + num := func(key string, dst *int, min, max int) { + v, ok := os.LookupEnv(key) + if !ok || v == "" { + return + } + n, e := strconv.Atoi(v) + if e != nil || n < min || n > max { + err = fmt.Errorf("%s: neispravna vrednost %q (ocekuje se broj %d-%d)", key, v, min, max) + return + } + *dst = n + } + str("LOGGER_UDP_HOST", &c.UDPHost) + str("LOGGER_HTTP_HOST", &c.HTTPHost) + str("LOGGER_DIR", &c.Dir) + str("LOGGER_TOKEN", &c.Token) + num("LOGGER_UDP_PORT", &c.UDPPort, 1, 65535) + num("LOGGER_HTTP_PORT", &c.HTTPPort, 1, 65535) + num("LOGGER_KEEP_DAYS", &c.KeepDays, 1, 3650) + num("LOGGER_MAX_FILE_MB", &c.MaxFileMB, 1, 10240) + num("LOGGER_MAX_DEVICE_MB", &c.MaxDeviceMB, 1, 102400) + num("LOGGER_MAX_TOTAL_MB", &c.MaxTotalMB, 1, 1024000) + num("LOGGER_MAX_DEVICES", &c.MaxDevices, 1, 100000) + num("LOGGER_HB_DEFAULT_SEC", &c.HbDefaultSec, 5, 86400) + if err != nil { + return c, err + } + if c.MaxDeviceMB > c.MaxTotalMB { + return c, fmt.Errorf("LOGGER_MAX_DEVICE_MB (%d) ne sme biti vece od LOGGER_MAX_TOTAL_MB (%d)", + c.MaxDeviceMB, c.MaxTotalMB) + } + return c, nil +} + +// Stats su brojaci prijemnika za prikaz na stranici. +type Stats struct { + Packets atomic.Uint64 + Heartbeat atomic.Uint64 + LogPkts atomic.Uint64 + Garbage atomic.Uint64 + WriteErrs atomic.Uint64 +} + +// Service spaja UDP prijem, registar uredjaja, pisanje u fajlove i web. +type Service struct { + cfg Config + reg *Registry + store *Store + stats *Stats + started time.Time + now func() time.Time // zamenjivo u testovima +} + +func New(cfg Config) (*Service, error) { + st, err := NewStore(StoreConfig{ + Dir: cfg.Dir, + MaxFileBytes: int64(cfg.MaxFileMB) << 20, + MaxDeviceBytes: int64(cfg.MaxDeviceMB) << 20, + MaxTotalBytes: int64(cfg.MaxTotalMB) << 20, + KeepDays: cfg.KeepDays, + }) + if err != nil { + return nil, err + } + return &Service{ + cfg: cfg, + reg: NewRegistry(cfg.MaxDevices, cfg.HbDefaultSec), + store: st, + stats: &Stats{}, + started: time.Now(), + now: time.Now, + }, nil +} + +// Handle obradjuje jedan primljen paket. Izdvojeno iz petlje da bi moglo da se +// testira bez mreze. +func (s *Service) Handle(data []byte, sourceIP string) { + s.stats.Packets.Add(1) + p := ParsePacket(data, sourceIP) + if p.Kind == KindGarbage { + s.stats.Garbage.Add(1) + return + } + now := s.now() + key, ok := s.reg.Update(p, now) + if !ok { + return + } + switch p.Kind { + case KindHeartbeat: + s.stats.Heartbeat.Add(1) + case KindLog: + s.stats.LogPkts.Add(1) + if err := s.store.Write(key, now, p.Lines); err != nil { + // Greska pisanja ne sme da obori prijem: paket se izgubi, a + // brojac gresaka se vidi na stranici. + s.stats.WriteErrs.Add(1) + log.Printf("prijemnik: %v", err) + } + } +} + +// Run pokrece UDP prijem, web i odrzavanje. Vraca se kad se kontekst otkaze. +func (s *Service) Run(ctx context.Context) error { + udpAddr := net.JoinHostPort(s.cfg.UDPHost, strconv.Itoa(s.cfg.UDPPort)) + pc, err := net.ListenPacket("udp", udpAddr) + if err != nil { + return fmt.Errorf("UDP %s: %w", udpAddr, err) + } + defer pc.Close() + + httpAddr := net.JoinHostPort(s.cfg.HTTPHost, strconv.Itoa(s.cfg.HTTPPort)) + srv := &http.Server{ + Addr: httpAddr, + Handler: s.Handler(), + ReadHeaderTimeout: 5 * time.Second, + } + ln, err := net.Listen("tcp", httpAddr) + if err != nil { + return fmt.Errorf("HTTP %s: %w", httpAddr, err) + } + + log.Printf("prijemnik dijagnostike: UDP %s, web http://%s, dnevnici u %s", + udpAddr, httpAddr, s.cfg.Dir) + if s.cfg.Token == "" && s.cfg.HTTPHost != "127.0.0.1" && s.cfg.HTTPHost != "localhost" { + log.Printf("UPOZORENJE: web je otvoren ka mrezi bez tokena (LOGGER_TOKEN), a dnevnik sadrzi brojeve karata") + } + + errCh := make(chan error, 1) + go func() { + if err := srv.Serve(ln); err != nil && !errors.Is(err, http.ErrServerClosed) { + errCh <- err + } + }() + go s.readLoop(pc, errCh) + go s.maintenance(ctx) + + select { + case <-ctx.Done(): + case err = <-errCh: + } + + shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + srv.Shutdown(shutCtx) + pc.Close() + s.store.Close() + return err +} + +func (s *Service) readLoop(pc net.PacketConn, errCh chan<- error) { + buf := make([]byte, MaxPacketBytes) + for { + n, addr, err := pc.ReadFrom(buf) + if err != nil { + // Zatvoren soket je uredan kraj; sve ostalo je greska mreze + // koju ne pokusavamo da lecimo. + if errors.Is(err, net.ErrClosed) { + return + } + errCh <- fmt.Errorf("UDP prijem: %w", err) + return + } + src := "" + if ua, ok := addr.(*net.UDPAddr); ok && ua.IP != nil { + src = ua.IP.String() + } + s.Handle(buf[:n], src) + } +} + +// maintenance povremeno cisti stare dnevnike i zatvara mirne fajlove. +func (s *Service) maintenance(ctx context.Context) { + if n, err := s.store.Cleanup(s.now()); err != nil { + log.Printf("prijemnik: ciscenje: %v", err) + } else if n > 0 { + log.Printf("prijemnik: obrisano %d starih dnevnika", n) + } + t := time.NewTicker(time.Minute) + defer t.Stop() + for { + select { + case <-ctx.Done(): + return + case <-t.C: + now := s.now() + s.store.CloseIdle(now, 5*time.Minute) + if n, err := s.store.Cleanup(now); err != nil { + log.Printf("prijemnik: ciscenje: %v", err) + } else if n > 0 { + log.Printf("prijemnik: obrisano %d starih dnevnika", n) + } + } + } +} diff --git a/logger/service_test.go b/logger/service_test.go new file mode 100644 index 0000000..8fd8f1d --- /dev/null +++ b/logger/service_test.go @@ -0,0 +1,200 @@ +package logger + +import ( + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" +) + +func testService(t *testing.T, prilagodi func(*Config)) *Service { + t.Helper() + cfg := DefaultConfig() + cfg.Dir = t.TempDir() + if prilagodi != nil { + prilagodi(&cfg) + } + s, err := New(cfg) + if err != nil { + t.Fatalf("New: %v", err) + } + t.Cleanup(s.store.Close) + return s +} + +func TestServiceHandleUpisujeDnevnikIPamtiHeartbeat(t *testing.T) { + s := testService(t, nil) + t0 := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + s.now = func() time.Time { return t0 } + + s.Handle([]byte(`{"app":"TERMINIA","name":"TERM-001","mac":"1C:DB:D4:A6:A4:60","ip":"192.168.0.20","fw":"0.13.4","hb":1}`), "192.168.0.20") + s.Handle([]byte("#TERMINIA TERM-001 1C:DB:D4:A6:A4:60 0.13.4\n[stampac] javio se\n"), "192.168.0.20") + s.Handle([]byte{0x00, 0x01, 0x02, 0xff, 0xfe}, "9.9.9.9") // binarno smece + + if got := s.stats.Heartbeat.Load(); got != 1 { + t.Errorf("heartbeata = %d", got) + } + if got := s.stats.LogPkts.Load(); got != 1 { + t.Errorf("paketa dnevnika = %d", got) + } + if got := s.stats.WriteErrs.Load(); got != 0 { + t.Errorf("gresaka upisa = %d", got) + } + lines, _ := s.store.Tail("1CDBD4A6A460", 10, "") + if len(lines) != 1 || !strings.Contains(lines[0], "[stampac] javio se") { + t.Errorf("dnevnik = %q", lines) + } +} + +func TestServiceStranicaUredjaja(t *testing.T) { + s := testService(t, nil) + t0 := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + s.now = func() time.Time { return t0 } + s.Handle([]byte("#TERMINIA TERM-001 1C:DB:D4:A6:A4:60 0.13.4\n[stampac] javio se\n[kapija] prolaz\n"), "192.168.0.20") + h := s.Handler() + + rec := httptest.NewRecorder() + h.ServeHTTP(rec, httptest.NewRequest("GET", "/", nil)) + if rec.Code != http.StatusOK { + t.Fatalf("/ = %d", rec.Code) + } + body := rec.Body.String() + if !strings.Contains(body, "TERM-001") || !strings.Contains(body, "1C:DB:D4:A6:A4:60") { + t.Errorf("spisak ne prikazuje uredjaj: %s", body) + } + + rec = httptest.NewRecorder() + h.ServeHTTP(rec, httptest.NewRequest("GET", "/uredjaj?k=1CDBD4A6A460&q=stampac", nil)) + if rec.Code != http.StatusOK { + t.Fatalf("/uredjaj = %d", rec.Code) + } + body = rec.Body.String() + if !strings.Contains(body, "[stampac] javio se") { + t.Error("pretraga ne vraca ocekivanu liniju") + } + if strings.Contains(body, "[kapija] prolaz") { + t.Error("pretraga vraca i linije koje ne odgovaraju upitu") + } +} + +func TestServiceNepoznatUredjajJeOznacen(t *testing.T) { + s := testService(t, nil) + s.Handle([]byte("goli tekst bez zaglavlja\n"), "192.168.0.77") + + rec := httptest.NewRecorder() + s.Handler().ServeHTTP(rec, httptest.NewRequest("GET", "/", nil)) + body := rec.Body.String() + if !strings.Contains(body, "nepoznat uređaj") { + t.Errorf("uredjaj bez zaglavlja nije oznacen kao nepoznat:\n%s", body) + } + if !strings.Contains(body, "192.168.0.77") { + t.Error("nije prikazana izvorna adresa nepoznatog uredjaja") + } +} + +func TestServiceTihUredjajJeCrven(t *testing.T) { + s := testService(t, nil) + t0 := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + s.now = func() time.Time { return t0 } + s.Handle([]byte(`{"name":"TERM-001","mac":"1C:DB:D4:A6:A4:60","hb":1}`), "192.168.0.20") + + s.now = func() time.Time { return t0.Add(10 * time.Minute) } // 3x60 s je davno proslo + rec := httptest.NewRecorder() + s.Handler().ServeHTTP(rec, httptest.NewRequest("GET", "/", nil)) + if !strings.Contains(rec.Body.String(), `class="tih"`) { + t.Error("tih uredjaj nije oznacen crvenim redom") + } +} + +func TestServiceHTMLEscapeIzPaketa(t *testing.T) { + // Sadrzaj paketa je tudji tekst - ne sme da postane HTML. + s := testService(t, nil) + s.Handle([]byte("#TERMINIA 1C:DB:D4:A6:A4:60 0.13.4\n\n"), "1.2.3.4") + + for _, url := range []string{"/", "/uredjaj?k=1CDBD4A6A460"} { + rec := httptest.NewRecorder() + s.Handler().ServeHTTP(rec, httptest.NewRequest("GET", url, nil)) + if strings.Contains(rec.Body.String(), "") || + strings.Contains(rec.Body.String(), "= s.cfg.MaxFileBytes { + s.closeLocked(key) // sledeci upis otvara naredni nastavak + } + return nil +} + +func (s *Store) fileForLocked(key string, when time.Time) (*openFile, error) { + day := when.Format(dayLayout) + if of, ok := s.open[key]; ok { + if of.day == day && of.size < s.cfg.MaxFileBytes { + return of, nil + } + s.closeLocked(key) + } + if len(s.open) >= maxOpenFiles { + s.closeIdlestLocked() + } + dir := filepath.Join(s.cfg.Dir, key) + if err := os.MkdirAll(dir, 0o750); err != nil { + return nil, fmt.Errorf("store: kreiranje %s: %w", dir, err) + } + path, size, seq, err := nextPath(dir, day, s.cfg.MaxFileBytes) + if err != nil { + return nil, err + } + f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o640) + if err != nil { + return nil, fmt.Errorf("store: otvaranje %s: %w", path, err) + } + of := &openFile{f: f, w: bufio.NewWriter(f), path: path, day: day, seq: seq, size: size, used: when} + s.open[key] = of + return of, nil +} + +// nextPath bira u koji fajl se pise za dati dan: osnovni ako jos nije pun, +// inace prvi slobodan nastavak. +func nextPath(dir, day string, maxBytes int64) (string, int64, int, error) { + base := filepath.Join(dir, day+logExt) + st, err := os.Stat(base) + switch { + case os.IsNotExist(err): + return base, 0, 0, nil + case err != nil: + return "", 0, 0, fmt.Errorf("store: stat %s: %w", base, err) + case st.Size() < maxBytes: + return base, st.Size(), 0, nil + } + for seq := 1; seq < 1000; seq++ { + p := fmt.Sprintf("%s.%03d", base, seq) + st, err := os.Stat(p) + if os.IsNotExist(err) { + return p, 0, seq, nil + } + if err != nil { + return "", 0, 0, fmt.Errorf("store: stat %s: %w", p, err) + } + if st.Size() < maxBytes { + return p, st.Size(), seq, nil + } + } + return "", 0, 0, fmt.Errorf("store: previse nastavaka za %s", base) +} + +func (s *Store) closeLocked(key string) { + if of, ok := s.open[key]; ok { + of.w.Flush() + of.f.Close() + delete(s.open, key) + } +} + +func (s *Store) closeIdlestLocked() { + var key string + var oldest time.Time + for k, of := range s.open { + if key == "" || of.used.Before(oldest) { + key, oldest = k, of.used + } + } + if key != "" { + s.closeLocked(key) + } +} + +// CloseIdle zatvara fajlove u koje se duze nista ne pise (da se ne drze +// otvoreni deskriptori i da rotacija/brisanje ne nailazi na njih). +func (s *Store) CloseIdle(now time.Time, idle time.Duration) { + s.mu.Lock() + defer s.mu.Unlock() + for k, of := range s.open { + if now.Sub(of.used) > idle { + s.closeLocked(k) + } + } +} + +// Close zatvara sve otvorene fajlove. +func (s *Store) Close() { + s.mu.Lock() + defer s.mu.Unlock() + for k := range s.open { + s.closeLocked(k) + } +} + +type logFile struct { + key string + path string + name string + day string + size int64 +} + +// list nabraja sve dnevnicke fajlove, sortirane od najstarijeg ka najnovijem. +func (s *Store) list() ([]logFile, error) { + dirs, err := os.ReadDir(s.cfg.Dir) + if err != nil { + return nil, fmt.Errorf("store: citanje %s: %w", s.cfg.Dir, err) + } + var out []logFile + for _, d := range dirs { + if !d.IsDir() { + continue + } + key := d.Name() + files, err := os.ReadDir(filepath.Join(s.cfg.Dir, key)) + if err != nil { + continue // uredjaj bez foldera nije razlog za pad prijemnika + } + for _, f := range files { + name := f.Name() + if f.IsDir() || !strings.Contains(name, logExt) { + continue + } + day := name + if i := strings.Index(name, logExt); i > 0 { + day = name[:i] + } + if _, err := time.Parse(dayLayout, day); err != nil { + continue // tudji fajl u folderu - ne diramo ga + } + info, err := f.Info() + if err != nil { + continue + } + out = append(out, logFile{ + key: key, + path: filepath.Join(s.cfg.Dir, key, name), + name: name, + day: day, + size: info.Size(), + }) + } + } + sort.Slice(out, func(i, j int) bool { + if out[i].name != out[j].name { + return out[i].name < out[j].name + } + return out[i].key < out[j].key + }) + return out, nil +} + +// Cleanup sprovodi tri granice redom: starost, kvota po uredjaju, ukupna +// kvota. Vraca broj obrisanih fajlova. +// +// Fajl koji je trenutno otvoren za pisanje se NIKAD ne brise: na Linuxu +// brisanje otvorenog fajla ne oslobodi prostor dok se ne zatvori, a upisi +// posle toga tiho odlaze u nepostojeci fajl. +func (s *Store) Cleanup(now time.Time) (int, error) { + files, err := s.list() + if err != nil { + return 0, err + } + s.mu.Lock() + busy := make(map[string]bool, len(s.open)) + for _, of := range s.open { + busy[of.path] = true + } + s.mu.Unlock() + + deleted := 0 + remove := func(f logFile) bool { + if busy[f.path] { + return false + } + if err := os.Remove(f.path); err != nil { + return false + } + deleted++ + return true + } + + // 1) starost + cutoff := now.AddDate(0, 0, -s.cfg.KeepDays) + kept := files[:0] + for _, f := range files { + d, err := time.Parse(dayLayout, f.day) + if err == nil && d.Before(cutoff) && remove(f) { + continue + } + kept = append(kept, f) + } + files = kept + + // 2) kvota po uredjaju + perDevice := map[string]int64{} + for _, f := range files { + perDevice[f.key] += f.size + } + kept = files[:0] + for _, f := range files { + if perDevice[f.key] > s.cfg.MaxDeviceBytes && remove(f) { + perDevice[f.key] -= f.size + continue + } + kept = append(kept, f) + } + files = kept + + // 3) ukupna kvota + var total int64 + for _, f := range files { + total += f.size + } + for _, f := range files { + if total <= s.cfg.MaxTotalBytes { + break + } + if remove(f) { + total -= f.size + } + } + + s.pruneEmptyDirs() + return deleted, nil +} + +// pruneEmptyDirs brise foldere uredjaja koji su ostali bez ijednog dnevnika. +func (s *Store) pruneEmptyDirs() { + dirs, err := os.ReadDir(s.cfg.Dir) + if err != nil { + return + } + for _, d := range dirs { + if !d.IsDir() { + continue + } + p := filepath.Join(s.cfg.Dir, d.Name()) + if entries, err := os.ReadDir(p); err == nil && len(entries) == 0 { + os.Remove(p) + } + } +} + +// Usage vraca ukupno zauzece na disku (bajtova) i broj fajlova. +func (s *Store) Usage() (int64, int) { + files, err := s.list() + if err != nil { + return 0, 0 + } + var total int64 + for _, f := range files { + total += f.size + } + return total, len(files) +} + +// Tail vraca poslednjih n linija dnevnika jednog uredjaja, opciono filtriranih +// po tekstu (bez razlike u velicini slova). Redosled je hronoloski. +func (s *Store) Tail(key string, n int, query string) ([]string, error) { + key = SanitizeSegment(key) + if n <= 0 || n > 5000 { + n = 200 + } + // Bafer se prazni da bi i linije primljene pre sekund bile vidljive. + s.mu.Lock() + if of, ok := s.open[key]; ok { + of.w.Flush() + } + s.mu.Unlock() + + files, err := s.list() + if err != nil { + return nil, err + } + var mine []logFile + for _, f := range files { + if f.key == key { + mine = append(mine, f) + } + } + q := strings.ToLower(strings.TrimSpace(query)) + + out := make([]string, 0, n) + // od najnovijeg fajla unazad, dok se ne skupi n linija + for i := len(mine) - 1; i >= 0 && len(out) < n; i-- { + lines, err := readTailLines(mine[i].path, tailReadCap) + if err != nil { + continue + } + for j := len(lines) - 1; j >= 0 && len(out) < n; j-- { + if q != "" && !strings.Contains(strings.ToLower(lines[j]), q) { + continue + } + out = append(out, lines[j]) + } + } + // out je od najnovije ka najstarijoj - okrecemo + for i, j := 0, len(out)-1; i < j; i, j = i+1, j-1 { + out[i], out[j] = out[j], out[i] + } + return out, nil +} + +// readTailLines cita poslednjih najvise cap bajtova fajla i vraca cele linije. +func readTailLines(path string, cap int64) ([]string, error) { + f, err := os.Open(path) + if err != nil { + return nil, err + } + defer f.Close() + st, err := f.Stat() + if err != nil { + return nil, err + } + size := st.Size() + off := int64(0) + if size > cap { + off = size - cap + } + if _, err := f.Seek(off, io.SeekStart); err != nil { + return nil, err + } + buf := make([]byte, size-off) + if _, err := io.ReadFull(f, buf); err != nil && err != io.ErrUnexpectedEOF { + return nil, err + } + text := string(buf) + if off > 0 { + // prva linija je verovatno presecena na pola - odbacujemo je + if i := strings.IndexByte(text, '\n'); i >= 0 { + text = text[i+1:] + } + } + lines := strings.Split(text, "\n") + out := lines[:0] + for _, ln := range lines { + if strings.TrimSpace(ln) != "" { + out = append(out, ln) + } + } + return out, nil +} + +// Days vraca datume za koje postoji dnevnik datog uredjaja (najnoviji prvi). +func (s *Store) Days(key string) []string { + key = SanitizeSegment(key) + files, err := s.list() + if err != nil { + return nil + } + seen := map[string]bool{} + var out []string + for i := len(files) - 1; i >= 0; i-- { + if files[i].key != key || seen[files[i].day] { + continue + } + seen[files[i].day] = true + out = append(out, files[i].day) + } + return out +} diff --git a/logger/store_test.go b/logger/store_test.go new file mode 100644 index 0000000..f682c27 --- /dev/null +++ b/logger/store_test.go @@ -0,0 +1,245 @@ +package logger + +import ( + "fmt" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +func testStore(t *testing.T, cfg StoreConfig) *Store { + t.Helper() + cfg.Dir = t.TempDir() + s, err := NewStore(cfg) + if err != nil { + t.Fatalf("NewStore: %v", err) + } + t.Cleanup(s.Close) + return s +} + +func TestStorePiseIVracaLinije(t *testing.T) { + s := testStore(t, StoreConfig{}) + t0 := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + + if err := s.Write("1CDBD4A6A460", t0, []string{"prva", "druga"}); err != nil { + t.Fatalf("Write: %v", err) + } + if err := s.Write("1CDBD4A6A460", t0.Add(time.Second), []string{"treca"}); err != nil { + t.Fatalf("Write: %v", err) + } + + lines, err := s.Tail("1CDBD4A6A460", 10, "") + if err != nil { + t.Fatalf("Tail: %v", err) + } + if len(lines) != 3 { + t.Fatalf("linija = %d: %q", len(lines), lines) + } + if !strings.HasSuffix(lines[0], "prva") || !strings.HasSuffix(lines[2], "treca") { + t.Errorf("redosled nije hronoloski: %q", lines) + } + if !strings.HasPrefix(lines[0], "2026-07-25 12:00:00") { + t.Errorf("nedostaje vremenska oznaka: %q", lines[0]) + } +} + +func TestStorePretraga(t *testing.T) { + s := testStore(t, StoreConfig{}) + t0 := time.Now() + s.Write("dev", t0, []string{"[stampac] greska", "[kapija] prolaz", "[stampac] ok"}) + + lines, _ := s.Tail("dev", 10, "STAMPAC") + if len(lines) != 2 { + t.Fatalf("pretraga = %d linija: %q", len(lines), lines) + } + lines, _ = s.Tail("dev", 10, "nepostojece") + if len(lines) != 0 { + t.Errorf("ocekivano 0 linija, dobijeno %q", lines) + } +} + +func TestStoreOgranicavaBrojLinija(t *testing.T) { + s := testStore(t, StoreConfig{}) + t0 := time.Now() + for i := 0; i < 100; i++ { + s.Write("dev", t0, []string{fmt.Sprintf("linija %d", i)}) + } + lines, _ := s.Tail("dev", 10, "") + if len(lines) != 10 { + t.Fatalf("linija = %d", len(lines)) + } + if !strings.HasSuffix(lines[9], "linija 99") { + t.Errorf("poslednja linija = %q", lines[9]) + } +} + +func TestStoreRotacijaPoVelicini(t *testing.T) { + s := testStore(t, StoreConfig{MaxFileBytes: 200}) + t0 := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + for i := 0; i < 30; i++ { + if err := s.Write("dev", t0, []string{fmt.Sprintf("popunjavanje %02d", i)}); err != nil { + t.Fatalf("Write: %v", err) + } + } + files, err := os.ReadDir(filepath.Join(s.cfg.Dir, "dev")) + if err != nil { + t.Fatalf("ReadDir: %v", err) + } + if len(files) < 2 { + t.Fatalf("ocekivana rotacija, fajlova = %d", len(files)) + } + var imena []string + for _, f := range files { + imena = append(imena, f.Name()) + if info, _ := f.Info(); info.Size() > 400 { + t.Errorf("fajl %s je %d B — granica se ne postuje", f.Name(), info.Size()) + } + } + if imena[0] != "2026-07-25.log" || imena[1] != "2026-07-25.log.001" { + t.Errorf("imena fajlova = %q", imena) + } + // Sve linije se i dalje citaju, hronoloski. + lines, _ := s.Tail("dev", 100, "") + if len(lines) != 30 { + t.Errorf("procitano %d linija, ocekivano 30", len(lines)) + } + if !strings.HasSuffix(lines[0], "popunjavanje 00") || !strings.HasSuffix(lines[29], "popunjavanje 29") { + t.Errorf("redosled kroz rotirane fajlove nije dobar: %q ... %q", lines[0], lines[29]) + } +} + +func TestStoreNoviDanNoviFajl(t *testing.T) { + s := testStore(t, StoreConfig{}) + t0 := time.Date(2026, 7, 25, 23, 59, 0, 0, time.UTC) + s.Write("dev", t0, []string{"juce"}) + s.Write("dev", t0.Add(2*time.Minute), []string{"danas"}) + + dani := s.Days("dev") + if len(dani) != 2 || dani[0] != "2026-07-26" || dani[1] != "2026-07-25" { + t.Errorf("dani = %q", dani) + } +} + +func TestCleanupBriseStareDnevnike(t *testing.T) { + s := testStore(t, StoreConfig{KeepDays: 3}) + now := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + for i := 0; i < 10; i++ { + s.Write("dev", now.AddDate(0, 0, -i), []string{"nesto"}) + } + s.Close() // da ciscenje sme da dira i tekuci fajl + + n, err := s.Cleanup(now) + if err != nil { + t.Fatalf("Cleanup: %v", err) + } + if n == 0 { + t.Fatal("nista nije obrisano") + } + dani := s.Days("dev") + for _, d := range dani { + dd, _ := time.Parse(dayLayout, d) + if dd.Before(now.AddDate(0, 0, -3)) { + t.Errorf("ostao stari dnevnik %s", d) + } + } +} + +func TestCleanupPostujeKvotuPoUredjaju(t *testing.T) { + s := testStore(t, StoreConfig{MaxFileBytes: 300, MaxDeviceBytes: 600, KeepDays: 3650}) + base := time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC) + for i := 0; i < 20; i++ { + s.Write("dev", base.AddDate(0, 0, i), []string{strings.Repeat("x", 200)}) + } + s.Close() + + if _, err := s.Cleanup(base.AddDate(0, 0, 20)); err != nil { + t.Fatalf("Cleanup: %v", err) + } + total, _ := s.Usage() + if total > 600 { + t.Errorf("posle ciscenja %d B, kvota je 600 B", total) + } + if total == 0 { + t.Error("ciscenje je obrisalo sve — najnoviji dnevnik mora ostati") + } +} + +func TestCleanupPostujeUkupnuKvotu(t *testing.T) { + s := testStore(t, StoreConfig{MaxFileBytes: 300, MaxDeviceBytes: 1 << 20, MaxTotalBytes: 800, KeepDays: 3650}) + base := time.Date(2026, 7, 1, 12, 0, 0, 0, time.UTC) + for i := 0; i < 10; i++ { + s.Write(fmt.Sprintf("dev%d", i), base.AddDate(0, 0, i), []string{strings.Repeat("y", 200)}) + } + s.Close() + + if _, err := s.Cleanup(base.AddDate(0, 0, 10)); err != nil { + t.Fatalf("Cleanup: %v", err) + } + total, _ := s.Usage() + if total > 800 { + t.Errorf("posle ciscenja %d B, ukupna kvota je 800 B", total) + } +} + +func TestCleanupNeBriseOtvoreniFajl(t *testing.T) { + // Fajl u koji se trenutno pise ne sme da nestane: na Linuxu bi upisi + // posle brisanja tiho odlazili u nepostojeci fajl. + s := testStore(t, StoreConfig{MaxTotalBytes: 1, MaxDeviceBytes: 1, KeepDays: 1}) + now := time.Date(2026, 7, 25, 12, 0, 0, 0, time.UTC) + s.Write("dev", now, []string{"vazna linija"}) + + if _, err := s.Cleanup(now); err != nil { + t.Fatalf("Cleanup: %v", err) + } + if _, err := os.Stat(filepath.Join(s.cfg.Dir, "dev", "2026-07-25.log")); err != nil { + t.Errorf("otvoreni fajl je obrisan: %v", err) + } +} + +func TestStoreOdbijaIzlazakIzFoldera(t *testing.T) { + s := testStore(t, StoreConfig{}) + now := time.Now() + if err := s.Write("../../../pobegao", now, []string{"linija"}); err != nil { + t.Fatalf("Write: %v", err) + } + // Sve mora ostati unutar Dir. + var vanFoldera []string + filepath.Walk(filepath.Dir(s.cfg.Dir), func(p string, info os.FileInfo, err error) error { + if err != nil || info.IsDir() { + return nil + } + if !strings.HasPrefix(p, s.cfg.Dir) { + vanFoldera = append(vanFoldera, p) + } + return nil + }) + if len(vanFoldera) > 0 { + t.Errorf("fajlovi napravljeni van foldera dnevnika: %q", vanFoldera) + } +} + +func TestTailPreseceneLinijeNaGranici(t *testing.T) { + s := testStore(t, StoreConfig{}) + now := time.Now() + s.Write("dev", now, []string{"cela linija jedan", "cela linija dva"}) + lines, _ := s.Tail("dev", 5, "") + for _, ln := range lines { + if !strings.Contains(ln, "cela linija") { + t.Errorf("krnja linija u ispisu: %q", ln) + } + } +} + +func TestTailNepostojeciUredjaj(t *testing.T) { + s := testStore(t, StoreConfig{}) + lines, err := s.Tail("nema-ga", 10, "") + if err != nil { + t.Fatalf("Tail: %v", err) + } + if len(lines) != 0 { + t.Errorf("linije = %q", lines) + } +} diff --git a/logger/web.go b/logger/web.go new file mode 100644 index 0000000..3de2036 --- /dev/null +++ b/logger/web.go @@ -0,0 +1,365 @@ +package logger + +import ( + "encoding/json" + "fmt" + "html/template" + "net/http" + "strconv" + "strings" + "time" +) + +// Handler vraca ceo web sloj prijemnika (bez baze, bez spoljnih biblioteka). +func (s *Service) Handler() http.Handler { + mux := http.NewServeMux() + mux.HandleFunc("GET /", s.pageIndex) + mux.HandleFunc("GET /uredjaj", s.pageDevice) + mux.HandleFunc("GET /dnevnik.txt", s.pageRaw) + mux.HandleFunc("GET /api/uredjaji", s.apiDevices) + return s.withToken(mux) +} + +// withToken je minimalna zastita kad web nije samo na 127.0.0.1. Token se +// prihvata iz ?t= ili iz zaglavlja X-Token; posle prvog uspesnog upita pamti se +// u kolacicu da linkovi ostanu kratki. +// +// Ovo NIJE zamena za nginx + TLS + lozinku - dnevnik sadrzi brojeve karata. +func (s *Service) withToken(next http.Handler) http.Handler { + if s.cfg.Token == "" { + return next + } + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + ok := r.URL.Query().Get("t") == s.cfg.Token || r.Header.Get("X-Token") == s.cfg.Token + if !ok { + if c, err := r.Cookie("logger_token"); err == nil && c.Value == s.cfg.Token { + ok = true + } + } + if !ok { + http.Error(w, "Pristup odbijen: nedostaje token.", http.StatusForbidden) + return + } + http.SetCookie(w, &http.Cookie{ + Name: "logger_token", Value: s.cfg.Token, Path: "/", + HttpOnly: true, SameSite: http.SameSiteLaxMode, MaxAge: 12 * 3600, + }) + next.ServeHTTP(w, r) + }) +} + +type deviceView struct { + Key string + Naziv string + MAC string + IP string + FW string + Serijski string + Poznat bool + Tih bool + Vidjen string + Heartbeat string + Interval int + Paketa uint64 + Linija uint64 + Poslednja string + IzvornaIP string + PrviPutStr string +} + +type indexData struct { + Uredjaji []deviceView + Ukupno int + Paketa uint64 + Heartbeata uint64 + Dnevnika uint64 + Smece uint64 + Greske uint64 + DiskMB string + Fajlova int + UDP string + Port int + Dir string + Vreme string + Token string +} + +type deviceData struct { + U deviceView + Linije []string + Q string + N int + Dani []string + Vreme string + Token string + Nadjen bool +} + +func (s *Service) view(d Device, now time.Time) deviceView { + naziv := d.Name + if naziv == "" { + naziv = "(bez naziva)" + } + hb := "nikad" + if !d.LastHeartbeat.IsZero() { + hb = humanAgo(now.Sub(d.LastHeartbeat)) + } + return deviceView{ + Key: d.Key, + Naziv: naziv, + MAC: d.MAC, + IP: d.IP, + FW: d.FW, + Serijski: d.Serial, + Poznat: d.Identified, + Tih: s.reg.Stale(&d, now), + Vidjen: humanAgo(now.Sub(d.LastSeen)), + Heartbeat: hb, + Interval: s.reg.Interval(&d), + Paketa: d.LogPackets + d.Heartbeats, + Linija: d.LogLines, + Poslednja: d.LastLine, + IzvornaIP: d.SourceIP, + PrviPutStr: d.FirstSeen.Format("02.01.2006. 15:04:05"), + } +} + +// humanAgo pise razmak vremena na srpskom, kratko. +func humanAgo(d time.Duration) string { + switch { + case d < 2*time.Second: + return "upravo sada" + case d < time.Minute: + return fmt.Sprintf("pre %d s", int(d.Seconds())) + case d < time.Hour: + return fmt.Sprintf("pre %d min", int(d.Minutes())) + case d < 48*time.Hour: + return fmt.Sprintf("pre %d h", int(d.Hours())) + default: + return fmt.Sprintf("pre %d dana", int(d.Hours()/24)) + } +} + +func (s *Service) pageIndex(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/" { + http.NotFound(w, r) + return + } + now := s.now() + devs := s.reg.Snapshot(now) + views := make([]deviceView, 0, len(devs)) + for _, d := range devs { + views = append(views, s.view(d, now)) + } + bytes, files := s.store.Usage() + data := indexData{ + Uredjaji: views, + Ukupno: len(views), + Paketa: s.stats.Packets.Load(), + Heartbeata: s.stats.Heartbeat.Load(), + Dnevnika: s.stats.LogPkts.Load(), + Smece: s.stats.Garbage.Load(), + Greske: s.stats.WriteErrs.Load(), + DiskMB: fmt.Sprintf("%.1f", float64(bytes)/(1<<20)), + Fajlova: files, + UDP: fmt.Sprintf("%s:%d", orAll(s.cfg.UDPHost), s.cfg.UDPPort), + Port: s.cfg.UDPPort, + Dir: s.cfg.Dir, + Vreme: now.Format("02.01.2006. 15:04:05"), + Token: s.cfg.Token, + } + render(w, tmplIndex, data) +} + +func (s *Service) pageDevice(w http.ResponseWriter, r *http.Request) { + key := SanitizeSegment(r.URL.Query().Get("k")) + q := r.URL.Query().Get("q") + if len(q) > 200 { + q = q[:200] + } + n := 200 + if v := r.URL.Query().Get("n"); v != "" { + if x, err := strconv.Atoi(v); err == nil && x > 0 && x <= 5000 { + n = x + } + } + now := s.now() + data := deviceData{Q: q, N: n, Vreme: now.Format("02.01.2006. 15:04:05"), Token: s.cfg.Token} + if d, ok := s.reg.Get(key); ok { + data.U = s.view(d, now) + data.Nadjen = true + } else { + // Uredjaj moze postojati samo na disku (posle restarta prijemnika). + data.U = deviceView{Key: key, Naziv: "(nije se javio od pokretanja prijemnika)"} + } + lines, err := s.store.Tail(key, n, q) + if err != nil { + lines = []string{"Greska pri citanju dnevnika: " + err.Error()} + } + data.Linije = lines + data.Dani = s.store.Days(key) + render(w, tmplDevice, data) +} + +// pageRaw daje isti sadrzaj kao stranica, ali kao cist tekst - za curl sa +// servera i za lepljenje u prepisku. +func (s *Service) pageRaw(w http.ResponseWriter, r *http.Request) { + key := SanitizeSegment(r.URL.Query().Get("k")) + n := 500 + if v := r.URL.Query().Get("n"); v != "" { + if x, err := strconv.Atoi(v); err == nil && x > 0 && x <= 5000 { + n = x + } + } + lines, err := s.store.Tail(key, n, r.URL.Query().Get("q")) + if err != nil { + http.Error(w, "Greska: "+err.Error(), http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "text/plain; charset=utf-8") + w.Write([]byte(strings.Join(lines, "\n") + "\n")) +} + +func (s *Service) apiDevices(w http.ResponseWriter, r *http.Request) { + now := s.now() + devs := s.reg.Snapshot(now) + out := make([]map[string]any, 0, len(devs)) + for _, d := range devs { + out = append(out, map[string]any{ + "kljuc": d.Key, + "naziv": d.Name, + "mac": d.MAC, + "ip": d.IP, + "izvorna_ip": d.SourceIP, + "fw": d.FW, + "serijski": d.Serial, + "identifikovan": d.Identified, + "tih": s.reg.Stale(&d, now), + "interval_sek": s.reg.Interval(&d), + "poslednji_put": d.LastSeen.Format(time.RFC3339), + "heartbeata": d.Heartbeats, + "paketa_dnevnik": d.LogPackets, + "linija": d.LogLines, + }) + } + w.Header().Set("Content-Type", "application/json; charset=utf-8") + enc := json.NewEncoder(w) + enc.SetIndent("", " ") + enc.Encode(map[string]any{"uredjaji": out, "vreme": now.Format(time.RFC3339)}) +} + +func orAll(h string) string { + if h == "" { + return "0.0.0.0" + } + return h +} + +func render(w http.ResponseWriter, t *template.Template, data any) { + w.Header().Set("Content-Type", "text/html; charset=utf-8") + if err := t.Execute(w, data); err != nil { + http.Error(w, "Greska pri prikazu: "+err.Error(), http.StatusInternalServerError) + } +} + +const cssCommon = ` +body{font-family:system-ui,Segoe UI,Arial,sans-serif;margin:0;background:#f4f6f8;color:#1c2733} +header{background:#12354f;color:#fff;padding:12px 18px;display:flex;gap:16px;align-items:baseline;flex-wrap:wrap} +header h1{font-size:18px;margin:0} +header a{color:#9fd3ff;text-decoration:none} +main{padding:16px 18px;max-width:1200px} +table{border-collapse:collapse;width:100%;background:#fff;box-shadow:0 1px 2px rgba(0,0,0,.12)} +th,td{padding:7px 10px;text-align:left;border-bottom:1px solid #e3e8ee;font-size:14px;vertical-align:top} +th{background:#eef2f6;font-weight:600} +tr.tih td{background:#ffe9e9} +.znak{display:inline-block;padding:1px 7px;border-radius:9px;font-size:12px;font-weight:600} +.ok{background:#d6f5dd;color:#0f6b2b} +.lose{background:#ffd2d2;color:#a11} +.nepoznat{background:#ffeccc;color:#8a5300} +.mali{color:#5a6b7b;font-size:12px} +pre{background:#0f1b26;color:#dbe7f0;padding:12px;border-radius:4px;overflow:auto;font-size:13px;line-height:1.45;max-height:70vh} +form.pretraga{margin:12px 0;display:flex;gap:8px;flex-wrap:wrap} +input,select{padding:6px 8px;border:1px solid #c6d0da;border-radius:4px;font-size:14px} +button{padding:6px 14px;border:0;border-radius:4px;background:#12354f;color:#fff;font-size:14px;cursor:pointer} +.prazno{padding:24px;background:#fff;color:#5a6b7b} +` + +var tmplIndex = template.Must(template.New("index").Parse(` + + + +Prijemnik dijagnostike — TERMINIA + + +

Prijemnik dijagnostike

+UDP {{.UDP}} · dnevnici: {{.Dir}} ({{.DiskMB}} MB / {{.Fajlova}} fajlova) +osveženo {{.Vreme}}
+
+

Uređaja: {{.Ukupno}} · primljeno paketa: {{.Paketa}} +(heartbeat {{.Heartbeata}}, dnevnik {{.Dnevnika}}, odbačeno {{.Smece}}) · greške upisa: {{.Greske}}

+{{if .Uredjaji}} + + +{{range .Uredjaji}} + + + + + + + + + + + +{{end}} +
UređajMACIPFWStanjePoslednji putHeartbeatLinijaPoslednja linija
{{.Naziv}} +{{if not .Poznat}}
nepoznat uređaj +FW stariji od 0.13.4 — prepoznat samo po adresi {{.IzvornaIP}}{{end}}
{{if .MAC}}{{.MAC}}{{else}}{{end}}{{if .IP}}{{.IP}}{{else}}{{.IzvornaIP}}{{end}}{{if .FW}}{{.FW}}{{else}}{{end}}{{if .Tih}}tih{{else}}javlja se{{end}}{{.Vidjen}}{{.Heartbeat}}
interval ~{{.Interval}} s
{{.Linija}}{{.Poslednja}}
+

Crveno = uređaj nije poslao ništa duže od 3× svog intervala javljanja.

+{{else}} +
Nijedan uređaj se još nije javio.

+Na uređaju: Podešavanja → ServissvcLogHost = <IP servera>:{{.Port}}, +uključi „Šalji dnevnik rada na server", heartbeat npr. 60 s.
+{{end}} +
`)) + +var tmplDevice = template.Must(template.New("device").Parse(` + + +{{.U.Naziv}} — dnevnik + + +

{{.U.Naziv}}

+{{if .U.MAC}}{{.U.MAC}} · {{end}}{{if .U.FW}}FW {{.U.FW}} · {{end}}{{if .U.IP}}{{.U.IP}}{{end}} +← svi uređaji +{{.Vreme}}
+
+{{if .Nadjen}} +

+{{if .U.Tih}}tih{{else}}javlja se{{end}} +· poslednji paket {{.U.Vidjen}} · heartbeat {{.U.Heartbeat}} (interval ~{{.U.Interval}} s) +· linija dnevnika: {{.U.Linija}} · prvi put viđen {{.U.PrviPutStr}} +{{if not .U.Poznat}}
nepoznat uređaj +identifikovan samo po izvornoj adresi {{.U.IzvornaIP}} — paket nije imao zaglavlje (FW stariji od 0.13.4){{end}} +

+{{else}} +

Uređaj se nije javio od pokretanja prijemnika — prikazano je ono što je ostalo na disku.

+{{end}} +
+ +{{if .Token}}{{end}} + + + +
+

Dani sa dnevnikom: {{if .Dani}}{{range .Dani}}{{.}} {{end}}{{else}}nema zapisa{{end}} +· preuzmi kao tekst

+{{if .Linije}}
{{range .Linije}}{{.}}
+{{end}}
{{else}}
Nema linija za zadate uslove.
{{end}} +
`))