瀏覽代碼

webapi: broadcast rate limit updates via per-account monitor

Mike 2 天之前
父節點
當前提交
0d131e250f

+ 74 - 44
foodgroup/oservice.go

@@ -21,6 +21,8 @@ type OServiceService struct {
 	logger           *slog.Logger
 	snacRateLimits   wire.SNACRateLimits
 	timeNow          func() time.Time
+	// how often MonitorRateLimits observes changes; a field so tests can shorten it
+	rateLimitMonitorInterval time.Duration
 
 	chatRoomManager       ChatRoomRegistry
 	cookieIssuer          CookieBaker
@@ -48,18 +50,19 @@ func NewOServiceService(
 	feedbagManager FeedbagManager,
 ) *OServiceService {
 	return &OServiceService{
-		cookieIssuer:          cookieIssuer,
-		messageRelayer:        messageRelayer,
-		buddyBroadcaster:      newBuddyNotifier(bartItemManager, relationshipFetcher, messageRelayer, sessionRetriever),
-		cfg:                   cfg,
-		logger:                logger,
-		snacRateLimits:        snacRateLimits,
-		timeNow:               time.Now,
-		chatRoomManager:       chatRoomManager,
-		chatMessageRelayer:    chatMessageRelayer,
-		profileManager:        profileManager,
-		offlineMessageManager: offlineMessageManager,
-		feedbagManager:        feedbagManager,
+		cookieIssuer:             cookieIssuer,
+		messageRelayer:           messageRelayer,
+		buddyBroadcaster:         newBuddyNotifier(bartItemManager, relationshipFetcher, messageRelayer, sessionRetriever),
+		cfg:                      cfg,
+		logger:                   logger,
+		snacRateLimits:           snacRateLimits,
+		timeNow:                  time.Now,
+		rateLimitMonitorInterval: time.Second,
+		chatRoomManager:          chatRoomManager,
+		chatMessageRelayer:       chatMessageRelayer,
+		profileManager:           profileManager,
+		offlineMessageManager:    offlineMessageManager,
+		feedbagManager:           feedbagManager,
 	}
 }
 
@@ -474,40 +477,67 @@ func (s OServiceService) HostOnline(service uint16) wire.SNACMessage {
 	}
 }
 
-// RateLimitUpdates produces update messages reflecting any recent changes in
-// rate limit class params or rate limit states for the current session.
-// Changes are reported relative to the previous invocation for this session.
-// Only newly observed transitions or updated rate parameters will be included.
-func (s OServiceService) RateLimitUpdates(ctx context.Context, instance *state.SessionInstance, now time.Time) []wire.SNACMessage {
-	msgs := make([]wire.SNACMessage, 0, 5)
-	classDelta, stateDelta := instance.Session().ObserveRateChanges(now)
-
-	for _, curRate := range classDelta {
-		s.logger.DebugContext(ctx, "rate limit class changed", "class", curRate.ID)
-		msgs = append(msgs, buildRateLimitUpdate(1, curRate, instance, now))
-	}
-
-	for _, curRate := range stateDelta {
-		s.logger.DebugContext(ctx, "rate limit state changed",
-			"class", curRate.ID,
-			"state", curRate.CurrentStatus)
-		var code uint16
-		switch curRate.CurrentStatus {
-		case wire.RateLimitStatusLimited:
-			code = 3
-		case wire.RateLimitStatusAlert:
-			code = 2
-		case wire.RateLimitStatusClear:
-			code = 4
-		case wire.RateLimitStatusDisconnect:
-			s.logger.DebugContext(ctx, "rate limit status disconnected, no point in returning status update")
-			continue
-		}
+// MonitorRateLimits observes account-wide rate limit changes on a fixed cadence
+// and broadcasts each transition to every instance in the session.
+//
+// Session.ObserveRateChanges is a single-consumer delta — it reports each
+// transition once, then overwrites its baseline — so one monitor per account is
+// the only correct consumer. The per-connection tickers it replaces raced for
+// that one delta and left all but one connection un-notified.
+//
+// Only classes a client on the account has subscribed to are broadcast. Fanout is
+// account-wide: the budget is shared, so a connection that spent nothing cannot
+// send either. That does mean an idle Web API tab raises its rate limit banner
+// when another tab spends the budget.
+//
+// Start it once per account from a Session.RunOnce block, with a server-lifetime
+// context; it runs until the session closes or that context is cancelled.
+func (s OServiceService) MonitorRateLimits(ctx context.Context, session *state.Session) {
+	ticker := time.NewTicker(s.rateLimitMonitorInterval)
+	defer ticker.Stop()
+
+	for {
+		select {
+		case <-ctx.Done(): // server shutdown
+			return
+		case <-session.Closed(): // account signed off; a later sign-on starts a fresh monitor
+			return
+		case <-ticker.C:
+			now := s.timeNow()
+			classDelta, stateDelta := session.ObserveRateChanges(now)
+			if len(classDelta) == 0 && len(stateDelta) == 0 {
+				continue
+			}
+			instances := session.Instances()
 
-		msgs = append(msgs, buildRateLimitUpdate(code, curRate, instance, now))
-	}
+			for _, curRate := range classDelta {
+				s.logger.DebugContext(ctx, "rate limit class changed", "class", curRate.ID)
+				for _, inst := range instances {
+					inst.RelayMessageToInstance(buildRateLimitUpdate(1, curRate, inst, now))
+				}
+			}
 
-	return msgs
+			for _, curRate := range stateDelta {
+				var code uint16
+				switch curRate.CurrentStatus {
+				case wire.RateLimitStatusLimited:
+					code = 3
+				case wire.RateLimitStatusAlert:
+					code = 2
+				case wire.RateLimitStatusClear:
+					code = 4
+				case wire.RateLimitStatusDisconnect:
+					continue // the connection is torn down anyway
+				}
+				s.logger.DebugContext(ctx, "rate limit state changed",
+					"class", curRate.ID,
+					"state", curRate.CurrentStatus)
+				for _, inst := range instances {
+					inst.RelayMessageToInstance(buildRateLimitUpdate(code, curRate, inst, now))
+				}
+			}
+		}
+	}
 }
 
 // buildRateLimitUpdate constructs a SNAC message notifying the client of a rate limit

+ 171 - 331
foodgroup/oservice_test.go

@@ -5,6 +5,7 @@ import (
 	"context"
 	"io"
 	"log/slog"
+	"sync"
 	"testing"
 	"time"
 
@@ -3059,361 +3060,200 @@ func TestOServiceService_SetPrivacyFlags(t *testing.T) {
 	svc.SetPrivacyFlags(context.Background(), body)
 }
 
-func TestOServiceService_RateLimitUpdates(t *testing.T) {
-	rateClasses := [5]wire.RateClass{
-		{
-			ID:              1,
-			WindowSize:      80,
-			ClearLevel:      2500,
-			AlertLevel:      2000,
-			LimitLevel:      1500,
-			DisconnectLevel: 800,
-			MaxLevel:        6000,
-		},
-		{
-			ID:              2,
-			WindowSize:      80,
-			ClearLevel:      3000,
-			AlertLevel:      2000,
-			LimitLevel:      1500,
-			DisconnectLevel: 1000,
-			MaxLevel:        6000,
-		},
-		{
-			ID:              3,
-			WindowSize:      20,
-			ClearLevel:      5100,
-			AlertLevel:      5000,
-			LimitLevel:      4000,
-			DisconnectLevel: 3000,
-			MaxLevel:        6000,
-		},
-		{
-			ID:              4,
-			WindowSize:      20,
-			ClearLevel:      5500,
-			AlertLevel:      5300,
-			LimitLevel:      4200,
-			DisconnectLevel: 3000,
-			MaxLevel:        8000,
-		},
-		{
-			ID:              5,
-			WindowSize:      10,
-			ClearLevel:      5500,
-			AlertLevel:      5300,
-			LimitLevel:      4200,
-			DisconnectLevel: 3000,
-			MaxLevel:        8000,
-		},
+func TestOServiceService_MonitorRateLimits(t *testing.T) {
+	now := time.Now()
+
+	sess := state.NewSession()
+	sess.SetRateClasses(now, wire.DefaultRateLimitClasses())
+	inst1 := sess.AddInstance()
+	inst2 := sess.AddInstance()
+
+	classID := wire.RateLimitClassID(3)
+	sess.SubscribeRateLimits([]wire.RateLimitClassID{classID})
+
+	// A clock shared with the monitor goroutine so the test controls when
+	// ObserveRateChanges sees recovery.
+	var clockMu sync.Mutex
+	clockNow := now
+	setClock := func(v time.Time) {
+		clockMu.Lock()
+		defer clockMu.Unlock()
+		clockNow = v
 	}
 
 	svc := OServiceService{
-		cfg:    config.Config{},
-		logger: slog.Default(),
+		logger:                   slog.New(slog.NewTextHandler(io.Discard, nil)),
+		rateLimitMonitorInterval: time.Millisecond,
+		timeNow: func() time.Time {
+			clockMu.Lock()
+			defer clockMu.Unlock()
+			return clockNow
+		},
 	}
 
-	t.Run("(win aim 1.x) transition state from clear > alert > limited > clear, then change rate limit param", func(t *testing.T) {
-		now := time.Now()
-		instance := newTestInstance("me")
-		instance.Session().SetRateClasses(now, wire.NewRateLimitClasses(rateClasses))
-
-		classId := wire.RateLimitClassID(3)
-		instance.Session().SubscribeRateLimits([]wire.RateLimitClassID{classId})
-
-		// get into an alert state
-		maxTries := 4
-		for i := 1; i <= maxTries; i++ {
-			now = now.Add(time.Millisecond)
-			if s := instance.Session().EvaluateRateLimit(now, classId); s == wire.RateLimitStatusAlert {
-				break
-			}
-			if i == maxTries {
-				t.Fail()
-				return
-			}
-		}
-
-		outputSNACs := svc.RateLimitUpdates(context.Background(), instance, now)
-		expect := wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 2,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 3000,
-					CurrentLevel:    4886,
-					MaxLevel:        6000,
-				},
-			},
-		}
-		assert.Equal(t, expect, outputSNACs[0])
-
-		// get into a rate-limited state
-		maxTries = 4
-		for i := 1; i <= maxTries; i++ {
-			now = now.Add(time.Millisecond)
-			if s := instance.Session().EvaluateRateLimit(now, classId); s == wire.RateLimitStatusLimited {
-				break
+	recv := func(inst *state.SessionInstance) wire.SNAC_0x01_0x0A_OServiceRateParamsChange {
+		t.Helper()
+		select {
+		case msg := <-inst.ReceiveMessage():
+			assert.Equal(t, wire.OServiceRateParamChange, msg.Frame.SubGroup)
+			body, ok := msg.Body.(wire.SNAC_0x01_0x0A_OServiceRateParamsChange)
+			if !ok {
+				t.Fatalf("unexpected body type %T", msg.Body)
 			}
-			if i == maxTries {
-				t.Fail()
-				return
-			}
-		}
-
-		outputSNACs = svc.RateLimitUpdates(context.Background(), instance, now)
-		expect = wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 3,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 3000,
-					CurrentLevel:    3978,
-					MaxLevel:        6000,
-				},
-			},
+			return body
+		case <-time.After(2 * time.Second):
+			t.Fatal("timed out waiting for a relayed rate limit SNAC")
+			return wire.SNAC_0x01_0x0A_OServiceRateParamsChange{}
 		}
-		assert.Equal(t, expect, outputSNACs[0])
-
-		// simulate waiting a minute for the clear threshold
-		now = now.Add(time.Minute)
+	}
 
-		// verify that the clear threshold has been reached
-		outputSNACs = svc.RateLimitUpdates(context.Background(), instance, now)
-		expect = wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 4,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 3000,
-					CurrentLevel:    6000,
-					MaxLevel:        6000,
-				},
-			},
+	// Drive the IM class into the limited state.
+	driveTime := now
+	var status wire.RateLimitStatus
+	for i := 0; status != wire.RateLimitStatusLimited; i++ {
+		if i > 100 {
+			t.Fatal("class never reached the limited state")
 		}
-		assert.Equal(t, expect, outputSNACs[0])
-
-		// verify rate class param changes are detected
-		classesCopy := rateClasses
-		classesCopy[2].DisconnectLevel--
-		instance.Session().SetRateClasses(now, wire.NewRateLimitClasses(classesCopy))
+		driveTime = driveTime.Add(time.Millisecond)
+		status = sess.EvaluateRateLimit(driveTime, classID)
+	}
+	setClock(driveTime)
 
-		outputSNACs = svc.RateLimitUpdates(context.Background(), instance, now)
-		expect = wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 1,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 2999,
-					CurrentLevel:    6000,
-					MaxLevel:        6000,
-				},
-			},
-		}
-		assert.Equal(t, expect, outputSNACs[0])
+	// Stop the monitor at test end by closing the account's session.
+	t.Cleanup(func() {
+		inst1.CloseInstance()
+		inst2.CloseInstance()
 	})
+	go svc.MonitorRateLimits(context.Background(), sess)
+
+	// The transition is observed once for the account and broadcast to every
+	// instance — the multi-connection bug the monitor fixes.
+	for _, inst := range []*state.SessionInstance{inst1, inst2} {
+		body := recv(inst)
+		assert.Equal(t, uint16(3), body.Code) // limited
+		assert.Equal(t, uint16(classID), body.Rate.ID)
+	}
 
-	t.Run("(win aim > 1.x) transition state from clear > alert > limited > clear", func(t *testing.T) {
-		now := time.Now()
-		instance := newTestInstance("me")
-		instance.Session().SetRateClasses(now, wire.NewRateLimitClasses(rateClasses))
+	// After a long idle gap the moving average clears; both instances are told.
+	setClock(driveTime.Add(time.Minute))
+	for _, inst := range []*state.SessionInstance{inst1, inst2} {
+		body := recv(inst)
+		assert.Equal(t, uint16(4), body.Code) // clear
+	}
+}
 
-		var versions [wire.MDir + 1]uint16
-		versions[wire.OService] = 3
-		instance.SetFoodGroupVersions(versions)
+// The monitor is started from a RunOnce block owned by whichever connection
+// signed on first. That connection departing must not take rate limit eventing
+// down with it.
+func TestOServiceService_MonitorRateLimits_outlivesTheInstanceThatStartedIt(t *testing.T) {
+	now := time.Now()
+
+	sess := state.NewSession()
+	sess.SetRateClasses(now, wire.DefaultRateLimitClasses())
+	inst1 := sess.AddInstance()
+
+	classID := wire.RateLimitClassID(3)
+	sess.SubscribeRateLimits([]wire.RateLimitClassID{classID})
+
+	var clockMu sync.Mutex
+	clockNow := now
+	setClock := func(v time.Time) {
+		clockMu.Lock()
+		defer clockMu.Unlock()
+		clockNow = v
+	}
 
-		classId := wire.RateLimitClassID(3)
-		instance.Session().SubscribeRateLimits([]wire.RateLimitClassID{classId})
+	svc := OServiceService{
+		logger:                   slog.New(slog.NewTextHandler(io.Discard, nil)),
+		rateLimitMonitorInterval: time.Millisecond,
+		timeNow: func() time.Time {
+			clockMu.Lock()
+			defer clockMu.Unlock()
+			return clockNow
+		},
+	}
 
-		// get into an alert state
-		maxTries := 4
-		for i := 1; i <= maxTries; i++ {
-			now = now.Add(time.Millisecond)
-			if s := instance.Session().EvaluateRateLimit(now, classId); s == wire.RateLimitStatusAlert {
-				break
-			}
-			if i == maxTries {
-				t.Fail()
-				return
+	// The first instance on the account starts the monitor, as RunOnce does.
+	go svc.MonitorRateLimits(context.Background(), sess)
+
+	// A second connection joins, then the one that started the monitor departs.
+	inst2 := sess.AddInstance()
+	t.Cleanup(inst2.CloseInstance)
+	inst1.CloseInstance()
+	require.False(t, sess.IsClosed(), "the account is still online via inst2")
+
+	recv := func() wire.SNAC_0x01_0x0A_OServiceRateParamsChange {
+		t.Helper()
+		select {
+		case msg := <-inst2.ReceiveMessage():
+			assert.Equal(t, wire.OServiceRateParamChange, msg.Frame.SubGroup)
+			body, ok := msg.Body.(wire.SNAC_0x01_0x0A_OServiceRateParamsChange)
+			if !ok {
+				t.Fatalf("unexpected body type %T", msg.Body)
 			}
+			return body
+		case <-time.After(2 * time.Second):
+			t.Fatal("surviving instance got no rate limit update after the starting instance left")
+			return wire.SNAC_0x01_0x0A_OServiceRateParamsChange{}
 		}
+	}
 
-		outputSNACs := svc.RateLimitUpdates(context.Background(), instance, now)
-		expect := wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 2,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 3000,
-					CurrentLevel:    4886,
-					MaxLevel:        6000,
-					V2Params: &struct {
-						LastTime      uint32
-						DroppingSNACs uint8
-					}{
-						DroppingSNACs: 0,
-					},
-				},
-			},
-		}
-		assert.Equal(t, expect, outputSNACs[0])
-
-		// get into a rate-limited state
-		maxTries = 4
-		for i := 1; i <= maxTries; i++ {
-			now = now.Add(time.Millisecond)
-			if s := instance.Session().EvaluateRateLimit(now, classId); s == wire.RateLimitStatusLimited {
-				break
-			}
-			if i == maxTries {
-				t.Fail()
-				return
-			}
+	driveTime := now
+	var status wire.RateLimitStatus
+	for i := 0; status != wire.RateLimitStatusLimited; i++ {
+		if i > 100 {
+			t.Fatal("class never reached the limited state")
 		}
+		driveTime = driveTime.Add(time.Millisecond)
+		status = sess.EvaluateRateLimit(driveTime, classID)
+	}
+	setClock(driveTime)
 
-		outputSNACs = svc.RateLimitUpdates(context.Background(), instance, now)
-		expect = wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 3,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 3000,
-					CurrentLevel:    3978,
-					MaxLevel:        6000,
-					V2Params: &struct {
-						LastTime      uint32
-						DroppingSNACs uint8
-					}{
-						DroppingSNACs: 1,
-					},
-				},
-			},
-		}
-		assert.Equal(t, expect, outputSNACs[0])
+	body := recv()
+	assert.Equal(t, uint16(3), body.Code) // limited
+	assert.Equal(t, uint16(classID), body.Rate.ID)
 
-		// simulate waiting a minute for the clear threshold
-		now = now.Add(time.Minute)
+	// Recovery reaches the surviving instance too.
+	setClock(driveTime.Add(time.Minute))
+	assert.Equal(t, uint16(4), recv().Code) // clear
+}
 
-		// verify that the clear threshold has been reached
-		outputSNACs = svc.RateLimitUpdates(context.Background(), instance, now)
-		expect = wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 4,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 3000,
-					CurrentLevel:    6000,
-					MaxLevel:        6000,
-					V2Params: &struct {
-						LastTime      uint32
-						DroppingSNACs uint8
-					}{
-						DroppingSNACs: 0,
-						LastTime:      60,
-					},
-				},
-			},
-		}
-		assert.Equal(t, expect, outputSNACs[0])
+// The monitor's lifetime tracks the account, not the connection that started it:
+// the WebAPI server runs it with a server-lifetime context, so closing the
+// session is the only thing that stops it once the account has signed off.
+func TestOServiceService_MonitorRateLimits_exitsWhenAccountEmpty(t *testing.T) {
+	sess := state.NewSession()
+	sess.SetRateClasses(time.Now(), wire.DefaultRateLimitClasses())
+	inst := sess.AddInstance()
 
-		// verify rate class param changes are detected
-		classesCopy := rateClasses
-		classesCopy[2].DisconnectLevel--
-		instance.Session().SetRateClasses(now, wire.NewRateLimitClasses(classesCopy))
+	svc := OServiceService{
+		logger:                   slog.New(slog.NewTextHandler(io.Discard, nil)),
+		rateLimitMonitorInterval: time.Millisecond,
+		timeNow:                  time.Now,
+	}
 
-		outputSNACs = svc.RateLimitUpdates(context.Background(), instance, now)
-		expect = wire.SNACMessage{
-			Frame: wire.SNACFrame{
-				FoodGroup: wire.OService,
-				SubGroup:  wire.OServiceRateParamChange,
-				RequestID: wire.ReqIDFromServer,
-			},
-			Body: wire.SNAC_0x01_0x0A_OServiceRateParamsChange{
-				Code: 1,
-				Rate: wire.RateParamsSNAC{
-					ID:              3,
-					WindowSize:      20,
-					ClearLevel:      5100,
-					AlertLevel:      5000,
-					LimitLevel:      4000,
-					DisconnectLevel: 2999,
-					CurrentLevel:    6000,
-					MaxLevel:        6000,
-					V2Params: &struct {
-						LastTime      uint32
-						DroppingSNACs uint8
-					}{
-						DroppingSNACs: 0,
-						LastTime:      0,
-					},
-				},
-			},
-		}
-		assert.Equal(t, expect, outputSNACs[0])
-	})
+	done := make(chan struct{})
+	go func() {
+		// A context that is never cancelled — like the WebAPI server's
+		// shutdownCtx before shutdown — so only Session.Closed() can stop it.
+		svc.MonitorRateLimits(context.Background(), sess)
+		close(done)
+	}()
+
+	// While an instance is present, the monitor keeps running.
+	select {
+	case <-done:
+		t.Fatal("monitor exited while an instance was still present")
+	case <-time.After(20 * time.Millisecond):
+	}
+
+	// Closing the account's last instance closes the session; the monitor exits.
+	inst.CloseInstance()
+	select {
+	case <-done:
+	case <-time.After(2 * time.Second):
+		t.Fatal("monitor did not exit after the account's last instance left")
+	}
 }
 
 func TestOServiceService_RateParamsSubAdd(t *testing.T) {

+ 17 - 38
server/oscar/mock_rate_limit_updater_test.go

@@ -6,10 +6,8 @@ package oscar
 
 import (
 	"context"
-	"time"
 
 	"github.com/mk6i/open-oscar-server/state"
-	"github.com/mk6i/open-oscar-server/wire"
 	mock "github.com/stretchr/testify/mock"
 )
 
@@ -40,67 +38,48 @@ func (_m *mockRateLimitUpdater) EXPECT() *mockRateLimitUpdater_Expecter {
 	return &mockRateLimitUpdater_Expecter{mock: &_m.Mock}
 }
 
-// RateLimitUpdates provides a mock function for the type mockRateLimitUpdater
-func (_mock *mockRateLimitUpdater) RateLimitUpdates(ctx context.Context, instance *state.SessionInstance, now time.Time) []wire.SNACMessage {
-	ret := _mock.Called(ctx, instance, now)
-
-	if len(ret) == 0 {
-		panic("no return value specified for RateLimitUpdates")
-	}
-
-	var r0 []wire.SNACMessage
-	if returnFunc, ok := ret.Get(0).(func(context.Context, *state.SessionInstance, time.Time) []wire.SNACMessage); ok {
-		r0 = returnFunc(ctx, instance, now)
-	} else {
-		if ret.Get(0) != nil {
-			r0 = ret.Get(0).([]wire.SNACMessage)
-		}
-	}
-	return r0
+// MonitorRateLimits provides a mock function for the type mockRateLimitUpdater
+func (_mock *mockRateLimitUpdater) MonitorRateLimits(ctx context.Context, session *state.Session) {
+	_mock.Called(ctx, session)
+	return
 }
 
-// mockRateLimitUpdater_RateLimitUpdates_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'RateLimitUpdates'
-type mockRateLimitUpdater_RateLimitUpdates_Call struct {
+// mockRateLimitUpdater_MonitorRateLimits_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'MonitorRateLimits'
+type mockRateLimitUpdater_MonitorRateLimits_Call struct {
 	*mock.Call
 }
 
-// RateLimitUpdates is a helper method to define mock.On call
+// MonitorRateLimits is a helper method to define mock.On call
 //   - ctx context.Context
-//   - instance *state.SessionInstance
-//   - now time.Time
-func (_e *mockRateLimitUpdater_Expecter) RateLimitUpdates(ctx interface{}, instance interface{}, now interface{}) *mockRateLimitUpdater_RateLimitUpdates_Call {
-	return &mockRateLimitUpdater_RateLimitUpdates_Call{Call: _e.mock.On("RateLimitUpdates", ctx, instance, now)}
+//   - session *state.Session
+func (_e *mockRateLimitUpdater_Expecter) MonitorRateLimits(ctx interface{}, session interface{}) *mockRateLimitUpdater_MonitorRateLimits_Call {
+	return &mockRateLimitUpdater_MonitorRateLimits_Call{Call: _e.mock.On("MonitorRateLimits", ctx, session)}
 }
 
-func (_c *mockRateLimitUpdater_RateLimitUpdates_Call) Run(run func(ctx context.Context, instance *state.SessionInstance, now time.Time)) *mockRateLimitUpdater_RateLimitUpdates_Call {
+func (_c *mockRateLimitUpdater_MonitorRateLimits_Call) Run(run func(ctx context.Context, session *state.Session)) *mockRateLimitUpdater_MonitorRateLimits_Call {
 	_c.Call.Run(func(args mock.Arguments) {
 		var arg0 context.Context
 		if args[0] != nil {
 			arg0 = args[0].(context.Context)
 		}
-		var arg1 *state.SessionInstance
+		var arg1 *state.Session
 		if args[1] != nil {
-			arg1 = args[1].(*state.SessionInstance)
-		}
-		var arg2 time.Time
-		if args[2] != nil {
-			arg2 = args[2].(time.Time)
+			arg1 = args[1].(*state.Session)
 		}
 		run(
 			arg0,
 			arg1,
-			arg2,
 		)
 	})
 	return _c
 }
 
-func (_c *mockRateLimitUpdater_RateLimitUpdates_Call) Return(sNACMessages []wire.SNACMessage) *mockRateLimitUpdater_RateLimitUpdates_Call {
-	_c.Call.Return(sNACMessages)
+func (_c *mockRateLimitUpdater_MonitorRateLimits_Call) Return() *mockRateLimitUpdater_MonitorRateLimits_Call {
+	_c.Call.Return()
 	return _c
 }
 
-func (_c *mockRateLimitUpdater_RateLimitUpdates_Call) RunAndReturn(run func(ctx context.Context, instance *state.SessionInstance, now time.Time) []wire.SNACMessage) *mockRateLimitUpdater_RateLimitUpdates_Call {
-	_c.Call.Return(run)
+func (_c *mockRateLimitUpdater_MonitorRateLimits_Call) RunAndReturn(run func(ctx context.Context, session *state.Session)) *mockRateLimitUpdater_MonitorRateLimits_Call {
+	_c.Run(run)
 	return _c
 }

+ 12 - 8
server/oscar/server.go

@@ -300,6 +300,8 @@ func (s oscarServer) connectToOSCARService(
 			}
 			// periodically decay warning level
 			go s.lowerWarnLevel(ctx, instance)
+			// broadcast rate limit transitions to every instance on the account
+			go s.rateLimitUpdater.MonitorRateLimits(ctx, instance.Session())
 			return nil
 		}); err != nil {
 			return err
@@ -351,6 +353,16 @@ func (s oscarServer) connectToOSCARService(
 			instance.CloseInstance()
 		}()
 
+		// A chat session is a Session of its own, with its own rate limit states
+		// and subscriptions, so it needs its own monitor — the BOS session's
+		// cannot see these states.
+		if err := instance.Session().RunOnce(func() error {
+			go s.rateLimitUpdater.MonitorRateLimits(ctx, instance.Session())
+			return nil
+		}); err != nil {
+			return err
+		}
+
 		go s.receiveSessMessages(ctx, instance, flapc)
 	default:
 		instance, err = s.authService.RetrieveBOSSession(ctx, cookie)
@@ -641,14 +653,6 @@ func (s oscarServer) dispatchIncomingMessages(
 			default:
 				return fmt.Errorf("got unknown FLAP frame type. flap: %v", flap)
 			}
-		case <-time.After(1 * time.Second):
-			updates := s.rateLimitUpdater.RateLimitUpdates(ctx, instance, time.Now())
-			for _, update := range updates {
-				if err := flapc.SendSNAC(update.Frame, update.Body); err != nil {
-					middleware.LogRequestError(ctx, s.logger, update.Frame, err)
-					return err
-				}
-			}
 		case <-instance.Closed():
 			// add logoff reason to clients that support multi-conn
 			if instance.MultiConnFlag() == wire.MultiConnFlagsOldClient {

+ 58 - 14
server/oscar/server_test.go

@@ -21,6 +21,27 @@ import (
 	"github.com/mk6i/open-oscar-server/wire"
 )
 
+// noopRateLimitUpdater satisfies RateLimitUpdater without running the monitor.
+type noopRateLimitUpdater struct{}
+
+func (noopRateLimitUpdater) MonitorRateLimits(context.Context, *state.Session) {}
+
+// recordingRateLimitUpdater reports which sessions the monitor was started for.
+type recordingRateLimitUpdater struct {
+	monitored chan *state.Session
+}
+
+func newRecordingRateLimitUpdater() recordingRateLimitUpdater {
+	return recordingRateLimitUpdater{monitored: make(chan *state.Session, 1)}
+}
+
+func (r recordingRateLimitUpdater) MonitorRateLimits(_ context.Context, session *state.Session) {
+	select {
+	case r.monitored <- session:
+	default:
+	}
+}
+
 func TestServer_ListenAndServeAndShutdown(t *testing.T) {
 	var mu sync.Mutex
 	var received []string
@@ -234,9 +255,10 @@ func TestOscarServer_RouteConnection_Auth_BUCP(t *testing.T) {
 		}, nil)
 
 	rt := oscarServer{
-		authService:   authService,
-		logger:        slog.Default(),
-		ipRateLimiter: NewIPRateLimiter(rate.Every(1*time.Minute), 10, 1*time.Minute),
+		rateLimitUpdater: noopRateLimitUpdater{},
+		authService:      authService,
+		logger:           slog.Default(),
+		ipRateLimiter:    NewIPRateLimiter(rate.Every(1*time.Minute), 10, 1*time.Minute),
 	}
 	assert.NoError(t, rt.routeConnection(context.Background(), clientFake, config.Listener{BOSAdvertisedHostPlain: "localhost:5190"}))
 
@@ -309,9 +331,10 @@ func TestOscarServer_RouteConnection_Auth_FLAP(t *testing.T) {
 		}, nil)
 
 	rt := oscarServer{
-		authService:   authService,
-		logger:        slog.Default(),
-		ipRateLimiter: NewIPRateLimiter(rate.Every(1*time.Minute), 10, 1*time.Minute),
+		rateLimitUpdater: noopRateLimitUpdater{},
+		authService:      authService,
+		logger:           slog.Default(),
+		ipRateLimiter:    NewIPRateLimiter(rate.Every(1*time.Minute), 10, 1*time.Minute),
 	}
 	assert.NoError(t, rt.routeConnection(context.Background(), clientFake, config.Listener{BOSAdvertisedHostPlain: "localhost:5190"}))
 
@@ -425,6 +448,7 @@ func TestOscarServer_RouteConnection_BOS(t *testing.T) {
 	}
 
 	rt := oscarServer{
+		rateLimitUpdater:   noopRateLimitUpdater{},
 		authService:        authService,
 		snacHandler:        handler,
 		logger:             slog.Default(),
@@ -538,6 +562,7 @@ func TestOscarServer_RouteConnection_BOS_MultiSessionSignoff(t *testing.T) {
 	}
 
 	rt := oscarServer{
+		rateLimitUpdater:   noopRateLimitUpdater{},
 		authService:        authService,
 		snacHandler:        handler,
 		logger:             slog.Default(),
@@ -611,8 +636,9 @@ func TestOscarServer_RouteConnection_BOS_MaxConcurrentSessionsReached(t *testing
 		Return(state.ServerCookie{Service: wire.BOS}, nil)
 
 	rt := oscarServer{
-		authService: authService,
-		logger:      slog.Default(),
+		rateLimitUpdater: noopRateLimitUpdater{},
+		authService:      authService,
+		logger:           slog.Default(),
 	}
 	assert.NoError(t, rt.routeConnection(context.Background(), clientFake, config.Listener{}))
 
@@ -622,6 +648,10 @@ func TestOscarServer_RouteConnection_BOS_MaxConcurrentSessionsReached(t *testing
 func TestOscarServer_RouteConnection_Chat(t *testing.T) {
 	instance := state.NewSession().AddInstance()
 
+	// Without a monitor of its own, a chat client that subscribes over this
+	// connection is never told its rate limit status changed.
+	rateLimitUpdater := newRecordingRateLimitUpdater()
+
 	clientConn, serverConn := net.Pipe()
 	addr, err := net.ResolveTCPAddr("tcp", "127.0.0.1:8080")
 	assert.NoError(t, err)
@@ -713,6 +743,7 @@ func TestOscarServer_RouteConnection_Chat(t *testing.T) {
 	}
 
 	rt := oscarServer{
+		rateLimitUpdater:   rateLimitUpdater,
 		authService:        authService,
 		snacHandler:        handler,
 		logger:             slog.Default(),
@@ -724,6 +755,13 @@ func TestOscarServer_RouteConnection_Chat(t *testing.T) {
 	assert.NoError(t, rt.routeConnection(context.Background(), clientFake, config.Listener{}))
 
 	wg.Wait()
+
+	select {
+	case monitored := <-rateLimitUpdater.monitored:
+		assert.Same(t, instance.Session(), monitored, "the monitor must run against the chat session")
+	case <-time.After(2 * time.Second):
+		t.Fatal("no rate limit monitor was started for the chat session")
+	}
 }
 
 func TestOscarServer_RouteConnection_Admin(t *testing.T) {
@@ -808,6 +846,7 @@ func TestOscarServer_RouteConnection_Admin(t *testing.T) {
 	}
 
 	rt := oscarServer{
+		rateLimitUpdater:   noopRateLimitUpdater{},
 		authService:        authService,
 		snacHandler:        handler,
 		logger:             slog.Default(),
@@ -832,7 +871,8 @@ func Test_oscarServer_dispatchIncomingMessages_shutdownSignoff(t *testing.T) {
 	go func() {
 		defer wg.Done()
 		srv := oscarServer{
-			logger: slog.Default(),
+			rateLimitUpdater: noopRateLimitUpdater{},
+			logger:           slog.Default(),
 		}
 		instance := state.NewSession().AddInstance()
 		instance.SetMultiConnFlag(wire.MultiConnFlagsRecentClient)
@@ -863,7 +903,8 @@ func Test_oscarServer_dispatchIncomingMessages_disconnect_old_client(t *testing.
 	go func() {
 		defer wg.Done()
 		srv := oscarServer{
-			logger: slog.Default(),
+			rateLimitUpdater: noopRateLimitUpdater{},
+			logger:           slog.Default(),
 		}
 		flapc := wire.NewFlapClient(0, serverConn, serverConn)
 		err := srv.dispatchIncomingMessages(ctx, wire.BOS, instance, flapc, serverConn, config.Listener{})
@@ -892,7 +933,8 @@ func Test_oscarServer_dispatchIncomingMessages_disconnect_new_client(t *testing.
 	go func() {
 		defer wg.Done()
 		srv := oscarServer{
-			logger: slog.Default(),
+			rateLimitUpdater: noopRateLimitUpdater{},
+			logger:           slog.Default(),
 		}
 		flapc := wire.NewFlapClient(0, serverConn, serverConn)
 		err := srv.dispatchIncomingMessages(ctx, wire.BOS, instance, flapc, serverConn, config.Listener{})
@@ -956,6 +998,7 @@ func Test_oscarServer_receiveSessMessages_BOS_integration(t *testing.T) {
 	chatSessionManager.EXPECT().RemoveUserFromAllChats(mock.Anything)
 
 	server := oscarServer{
+		rateLimitUpdater:   noopRateLimitUpdater{},
 		authService:        authService,
 		buddyListRegistry:  buddyListRegistry,
 		chatSessionManager: chatSessionManager,
@@ -1097,9 +1140,10 @@ func Test_oscarServer_receiveSessMessages_Chat_integration(t *testing.T) {
 		})
 
 	server := oscarServer{
-		authService:    authService,
-		onlineNotifier: onlineNotifier,
-		logger:         slog.New(slog.NewTextHandler(io.Discard, nil)),
+		rateLimitUpdater: noopRateLimitUpdater{},
+		authService:      authService,
+		onlineNotifier:   onlineNotifier,
+		logger:           slog.New(slog.NewTextHandler(io.Discard, nil)),
 	}
 
 	// Fake client connection with address

+ 3 - 3
server/oscar/types.go

@@ -2,7 +2,6 @@ package oscar
 
 import (
 	"context"
-	"time"
 
 	"github.com/google/uuid"
 
@@ -40,9 +39,10 @@ type ChatSessionManager interface {
 	RemoveUserFromAllChats(user state.IdentScreenName)
 }
 
-// RateLimitUpdater provides rate limit updates for subscribed rate limit classes.
+// RateLimitUpdater runs the per-account rate limit monitor that broadcasts rate
+// limit transitions to every instance in a session.
 type RateLimitUpdater interface {
-	RateLimitUpdates(ctx context.Context, instance *state.SessionInstance, now time.Time) []wire.SNACMessage
+	MonitorRateLimits(ctx context.Context, session *state.Session)
 }
 
 type AuthService interface {

+ 2 - 0
server/toc/cmd_client.go

@@ -2390,6 +2390,8 @@ func (s OSCARProxy) Signon(ctx context.Context, args []byte, recalcWarning func(
 		}
 		// periodically decay warning level
 		go lowerWarnLevel(ctx, instance)
+		// broadcast rate limit transitions to every instance on the account
+		go s.OServiceService.MonitorRateLimits(ctx, instance.Session())
 		return nil
 	}); err != nil {
 		return nil, s.runtimeErr(ctx, fmt.Errorf("Session.RunOnce: %w", err))

+ 19 - 0
server/toc/cmd_client_test.go

@@ -18,6 +18,24 @@ import (
 	"github.com/mk6i/open-oscar-server/wire"
 )
 
+// nopOServiceService satisfies OServiceService without running the rate limit
+// monitor, for signon tests that only need Signon's RunOnce to not panic.
+type nopOServiceService struct{}
+
+func (nopOServiceService) ClientOnline(context.Context, uint16, wire.SNAC_0x01_0x02_OServiceClientOnline, *state.SessionInstance) error {
+	return nil
+}
+
+func (nopOServiceService) IdleNotification(context.Context, *state.SessionInstance, wire.SNAC_0x01_0x11_OServiceIdleNotification) error {
+	return nil
+}
+
+func (nopOServiceService) MonitorRateLimits(context.Context, *state.Session) {}
+
+func (nopOServiceService) ServiceRequest(context.Context, uint16, *state.SessionInstance, wire.SNACFrame, wire.SNAC_0x01_0x04_OServiceServiceRequest, config.Listener) (wire.SNACMessage, error) {
+	return wire.SNACMessage{}, nil
+}
+
 func TestOSCARProxy_RecvClientCmd_AddBuddy(t *testing.T) {
 	cases := []struct {
 		// name is the unit test name
@@ -6983,6 +7001,7 @@ func TestOSCARProxy_Signon(t *testing.T) {
 				TOCConfigStore:     tocCfg,
 				FeedbagService:     fbSvc,
 				FeedbagManager:     fbMgr,
+				OServiceService:    nopOServiceService{},
 			}
 			sess, msg := svc.Signon(ctx, tc.givenCmd,
 				func(ctx context.Context, instance *state.SessionInstance) error { return nil },

+ 46 - 0
server/toc/mock_oservice_service_test.go

@@ -172,6 +172,52 @@ func (_c *mockOServiceService_IdleNotification_Call) RunAndReturn(run func(ctx c
 	return _c
 }
 
+// MonitorRateLimits provides a mock function for the type mockOServiceService
+func (_mock *mockOServiceService) MonitorRateLimits(ctx context.Context, session *state.Session) {
+	_mock.Called(ctx, session)
+	return
+}
+
+// mockOServiceService_MonitorRateLimits_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'MonitorRateLimits'
+type mockOServiceService_MonitorRateLimits_Call struct {
+	*mock.Call
+}
+
+// MonitorRateLimits is a helper method to define mock.On call
+//   - ctx context.Context
+//   - session *state.Session
+func (_e *mockOServiceService_Expecter) MonitorRateLimits(ctx interface{}, session interface{}) *mockOServiceService_MonitorRateLimits_Call {
+	return &mockOServiceService_MonitorRateLimits_Call{Call: _e.mock.On("MonitorRateLimits", ctx, session)}
+}
+
+func (_c *mockOServiceService_MonitorRateLimits_Call) Run(run func(ctx context.Context, session *state.Session)) *mockOServiceService_MonitorRateLimits_Call {
+	_c.Call.Run(func(args mock.Arguments) {
+		var arg0 context.Context
+		if args[0] != nil {
+			arg0 = args[0].(context.Context)
+		}
+		var arg1 *state.Session
+		if args[1] != nil {
+			arg1 = args[1].(*state.Session)
+		}
+		run(
+			arg0,
+			arg1,
+		)
+	})
+	return _c
+}
+
+func (_c *mockOServiceService_MonitorRateLimits_Call) Return() *mockOServiceService_MonitorRateLimits_Call {
+	_c.Call.Return()
+	return _c
+}
+
+func (_c *mockOServiceService_MonitorRateLimits_Call) RunAndReturn(run func(ctx context.Context, session *state.Session)) *mockOServiceService_MonitorRateLimits_Call {
+	_c.Run(run)
+	return _c
+}
+
 // ServiceRequest provides a mock function for the type mockOServiceService
 func (_mock *mockOServiceService) ServiceRequest(ctx context.Context, service uint16, instance *state.SessionInstance, inFrame wire.SNACFrame, inBody wire.SNAC_0x01_0x04_OServiceServiceRequest, listener config.Listener) (wire.SNACMessage, error) {
 	ret := _mock.Called(ctx, service, instance, inFrame, inBody, listener)

+ 1 - 0
server/toc/types.go

@@ -41,6 +41,7 @@ type ICBMService interface {
 type OServiceService interface {
 	ClientOnline(ctx context.Context, service uint16, inBody wire.SNAC_0x01_0x02_OServiceClientOnline, instance *state.SessionInstance) error
 	IdleNotification(ctx context.Context, instance *state.SessionInstance, inBody wire.SNAC_0x01_0x11_OServiceIdleNotification) error
+	MonitorRateLimits(ctx context.Context, session *state.Session)
 	ServiceRequest(ctx context.Context, service uint16, instance *state.SessionInstance, inFrame wire.SNACFrame, inBody wire.SNAC_0x01_0x04_OServiceServiceRequest, listener config.Listener) (wire.SNACMessage, error)
 }
 

+ 1 - 0
server/webapi/handlers/auth.go

@@ -26,6 +26,7 @@ type AuthHandler struct {
 
 type OServiceService interface {
 	ClientOnline(ctx context.Context, service uint16, inBody wire.SNAC_0x01_0x02_OServiceClientOnline, instance *state.SessionInstance) error
+	RateParamsSubAdd(ctx context.Context, instance *state.SessionInstance, inBody wire.SNAC_0x01_0x08_OServiceRateParamsSubAdd)
 }
 
 // ClientLoginRequest represents the request body for clientLogin.

+ 37 - 77
server/webapi/handlers/ratelimit.go

@@ -38,27 +38,20 @@ type SessionHandlerFunc = func(http.ResponseWriter, *http.Request, *state.WebAPI
 // It lives alongside the handlers (rather than in the middleware package) so that
 // its rejection can be encoded through the same SendResponse path the handlers
 // use, honoring the request's JSON/JSONP/XML/AMF format.
+//
+// It only enforces the limit (the 430 rejection). Telling the client its status
+// changed is the job of OServiceService.MonitorRateLimits.
 type RateLimitMiddleware struct {
 	snacRateLimits wire.SNACRateLimits
-	// alertRateClassID is the only rate class that pushes a rateLimit event; see
-	// syncRateAlert. Zero when the SNAC table has no mapping for sending an IM,
-	// which disables the push rather than indexing a rate class that isn't there.
-	alertRateClassID wire.RateLimitClassID
-	logger           *slog.Logger
+	logger         *slog.Logger
 }
 
 // NewRateLimitMiddleware creates a RateLimitMiddleware. snacRateLimits is the
 // same SNAC-to-rate-class mapping the OSCAR and TOC servers use.
 func NewRateLimitMiddleware(snacRateLimits wire.SNACRateLimits, logger *slog.Logger) *RateLimitMiddleware {
-	alertRateClassID, ok := snacRateLimits.RateClassLookup(wire.ICBM, wire.ICBMChannelMsgToHost)
-	if !ok {
-		logger.Error("no rate class maps to sending an IM, rate limit events are disabled")
-	}
-
 	return &RateLimitMiddleware{
-		snacRateLimits:   snacRateLimits,
-		alertRateClassID: alertRateClassID,
-		logger:           logger,
+		snacRateLimits: snacRateLimits,
+		logger:         logger,
 	}
 }
 
@@ -85,7 +78,6 @@ func (l *RateLimitMiddleware) OSCAR(foodGroup uint16, subGroup uint16) func(Sess
 
 			sess := session.OSCARSession.Session()
 			status := sess.EvaluateRateLimit(time.Now(), rateClassID)
-			l.syncRateAlert(session, rateClassID, status)
 
 			// Disconnect is rejected alongside Limited: EvaluateRateLimit has
 			// already closed the account's OSCAR session by the time it returns,
@@ -114,69 +106,6 @@ func (l *RateLimitMiddleware) OSCAR(foodGroup uint16, subGroup uint16) func(Sess
 	}
 }
 
-// RecoverOnPoll re-evaluates rate limit recovery before serving a long poll and
-// pushes the "clear" that dismisses the client's alert once it recovers.
-//
-// OSCAR only recomputes rate limit status on a charged request (see OSCAR), so a
-// session that trips the limit and then goes idle would never learn it recovered,
-// and the client's alert is sticky until a "clear" arrives. fetchEvents is polled
-// continuously, so recovering here surfaces the clear without the user sending
-// anything; the push happens before next() so it rides out on this poll.
-//
-// Only a clear is pushed. A poll is not a request the user made, so it must not
-// raise an alert: the budget is account-wide, and syncing an idle tab up to a
-// limit another tab tripped would tell a user doing nothing to stop chatting. The
-// limit reaches a client on the request it actually rejects, in OSCAR.
-func (l *RateLimitMiddleware) RecoverOnPoll(next SessionHandlerFunc) SessionHandlerFunc {
-	return func(w http.ResponseWriter, r *http.Request, session *state.WebAPISession) {
-		if l.alertRateClassID == 0 {
-			next(w, r, session)
-			return
-		}
-
-		sess := session.OSCARSession.Session()
-		sess.RecoverRateLimits(time.Now())
-		if status := sess.RateLimitStates()[l.alertRateClassID-1].CurrentStatus; status == wire.RateLimitStatusClear {
-			l.syncRateAlert(session, l.alertRateClassID, status)
-		}
-
-		next(w, r, session)
-	}
-}
-
-// syncRateAlert notifies the client that its rate limit status changed so it can
-// show (or dismiss) the rate limit alert. The disconnect push is best-effort:
-// EvaluateRateLimit closes the session before the queue drains.
-//
-// Only alertRateClassID pushes. The client throws away the class id the event
-// carries and feeds the status into a single alert rendered inside a
-// conversation window ("You have been rate limited. Wait for a few moments until
-// you can chat again."), so the one class it can truthfully describe is the one
-// sending an IM spends. Pushing for the others would pop a chat alert for a
-// buddy list edit while IMs still work, and let a "clear" on an unrelated class
-// dismiss an alert the IM class raised. They stay enforced, silently.
-func (l *RateLimitMiddleware) syncRateAlert(session *state.WebAPISession, classID wire.RateLimitClassID, status wire.RateLimitStatus) {
-	if classID != l.alertRateClassID {
-		return
-	}
-	name := rateLimitStatusName(status)
-	if name == "" {
-		return
-	}
-	if !session.ObserveRateAlert(status) {
-		return
-	}
-
-	session.EventQueue.Push(types.EventTypeRateLimit, types.RateLimitEvent{
-		Classes: []types.RateLimitClass{
-			{
-				ID:     int(classID),
-				Status: name,
-			},
-		},
-	})
-}
-
 // retryAfterFor returns how long the client must wait for its next request on
 // this class to clear the limit.
 //
@@ -201,6 +130,37 @@ func retryAfterFor(rcs state.RateClassState) time.Duration {
 	return max(time.Duration((neededMs+999)/1000)*time.Second, minRetryAfter)
 }
 
+// seedRateLimitAlert raises the client's rate limit alert when a session starts
+// on an account that is already rate limited.
+//
+// The monitor broadcasts transitions, not current state, so a session signing on
+// mid-limit missed the one that raised the alert — and the client's alert is
+// sticky, so the eventual "clear" would arrive with nothing to dismiss. An OSCAR
+// client learns the current state from the rate params it gets at handshake; this
+// is the Web API's equivalent.
+//
+// Only the limited state is seeded: alert is a warning the user cannot act on,
+// and seeding clear would render nothing.
+func seedRateLimitAlert(session *state.WebAPISession, classID wire.RateLimitClassID) {
+	if classID == 0 {
+		return
+	}
+
+	status := session.OSCARSession.Session().RateLimitStates()[classID-1].CurrentStatus
+	if status != wire.RateLimitStatusLimited {
+		return
+	}
+
+	session.EventQueue.Push(types.EventTypeRateLimit, types.RateLimitEvent{
+		Classes: []types.RateLimitClass{
+			{
+				ID:     int(classID),
+				Status: rateLimitStatusName(status),
+			},
+		},
+	})
+}
+
 // rateLimitStatusName maps an OSCAR rate limit status onto the status string the
 // web client switches on. It returns "" for a status the client does not know.
 func rateLimitStatusName(status wire.RateLimitStatus) string {

+ 37 - 104
server/webapi/handlers/ratelimit_test.go

@@ -132,8 +132,6 @@ func TestRateLimitMiddleware_OSCAR(t *testing.T) {
 		requests int
 		// wantCalls is how many of those requests reach the handler
 		wantCalls int
-		// wantEvents are the rateLimit event statuses pushed to the session
-		wantEvents []string
 		// wantRetryAfter is the Retry-After header on the final rejection, "" for
 		// a rejection that carries none
 		wantRetryAfter string
@@ -151,9 +149,6 @@ func TestRateLimitMiddleware_OSCAR(t *testing.T) {
 			subGroup:  wire.ICBMChannelMsgToHost,
 			requests:  5,
 			wantCalls: 1,
-			// clear -> limit on the second request; the rest stay limited and
-			// so push nothing further.
-			wantEvents: []string{"limit"},
 			// retryAfterFor rounds the sub-second wait these tight classes need
 			// up to the minRetryAfter floor.
 			wantRetryAfter: "1",
@@ -165,22 +160,18 @@ func TestRateLimitMiddleware_OSCAR(t *testing.T) {
 			// the 7th request drives the average below DisconnectLevel
 			requests:  7,
 			wantCalls: 1,
-			// EvaluateRateLimit closes the session on disconnect, so the event
-			// is best-effort; it is queued but the client may never fetch it.
-			wantEvents: []string{"limit", "disconnect"},
 			// a disconnected session has no aimsid left to retry with
 			wantRetryAfter: "",
 		},
 		{
-			// The client renders the rateLimit event as an IM alert inside a
-			// conversation window, so a class the IM path does not spend is
-			// enforced without one.
-			name:           "a class the alert cannot describe is limited silently",
+			// A non-IM class is enforced the same way. The middleware only ever
+			// rejects; notifying the client is the per-account monitor's job, and
+			// it surfaces only the IM class to the alert.
+			name:           "a non-IM class is still enforced",
 			foodGroup:      wire.Feedbag,
 			subGroup:       wire.FeedbagInsertItem,
 			requests:       5,
 			wantCalls:      1,
-			wantEvents:     nil,
 			wantRetryAfter: "1",
 		},
 		{
@@ -212,7 +203,9 @@ func TestRateLimitMiddleware_OSCAR(t *testing.T) {
 			}
 
 			assert.Equal(t, tt.wantCalls, calls)
-			assert.Equal(t, tt.wantEvents, rateLimitEventStatuses(t, session))
+			// The middleware enforces the limit but never pushes rate limit
+			// events; that is the per-account monitor's responsibility.
+			assert.Empty(t, rateLimitEventStatuses(t, session))
 
 			if tt.wantCalls < tt.requests {
 				assertRateLimited(t, last, tt.wantRetryAfter)
@@ -223,109 +216,49 @@ func TestRateLimitMiddleware_OSCAR(t *testing.T) {
 	}
 }
 
-// The client only dismisses its rate limit alert when it receives a "clear"
-// event, and only pushes on a transition, so a recovering session must produce
-// exactly one clear.
-func TestRateLimitMiddleware_OSCAR_pushesClearOnRecovery(t *testing.T) {
-	session := newTestWebAPISession(t, tightRateLimitClasses())
-	middleware := newTestRateLimitMiddleware()
+// The monitor broadcasts transitions, not current state, so without a seed a
+// session signing on mid-limit shows no banner while its sends are rejected — and
+// the client's alert is sticky, so the eventual "clear" has nothing to dismiss.
+func TestSeedRateLimitAlert(t *testing.T) {
+	imClass, ok := wire.DefaultSNACRateLimits().RateClassLookup(wire.ICBM, wire.ICBMChannelMsgToHost)
+	require.True(t, ok)
 
-	handler := middleware.OSCAR(wire.ICBM, wire.ICBMChannelMsgToHost)(
-		func(w http.ResponseWriter, r *http.Request, s *state.WebAPISession) {})
+	// limitedSession returns a session on an account already in the limited state.
+	limitedSession := func(t *testing.T) *state.WebAPISession {
+		t.Helper()
 
-	// Trip the limit.
-	for range 3 {
-		handler(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/im/sendIM", nil), session)
+		session := newTestWebAPISession(t, tightRateLimitClasses())
+		sess := session.OSCARSession.Session()
+		for i := 0; sess.RateLimitStates()[imClass-1].CurrentStatus != wire.RateLimitStatusLimited; i++ {
+			require.Less(t, i, 100, "class never reached the limited state")
+			sess.EvaluateRateLimit(time.Now(), imClass)
+		}
+		return session
 	}
-	assert.Equal(t, []string{"limit"}, rateLimitEventStatuses(t, session))
 
-	// Let the moving average recover past ClearLevel. See tightRateLimitClasses
-	// for why this is a short wait rather than a stubbed clock.
-	time.Sleep(250 * time.Millisecond)
+	t.Run("a session starting on a limited account is told", func(t *testing.T) {
+		session := limitedSession(t)
 
-	rec := httptest.NewRecorder()
-	handler(rec, httptest.NewRequest(http.MethodGet, "/im/sendIM", nil), session)
+		seedRateLimitAlert(session, imClass)
 
-	assert.Equal(t, http.StatusOK, rec.Code)
-	assert.Equal(t, []string{"limit", "clear"}, rateLimitEventStatuses(t, session))
-}
-
-// A session that trips the limit and then goes idle must still be told it
-// recovered: the fetchEvents poll carries the clear even though no charged
-// request arrives to recompute the status.
-func TestRateLimitMiddleware_RecoverOnPoll(t *testing.T) {
-	session := newTestWebAPISession(t, tightRateLimitClasses())
-	middleware := newTestRateLimitMiddleware()
-
-	oscar := middleware.OSCAR(wire.ICBM, wire.ICBMChannelMsgToHost)(
-		func(w http.ResponseWriter, r *http.Request, s *state.WebAPISession) {})
+		assert.Equal(t, []string{"limit"}, rateLimitEventStatuses(t, session))
+	})
 
-	// Trip the limit.
-	for range 3 {
-		oscar(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/im/sendIM", nil), session)
-	}
-	assert.Equal(t, []string{"limit"}, rateLimitEventStatuses(t, session))
+	t.Run("a session starting on a clear account is told nothing", func(t *testing.T) {
+		session := newTestWebAPISession(t, tightRateLimitClasses())
 
-	// Let the moving average recover past ClearLevel without sending anything.
-	// See tightRateLimitClasses for why this is a short wait, not a stubbed clock.
-	time.Sleep(250 * time.Millisecond)
+		seedRateLimitAlert(session, imClass)
 
-	calls := 0
-	poll := middleware.RecoverOnPoll(func(w http.ResponseWriter, r *http.Request, s *state.WebAPISession) {
-		calls++
+		assert.Empty(t, rateLimitEventStatuses(t, session))
 	})
 
-	// The poll surfaces the clear and still serves the handler.
-	poll(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/aim/fetchEvents", nil), session)
-	assert.Equal(t, 1, calls)
-	assert.Equal(t, []string{"limit", "clear"}, rateLimitEventStatuses(t, session))
-
-	// A second poll after recovery pushes nothing further.
-	poll(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/aim/fetchEvents", nil), session)
-	assert.Equal(t, 2, calls)
-	assert.Equal(t, []string{"limit", "clear"}, rateLimitEventStatuses(t, session))
-}
-
-// Two browser tabs on one account share a single set of rate limit states, so
-// whichever of them polls first observes the underlying limited -> clear
-// transition. The tab that never showed an alert must not be able to consume the
-// recovery the limited tab is waiting for.
-func TestRateLimitMiddleware_RecoverOnPoll_perClientTransitions(t *testing.T) {
-	instance := newTestOSCARInstance(t, tightRateLimitClasses())
-	tabA := newTestWebAPISessionOn("aimsid-a", instance)
-	tabB := newTestWebAPISessionOn("aimsid-b", instance)
+	t.Run("a zero class id disables the alert", func(t *testing.T) {
+		session := limitedSession(t)
 
-	middleware := newTestRateLimitMiddleware()
-	send := middleware.OSCAR(wire.ICBM, wire.ICBMChannelMsgToHost)(
-		func(w http.ResponseWriter, r *http.Request, s *state.WebAPISession) {})
-	poll := middleware.RecoverOnPoll(
-		func(w http.ResponseWriter, r *http.Request, s *state.WebAPISession) {})
+		seedRateLimitAlert(session, 0)
 
-	// Tab A trips the limit and raises the alert.
-	for range 3 {
-		send(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/im/sendIM", nil), tabA)
-	}
-	assert.Equal(t, []string{"limit"}, rateLimitEventStatuses(t, tabA))
-
-	// Tab B polls while the account is still limited. The budget is shared, so
-	// tab B really cannot send either — but it is idle, and a poll is not a
-	// request it made, so it is not told to stop chatting.
-	poll(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/aim/fetchEvents", nil), tabB)
-	assert.Empty(t, rateLimitEventStatuses(t, tabB))
-
-	// Let the moving average recover past ClearLevel. See tightRateLimitClasses
-	// for why this is a short wait rather than a stubbed clock.
-	time.Sleep(250 * time.Millisecond)
-
-	// Tab B's idle poll observes the recovery first. It has no alert up, so there
-	// is still nothing to tell it.
-	poll(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/aim/fetchEvents", nil), tabB)
-	assert.Empty(t, rateLimitEventStatuses(t, tabB))
-
-	// Tab A still gets its clear on the next poll, even though tab B already
-	// consumed the transition it would otherwise have been derived from.
-	poll(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "/aim/fetchEvents", nil), tabA)
-	assert.Equal(t, []string{"limit", "clear"}, rateLimitEventStatuses(t, tabA))
+		assert.Empty(t, rateLimitEventStatuses(t, session))
+	})
 }
 
 // Retry-After must name a wait that actually clears the limit. A rejected

+ 29 - 6
server/webapi/handlers/session.go

@@ -28,9 +28,12 @@ type SessionHandler struct {
 	IconSource       BuddyIconSource
 	Logger           *slog.Logger
 	OServiceService  OServiceService
-	FnSessCfg        func(sess *state.Session)
-	FnSessInit       func(instance *state.SessionInstance) func() error
-	FnInstanceClose  func(instance *state.SessionInstance) func()
+	// the same SNAC-to-rate-class mapping RateLimitMiddleware enforces against, so
+	// the class a session alerts on cannot drift from the one it is charged
+	SNACRateLimits  wire.SNACRateLimits
+	FnSessCfg       func(sess *state.Session)
+	FnSessInit      func(instance *state.SessionInstance) func() error
+	FnInstanceClose func(instance *state.SessionInstance) func()
 }
 
 // AuthService defines methods needed for authentication.
@@ -224,9 +227,8 @@ func (h *SessionHandler) StartSession(w http.ResponseWriter, r *http.Request) {
 
 	if err = instance.Session().RunOnce(h.FnSessInit(instance)); err != nil {
 		h.Logger.ErrorContext(context.Background(), "failed to init session", "err", err.Error())
-		// Nothing owns the instance yet, so close it here. Left open, it stays
-		// registered in the session manager, counting against the user's
-		// concurrent-instance budget until the server restarts.
+		// RunOnce has already closed the whole session; this is belt-and-braces,
+		// since CloseInstance is idempotent and nothing else owns the instance yet.
 		instance.CloseInstance()
 		h.sendError(w, r, http.StatusInternalServerError, "internal server error")
 		return
@@ -255,6 +257,23 @@ func (h *SessionHandler) StartSession(w http.ResponseWriter, r *http.Request) {
 		return
 	}
 
+	// The rate class sending an IM spends. Only its updates surface to the client's
+	// conversation-window alert, which renders any rateLimit event as the IM
+	// banner. A miss yields zero, which disables the alert rather than indexing a
+	// rate class that isn't there.
+	imRateClassID, ok := h.SNACRateLimits.RateClassLookup(wire.ICBM, wire.ICBMChannelMsgToHost)
+	if !ok {
+		h.Logger.ErrorContext(ctx, "no rate class maps to sending an IM, rate limit events are disabled")
+	}
+
+	// Subscribe the same way an OSCAR client does at handshake, so the per-account
+	// monitor broadcasts this class's transitions to this session.
+	if imRateClassID != 0 {
+		h.OServiceService.RateParamsSubAdd(ctx, instance, wire.SNAC_0x01_0x08_OServiceRateParamsSubAdd{
+			ClassIDs: []uint16{uint16(imRateClassID)},
+		})
+	}
+
 	// Create WebAPI session
 	// Record the origin the client reached us on. Asset URLs published to the
 	// client (buddy icons) must be absolute, since the client page is served from
@@ -323,6 +342,10 @@ func (h *SessionHandler) StartSession(w http.ResponseWriter, r *http.Request) {
 		return permitDenyData(reply.Items), nil
 	}
 
+	// Only IM-class rate limit updates should surface to the client alert.
+	session.IMRateClassID = imRateClassID
+	seedRateLimitAlert(session, imRateClassID)
+
 	// Now that every refresher callback is wired, start the OSCAR listener. Doing
 	// this inside CreateSession would race these assignments, since the goroutine
 	// reads the callbacks as it converts SNACs into events.

+ 5 - 6
server/webapi/server.go

@@ -36,6 +36,7 @@ func NewServer(listeners []string, logger *slog.Logger, handler Handler, apiKeyV
 		IconSource:       handler.IconSource,
 		Logger:           logger,
 		OServiceService:  handler.OServiceService,
+		SNACRateLimits:   handler.SNACRateLimits,
 	}
 
 	eventsHandler := &handlers.EventsHandler{
@@ -165,12 +166,8 @@ func NewServer(listeners []string, logger *slog.Logger, handler Handler, apiKeyV
 		mux.Handle("GET /aim/endSession", sessionRoute(sessionHandler.EndSession))
 
 		// Event fetching - uses aimsid for auth, no k required. This is the
-		// long-poll loop the client runs continuously; RecoverOnPoll uses it to
-		// surface a rate limit "clear" to an idle session that already recovered.
-		mux.Handle("GET /aim/fetchEvents", authMiddleware.AuthenticateFlexible(
-			authMiddleware.CORSMiddleware(
-				authMiddleware.RequireSession(sessionManager,
-					rateLimiter.RecoverOnPoll(eventsHandler.FetchEvents)))))
+		// long-poll loop the client runs continuously.
+		mux.Handle("GET /aim/fetchEvents", sessionRoute(eventsHandler.FetchEvents))
 
 		// Temp buddies are session-local rather than feedbag-backed, but they
 		// are the Web API's equivalent of the BUDDY temp buddy SNACs and are
@@ -309,6 +306,8 @@ func NewServer(listeners []string, logger *slog.Logger, handler Handler, apiKeyV
 			}
 			// periodically decay warning level
 			go handler.LowerWarnLevel(shutdownCtx, instance)
+			// broadcast rate limit transitions to every instance on the account
+			go handler.OServiceService.MonitorRateLimits(shutdownCtx, instance.Session())
 			return nil
 		}
 	}

+ 2 - 0
server/webapi/types.go

@@ -21,6 +21,8 @@ type ICBMService interface {
 type OServiceService interface {
 	ClientOnline(ctx context.Context, service uint16, inBody wire.SNAC_0x01_0x02_OServiceClientOnline, instance *state.SessionInstance) error
 	IdleNotification(ctx context.Context, instance *state.SessionInstance, inBody wire.SNAC_0x01_0x11_OServiceIdleNotification) error
+	MonitorRateLimits(ctx context.Context, session *state.Session)
+	RateParamsSubAdd(ctx context.Context, instance *state.SessionInstance, inBody wire.SNAC_0x01_0x08_OServiceRateParamsSubAdd)
 	ServiceRequest(ctx context.Context, service uint16, instance *state.SessionInstance, inFrame wire.SNACFrame, inBody wire.SNAC_0x01_0x04_OServiceServiceRequest, listener config.Listener) (wire.SNACMessage, error)
 }
 

+ 89 - 55
state/session.go

@@ -95,9 +95,16 @@ type Session struct {
 	instances        map[uint8]*SessionInstance
 	instancesOrdered []*SessionInstance
 
-	initOnce      sync.Once
+	initOnce sync.Once
+	// what RunOnce's fn returned, so every caller learns the outcome, not just
+	// the one that ran fn
+	initErr       error
 	onSessCloseFn func()
-	nowFn         func() time.Time
+	// closeCh is closed once the account's last instance departs, so per-account
+	// goroutines can stop with the session. closed guards its one-time close.
+	closeCh chan struct{}
+	closed  bool
+	nowFn   func() time.Time
 }
 
 // NewSession creates a new Session for a user.
@@ -107,6 +114,7 @@ func NewSession() *Session {
 		instances:        make(map[uint8]*SessionInstance),
 		instancesOrdered: make([]*SessionInstance, 0),
 		onSessCloseFn:    func() {},
+		closeCh:          make(chan struct{}),
 		nowFn:            time.Now,
 	}
 }
@@ -121,6 +129,10 @@ func (s *Session) AddInstance() *SessionInstance {
 	s.mutex.Lock()
 	defer s.mutex.Unlock()
 
+	if s.closed {
+		return nil
+	}
+
 	instance := &SessionInstance{
 		session:           s,
 		instanceNum:       s.generateInstanceNum(),
@@ -169,10 +181,36 @@ func (s *Session) Instances() []*SessionInstance {
 	return instances
 }
 
+// removeInstanceAndMaybeClose removes instance and reports whether it was the
+// last one, marking the session closed in the same critical section that made it
+// empty. Doing both under one lock is what makes AddInstance's closed check
+// meaningful: a sign-on racing the departure either takes the lock first (and
+// keeps the session alive) or finds it closed and is refused.
+func (s *Session) removeInstanceAndMaybeClose(instance *SessionInstance) bool {
+	s.mutex.Lock()
+	defer s.mutex.Unlock()
+
+	s.removeInstanceLocked(instance)
+
+	if len(s.instances) > 0 {
+		return false
+	}
+	if !s.closed {
+		s.closed = true
+		close(s.closeCh)
+	}
+	return true
+}
+
 // RemoveInstance removes an instance from the session group.
 func (s *Session) RemoveInstance(instance *SessionInstance) {
 	s.mutex.Lock()
 	defer s.mutex.Unlock()
+	s.removeInstanceLocked(instance)
+}
+
+// removeInstanceLocked removes an instance. The caller must hold s.mutex.
+func (s *Session) removeInstanceLocked(instance *SessionInstance) {
 	delete(s.instances, instance.instanceNum)
 	for i, inst := range s.instancesOrdered {
 		if inst == instance {
@@ -424,50 +462,6 @@ func (s *Session) EvaluateRateLimit(now time.Time, rateClassID wire.RateLimitCla
 	return status
 }
 
-// RecoverRateLimits recomputes, for each rate class the session is currently
-// limited on, whether enough time has elapsed since the last charged request for
-// the moving average to climb back above the clear threshold, and returns the
-// classes whose status changed (in practice, limited -> clear).
-//
-// Unlike EvaluateRateLimit it does not treat the call as a new request: LastTime
-// is preserved, so successive calls measure the full elapsed time since the last
-// real request rather than the (small) gap between calls. This mirrors how OSCAR's
-// ObserveRateChanges detects recovery on its ticker, and is what lets a repeated
-// poll surface the clear. Classes that are not currently limited are left
-// untouched, so a poll never perturbs a class the user is actively spending.
-func (s *Session) RecoverRateLimits(now time.Time) []RateClassState {
-	s.mutex.Lock()
-	defer s.mutex.Unlock()
-
-	var changed []RateClassState
-	for i := range s.rateLimitStates {
-		rc := &s.rateLimitStates[i]
-		if !rc.LimitedNow {
-			continue
-		}
-
-		status, newLevel := wire.CheckRateLimit(rc.LastTime, now, rc.RateClass, rc.CurrentLevel, rc.LimitedNow)
-		// Recovery only. A poll must never escalate a limited class: a short
-		// elapsed gap can make CheckRateLimit report disconnect, but an idle
-		// session recovering is not abuse, and disconnect is enforced by
-		// EvaluateRateLimit on real requests instead. While limitedNow the only
-		// non-escalating outcome CheckRateLimit yields is clear, so that is the
-		// sole transition acted on; a still-limited class is left untouched (its
-		// CurrentLevel keeps recovering from the same base, as in
-		// ObserveRateChanges).
-		if status != wire.RateLimitStatusClear {
-			continue
-		}
-
-		rc.CurrentStatus = status
-		rc.CurrentLevel = newLevel
-		rc.LimitedNow = false
-		changed = append(changed, *rc)
-	}
-
-	return changed
-}
-
 // ObserveRateChanges updates rate limit states and returns changes.
 func (s *Session) ObserveRateChanges(now time.Time) (classDelta []RateClassState, stateDelta []RateClassState) {
 	s.mutex.Lock()
@@ -625,6 +619,11 @@ func (s *Session) WarningCh() chan uint16 {
 
 // CloseSession closes all instances in the session.
 func (s *Session) CloseSession() {
+	// Mark closed before tearing the instances down: AddInstance refuses a closed
+	// session, so a sign-on landing mid-teardown builds a fresh Session rather
+	// than joining this one.
+	s.markClosed()
+
 	s.mutex.RLock()
 	instances := make([]*SessionInstance, 0, len(s.instances))
 	for _, instance := range s.instances {
@@ -644,15 +643,54 @@ func (s *Session) OnSessionClose(fn func()) {
 	s.onSessCloseFn = fn
 }
 
+// Closed returns a channel that is closed once the account's last instance has
+// departed. Unlike OnSessionClose (a single-slot callback), any number of
+// per-account goroutines can select on it.
+func (s *Session) Closed() <-chan struct{} {
+	return s.closeCh
+}
+
+// IsClosed is the non-blocking counterpart of Closed.
+func (s *Session) IsClosed() bool {
+	select {
+	case <-s.closeCh:
+		return true
+	default:
+		return false
+	}
+}
+
+// markClosed closes the session's Closed channel exactly once.
+func (s *Session) markClosed() {
+	s.mutex.Lock()
+	defer s.mutex.Unlock()
+	if s.closed {
+		return
+	}
+	s.closed = true
+	close(s.closeCh)
+}
+
 // RunOnce executes the given function once across all invocations. Used to
 // run arbitrary code that must only run once when the first session instance
 // connects. The function must not block.
+//
+// If fn returns an error the session is closed, because init is one-shot: fn
+// starts the account's per-session goroutines (rate limit monitor, warning decay)
+// and will not run again, so a half-initialized session that later instances
+// attach to would run without them for its whole life. Closing it makes the next
+// sign-on build a fresh Session.
+//
+// The error is reported to every caller, not just the one that ran fn — a
+// concurrent caller must not read a failed init as success. sync.Once establishes
+// the happens-before edge that makes initErr safe to read after Do returns.
 func (s *Session) RunOnce(fn func() error) error {
-	var err error
 	s.initOnce.Do(func() {
-		err = fn()
+		if s.initErr = fn(); s.initErr != nil {
+			s.CloseSession()
+		}
 	})
-	return err
+	return s.initErr
 }
 
 // SetNowFn sets the function used to get the current time. This is useful for testing.
@@ -1188,10 +1226,7 @@ func (s *SessionInstance) CloseInstance() {
 	s.mutex.Unlock()
 
 	// remove the instance now so that the function has an updated view of the world
-	s.session.RemoveInstance(s)
-
-	count := s.session.InstanceCount()
-	if count == 0 {
+	if s.session.removeInstanceAndMaybeClose(s) {
 		s.session.mutex.RLock()
 		onSessCloseFn := s.session.onSessCloseFn
 		s.session.mutex.RUnlock()
@@ -1304,8 +1339,7 @@ func (s *SessionInstance) closeOnly() {
 	s.closed = true
 	s.mutex.Unlock()
 
-	s.session.RemoveInstance(s)
-	if s.session.InstanceCount() == 0 {
+	if s.session.removeInstanceAndMaybeClose(s) {
 		s.session.mutex.RLock()
 		onSessCloseFn := s.session.onSessCloseFn
 		s.session.mutex.RUnlock()

+ 34 - 3
state/session_manager.go

@@ -210,6 +210,28 @@ func (s *InMemorySessionManager) AddSession(ctx context.Context, screenName Disp
 	active := s.findRec(screenName.IdentScreenName())
 	s.mapMutex.Unlock()
 
+	// A closed session is a tombstone: its last instance has departed (or its
+	// RunOnce init failed), but RemoveSession has not run yet — Signout only
+	// reaches it at the end of onSessCloseFn.
+	//
+	// Attaching to it is wrong: its RunOnce is spent and its per-account
+	// goroutines have exited, so the new instance would get no rate limit monitor
+	// or warning decay. Evicting it and racing ahead is wrong too: the teardown's
+	// UnregisterBuddyList and RemoveUserFromAllChats would land after the fresh
+	// session registered itself, leaving a signed-on user invisible to buddies.
+	//
+	// So wait the teardown out, as the single-session displacement path below
+	// does, then build a fresh session. RemoveSession is the last act of every
+	// onSessCloseFn, so the wait is bounded by the teardown and by ctx.
+	if active != nil && active.session.IsClosed() {
+		select {
+		case <-active.removed: // wait for RemoveSession to be called
+		case <-ctx.Done():
+			return nil, fmt.Errorf("waiting for closed session to be torn down: %w", ctx.Err())
+		}
+		active = nil
+	}
+
 	if active != nil {
 		if doMultiSess {
 			if !active.multiSession {
@@ -222,9 +244,18 @@ func (s *InMemorySessionManager) AddSession(ctx context.Context, screenName Disp
 				return nil, fmt.Errorf("%w: max instance(s) = %d", ErrMaxConcurrentSessionsReached, s.maxConcurrentSessions)
 			}
 
-			instance := active.session.AddInstance()
-
-			return instance, nil
+			// AddInstance refuses a session that closed after the tombstone
+			// check above — the account's last instance departed in the
+			// interim. Wait that teardown out and build a fresh session, as
+			// the tombstone path does.
+			if instance := active.session.AddInstance(); instance != nil {
+				return instance, nil
+			}
+			select {
+			case <-active.removed: // wait for RemoveSession to be called
+			case <-ctx.Done():
+				return nil, fmt.Errorf("waiting for closed session to be torn down: %w", ctx.Err())
+			}
 		} else {
 			// signal to callers that this session group has to go
 			active.session.CloseSession()

+ 180 - 0
state/session_manager_test.go

@@ -2,6 +2,7 @@ package state
 
 import (
 	"context"
+	"io"
 	"log/slog"
 	"math/rand"
 	"sync"
@@ -36,6 +37,185 @@ func TestInMemorySessionManager_AddSession(t *testing.T) {
 	assert.Contains(t, sm.AllSessions(), sess2.Session())
 }
 
+// A closed session lingers in the store until RemoveSession runs, which Signout
+// only reaches at the END of onSessCloseFn. A sign-on landing in that window must
+// neither attach to the dead session (its RunOnce is spent, so no rate limit
+// monitor or warning decay) nor evict it and race ahead (the in-flight
+// UnregisterBuddyList would land after the fresh session registered itself). It
+// has to wait the teardown out, then build a fresh session.
+func TestInMemorySessionManager_AddSession_WaitsForClosedSessionTeardown(t *testing.T) {
+	for _, doMultiSess := range []bool{true, false} {
+		name := "multi-session client"
+		if !doMultiSess {
+			name = "single-session client"
+		}
+		t.Run(name, func(t *testing.T) {
+			synctest.Test(t, func(t *testing.T) {
+				sm := NewInMemorySessionManager(slog.Default())
+				ctx := context.Background()
+
+				var mu sync.Mutex
+				var events []string
+				record := func(event string) {
+					mu.Lock()
+					defer mu.Unlock()
+					events = append(events, event)
+				}
+
+				// Stands in for the servers' OnSessionClose: announce the
+				// departure, tear the account's registrations down, and only
+				// then Signout, whose last act is RemoveSession.
+				teardown := func(sess *Session) {
+					sess.OnSessionClose(func() {
+						record("teardown begin")
+						time.Sleep(time.Second) // BroadcastBuddyDeparted
+						record("unregister buddy list")
+						sm.RemoveSession(sess) // Signout
+					})
+				}
+
+				inst1, err := sm.AddSession(ctx, "user-screen-name", doMultiSess, teardown)
+				require.NoError(t, err)
+				inst1.SetSignonComplete()
+				dead := inst1.Session()
+
+				// The account's last instance departs, running onSessCloseFn on
+				// its own goroutine as a closing connection does.
+				wg := &sync.WaitGroup{}
+				wg.Add(1)
+				go func() {
+					defer wg.Done()
+					inst1.CloseInstance()
+				}()
+
+				// Let the teardown get under way and park mid-flight. The
+				// session is closed the moment its last instance departs, but
+				// the store still holds it — the tombstone window. It is not
+				// observable through AllSessions/RetrieveSession (both filter on
+				// live instances), which is why AddSession checks IsClosed.
+				synctest.Wait()
+				require.True(t, dead.IsClosed())
+				sm.mapMutex.RLock()
+				require.NotNil(t, sm.findRec(dead.IdentScreenName()), "tombstone should still be in the store")
+				sm.mapMutex.RUnlock()
+
+				inst2, err := sm.AddSession(ctx, "user-screen-name", doMultiSess, teardown)
+				record("signon complete")
+				require.NoError(t, err)
+				inst2.SetSignonComplete()
+				wg.Wait()
+
+				// The sign-on waited the whole teardown out. Without the wait,
+				// "signon complete" lands before "unregister buddy list" and
+				// that unregister wipes the buddy list the new session has
+				// already registered.
+				mu.Lock()
+				assert.Equal(t, []string{
+					"teardown begin",
+					"unregister buddy list",
+					"signon complete",
+				}, events)
+				mu.Unlock()
+
+				// And it got a fresh session rather than the tombstone, so its
+				// RunOnce is unspent and its per-account goroutines start.
+				assert.NotSame(t, dead, inst2.Session(), "must not attach to the closed session")
+				assert.False(t, inst2.Session().IsClosed())
+				assert.Contains(t, sm.AllSessions(), inst2.Session())
+
+				ran := false
+				require.NoError(t, inst2.Session().RunOnce(func() error { ran = true; return nil }))
+				assert.True(t, ran)
+			})
+		})
+	}
+}
+
+// The tombstone gate reads IsClosed and then adds an instance to the session it
+// just cleared, so a sign-on whose AddInstance lands as the last connection
+// departs can pass the gate and have the session close underneath it. The
+// instance is then live but orphaned: Closed() has fired, so the rate limit
+// monitor has exited, and RunOnce is spent so nothing restarts it.
+//
+// AddSession must never return an instance on a closed session.
+func TestInMemorySessionManager_AddSession_RacesLastInstanceDeparture(t *testing.T) {
+	const iterations = 20_000
+
+	logger := slog.New(slog.NewTextHandler(io.Discard, nil))
+
+	for i := range iterations {
+		sm := NewInMemorySessionManager(logger)
+
+		ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
+
+		// Stands in for the servers' OnSessionClose, whose last act is
+		// RemoveSession. Without it the tombstone wait has nothing to wake on.
+		teardown := func(sess *Session) {
+			sess.OnSessionClose(func() { sm.RemoveSession(sess) })
+		}
+
+		inst1, err := sm.AddSession(ctx, "user-screen-name", true, teardown)
+		require.NoError(t, err)
+		inst1.SetSignonComplete()
+		departing := inst1.Session()
+
+		// The account's last connection drops at the same moment a new one signs
+		// on. Neither path takes a lock the other holds: CloseInstance runs on a
+		// connection goroutine, AddSession under the per-user lock.
+		start := make(chan struct{})
+		var wg sync.WaitGroup
+		wg.Add(2)
+
+		go func() {
+			defer wg.Done()
+			<-start
+			inst1.CloseInstance()
+		}()
+
+		var inst2 *SessionInstance
+		var addErr error
+		go func() {
+			defer wg.Done()
+			<-start
+			inst2, addErr = sm.AddSession(ctx, "user-screen-name", true, teardown)
+		}()
+
+		close(start)
+		wg.Wait()
+		cancel()
+
+		require.NoError(t, addErr, "iteration %d", i)
+		require.False(t, inst2.Session().IsClosed(),
+			"iteration %d: sign-on returned an instance on a closed session "+
+				"(same session as the departing connection: %t) — its rate limit "+
+				"monitor has already exited and RunOnce is spent",
+			i, inst2.Session() == departing)
+	}
+}
+
+// A teardown that never reaches RemoveSession must not wedge the sign-on
+// forever; the wait is bounded by the caller's context, as the single-session
+// displacement path already was.
+func TestInMemorySessionManager_AddSession_ClosedSessionTeardownTimesOut(t *testing.T) {
+	synctest.Test(t, func(t *testing.T) {
+		sm := NewInMemorySessionManager(slog.Default())
+
+		inst1, err := sm.AddSession(context.Background(), "user-screen-name", true)
+		require.NoError(t, err)
+		inst1.SetSignonComplete()
+
+		// No OnSessionClose is registered, so RemoveSession is never called and
+		// the tombstone stays in the store for good.
+		inst1.CloseInstance()
+
+		ctx, cancel := context.WithTimeout(context.Background(), time.Second)
+		defer cancel()
+
+		_, err = sm.AddSession(ctx, "user-screen-name", true)
+		assert.ErrorIs(t, err, context.DeadlineExceeded)
+	})
+}
+
 func TestInMemorySessionManager_AddSession_AppliesCfgToSession(t *testing.T) {
 	sm := NewInMemorySessionManager(slog.Default())
 	wantUIN := uint32(424242)

+ 100 - 39
state/session_test.go

@@ -7,6 +7,7 @@ import (
 	"net/netip"
 	"sync"
 	"testing"
+	"testing/synctest"
 	"time"
 
 	"github.com/mk6i/open-oscar-server/wire"
@@ -791,45 +792,6 @@ func TestSession_EvaluateRateLimit_ObserveRateChanges(t *testing.T) {
 			assert.Equal(t, wire.RateLimitStatusClear, have)
 		}
 	})
-
-	t.Run("RecoverRateLimits clears an idle limited class without counting as a request", func(t *testing.T) {
-		now := time.Now()
-
-		sess := NewSession()
-		sess.AddInstance()
-		sess.SetRateClasses(now, rateClasses)
-
-		rateClass := rateClasses.Get(3)
-
-		// Drive the class into the limited state.
-		var status wire.RateLimitStatus
-		for status != wire.RateLimitStatusLimited {
-			now = now.Add(1 * time.Second)
-			status = sess.EvaluateRateLimit(now, rateClass.ID)
-		}
-
-		// Shortly after, the class is still limited: RecoverRateLimits reports no
-		// transition and, crucially, does not advance LastTime (so it is not
-		// charged as a request).
-		now = now.Add(1 * time.Second)
-		assert.Empty(t, sess.RecoverRateLimits(now))
-		assert.Equal(t, wire.RateLimitStatusLimited, sess.RateLimitStates()[rateClass.ID-1].CurrentStatus)
-
-		// After a long idle gap the moving average clears; the transition is
-		// reported exactly once.
-		now = now.Add(1 * time.Hour)
-		changed := sess.RecoverRateLimits(now)
-		if assert.Len(t, changed, 1) {
-			assert.Equal(t, rateClass.ID, changed[0].ID)
-			assert.Equal(t, wire.RateLimitStatusClear, changed[0].CurrentStatus)
-			assert.False(t, changed[0].LimitedNow)
-		}
-
-		// Once cleared, further polls report nothing.
-		now = now.Add(1 * time.Second)
-		assert.Empty(t, sess.RecoverRateLimits(now))
-		assert.Equal(t, wire.RateLimitStatusClear, sess.RateLimitStates()[rateClass.ID-1].CurrentStatus)
-	})
 }
 
 func TestSession_SetAndGetFoodGroupVersions(t *testing.T) {
@@ -1908,6 +1870,105 @@ func TestSession_RunOnce(t *testing.T) {
 		assert.Error(t, err)
 		assert.Equal(t, expectedErr, err)
 	})
+
+	// Init is one-shot, so a session that failed to initialize must not linger for
+	// later instances to attach to.
+	t.Run("closes the session when the function fails", func(t *testing.T) {
+		s := NewSession()
+		instance := s.AddInstance()
+
+		err := s.RunOnce(func() error {
+			return assert.AnError
+		})
+
+		assert.Error(t, err)
+		assert.True(t, s.IsClosed(), "session should be closed after a failed init")
+		assert.True(t, instance.IsClosed(), "instances should be torn down with the session")
+	})
+
+	t.Run("leaves the session open when the function succeeds", func(t *testing.T) {
+		s := NewSession()
+		s.AddInstance()
+
+		assert.NoError(t, s.RunOnce(func() error { return nil }))
+		assert.False(t, s.IsClosed())
+	})
+
+	// A failed init tears the session down, so a caller that did not run fn must
+	// still learn it failed rather than sign a client on to a session whose
+	// instances have just been closed underneath it.
+	t.Run("reports the init error to every caller", func(t *testing.T) {
+		s := NewSession()
+		s.AddInstance()
+
+		first := s.RunOnce(func() error { return assert.AnError })
+		second := s.RunOnce(func() error { return nil })
+
+		assert.ErrorIs(t, first, assert.AnError)
+		assert.ErrorIs(t, second, assert.AnError, "the caller that did not run fn must not read failure as success")
+	})
+
+	// Two instances can be added to a fresh session and reach RunOnce
+	// concurrently. Whichever one loses the race still has a client to sign on,
+	// and must be told the session is not usable.
+	t.Run("reports the init error to a concurrent caller", func(t *testing.T) {
+		synctest.Test(t, func(t *testing.T) {
+			s := NewSession()
+			s.AddInstance()
+			s.AddInstance()
+
+			var wg sync.WaitGroup
+			errs := make([]error, 2)
+			for i := range errs {
+				wg.Add(1)
+				go func() {
+					defer wg.Done()
+					errs[i] = s.RunOnce(func() error { return assert.AnError })
+				}()
+			}
+			wg.Wait()
+
+			for i, err := range errs {
+				assert.ErrorIs(t, err, assert.AnError, "caller %d", i)
+			}
+		})
+	})
+}
+
+// An instance added to a closed session is live but orphaned: Closed() has
+// already fired, so the per-account goroutines that select on it — the rate limit
+// monitor above all — have exited, and RunOnce is spent so nothing restarts them.
+// AddInstance must refuse, so AddSession can fall through to the tombstone wait it
+// already implements.
+func TestSession_AddInstance_RefusesClosedSession(t *testing.T) {
+	s := NewSession()
+	first := s.AddInstance()
+
+	// Stand in for OServiceService.MonitorRateLimits: a per-account goroutine
+	// started from RunOnce that runs until the session closes.
+	require.NoError(t, s.RunOnce(func() error { return nil }))
+	monitorDone := make(chan struct{})
+	go func() {
+		defer close(monitorDone)
+		<-s.Closed()
+	}()
+
+	// The account's last instance departs, closing the session and stopping the
+	// monitor. This is the moment AddSession's IsClosed check races.
+	first.CloseInstance()
+	select {
+	case <-monitorDone:
+	case <-time.After(2 * time.Second):
+		t.Fatal("monitor did not exit when the session closed")
+	}
+
+	assert.Nil(t, s.AddInstance(), "AddInstance must refuse a closed session")
+
+	// RunOnce is spent, so an instance that did attach could never start a
+	// replacement monitor.
+	ran := false
+	assert.NoError(t, s.RunOnce(func() error { ran = true; return nil }))
+	assert.False(t, ran)
 }
 
 func TestSession_CloseInstance(t *testing.T) {

+ 55 - 41
state/webapi_session.go

@@ -86,12 +86,12 @@ type WebAPISession struct {
 	aliasMu      sync.Mutex
 	imLog        map[string][]WebAPIStoredIM
 	imLogMu      sync.Mutex
-	// rateAlertStatus is the rate limit status this client was last told about,
-	// zero until it has been told anything. See ObserveRateAlert.
-	rateAlertStatus wire.RateLimitStatus
-	rateAlertMu     sync.Mutex
-	logger          *slog.Logger // Logger for debugging
-	listeners       sync.WaitGroup
+	// IMRateClassID is the rate class that sending an IM spends. The web client
+	// renders any rate limit event as the IM banner, so only this class's updates
+	// may reach it. Zero disables the alert.
+	IMRateClassID wire.RateLimitClassID
+	logger        *slog.Logger // Logger for debugging
+	listeners     sync.WaitGroup
 
 	ctx    context.Context
 	cancel context.CancelFunc
@@ -105,35 +105,6 @@ func (s *WebAPISession) IsExpired() bool {
 	return time.Now().After(s.ExpiresAt)
 }
 
-// ObserveRateAlert records status as the rate limit status this web client was
-// last told about, and reports whether it differs from the one before it.
-//
-// The client's rate limit alert is sticky — it dismisses only on a "clear" that
-// follows a "limit" — so the server has to push exactly the transitions, and a
-// transition has to be measured per web client rather than off the OSCAR rate
-// state. That state belongs to the account, not to this session: a second
-// browser tab holds its own aimsid and its own WebAPISession over the same
-// state.Session, and an OSCAR client's ObserveRateChanges ticker runs against it
-// too. Any of them can consume a limited -> clear transition before this client
-// hears about it, which would leave this client's alert stuck up forever.
-func (s *WebAPISession) ObserveRateAlert(status wire.RateLimitStatus) bool {
-	s.rateAlertMu.Lock()
-	defer s.rateAlertMu.Unlock()
-
-	prev := s.rateAlertStatus
-	if prev == 0 {
-		// A client that has been told nothing renders no alert, which is what
-		// clear means.
-		prev = wire.RateLimitStatusClear
-	}
-	if prev == status {
-		return false
-	}
-
-	s.rateAlertStatus = status
-	return true
-}
-
 // Aliases returns this session owner's private buddy aliases, keyed by normalized
 // screen name. Aliases live in the owner's feedbag, so the map is loaded once and
 // cached until a feedbag change invalidates it: a signon that brings a large buddy
@@ -276,14 +247,21 @@ func (s *WebAPISession) handleSNACMessage(msg wire.SNACMessage) {
 }
 
 // handleOServiceMessage handles OService SNAC messages relayed to the session's
-// own OSCAR instance. The only one we surface is OServiceUserInfoUpdate, which the
-// server relays to a user when their own user info changes (notably a buddy icon
-// upload or clear). The client re-renders its identity badge from myInfo events
-// only, so we translate this into a fresh myInfo.
+// own OSCAR instance.
 func (s *WebAPISession) handleOServiceMessage(msg wire.SNACMessage) {
-	if msg.Frame.SubGroup != wire.OServiceUserInfoUpdate {
-		return
+	switch msg.Frame.SubGroup {
+	case wire.OServiceUserInfoUpdate:
+		s.handleUserInfoUpdate(msg)
+	case wire.OServiceRateParamChange:
+		s.handleRateLimitUpdate(msg)
 	}
+}
+
+// handleUserInfoUpdate surfaces OServiceUserInfoUpdate, which the server relays to
+// a user when their own user info changes (notably a buddy icon upload or clear).
+// The client re-renders its identity badge from myInfo events only, so we
+// translate this into a fresh myInfo.
+func (s *WebAPISession) handleUserInfoUpdate(msg wire.SNACMessage) {
 	if !s.IsSubscribedTo("myInfo") && !s.IsSubscribedTo("presence") {
 		return
 	}
@@ -298,6 +276,42 @@ func (s *WebAPISession) handleOServiceMessage(msg wire.SNACMessage) {
 	s.EventQueue.Push(types.EventType("myInfo"), data)
 }
 
+// handleRateLimitUpdate translates a rate limit status change — broadcast by the
+// account's rate limit monitor — into a rateLimit event. Only the IM class is
+// surfaced, since the client feeds any rateLimit event into the
+// conversation-window alert. Code 1 is a class-params change, not a status
+// transition, and is ignored.
+func (s *WebAPISession) handleRateLimitUpdate(msg wire.SNACMessage) {
+	if s.IMRateClassID == 0 {
+		return
+	}
+	body, ok := msg.Body.(wire.SNAC_0x01_0x0A_OServiceRateParamsChange)
+	if !ok {
+		return
+	}
+	if wire.RateLimitClassID(body.Rate.ID) != s.IMRateClassID {
+		return
+	}
+
+	var status string
+	switch body.Code {
+	case 2:
+		status = "warn"
+	case 3:
+		status = "limit"
+	case 4:
+		status = "clear"
+	default:
+		return
+	}
+
+	s.EventQueue.Push(types.EventTypeRateLimit, types.RateLimitEvent{
+		Classes: []types.RateLimitClass{
+			{ID: int(body.Rate.ID), Status: status},
+		},
+	})
+}
+
 // handleICBMMessage handles ICBM (instant messaging) SNAC messages.
 func (s *WebAPISession) handleICBMMessage(msg wire.SNACMessage) {
 	switch msg.Frame.SubGroup {

+ 40 - 23
state/webapi_session_test.go

@@ -291,35 +291,52 @@ func TestWebAPISessionManager_CreateAfterShutdown(t *testing.T) {
 	assert.ErrorIs(t, err, ErrWebAPISessionManagerClosed)
 }
 
-// The rate limit alert is per web client, so the transition that drives it has
-// to be tracked per web client too: the OSCAR rate state these sessions sit on
-// is shared by the whole account, and whoever reads it first would otherwise
-// consume the transition the other one is waiting for.
-func TestWebAPISession_ObserveRateAlert(t *testing.T) {
-	t.Run("a client that has been told nothing is already clear", func(t *testing.T) {
-		sess := &WebAPISession{}
-		assert.False(t, sess.ObserveRateAlert(wire.RateLimitStatusClear))
-	})
+// A broadcast rate limit SNAC surfaces to the client only for the IM class: the
+// web client renders any rateLimit event as the conversation-window alert. Code 1
+// (a class-params change) is not a status transition and is dropped.
+func TestWebAPISession_handleRateLimitUpdate(t *testing.T) {
+	const imClass = wire.RateLimitClassID(3)
 
-	t.Run("only changes are reported", func(t *testing.T) {
-		sess := &WebAPISession{}
+	newSession := func() *WebAPISession {
+		return &WebAPISession{
+			IMRateClassID: imClass,
+			EventQueue:    types.NewEventQueue(10),
+			logger:        slog.New(slog.NewTextHandler(io.Discard, nil)),
+		}
+	}
 
-		assert.True(t, sess.ObserveRateAlert(wire.RateLimitStatusLimited))
-		assert.False(t, sess.ObserveRateAlert(wire.RateLimitStatusLimited))
-		assert.True(t, sess.ObserveRateAlert(wire.RateLimitStatusClear))
-		assert.False(t, sess.ObserveRateAlert(wire.RateLimitStatusClear))
+	rateSNAC := func(classID uint16, code uint16) wire.SNACMessage {
+		return wire.SNACMessage{
+			Frame: wire.SNACFrame{FoodGroup: wire.OService, SubGroup: wire.OServiceRateParamChange},
+			Body:  wire.SNAC_0x01_0x0A_OServiceRateParamsChange{Code: code, Rate: wire.RateParamsSNAC{ID: classID}},
+		}
+	}
+
+	t.Run("IM-class transitions become rateLimit events", func(t *testing.T) {
+		sess := newSession()
+		sess.handleSNACMessage(rateSNAC(uint16(imClass), 3)) // limited
+		sess.handleSNACMessage(rateSNAC(uint16(imClass), 4)) // clear
+
+		events := sess.EventQueue.GetAllEvents()
+		require.Len(t, events, 2)
+		assert.Equal(t, "limit", events[0].Data.(types.RateLimitEvent).Classes[0].Status)
+		assert.Equal(t, "clear", events[1].Data.(types.RateLimitEvent).Classes[0].Status)
 	})
 
-	t.Run("two sessions on one account track their alerts independently", func(t *testing.T) {
-		tabA, tabB := &WebAPISession{}, &WebAPISession{}
+	t.Run("other classes and non-status codes are ignored", func(t *testing.T) {
+		sess := newSession()
+		sess.handleSNACMessage(rateSNAC(1, 3))               // class 1 limited: not the IM class
+		sess.handleSNACMessage(rateSNAC(uint16(imClass), 1)) // IM class param change, not a status
+
+		assert.Empty(t, sess.EventQueue.GetAllEvents())
+	})
 
-		// Only tab A is told it is limited.
-		assert.True(t, tabA.ObserveRateAlert(wire.RateLimitStatusLimited))
+	t.Run("a session with no IM class disables the alert", func(t *testing.T) {
+		sess := newSession()
+		sess.IMRateClassID = 0
+		sess.handleSNACMessage(rateSNAC(uint16(imClass), 3))
 
-		// The account recovers. Tab B has no alert up, so the clear is not news
-		// to it — and consuming it must not cost tab A its own clear.
-		assert.False(t, tabB.ObserveRateAlert(wire.RateLimitStatusClear))
-		assert.True(t, tabA.ObserveRateAlert(wire.RateLimitStatusClear))
+		assert.Empty(t, sess.EventQueue.GetAllEvents())
 	})
 }