From 29583954eba742633de2ec533e44ba40903d9c35 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Mon, 11 May 2026 16:49:28 +0000 Subject: [PATCH] knotmirror: add inflight api Signed-off-by: Seongmin Lee --- knotmirror/adminpage.go | 15 ++++++++++++++- knotmirror/knotmirror.go | 2 +- knotmirror/xrpc/inflight.go | 91 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotmirror/xrpc/xrpc.go | 36 +++++++++++++++++++++--------------- 4 file(s) changed, 127 insertion(s)(+), 17 deletion(s)(-) diff --git a/knotmirror/adminpage.go b/knotmirror/adminpage.go --- a/knotmirror/adminpage.go +++ b/knotmirror/adminpage.go @@ -3,6 +3,7 @@ import ( "database/sql" "embed" + "encoding/json" "fmt" "html" "html/template" @@ -16,6 +17,7 @@ "github.com/go-chi/chi/v5" "tangled.org/core/appview/pagination" "tangled.org/core/knotmirror/db" "tangled.org/core/knotmirror/models" + "tangled.org/core/knotmirror/xrpc" ) //go:embed templates/*.html @@ -26,13 +28,15 @@ type AdminServer struct { db *sql.DB resyncer *Resyncer + xrpc *xrpc.Xrpc logger *slog.Logger } -func NewAdminServer(l *slog.Logger, database *sql.DB, resyncer *Resyncer) *AdminServer { +func NewAdminServer(l *slog.Logger, database *sql.DB, resyncer *Resyncer, x *xrpc.Xrpc) *AdminServer { return &AdminServer{ db: database, resyncer: resyncer, + xrpc: x, logger: l, } } @@ -44,7 +48,16 @@ r.Get("/hosts", s.handleHosts()) r.Post("/api/triggerRepoResync", s.handleRepoResyncTrigger()) r.Post("/api/cancelRepoResync", s.handleRepoResyncCancel()) + r.Get("/api/inflight", s.handleInflight()) return r +} + +func (s *AdminServer) handleInflight() http.HandlerFunc { + return func(w http.ResponseWriter, r *http.Request) { + entries := s.xrpc.Inflight() + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(entries) + } } func funcmap() template.FuncMap { diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go --- a/knotmirror/knotmirror.go +++ b/knotmirror/knotmirror.go @@ -55,8 +55,8 @@ knotstream := knotstream.NewKnotStream(logger, db, cfg) crawler := NewCrawler(logger, db) resyncer := NewResyncer(logger, db, gitm, cfg) - adminpage := NewAdminServer(logger, db, resyncer) xrpc := xrpc.New(logger, cfg, db, rdb, resolver, knotstream) + adminpage := NewAdminServer(logger, db, resyncer, xrpc) // maintain repository list with tap // NOTE: this can be removed once we introduce did-for-repo because then we can just listen to KnotStream for #identity events. diff --git a/knotmirror/xrpc/inflight.go b/knotmirror/xrpc/inflight.go new file mode 100644 --- /dev/null +++ b/knotmirror/xrpc/inflight.go @@ -0,0 +1,91 @@ +package xrpc + +import ( + "net/http" + "strings" + "sync" + "sync/atomic" + "time" +) + +type InflightEntry struct { + ID uint64 `json:"id"` + Method string `json:"method"` + Path string `json:"path"` + RawQuery string `json:"raw_query,omitempty"` + URL string `json:"url"` + Repo string `json:"repo,omitempty"` + RemoteAddr string `json:"remote_addr,omitempty"` + StartedAt time.Time `json:"started_at"` + DurationMs int64 `json:"duration_ms"` +} + +type inflightTracker struct { + mu sync.Mutex + nextID atomic.Uint64 + active map[uint64]*InflightEntry +} + +func newInflightTracker() *inflightTracker { + return &inflightTracker{active: make(map[uint64]*InflightEntry)} +} + +func (t *inflightTracker) add(r *http.Request) uint64 { + id := t.nextID.Add(1) + q := r.URL.Query() + entry := &InflightEntry{ + ID: id, + Method: r.Method, + Path: r.URL.Path, + RawQuery: r.URL.RawQuery, + URL: r.URL.RequestURI(), + Repo: q.Get("repo"), + RemoteAddr: clientAddr(r), + StartedAt: time.Now(), + } + t.mu.Lock() + t.active[id] = entry + t.mu.Unlock() + return id +} + +func (t *inflightTracker) remove(id uint64) { + t.mu.Lock() + delete(t.active, id) + t.mu.Unlock() +} + +func (t *inflightTracker) snapshot() []InflightEntry { + t.mu.Lock() + defer t.mu.Unlock() + out := make([]InflightEntry, 0, len(t.active)) + now := time.Now() + for _, e := range t.active { + cp := *e + cp.DurationMs = now.Sub(e.StartedAt).Milliseconds() + out = append(out, cp) + } + return out +} + +func clientAddr(r *http.Request) string { + if xff := r.Header.Get("X-Forwarded-For"); xff != "" { + if i := strings.IndexByte(xff, ','); i >= 0 { + return strings.TrimSpace(xff[:i]) + } + return strings.TrimSpace(xff) + } + return r.RemoteAddr +} + +func (t *inflightTracker) middleware(next http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + id := t.add(r) + defer t.remove(id) + next.ServeHTTP(w, r) + }) +} + +func (x *Xrpc) Inflight() []InflightEntry { + return x.inflight.snapshot() +} diff --git a/knotmirror/xrpc/xrpc.go b/knotmirror/xrpc/xrpc.go --- a/knotmirror/xrpc/xrpc.go +++ b/knotmirror/xrpc/xrpc.go @@ -26,6 +26,7 @@ resolver *idresolver.Resolver ks *knotstream.KnotStream logger *slog.Logger httpClient *http.Client + inflight *inflightTracker } func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, rdb *redis.Client, resolver *idresolver.Resolver, ks *knotstream.KnotStream) *Xrpc { @@ -39,27 +40,32 @@ logger: log.SubLogger(logger, "xrpc"), httpClient: &http.Client{ Timeout: 30 * time.Second, }, + inflight: newInflightTracker(), } } func (x *Xrpc) Router() http.Handler { r := chi.NewRouter() - r.Get("/"+tangled.GitTempGetArchiveNSID, x.GetArchive) - r.Get("/"+tangled.GitTempGetBlobNSID, x.GetBlob) - r.Get("/"+tangled.GitTempGetBranchNSID, x.GetBranch) - // r.Get("/"+tangled.GitTempGetCommitNSID, x.GetCommit) // todo - // r.Get("/"+tangled.GitTempGetDiffNSID, x.GetDiff) // todo - // r.Get("/"+tangled.GitTempGetEntityNSID, x.GetEntity) // todo - // r.Get("/"+tangled.GitTempGetHeadNSID, x.GetHead) // todo - r.Get("/"+tangled.GitTempGetTagNSID, x.GetTag) // using types.Response - r.Get("/"+tangled.GitTempGetTreeNSID, x.GetTree) - r.Get("/"+tangled.GitTempListBranchesNSID, x.ListBranches) // wip, unknown output - r.Get("/"+tangled.GitTempListCommitsNSID, x.ListCommits) - r.Get("/"+tangled.GitTempListLanguagesNSID, x.ListLanguages) - r.Get("/"+tangled.GitTempListTagsNSID, x.ListTags) - r.Get("/"+tangled.RepoBlobNSID, x.RepoBlob) - r.Post("/"+tangled.SyncRequestCrawlNSID, x.RequestCrawl) + r.Group(func(r chi.Router) { + r.Use(x.inflight.middleware) + + r.Get("/"+tangled.GitTempGetArchiveNSID, x.GetArchive) + r.Get("/"+tangled.GitTempGetBlobNSID, x.GetBlob) + r.Get("/"+tangled.GitTempGetBranchNSID, x.GetBranch) + // r.Get("/"+tangled.GitTempGetCommitNSID, x.GetCommit) // todo + // r.Get("/"+tangled.GitTempGetDiffNSID, x.GetDiff) // todo + // r.Get("/"+tangled.GitTempGetEntityNSID, x.GetEntity) // todo + // r.Get("/"+tangled.GitTempGetHeadNSID, x.GetHead) // todo + r.Get("/"+tangled.GitTempGetTagNSID, x.GetTag) // using types.Response + r.Get("/"+tangled.GitTempGetTreeNSID, x.GetTree) + r.Get("/"+tangled.GitTempListBranchesNSID, x.ListBranches) // wip, unknown output + r.Get("/"+tangled.GitTempListCommitsNSID, x.ListCommits) + r.Get("/"+tangled.GitTempListLanguagesNSID, x.ListLanguages) + r.Get("/"+tangled.GitTempListTagsNSID, x.ListTags) + r.Get("/"+tangled.RepoBlobNSID, x.RepoBlob) + r.Post("/"+tangled.SyncRequestCrawlNSID, x.RequestCrawl) + }) return r } -- tangled.sh