This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knotmirror / knotmirror.go
4.2 kB 145 lines
1package knotmirror 2 3import ( 4 "context" 5 "fmt" 6 "net/http" 7 _ "net/http/pprof" 8 "time" 9 10 "github.com/go-chi/chi/v5" 11 "github.com/prometheus/client_golang/prometheus/promhttp" 12 "github.com/redis/go-redis/v9" 13 "tangled.org/core/idresolver" 14 "tangled.org/core/knotmirror/config" 15 "tangled.org/core/knotmirror/db" 16 "tangled.org/core/knotmirror/knotstream" 17 "tangled.org/core/knotmirror/models" 18 "tangled.org/core/knotmirror/xrpc" 19 "tangled.org/core/log" 20) 21 22func Run(ctx context.Context, cfg *config.Config) error { 23 // make sure every services are cleaned up on fast return 24 ctx, cancel := context.WithCancel(ctx) 25 defer cancel() 26 27 logger := log.FromContext(ctx) 28 29 db, err := db.Make(ctx, cfg.DbUrl, 32) 30 if err != nil { 31 return fmt.Errorf("initializing db: %w", err) 32 } 33 34 rdb := redis.NewClient(&redis.Options{Addr: cfg.RedisAddr}) 35 36 resolver, err := idresolver.RedisResolver("redis://"+cfg.RedisAddr, cfg.PlcUrl) 37 if err != nil { 38 logger.Error("failed to create redis resolver for admin, falling back to default", "err", err) 39 resolver = idresolver.DefaultResolver(cfg.PlcUrl) 40 } 41 42 // NOTE: using plain git-cli for clone/fetch as go-git is too memory-intensive. 43 gitm := NewCliGitMirrorManager(cfg.GitRepoBasePath, cfg.KnotUseSSL) 44 45 res, err := db.ExecContext(ctx, 46 `update repos set state = $1 where state = $2`, 47 models.RepoStateDesynchronized, 48 models.RepoStateResyncing, 49 ) 50 if err != nil { 51 return fmt.Errorf("clearing resyning states: %w", err) 52 } 53 rows, err := res.RowsAffected() 54 if err != nil { 55 return fmt.Errorf("getting affected rows: %w", err) 56 } 57 logger.Info(fmt.Sprintf("clearing resyning states: %d records updated", rows)) 58 59 knotstream := knotstream.NewKnotStream(logger, db, cfg) 60 crawler := NewCrawler(logger, db) 61 resyncer := NewResyncer(logger, db, gitm, cfg) 62 xrpc := xrpc.New(logger, cfg, db, rdb, resolver, knotstream) 63 adminpage := NewAdminServer(logger, db, resyncer, xrpc, resolver) 64 65 // maintain repository list with tap 66 // NOTE: this can be removed once we introduce did-for-repo because then we can just listen to KnotStream for #identity events. 67 tap := NewTapClient(logger, cfg, db, gitm, knotstream) 68 69 // start http server 70 go func() { 71 logger.Info("starting http server", "addr", cfg.Listen) 72 73 mux := chi.NewRouter() 74 mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { 75 w.Write([]byte("Welcome to a knotmirror server.\n")) 76 }) 77 mux.Mount("/xrpc", xrpc.Router()) 78 79 if err := http.ListenAndServe(cfg.Listen, mux); err != nil { 80 logger.Error("xrpc server failed", "error", err) 81 } 82 }() 83 84 // start metrics endpoint 85 go func() { 86 metricsAddr := cfg.MetricsListen 87 logger.Info("starting metrics server", "addr", metricsAddr) 88 http.Handle("/metrics", promhttp.Handler()) 89 if err := http.ListenAndServe(metricsAddr, nil); err != nil { 90 logger.Error("metrics server failed", "error", err) 91 } 92 }() 93 94 // start admin page endpoint 95 go func() { 96 logger.Info("starting admin server", "addr", cfg.AdminListen) 97 if err := http.ListenAndServe(cfg.AdminListen, adminpage.Router()); err != nil { 98 logger.Error("admin server failed", "error", err) 99 } 100 }() 101 102 tap.Start(ctx) 103 104 resyncer.Start(ctx) 105 106 // periodically crawl the entire network to mirror the repositories 107 crawler.Start(ctx) 108 109 // listen to knotstream (currently we don't have relay for knots, so subscribe every known knots) 110 knotstream.Start(ctx) 111 112 svcErr := make(chan error, 1) 113 if err := knotstream.ResubscribeAllHosts(ctx); err != nil { 114 svcErr <- fmt.Errorf("resubscribing known hosts: %w", err) 115 } 116 117 logger.Info("startup complete") 118 select { 119 case <-ctx.Done(): 120 logger.Info("received shutdown signal", "reason", ctx.Err()) 121 case err := <-svcErr: 122 if err != nil { 123 logger.Error("service error", "error", err) 124 } 125 cancel() 126 } 127 128 logger.Info("shutting down knotmirror") 129 shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second) 130 defer shutdownCancel() 131 132 var errs []error 133 if err := knotstream.Shutdown(shutdownCtx); err != nil { 134 errs = append(errs, err) 135 } 136 if err := db.Close(); err != nil { 137 errs = append(errs, err) 138 } 139 for _, err := range errs { 140 logger.Error("error during shutdown", "err", err) 141 } 142 143 logger.Info("shutdown complete") 144 return nil 145}