perf(server): batch achievement updates and eliminate redundant snapshots, writes and JSON round trips
This commit is contained in:
@@ -7,6 +7,14 @@ import (
|
||||
)
|
||||
|
||||
func (s *Service) BeginSession(id string) { s.SetSession(id) }
|
||||
func (s *Service) inventorySnapshot() (world.GameplayAchievementSnapshot, error) {
|
||||
if provider, ok := s.provider.(interface {
|
||||
InventorySnapshot() (world.GameplayAchievementSnapshot, error)
|
||||
}); ok {
|
||||
return provider.InventorySnapshot()
|
||||
}
|
||||
return s.provider.Snapshot()
|
||||
}
|
||||
func (s *Service) BeforeDispatch(string, []byte) error {
|
||||
s.mu.Lock()
|
||||
s.beforeMissions = s.visibleMissionValues()
|
||||
@@ -15,13 +23,13 @@ func (s *Service) BeforeDispatch(string, []byte) error {
|
||||
return nil
|
||||
}
|
||||
var e error
|
||||
s.before, e = s.provider.Snapshot()
|
||||
s.before, e = s.inventorySnapshot()
|
||||
return e
|
||||
}
|
||||
func (s *Service) AttachGameplayProvider(p world.GameplayAchievementProvider) { s.provider = p }
|
||||
func (s *Service) AfterDispatch(path string, request, response []byte) ([]byte, error) {
|
||||
if s.provider != nil {
|
||||
after, e := s.provider.Snapshot()
|
||||
after, e := s.inventorySnapshot()
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
@@ -30,6 +38,8 @@ func (s *Service) AfterDispatch(path string, request, response []byte) ([]byte,
|
||||
_, seen := s.state.Receipts[rk]
|
||||
s.mu.Unlock()
|
||||
if !seen {
|
||||
type delta struct{ condition, sub, count uint64 }
|
||||
var deltas []delta
|
||||
for kind, old := range s.before.Items {
|
||||
current := after.Items[kind]
|
||||
if current < old {
|
||||
@@ -37,41 +47,47 @@ func (s *Service) AfterDispatch(path string, request, response []byte) ([]byte,
|
||||
if kind[1] == 0 {
|
||||
condition = 11
|
||||
}
|
||||
if e = s.RecordEvent(condition, kind[0], old-current, s.unlocked); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
deltas = append(deltas, delta{condition, kind[0], old - current})
|
||||
}
|
||||
}
|
||||
for kind, current := range after.Items {
|
||||
old := s.before.Items[kind]
|
||||
if current > old {
|
||||
if e = s.RecordEvent(32, kind[0], current-old, s.unlocked); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
deltas = append(deltas, delta{32, kind[0], current - old})
|
||||
}
|
||||
}
|
||||
for idx, current := range after.Equipment {
|
||||
old, ok := s.before.Equipment[idx]
|
||||
if ok && current.Level > old.Level {
|
||||
if e = s.RecordEvent(14, 0, current.Level-old.Level, s.unlocked); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
deltas = append(deltas, delta{14, 0, current.Level - old.Level})
|
||||
}
|
||||
}
|
||||
for idx, current := range after.Costumes {
|
||||
old, ok := s.before.Costumes[idx]
|
||||
if ok && current.Level > old.Level {
|
||||
if e = s.RecordEvent(104, current.ID, current.Level-old.Level, s.unlocked); e != nil {
|
||||
return nil, e
|
||||
}
|
||||
deltas = append(deltas, delta{104, current.ID, current.Level - old.Level})
|
||||
}
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.state.Receipts[rk] = receipt{Digest: "observer"}
|
||||
e = s.save()
|
||||
s.mu.Unlock()
|
||||
if e != nil {
|
||||
return nil, e
|
||||
if len(deltas) > 0 {
|
||||
s.mu.Lock()
|
||||
before, e := json.Marshal(s.state)
|
||||
if e != nil {
|
||||
s.mu.Unlock()
|
||||
return nil, e
|
||||
}
|
||||
for _, d := range deltas {
|
||||
s.recordEventLocked(d.condition, d.sub, d.count, s.unlocked)
|
||||
}
|
||||
s.state.Receipts[rk] = receipt{Digest: "observer"}
|
||||
e = s.save()
|
||||
if e != nil {
|
||||
s.state = snapshot{}
|
||||
_ = json.Unmarshal(before, &s.state)
|
||||
}
|
||||
s.mu.Unlock()
|
||||
if e != nil {
|
||||
return nil, e
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,142 @@
|
||||
package eventtasks
|
||||
|
||||
import (
|
||||
"bd2server/internal/server/stateio"
|
||||
"bd2server/internal/server/world"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
type narrowNoticeProvider struct {
|
||||
noticeProvider
|
||||
fullCalls, narrowCalls int
|
||||
}
|
||||
|
||||
func (p *narrowNoticeProvider) Snapshot() (world.GameplayAchievementSnapshot, error) {
|
||||
p.fullCalls++
|
||||
return p.noticeProvider.Snapshot()
|
||||
}
|
||||
func (p *narrowNoticeProvider) InventorySnapshot() (world.GameplayAchievementSnapshot, error) {
|
||||
p.narrowCalls++
|
||||
return p.noticeProvider.Snapshot()
|
||||
}
|
||||
|
||||
type failedObserverStore struct{ stateio.Store }
|
||||
|
||||
func (s failedObserverStore) Save(string, []byte) error { return errors.New("observer save failure") }
|
||||
|
||||
func TestObserverUsesInventoryProjectionAndRestoresMemoryOnSaveFailure(t *testing.T) {
|
||||
s, _, store := setup(t)
|
||||
p := &narrowNoticeProvider{noticeProvider: noticeProvider{count: 1}}
|
||||
s.AttachGameplayProvider(p)
|
||||
task := s.design.Missions[10]
|
||||
task.Type = 32
|
||||
s.design.Missions[10] = task
|
||||
if e := s.BeforeDispatch("/grant", req(1)); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
p.count = 2
|
||||
s.store = failedObserverStore{store}
|
||||
if _, e := s.AfterDispatch("/grant", req(1), nil); e == nil {
|
||||
t.Fatal("observer persistence failure ignored")
|
||||
}
|
||||
if len(s.state.Missions) != 0 || len(s.state.Receipts) != 0 {
|
||||
t.Fatal("failed observer save retained partial tasks or receipt")
|
||||
}
|
||||
if p.fullCalls != 0 || p.narrowCalls != 2 {
|
||||
t.Fatalf("observer requested expensive full snapshot: full=%d narrow=%d", p.fullCalls, p.narrowCalls)
|
||||
}
|
||||
s.store = store
|
||||
p.count = 1
|
||||
s.BeforeDispatch("/grant", req(1))
|
||||
p.count = 2
|
||||
b, e := s.AfterDispatch("/grant", req(1), nil)
|
||||
if e != nil || len(b) == 0 {
|
||||
t.Fatal("retry after failed observer persistence lost progress")
|
||||
}
|
||||
}
|
||||
|
||||
type writeCountingStore struct {
|
||||
stateio.Store
|
||||
writes int
|
||||
bytes int
|
||||
}
|
||||
|
||||
func (s *writeCountingStore) Save(name string, b []byte) error {
|
||||
s.writes++
|
||||
s.bytes += len(b)
|
||||
return s.Store.Save(name, b)
|
||||
}
|
||||
|
||||
func TestReadOnlyBatchHasNoObserverWritesButRealDeltaNotifiesOnce(t *testing.T) {
|
||||
s, _, store := setup(t)
|
||||
counter := &writeCountingStore{Store: store}
|
||||
s.store = counter
|
||||
p := ¬iceProvider{count: 1}
|
||||
s.AttachGameplayProvider(p)
|
||||
start := time.Now()
|
||||
for seq := uint64(1); seq <= 57; seq++ {
|
||||
if e := s.BeforeDispatch("/read", req(seq)); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
b, e := s.AfterDispatch("/read", req(seq), nil)
|
||||
if e != nil || len(b) > 0 {
|
||||
t.Fatalf("read-only changed missions: %x %v", b, e)
|
||||
}
|
||||
}
|
||||
if counter.writes != 0 || len(s.state.Receipts) != 0 {
|
||||
t.Fatalf("read-only 57 packets made %d writes/%d receipts", counter.writes, len(s.state.Receipts))
|
||||
}
|
||||
t.Logf("57 unchanged observer boundaries: %s, writes=%d", time.Since(start), counter.writes)
|
||||
task := s.design.Missions[10]
|
||||
task.Type = 32
|
||||
s.design.Missions[10] = task
|
||||
if e := s.BeforeDispatch("/grant", req(58)); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
p.count = 2
|
||||
b, e := s.AfterDispatch("/grant", req(58), nil)
|
||||
if e != nil || len(b) == 0 || counter.writes != 1 {
|
||||
t.Fatalf("real delta not persisted/notified once: writes=%d body=%x error=%v", counter.writes, b, e)
|
||||
}
|
||||
// Even if an upstream replay temporarily exposes the same before/after
|
||||
// delta, the committed request receipt must not increment tasks twice.
|
||||
p.count = 1
|
||||
s.BeforeDispatch("/grant", req(58))
|
||||
p.count = 2
|
||||
b, e = s.AfterDispatch("/grant", req(58), nil)
|
||||
if e != nil || len(b) != 0 || counter.writes != 1 {
|
||||
t.Fatal("replay repeated mission increment or write")
|
||||
}
|
||||
if e = s.RecordEvent(999999, 0, 1, nil); e != nil || counter.writes != 1 {
|
||||
t.Fatal("irrelevant condition wrote state")
|
||||
}
|
||||
if e = s.RecordEvent(32, 0, 100, nil); e != nil {
|
||||
t.Fatal(e)
|
||||
}
|
||||
writes := counter.writes
|
||||
if e = s.RecordEvent(32, 0, 100, nil); e != nil || counter.writes != writes {
|
||||
t.Fatal("capped task still wrote whole snapshot")
|
||||
}
|
||||
}
|
||||
|
||||
func BenchmarkUnchanged57PacketObserverBatch(b *testing.B) {
|
||||
// Snapshot costs belong to the provider; this benchmark isolates event
|
||||
// mission observation and persistence decisions without a user's database.
|
||||
s, _, _ := setup(b)
|
||||
p := ¬iceProvider{count: 1}
|
||||
s.AttachGameplayProvider(p)
|
||||
b.ReportAllocs()
|
||||
b.ResetTimer()
|
||||
for n := 0; n < b.N; n++ {
|
||||
for seq := uint64(1); seq <= 57; seq++ {
|
||||
if e := s.BeforeDispatch("/read", req(seq)); e != nil {
|
||||
b.Fatal(e)
|
||||
}
|
||||
if _, e := s.AfterDispatch("/read", req(seq), nil); e != nil {
|
||||
b.Fatal(e)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -885,9 +885,20 @@ func (s *Service) handle(path string, b []byte, identity string) ([]byte, error)
|
||||
func (s *Service) RecordEvent(condition, sub, count uint64, unlocked func(uint64, uint64) bool) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if count == 0 {
|
||||
if !s.recordEventLocked(condition, sub, count, unlocked) {
|
||||
return nil
|
||||
}
|
||||
return s.save()
|
||||
}
|
||||
|
||||
// recordEventLocked reports every persistent mutation, including initializing
|
||||
// a matching task or rolling its period. Irrelevant and already capped events
|
||||
// leave the account snapshot untouched.
|
||||
func (s *Service) recordEventLocked(condition, sub, count uint64, unlocked func(uint64, uint64) bool) bool {
|
||||
changed := false
|
||||
if count == 0 {
|
||||
return false
|
||||
}
|
||||
for _, v := range s.taskSchedules() {
|
||||
if v.Type != 4 || !s.active(v) {
|
||||
continue
|
||||
@@ -908,7 +919,15 @@ func (s *Service) RecordEvent(condition, sub, count uint64, unlocked func(uint64
|
||||
if !match {
|
||||
continue
|
||||
}
|
||||
previous := s.state.Missions[scheduleKey(v)+"/"+key(t.ID)]
|
||||
var prior mission
|
||||
if previous != nil {
|
||||
prior = *previous
|
||||
}
|
||||
m := s.mission(v, t.ID)
|
||||
if previous == nil || prior != *m {
|
||||
changed = true
|
||||
}
|
||||
if m.Claimed {
|
||||
continue
|
||||
}
|
||||
@@ -920,9 +939,10 @@ func (s *Service) RecordEvent(condition, sub, count uint64, unlocked func(uint64
|
||||
} else {
|
||||
m.Value += count
|
||||
}
|
||||
changed = true
|
||||
}
|
||||
}
|
||||
return s.save()
|
||||
return changed
|
||||
}
|
||||
func (s *Service) Notify() ([]byte, error) {
|
||||
s.mu.Lock()
|
||||
|
||||
@@ -51,7 +51,7 @@ func (m *attendanceMailStub) IssueAttachmentsOnce(identity, title, body string,
|
||||
m.identity, m.title, m.body, m.sentAt = identity, title, body, sentAt
|
||||
return nil
|
||||
}
|
||||
func setup(t *testing.T) (*Service, *economyStub, stateio.Store) {
|
||||
func setup(t testing.TB) (*Service, *economyStub, stateio.Store) {
|
||||
t.Helper()
|
||||
now := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC)
|
||||
registry := events.NewRegistry()
|
||||
|
||||
Reference in New Issue
Block a user