/
t3
/
cli
Обзор
Документация
Войти
/
t3
/
cli
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
internal/agent/agent.go
137 строк
3 KB
Ivan Shibkikh
fix console.log issue
06 июл 2026, 17:09
06 июл 2026, 17:09
5e5a141
Код
Авторство
О чём код?
package agent import ( "context" "fmt" "log/slog" "net" "net/http" "os" "sync" "time" "gitverse.ru/t3/cli/internal/events" "gitverse.ru/t3/cli/internal/lifecycle" "gitverse.ru/t3/cli/internal/metrics" "gitverse.ru/t3/cli/pkg/pb" "google.golang.org/grpc" ) type Agent struct { pb.UnimplementedAgentServiceServer ctx context.Context addr string token string hostname string launch *lifecycle.Launch registry *metrics.MetricsRegistry metricsAddr string currentRunID string runIDMu sync.Mutex cacheDir string cacheMaxSize int64 sender events.LogSender // живет от старта до завершения стрима senderMu sync.Mutex } func NewAgent(ctx context.Context, addr string, token *string, metricsAddr string, cacheDir string, cacheMaxSize int64) *Agent { t := "" if token != nil { t = *token } hostname, err := os.Hostname() if err != nil { panic(err) } reg := metrics.NewRegistry(hostname) launch := lifecycle.NewLaunch() launch.SetRegistry(reg) return &Agent{ ctx: ctx, addr: addr, token: t, hostname: hostname, launch: launch, registry: reg, metricsAddr: metricsAddr, cacheDir: cacheDir, cacheMaxSize: cacheMaxSize, } } func (ag *Agent) Run() error { lis, err := net.Listen("tcp", ag.addr) if err != nil { return fmt.Errorf("agent failed to listen: %w", err) } grpcServer := grpc.NewServer() pb.RegisterAgentServiceServer(grpcServer, ag) if ag.token != "" { slog.Info("agent started", "addr", ag.addr, "hostname", ag.hostname, "token", ag.token) } else { slog.Info("agent started", "addr", ag.addr, "hostname", ag.hostname) } metricsServer := ag.startMetricsServer() go func() { <-ag.ctx.Done() if metricsServer != nil { shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) metricsServer.Shutdown(shutdownCtx) //nolint:errcheck cancel() } stopped := make(chan struct{}) go func() { grpcServer.GracefulStop() close(stopped) }() shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() select { case <-stopped: case <-shutdownCtx.Done(): grpcServer.Stop() } slog.Info("agent stopped", "hostname", ag.hostname) }() err = grpcServer.Serve(lis) if err != nil && err != grpc.ErrServerStopped { return err } return nil } func (ag *Agent) startMetricsServer() *http.Server { if ag.metricsAddr == "" { return nil } mux := http.NewServeMux() mux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) { ag.registry.VMSet().WritePrometheus(w) }) srv := &http.Server{ Addr: ag.metricsAddr, Handler: mux, } go func() { slog.Info("metrics server started", "addr", ag.metricsAddr) if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { slog.Error("metrics server error", "error", err) } }() return srv }