From 17bb59e37536aa476512ee2a4935da3d964a590d Mon Sep 17 00:00:00 2001 From: Akshay Date: Wed, 12 Feb 2025 17:20:49 +0000 Subject: [PATCH] fix jetstream not reconnecting --- appview/state/state.go | 1 + knotserver/jetstream.go | 9 ++++++++- knotserver/routes.go | 4 ++-- 3 files changed, 11 insertions(+), 3 deletions(-) diff --git a/appview/state/state.go b/appview/state/state.go index 65a126d8..6acf47ef 100644 --- a/appview/state/state.go +++ b/appview/state/state.go @@ -467,6 +467,7 @@ func (s *State) AddMember(w http.ResponseWriter, r *http.Request) { ksResp, err := ksClient.AddMember(memberIdent.DID.String()) if err != nil { log.Printf("failed to make request to %s: %s", domain, err) + return } if ksResp.StatusCode != http.StatusNoContent { diff --git a/knotserver/jetstream.go b/knotserver/jetstream.go index a0bc3d88..4ca30632 100644 --- a/knotserver/jetstream.go +++ b/knotserver/jetstream.go @@ -47,7 +47,7 @@ func (h *Handle) StartJetstream(ctx context.Context) error { jc := &JetstreamClient{ cfg: cfg, client: client, - reconnectCh: make(chan struct{}), + reconnectCh: make(chan struct{}, 1), } h.jc = jc @@ -80,6 +80,13 @@ func (h *Handle) connectAndRead(ctx context.Context, cursor *int64) { } } +func (j *JetstreamClient) AddDid(did string) { + j.mu.Lock() + j.cfg.WantedDids = append(j.cfg.WantedDids, did) + j.mu.Unlock() + j.reconnectCh <- struct{}{} +} + func (j *JetstreamClient) UpdateDids(dids []string) { j.mu.Lock() j.cfg.WantedDids = dids diff --git a/knotserver/routes.go b/knotserver/routes.go index 52a3209b..a3e894b5 100644 --- a/knotserver/routes.go +++ b/knotserver/routes.go @@ -484,7 +484,7 @@ func (h *Handle) AddMember(w http.ResponseWriter, r *http.Request) { return } - h.jc.UpdateDids([]string{did}) + h.jc.AddDid(did) if err := h.e.AddMember(ThisServer, did); err != nil { l.Error("adding member", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) @@ -520,7 +520,7 @@ func (h *Handle) AddRepoCollaborator(w http.ResponseWriter, r *http.Request) { writeError(w, err.Error(), http.StatusInternalServerError) return } - h.jc.UpdateDids([]string{data.Did}) + h.jc.AddDid(data.Did) repoName := filepath.Join(ownerDid, repo) if err := h.e.AddRepo(data.Did, ThisServer, repoName); err != nil { -- 2.51.2