diff --git a/flake.nix b/flake.nix index 88a6ba92..9e66decc 100644 --- a/flake.nix +++ b/flake.nix @@ -107,7 +107,7 @@ knot = self.callPackage ./nix/pkgs/knot.nix {}; dolly = self.callPackage ./nix/pkgs/dolly.nix {}; tap = self.callPackage ./nix/pkgs/tap.nix {}; - knotmirror = self.callPackage ./nix/pkgs/knot-mirror.nix {}; + knotmirror = self.callPackage ./nix/pkgs/knotmirror.nix {}; }); in { overlays.default = final: prev: { @@ -133,6 +133,7 @@ docs dolly tap + knotmirror ; pkgsStatic-appview = staticPackages.appview; diff --git a/knotmirror/knotstream/slurper.go b/knotmirror/knotstream/slurper.go index 5b6c2cae..a19a3d98 100644 --- a/knotmirror/knotstream/slurper.go +++ b/knotmirror/knotstream/slurper.go @@ -263,15 +263,17 @@ func (s *KnotSlurper) ProcessEvent(ctx context.Context, task *Task) error { return fmt.Errorf("unmarshaling message: %w", err) } - if err := s.ProcessLegacyGitRefUpdate(ctx, &legacyMessage); err != nil { + if err := s.ProcessLegacyGitRefUpdate(ctx, task.key, &legacyMessage); err != nil { return fmt.Errorf("processing gitRefUpdate: %w", err) } return nil } -func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, evt *LegacyGitEvent) error { +func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, source string, evt *LegacyGitEvent) error { knotstreamEventsReceived.Inc() + l := s.logger.With("src", source) + curr, err := db.GetRepoByName(ctx, s.db, syntax.DID(evt.Event.RepoDid), evt.Event.RepoName) if err != nil { return fmt.Errorf("failed to get repo '%s': %w", evt.Event.RepoDid+"/"+evt.Event.RepoName, err) @@ -284,11 +286,11 @@ func (s *KnotSlurper) ProcessLegacyGitRefUpdate(ctx context.Context, evt *Legacy // But we want to store that in did/rkey in knot-mirror. // Therefore, we should ignore when the repository is unknown. // Hopefully crawler will sync it later. - s.logger.Warn("skipping event from unknown repo", "did/repo", evt.Event.RepoDid+"/"+evt.Event.RepoName) + l.Warn("skipping event from unknown repo", "did/name", evt.Event.RepoDid+"/"+evt.Event.RepoName) knotstreamEventsSkipped.Inc() return nil } - l := s.logger.With("repoAt", curr.AtUri()) + l = l.With("repoAt", curr.AtUri()) // TODO: should plan resync to resyncBuffer on RepoStateResyncing if curr.State != models.RepoStateActive { diff --git a/knotmirror/models/models.go b/knotmirror/models/models.go index a1114590..8dbc27e9 100644 --- a/knotmirror/models/models.go +++ b/knotmirror/models/models.go @@ -85,6 +85,22 @@ var AllHostStatuses = []HostStatus{ HostStatusBanned, } +func (h *Host) URL() string { + if h.NoSSL { + return fmt.Sprintf("http://%s", h.Hostname) + } else { + return fmt.Sprintf("https://%s", h.Hostname) + } +} + +func (h *Host) WsURL() string { + if h.NoSSL { + return fmt.Sprintf("ws://%s", h.Hostname) + } else { + return fmt.Sprintf("wss://%s", h.Hostname) + } +} + // func (h *Host) SubscribeGitRefsURL(cursor int64) string { // scheme := "wss" // if h.NoSSL { @@ -98,11 +114,7 @@ var AllHostStatuses = []HostStatus{ // } func (h *Host) LegacyEventsURL(cursor int64) string { - scheme := "wss" - if h.NoSSL { - scheme = "ws" - } - u := fmt.Sprintf("%s://%s/events", scheme, h.Hostname) + u := fmt.Sprintf("%s/events", h.WsURL()) if cursor > 0 { u = fmt.Sprintf("%s?cursor=%d", u, cursor) } diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go index 6e39d595..1cf90397 100644 --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -24,6 +24,7 @@ type Resyncer struct { logger *slog.Logger db *sql.DB gitm GitMirrorManager + cfg *config.Config claimJobMu sync.Mutex @@ -43,6 +44,7 @@ func NewResyncer(l *slog.Logger, db *sql.DB, gitm GitMirrorManager, cfg *config. logger: log.SubLogger(l, "resyncer"), db: db, gitm: gitm, + cfg: cfg, runningJobs: make(map[syntax.ATURI]context.CancelFunc), @@ -272,7 +274,7 @@ func isRateLimitError(err error) bool { // checkKnotReachability checks if Knot is reachable and is valid git remote server func (r *Resyncer) checkKnotReachability(ctx context.Context, repo *models.Repo) error { - repoUrl, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), true) + repoUrl, err := makeRepoRemoteUrl(repo.KnotDomain, repo.DidSlashRepo(), r.cfg.KnotUseSSL) if err != nil { return err } diff --git a/knotmirror/tapclient.go b/knotmirror/tapclient.go index 67e5242f..14ba6d20 100644 --- a/knotmirror/tapclient.go +++ b/knotmirror/tapclient.go @@ -8,6 +8,7 @@ import ( "log/slog" "net/netip" "net/url" + "strings" "time" "tangled.org/core/api/tangled" @@ -78,9 +79,23 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error return fmt.Errorf("parsing record: %w", err) } + knotUrl := record.Knot + if !strings.Contains(record.Knot, "://") { + if host, _ := db.GetHost(ctx, t.db, record.Knot); host != nil { + knotUrl = host.URL() + } else { + t.logger.Warn("repo is from unknown knot") + if t.cfg.KnotUseSSL { + knotUrl = "https://" + knotUrl + } else { + knotUrl = "http://" + knotUrl + } + } + } + status := models.RepoStatePending errMsg := "" - u, err := url.Parse("http://" + record.Knot) // parsing with fake scheme + u, err := url.Parse(knotUrl) if err != nil { status = models.RepoStateSuspended errMsg = "failed to parse knot url" @@ -94,7 +109,7 @@ func (t *Tap) processRepo(ctx context.Context, evt *tapc.RecordEventData) error Rkey: evt.Rkey, Cid: evt.CID, Name: record.Name, - KnotDomain: record.Knot, + KnotDomain: knotUrl, State: status, ErrorMsg: errMsg, RetryAfter: 0, // clear retry info diff --git a/nix/modules/knotmirror.nix b/nix/modules/knotmirror.nix index 9798687c..b9f2ed28 100644 --- a/nix/modules/knotmirror.nix +++ b/nix/modules/knotmirror.nix @@ -66,6 +66,18 @@ in description = "Whether to automatically mirror from entire network"; }; + knotUseSSL = mkOption { + type = types.bool; + default = true; + description = "Use SSL for knot connection"; + }; + + knotSSRF = mkOption { + type = types.bool; + default = true; + description = "enable SSRF protection for knots"; + }; + tap = { port = mkOption { type = types.port; @@ -128,8 +140,8 @@ in "MIRROR_TAP_URL=http://localhost:${toString cfg.tap.port}" "MIRROR_DB_URL=${cfg.dbUrl}" "MIRROR_GIT_BASEPATH=/var/lib/knotmirror/repos" - "MIRROR_KNOT_USE_SSL=true" - "MIRROR_KNOT_SSRF=true" + "MIRROR_KNOT_USE_SSL=${boolToString cfg.knotUseSSL}" + "MIRROR_KNOT_SSRF=${boolToString cfg.knotSSRF}" "MIRROR_RESYNC_PARALLELISM=12" "MIRROR_METRICS_LISTEN=127.0.0.1:7100" "MIRROR_ADMIN_LISTEN=${cfg.adminListenAddr}" diff --git a/nix/pkgs/knot-mirror.nix b/nix/pkgs/knotmirror.nix similarity index 100% rename from nix/pkgs/knot-mirror.nix rename to nix/pkgs/knotmirror.nix diff --git a/nix/vm.nix b/nix/vm.nix index 461dffdc..35eda3de 100644 --- a/nix/vm.nix +++ b/nix/vm.nix @@ -25,6 +25,7 @@ in modules = [ self.nixosModules.knot self.nixosModules.spindle + self.nixosModules.knotmirror ({ lib, config, @@ -57,6 +58,24 @@ in host.port = 6555; guest.port = 6555; } + # knotmirror + { + from = "host"; + host.port = 7007; # 7000 is deserved in macos for Airplay + guest.port = 7000; + } + # knotmirror-tap + { + from = "host"; + host.port = 7480; + guest.port = 7480; + } + # knotmirror-admin + { + from = "host"; + host.port = 7200; + guest.port = 7200; + } ]; sharedDirectories = { # We can't use the 9p mounts directly for most of these @@ -81,7 +100,7 @@ in networking.firewall.enable = false; time.timeZone = "Europe/London"; services.getty.autologinUser = "root"; - environment.systemPackages = with pkgs; [curl vim git sqlite litecli]; + environment.systemPackages = with pkgs; [curl vim git sqlite litecli postgresql_14]; services.tangled.knot = { enable = true; motd = "Welcome to the development knot!\n"; @@ -109,6 +128,27 @@ in }; }; }; + services.postgresql = { + enable = true; + package = pkgs.postgresql_14; + ensureDatabases = ["mirror" "tap"]; + ensureUsers = [ + {name = "tnglr";} + ]; + authentication = '' + local all tnglr trust + host all tnglr 127.0.0.1/32 trust + ''; + }; + services.tangled.knotmirror = { + enable = true; + listenAddr = "0.0.0.0:7000"; + adminListenAddr = "0.0.0.0:7200"; + hostname = "localhost:7000"; + dbUrl = "postgresql://tnglr@127.0.0.1:5432/mirror"; + fullNetwork = false; + tap.dbUrl = "postgresql://tnglr@127.0.0.1:5432/tap"; + }; users = { # So we don't have to deal with permission clashing between # blank disk VMs and existing state @@ -135,6 +175,8 @@ in in { knot = mkDataSyncScripts "/mnt/knot-data" config.services.tangled.knot.stateDir; spindle = mkDataSyncScripts "/mnt/spindle-data" (builtins.dirOf config.services.tangled.spindle.server.dbPath); + knotmirror.after = ["postgresql.target"]; + tap-knotmirror.after = ["postgresql.target"]; }; }) ];