Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 57 additions & 6 deletions cmd/mxcli/docker/hubclient.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,11 @@ type HubRegistration struct {
MultiTenant bool // false when we fell back to a slice-1 single-app hub

hubURL string
// Inputs kept so the heartbeat can re-register if the hub forgets us (e.g. the
// hub restarted and lost its in-memory registry — /api/status then 404s).
secret string
meta HubMeta
appPort int
}

type registerResponse struct {
Expand Down Expand Up @@ -94,6 +99,9 @@ func RegisterWithHub(hubURL, secret string, meta HubMeta, appPort int) (*HubRegi
HeartbeatInterval: time.Duration(rr.HeartbeatIntervalSec) * time.Second,
MultiTenant: true,
hubURL: hubURL,
secret: secret,
meta: meta,
appPort: appPort,
}, nil
case http.StatusNotFound:
// No registration API — a single-app hub. Serve directly at the hub URL.
Expand All @@ -120,7 +128,16 @@ type Heartbeat struct {
}

// StartHeartbeat begins heartbeating (only for a multi-tenant registration).
func StartHeartbeat(reg *HubRegistration) *Heartbeat {
//
// If the hub forgets us — it restarts and loses its in-memory registry, so
// /api/status returns 404 "unknown token" — the heartbeat re-registers in place
// so the preview reappears in the hub and the public URL works again (previously
// the tunnel stayed connected but the hub listed nothing, silently dead). When a
// re-register lands on a different reverse port, onReRegister (if set) is invoked
// with the refreshed registration so the caller can restart the tunnel; the same
// identity usually re-registers to the same port, in which case the existing
// tunnel keeps working and onReRegister is not called.
func StartHeartbeat(reg *HubRegistration, onReRegister func(*HubRegistration)) *Heartbeat {
if !reg.MultiTenant || reg.HeartbeatInterval <= 0 {
return &Heartbeat{reg: reg}
}
Expand All @@ -135,13 +152,42 @@ func StartHeartbeat(reg *HubRegistration) *Heartbeat {
case <-ctx.Done():
return
case <-t.C:
postToken(ctx, reg.hubURL+"/api/status", reg.Token)
code := postToken(ctx, reg.hubURL+"/api/status", reg.Token)
if code != http.StatusNotFound && code != http.StatusUnauthorized {
continue
}
// The hub no longer knows this token — re-register in place.
portChanged, err := reg.reRegister()
if err != nil {
continue // transient; try again next tick
}
if portChanged && onReRegister != nil {
onReRegister(reg)
}
}
}
}()
return h
}

// reRegister re-POSTs /api/register with the original identity after the hub has
// forgotten this preview, updating the registration in place. Returns whether the
// assigned reverse port changed (the caller must then restart the tunnel).
func (reg *HubRegistration) reRegister() (portChanged bool, err error) {
fresh, err := RegisterWithHub(reg.hubURL, reg.secret, reg.meta, reg.appPort)
if err != nil {
return false, err
}
portChanged = fresh.ReversePort != reg.ReversePort
reg.URL = fresh.URL
reg.Subdomain = fresh.Subdomain
reg.ControlURL = fresh.ControlURL
reg.ReversePort = fresh.ReversePort
reg.Token = fresh.Token
reg.TunnelAuth = fresh.TunnelAuth
return portChanged, nil
}

// Stop ends heartbeating and best-effort deregisters the preview from the hub.
func (h *Heartbeat) Stop() {
if h.cancel != nil {
Expand All @@ -155,15 +201,20 @@ func (h *Heartbeat) Stop() {
}
}

func postToken(ctx context.Context, url, token string) {
// postToken POSTs to a hub endpoint with the bearer token and returns the HTTP
// status code (0 on a transport error, so callers treat it as "try again").
func postToken(ctx context.Context, url, token string) int {
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, nil)
if err != nil {
return
return 0
}
req.Header.Set("Authorization", "Bearer "+token)
if resp, err := http.DefaultClient.Do(req); err == nil {
_ = resp.Body.Close()
resp, err := http.DefaultClient.Do(req)
if err != nil {
return 0
}
defer resp.Body.Close()
return resp.StatusCode
}

// DetectHubMeta fills a HubMeta from the project path + git, with explicit
Expand Down
70 changes: 70 additions & 0 deletions cmd/mxcli/docker/hubclient_reregister_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
// SPDX-License-Identifier: Apache-2.0

package docker

import (
"encoding/json"
"net/http"
"net/http/httptest"
"sync/atomic"
"testing"
"time"
)

// TestHeartbeatReRegistersAfterHubForgets simulates a hub that restarts and loses
// its registry: /api/status then 404s. The heartbeat must re-register in place
// (findings #16 — previously it pinged once at startup and never recovered).
func TestHeartbeatReRegistersAfterHubForgets(t *testing.T) {
var registerCalls int32
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api/register":
n := atomic.AddInt32(&registerCalls, 1)
// First registration: port 40001. Re-registration: a different port
// (40002) so the onReRegister callback fires.
port := 40001
token := "tok-1"
if n > 1 {
port = 40002
token = "tok-2"
}
_ = json.NewEncoder(w).Encode(map[string]any{
"subdomain": "app", "url": srvURL(r), "reversePort": port,
"controlUrl": srvURL(r), "token": token, "tunnelAuth": "auth",
"heartbeatIntervalSec": 1,
})
case "/api/status":
// The hub has forgotten us.
w.WriteHeader(http.StatusNotFound)
default:
w.WriteHeader(http.StatusNotFound)
}
}))
defer srv.Close()

reg, err := RegisterWithHub(srv.URL, "", HubMeta{Project: "p"}, 8080)
if err != nil {
t.Fatalf("initial register: %v", err)
}
if reg.ReversePort != 40001 || reg.Token != "tok-1" {
t.Fatalf("unexpected initial reg: port=%d token=%s", reg.ReversePort, reg.Token)
}

changed := make(chan *HubRegistration, 1)
hb := StartHeartbeat(reg, func(r *HubRegistration) { changed <- r })
defer hb.Stop()

select {
case r := <-changed:
if r.ReversePort != 40002 || r.Token != "tok-2" {
t.Errorf("re-register did not update reg: port=%d token=%s", r.ReversePort, r.Token)
}
if atomic.LoadInt32(&registerCalls) < 2 {
t.Errorf("expected a second /api/register call, got %d", registerCalls)
}
case <-time.After(5 * time.Second):
t.Fatal("heartbeat did not re-register within 5s after the hub 404'd")
}
}

func srvURL(r *http.Request) string { return "http://" + r.Host }
29 changes: 26 additions & 3 deletions cmd/mxcli/docker/runlocal.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (
"os/signal"
"path/filepath"
"strings"
"sync"
"syscall"
"time"

Expand Down Expand Up @@ -492,9 +493,31 @@ func RunLocal(opts LocalRunOptions) error {
if err != nil {
return fmt.Errorf("starting hub tunnel: %w", err)
}
defer tunnel.Stop()

hb := StartHeartbeat(hubReg)
// The heartbeat re-registers if the hub restarts; when that lands on a new
// reverse port, restart the tunnel to the new port. Guard the handle since
// the callback runs on the heartbeat goroutine.
var tunMu sync.Mutex
defer func() { tunMu.Lock(); tunnel.Stop(); tunMu.Unlock() }()

hb := StartHeartbeat(hubReg, func(reg *HubRegistration) {
tunMu.Lock()
defer tunMu.Unlock()
tunnel.Stop()
nt, err := StartTunnel(TunnelOptions{
HubURL: reg.ControlURL,
LocalPort: opts.AppPort,
RemotePort: reg.ReversePort,
Secret: reg.TunnelAuth,
PublicURL: reg.URL,
Stdout: w,
})
if err != nil {
fmt.Fprintf(stderr, "hub re-registered but restarting the tunnel failed: %v\n", err)
return
}
tunnel = nt
fmt.Fprintf(w, "Re-registered with hub after restart; preview available at %s\n", reg.URL)
})
defer hb.Stop()

fmt.Fprintf(w, "Preview available at %s\n", hubReg.URL)
Expand Down
Loading