mirror of
https://github.com/p1neappleXpress/OpenFlux.git
synced 2026-10-02 05:04:39 +08:00
Name every carrier data is spread over, not just the first
Session.Send spreads flows over all live carriers sharing the highest priority, but ActiveTransport - and so the apps' "current transport" - named only the first of them. Session.ActiveTransports returns the whole group, computed by the same topLinks helper Send now uses, so the two cannot drift apart. The IPC status gains active_all and the mobile bridge CurrentTransports; active and CurrentTransport are unchanged for existing clients. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
51c6f65ca8
commit
3c605a1c1e
@@ -945,6 +945,7 @@ func ipcStatusLoop(srv *ipc.Server, m *manager.Manager) {
|
||||
BytesOut: st.BytesSent,
|
||||
UptimeMs: time.Since(started).Milliseconds(),
|
||||
Active: m.Session().ActiveTransport(),
|
||||
ActiveAll: m.Session().ActiveTransports(),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package mobile
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"openflux/transport"
|
||||
@@ -47,3 +48,20 @@ func CurrentTransport() string {
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// CurrentTransports names every carrier traffic currently goes through,
|
||||
// joined by ",": several in Session mode when they share the highest
|
||||
// priority. Session carriers go by name ("boards", "boards-2"), a classic
|
||||
// connection by its type; "" when nothing is connected.
|
||||
func CurrentTransports() string {
|
||||
route.mu.Lock()
|
||||
s, classic := route.session, route.classic
|
||||
route.mu.Unlock()
|
||||
if s != nil {
|
||||
return strings.Join(s.ActiveTransports(), ",")
|
||||
}
|
||||
if classic != "" && (IsConnected() || ProxyIsConnected() || ExitIsConnected()) {
|
||||
return classic
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -61,6 +61,9 @@ type StatusPayload struct {
|
||||
// Active names the carrier data currently goes through ("" when none
|
||||
// reaches the peer). Sessions only.
|
||||
Active string `json:"active,omitempty"`
|
||||
// ActiveAll names every carrier data is spread over: several when they
|
||||
// share the highest priority. Active is the first. Sessions only.
|
||||
ActiveAll []string `json:"active_all,omitempty"`
|
||||
}
|
||||
|
||||
type CommandPayload struct {
|
||||
|
||||
+33
-7
@@ -529,6 +529,38 @@ func (s *Session) ActiveTransport() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
// ActiveTransports names every carrier data goes through now: the
|
||||
// highest-priority live ones, which Send spreads flows over. ActiveTransport
|
||||
// is the first of them. Empty before the handshake or after Stop.
|
||||
func (s *Session) ActiveTransports() []string {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if !s.ready || s.stopped {
|
||||
return nil
|
||||
}
|
||||
var names []string
|
||||
for _, l := range topLinks(s.liveLinksLocked()) {
|
||||
names = append(names, l.name)
|
||||
}
|
||||
return names
|
||||
}
|
||||
|
||||
// topLinks is the group of live links Send uses: the first one and every
|
||||
// following one of the same priority.
|
||||
func topLinks(links []*transportLink) []*transportLink {
|
||||
if len(links) == 0 {
|
||||
return nil
|
||||
}
|
||||
top := links[:1:1]
|
||||
for _, l := range links[1:] {
|
||||
if l.priority != links[0].priority {
|
||||
break
|
||||
}
|
||||
top = append(top, l)
|
||||
}
|
||||
return top
|
||||
}
|
||||
|
||||
// LiveTransports names the carriers that currently reach the peer, highest
|
||||
// priority first. Empty before the handshake or after Stop.
|
||||
func (s *Session) LiveTransports() []string {
|
||||
@@ -651,13 +683,7 @@ func (s *Session) Send(p []byte) error {
|
||||
}
|
||||
raw = append(raw, p...)
|
||||
|
||||
top := links[:1]
|
||||
for _, l := range links[1:] {
|
||||
if l.priority != links[0].priority {
|
||||
break
|
||||
}
|
||||
top = append(top, l)
|
||||
}
|
||||
top := topLinks(links)
|
||||
key := extractFlowKeyBytes(p)
|
||||
idx := int(flowHashBytes(key) % uint64(len(top)))
|
||||
chosen := top[idx]
|
||||
|
||||
@@ -15,6 +15,12 @@ var testParams = PeerParameters{
|
||||
// linkedSessions builds a client and an exit joined by one wire pair per
|
||||
// name, in the given priority order.
|
||||
func linkedSessions(t *testing.T, names ...string) (client, exit *Session, cw, ew map[string]*startCountingWire) {
|
||||
t.Helper()
|
||||
return linkedSessionsWith(t, func(i int) int { return 100 - i*10 }, names...)
|
||||
}
|
||||
|
||||
// linkedSessionsWith is linkedSessions with the i-th carrier at priority(i).
|
||||
func linkedSessionsWith(t *testing.T, priority func(i int) int, names ...string) (client, exit *Session, cw, ew map[string]*startCountingWire) {
|
||||
t.Helper()
|
||||
var err error
|
||||
if client, err = NewSession(testParams, false); err != nil {
|
||||
@@ -33,11 +39,10 @@ func linkedSessions(t *testing.T, names ...string) (client, exit *Session, cw, e
|
||||
a, b := &startCountingWire{}, &startCountingWire{}
|
||||
a.peer, b.peer = &b.negotiationWire, &a.negotiationWire
|
||||
cw[name], ew[name] = a, b
|
||||
priority := 100 - i*10
|
||||
if err := client.AddTransport(name, a, testSessionSecret, testSessionCtx, priority); err != nil {
|
||||
if err := client.AddTransport(name, a, testSessionSecret, testSessionCtx, priority(i)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := exit.AddTransport(name, b, testSessionSecret, testSessionCtx, priority); err != nil {
|
||||
if err := exit.AddTransport(name, b, testSessionSecret, testSessionCtx, priority(i)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -106,3 +106,25 @@ func TestSessionToleratesPeerWithoutKeepalive(t *testing.T) {
|
||||
}
|
||||
eventually(t, "data delivery", func() bool { return got.Load() == 1 })
|
||||
}
|
||||
|
||||
// Carriers of equal top priority share the traffic, so the app must be told
|
||||
// all of them, and not the lower-priority standby.
|
||||
func TestSessionActiveTransportsNameTheWholeTopGroup(t *testing.T) {
|
||||
priorities := []int{100, 100, 50}
|
||||
client, exit, _, _ := linkedSessionsWith(t, func(i int) int { return priorities[i] }, "boards", "yandex", "direct")
|
||||
fastKeepalive(client, exit)
|
||||
startPair(t, client, exit)
|
||||
|
||||
eventually(t, "both top-priority carriers to be active", func() bool {
|
||||
got := client.ActiveTransports()
|
||||
return len(got) == 2 && got[0] == "boards" && got[1] == "yandex"
|
||||
})
|
||||
if got := client.ActiveTransport(); got != "boards" {
|
||||
t.Fatalf("ActiveTransport = %q, want the first of the group", got)
|
||||
}
|
||||
|
||||
_ = client.Stop()
|
||||
if got := client.ActiveTransports(); got != nil {
|
||||
t.Fatalf("ActiveTransports after Stop = %v, want none", got)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user