From f308060f47abd0c4ff8ddf8ea93f746a5993e44b Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Thu, 29 Jan 2026 10:32:42 -0800 Subject: [PATCH 1/3] gosky: refactor log setup --- cmd/gosky/car.go | 1 + cmd/gosky/debug.go | 3 +++ cmd/gosky/main.go | 38 +++++++++++++++++++++++++++++--------- cmd/gosky/streamdiff.go | 2 ++ cmd/gosky/sync.go | 1 + 5 files changed, 36 insertions(+), 9 deletions(-) diff --git a/cmd/gosky/car.go b/cmd/gosky/car.go index 2098a2cd..42b97d6e 100644 --- a/cmd/gosky/car.go +++ b/cmd/gosky/car.go @@ -38,6 +38,7 @@ var carUnpackCmd = &cli.Command{ }, ArgsUsage: ``, Action: func(cctx *cli.Context) error { + log := configLogger(cctx, os.Stderr) ctx := context.Background() arg := cctx.Args().First() if arg == "" { diff --git a/cmd/gosky/debug.go b/cmd/gosky/debug.go index 8eb0110d..4ac54c74 100644 --- a/cmd/gosky/debug.go +++ b/cmd/gosky/debug.go @@ -301,6 +301,8 @@ var compareStreamsCmd = &cli.Command{ }, ArgsUsage: ``, Action: func(cctx *cli.Context) error { + log := configLogger(cctx, os.Stderr) + h1 := cctx.String("host1") h2 := cctx.String("host2") @@ -819,6 +821,7 @@ var debugCompareReposCmd = &cli.Command{ }, ArgsUsage: ``, Action: func(cctx *cli.Context) error { + log := configLogger(cctx, os.Stderr) ctx := cctx.Context did, err := syntax.ParseAtIdentifier(cctx.Args().First()) if err != nil { diff --git a/cmd/gosky/main.go b/cmd/gosky/main.go index 1c6416a3..f01ed580 100644 --- a/cmd/gosky/main.go +++ b/cmd/gosky/main.go @@ -44,8 +44,6 @@ import ( "golang.org/x/time/rate" ) -var log = slog.Default().With("system", "gosky") - func main() { run(os.Args) } @@ -77,13 +75,12 @@ func run(args []string) { Value: "https://plc.directory", EnvVars: []string{"ATP_PLC_HOST"}, }, - } - - _, _, err := cliutil.SetupSlog(cliutil.LogOptions{}) - if err != nil { - fmt.Fprintf(os.Stderr, "logging setup error: %s\n", err.Error()) - os.Exit(1) - return + &cli.StringFlag{ + Name: "log-level", + Usage: "log verbosity level (debug, info, warn, error)", + Value: "info", + EnvVars: []string{"GOSKY_LOG_LEVEL", "LOG_LEVEL", "BSKYLOG_LOG_LEVEL", "GOLOG_LOG_LEVEL"}, + }, } app.Commands = []*cli.Command{ @@ -161,6 +158,7 @@ var readRepoStreamCmd = &cli.Command{ }, ArgsUsage: `[ [cursor]]`, Action: func(cctx *cli.Context) error { + log := configLogger(cctx, os.Stderr) ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT) defer stop() @@ -849,3 +847,25 @@ var verifyUserCmd = &cli.Command{ return nil }, } + +func configLogger(cmd *cli.Context, writer *os.File) *slog.Logger { + var level slog.Level + switch cmd.String("log-level") { + case "debug": + level = slog.LevelDebug + case "info": + level = slog.LevelInfo + case "warn": + level = slog.LevelWarn + case "error": + level = slog.LevelError + default: + level = slog.LevelInfo + } + + logger := slog.New(slog.NewJSONHandler(writer, &slog.HandlerOptions{ + Level: level, + })) + + return logger +} diff --git a/cmd/gosky/streamdiff.go b/cmd/gosky/streamdiff.go index e1ab67af..845787e0 100644 --- a/cmd/gosky/streamdiff.go +++ b/cmd/gosky/streamdiff.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "net/http" + "os" comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/events" @@ -20,6 +21,7 @@ var streamCompareCmd = &cli.Command{ Flags: []cli.Flag{}, ArgsUsage: ` `, Action: func(cctx *cli.Context) error { + log := configLogger(cctx, os.Stderr) d := websocket.DefaultDialer args, err := needArgs(cctx, "hostA", "hostB") diff --git a/cmd/gosky/sync.go b/cmd/gosky/sync.go index 13c60f42..e401078d 100644 --- a/cmd/gosky/sync.go +++ b/cmd/gosky/sync.go @@ -33,6 +33,7 @@ var syncGetRepoCmd = &cli.Command{ }, }, Action: func(cctx *cli.Context) error { + log := configLogger(cctx, os.Stderr) ctx := context.Background() arg := cctx.Args().First() if arg == "" { -- 2.51.2 From 55cb019166da572f2963bf0ce8cc598dbfbdc472 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Thu, 29 Jan 2026 10:36:24 -0800 Subject: [PATCH 2/3] remove SetupSlog and supporting code --- util/cliutil/ipfslog.go | 35 ---- util/cliutil/util.go | 349 ---------------------------------------- 2 files changed, 384 deletions(-) delete mode 100644 util/cliutil/ipfslog.go diff --git a/util/cliutil/ipfslog.go b/util/cliutil/ipfslog.go deleted file mode 100644 index 4f3a0f1c..00000000 --- a/util/cliutil/ipfslog.go +++ /dev/null @@ -1,35 +0,0 @@ -package cliutil - -import ( - "io" - - ipfslog "github.com/ipfs/go-log/v2" - "go.uber.org/zap/zapcore" -) - -func SetIpfsWriter(out io.Writer, format string, level string) { - var ze zapcore.Encoder - switch format { - case "json": - ze = zapcore.NewJSONEncoder(zapcore.EncoderConfig{}) - case "text": - ze = zapcore.NewConsoleEncoder(zapcore.EncoderConfig{}) - default: - ze = zapcore.NewConsoleEncoder(zapcore.EncoderConfig{}) - } - var zl zapcore.LevelEnabler - switch level { - case "debug": - zl = zapcore.DebugLevel - case "info": - zl = zapcore.InfoLevel - case "warn": - zl = zapcore.WarnLevel - case "error": - zl = zapcore.ErrorLevel - default: - zl = zapcore.InfoLevel - } - nc := zapcore.NewCore(ze, zapcore.AddSync(out), zl) - ipfslog.SetPrimaryCore(nc) -} diff --git a/util/cliutil/util.go b/util/cliutil/util.go index 47973796..37a9cedc 100644 --- a/util/cliutil/util.go +++ b/util/cliutil/util.go @@ -2,17 +2,10 @@ package cliutil import ( "encoding/json" - "errors" "fmt" - "io" - "io/fs" - "log/slog" "net/http" "os" "path/filepath" - "regexp" - "sort" - "strconv" "strings" "time" @@ -189,345 +182,3 @@ func SetupDatabase(dburl string, maxConnections int) (*gorm.DB, error) { return db, nil } - -type LogOptions struct { - // e.g. 1_000_000_000 - LogRotateBytes int64 - - // path to write to, if rotating, %T gets UnixMilli at file open time - // NOTE: substitution is simple replace("%T", "") - LogPath string - - // text|json - LogFormat string - - // info|debug|warn|error - LogLevel string - - // Keep N old logs (not including current); <0 disables removal, 0==remove all old log files immediately - KeepOld int -} - -func firstenv(env_var_names ...string) string { - for _, env_var_name := range env_var_names { - val := os.Getenv(env_var_name) - if val != "" { - return val - } - } - return "" -} - -// SetupSlog integrates passed in options and env vars. -// -// passing default cliutil.LogOptions{} is ok. -// -// BSKYLOG_LOG_LEVEL=info|debug|warn|error -// -// BSKYLOG_LOG_FMT=text|json -// -// BSKYLOG_FILE=path (or "-" or "" for stdout), %T gets UnixMilli; if a path with '/', {prefix}/current becomes a link to active log file -// -// BSKYLOG_ROTATE_BYTES=int maximum size of log chunk before rotating -// -// BSKYLOG_ROTATE_KEEP=int keep N olg logs (not including current) -// -// The env vars were derived from ipfs logging library, and also respond to some GOLOG_ vars from that library, -// but BSKYLOG_ variables are preferred because imported code still using the ipfs log library may misbehave -// if some GOLOG values are set, especially GOLOG_FILE. -func SetupSlog(options LogOptions) (*slog.Logger, io.Writer, error) { - fmt.Fprintf(os.Stderr, "SetupSlog\n") - var hopts slog.HandlerOptions - hopts.Level = slog.LevelInfo - hopts.AddSource = true - if options.LogLevel == "" { - options.LogLevel = firstenv("BSKYLOG_LOG_LEVEL", "GOLOG_LOG_LEVEL") - } - if options.LogLevel == "" { - hopts.Level = slog.LevelInfo - options.LogLevel = "info" - } else { - level := strings.ToLower(options.LogLevel) - switch level { - case "debug": - hopts.Level = slog.LevelDebug - case "info": - hopts.Level = slog.LevelInfo - case "warn": - hopts.Level = slog.LevelWarn - case "error": - hopts.Level = slog.LevelError - default: - return nil, nil, fmt.Errorf("unknown log level: %#v", options.LogLevel) - } - } - if options.LogFormat == "" { - options.LogFormat = firstenv("BSKYLOG_LOG_FMT", "GOLOG_LOG_FMT") - } - if options.LogFormat == "" { - options.LogFormat = "text" - } else { - format := strings.ToLower(options.LogFormat) - if format == "json" || format == "text" { - // ok - } else { - return nil, nil, fmt.Errorf("invalid log format: %#v", options.LogFormat) - } - options.LogFormat = format - } - - if options.LogPath == "" { - options.LogPath = firstenv("BSKYLOG_FILE", "GOLOG_FILE") - } - if options.LogRotateBytes == 0 { - rotateBytesStr := os.Getenv("BSKYLOG_ROTATE_BYTES") // no GOLOG equivalent - if rotateBytesStr != "" { - rotateBytes, err := strconv.ParseInt(rotateBytesStr, 10, 64) - if err != nil { - return nil, nil, fmt.Errorf("invalid BSKYLOG_ROTATE_BYTES value: %w", err) - } - options.LogRotateBytes = rotateBytes - } - } - if options.KeepOld == 0 { - keepOldUnset := true - keepOldStr := os.Getenv("BSKYLOG_ROTATE_KEEP") // no GOLOG equivalent - if keepOldStr != "" { - keepOld, err := strconv.ParseInt(keepOldStr, 10, 64) - if err != nil { - return nil, nil, fmt.Errorf("invalid BSKYLOG_ROTATE_KEEP value: %w", err) - } - keepOldUnset = false - options.KeepOld = int(keepOld) - } - if keepOldUnset { - options.KeepOld = 2 - } - } - logaround := make(chan string, 100) - go logbouncer(logaround) - var out io.Writer - if (options.LogPath == "") || (options.LogPath == "-") { - out = os.Stdout - } else if options.LogRotateBytes != 0 { - out = &logRotateWriter{ - rotateBytes: options.LogRotateBytes, - outPathTemplate: options.LogPath, - keep: options.KeepOld, - logaround: logaround, - } - } else { - var err error - out, err = os.Create(options.LogPath) - if err != nil { - return nil, nil, fmt.Errorf("%s: %w", options.LogPath, err) - } - fmt.Fprintf(os.Stderr, "SetupSlog create %#v\n", options.LogPath) - } - var handler slog.Handler - switch options.LogFormat { - case "text": - handler = slog.NewTextHandler(out, &hopts) - case "json": - handler = slog.NewJSONHandler(out, &hopts) - default: - return nil, nil, fmt.Errorf("unknown log format: %#v", options.LogFormat) - } - logger := slog.New(handler) - slog.SetDefault(logger) - templateDirPart, _ := filepath.Split(options.LogPath) - ents, _ := os.ReadDir(templateDirPart) - for _, ent := range ents { - fmt.Fprintf(os.Stdout, "%s\n", filepath.Join(templateDirPart, ent.Name())) - } - SetIpfsWriter(out, options.LogFormat, options.LogLevel) - return logger, out, nil -} - -type logRotateWriter struct { - currentWriter io.WriteCloser - - // how much has been written to current log file - currentBytes int64 - - // e.g. path/to/logs/foo%T - currentPath string - - // e.g. path/to/logs/current - currentPathCurrent string - - rotateBytes int64 - - outPathTemplate string - - // keep the most recent N log files (not including current) - keep int - - // write strings to this from inside the log system, a task outside the log system hands them to slog.Info() - logaround chan<- string -} - -func logbouncer(out <-chan string) { - var logger *slog.Logger - for line := range out { - fmt.Fprintf(os.Stderr, "ll %s\n", line) - if logger == nil { - // lazy to make sure it crops up after slog Default has been set - logger = slog.Default().With("system", "logging") - } - logger.Info(line) - } -} - -var currentMatcher = regexp.MustCompile("current_\\d+") - -func (w *logRotateWriter) cleanOldLogs() { - if w.keep < 0 { - // old log removal is disabled - return - } - // w.currentPath was recently set as the new log - dirpart, _ := filepath.Split(w.currentPath) - // find old logs - templateDirPart, templateNamePart := filepath.Split(w.outPathTemplate) - if dirpart != templateDirPart { - w.logaround <- fmt.Sprintf("current dir part %#v != template dir part %#v\n", w.currentPath, w.outPathTemplate) - return - } - // build a regexp that is string literal parts with \d+ replacing the UnixMilli part - templateNameParts := strings.Split(templateNamePart, "%T") - var sb strings.Builder - first := true - for _, part := range templateNameParts { - if first { - first = false - } else { - sb.WriteString("\\d+") - } - sb.WriteString(regexp.QuoteMeta(part)) - } - tmre, err := regexp.Compile(sb.String()) - if err != nil { - w.logaround <- fmt.Sprintf("failed to compile old log template regexp: %#v\n", err) - return - } - dir, err := os.ReadDir(dirpart) - if err != nil { - w.logaround <- fmt.Sprintf("failed to read old log template dir: %#v\n", err) - return - } - var found []fs.FileInfo - for _, ent := range dir { - name := ent.Name() - if tmre.MatchString(name) || currentMatcher.MatchString(name) { - fi, err := ent.Info() - if err != nil { - continue - } - found = append(found, fi) - } - } - if len(found) <= w.keep { - // not too many, nothing to do - return - } - foundMtimeLess := func(i, j int) bool { - return found[i].ModTime().Before(found[j].ModTime()) - } - sort.Slice(found, foundMtimeLess) - drops := found[:len(found)-w.keep] - for _, fi := range drops { - fullpath := filepath.Join(dirpart, fi.Name()) - err = os.Remove(fullpath) - if err != nil { - w.logaround <- fmt.Sprintf("failed to rm old log: %#v\n", err) - // but keep going - } - // maybe it would be safe to debug-log old log removal from within the logging infrastructure? - } -} - -func (w *logRotateWriter) closeOldLog() []error { - if w.currentWriter == nil { - return nil - } - var earlyWeakErrors []error - err := w.currentWriter.Close() - if err != nil { - earlyWeakErrors = append(earlyWeakErrors, err) - } - w.currentWriter = nil - w.currentBytes = 0 - w.currentPath = "" - if w.currentPathCurrent != "" { - err = os.Remove(w.currentPathCurrent) // not really an error until something else goes wrong - if err != nil { - earlyWeakErrors = append(earlyWeakErrors, err) - } - w.currentPathCurrent = "" - } - return earlyWeakErrors -} - -func (w *logRotateWriter) openNewLog(earlyWeakErrors []error) (badErr error, weakErrors []error) { - nowMillis := time.Now().UnixMilli() - nows := strconv.FormatInt(nowMillis, 10) - w.currentPath = strings.Replace(w.outPathTemplate, "%T", nows, -1) - var err error - w.currentWriter, err = os.Create(w.currentPath) - if err != nil { - earlyWeakErrors = append(earlyWeakErrors, err) - return errors.Join(earlyWeakErrors...), nil - } - w.logaround <- fmt.Sprintf("new log file %#v", w.currentPath) - w.cleanOldLogs() - dirpart, _ := filepath.Split(w.currentPath) - if dirpart != "" { - w.currentPathCurrent = filepath.Join(dirpart, "current") - fi, err := os.Stat(w.currentPathCurrent) - if err == nil && fi.Mode().IsRegular() { - // move aside unknown "current" from a previous run - // see also currentMatcher regexp current_\d+ - err = os.Rename(w.currentPathCurrent, w.currentPathCurrent+"_"+nows) - if err != nil { - // not crucial if we can't move aside "current" - // TODO: log warning ... but not from inside log writer? - earlyWeakErrors = append(earlyWeakErrors, err) - } - } - err = os.Link(w.currentPath, w.currentPathCurrent) - if err != nil { - // not crucial if we can't make "current" link - // TODO: log warning ... but not from inside log writer? - earlyWeakErrors = append(earlyWeakErrors, err) - } - } - return nil, earlyWeakErrors -} - -func (w *logRotateWriter) Write(p []byte) (n int, err error) { - var earlyWeakErrors []error - if int64(len(p))+w.currentBytes > w.rotateBytes { - // next write would be over the limit - earlyWeakErrors = w.closeOldLog() - } - if w.currentWriter == nil { - // start new log file - var err error - err, earlyWeakErrors = w.openNewLog(earlyWeakErrors) - if err != nil { - return 0, err - } - } - var wrote int - wrote, err = w.currentWriter.Write(p) - w.currentBytes += int64(wrote) - if err != nil { - earlyWeakErrors = append(earlyWeakErrors, err) - return wrote, errors.Join(earlyWeakErrors...) - } - if earlyWeakErrors != nil { - w.logaround <- fmt.Sprintf("ok, but: %s", errors.Join(earlyWeakErrors...).Error()) - } - return wrote, nil -} -- 2.51.2 From b1401ec901f4114f22c744afa749088804cb49b3 Mon Sep 17 00:00:00 2001 From: bryan newbold Date: Thu, 29 Jan 2026 10:36:34 -0800 Subject: [PATCH 3/3] go mod tidy --- go.mod | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/go.mod b/go.mod index 33b11779..b358cda7 100644 --- a/go.mod +++ b/go.mod @@ -29,7 +29,6 @@ require ( github.com/ipfs/go-ipld-cbor v0.1.0 github.com/ipfs/go-ipld-format v0.6.0 github.com/ipfs/go-libipfs v0.7.0 - github.com/ipfs/go-log/v2 v2.5.1 github.com/ipld/go-car v0.6.1-0.20230509095817-92d28eb23ba4 github.com/ipld/go-car/v2 v2.13.1 github.com/jackc/pgx/v5 v5.5.0 @@ -64,7 +63,6 @@ require ( go.opentelemetry.io/otel/sdk v1.21.0 go.opentelemetry.io/otel/trace v1.21.0 go.uber.org/automaxprocs v1.5.3 - go.uber.org/zap v1.26.0 golang.org/x/sync v0.7.0 golang.org/x/text v0.14.0 golang.org/x/time v0.3.0 @@ -93,6 +91,7 @@ require ( github.com/hailocab/go-hostpool v0.0.0-20160125115350-e80d13ce29ed // indirect github.com/hashicorp/golang-lru v1.0.2 // indirect github.com/ipfs/go-log v1.0.5 // indirect + github.com/ipfs/go-log/v2 v2.5.1 // indirect github.com/jackc/puddle/v2 v2.2.1 // indirect github.com/klauspost/compress v1.17.3 // indirect github.com/kr/pretty v0.3.1 // indirect @@ -106,6 +105,7 @@ require ( github.com/vmihailenco/msgpack/v5 v5.4.1 // indirect github.com/vmihailenco/tagparser/v2 v2.0.0 // indirect github.com/whyrusleeping/cbor v0.0.0-20171005072247-63513f603b11 // indirect + go.uber.org/zap v1.26.0 // indirect golang.org/x/crypto v0.21.0 // indirect golang.org/x/exp v0.0.0-20231110203233-9a3e6036ecaa // indirect gopkg.in/inf.v0 v0.9.1 // indirect -- 2.51.2