b07de747b54ae033d6f2cf2a061f852034a39911

Author
Ayman Bagabas <ayman.bagabas@gmail.com>
Committer
GitHub <noreply@github.com>
Date

Message

fix(log): respect log config (#280)

* fix(log): respect log config

Parse log config

fix(git): ensure os envs are present

feat: set the default time format to dateTime

fix(log): change update mirror log message into debug

fix(config): rename log config struct

fix(config): always return cfg

* perf: update mirrors in a workpool (#285)

* perf: update mirrors in a workqueue

Implement a simple chunked workqueue to queue updating mirrors. We use
the number of cpus to calculate the number of workers to distribute the
work to.

* fix: use automaxprocs

Signed-off-by: Carlos Alexandro Becker <caarlos0@users.noreply.github.com>

* fix: set maxprocs in main

* feat(wp): use a workpool impl

Use semaphores to implement a workpool of n workers
and use that to run the mirroring job

---------

Signed-off-by: Carlos Alexandro Becker <caarlos0@users.noreply.github.com>
Co-authored-by: Carlos Alexandro Becker <caarlos0@users.noreply.github.com>

---------

Signed-off-by: Carlos Alexandro Becker <caarlos0@users.noreply.github.com>
Co-authored-by: Carlos Alexandro Becker <caarlos0@users.noreply.github.com>

Diff

  1diff --git a/cmd/soft/root.go b/cmd/soft/root.go
  2index 1dc8e1883c9bac8c1b38b00c422cbc4bb7037532..273c9d6a3e7f54649bf1dab8927724cec5665567 100644
  3--- a/cmd/soft/root.go
  4+++ b/cmd/soft/root.go
  5@@ -8,6 +8,7 @@ import (
  6 	"github.com/charmbracelet/log"
  7 	. "github.com/charmbracelet/soft-serve/internal/log"
  8 	"github.com/spf13/cobra"
  9+	"go.uber.org/automaxprocs/maxprocs"
 10 )
 11 
 12 var (
 13@@ -52,6 +53,13 @@ func init() {
 14 
 15 func main() {
 16 	logger := NewDefaultLogger()
 17+
 18+	// Set the max number of processes to the number of CPUs
 19+	// This is useful when running soft serve in a container
 20+	if _, err := maxprocs.Set(maxprocs.Logger(logger.Debugf)); err != nil {
 21+		logger.Warn("couldn't set automaxprocs", "error", err)
 22+	}
 23+
 24 	ctx := log.WithContext(context.Background(), logger)
 25 	if err := rootCmd.ExecuteContext(ctx); err != nil {
 26 		os.Exit(1)
 27diff --git a/go.mod b/go.mod
 28index e7ca5b691bd1264b1114d256ca39ee868b16f46f..bdd5d6109e36bd4ba33daa721386777ab263197f 100644
 29--- a/go.mod
 30+++ b/go.mod
 31@@ -32,6 +32,7 @@ require (
 32 	github.com/prometheus/client_golang v1.15.1
 33 	github.com/robfig/cron/v3 v3.0.1
 34 	github.com/spf13/cobra v1.7.0
 35+	go.uber.org/automaxprocs v1.5.2
 36 	goji.io v2.0.2+incompatible
 37 	golang.org/x/crypto v0.9.0
 38 	golang.org/x/sync v0.2.0
 39diff --git a/go.sum b/go.sum
 40index 3cbdc98fb3c8fb6879af328bfd44a958c92dd536..29e1ad2153c2c2fa4eae1b0f68884ee53435af72 100644
 41--- a/go.sum
 42+++ b/go.sum
 43@@ -161,6 +161,7 @@ github.com/pjbgf/sha1cd v0.3.0/go.mod h1:nZ1rrWOcGJ5uZgEEVL1VUM9iRQiZvWdbZjkKyFz
 44 github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
 45 github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
 46 github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
 47+github.com/prashantv/gostub v1.1.0 h1:BTyx3RfQjRHnUWaGF9oQos79AlQ5k8WNktv7VGvVH4g=
 48 github.com/prometheus/client_golang v1.15.1 h1:8tXpTmJbyH5lydzFPoxSIJ0J46jdh3tylbvM1xCv0LI=
 49 github.com/prometheus/client_golang v1.15.1/go.mod h1:e9yaBhRPU2pPNsZwE+JdQl0KEt1N9XgF6zxWmaC0xOk=
 50 github.com/prometheus/client_model v0.3.0 h1:UBgGFHqYdG/TPFD1B1ogZywDqEkwp3fBMvqdiQ7Xew4=
 51@@ -207,6 +208,8 @@ github.com/yuin/goldmark v1.5.2 h1:ALmeCk/px5FSm1MAcFBAsVKZjDuMVj8Tm7FFIlMJnqU=
 52 github.com/yuin/goldmark v1.5.2/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY=
 53 github.com/yuin/goldmark-emoji v1.0.1 h1:ctuWEyzGBwiucEqxzwe0SOYDXPAucOrE9NQC18Wa1os=
 54 github.com/yuin/goldmark-emoji v1.0.1/go.mod h1:2w1E6FEWLcDQkoTE+7HU6QF1F6SLlNGjRIBbIZQFqkQ=
 55+go.uber.org/automaxprocs v1.5.2 h1:2LxUOGiR3O6tw8ui5sZa2LAaHnsviZdVOUZw4fvbnME=
 56+go.uber.org/automaxprocs v1.5.2/go.mod h1:eRbA25aqJrxAbsLO0xy5jVwPt7FQnRgjW+efnwa1WM0=
 57 goji.io v2.0.2+incompatible h1:uIssv/elbKRLznFUy3Xj4+2Mz/qKhek/9aZQDUMae7c=
 58 goji.io v2.0.2+incompatible/go.mod h1:sbqFwrtqZACxLBTQcdgVjFh54yGVCvwq8+w49MVMMIk=
 59 golang.org/x/arch v0.1.0/go.mod h1:5om86z9Hs0C8fWVUuoMHwpExlXzs5Tkyp9hOrfG7pp8=
 60diff --git a/internal/log/log.go b/internal/log/log.go
 61index a1184153b81da766f76fc10ca7f533c0a486b710..7389b80fb9e073174653271fe505282369df91b5 100644
 62--- a/internal/log/log.go
 63+++ b/internal/log/log.go
 64@@ -2,17 +2,29 @@ package log
 65 
 66 import (
 67 	"os"
 68+	"path/filepath"
 69 	"strconv"
 70 	"strings"
 71 	"time"
 72 
 73 	"github.com/charmbracelet/log"
 74+	"github.com/charmbracelet/soft-serve/server/config"
 75 )
 76 
 77 var contextKey = &struct{ string }{"logger"}
 78 
 79 // NewDefaultLogger returns a new logger with default settings.
 80 func NewDefaultLogger() *log.Logger {
 81+	dp := os.Getenv("SOFT_SERVE_DATA_PATH")
 82+	if dp == "" {
 83+		dp = "data"
 84+	}
 85+
 86+	cfg, err := config.ParseConfig(filepath.Join(dp, "config.yaml"))
 87+	if err != nil {
 88+		log.Errorf("failed to parse config: %v", err)
 89+	}
 90+
 91 	logger := log.NewWithOptions(os.Stderr, log.Options{
 92 		ReportTimestamp: true,
 93 		TimeFormat:      time.DateOnly,
 94@@ -22,11 +34,9 @@ func NewDefaultLogger() *log.Logger {
 95 		logger.SetLevel(log.DebugLevel)
 96 	}
 97 
 98-	if tsfmt := os.Getenv("SOFT_SERVE_LOG_TIME_FORMAT"); tsfmt != "" {
 99-		logger.SetTimeFormat(tsfmt)
100-	}
101+	logger.SetTimeFormat(cfg.Log.TimeFormat)
102 
103-	switch strings.ToLower(os.Getenv("SOFT_SERVE_LOG_FORMAT")) {
104+	switch strings.ToLower(cfg.Log.Format) {
105 	case "json":
106 		logger.SetFormatter(log.JSONFormatter)
107 	case "logfmt":
108diff --git a/internal/sync/workqueue.go b/internal/sync/workqueue.go
109new file mode 100644
110index 0000000000000000000000000000000000000000..16eb442a3c5588ca7be41cfd47281c1f90109629
111--- /dev/null
112+++ b/internal/sync/workqueue.go
113@@ -0,0 +1,98 @@
114+package sync
115+
116+import (
117+	"context"
118+	"sync"
119+
120+	"golang.org/x/sync/semaphore"
121+)
122+
123+// WorkPool is a pool of work to be done.
124+type WorkPool struct {
125+	workers int
126+	work    map[string]func()
127+	mu      sync.RWMutex
128+	sem     *semaphore.Weighted
129+	ctx     context.Context
130+	logger  func(string, ...interface{})
131+}
132+
133+// WorkPoolOption is a function that configures a WorkPool.
134+type WorkPoolOption func(*WorkPool)
135+
136+// WithWorkPoolLogger sets the logger to use.
137+func WithWorkPoolLogger(logger func(string, ...interface{})) WorkPoolOption {
138+	return func(wq *WorkPool) {
139+		wq.logger = logger
140+	}
141+}
142+
143+// NewWorkPool creates a new work pool. The workers argument specifies the
144+// number of concurrent workers to run the work.
145+// The queue will chunk the work into batches of workers size.
146+func NewWorkPool(ctx context.Context, workers int, opts ...WorkPoolOption) *WorkPool {
147+	wq := &WorkPool{
148+		workers: workers,
149+		work:    make(map[string]func()),
150+		ctx:     ctx,
151+	}
152+
153+	for _, opt := range opts {
154+		opt(wq)
155+	}
156+
157+	if wq.workers <= 0 {
158+		wq.workers = 1
159+	}
160+
161+	wq.sem = semaphore.NewWeighted(int64(wq.workers))
162+
163+	return wq
164+}
165+
166+// Run starts the workers and waits for them to finish.
167+func (wq *WorkPool) Run() {
168+	for id, fn := range wq.work {
169+		if err := wq.sem.Acquire(wq.ctx, 1); err != nil {
170+			wq.logf("workpool: %v", err)
171+			return
172+		}
173+
174+		go func(id string, fn func()) {
175+			defer wq.sem.Release(1)
176+			fn()
177+			wq.mu.Lock()
178+			delete(wq.work, id)
179+			wq.mu.Unlock()
180+		}(id, fn)
181+	}
182+
183+	if err := wq.sem.Acquire(wq.ctx, int64(wq.workers)); err != nil {
184+		wq.logf("workpool: %v", err)
185+	}
186+}
187+
188+// Add adds a new job to the pool.
189+// If the job already exists, it is a no-op.
190+func (wq *WorkPool) Add(id string, fn func()) {
191+	wq.mu.Lock()
192+	defer wq.mu.Unlock()
193+	if _, ok := wq.work[id]; ok {
194+		return
195+	}
196+	wq.work[id] = fn
197+}
198+
199+// Status checks if a job is in the queue.
200+func (wq *WorkPool) Status(id string) bool {
201+	wq.mu.RLock()
202+	defer wq.mu.RUnlock()
203+	_, ok := wq.work[id]
204+	return ok
205+}
206+
207+func (wq *WorkPool) logf(format string, args ...interface{}) {
208+	if wq.logger != nil {
209+		wq.logger(format, args...)
210+	}
211+}
212diff --git a/server/config/config.go b/server/config/config.go
213index f75a7a6d20a9ca062f5567876f84effae3745561..a96e69c75f5bc2e3944160aa86831236a177b71c 100644
214--- a/server/config/config.go
215+++ b/server/config/config.go
216@@ -73,6 +73,17 @@ type StatsConfig struct {
217 	ListenAddr string `env:"LISTEN_ADDR" yaml:"listen_addr"`
218 }
219 
220+// LogConfig is the logger configuration.
221+type LogConfig struct {
222+	// Format is the format of the logs.
223+	// Valid values are "json", "logfmt", and "text".
224+	Format string `env:"FORMAT" yaml:"format"`
225+
226+	// Time format for the log `ts` field.
227+	// Format must be described in Golang's time format.
228+	TimeFormat string `env:"TIME_FORMAT" yaml:"time_format"`
229+}
230+
231 // Config is the configuration for Soft Serve.
232 type Config struct {
233 	// Name is the name of the server.
234@@ -90,13 +101,8 @@ type Config struct {
235 	// Stats is the configuration for the stats server.
236 	Stats StatsConfig `envPrefix:"STATS_" yaml:"stats"`
237 
238-	// LogFormat is the format of the logs.
239-	// Valid values are "json", "logfmt", and "text".
240-	LogFormat string `env:"LOG_FORMAT" yaml:"log_format"`
241-
242-	// Time format for the log `ts` field.
243-	// Format must be described in Golang's time format.
244-	LogTimeFormat string `env:"LOG_TIME_FORMAT" yaml:"log_time_format"`
245+	// Log is the logger configuration.
246+	Log LogConfig `envPrefix:"LOG_" yaml:"log"`
247 
248 	// InitialAdminKeys is a list of public keys that will be added to the list of admins.
249 	InitialAdminKeys []string `env:"INITIAL_ADMIN_KEYS" envSeparator:"\n" yaml:"initial_admin_keys"`
250@@ -111,10 +117,8 @@ type Config struct {
251 func parseConfig(path string) (*Config, error) {
252 	dataPath := filepath.Dir(path)
253 	cfg := &Config{
254-		Name:          "Soft Serve",
255-		LogFormat:     "text",
256-		LogTimeFormat: time.DateOnly,
257-		DataPath:      dataPath,
258+		Name:     "Soft Serve",
259+		DataPath: dataPath,
260 		SSH: SSHConfig{
261 			ListenAddr:    ":23231",
262 			PublicURL:     "ssh://localhost:23231",
263@@ -136,6 +140,10 @@ func parseConfig(path string) (*Config, error) {
264 		Stats: StatsConfig{
265 			ListenAddr: "localhost:23233",
266 		},
267+		Log: LogConfig{
268+			Format:     "text",
269+			TimeFormat: time.DateTime,
270+		},
271 	}
272 
273 	f, err := os.Open(path)
274@@ -182,11 +190,11 @@ func parseConfig(path string) (*Config, error) {
275 func ParseConfig(path string) (*Config, error) {
276 	cfg, err := parseConfig(path)
277 	if err != nil {
278-		return nil, err
279+		return cfg, err
280 	}
281 
282 	if err := cfg.validate(); err != nil {
283-		return nil, err
284+		return cfg, err
285 	}
286 
287 	return cfg, nil
288diff --git a/server/config/file.go b/server/config/file.go
289index 36c0a7a311bd96d924c33eeeea33a701d37e8d0d..09e5ce2e00dec68891f27e63bd5cf14fa1faca2c 100644
290--- a/server/config/file.go
291+++ b/server/config/file.go
292@@ -11,9 +11,13 @@ var configFileTmpl = template.Must(template.New("config").Parse(`# Soft Serve Se
293 # This is the name that will be displayed in the UI.
294 name: "{{ .Name }}"
295 
296-# Log format to use. Valid values are "json", "logfmt", and "text".
297-log_format: "{{ .LogFormat }}"
298-log_time_format: "{{ .LogTimeFormat }}"
299+# Logging configuration.
300+log:
301+  # Log format to use. Valid values are "json", "logfmt", and "text".
302+  format: "{{ .Log.Format }}"
303+  # Time format for the log "timestamp" field.
304+  # Should be described in Golang's time format.
305+  time_format: "{{ .Log.TimeFormat }}"
306 
307 # The SSH server configuration.
308 ssh:
309diff --git a/server/git/git.go b/server/git/git.go
310index 85e2cf068e5f24e141c0455db2d8a1a5aa01121c..ef8affe207a0cf4993976cbcc28fe91e5e9d512f 100644
311--- a/server/git/git.go
312+++ b/server/git/git.go
313@@ -83,12 +83,12 @@ func RunGit(ctx context.Context, in io.Reader, out io.Writer, er io.Writer, dir
314 	logger := log.FromContext(ctx).WithPrefix("rungit")
315 	c := exec.CommandContext(ctx, "git", args...)
316 	c.Dir = dir
317-	c.Env = append(c.Env, envs...)
318+	c.Env = append(os.Environ(), envs...)
319 	c.Env = append(c.Env, "PATH="+os.Getenv("PATH"))
320 	c.Env = append(c.Env, "SOFT_SERVE_DEBUG="+os.Getenv("SOFT_SERVE_DEBUG"))
321 	if cfg != nil {
322-		c.Env = append(c.Env, "SOFT_SERVE_LOG_FORMAT="+cfg.LogFormat)
323-		c.Env = append(c.Env, "SOFT_SERVE_LOG_TIME_FORMAT="+cfg.LogTimeFormat)
324+		c.Env = append(c.Env, "SOFT_SERVE_LOG_FORMAT="+cfg.Log.Format)
325+		c.Env = append(c.Env, "SOFT_SERVE_LOG_TIME_FORMAT="+cfg.Log.TimeFormat)
326 	}
327 
328 	stdin, err := c.StdinPipe()
329diff --git a/server/jobs.go b/server/jobs.go
330index 06656d1ccec5b53b6868fb613168b170802962db..9cd23a23f1ad78fc1eed2f32ae8f63e6f5b497bb 100644
331--- a/server/jobs.go
332+++ b/server/jobs.go
333@@ -3,15 +3,15 @@ package server
334 import (
335 	"fmt"
336 	"path/filepath"
337+	"runtime"
338 
339 	"github.com/charmbracelet/soft-serve/git"
340+	"github.com/charmbracelet/soft-serve/internal/sync"
341 )
342 
343-var (
344-	jobSpecs = map[string]string{
345-		"mirror": "@every 10m",
346-	}
347-)
348+var jobSpecs = map[string]string{
349+	"mirror": "@every 10m",
350+}
351 
352 // mirrorJob runs the (pull) mirror job task.
353 func (s *Server) mirrorJob() func() {
354@@ -25,26 +25,37 @@ func (s *Server) mirrorJob() func() {
355 			return
356 		}
357 
358+		// Divide the work up among the number of CPUs.
359+		wq := sync.NewWorkPool(s.ctx, runtime.GOMAXPROCS(0),
360+			sync.WithWorkPoolLogger(logger.Errorf),
361+		)
362+
363+		logger.Debug("updating mirror repos")
364 		for _, repo := range repos {
365 			if repo.IsMirror() {
366-				logger.Info("updating mirror", "repo", repo.Name())
367 				r, err := repo.Open()
368 				if err != nil {
369 					logger.Error("error opening repository", "repo", repo.Name(), "err", err)
370 					continue
371 				}
372 
373-				cmd := git.NewCommand("remote", "update", "--prune")
374-				cmd.AddEnvs(
375-					fmt.Sprintf(`GIT_SSH_COMMAND=ssh -o UserKnownHostsFile="%s" -o StrictHostKeyChecking=no -i "%s"`,
376-						filepath.Join(cfg.DataPath, "ssh", "known_hosts"),
377-						cfg.SSH.ClientKeyPath,
378-					),
379-				)
380-				if _, err := cmd.RunInDir(r.Path); err != nil {
381-					logger.Error("error running git remote update", "repo", repo.Name(), "err", err)
382-				}
383+				name := repo.Name()
384+				wq.Add(name, func() {
385+					cmd := git.NewCommand("remote", "update", "--prune")
386+					cmd.AddEnvs(
387+						fmt.Sprintf(`GIT_SSH_COMMAND=ssh -o UserKnownHostsFile="%s" -o StrictHostKeyChecking=no -i "%s"`,
388+							filepath.Join(cfg.DataPath, "ssh", "known_hosts"),
389+							cfg.SSH.ClientKeyPath,
390+						),
391+					)
392+					if _, err := cmd.RunInDir(r.Path); err != nil {
393+						logger.Error("error running git remote update", "repo", name, "err", err)
394+					}
395+
396+				})
397 			}
398 		}
399+
400+		wq.Run()
401 	}
402 }
403diff --git a/server/server.go b/server/server.go
404index 3b9554b2bd83069aa890348fd67df8fea2d232c3..690c9650ba37c47e50780704225c0b05ab625331 100644
405--- a/server/server.go
406+++ b/server/server.go
407@@ -62,7 +62,7 @@ func NewServer(ctx context.Context) (*Server, error) {
408 	}
409 
410 	// Add cron jobs.
411-	srv.Cron.AddFunc(jobSpecs["mirror"], srv.mirrorJob())
412+	_, _ = srv.Cron.AddFunc(jobSpecs["mirror"], srv.mirrorJob())
413 
414 	srv.SSHServer, err = sshsrv.NewSSHServer(ctx)
415 	if err != nil {