package 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 }