From 01cee5e83f8c90b8b7fe50e0dce99197e1919647 Mon Sep 17 00:00:00 2001 From: oppiliappan Date: Tue, 27 May 2025 16:40:14 +0100 Subject: [PATCH] knotserver: implement internal endpoint to record git push records git operations in an "oplog" table. this table can be routinely cleaned up as necessary. Signed-off-by: oppiliappan --- knotserver/db/init.go | 9 ++++++ knotserver/db/oplog.go | 63 ++++++++++++++++++++++++++++++++++++++++++ knotserver/internal.go | 56 +++++++++++++++++++++++++++++++++++++ knotserver/util.go | 7 +++++ 4 files changed, 135 insertions(+) create mode 100644 knotserver/db/oplog.go diff --git a/knotserver/db/init.go b/knotserver/db/init.go index 0f74fd20..ec7233b1 100644 --- a/knotserver/db/init.go +++ b/knotserver/db/init.go @@ -43,6 +43,15 @@ func Setup(dbPath string) (*DB, error) { id integer primary key autoincrement, last_time_us integer not null ); + + create table if not exists oplog ( + tid text primary key, + did text not null, + repo text not null, + old_sha text not null, + new_sha text not null, + ref text not null + ); `) if err != nil { return nil, err diff --git a/knotserver/db/oplog.go b/knotserver/db/oplog.go new file mode 100644 index 00000000..7672a9d5 --- /dev/null +++ b/knotserver/db/oplog.go @@ -0,0 +1,63 @@ +package db + +import ( + "fmt" +) + +type Op struct { + Tid string // time based ID, easy to enumerate & monotonic + Did string // did of pusher + Repo string // fully qualified repo + OldSha string // old sha of reference being updated + NewSha string // new sha of reference being updated + Ref string // the reference being updated +} + +func (d *DB) InsertOp(op Op) error { + _, err := d.db.Exec( + `insert into oplog (tid, did, repo, old_sha, new_sha, ref) values (?, ?, ?, ?, ?, ?)`, + op.Tid, + op.Did, + op.Repo, + op.OldSha, + op.NewSha, + op.Ref, + ) + return err +} + +func (d *DB) GetOps(cursor string) ([]Op, error) { + whereClause := "" + args := []any{} + if cursor != "" { + whereClause = "where tid > ?" + args = append(args, cursor) + } + + query := fmt.Sprintf(` + select tid, did, repo, old_sha, new_sha, ref + from oplog + %s + order by tid asc + limit 100 + `, whereClause) + + rows, err := d.db.Query(query, args...) + if err != nil { + return nil, err + } + defer rows.Close() + + var ops []Op + for rows.Next() { + var op Op + rows.Scan(&op.Tid, &op.Did, &op.Repo, &op.OldSha, &op.NewSha, &op.Ref) + ops = append(ops, op) + } + + if err := rows.Err(); err != nil { + return nil, err + } + + return ops, nil +} diff --git a/knotserver/internal.go b/knotserver/internal.go index 6359852c..0fcfb2bf 100644 --- a/knotserver/internal.go +++ b/knotserver/internal.go @@ -1,9 +1,12 @@ package knotserver import ( + "bufio" "context" "log/slog" "net/http" + "path/filepath" + "strings" "github.com/go-chi/chi/v5" "github.com/go-chi/chi/v5/middleware" @@ -54,6 +57,58 @@ func (h *InternalHandle) InternalKeys(w http.ResponseWriter, r *http.Request) { return } +func (h *InternalHandle) PostReceiveHook(w http.ResponseWriter, r *http.Request) { + l := h.l.With("handler", "PostReceiveHook") + + gitAbsoluteDir := r.Header.Get("X-Git-Dir") + gitRelativeDir, err := filepath.Rel(h.c.Repo.ScanPath, gitAbsoluteDir) + if err != nil { + l.Error("failed to calculate relative git dir", "scanPath", h.c.Repo.ScanPath, "gitAbsoluteDir", gitAbsoluteDir) + return + } + gitUserDid := r.Header.Get("X-Git-User-Did") + + var ops []db.Op + scanner := bufio.NewScanner(r.Body) + for scanner.Scan() { + line := scanner.Text() + parts := strings.SplitN(line, " ", 3) + if len(parts) != 3 { + l.Error("invalid payload", "parts", parts) + continue + } + + tid := TID() + oldSha := parts[0] + newSha := parts[1] + ref := parts[2] + op := db.Op{ + Tid: tid, + Did: gitUserDid, + Repo: gitRelativeDir, + OldSha: oldSha, + NewSha: newSha, + Ref: ref, + } + ops = append(ops, op) + } + + if err := scanner.Err(); err != nil { + l.Error("failed to read payload", "err", err) + return + } + + for _, op := range ops { + err := h.db.InsertOp(op) + if err != nil { + l.Error("failed to insert op", "err", err, "op", op) + continue + } + } + + return +} + func Internal(ctx context.Context, c *config.Config, db *db.DB, e *rbac.Enforcer, l *slog.Logger) http.Handler { r := chi.NewRouter() @@ -66,6 +121,7 @@ func Internal(ctx context.Context, c *config.Config, db *db.DB, e *rbac.Enforcer r.Get("/push-allowed", h.PushAllowed) r.Get("/keys", h.InternalKeys) + r.Post("/hooks/post-receive", h.PostReceiveHook) r.Mount("/debug", middleware.Profiler()) return r diff --git a/knotserver/util.go b/knotserver/util.go index 947a7d6d..c7e8fa3f 100644 --- a/knotserver/util.go +++ b/knotserver/util.go @@ -5,6 +5,7 @@ import ( "os" "path/filepath" + "github.com/bluesky-social/indigo/atproto/syntax" securejoin "github.com/cyphar/filepath-securejoin" "github.com/go-chi/chi/v5" "github.com/microcosm-cc/bluemonday" @@ -43,3 +44,9 @@ func setGZipMIME(w http.ResponseWriter) { func setMIME(w http.ResponseWriter, mime string) { w.Header().Add("Content-Type", mime) } + +var TIDClock = syntax.NewTIDClock(0) + +func TID() string { + return TIDClock.Next().String() +} -- 2.51.2