From cfbcc31012529bcaec101039e1b70668538002b0 Mon Sep 17 00:00:00 2001 From: hailey Date: Sat, 25 Oct 2025 12:13:51 -0700 Subject: [PATCH] request crawl when websocket dies (#31) --- server/handle_sync_subscribe_repos.go | 42 ++++++++++++++++++++------- server/server.go | 41 ++++++++++++++++++++++---- 2 files changed, 67 insertions(+), 16 deletions(-) diff --git a/server/handle_sync_subscribe_repos.go b/server/handle_sync_subscribe_repos.go index 6423ce6..d955aa4 100644 --- a/server/handle_sync_subscribe_repos.go +++ b/server/handle_sync_subscribe_repos.go @@ -1,7 +1,8 @@ package server import ( - "fmt" + "context" + "time" "github.com/bluesky-social/indigo/events" "github.com/bluesky-social/indigo/lex/util" @@ -10,16 +11,18 @@ import ( ) func (s *Server) handleSyncSubscribeRepos(e echo.Context) error { + ctx := e.Request().Context() + logger := s.logger.With("component", "subscribe-repos-websocket") + conn, err := websocket.Upgrade(e.Response().Writer, e.Request(), e.Response().Header(), 1<<10, 1<<10) if err != nil { + logger.Error("unable to establish websocket with relay", "err", err) return err } - s.logger.Info("new connection", "ua", e.Request().UserAgent()) - - ctx := e.Request().Context() - ident := e.RealIP() + "-" + e.Request().UserAgent() + logger = logger.With("ident", ident) + logger.Info("new connection established") evts, cancel, err := s.evtman.Subscribe(ctx, ident, func(evt *events.XRPCStreamEvent) bool { return true @@ -33,11 +36,16 @@ func (s *Server) handleSyncSubscribeRepos(e echo.Context) error { for evt := range evts { wc, err := conn.NextWriter(websocket.BinaryMessage) if err != nil { - return err + logger.Error("error writing message to relay", "err", err) + break } - var obj util.CBOR + if ctx.Err() != nil { + logger.Error("context error", "err", err) + break + } + var obj util.CBOR switch { case evt.Error != nil: header.Op = events.EvtKindErrorFrame @@ -55,21 +63,33 @@ func (s *Server) handleSyncSubscribeRepos(e echo.Context) error { header.MsgType = "#info" obj = evt.RepoInfo default: - return fmt.Errorf("unrecognized event kind") + logger.Warn("unrecognized event kind") + return nil } if err := header.MarshalCBOR(wc); err != nil { - return fmt.Errorf("failed to write header: %w", err) + logger.Error("failed to write header to relay", "err", err) + break } if err := obj.MarshalCBOR(wc); err != nil { - return fmt.Errorf("failed to write event: %w", err) + logger.Error("failed to write event to relay", "err", err) + break } if err := wc.Close(); err != nil { - return fmt.Errorf("failed to flush-close our event write: %w", err) + logger.Error("failed to flush-close our event write", "err", err) + break } } + // we should tell the relay to request a new crawl at this point if we got disconnected + // use a new context since the old one might be cancelled at this point + ctx, cancel = context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if err := s.requestCrawl(ctx); err != nil { + logger.Error("error requesting crawls", "err", err) + } + return nil } diff --git a/server/server.go b/server/server.go index 07f3ae7..85d1f99 100644 --- a/server/server.go +++ b/server/server.go @@ -77,6 +77,9 @@ type Server struct { passport *identity.Passport fallbackProxy string + lastRequestCrawl time.Time + requestCrawlMu sync.Mutex + dbName string s3Config *S3Config } @@ -518,16 +521,44 @@ func (s *Server) Serve(ctx context.Context) error { go s.backupRoutine() + go func() { + if err := s.requestCrawl(ctx); err != nil { + s.logger.Error("error requesting crawls", "err", err) + } + }() + + <-ctx.Done() + + fmt.Println("shut down") + + return nil +} + +func (s *Server) requestCrawl(ctx context.Context) error { + logger := s.logger.With("component", "request-crawl") + s.requestCrawlMu.Lock() + defer s.requestCrawlMu.Unlock() + + logger.Info("requesting crawl with configured relays") + + if time.Now().Sub(s.lastRequestCrawl) <= 1*time.Minute { + return fmt.Errorf("a crawl request has already been made within the last minute") + } + for _, relay := range s.config.Relays { + logger := logger.With("relay", relay) + logger.Info("requesting crawl from relay") cli := xrpc.Client{Host: relay} - atproto.SyncRequestCrawl(ctx, &cli, &atproto.SyncRequestCrawl_Input{ + if err := atproto.SyncRequestCrawl(ctx, &cli, &atproto.SyncRequestCrawl_Input{ Hostname: s.config.Hostname, - }) + }); err != nil { + logger.Error("error requesting crawl", "err", err) + } else { + logger.Info("crawl requested successfully") + } } - <-ctx.Done() - - fmt.Println("shut down") + s.lastRequestCrawl = time.Now() return nil } -- 2.51.2