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) } } } }