Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135package migrator
import ( "context" "fmt" "net/http" "time"
"tangled.org/core/idresolver" "tangled.org/core/log" "tangled.org/core/migrator/api" "tangled.org/core/migrator/config" "tangled.org/core/migrator/db" "tangled.org/core/migrator/git" "tangled.org/core/migrator/gitserver" migratoroauth "tangled.org/core/migrator/oauth" "tangled.org/core/netutil" "tangled.org/core/migrator/worker" "tangled.org/core/xrpc/serviceauth")
func Run(ctx context.Context, cfg *config.Config) error { ctx, cancel := context.WithCancel(ctx) defer cancel()
logger := log.FromContext(ctx)
database, err := db.Make(ctx, cfg.DbPath) if err != nil { return fmt.Errorf("initializing database: %w", err) } defer database.Close() database.MaxActivePerOwner = cfg.MaxActivePerOwner
connectProxy := netutil.NewSafeConnectProxy(nil, nil) proxyAddr, err := connectProxy.Start() if err != nil { return fmt.Errorf("starting connect proxy: %w", err) } defer connectProxy.Close() logger.Info("started safe connect proxy", "addr", proxyAddr)
signer, err := serviceauth.NewSigner(cfg.ServiceDid, cfg.PrivateKey) if err != nil { return fmt.Errorf("initializing signer: %w", err) }
resolver := idresolver.DefaultResolver(cfg.PlcUrl) serviceAuth := serviceauth.NewServiceAuth(logger, resolver.Directory(), cfg.ServiceDid.String())
oauthClient, err := migratoroauth.NewClient(cfg, database, resolver.Directory(), logger) if err != nil { return fmt.Errorf("initializing oauth client: %w", err) }
workerPool, err := worker.NewWorkerPool( database, cfg, connectProxy, git.RealCommandRunner{ MaxDiskBytes: cfg.MaxDiskBytes, }, oauthClient, logger, ) if err != nil { return fmt.Errorf("initializing worker pool: %w", err) }
workerPool.Start(ctx)
go func() { ticker := time.NewTicker(15 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: scrubbed, err := database.ScrubExpiredCredentials(ctx) if err != nil { logger.Error("error scrubbing expired credentials", "err", err) } else if scrubbed > 0 { logger.Info("scrubbed expired GitHub credentials past 1h window", "count", scrubbed) } } } }()
apiServer := api.NewServer(database, cfg, signer, serviceAuth, oauthClient, logger, nil, proxyAddr)
// short timeout on json routes so a hung knot or pds cannot pin connections short := func(h http.Handler) http.Handler { return http.TimeoutHandler(h, cfg.APITimeout, "api timeout") }
mux := http.NewServeMux() mux.Handle("/git/", gitserver.New(database, cfg)) mux.Handle("/oauth-client-metadata.json", short(oauthClient.Routes())) mux.Handle("/oauth/", short(oauthClient.Routes())) mux.Handle("/", short(apiServer.Routes()))
srv := &http.Server{ Addr: cfg.ListenAddr, Handler: mux, ReadHeaderTimeout: 5 * time.Second, ReadTimeout: cfg.JobTimeout, // streaming budget for /git/ WriteTimeout: cfg.JobTimeout, IdleTimeout: 60 * time.Second, }
go func() { logger.Info("starting migrator http server", "addr", cfg.ListenAddr, "serviceDid", cfg.ServiceDid) if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { logger.Error("http server stopped with error", "err", err) cancel() } }()
<-ctx.Done() logger.Info("shutting down migrator daemon")
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 15*time.Second) defer shutdownCancel()
if err := srv.Shutdown(shutdownCtx); err != nil { logger.Error("error shutting down http server", "err", err) }
workerPool.Wait() logger.Info("migrator shutdown complete")
return nil}