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
53 changes: 9 additions & 44 deletions cmd/agent/main.go
Original file line number Diff line number Diff line change
@@ -1,68 +1,33 @@
package main

import (
"flag"
"fmt"
"context"
"github.com/iliaonishchenko/aggreg8/internal/agent"
agentconfig "github.com/iliaonishchenko/aggreg8/internal/config/agent"
models "github.com/iliaonishchenko/aggreg8/internal/model"
"log"
"os"
"os/signal"
"syscall"
"time"
)

func main() {

defaultServerAddress := "localhost:8080"
defaultReportInterval := 10
defaultPollInterval := 2
defaultRateLimit := 5

config, err := agentconfig.LoadConfig()

if err != nil {
log.Fatalf("error loading config: %v", err)
}

if config.ServerAddress == "" {
flag.StringVar(&config.ServerAddress, "a", defaultServerAddress, "address and port to run server")
}
if config.ReportInterval == 0 {
flag.IntVar(&config.ReportInterval, "r", defaultReportInterval, "report interval in seconds")
}
if config.PollInterval == 0 {
flag.IntVar(&config.PollInterval, "p", defaultPollInterval, "poll interval in seconds")
}

flag.Parse()
agent.ParseFlags(config, defaultServerAddress, "", defaultReportInterval, defaultPollInterval, defaultRateLimit)

collector := agent.NewCollector()
errorClassifier := agent.NewAgentErrorClassifier()
sender := agent.NewSender(config.ServerAddress, errorClassifier)
collectTicker := time.NewTicker(time.Duration(config.PollInterval) * time.Second)
defer collectTicker.Stop()
sendTicker := time.NewTicker(time.Duration(config.ReportInterval) * time.Second)
defer sendTicker.Stop()
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
initLogger(config)

for {
select {
case <-collectTicker.C:
collector.Collect()
case <-sendTicker.C:
metrics := collector.GetMetrics()
if err = sendMetrics(metrics, sender); err != nil {
log.Printf("error sending metrics: %v", err)
}
case <-sigChan:
fmt.Println("Shutting down...")
return
}
}
}
context := context.Background()

func sendMetrics(metrics []*models.Metrics, sender *agent.Sender) error {
return sender.SendJSONWithRetries(metrics...)
collector := agent.NewCollector()
sender := initSender(config)
app := agent.NewAgent(collector, sender, config.PollInterval, config.ReportInterval, config.RateLimit)
app.Run(context)
}
30 changes: 30 additions & 0 deletions cmd/agent/setup.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
package main

import (
"github.com/iliaonishchenko/aggreg8/internal/agent"
agentconfig "github.com/iliaonishchenko/aggreg8/internal/config/agent"
"github.com/iliaonishchenko/aggreg8/internal/logger"
"github.com/iliaonishchenko/aggreg8/internal/signature"
"log"
)

func initSender(config *agentconfig.Config) *agent.Sender {
var sender *agent.Sender
var sign *signature.Signature
errorClassifier := agent.NewAgentErrorClassifier()

if config.Key != "" {
sign = signature.NewSignature(config.Key)
sender = agent.NewSender(config.ServerAddress, errorClassifier, sign)
} else {
sender = agent.NewSender(config.ServerAddress, errorClassifier, nil)
}

return sender
}

func initLogger(config *agentconfig.Config) {
if err := logger.Initialize(config.LogLevel); err != nil {
log.Printf("error initializing logger: %v", err)
}
}
9 changes: 7 additions & 2 deletions cmd/server/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,21 +5,23 @@ import (
"github.com/iliaonishchenko/aggreg8/internal/config/server"
)

func parseFlags(cfg *server.Config, defaultServerAddress string, defaultStoreInterval int, defaultFileStoragePath string, defaultRestore bool) {
func parseFlags(cfg *server.Config, defaultServerAddress string, defaultStoreInterval int, defaultFileStoragePath string, defaultRestore bool, defaultKey string) {
var serverAddress string
var storeInterval int
var fileStoragePath string
var restore bool
var databaseDSN string
var key string

flag.StringVar(&serverAddress, "a", defaultServerAddress, "address and port to run server")
flag.IntVar(&storeInterval, "i", defaultStoreInterval, "interval to store aggregator data")
flag.StringVar(&fileStoragePath, "f", defaultFileStoragePath, "file path to store aggregator data")
flag.BoolVar(&restore, "r", defaultRestore, "restore aggregator data from file on startup")
flag.StringVar(&databaseDSN, "d", "", "database DSN")
flag.StringVar(&key, "k", "", "key for signature")

flag.Parse()

if cfg.ServerAddress == "" {
cfg.ServerAddress = serverAddress
}
Expand All @@ -39,4 +41,7 @@ func parseFlags(cfg *server.Config, defaultServerAddress string, defaultStoreInt
if cfg.DatabaseDSN == "" {
cfg.DatabaseDSN = databaseDSN
}
if cfg.Key == "" {
cfg.Key = key
}
}
7 changes: 6 additions & 1 deletion cmd/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
"github.com/iliaonishchenko/aggreg8/internal/service/memory"
"github.com/iliaonishchenko/aggreg8/internal/service/pg"
"github.com/iliaonishchenko/aggreg8/internal/service/sync"
"github.com/iliaonishchenko/aggreg8/internal/signature"
_ "github.com/jackc/pgx/v5/stdlib"
"go.uber.org/zap"
"log"
Expand Down Expand Up @@ -72,12 +73,14 @@ func main() {
defaultStoreInterval := 300
defaultFileStoragePath := "./snapshot.json"
defaultRestore := false
defaultKey := ""

cfg, err := server.LoadConfig()
if err != nil {
log.Fatalf("error loading config: %v", err)
}

parseFlags(cfg, defaultServerAddress, defaultStoreInterval, defaultFileStoragePath, defaultRestore)
parseFlags(cfg, defaultServerAddress, defaultStoreInterval, defaultFileStoragePath, defaultRestore, defaultKey)

if err := logger.Initialize(cfg.LogLevel); err != nil {
log.Fatalf("error initializing logger: %v", err)
Expand All @@ -101,9 +104,11 @@ func run(cfg server.Config, storage service.MetricStorage, repo *repository.Metr
allHandler := handler.NewAllMetricsHandler(storage)
getMetricHandler := handler.NewGetMetricHandler(storage)
pingHandler := handler.NewPingHandler(repo)
sign := signature.NewSignature(cfg.Key)

r.Use(logger.WithLogger)
r.Use(router.WithCompression)
r.Use(router.WithHash(sign))

r.Route("/", func(r chi.Router) {
r.Get("/ping", pingHandler.HandlePing)
Expand Down
7 changes: 6 additions & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -8,21 +8,26 @@ require (
github.com/golang/mock v1.6.0
github.com/jackc/pgx/v5 v5.8.0
github.com/pressly/goose/v3 v3.26.0
github.com/shirou/gopsutil v3.21.11+incompatible
github.com/stretchr/testify v1.11.1
go.uber.org/zap v1.27.1
)

require (
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/jackc/pgerrcode v0.0.0-20250907135507-afb5586c32a6 // indirect
github.com/go-ole/go-ole v1.2.6 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/mfridman/interpolate v0.0.2 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/sethvargo/go-retry v0.3.0 // indirect
github.com/tklauser/go-sysconf v0.3.16 // indirect
github.com/tklauser/numcpus v0.11.0 // indirect
github.com/yusufpapurcu/wmi v1.2.4 // indirect
go.uber.org/multierr v1.11.0 // indirect
golang.org/x/sync v0.17.0 // indirect
golang.org/x/sys v0.38.0 // indirect
golang.org/x/text v0.29.0 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)
17 changes: 13 additions & 4 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,12 @@ github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkp
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
github.com/go-chi/chi/v5 v5.2.3 h1:WQIt9uxdsAbgIYgid+BpYc+liqQZGMHRaUwp0JUcvdE=
github.com/go-chi/chi/v5 v5.2.3/go.mod h1:L2yAIGWB3H+phAw1NxKwWM+7eUH/lU8pOMm5hHcoops=
github.com/go-ole/go-ole v1.2.6 h1:/Fpf6oFPoeFik9ty7siob0G6Ke8QvQEuVcuChpwXzpY=
github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0=
github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc=
github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/jackc/pgerrcode v0.0.0-20250907135507-afb5586c32a6 h1:D/V0gu4zQ3cL2WKeVNVM4r2gLxGGf6McLwgXzRTo2RQ=
github.com/jackc/pgerrcode v0.0.0-20250907135507-afb5586c32a6/go.mod h1:a/s9Lp5W7n/DD0VrVoyJ00FbP2ytTPDVOivvn2bMlds=
github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM=
github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg=
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo=
Expand Down Expand Up @@ -41,12 +41,20 @@ github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0t
github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc=
github.com/sethvargo/go-retry v0.3.0 h1:EEt31A35QhrcRZtrYFDTBg91cqZVnFL2navjDrah2SE=
github.com/sethvargo/go-retry v0.3.0/go.mod h1:mNX17F0C/HguQMyMyJxcnU471gOZGxCLyYaFyAZraas=
github.com/shirou/gopsutil v3.21.11+incompatible h1:+1+c1VGhc88SSonWP6foOcLhvnKlUeu/erjjvaPEYiI=
github.com/shirou/gopsutil v3.21.11+incompatible/go.mod h1:5b4v6he4MtMOwMlS0TUMTu2PcXUg8+E1lC7eC3UO/RA=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/tklauser/go-sysconf v0.3.16 h1:frioLaCQSsF5Cy1jgRBrzr6t502KIIwQ0MArYICU0nA=
github.com/tklauser/go-sysconf v0.3.16/go.mod h1:/qNL9xxDhc7tx3HSRsLWNnuzbVfh3e7gh/BmM179nYI=
github.com/tklauser/numcpus v0.11.0 h1:nSTwhKH5e1dMNsCdVBukSZrURJRoHbSEQjdEbY+9RXw=
github.com/tklauser/numcpus v0.11.0/go.mod h1:z+LwcLq54uWZTX0u/bGobaV34u6V7KNlTZejzM6/3MQ=
github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k=
github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0=
github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0=
go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto=
go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE=
go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0=
Expand All @@ -67,11 +75,12 @@ golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug=
golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210510120138-977fb7262007/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.34.0 h1:H5Y5sJ2L2JRdyv7ROF1he/lPdvFsd0mJHFw2ThKHxLA=
golang.org/x/sys v0.34.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
golang.org/x/sys v0.38.0 h1:3yZWxaJjBmCWXqhN1qh02AkOnCQ1poK6oF+a7xWL6Gc=
golang.org/x/sys v0.38.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
Expand Down
77 changes: 77 additions & 0 deletions internal/agent/agent.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
package agent

import (
"context"
"github.com/iliaonishchenko/aggreg8/internal/logger"
models "github.com/iliaonishchenko/aggreg8/internal/model"
"time"
)

type CollectorService interface {
CollectRuntime()
CollectSystem()
GetMetrics() [][]*models.Metrics
}

type SenderService interface {
SendJSONWithRetries(metrics ...*models.Metrics) error
}

type Agent struct {
collector CollectorService
sender SenderService
pollInterval int
reportInterval int
rateLimit int
}

func NewAgent(collector CollectorService, sender SenderService, pollInterval, reportInterval, rateLimit int) *Agent {
return &Agent{
collector: collector,
sender: sender,
pollInterval: pollInterval,
reportInterval: reportInterval,
rateLimit: rateLimit,
}
}

func (a *Agent) Run(ctx context.Context) {
collectTicker := time.NewTicker(time.Duration(a.pollInterval) * time.Second)
defer collectTicker.Stop()
reportTicker := time.NewTicker(time.Duration(a.reportInterval) * time.Second)
defer reportTicker.Stop()
jobsCh := make(chan []*models.Metrics, a.rateLimit)

for i := 0; i < a.rateLimit; i++ {
go a.worker(jobsCh)
}

for {
select {
case <-collectTicker.C:
go a.collector.CollectRuntime()
go a.collector.CollectSystem()
case <-reportTicker.C:
go func() {
metrics := a.collector.GetMetrics()
for _, metricBatch := range metrics {
jobsCh <- metricBatch
}
}()
case <-ctx.Done():
logger.Log.Info("Shutting down agent...")
close(jobsCh)
return
}
}
}

func (a *Agent) worker(jobs <-chan []*models.Metrics) {
for metrics := range jobs {
logger.Log.Info("Sending metrics %v", logger.Field("metrics", metrics[0]))
err := a.sender.SendJSONWithRetries(metrics...)
if err != nil {
logger.Log.Error("error sending metrics: %v", logger.Err(err))
}
}
}
50 changes: 50 additions & 0 deletions internal/agent/agent_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
package agent

import (
"context"
"github.com/golang/mock/gomock"
"github.com/iliaonishchenko/aggreg8/internal/agent/mocks"
"testing"
"time"
)

func TestRun(t *testing.T) {
t.Run("collects on poll interval", func(t *testing.T) {
ctrl := gomock.NewController(t)
defer ctrl.Finish()
mockCollector := mocks.NewMockCollectorService(ctrl)
mockSender := mocks.NewMockSenderService(ctrl)
agent := NewAgent(mockCollector, mockSender, 1, 2, 2)
ctx, cancel := context.WithCancel(context.Background())

mockCollector.EXPECT().CollectSystem().MinTimes(2)
mockCollector.EXPECT().CollectRuntime().MinTimes(2)
mockCollector.EXPECT().GetMetrics().AnyTimes()

go agent.Run(ctx)

time.Sleep(3 * time.Second)
cancel()
time.Sleep(100 * time.Millisecond)
})

t.Run("sends on report interval", func(t *testing.T) {
ctrl := gomock.NewController(t)
defer ctrl.Finish()
mockCollector := mocks.NewMockCollectorService(ctrl)
mockSender := mocks.NewMockSenderService(ctrl)
agent := NewAgent(mockCollector, mockSender, 1, 2, 2)
ctx, cancel := context.WithCancel(context.Background())

mockCollector.EXPECT().CollectSystem().AnyTimes()
mockCollector.EXPECT().CollectRuntime().AnyTimes()
mockCollector.EXPECT().GetMetrics().AnyTimes()
mockSender.EXPECT().SendJSONWithRetries().AnyTimes()

go agent.Run(ctx)

time.Sleep(3 * time.Second)
cancel()
time.Sleep(100 * time.Millisecond)
})
}
Loading