diff --git a/knotserver/file.go b/knotserver/file.go --- a/knotserver/file.go +++ b/knotserver/file.go @@ -3,7 +3,7 @@ import ( "bytes" "io" - "log" + "log/slog" "net/http" "strings" @@ -43,11 +43,11 @@ } } -func (h *Handle) showFile(content string, data map[string]any, w http.ResponseWriter) { +func (h *Handle) showFile(content string, data map[string]any, w http.ResponseWriter, l *slog.Logger) { lc, err := countLines(strings.NewReader(content)) if err != nil { // Non-fatal, we'll just skip showing line numbers in the template. - log.Printf("counting lines: %s", err) + l.Warn("counting lines", "error", err) } lines := make([]int, lc) diff --git a/knotserver/git.go b/knotserver/git.go --- a/knotserver/git.go +++ b/knotserver/git.go @@ -3,7 +3,6 @@ import ( "compress/gzip" "io" - "log" "net/http" "path/filepath" @@ -26,7 +25,7 @@ if err := cmd.InfoRefs(); err != nil { http.Error(w, err.Error(), 500) - log.Printf("git: failed to execute git-upload-pack (info/refs) %s", err) + d.l.Error("git: failed to execute git-upload-pack (info/refs)", "handler", "InfoRefs", "error", err) return } } @@ -53,7 +52,7 @@ reader, err := gzip.NewReader(r.Body) if err != nil { http.Error(w, err.Error(), 500) - log.Printf("git: failed to create gzip reader: %s", err) + d.l.Error("git: failed to create gzip reader", "handler", "UploadPack", "error", err) return } defer reader.Close() @@ -62,7 +61,7 @@ cmd.Stdin = reader if err := cmd.UploadPack(); err != nil { http.Error(w, err.Error(), 500) - log.Printf("git: failed to execute git-upload-pack %s", err) + d.l.Error("git: failed to execute git-upload-pack", "handler", "UploadPack", "error", err) return } } diff --git a/knotserver/handler.go b/knotserver/handler.go --- a/knotserver/handler.go +++ b/knotserver/handler.go @@ -3,6 +3,7 @@ import ( "context" "fmt" + "log/slog" "net/http" "github.com/go-chi/chi/v5" @@ -21,6 +22,7 @@ db *db.DB js *jsclient.JetstreamClient e *rbac.Enforcer + l *slog.Logger // init is a channel that is closed when the knot has been initailized // i.e. when the first user (knot owner) has been added. @@ -28,13 +30,14 @@ knotInitialized bool } -func Setup(ctx context.Context, c *config.Config, db *db.DB, e *rbac.Enforcer) (http.Handler, error) { +func Setup(ctx context.Context, c *config.Config, db *db.DB, e *rbac.Enforcer, l *slog.Logger) (http.Handler, error) { r := chi.NewRouter() h := Handle{ c: c, db: db, e: e, + l: l, init: make(chan struct{}), } diff --git a/knotserver/jetstream.go b/knotserver/jetstream.go --- a/knotserver/jetstream.go +++ b/knotserver/jetstream.go @@ -5,7 +5,6 @@ "encoding/json" "fmt" "io" - "log" "net/http" "net/url" "strings" @@ -14,13 +13,16 @@ "github.com/sotangled/tangled/api/tangled" "github.com/sotangled/tangled/knotserver/db" "github.com/sotangled/tangled/knotserver/jsclient" + "github.com/sotangled/tangled/log" ) func (h *Handle) StartJetstream(ctx context.Context) error { + l := h.l.With("component", "jetstream") + ctx = log.IntoContext(ctx, l) collections := []string{tangled.PublicKeyNSID, tangled.KnotMemberNSID} dids := []string{} - lastTimeUs, err := h.getLastTimeUs() + lastTimeUs, err := h.getLastTimeUs(ctx) if err != nil { return err } @@ -31,58 +33,63 @@ return fmt.Errorf("failed to read from jetstream: %w", err) } - go h.processMessages(messages) + go h.processMessages(ctx, messages) return nil } -func (h *Handle) getLastTimeUs() (int64, error) { +func (h *Handle) getLastTimeUs(ctx context.Context) (int64, error) { + l := log.FromContext(ctx) lastTimeUs, err := h.db.GetLastTimeUs() if err != nil { - log.Println("couldn't get last time us, starting from now") + l.Info("couldn't get last time us, starting from now") lastTimeUs = time.Now().UnixMicro() } // If last time is older than a week, start from now if time.Now().UnixMicro()-lastTimeUs > 7*24*60*60*1000*1000 { lastTimeUs = time.Now().UnixMicro() - log.Printf("last time us is older than a week. discarding that and starting from now.") + l.Info("last time us is older than a week. discarding that and starting from now") err = h.db.SaveLastTimeUs(lastTimeUs) if err != nil { - log.Println("failed to save last time us") + l.Error("failed to save last time us") } } - log.Printf("found last time_us %d", lastTimeUs) + l.Info("found last time_us", "time_us", lastTimeUs) return lastTimeUs, nil } -func (h *Handle) processPublicKey(did string, record map[string]interface{}) { +func (h *Handle) processPublicKey(ctx context.Context, did string, record map[string]interface{}) error { + l := log.FromContext(ctx) if err := h.db.AddPublicKeyFromRecord(did, record); err != nil { - log.Printf("failed to add public key: %v", err) - } else { - log.Printf("added public key from firehose: %s", did) + l.Error("failed to add public key", "error", err) + return fmt.Errorf("failed to add public key: %w", err) } + l.Info("added public key from firehose", "did", did) + return nil } -func (h *Handle) fetchAndAddKeys(did string) { +func (h *Handle) fetchAndAddKeys(ctx context.Context, did string) error { + l := log.FromContext(ctx) + keysEndpoint, err := url.JoinPath(h.c.AppViewEndpoint, "keys", did) if err != nil { - log.Printf("error building endpoint url: %s: %v", did, err) - return + l.Error("error building endpoint url", "did", did, "error", err.Error()) + return fmt.Errorf("error building endpoint url: %w", err) } resp, err := http.Get(keysEndpoint) if err != nil { - log.Printf("error getting keys for %s: %v", did, err) - return + l.Error("error getting keys", "did", did, "error", err) + return fmt.Errorf("error getting keys: %w", err) } defer resp.Body.Close() plaintext, err := io.ReadAll(resp.Body) if err != nil { - log.Printf("error reading response body: %v", err) - return + l.Error("error reading response body", "error", err) + return fmt.Errorf("error reading response body: %w", err) } for _, key := range strings.Split(string(plaintext), "\n") { @@ -94,37 +101,51 @@ } pk.Key = key if err := h.db.AddPublicKey(pk); err != nil { - log.Printf("failed to add public key: %v", err) + l.Error("failed to add public key", "error", err) + return fmt.Errorf("failed to add public key: %w", err) } } + return nil } -func (h *Handle) processKnotMember(did string, record map[string]interface{}) { +func (h *Handle) processKnotMember(ctx context.Context, did string, record map[string]interface{}) error { + l := log.FromContext(ctx) ok, err := h.e.E.Enforce(did, ThisServer, ThisServer, "server:invite") if err != nil || !ok { - log.Printf("failed to add member from did %s", did) - return + l.Error("failed to add member", "did", did) + return fmt.Errorf("failed to enforce permissions: %w", err) } - log.Printf("adding member") + l.Info("adding member") if err := h.e.AddMember(ThisServer, record["member"].(string)); err != nil { - log.Printf("failed to add member: %v", err) - } else { - log.Printf("added member from firehose: %s", record["member"]) + l.Error("failed to add member", "error", err) + return fmt.Errorf("failed to add member: %w", err) + } + l.Info("added member from firehose", "member", record["member"]) + + if err := h.db.AddDid(did); err != nil { + l.Error("failed to add did", "error", err) + return fmt.Errorf("failed to add did: %w", err) } - h.fetchAndAddKeys(did) + if err := h.fetchAndAddKeys(ctx, did); err != nil { + return fmt.Errorf("failed to fetch and add keys: %w", err) + } + h.js.UpdateDids([]string{did}) + return nil } -func (h *Handle) processMessages(messages <-chan []byte) { +func (h *Handle) processMessages(ctx context.Context, messages <-chan []byte) { + l := log.FromContext(ctx) + l.Info("waiting for knot to be initialized") <-h.init - log.Println("initalized jetstream watcher") + l.Info("initialized jetstream watcher") for msg := range messages { var data map[string]interface{} if err := json.Unmarshal(msg, &data); err != nil { - log.Printf("error unmarshaling message: %v", err) + l.Error("error unmarshaling message", "error", err) continue } @@ -133,16 +154,27 @@ did := data["did"].(string) record := commit["record"].(map[string]interface{}) + var processErr error switch commit["collection"].(string) { case tangled.PublicKeyNSID: - h.processPublicKey(did, record) + if err := h.processPublicKey(ctx, did, record); err != nil { + processErr = fmt.Errorf("failed to process public key: %w", err) + } case tangled.KnotMemberNSID: - h.processKnotMember(did, record) + if err := h.processKnotMember(ctx, did, record); err != nil { + processErr = fmt.Errorf("failed to process knot member: %w", err) + } + } + + if processErr != nil { + l.Error("error processing message", "error", processErr) + continue } lastTimeUs := int64(data["time_us"].(float64)) if err := h.db.SaveLastTimeUs(lastTimeUs); err != nil { - log.Printf("failed to save last time us: %v", err) + l.Error("failed to save last time us", "error", err) + continue } } } diff --git a/knotserver/routes.go b/knotserver/routes.go --- a/knotserver/routes.go +++ b/knotserver/routes.go @@ -9,7 +9,6 @@ "errors" "fmt" "html/template" - "log" "net/http" "path/filepath" "strconv" @@ -30,6 +29,7 @@ func (h *Handle) RepoIndex(w http.ResponseWriter, r *http.Request) { path := filepath.Join(h.c.Repo.ScanPath, didPath(r)) + l := h.l.With("path", path, "handler", "RepoIndex") gr, err := git.Open(path, "") if err != nil { @@ -37,7 +37,7 @@ writeMsg(w, "repo empty") return } else { - log.Println(err) + l.Error("opening repo", "error", err.Error()) notFound(w) return } @@ -45,7 +45,7 @@ commits, err := gr.Commits() if err != nil { writeError(w, err.Error(), http.StatusInternalServerError) - log.Println(err) + l.Error("fetching commits", "error", err.Error()) return } @@ -73,13 +73,13 @@ } if readmeContent == "" { - log.Printf("no readme found for %s", path) + l.Warn("no readme found") } mainBranch, err := gr.FindMainBranch(h.c.Repo.MainBranch) if err != nil { writeError(w, err.Error(), http.StatusInternalServerError) - log.Println(err) + l.Error("finding main branch", "error", err.Error()) return } @@ -100,6 +100,8 @@ treePath := chi.URLParam(r, "*") ref := chi.URLParam(r, "ref") + l := h.l.With("handler", "RepoTree", "ref", ref, "treePath", treePath) + path := filepath.Join(h.c.Repo.ScanPath, didPath(r)) gr, err := git.Open(path, ref) if err != nil { @@ -110,7 +112,7 @@ files, err := gr.FileTree(treePath) if err != nil { writeError(w, err.Error(), http.StatusInternalServerError) - log.Println(err) + l.Error("file tree", "error", err.Error()) return } @@ -132,6 +134,8 @@ treePath := chi.URLParam(r, "*") ref := chi.URLParam(r, "ref") + + l := h.l.With("handler", "FileContent", "ref", ref, "treePath", treePath) path := filepath.Join(h.c.Repo.ScanPath, didPath(r)) gr, err := git.Open(path, ref) @@ -155,13 +159,15 @@ if raw { h.showRaw(string(safe), w) } else { - h.showFile(string(safe), data, w) + h.showFile(string(safe), data, w, l) } } func (h *Handle) Archive(w http.ResponseWriter, r *http.Request) { name := chi.URLParam(r, "name") file := chi.URLParam(r, "file") + + l := h.l.With("handler", "Archive", "name", name, "file", file) // TODO: extend this to add more files compression (e.g.: xz) if !strings.HasSuffix(file, ".tar.gz") { @@ -192,7 +198,7 @@ if err != nil { // once we start writing to the body we can't report error anymore // so we are only left with printing the error. - log.Println(err) + l.Error("writing tar file", "error", err.Error()) return } @@ -200,16 +206,17 @@ if err != nil { // once we start writing to the body we can't report error anymore // so we are only left with printing the error. - log.Println(err) + l.Error("flushing?", "error", err.Error()) return } } func (h *Handle) Log(w http.ResponseWriter, r *http.Request) { - fmt.Println(r.URL.Path) ref := chi.URLParam(r, "ref") - path := filepath.Join(h.c.Repo.ScanPath, didPath(r)) + + l := h.l.With("handler", "Log", "ref", ref, "path", path) + gr, err := git.Open(path, ref) if err != nil { notFound(w) @@ -219,7 +226,7 @@ commits, err := gr.Commits() if err != nil { writeError(w, err.Error(), http.StatusInternalServerError) - log.Println(err) + l.Error("fetching commits", "error", err.Error()) return } @@ -269,6 +276,8 @@ func (h *Handle) Diff(w http.ResponseWriter, r *http.Request) { ref := chi.URLParam(r, "ref") + l := h.l.With("handler", "Diff", "ref", ref) + path := filepath.Join(h.c.Repo.ScanPath, didPath(r)) gr, err := git.Open(path, ref) if err != nil { @@ -279,7 +288,7 @@ diff, err := gr.Diff() if err != nil { writeError(w, err.Error(), http.StatusInternalServerError) - log.Println(err) + l.Error("getting diff", "error", err.Error()) return } @@ -297,6 +306,8 @@ func (h *Handle) Refs(w http.ResponseWriter, r *http.Request) { path := filepath.Join(h.c.Repo.ScanPath, didPath(r)) + l := h.l.With("handler", "Refs") + gr, err := git.Open(path, "") if err != nil { notFound(w) @@ -306,12 +317,12 @@ tags, err := gr.Tags() if err != nil { // Non-fatal, we *should* have at least one branch to show. - log.Println(err) + l.Error("getting tags", "error", err.Error()) } branches, err := gr.Branches() if err != nil { - log.Println(err) + l.Error("getting branches", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) return } @@ -327,12 +338,14 @@ } func (h *Handle) Keys(w http.ResponseWriter, r *http.Request) { + l := h.l.With("handler", "Keys") + switch r.Method { case http.MethodGet: keys, err := h.db.GetAllPublicKeys() if err != nil { writeError(w, err.Error(), http.StatusInternalServerError) - log.Println(err) + l.Error("getting public keys", "error", err.Error()) return } @@ -358,7 +371,7 @@ if err := h.db.AddPublicKey(pk); err != nil { writeError(w, err.Error(), http.StatusInternalServerError) - log.Printf("adding public key: %s", err) + l.Error("adding public key", "error", err.Error()) return } @@ -368,6 +381,8 @@ } func (h *Handle) NewRepo(w http.ResponseWriter, r *http.Request) { + l := h.l.With("handler", "NewRepo") + data := struct { Did string `json:"did"` Name string `json:"name"` @@ -385,7 +400,7 @@ repoPath := filepath.Join(h.c.Repo.ScanPath, relativeRepoPath) err := git.InitBare(repoPath) if err != nil { - log.Println(err) + l.Error("initializing bare repo", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) return } @@ -393,7 +408,7 @@ // add perms for this user to access the repo err = h.e.AddRepo(did, ThisServer, relativeRepoPath) if err != nil { - log.Println(err) + l.Error("adding repo permissions", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) return } @@ -402,6 +417,8 @@ } func (h *Handle) AddMember(w http.ResponseWriter, r *http.Request) { + l := h.l.With("handler", "AddMember") + data := struct { Did string `json:"did"` PublicKeys []string `json:"keys"` @@ -427,7 +444,7 @@ h.js.UpdateDids([]string{did}) if err := h.e.AddMember(ThisServer, did); err != nil { - log.Println(err) + l.Error("adding member", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) return } @@ -436,6 +453,8 @@ } func (h *Handle) Init(w http.ResponseWriter, r *http.Request) { + l := h.l.With("handler", "Init") + if h.knotInitialized { writeError(w, "knot already initialized", http.StatusConflict) return @@ -447,11 +466,13 @@ }{} if err := json.NewDecoder(r.Body).Decode(&data); err != nil { + l.Error("failed to decode request body", "error", err.Error()) writeError(w, "invalid request body", http.StatusBadRequest) return } if data.Did == "" { + l.Error("empty DID in request") writeError(w, "did is empty", http.StatusBadRequest) return } @@ -464,22 +485,24 @@ pk.Key = k err := h.db.AddPublicKey(pk) if err != nil { + l.Error("failed to add public key", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) return } } } else { + l.Error("failed to add DID", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) return } h.js.UpdateDids([]string{data.Did}) if err := h.e.AddOwner(ThisServer, data.Did); err != nil { - log.Println(err) + l.Error("adding owner", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) return } - // Signal that the knot is ready + close(h.init) mac := hmac.New(sha256.New, []byte(h.c.Server.Secret)) diff --git a/log/log.go b/log/log.go new file mode 100644 --- /dev/null +++ b/log/log.go @@ -0,0 +1,49 @@ +package log + +import ( + "context" + "log/slog" + "os" +) + +// NewHandler sets up a new slog.Handler with the service name +// as an attribute +func NewHandler(name string) slog.Handler { + handler := slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{}) + + var attrs []slog.Attr + attrs = append(attrs, slog.Attr{Key: "service", Value: slog.StringValue(name)}) + handler.WithAttrs(attrs) + return handler +} + +func New(name string) *slog.Logger { + return slog.New(NewHandler(name)) +} + +func NewContext(ctx context.Context, name string) context.Context { + return IntoContext(ctx, New(name)) +} + +type ctxKey struct{} + +// IntoContext adds a logger to a context. Use FromContext to +// pull the logger out. +func IntoContext(ctx context.Context, l *slog.Logger) context.Context { + return context.WithValue(ctx, ctxKey{}, l) +} + +// FromContext returns a logger from a context.Context; +// if the passed context is nil, we return the default slog +// logger. +func FromContext(ctx context.Context) *slog.Logger { + if ctx != nil { + v := ctx.Value(ctxKey{}) + if v == nil { + return slog.Default() + } + return v.(*slog.Logger) + } + + return slog.Default() +} diff --git a/cmd/knotserver/main.go b/cmd/knotserver/main.go --- a/cmd/knotserver/main.go +++ b/cmd/knotserver/main.go @@ -3,14 +3,12 @@ import ( "context" "fmt" - "log" - "log/slog" "net/http" - "os" "github.com/sotangled/tangled/knotserver" "github.com/sotangled/tangled/knotserver/config" "github.com/sotangled/tangled/knotserver/db" + "github.com/sotangled/tangled/log" "github.com/sotangled/tangled/rbac" ) @@ -19,34 +17,39 @@ // ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) // defer stop() - slog.SetDefault(slog.New(slog.NewTextHandler(os.Stdout, nil))) + l := log.New("knotserver") c, err := config.Load(ctx) if err != nil { - log.Fatal(err) + l.Error("failed to load config", "error", err) + return } if c.Server.Dev { - log.Println("running in dev mode, signature verification is disabled") + l.Info("running in dev mode, signature verification is disabled") } db, err := db.Setup(c.Server.DBPath) if err != nil { - log.Fatalf("failed to setup db: %s", err) + l.Error("failed to setup db", "error", err) + return } e, err := rbac.NewEnforcer(c.Server.DBPath) if err != nil { - log.Fatalf("failed to setup rbac enforcer: %s", err) + l.Error("failed to setup rbac enforcer", "error", err) + return } - mux, err := knotserver.Setup(ctx, c, db, e) + mux, err := knotserver.Setup(ctx, c, db, e, l) if err != nil { - log.Fatal(err) + l.Error("failed to setup server", "error", err) + return } addr := fmt.Sprintf("%s:%d", c.Server.Host, c.Server.Port) - log.Println("starting main server on", addr) - log.Fatal(http.ListenAndServe(addr, mux)) + l.Info("starting main server", "address", addr) + l.Error("server error", "error", http.ListenAndServe(addr, mux)) + return }