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 {