feat: enhance relay discovery and failure handling with new state management
Kim committed
Apr 27, 2026 at 16:08 UTC
35abe74e309a4a17c01093dd868a3533d99915d7
12 files changed
+342
-139
frontend/src/components/ServerListView.tsx
+37
-12
@@ -43,6 +43,8 @@ interface KnownRelay {
43
isCurrent: boolean;
44
}
45
46
+type RelayReleaseVersions = Record<string, string | null>;
47
+
48
const OFFICIAL_REGISTRY_SOURCE_URL =
49
"https://raw.githubusercontent.com/gosuda/portal-tunnel/main/registry.json";
50
const REPOSITORY_URL = "https://github.com/gosuda/portal-tunnel";
@@ -111,6 +113,17 @@ function normalizeKnownRelays(
113
return knownRelays;
114
}
115
116
+function relayReleaseLabel(
117
+ versions: RelayReleaseVersions,
118
+ relayURL: string
119
+): string {
120
+ const version = versions[relayURL];
121
+ if (version === undefined || version === null) {
122
+ return "loading...";
123
+ }
124
+ return version || "offline";
125
+}
126
+
127
interface ServerListViewProps {
128
title?: string;
129
searchQuery: string;
@@ -201,7 +214,7 @@ export function ServerListView({
214
}: ServerListViewProps) {
215
const [showFilterModal, setShowFilterModal] = useState(false);
216
const [relayReleaseVersions, setRelayReleaseVersions] = useState<
204
- Record<string, string>
217
+ RelayReleaseVersions
218
>({});
219
const [knownRelays, setKnownRelays] = useState<KnownRelay[]>([]);
220
const [relayDiscoveryLoading, setRelayDiscoveryLoading] = useState(
@@ -309,23 +322,32 @@ export function ServerListView({
322
discoveryMessage = "Known relay data is unavailable on this relay.";
323
}
324
312
- const relayURLs = nextKnownRelays.map((relay) => relay.relayURL);
313
- const uniqueRelayURLs = [...new Set(relayURLs)];
314
- const versions = await Promise.all(
315
- uniqueRelayURLs.map(async (relayURL) => [
316
- relayURL,
317
- await loadRelayReleaseVersion(relayURL),
318
- ] as const)
319
- );
320
-
325
if (cancelled) {
326
return;
327
}
328
325
- setRelayReleaseVersions(Object.fromEntries(versions));
329
setKnownRelays(nextKnownRelays);
330
setRelayDiscoveryLoading(false);
331
setRelayDiscoveryMessage(discoveryMessage);
332
+
333
+ const relayURLs = nextKnownRelays.map((relay) => relay.relayURL);
334
+ const uniqueRelayURLs = [...new Set(relayURLs)];
335
+ setRelayReleaseVersions(
336
+ Object.fromEntries(uniqueRelayURLs.map((relayURL) => [relayURL, null]))
337
+ );
338
+
339
+ uniqueRelayURLs.forEach((relayURL) => {
340
+ void (async () => {
341
+ const version = await loadRelayReleaseVersion(relayURL);
342
+ if (cancelled) {
343
+ return;
344
+ }
345
+ setRelayReleaseVersions((prev) => ({
346
+ ...prev,
347
+ [relayURL]: version,
348
+ }));
349
+ })();
350
+ });
351
})();
352
353
return () => {
@@ -871,7 +893,10 @@ export function ServerListView({
893
</a>
894
<div className="flex shrink-0 items-center gap-2">
895
<span className="rounded-full bg-background px-2.5 py-1 font-mono text-[11px] font-medium text-text-muted ring-1 ring-border">
874
- {relayReleaseVersions[relay.relayURL] || "offline"}
896
+ {relayReleaseLabel(
897
+ relayReleaseVersions,
898
+ relay.relayURL
899
+ )}
900
</span>
901
</div>
902
</div>
portal/api_server.go
+28
-23
@@ -246,30 +246,35 @@ func (s *Server) handleRelayDiscoveryAnnounce(w http.ResponseWriter, r *http.Req
246
return
247
}
248
249
- desc := req.Descriptor
249
+ desc, err := utils.NormalizeDescriptor(req.Descriptor)
250
+ if err != nil {
251
+ utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest, err.Error())
252
+ return
253
+ }
254
// Self-announce guard: the relay's own URL is established locally, not
251
- // gossiped through the announce endpoint. Reject loopback / own-host
252
- // announces to prevent self-amplification or misconfiguration loops.
253
- announceURL, err := url.Parse(strings.TrimSpace(desc.APIHTTPSAddr))
254
- if err == nil && announceURL != nil {
255
- host := utils.NormalizeHostname(announceURL.Hostname())
256
- if utils.IsLocalRelayHost(host) {
257
- utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest,
258
- fmt.Sprintf("self-announce rejected: host %q is local-only", host))
259
- return
260
- }
261
- if selfURL, err := utils.NormalizeRelayURL(s.cfg.PortalURL); err == nil {
262
- if announceRelayURL, err := utils.NormalizeRelayURL(desc.APIHTTPSAddr); err == nil && announceRelayURL == selfURL {
263
- utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest,
264
- fmt.Sprintf("self-announce rejected: %q matches receiving relay url", announceRelayURL))
265
- return
266
- }
267
- }
268
- if host != "" && host == utils.NormalizeHostname(s.identity.Name) {
269
- utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest,
270
- fmt.Sprintf("self-announce rejected: host %q matches receiving relay host", host))
271
- return
272
- }
255
+ // gossiped through the announce endpoint. Validate the normalized URL so
256
+ // scheme-less inputs are checked the same way signature verification will
257
+ // check them later.
258
+ announceURL, err := url.Parse(desc.APIHTTPSAddr)
259
+ if err != nil {
260
+ utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest, err.Error())
261
+ return
262
+ }
263
+ host := utils.NormalizeHostname(announceURL.Hostname())
264
+ if utils.IsLocalRelayHost(host) {
265
+ utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest,
266
+ fmt.Sprintf("self-announce rejected: host %q is local-only", host))
267
+ return
268
+ }
269
+ if selfURL, err := utils.NormalizeRelayURL(s.cfg.PortalURL); err == nil && desc.APIHTTPSAddr == selfURL {
270
+ utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest,
271
+ fmt.Sprintf("self-announce rejected: %q matches receiving relay url", desc.APIHTTPSAddr))
272
+ return
273
+ }
274
+ if host != "" && host == utils.NormalizeHostname(s.identity.Name) {
275
+ utils.WriteAPIError(w, http.StatusBadRequest, types.APIErrorCodeInvalidRequest,
276
+ fmt.Sprintf("self-announce rejected: host %q matches receiving relay host", host))
277
+ return
278
}
279
280
now := time.Now().UTC()
portal/discovery/mols.go
+82
-43
@@ -1,16 +1,15 @@
1
package discovery
2
3
-// MOLSRelayPolicy implements RelayPolicy using a Multi-path Orthogonal Latin
4
-// Squares (MOLS) engine over GF(2^6). SelectPriority provides deterministic,
5
-// load-balanced, and collision-resistant relay scoring without requiring a
6
-// central coordinator.
3
+// MOLSRelayPolicy uses a GF(64) MOLS-derived score as the primary
4
+// deterministic ordering for eligible relays. Health and freshness gates decide
5
+// eligibility before the MOLS score is applied.
6
7
// # Core Design
8
//
10
-// The engine uses an order-64 MOLS grid derived from Galois Field GF(64).
11
-// To achieve Magic Square properties (Sum_row = Sum_col = Sum_diag), the
12
-// construction follows a structured mapping where the composite score
13
-// balances the field elements across the 4096-element space.
9
+// The engine uses an order-64 grid derived from Galois Field GF(64). The
10
+// composite score is deterministic for a (client identity, relay URL) pair and
11
+// drives ordering after freshness and failure-suppression gates. Confirmation
12
+// and RTT remain tie-breakers for equal scores.
13
//
14
// L_m[i][j] = gf64Mul(m, i) XOR j (Latin-square row for multiplier m)
15
// score(i, j) = L_m1[i][j] * 64 + L_m2[i][j] + 1 (composite, range 1..4096)
@@ -22,20 +21,22 @@ package discovery
21
//
22
// congestionScore(i, j) = (n^2+1) - score(i, 63-j)
23
//
25
-// This mirrors the priority ordering so underutilised paths move to the front.
24
+// This mirrors the deterministic tie-break order when the whole observed pool
25
+// appears slow.
26
//
27
// # Non-Linear Load (Variant Grid)
28
//
29
-// When the coefficient of variation of per-relay RTTs exceeds molsCVThreshold
30
-// (indicating bursty load), the engine switches multipliers from (3, 5) to
31
-// (7, 11). Non-linear detection takes precedence over congestion switching.
29
+// When the coefficient of variation of per-relay discovery RTTs exceeds
30
+// molsCVThreshold, the engine switches multipliers from (3, 5) to (7, 11).
31
+// Non-linear detection takes precedence over congestion switching.
32
//
33
// # Health & Fallback
34
//
35
// Relays whose measured discovery RTT exceeds molsFallbackRTTThreshold are
36
-// treated as Fallback and placed at the end of the priority queue. The engine
37
-// ensures at least molsMinActiveNodes non-fallback relays remain reachable; if
38
-// fewer are available, Fallback relays are promoted to meet the minimum.
36
+// treated as Fallback and placed at the end of the priority queue. Discovery
37
+// polling failures and SDK listener failures are tracked separately so a
38
+// discovery retry delay does not by itself remove an otherwise active relay
39
+// candidate.
40
41
import (
42
"hash/fnv"
@@ -58,6 +59,7 @@ const (
59
molsCVThreshold = 0.5
60
molsFallbackRTTThreshold = 2 * time.Second
61
molsMinActiveNodes = 2
62
+ defaultMaxActiveRelays = 3
63
)
64
65
// gf64Mul performs multiplication in GF(2^6) with primitive polynomial x^6 + x + 1 (0x43).
@@ -84,15 +86,7 @@ func molsScore(i, j, m1, m2 uint8) int {
86
l1 := gf64Mul(m1, i) ^ j
87
l2 := gf64Mul(m2, i) ^ j
88
87
- // Magic Square Diagonal correction:
88
- // To ensure diagonal sums match row/col sums (131,104), we apply a
89
- // deterministic permutation based on the field property of GF(64).
89
score := int(l1)*molsOrder + int(l2) + 1
91
-
92
- // Semi-magic to Magic conversion for GF(2^n) grids
93
- if i == j {
94
- return score
95
- }
90
return score
91
}
92
@@ -163,9 +157,10 @@ func (p MOLSRelayPolicy) SelectConfirmed(states []RelayState) []RelayState {
157
return out
158
}
159
166
-func (p MOLSRelayPolicy) OnConfirmed(state RelayState) RelayState {
160
+func (p MOLSRelayPolicy) OnActiveConfirmed(state RelayState) RelayState {
161
state.Confirmed = true
168
- state.consecutiveFailures = 0 // Critical fix: reset failures on success
162
+ state.activeFailures = 0
163
+ state.suppressActiveUntil = time.Time{}
164
return state
165
}
166
@@ -174,20 +169,40 @@ func (p MOLSRelayPolicy) OnUnconfirmed(state RelayState) RelayState {
169
return state
170
}
171
177
-func (p MOLSRelayPolicy) OnFailure(state RelayState, err error, recoveryFailures int) (RelayState, bool, string) {
178
- state.consecutiveFailures++
172
+func (p MOLSRelayPolicy) OnDiscoveryConfirmed(state RelayState) RelayState {
173
+ state.discoveryFailures = 0
174
+ state.nextDiscoveryRefreshAt = time.Time{}
175
+ return state
176
+}
177
+
178
+func (p MOLSRelayPolicy) OnDiscoveryFailure(state RelayState, err error, recoveryFailures int) (RelayState, bool, string) {
179
+ state.discoveryFailures++
180
180
- // Exponential backoff
181
- backoff := 1 * time.Second << min(state.consecutiveFailures, 6)
182
- if backoff > 60*time.Second {
183
- backoff = 60 * time.Second
181
+ if recoveryFailures <= 0 || state.discoveryFailures < recoveryFailures {
182
+ return state, false, "retry"
183
}
185
- state.nextDirectRefreshAt = time.Now().Add(backoff)
184
+ failuresOverBudget := state.discoveryFailures - recoveryFailures
185
+ backoff := defaultDirectRecoveryBackoff << min(failuresOverBudget, 3)
186
+ if backoff > maxDirectRecoveryBackoff {
187
+ backoff = maxDirectRecoveryBackoff
188
+ }
189
+ state.nextDiscoveryRefreshAt = time.Now().Add(backoff)
190
+ return state, true, "discovery"
191
+}
192
187
- if state.consecutiveFailures < recoveryFailures {
193
+func (p MOLSRelayPolicy) OnActiveFailure(state RelayState, err error, recoveryFailures int) (RelayState, bool, string) {
194
+ state.activeFailures++
195
+
196
+ if recoveryFailures <= 0 || state.activeFailures < recoveryFailures {
197
return state, false, "retry"
198
}
190
- return state, true, "recovery"
199
+ failuresOverBudget := state.activeFailures - recoveryFailures
200
+ backoff := defaultDirectRecoveryBackoff << min(failuresOverBudget, 3)
201
+ if backoff > maxDirectRecoveryBackoff {
202
+ backoff = maxDirectRecoveryBackoff
203
+ }
204
+ state.suppressActiveUntil = time.Now().Add(backoff)
205
+ return state, true, "active"
206
}
207
208
func (p MOLSRelayPolicy) OnBanned(state RelayState) RelayState {
@@ -288,27 +303,48 @@ func (p MOLSRelayPolicy) SelectPriority(states []RelayState, clientState ClientS
303
return nil
304
}
305
306
+ now := time.Now().UTC()
307
explicit := make([]string, 0)
308
autoPool := make([]RelayState, 0, len(selected))
309
for _, state := range selected {
294
- if clientState.RequireUDP && state.hasObservedDescriptor() && !state.Descriptor.SupportsUDP {
295
- continue
296
- }
297
- if clientState.RequireTCP && state.hasObservedDescriptor() && !state.Descriptor.SupportsTCP {
298
- continue
299
- }
300
-
310
relayURL := state.Descriptor.APIHTTPSAddr
311
if slices.Contains(clientState.ExplicitRelayURLs, relayURL) {
312
+ if state.hasObservedDescriptor() && state.Descriptor.ExpiresAt.After(now) {
313
+ if clientState.RequireUDP && !state.Descriptor.SupportsUDP {
314
+ continue
315
+ }
316
+ if clientState.RequireTCP && !state.Descriptor.SupportsTCP {
317
+ continue
318
+ }
319
+ }
320
explicit = append(explicit, relayURL)
321
continue
322
}
323
+
324
+ if state.hasObservedDescriptor() {
325
+ if !state.Descriptor.ExpiresAt.After(now) {
326
+ continue
327
+ }
328
+ if clientState.RequireUDP && !state.Descriptor.SupportsUDP {
329
+ continue
330
+ }
331
+ if clientState.RequireTCP && !state.Descriptor.SupportsTCP {
332
+ continue
333
+ }
334
+ }
335
+ if !state.suppressActiveUntil.IsZero() && state.suppressActiveUntil.After(now) {
336
+ continue
337
+ }
338
autoPool = append(autoPool, state)
339
}
340
341
autoURLs := p.rankRelayPool(autoPool, clientState.LocalAddress)
310
- if clientState.MaxActiveRelays > 0 && len(autoURLs) > clientState.MaxActiveRelays {
311
- autoURLs = autoURLs[:clientState.MaxActiveRelays]
342
+ maxActiveRelays := clientState.MaxActiveRelays
343
+ if maxActiveRelays <= 0 {
344
+ maxActiveRelays = defaultMaxActiveRelays
345
+ }
346
+ if len(autoURLs) > maxActiveRelays {
347
+ autoURLs = autoURLs[:maxActiveRelays]
348
}
349
return append(explicit, autoURLs...)
350
}
@@ -335,6 +371,9 @@ func (p MOLSRelayPolicy) SelectMultiHop(states []RelayState, clientState ClientS
371
if !state.hasObservedDescriptor() || !state.Descriptor.ExpiresAt.After(now) || !state.Descriptor.HasOverlayPeer() {
372
continue
373
}
374
+ if !state.suppressActiveUntil.IsZero() && state.suppressActiveUntil.After(now) {
375
+ continue
376
+ }
377
autoPool = append(autoPool, state)
378
}
379
portal/discovery/mols_test.go
+70
@@ -460,6 +460,76 @@ func TestMOLSSelectPriorityMaxActiveRelaysLimitsAutoPool(t *testing.T) {
460
}
461
}
462
463
+func TestMOLSSelectPriorityZeroMaxActiveRelaysUsesDefault(t *testing.T) {
464
+ policy := MOLSRelayPolicy{}
465
+
466
+ relays := make([]RelayState, 10)
467
+ for i := range relays {
468
+ relays[i] = confirmedPolicyRelayState(t, fmt.Sprintf("https://relay-default-%d.example", i))
469
+ }
470
+
471
+ selected := policy.SelectPriority(relays, ClientState{MaxActiveRelays: 0})
472
+ if len(selected) != defaultMaxActiveRelays {
473
+ t.Fatalf("len(selected) = %d, want %d", len(selected), defaultMaxActiveRelays)
474
+ }
475
+}
476
+
477
+func TestMOLSSelectPrioritySkipsExpiredAutoRelay(t *testing.T) {
478
+ policy := MOLSRelayPolicy{}
479
+ expired := confirmedPolicyRelayState(t, "https://relay-expired.example")
480
+ expired.Descriptor.ExpiresAt = time.Now().UTC().Add(-time.Minute)
481
+
482
+ if selected := policy.SelectPriority([]RelayState{expired}, ClientState{}); len(selected) != 0 {
483
+ t.Fatalf("SelectPriority(expired auto) = %v, want empty", selected)
484
+ }
485
+}
486
+
487
+func TestMOLSSelectPriorityKeepsExpiredExplicitRelay(t *testing.T) {
488
+ policy := MOLSRelayPolicy{}
489
+ relayURL := "https://relay-explicit-expired.example"
490
+ expired := confirmedPolicyRelayState(t, relayURL)
491
+ expired.Descriptor.ExpiresAt = time.Now().UTC().Add(-time.Minute)
492
+
493
+ selected := policy.SelectPriority([]RelayState{expired}, ClientState{
494
+ ExplicitRelayURLs: []string{relayURL},
495
+ })
496
+ if len(selected) != 1 || selected[0] != relayURL {
497
+ t.Fatalf("SelectPriority(expired explicit) = %v, want [%q]", selected, relayURL)
498
+ }
499
+}
500
+
501
+func TestMOLSSelectPrioritySkipsAutoRelayInBackoff(t *testing.T) {
502
+ policy := MOLSRelayPolicy{}
503
+ backingOff := confirmedPolicyRelayState(t, "https://relay-backoff.example")
504
+ backingOff.suppressActiveUntil = time.Now().UTC().Add(time.Minute)
505
+
506
+ if selected := policy.SelectPriority([]RelayState{backingOff}, ClientState{}); len(selected) != 0 {
507
+ t.Fatalf("SelectPriority(backing off auto) = %v, want empty", selected)
508
+ }
509
+}
510
+
511
+func TestMOLSSelectPriorityKeepsDiscoveryBackoffRelay(t *testing.T) {
512
+ policy := MOLSRelayPolicy{}
513
+ relayURL := "https://relay-discovery-backoff.example"
514
+ backingOff := confirmedPolicyRelayState(t, relayURL)
515
+ backingOff.nextDiscoveryRefreshAt = time.Now().UTC().Add(time.Minute)
516
+
517
+ selected := policy.SelectPriority([]RelayState{backingOff}, ClientState{})
518
+ if len(selected) != 1 || selected[0] != relayURL {
519
+ t.Fatalf("SelectPriority(discovery backoff) = %v, want [%q]", selected, relayURL)
520
+ }
521
+}
522
+
523
+func TestMOLSSelectPriorityKeepsUnobservedAutoSeed(t *testing.T) {
524
+ policy := MOLSRelayPolicy{}
525
+ relayURL := "https://relay-seed.example"
526
+
527
+ selected := policy.SelectPriority([]RelayState{bootstrapPolicyRelayState(relayURL)}, ClientState{})
528
+ if len(selected) != 1 || selected[0] != relayURL {
529
+ t.Fatalf("SelectPriority(unobserved seed) = %v, want [%q]", selected, relayURL)
530
+ }
531
+}
532
+
533
// TestMOLSMagicRowSum verifies that each row of the base MOLS score grid sums
534
// to the magic constant n*(n²+1)/2 = 131104.
535
func TestMOLSMagicRowSum(t *testing.T) {
portal/discovery/policy_test.go
+35
-10
@@ -118,7 +118,7 @@ func TestSelectPriorityCongestionInversion(t *testing.T) {
118
r1, r2 := "https://r1.net", "https://r2.net"
119
states := []RelayState{
120
confirmedPolicyRelayStateWithRTT(t, r1, 800*time.Millisecond),
121
- confirmedPolicyRelayStateWithRTT(t, r2, 900*time.Millisecond),
121
+ confirmedPolicyRelayStateWithRTT(t, r2, 800*time.Millisecond),
122
}
123
124
selected := policy.SelectPriority(states, ClientState{LocalAddress: clientAddr})
@@ -164,24 +164,28 @@ func TestSelectPriorityFallbackPromotion(t *testing.T) {
164
}
165
}
166
167
-func TestOnConfirmedResetsFailures(t *testing.T) {
167
+func TestOnActiveConfirmedResetsActiveFailures(t *testing.T) {
168
policy := MOLSRelayPolicy{}
169
state := RelayState{
170
- consecutiveFailures: 5,
170
+ activeFailures: 5,
171
+ suppressActiveUntil: time.Now().UTC().Add(time.Minute),
172
Confirmed: false,
173
}
174
174
- state = policy.OnConfirmed(state)
175
+ state = policy.OnActiveConfirmed(state)
176
177
if !state.Confirmed {
178
t.Fatal("Confirmed should be true")
179
}
179
- if state.consecutiveFailures != 0 {
180
- t.Errorf("consecutiveFailures = %d, want 0", state.consecutiveFailures)
180
+ if state.activeFailures != 0 {
181
+ t.Errorf("activeFailures = %d, want 0", state.activeFailures)
182
+ }
183
+ if !state.suppressActiveUntil.IsZero() {
184
+ t.Errorf("suppressActiveUntil = %v, want zero", state.suppressActiveUntil)
185
}
186
}
187
184
-func TestOnFailureBackoff(t *testing.T) {
188
+func TestOnDiscoveryFailureBackoff(t *testing.T) {
189
policy := MOLSRelayPolicy{}
190
state := confirmedPolicyRelayState(t, "https://error.io")
191
budget := 3
@@ -189,13 +193,34 @@ func TestOnFailureBackoff(t *testing.T) {
193
start := time.Now()
194
for i := 0; i < budget; i++ {
195
var backed bool
192
- state, backed, _ = policy.OnFailure(state, errors.New("err"), budget)
196
+ state, backed, _ = policy.OnDiscoveryFailure(state, errors.New("err"), budget)
197
if i < budget-1 && backed {
198
t.Fatal("Premature backoff")
199
}
200
}
201
198
- if !state.nextDirectRefreshAt.After(start) {
199
- t.Fatal("Retry timer not scheduled")
202
+ if !state.nextDiscoveryRefreshAt.After(start) {
203
+ t.Fatal("discovery retry timer not scheduled")
204
+ }
205
+ if !state.suppressActiveUntil.IsZero() {
206
+ t.Fatalf("suppressActiveUntil = %v, want zero", state.suppressActiveUntil)
207
+ }
208
+}
209
+
210
+func TestOnActiveFailureBackoff(t *testing.T) {
211
+ policy := MOLSRelayPolicy{}
212
+ state := confirmedPolicyRelayState(t, "https://error.io")
213
+ start := time.Now()
214
+
215
+ var backed bool
216
+ state, backed, _ = policy.OnActiveFailure(state, errors.New("err"), 1)
217
+ if !backed {
218
+ t.Fatal("active failure should back off at budget")
219
+ }
220
+ if !state.suppressActiveUntil.After(start) {
221
+ t.Fatal("active suppression timer not scheduled")
222
+ }
223
+ if !state.nextDiscoveryRefreshAt.IsZero() {
224
+ t.Fatalf("nextDiscoveryRefreshAt = %v, want zero", state.nextDiscoveryRefreshAt)
225
}
226
}
portal/discovery/refresher.go
+7
-4
@@ -153,7 +153,7 @@ func (r *Refresher) refreshHTTPS(ctx context.Context) error {
153
continue
154
}
155
} else if !state.Bootstrap {
156
- if !state.nextDirectRefreshAt.IsZero() && state.nextDirectRefreshAt.After(now) {
156
+ if !state.nextDiscoveryRefreshAt.IsZero() && state.nextDiscoveryRefreshAt.After(now) {
157
continue
158
}
159
}
@@ -209,6 +209,9 @@ func (r *Refresher) refreshOverlay(ctx context.Context) error {
209
}
210
relaySetChanged := false
211
for _, state := range states {
212
+ if !state.nextDiscoveryRefreshAt.IsZero() && state.nextDiscoveryRefreshAt.After(time.Now().UTC()) {
213
+ continue
214
+ }
215
relay := state.Descriptor
216
recoveryFailures := r.directRecoveryFailures
217
if state.Bootstrap {
@@ -249,7 +252,7 @@ func (r *Refresher) refreshOverlay(ctx context.Context) error {
252
}
253
254
func (r *Refresher) logDiscoveryFailure(targetRelayURL, sourceURL string, recoveryFailures int, err error) {
252
- backedOff, backoffReason, consecutiveFailures := r.relaySet.RecordRelayFailure(targetRelayURL, err, recoveryFailures)
255
+ backedOff, backoffReason, failureCount := r.relaySet.RecordDiscoveryFailure(targetRelayURL, err, recoveryFailures)
256
if !backedOff {
257
return
258
}
@@ -259,8 +262,8 @@ func (r *Refresher) logDiscoveryFailure(targetRelayURL, sourceURL string, recove
262
Str("relay", sourceURL).
263
Bool("backed_off", true).
264
Str("reason", backoffReason)
262
- if consecutiveFailures > 0 {
263
- event = event.Int("consecutive_failures", consecutiveFailures)
265
+ if failureCount > 0 {
266
+ event = event.Int("discovery_failures", failureCount)
267
}
268
event.Msg("discovery source retry delayed")
269
}
portal/discovery/relayset.go
+40
-18
@@ -175,7 +175,9 @@ func (s *RelaySet) SetBootstrapRelayURLs(inputs []string) {
175
for key, state := range s.relays {
176
_, bootstrap := keep[key]
177
state.Bootstrap = bootstrap
178
- if !state.Bootstrap && !state.hasObservedDescriptor() && !state.Banned && state.consecutiveFailures == 0 {
178
+ if !state.Bootstrap && !state.hasObservedDescriptor() && !state.Banned &&
179
+ state.discoveryFailures == 0 && state.activeFailures == 0 &&
180
+ state.nextDiscoveryRefreshAt.IsZero() && state.suppressActiveUntil.IsZero() {
181
delete(s.relays, key)
182
continue
183
}
@@ -354,7 +356,7 @@ func (s *RelaySet) ConfirmRelayURL(relayURL string) {
356
if !ok {
357
state = newRelayState(relayURL)
358
}
357
- state = s.policy.OnConfirmed(state)
359
+ state = s.policy.OnActiveConfirmed(state)
360
s.relays[relayURL] = state
361
}
362
@@ -425,10 +427,14 @@ func (s *RelaySet) ApplyRelayDiscoveryResponse(targetURL string, resp types.Disc
427
record.Bootstrap = record.Bootstrap || existingAtURL.Bootstrap
428
record.Confirmed = record.Confirmed || existingAtURL.Confirmed
429
record.Banned = record.Banned || existingAtURL.Banned
428
- if record.consecutiveFailures < existingAtURL.consecutiveFailures {
429
- record.consecutiveFailures = existingAtURL.consecutiveFailures
430
+ if record.discoveryFailures < existingAtURL.discoveryFailures {
431
+ record.discoveryFailures = existingAtURL.discoveryFailures
432
}
431
- record.nextDirectRefreshAt = existingAtURL.nextDirectRefreshAt
433
+ if record.activeFailures < existingAtURL.activeFailures {
434
+ record.activeFailures = existingAtURL.activeFailures
435
+ }
436
+ record.nextDiscoveryRefreshAt = existingAtURL.nextDiscoveryRefreshAt
437
+ record.suppressActiveUntil = existingAtURL.suppressActiveUntil
438
if record.DiscoveryRTTAt.IsZero() || (!existingAtURL.DiscoveryRTTAt.IsZero() && existingAtURL.DiscoveryRTTAt.After(record.DiscoveryRTTAt)) {
439
record.DiscoveryRTT = existingAtURL.DiscoveryRTT
440
record.DiscoveryRTTAt = existingAtURL.DiscoveryRTTAt
@@ -436,8 +442,7 @@ func (s *RelaySet) ApplyRelayDiscoveryResponse(targetURL string, resp types.Disc
442
443
isAuthoritativeTarget := !protocolMismatch && !missingTarget && authoritative && relayURL == targetURL
444
if isAuthoritativeTarget {
439
- record.consecutiveFailures = 0
440
- record.nextDirectRefreshAt = time.Time{}
445
+ record = s.policy.OnDiscoveryConfirmed(record)
446
}
447
448
if upsert := s.upsertDescriptorLocked(record, now, isAuthoritativeTarget); upsert != upsertAccepted {
@@ -448,9 +453,8 @@ func (s *RelaySet) ApplyRelayDiscoveryResponse(targetURL string, resp types.Disc
453
// authoritative target we should still credit it as alive on its
454
// existing URL slot.
455
if isAuthoritativeTarget && hasExistingAtURL {
451
- if existingAtURL.consecutiveFailures != 0 || !existingAtURL.nextDirectRefreshAt.IsZero() {
452
- existingAtURL.consecutiveFailures = 0
453
- existingAtURL.nextDirectRefreshAt = time.Time{}
456
+ if existingAtURL.discoveryFailures != 0 || !existingAtURL.nextDiscoveryRefreshAt.IsZero() {
457
+ existingAtURL = s.policy.OnDiscoveryConfirmed(existingAtURL)
458
s.relays[relayURL] = existingAtURL
459
relaySetChanged = true
460
}
@@ -497,8 +501,9 @@ func (s *RelaySet) RecordDiscoveryRTT(relayURL string, rtt time.Duration, measur
501
// future) and not significantly clock-skewed (IssuedAt no further into
502
// the future than AnnounceClockSkewTolerance, validity window no longer
503
// than AnnounceMaxValidity).
500
-// 3. Local merge preserves Bootstrap, Confirmed, Banned, telemetry, and
501
-// direct-refresh retry state from any pre-existing entry at the same URL.
504
+// 3. Local merge preserves Bootstrap, Confirmed, Banned, discovery retry
505
+// state, active suppression state, and telemetry from any pre-existing
506
+// entry at the same URL.
507
// 4. The shared upsertDescriptorLocked method enforces the
508
// monotonic-IssuedAt-per-key rollback guard and the cross-identity
509
// URL-takeover guard. Announce never grants takeover authority; only
@@ -536,10 +541,14 @@ func (s *RelaySet) InsertAnnounced(desc types.RelayDescriptor, now time.Time) er
541
record.Bootstrap = record.Bootstrap || existing.Bootstrap
542
record.Confirmed = record.Confirmed || existing.Confirmed
543
record.Banned = record.Banned || existing.Banned
539
- if record.consecutiveFailures < existing.consecutiveFailures {
540
- record.consecutiveFailures = existing.consecutiveFailures
544
+ if record.discoveryFailures < existing.discoveryFailures {
545
+ record.discoveryFailures = existing.discoveryFailures
546
+ }
547
+ if record.activeFailures < existing.activeFailures {
548
+ record.activeFailures = existing.activeFailures
549
}
542
- record.nextDirectRefreshAt = existing.nextDirectRefreshAt
550
+ record.nextDiscoveryRefreshAt = existing.nextDiscoveryRefreshAt
551
+ record.suppressActiveUntil = existing.suppressActiveUntil
552
if record.DiscoveryRTTAt.IsZero() || (!existing.DiscoveryRTTAt.IsZero() && existing.DiscoveryRTTAt.After(record.DiscoveryRTTAt)) {
553
record.DiscoveryRTT = existing.DiscoveryRTT
554
record.DiscoveryRTTAt = existing.DiscoveryRTTAt
@@ -625,7 +634,20 @@ func (s *RelaySet) enforceCapLocked() {
634
}
635
}
636
628
-func (s *RelaySet) RecordRelayFailure(relayURL string, err error, recoveryFailures int) (backedOff bool, backoffReason string, consecutiveFailures int) {
637
+func (s *RelaySet) RecordDiscoveryFailure(relayURL string, err error, recoveryFailures int) (backedOff bool, backoffReason string, failureCount int) {
638
+ s.mu.Lock()
639
+ defer s.mu.Unlock()
640
+
641
+ state, ok := s.relays[relayURL]
642
+ if !ok {
643
+ return false, "", 0
644
+ }
645
+ state, backedOff, backoffReason = s.policy.OnDiscoveryFailure(state, err, recoveryFailures)
646
+ s.relays[relayURL] = state
647
+ return backedOff, backoffReason, state.discoveryFailures
648
+}
649
+
650
+func (s *RelaySet) RecordActiveFailure(relayURL string, err error, recoveryFailures int) (backedOff bool, backoffReason string, failureCount int) {
651
s.mu.Lock()
652
defer s.mu.Unlock()
653
@@ -633,7 +655,7 @@ func (s *RelaySet) RecordRelayFailure(relayURL string, err error, recoveryFailur
655
if !ok {
656
return false, "", 0
657
}
636
- state, backedOff, backoffReason = s.policy.OnFailure(state, err, recoveryFailures)
658
+ state, backedOff, backoffReason = s.policy.OnActiveFailure(state, err, recoveryFailures)
659
s.relays[relayURL] = state
638
- return backedOff, backoffReason, state.consecutiveFailures
660
+ return backedOff, backoffReason, state.activeFailures
661
}
portal/discovery/relayset_test.go
+24
-16
@@ -102,17 +102,19 @@ func TestApplyRelayDiscoveryResponseCollectsHintsWhenTargetDescriptorIsMissing(t
102
}
103
}
104
105
-func TestApplyRelayDiscoveryResponseClearsDirectRetryOnAuthoritativeSuccess(t *testing.T) {
105
+func TestApplyRelayDiscoveryResponseClearsDiscoveryRetryOnAuthoritativeSuccess(t *testing.T) {
106
set := NewRelaySet(nil)
107
108
relayURL := "https://relay-source.example"
109
desc := mustPolicyRelayDescriptor(t, relayURL)
110
set.mu.Lock()
111
state := RelayState{
112
- Descriptor: desc,
113
- LastSeenAt: time.Now().UTC(),
114
- consecutiveFailures: defaultRecoveryFailures,
115
- nextDirectRefreshAt: time.Now().UTC().Add(time.Minute),
112
+ Descriptor: desc,
113
+ LastSeenAt: time.Now().UTC(),
114
+ discoveryFailures: defaultRecoveryFailures,
115
+ nextDiscoveryRefreshAt: time.Now().UTC().Add(time.Minute),
116
+ activeFailures: 1,
117
+ suppressActiveUntil: time.Now().UTC().Add(time.Minute),
118
}
119
set.relays[relayURL] = state
120
set.mu.Unlock()
@@ -127,25 +129,31 @@ func TestApplyRelayDiscoveryResponseClearsDirectRetryOnAuthoritativeSuccess(t *t
129
set.mu.RLock()
130
refreshed := set.relays[relayURL]
131
set.mu.RUnlock()
130
- if refreshed.consecutiveFailures != 0 {
131
- t.Fatalf("consecutiveFailures = %d, want 0", refreshed.consecutiveFailures)
132
+ if refreshed.discoveryFailures != 0 {
133
+ t.Fatalf("discoveryFailures = %d, want 0", refreshed.discoveryFailures)
134
}
133
- if !refreshed.nextDirectRefreshAt.IsZero() {
134
- t.Fatalf("nextDirectRefreshAt = %v, want zero time", refreshed.nextDirectRefreshAt)
135
+ if !refreshed.nextDiscoveryRefreshAt.IsZero() {
136
+ t.Fatalf("nextDiscoveryRefreshAt = %v, want zero time", refreshed.nextDiscoveryRefreshAt)
137
+ }
138
+ if refreshed.activeFailures != 1 {
139
+ t.Fatalf("activeFailures = %d, want 1", refreshed.activeFailures)
140
+ }
141
+ if refreshed.suppressActiveUntil.IsZero() {
142
+ t.Fatal("suppressActiveUntil was cleared by discovery success")
143
}
144
}
145
138
-func TestApplyRelayDiscoveryResponsePreservesDirectRetryOnHint(t *testing.T) {
146
+func TestApplyRelayDiscoveryResponsePreservesDiscoveryRetryOnHint(t *testing.T) {
147
set := NewRelaySet(nil)
148
149
relayURL := "https://relay-hinted.example"
150
desc := mustPolicyRelayDescriptor(t, relayURL)
143
- nextDirectRefreshAt := time.Now().UTC().Add(time.Minute)
151
+ nextDiscoveryRefreshAt := time.Now().UTC().Add(time.Minute)
152
set.mu.Lock()
153
state := RelayState{
146
- Descriptor: desc,
147
- LastSeenAt: time.Now().UTC(),
148
- nextDirectRefreshAt: nextDirectRefreshAt,
154
+ Descriptor: desc,
155
+ LastSeenAt: time.Now().UTC(),
156
+ nextDiscoveryRefreshAt: nextDiscoveryRefreshAt,
157
}
158
set.relays[relayURL] = state
159
set.mu.Unlock()
@@ -160,8 +168,8 @@ func TestApplyRelayDiscoveryResponsePreservesDirectRetryOnHint(t *testing.T) {
168
set.mu.RLock()
169
refreshed := set.relays[relayURL]
170
set.mu.RUnlock()
163
- if !refreshed.nextDirectRefreshAt.Equal(nextDirectRefreshAt) {
164
- t.Fatalf("nextDirectRefreshAt = %v, want %v", refreshed.nextDirectRefreshAt, nextDirectRefreshAt)
171
+ if !refreshed.nextDiscoveryRefreshAt.Equal(nextDiscoveryRefreshAt) {
172
+ t.Fatalf("nextDiscoveryRefreshAt = %v, want %v", refreshed.nextDiscoveryRefreshAt, nextDiscoveryRefreshAt)
173
}
174
}
175
portal/discovery/relaystate.go
+10
-6
@@ -39,8 +39,10 @@ type RelayState struct {
39
DiscoveryRTT time.Duration
40
DiscoveryRTTAt time.Time
41
42
- consecutiveFailures int
43
- nextDirectRefreshAt time.Time
42
+ discoveryFailures int
43
+ activeFailures int
44
+ nextDiscoveryRefreshAt time.Time
45
+ suppressActiveUntil time.Time
46
}
47
48
func newRelayState(relayURL string) RelayState {
@@ -57,10 +59,12 @@ func (state RelayState) hasObservedDescriptor() bool {
59
60
type ClientState struct {
61
ExplicitRelayURLs []string
60
- MaxActiveRelays int
61
- MultiHopDepth int
62
- RequireUDP bool
63
- RequireTCP bool
62
+ // MaxActiveRelays caps auto-selected relays. Zero or negative values use
63
+ // the policy default of 3.
64
+ MaxActiveRelays int
65
+ MultiHopDepth int
66
+ RequireUDP bool
67
+ RequireTCP bool
68
// LocalAddress is the ingress identity address used by MOLSRelayPolicy to
69
// derive a deterministic row index into the GF(64) MOLS grid.
70
LocalAddress string
sdk/expose.go
+5
-5
@@ -112,11 +112,11 @@ func Expose(ctx context.Context, cfg ExposeConfig) (*Exposure, error) {
112
return nil, err
113
}
114
} else {
115
- listenerRelayURLs, err = utils.ResolvePortalRelayURLs(ctx, explicitRelayURLs, cfg.Discovery)
115
+ relaySetURLs, err = utils.ResolvePortalRelayURLs(ctx, explicitRelayURLs, cfg.Discovery)
116
if err != nil {
117
return nil, err
118
}
119
- relaySetURLs = listenerRelayURLs
119
+ listenerRelayURLs = append([]string(nil), explicitRelayURLs...)
120
}
121
122
identity, createdIdentity, err := utils.ResolveListenerIdentity(
@@ -166,15 +166,15 @@ func Expose(ctx context.Context, cfg ExposeConfig) (*Exposure, error) {
166
relayListeners: make(map[string]*listener, initialRouteCapacity(listenerRelayURLs, cfg.MultiHopDepth)),
167
}
168
169
- if len(multiHop) > 0 || cfg.MultiHopDepth > 1 {
169
+ if cfg.Discovery || len(multiHop) > 0 || cfg.MultiHopDepth > 1 {
170
refresher := discovery.NewRefresher(exposure.relaySet, nil)
171
if err := refresher.Refresh(ctx, nil); err != nil {
172
_ = exposure.Close()
173
- return nil, fmt.Errorf("discover multi-hop relays: %w", err)
173
+ return nil, fmt.Errorf("discover relays: %w", err)
174
}
175
}
176
177
- if len(listenerRelayURLs) > 0 || cfg.MultiHopDepth > 1 {
177
+ if len(listenerRelayURLs) > 0 || cfg.Discovery || cfg.MultiHopDepth > 1 {
178
if err := exposure.reconcileRelayListeners(true); err != nil {
179
_ = exposure.Close()
180
return nil, err
sdk/expose_test.go
+2
@@ -28,6 +28,7 @@ func TestExposureReconcileRemovesBannedRelayFromActiveSet(t *testing.T) {
28
}
29
30
exposure := &Exposure{
31
+ explicitRelays: []string{relayA, relayB},
32
relaySet: mustRelaySet(t, relayA, relayB),
33
relayListeners: make(map[string]*listener, 2),
34
}
@@ -83,6 +84,7 @@ func TestExposureReconcileRemovesStaleListener(t *testing.T) {
84
85
relayAClosed := make(chan struct{})
86
exposure := &Exposure{
87
+ explicitRelays: []string{relayA, relayB},
88
relaySet: mustRelaySet(t, relayA, relayB),
89
relayListeners: make(map[string]*listener, 2),
90
}
sdk/listener.go
+2
-2
@@ -151,7 +151,7 @@ func (l *listener) run(ctx context.Context) {
151
relayURL := l.relayURL.String()
152
if l.relaySet != nil && relayURL != "" {
153
l.relaySet.UnconfirmRelayURL(relayURL)
154
- l.relaySet.RecordRelayFailure(relayURL, err, 1)
154
+ l.relaySet.RecordActiveFailure(relayURL, err, 1)
155
}
156
log.Error().
157
Err(err).
@@ -839,7 +839,7 @@ func (l *listener) waitRetry(ctx context.Context, operation string, err error, r
839
if l.retryCount > 0 && retries > l.retryCount {
840
if l.relaySet != nil && relayURL != "" {
841
l.relaySet.UnconfirmRelayURL(relayURL)
842
- l.relaySet.RecordRelayFailure(relayURL, err, 1)
842
+ l.relaySet.RecordActiveFailure(relayURL, err, 1)
843
}
844
logger.Error().
845
Err(err).