Skip to content

Commit a709fe0

Browse files
author
Marko Petzold
committed
router/dealer: implement force_reregister RegisterOption
Honor the WAMP advanced-profile force_reregister option on REGISTER: when set to true and an existing single-policy registration is present, forcibly evict the prior callee(s) and install the new registration instead of returning wamp.error.procedure_already_exists. - Advertise force_reregister=true in dealer role features. - Add wamp.OptForceReregister option key + wamp.FeatureForceReregister. - Add wamp.ErrUnregistered (wamp.error.unregistered) revocation reason. - Extend Unregistered with an optional Details dict so the dealer can emit the unsolicited [UNREGISTERED, 0, Details] revocation form carrying Details.registration and Details.reason to evicted callees. Normal [UNREGISTERED, Request] replies stay byte-for-byte identical (omitempty). - syncEvictRegistration tears down the existing registration, notifies every attached callee, scrubs calleeRegIDSet / procRegMap / registrations, and emits on_unregister (per callee) + on_delete meta events. The new REGISTER then falls through the fresh-registration branch, producing a new regID and the usual on_create + on_register meta events. - force_reregister applies only when the existing registration's invocation policy is single (or empty); shared registrations are untouched, falling through to the existing policy-conflict path. Tests: - router/dealer_test.go: happy-path eviction, shared-reg ignore-path, feature-advertisement pin. - test/spec_force_reregister_test.go: end-to-end across the transport × serializer matrix; verifies subsequent CALLs route to the new callee.
1 parent a07a23d commit a709fe0

9 files changed

Lines changed: 299 additions & 7 deletions

File tree

README.md

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
# WAMP v2 router library, client library and router service
44

55
[![Main CI](https://github.com/gammazero/nexus/actions/workflows/main-golint.yml/badge.svg)](https://github.com/gammazero/nexus/actions/workflows/main-golint.yml)
6-
[![Coverage](https://img.shields.io/badge/coverage-62.4%25-orange)](https://github.com/gammazero/nexus/actions/workflows/main-golint.yml)
6+
[![Coverage](https://img.shields.io/badge/coverage-62.6%25-orange)](https://github.com/gammazero/nexus/actions/workflows/main-golint.yml)
77
[![License](https://img.shields.io/badge/License-MIT-blue.svg)](LICENSE)
88
[![GoDoc](https://godoc.org/github.com/gammazero/nexus?status.svg)](https://godoc.org/github.com/gammazero/nexus)
99

@@ -171,6 +171,7 @@ The currently maintained version of this module is 3.x. Earlier major versions a
171171
| registration_meta_procedures | Yes |
172172
| pattern_based_registration | Yes |
173173
| shared_registration | Yes |
174+
| force_reregister | Yes |
174175
| sharded_registration | No |
175176
| registration_revocation | No |
176177
| procedure_reflection | No |

router/dealer.go

Lines changed: 70 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ var dealerRole = wamp.Dict{ //nolint:gochecknoglobals
3333
wamp.FeatureProgCallInvocations: true,
3434
wamp.FeatureSessionMetaAPI: true,
3535
wamp.FeatureSharedReg: true,
36+
wamp.FeatureForceReregister: true,
3637
wamp.FeatureRegMetaAPI: true,
3738
wamp.FeatureTestamentMetaAPI: true,
3839
wamp.FeaturePayloadPassthruMode: true,
@@ -260,10 +261,11 @@ func (d *dealer) Register(callee *wamp.Session, msg *wamp.Register) {
260261

261262
invoke, _ := wamp.AsString(msg.Options[wamp.OptInvoke])
262263
forwardTimeout, _ := msg.Options[wamp.OptForwardTimeout].(bool)
264+
forceReregister, _ := msg.Options[wamp.OptForceReregister].(bool)
263265
var metaPubs []*wamp.Publish
264266
done := make(chan struct{})
265267
d.actionChan <- func() {
266-
metaPubs = d.syncRegister(callee, msg, match, invoke, disclose, forwardTimeout, wampURI)
268+
metaPubs = d.syncRegister(callee, msg, match, invoke, disclose, forwardTimeout, forceReregister, wampURI)
267269
close(done)
268270
}
269271
<-done
@@ -555,7 +557,7 @@ func (d *dealer) run() {
555557
close(d.stopped)
556558
}
557559

558-
func (d *dealer) syncRegister(callee *wamp.Session, msg *wamp.Register, match, invokePolicy string, disclose, forwardTimeout, wampURI bool) []*wamp.Publish { //nolint:lll
560+
func (d *dealer) syncRegister(callee *wamp.Session, msg *wamp.Register, match, invokePolicy string, disclose, forwardTimeout, forceReregister, wampURI bool) []*wamp.Publish { //nolint:lll
559561
var metaPubs []*wamp.Publish
560562
var reg *registration
561563
switch match {
@@ -567,6 +569,15 @@ func (d *dealer) syncRegister(callee *wamp.Session, msg *wamp.Register, match, i
567569
reg = d.wcProcRegMap[msg.Procedure]
568570
}
569571

572+
// force_reregister lets a new callee forcibly take over a procedure
573+
// currently held under invoke=single. Evict the prior callee(s), send
574+
// them an unsolicited UNREGISTERED, fire meta events, then fall through
575+
// to the fresh-registration branch below.
576+
if reg != nil && forceReregister && (reg.policy == "" || reg.policy == wamp.InvokeSingle) {
577+
metaPubs = append(metaPubs, d.syncEvictRegistration(reg, wampURI)...)
578+
reg = nil
579+
}
580+
570581
var created string
571582
var regID wamp.ID
572583
// If no existing registration found for the procedure, then create a new
@@ -1459,6 +1470,63 @@ func (d *dealer) syncRemoveSession(sess *wamp.Session) []*wamp.Publish {
14591470
return metaPubs
14601471
}
14611472

1473+
// syncEvictRegistration tears down an existing registration on behalf of a
1474+
// force_reregister request. It sends an unsolicited UNREGISTERED message to
1475+
// every callee currently attached to reg, drops reg from all dealer indexes
1476+
// and per-callee reg sets, and returns the meta events that should be
1477+
// published (one on_unregister per callee, plus a single on_delete).
1478+
//
1479+
// Callers run on the dealer actor goroutine.
1480+
func (d *dealer) syncEvictRegistration(reg *registration, wampURI bool) []*wamp.Publish {
1481+
var metaPubs []*wamp.Publish
1482+
prevRegID := reg.id
1483+
for _, prev := range reg.callees {
1484+
d.trySend(prev, &wamp.Unregistered{
1485+
Details: wamp.Dict{
1486+
"registration": prevRegID,
1487+
"reason": string(wamp.ErrUnregistered),
1488+
},
1489+
})
1490+
if set, ok := d.calleeRegIDSet[prev]; ok {
1491+
delete(set, prevRegID)
1492+
if len(set) == 0 {
1493+
delete(d.calleeRegIDSet, prev)
1494+
}
1495+
}
1496+
if !wampURI && d.metaPeer != nil {
1497+
metaPubs = append(metaPubs, &wamp.Publish{
1498+
Request: wamp.GlobalID(),
1499+
Topic: wamp.MetaEventRegOnUnregister,
1500+
Arguments: wamp.List{prev.ID, prevRegID},
1501+
})
1502+
}
1503+
}
1504+
delete(d.registrations, prevRegID)
1505+
switch reg.match {
1506+
default:
1507+
delete(d.procRegMap, reg.procedure)
1508+
case wamp.MatchPrefix:
1509+
delete(d.pfxProcRegMap, reg.procedure)
1510+
case wamp.MatchWildcard:
1511+
delete(d.wcProcRegMap, reg.procedure)
1512+
}
1513+
if !wampURI && d.metaPeer != nil && len(reg.callees) > 0 {
1514+
// on_delete uses the last callee's session ID, mirroring the order
1515+
// upstream uses elsewhere (see syncRemoveSession).
1516+
last := reg.callees[len(reg.callees)-1]
1517+
metaPubs = append(metaPubs, &wamp.Publish{
1518+
Request: wamp.GlobalID(),
1519+
Topic: wamp.MetaEventRegOnDelete,
1520+
Arguments: wamp.List{last.ID, prevRegID},
1521+
})
1522+
}
1523+
if d.debug {
1524+
d.log.Printf("Evicted registration %v for procedure %v (force_reregister)",
1525+
prevRegID, reg.procedure)
1526+
}
1527+
return metaPubs
1528+
}
1529+
14621530
// syncDelCalleeReg deletes the the callee from the specified registration and
14631531
// deletes the registration from the set of registrations for the callee.
14641532
//

router/dealer_test.go

Lines changed: 149 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1527,3 +1527,152 @@ func TestReceiveProgressForwardedWithoutCallCanceling(t *testing.T) {
15271527
"dealer should forward receive_progress when callee declares "+
15281528
"progressive_call_results, regardless of call_canceling")
15291529
}
1530+
1531+
// TestForceReregister pins the WAMP advanced-profile force_reregister
1532+
// behavior: a second REGISTER for the same procedure with
1533+
// force_reregister=true must evict the prior callee, send it an
1534+
// unsolicited UNREGISTERED with Details.{registration,reason}, fire the
1535+
// matching meta events, and install the new callee.
1536+
func TestForceReregister(t *testing.T) {
1537+
dealer, metaClient := newTestDealer(t)
1538+
1539+
// Callee1 registers normally (default invoke=single).
1540+
callee1 := newTestPeer()
1541+
sess1 := wamp.NewSession(callee1, 0, nil, nil)
1542+
dealer.Register(sess1, &wamp.Register{Request: 1, Procedure: testProcedure})
1543+
regID1 := (<-callee1.Recv()).(*wamp.Registered).Registration
1544+
checkMetaReg(t, metaClient, sess1.ID) // on_create
1545+
checkMetaReg(t, metaClient, sess1.ID) // on_register
1546+
1547+
// Without force_reregister, a second REGISTER must fail.
1548+
callee2 := newTestPeer()
1549+
sess2 := wamp.NewSession(callee2, 0, nil, nil)
1550+
dealer.Register(sess2, &wamp.Register{Request: 2, Procedure: testProcedure})
1551+
errMsg, ok := (<-callee2.Recv()).(*wamp.Error)
1552+
require.True(t, ok, "expected ERROR without force_reregister")
1553+
require.Equal(t, wamp.ErrProcedureAlreadyExists, errMsg.Error)
1554+
1555+
// Now with force_reregister=true: callee2 should win.
1556+
dealer.Register(sess2, &wamp.Register{
1557+
Request: 3,
1558+
Procedure: testProcedure,
1559+
Options: wamp.SetOption(nil, wamp.OptForceReregister, true),
1560+
})
1561+
1562+
// Callee1 receives an unsolicited UNREGISTERED carrying revocation details.
1563+
rev, ok := (<-callee1.Recv()).(*wamp.Unregistered)
1564+
require.True(t, ok, "evicted callee should receive UNREGISTERED")
1565+
require.Equal(t, wamp.ID(0), rev.Request, "revocation UNREGISTERED has Request=0")
1566+
require.NotNil(t, rev.Details, "revocation UNREGISTERED must carry Details")
1567+
revRegID, _ := wamp.AsID(rev.Details["registration"])
1568+
require.Equal(t, regID1, revRegID, "revocation Details.registration must match evicted regID")
1569+
require.Equal(t, string(wamp.ErrUnregistered), rev.Details["reason"])
1570+
1571+
// Callee2 receives a fresh REGISTERED.
1572+
reg2, ok := (<-callee2.Recv()).(*wamp.Registered)
1573+
require.True(t, ok, "callee2 should receive REGISTERED")
1574+
require.NotEqual(t, regID1, reg2.Registration,
1575+
"force_reregister should produce a new registration ID")
1576+
1577+
// Meta events: on_unregister for the evicted callee, on_delete for
1578+
// the torn-down registration, then on_create + on_register for the new
1579+
// registration installed for callee2.
1580+
checkMetaReg(t, metaClient, sess1.ID) // on_unregister (evicted)
1581+
checkMetaReg(t, metaClient, sess1.ID) // on_delete (last-seen sess id)
1582+
checkMetaReg(t, metaClient, sess2.ID) // on_create
1583+
checkMetaReg(t, metaClient, sess2.ID) // on_register
1584+
1585+
// Dealer state: only callee2 should own the procedure now.
1586+
reg, ok := dealer.procRegMap[testProcedure]
1587+
require.True(t, ok, "procedure registration missing after force_reregister")
1588+
require.Equal(t, reg2.Registration, reg.id)
1589+
require.Equal(t, 1, len(reg.callees))
1590+
require.Same(t, sess2, reg.callees[0])
1591+
1592+
// Callee1 must no longer appear in the dealer's reg-set bookkeeping.
1593+
_, stillHasCallee1 := dealer.calleeRegIDSet[sess1]
1594+
require.False(t, stillHasCallee1, "evicted callee should be removed from calleeRegIDSet")
1595+
1596+
// Calls should now invoke callee2.
1597+
caller := newTestPeer()
1598+
callerSess := wamp.NewSession(caller, 0, nil, nil)
1599+
dealer.Call(callerSess, &wamp.Call{Request: 99, Procedure: testProcedure})
1600+
select {
1601+
case msg := <-callee2.Recv():
1602+
_, ok := msg.(*wamp.Invocation)
1603+
require.True(t, ok, "expected INVOCATION on callee2")
1604+
case <-callee1.Recv():
1605+
require.FailNow(t, "evicted callee1 must not receive INVOCATION")
1606+
case <-time.After(time.Second):
1607+
require.FailNow(t, "timed out waiting for INVOCATION")
1608+
}
1609+
}
1610+
1611+
// TestForceReregisterIgnoredForSharedRegistration confirms force_reregister
1612+
// does not tear down a registration whose invocation policy is shared
1613+
// (anything other than single). The new caller should fall through to the
1614+
// normal policy-conflict / shared-add path.
1615+
func TestForceReregisterIgnoredForSharedRegistration(t *testing.T) {
1616+
dealer, metaClient := newTestDealer(t)
1617+
1618+
calleeRoles := wamp.Dict{
1619+
"roles": wamp.Dict{
1620+
"callee": wamp.Dict{
1621+
"features": wamp.Dict{
1622+
"shared_registration": true,
1623+
},
1624+
},
1625+
},
1626+
}
1627+
1628+
// Callee1 holds a roundrobin shared registration.
1629+
callee1 := newTestPeer()
1630+
sess1 := wamp.NewSession(callee1, 0, nil, calleeRoles)
1631+
dealer.Register(sess1, &wamp.Register{
1632+
Request: 1,
1633+
Procedure: testProcedure,
1634+
Options: wamp.SetOption(nil, wamp.OptInvoke, wamp.InvokeRoundRobin),
1635+
})
1636+
regID1 := (<-callee1.Recv()).(*wamp.Registered).Registration
1637+
checkMetaReg(t, metaClient, sess1.ID)
1638+
checkMetaReg(t, metaClient, sess1.ID)
1639+
1640+
// Callee2 requests force_reregister=true but invoke=single — the
1641+
// existing registration is shared, so eviction must be skipped and
1642+
// the policy-conflict error must fire.
1643+
callee2 := newTestPeer()
1644+
sess2 := wamp.NewSession(callee2, 0, nil, calleeRoles)
1645+
opts := wamp.SetOption(nil, wamp.OptForceReregister, true)
1646+
opts = wamp.SetOption(opts, wamp.OptInvoke, wamp.InvokeSingle)
1647+
dealer.Register(sess2, &wamp.Register{
1648+
Request: 2,
1649+
Procedure: testProcedure,
1650+
Options: opts,
1651+
})
1652+
errMsg, ok := (<-callee2.Recv()).(*wamp.Error)
1653+
require.True(t, ok, "expected ERROR — force_reregister must not evict shared regs")
1654+
require.Equal(t, wamp.ErrProcedureAlreadyExists, errMsg.Error)
1655+
1656+
// Original registration must still exist with callee1 attached.
1657+
reg, ok := dealer.registrations[regID1]
1658+
require.True(t, ok, "original shared registration was unexpectedly removed")
1659+
require.Equal(t, 1, len(reg.callees))
1660+
require.Same(t, sess1, reg.callees[0])
1661+
1662+
// Callee1 should not have received any UNREGISTERED.
1663+
select {
1664+
case msg := <-callee1.Recv():
1665+
require.FailNow(t, "callee1 should not receive any message",
1666+
"got %T", msg)
1667+
default:
1668+
}
1669+
}
1670+
1671+
// TestForceReregisterAdvertisedFeature ensures the dealer announces
1672+
// force_reregister in its role features so clients can discover support.
1673+
func TestForceReregisterAdvertisedFeature(t *testing.T) {
1674+
features, ok := dealerRole["features"].(wamp.Dict)
1675+
require.True(t, ok, "dealer role missing features dict")
1676+
supported, _ := features[wamp.FeatureForceReregister].(bool)
1677+
require.True(t, supported, "dealer must advertise force_reregister=true")
1678+
}

test/spec_force_reregister_test.go

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,64 @@
1+
package test_test
2+
3+
import (
4+
"context"
5+
"sync/atomic"
6+
"testing"
7+
"time"
8+
9+
"github.com/stretchr/testify/require"
10+
11+
"github.com/gammazero/nexus/v3/client"
12+
"github.com/gammazero/nexus/v3/wamp"
13+
)
14+
15+
// TestSpecForceReregister exercises the WAMP advanced-profile
16+
// force_reregister option end-to-end across whichever transport the
17+
// test matrix is currently running on. The second callee must forcibly
18+
// take over the procedure from the first, and subsequent CALLs must be
19+
// routed to the new callee.
20+
func TestSpecForceReregister(t *testing.T) {
21+
checkGoLeaks(t)
22+
23+
const procName = "nexus.test.force_reregister"
24+
25+
callee1 := connectClient(t)
26+
callee2 := connectClient(t)
27+
caller := connectClient(t)
28+
29+
var hits1, hits2 atomic.Int32
30+
handler1 := func(_ context.Context, _ *wamp.Invocation) client.InvokeResult {
31+
hits1.Add(1)
32+
return client.InvokeResult{Args: wamp.List{"callee1"}}
33+
}
34+
handler2 := func(_ context.Context, _ *wamp.Invocation) client.InvokeResult {
35+
hits2.Add(1)
36+
return client.InvokeResult{Args: wamp.List{"callee2"}}
37+
}
38+
39+
require.NoError(t, callee1.Register(procName, handler1, nil),
40+
"callee1 should be able to register the procedure")
41+
42+
// Second REGISTER without force_reregister must fail with
43+
// procedure_already_exists.
44+
err := callee2.Register(procName, handler2, nil)
45+
require.Error(t, err,
46+
"second register without force_reregister must fail")
47+
48+
// Now callee2 force-reregisters; should succeed and evict callee1.
49+
require.NoError(t,
50+
callee2.Register(procName, handler2,
51+
wamp.Dict{wamp.OptForceReregister: true}),
52+
"force_reregister must let callee2 take over")
53+
54+
// A CALL must now be routed to callee2.
55+
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
56+
defer cancel()
57+
res, err := caller.Call(ctx, procName, nil, nil, nil, nil)
58+
require.NoError(t, err, "post-force_reregister CALL should succeed")
59+
require.NotEmpty(t, res.Arguments)
60+
require.Equal(t, "callee2", res.Arguments[0],
61+
"force_reregister should route subsequent CALLs to the new callee")
62+
require.Equal(t, int32(1), hits2.Load(), "callee2 should have served the call")
63+
require.Equal(t, int32(0), hits1.Load(), "callee1 must not have served the call")
64+
}

test/spec_unimplemented_test.go

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -36,10 +36,6 @@ func TestSpecUnimplementedTopicReflection(t *testing.T) {
3636
t.Skip("pending: topic_reflection — wamp.reflection.topic.* not implemented")
3737
}
3838

39-
func TestSpecUnimplementedForceReregister(t *testing.T) {
40-
t.Skip("pending: force_reregister — option not honored by dealer.Register")
41-
}
42-
4339
func TestSpecUnimplementedBatchedWSTransport(t *testing.T) {
4440
t.Skip("pending: batched WebSocket transport (wamp.2.json.batched / msgpack.batched)")
4541
}

wamp/message.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -335,8 +335,15 @@ func (msg *Unregister) MessageType() MessageType { return UNREGISTER }
335335
// the Callee:
336336
//
337337
// [UNREGISTERED, UNREGISTER.Request|id]
338+
//
339+
// The Dealer also sends UNREGISTERED unsolicited when a registration is
340+
// revoked (e.g. force_reregister); the extended form carries a Details dict
341+
// with "registration" and "reason":
342+
//
343+
// [UNREGISTERED, 0, Details|dict]
338344
type Unregistered struct {
339345
Request ID
346+
Details Dict `wamp:"omitempty"`
340347
}
341348

342349
func (msg *Unregistered) MessageType() MessageType { return UNREGISTERED }

wamp/options.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ const (
2323
OptPPTKeyId = "ppt_keyid"
2424
OptSticky = "sticky"
2525
OptForwardTimeout = "forward_timeout"
26+
OptForceReregister = "force_reregister"
2627

2728
// Values for URI matching mode.
2829
MatchExact = "exact"

wamp/roles_reatures.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ const (
1818
FeatureProgCallInvocations = "progressive_call_invocations"
1919
FeatureSessionMetaAPI = "session_meta_api"
2020
FeatureSharedReg = "shared_registration"
21+
FeatureForceReregister = "force_reregister"
2122
FeatureRegMetaAPI = "registration_meta_api"
2223
FeatureTestamentMetaAPI = "testament_meta_api"
2324

wamp/uris.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,11 @@ const (
2222
// is not active.
2323
ErrNoSuchRegistration = URI("wamp.error.no_such_registration")
2424

25+
// Sent by a Dealer in an unsolicited UNREGISTERED message to a Callee
26+
// whose registration has been revoked — e.g. because another peer forcibly
27+
// took it over with force_reregister=true.
28+
ErrUnregistered = URI("wamp.error.unregistered")
29+
2530
// A Broker could not perform an unsubscribe, since the given subscription
2631
// is not active.
2732
ErrNoSuchSubscription = URI("wamp.error.no_such_subscription")

0 commit comments

Comments
 (0)