diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go --- a/knotmirror/knotmirror.go +++ b/knotmirror/knotmirror.go @@ -49,11 +49,11 @@ } logger.Info(fmt.Sprintf("clearing resyning states: %d records updated", rows)) - xrpc := xrpc.New(logger, cfg, db, resolver) knotstream := knotstream.NewKnotStream(logger, db, cfg) crawler := NewCrawler(logger, db) resyncer := NewResyncer(logger, db, gitm, cfg) adminpage := NewAdminServer(logger, db, resyncer) + xrpc := xrpc.New(logger, cfg, db, resolver, knotstream) // maintain repository list with tap // NOTE: this can be removed once we introduce did-for-repo because then we can just listen to KnotStream for #identity events. diff --git a/knotmirror/resyncer.go b/knotmirror/resyncer.go --- a/knotmirror/resyncer.go +++ b/knotmirror/resyncer.go @@ -281,6 +281,8 @@ repoUrl += "/info/refs?service=git-upload-pack" + r.logger.Debug("checking knot reachability", "url", repoUrl) + client := http.Client{ Timeout: 30 * time.Second, } diff --git a/knotserver/events.go b/knotserver/events.go --- a/knotserver/events.go +++ b/knotserver/events.go @@ -7,7 +7,9 @@ "strconv" "time" + "github.com/bluesky-social/indigo/xrpc" "github.com/gorilla/websocket" + "tangled.org/core/api/tangled" "tangled.org/core/log" ) @@ -60,6 +62,17 @@ l.Error("failed to backfill", "err", err) return } + + // try request crawl when connection closed + defer func() { + go func() { + retryCtx, retryCancel := context.WithTimeout(context.Background(), 10*time.Second) + defer retryCancel() + if err := h.requestCrawl(retryCtx); err != nil { + l.Error("error requesting crawls", "err", err) + } + }() + }() for { // wait for new data or timeout @@ -116,5 +129,21 @@ *cursor = event.Created } + return nil +} + +func (h *Knot) requestCrawl(ctx context.Context) error { + h.l.Info("requesting crawl", "mirrors", h.c.KnotMirrors) + input := &tangled.SyncRequestCrawl_Input{ + Hostname: h.c.Server.Hostname, + } + for _, knotmirror := range h.c.KnotMirrors { + xrpcc := xrpc.Client{Host: knotmirror} + if err := tangled.SyncRequestCrawl(ctx, &xrpcc, input); err != nil { + h.l.Error("error requesting crawl", "err", err) + } else { + h.l.Info("crawl requested successfully") + } + } return nil } diff --git a/knotserver/server.go b/knotserver/server.go --- a/knotserver/server.go +++ b/knotserver/server.go @@ -5,6 +5,7 @@ "fmt" "net/http" + "github.com/bluesky-social/indigo/xrpc" "github.com/urfave/cli/v3" "tangled.org/core/api/tangled" "tangled.org/core/hook" @@ -97,6 +98,21 @@ logger.Info("starting internal server", "address", c.Server.InternalListenAddr) go http.ListenAndServe(c.Server.InternalListenAddr, imux) + + // TODO(boltless): too lazy here. should clear this up + go func() { + input := &tangled.SyncRequestCrawl_Input{ + Hostname: c.Server.Hostname, + } + for _, knotmirror := range c.KnotMirrors { + xrpcc := xrpc.Client{Host: knotmirror} + if err := tangled.SyncRequestCrawl(ctx, &xrpcc, input); err != nil { + logger.Error("error requesting crawl", "err", err) + } else { + logger.Info("crawl requested successfully") + } + } + }() logger.Info("starting main server", "address", c.Server.ListenAddr) logger.Error("server error", "error", http.ListenAndServe(c.Server.ListenAddr, mux)) diff --git a/nix/vm.nix b/nix/vm.nix --- a/nix/vm.nix +++ b/nix/vm.nix @@ -110,7 +110,11 @@ plcUrl = plcUrl; jetstreamEndpoint = jetstream; listenAddr = "0.0.0.0:6444"; + dev = true; }; + knotmirrors = [ + "http://localhost:7000" + ]; }; services.tangled.spindle = { enable = true; diff --git a/knotmirror/hostutil/hostutil.go b/knotmirror/hostutil/hostutil.go new file mode 100644 --- /dev/null +++ b/knotmirror/hostutil/hostutil.go @@ -0,0 +1,56 @@ +package hostutil + +import ( + "fmt" + "net/url" + "strings" + + "github.com/bluesky-social/indigo/atproto/syntax" +) + +func ParseHostname(raw string) (hostname string, noSSL bool, err error) { + // handle case of bare hostname + if !strings.Contains(raw, "://") { + if strings.HasPrefix(raw, "localhost:") { + raw = "http://" + raw + } else { + raw = "https://" + raw + } + } + + u, err := url.Parse(raw) + if err != nil { + return "", false, fmt.Errorf("not a valid host URL: %w", err) + } + + switch u.Scheme { + case "https", "wss": + noSSL = false + case "http", "ws": + noSSL = true + default: + return "", false, fmt.Errorf("unsupported URL scheme: %s", u.Scheme) + } + + // 'localhost' (exact string) is allowed *with* a required port number; SSL is optional + if u.Hostname() == "localhost" { + if u.Port() == "" || !strings.HasPrefix(u.Host, "localhost:") { + return "", false, fmt.Errorf("port number is required for localhost") + } + return u.Host, noSSL, nil + } + + // port numbers not allowed otherwise + if u.Port() != "" { + return "", false, fmt.Errorf("port number not allowed for non-local names") + } + + // check it is a real hostname (eg, not IP address or single-word alias) + h, err := syntax.ParseHandle(u.Host) + if err != nil { + return "", false, fmt.Errorf("not a public hostname") + } + + // lower-case in response + return h.Normalize().String(), noSSL, nil +} diff --git a/knotmirror/xrpc/sync_requestCrawl.go b/knotmirror/xrpc/sync_requestCrawl.go new file mode 100644 --- /dev/null +++ b/knotmirror/xrpc/sync_requestCrawl.go @@ -0,0 +1,104 @@ +package xrpc + +import ( + "encoding/json" + "fmt" + "net/http" + "strings" + + "github.com/bluesky-social/indigo/api/atproto" + "github.com/bluesky-social/indigo/atproto/atclient" + "github.com/bluesky-social/indigo/atproto/syntax" + "github.com/bluesky-social/indigo/xrpc" + "tangled.org/core/api/tangled" + "tangled.org/core/knotmirror/db" + "tangled.org/core/knotmirror/hostutil" + "tangled.org/core/knotmirror/models" +) + +func (x *Xrpc) RequestCrawl(w http.ResponseWriter, r *http.Request) { + var input tangled.SyncRequestCrawl_Input + if err := json.NewDecoder(r.Body).Decode(&input); err != nil { + writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: "failed to decode json body"}) + return + } + + ctx := r.Context() + + l := x.logger.With("input", input) + + hostname, noSSL, err := hostutil.ParseHostname(input.Hostname) + if err != nil { + l.Error("invalid hostname", "err", err) + writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: fmt.Sprintf("hostname field empty or invalid: %s", input.Hostname)}) + return + } + + // TODO: check if host is Knot with knot.describeServer + + // store given repoAt to db + // this will allow knotmirror to ingest repo creation event bypassing tap. + // this step won't be needed once we introduce did-for-repo + // TODO(boltless): remove this section + if input.EnsureRepo != nil { + repoAt, err := syntax.ParseATURI(*input.EnsureRepo) + if err != nil { + l.Error("invalid repo at-uri", "err", err) + writeJson(w, http.StatusBadRequest, atclient.ErrorBody{Name: "BadRequest", Message: fmt.Sprintf("repo parameter invalid: %s", *input.EnsureRepo)}) + return + } + owner, err := x.resolver.ResolveIdent(ctx, repoAt.Authority().String()) + if err != nil || owner.Handle.IsInvalidHandle() { + l.Error("failed to resolve ident", "err", err, "owner", repoAt.Authority().String()) + writeErr(w, fmt.Errorf("failed to resolve repo owner")) + return + } + xrpcc := xrpc.Client{Host: owner.PDSEndpoint()} + out, err := atproto.RepoGetRecord(ctx, &xrpcc, "", tangled.RepoNSID, repoAt.Authority().String(), repoAt.RecordKey().String()) + if err != nil { + l.Error("failed to get repo record", "err", err, "repo", repoAt) + writeErr(w, fmt.Errorf("failed to get repo record")) + return + } + record := out.Value.Val.(*tangled.Repo) + + knotUrl := record.Knot + if !strings.Contains(record.Knot, "://") { + if noSSL { + knotUrl = "http://" + knotUrl + } else { + knotUrl = "https://" + knotUrl + } + } + + repo := &models.Repo{ + Did: owner.DID, + Rkey: repoAt.RecordKey(), + Cid: (*syntax.CID)(out.Cid), + Name: record.Name, + KnotDomain: knotUrl, + State: models.RepoStatePending, + ErrorMsg: "", + RetryAfter: 0, + RetryCount: 0, + } + + if err := db.UpsertRepo(ctx, x.db, repo); err != nil { + l.Error("failed to upsert repo", "err", err) + writeErr(w, err) + return + } + } + + // subscribe to requested host + if !x.ks.CheckIfSubscribed(hostname) { + if err := x.ks.SubscribeHost(ctx, hostname, noSSL); err != nil { + // TODO(boltless): return HostBanned on banned hosts + l.Error("failed to subscribe host", "err", err) + writeErr(w, err) + return + } + } + + w.WriteHeader(http.StatusOK) +} diff --git a/knotmirror/xrpc/xrpc.go b/knotmirror/xrpc/xrpc.go --- a/knotmirror/xrpc/xrpc.go +++ b/knotmirror/xrpc/xrpc.go @@ -12,6 +12,7 @@ "tangled.org/core/api/tangled" "tangled.org/core/idresolver" "tangled.org/core/knotmirror/config" + "tangled.org/core/knotmirror/knotstream" "tangled.org/core/log" ) @@ -19,14 +20,16 @@ cfg *config.Config db *sql.DB resolver *idresolver.Resolver + ks *knotstream.KnotStream logger *slog.Logger } -func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, resolver *idresolver.Resolver) *Xrpc { +func New(logger *slog.Logger, cfg *config.Config, db *sql.DB, resolver *idresolver.Resolver, ks *knotstream.KnotStream) *Xrpc { return &Xrpc{ cfg, db, resolver, + ks, log.SubLogger(logger, "xrpc"), } } @@ -47,6 +50,7 @@ r.Get("/"+tangled.GitTempListCommitsNSID, x.ListCommits) r.Get("/"+tangled.GitTempListLanguagesNSID, x.ListLanguages) r.Get("/"+tangled.GitTempListTagsNSID, x.ListTags) + r.Post("/"+tangled.SyncRequestCrawlNSID, x.RequestCrawl) return r } diff --git a/knotserver/config/config.go b/knotserver/config/config.go --- a/knotserver/config/config.go +++ b/knotserver/config/config.go @@ -39,10 +39,11 @@ } type Config struct { - Repo Repo `env:",prefix=KNOT_REPO_"` - Server Server `env:",prefix=KNOT_SERVER_"` - Git Git `env:",prefix=KNOT_GIT_"` - AppViewEndpoint string `env:"APPVIEW_ENDPOINT, default=https://tangled.org"` + Repo Repo `env:",prefix=KNOT_REPO_"` + Server Server `env:",prefix=KNOT_SERVER_"` + Git Git `env:",prefix=KNOT_GIT_"` + AppViewEndpoint string `env:"APPVIEW_ENDPOINT, default=https://tangled.org"` + KnotMirrors []string `env:"KNOT_MIRRORS, default=https://mirror.tangled.network"` } func Load(ctx context.Context) (*Config, error) { diff --git a/knotserver/xrpc/create_repo.go b/knotserver/xrpc/create_repo.go --- a/knotserver/xrpc/create_repo.go +++ b/knotserver/xrpc/create_repo.go @@ -1,12 +1,14 @@ package xrpc import ( + "context" "encoding/json" "errors" "fmt" "net/http" "path/filepath" "strings" + "time" comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" @@ -120,7 +122,35 @@ repoPath, ) + // HACK: request crawl for this repository + // Users won't want to sync entire network from their local knotmirror. + // Therefore, to bypass the local tap, requestCrawl directly to the knotmirror. + go func() { + if h.Config.Server.Dev { + repoAt := fmt.Sprintf("at://%s/%s/%s", actorDid, tangled.RepoNSID, rkey) + rCtx, rCancel := context.WithTimeout(context.Background(), 10*time.Second) + defer rCancel() + h.requestCrawl(rCtx, &tangled.SyncRequestCrawl_Input{ + Hostname: h.Config.Server.Hostname, + EnsureRepo: &repoAt, + }) + } + }() + w.WriteHeader(http.StatusOK) +} + +func (h *Xrpc) requestCrawl(ctx context.Context, input *tangled.SyncRequestCrawl_Input) error { + h.Logger.Info("requesting crawl", "mirrors", h.Config.KnotMirrors) + for _, knotmirror := range h.Config.KnotMirrors { + xrpcc := xrpc.Client{Host: knotmirror} + if err := tangled.SyncRequestCrawl(ctx, &xrpcc, input); err != nil { + h.Logger.Error("error requesting crawl", "err", err) + } else { + h.Logger.Info("crawl requested successfully") + } + } + return nil } func validateRepoName(name string) error { diff --git a/nix/modules/knot.nix b/nix/modules/knot.nix --- a/nix/modules/knot.nix +++ b/nix/modules/knot.nix @@ -115,6 +115,14 @@ ''; }; + knotmirrors = mkOption { + type = types.listOf types.str; + default = [ + "https://mirror.tangled.network" + ]; + description = "List of knotmirror hosts to request crawl"; + }; + server = { listenAddr = mkOption { type = types.str; @@ -263,6 +271,7 @@ "KNOT_SERVER_PLC_URL=${cfg.server.plcUrl}" "KNOT_SERVER_JETSTREAM_ENDPOINT=${cfg.server.jetstreamEndpoint}" "KNOT_SERVER_OWNER=${cfg.server.owner}" + "KNOT_MIRRORS=${concatStringsSep "," cfg.knotmirrors}" "KNOT_SERVER_LOG_DIDS=${ if cfg.server.logDids then "true"