This repository has no description
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}