Skip to content

Commit 415147d

Browse files
committed
observability: add compression/fragmentation samples
1 parent 9ceb85f commit 415147d

7 files changed

Lines changed: 312 additions & 21 deletions

File tree

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
package messaging
2+
3+
import (
4+
crand "crypto/rand"
5+
"math/rand"
6+
"sort"
7+
"time"
8+
9+
"ergo.services/ergo/act"
10+
"ergo.services/ergo/gen"
11+
)
12+
13+
type bulkSender struct {
14+
act.Actor
15+
16+
target int
17+
}
18+
19+
func factoryBulkSender() gen.ProcessBehavior {
20+
return &bulkSender{}
21+
}
22+
23+
func (s *bulkSender) Init(args ...any) error {
24+
s.Log().Info("bulk sender started on %s", s.Node().Name())
25+
s.SetCompression(true)
26+
s.SendAfter(s.PID(), messageBulkBurst{}, startDelay+3*time.Second)
27+
return nil
28+
}
29+
30+
func (s *bulkSender) HandleMessage(from gen.PID, message any) error {
31+
switch message.(type) {
32+
case messageBulkBurst:
33+
s.doBurst()
34+
wait := time.Duration(3+rand.Intn(15)) * time.Second
35+
s.SendAfter(s.PID(), messageBulkBurst{}, wait)
36+
}
37+
return nil
38+
}
39+
40+
func (s *bulkSender) doBurst() {
41+
registrar, err := s.Node().Network().Registrar()
42+
if err != nil {
43+
return
44+
}
45+
46+
routes, err := registrar.Resolver().ResolveApplication(appName)
47+
if err != nil {
48+
return
49+
}
50+
51+
myName := s.Node().Name()
52+
remotes := make([]gen.Atom, 0, len(routes))
53+
for _, route := range routes {
54+
if route.Node == myName {
55+
continue
56+
}
57+
remotes = append(remotes, route.Node)
58+
}
59+
60+
if len(remotes) == 0 {
61+
return
62+
}
63+
64+
sort.Slice(remotes, func(i, j int) bool {
65+
return string(remotes[i]) < string(remotes[j])
66+
})
67+
68+
node := remotes[s.target%len(remotes)]
69+
s.target++
70+
71+
count := 5 + rand.Intn(20) // 5..24 large messages
72+
to := gen.ProcessID{Name: poolName, Node: node}
73+
for i := 0; i < count; i++ {
74+
// 100KB..300KB -- large enough to trigger both compression and fragmentation
75+
length := 100000 + rand.Intn(200001)
76+
data := make([]byte, length)
77+
crand.Read(data)
78+
msg := MessageBulkPayload{Data: data}
79+
if err := s.Send(to, msg); err != nil {
80+
s.Log().Warning("bulk sender: send error at %d: %s", i, err)
81+
return
82+
}
83+
}
84+
85+
s.Log().Debug("bulk sender: burst %d large messages -> %s", count, node)
86+
}
87+
88+
func (s *bulkSender) Terminate(reason error) {
89+
s.Log().Info("bulk sender terminated: %s", reason)
90+
}

observability/apps/messaging/sup.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,10 @@ func (s *messagingSup) Init(args ...any) (act.SupervisorSpec, error) {
2828
Name: "messaging_sender",
2929
Factory: factorySender,
3030
},
31+
{
32+
Name: "messaging_bulk_sender",
33+
Factory: factoryBulkSender,
34+
},
3135
},
3236
}, nil
3337
}

observability/apps/messaging/types.go

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ type TestOrder struct {
2525

2626
func init() {
2727
edf.RegisterTypeOf(MessagePayload{})
28+
edf.RegisterTypeOf(MessageBulkPayload{})
2829
edf.RegisterTypeOf(OrderSide(""))
2930
edf.RegisterTypeOf(OrderTag{})
3031
edf.RegisterTypeOf(TestOrder{})
@@ -37,3 +38,11 @@ type MessagePayload struct {
3738

3839
// messageBurst is an internal trigger for the sender to fire a burst
3940
type messageBurst struct{}
41+
42+
// MessageBulkPayload is sent by bulk sender -- large enough to trigger fragmentation
43+
type MessageBulkPayload struct {
44+
Data []byte
45+
}
46+
47+
// messageBulkBurst is an internal trigger for the bulk sender
48+
type messageBulkBurst struct{}

observability/apps/messaging/worker.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@ func (w *worker) HandleMessage(from gen.PID, message any) error {
2121
switch m := message.(type) {
2222
case MessagePayload:
2323
w.Log().Debug("received payload %d bytes from %s", len(m.Data), from)
24+
case MessageBulkPayload:
25+
w.Log().Debug("received bulk payload %d bytes from %s", len(m.Data), from)
2426
}
2527
return nil
2628
}

observability/go.mod

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,17 +3,17 @@ module observability
33
go 1.24.0
44

55
require (
6-
ergo.services/application/mcp v0.0.0-20260305225126-d8948924cf56
6+
ergo.services/application/mcp v0.0.0-20260310144611-bcf83e97036a
77
ergo.services/application/observer v0.1.0
8-
ergo.services/application/radar v0.0.0-20260305212329-10430de2a5c3
9-
ergo.services/ergo v1.999.321-0.20260305211829-909f6f11d916
8+
ergo.services/application/radar v0.0.0-20260310144611-bcf83e97036a
9+
ergo.services/ergo v1.999.321-0.20260310144222-83254b31cb81
1010
ergo.services/logger/colored v0.1.0
1111
ergo.services/registrar/etcd v0.1.0
1212
)
1313

1414
require (
15-
ergo.services/actor/health v0.0.0-20260305212201-8634a257254b // indirect
16-
ergo.services/actor/metrics v0.2.2-0.20260305212201-8634a257254b // indirect
15+
ergo.services/actor/health v0.0.0-20260310144405-e07a5309c0ae // indirect
16+
ergo.services/actor/metrics v0.2.2-0.20260310144405-e07a5309c0ae // indirect
1717
ergo.services/meta/websocket v0.1.0 // indirect
1818
github.com/beorn7/perks v1.0.1 // indirect
1919
github.com/cespare/xxhash/v2 v2.3.0 // indirect

observability/go.sum

Lines changed: 10 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,15 @@
1-
ergo.services/actor/health v0.0.0-20260305212201-8634a257254b h1:/nBphGD4sVSTfz0DSccr/XumvXY+iUYK8BfOIMQkmEk=
2-
ergo.services/actor/health v0.0.0-20260305212201-8634a257254b/go.mod h1:kZpT6IDcGic/6LK1v3ntLK5REcLS0b3wX8gEf1PBXPo=
3-
ergo.services/actor/metrics v0.2.2-0.20260305212201-8634a257254b h1:bT9dv3gMiFCtVG6di8w1tWjx5S2fatnwuAGl/ADlFks=
4-
ergo.services/actor/metrics v0.2.2-0.20260305212201-8634a257254b/go.mod h1:ly0cXPlL5KfI7+UV+fS1QA82/rKRnTv7N0GJWlDgGcg=
5-
ergo.services/application/mcp v0.0.0-20260305225126-d8948924cf56 h1:sBaLjUDP/CJt0LiOw5QJeZXw3+KcPBy61sg3ThhoBtU=
6-
ergo.services/application/mcp v0.0.0-20260305225126-d8948924cf56/go.mod h1:mGRA9WHXqS1B1hJMod+09bCZmjiEIUNGBW0l4qGW0SU=
1+
ergo.services/actor/health v0.0.0-20260310144405-e07a5309c0ae h1:bCp6CWmvy6zG8gYoCiD9jD2SVHSCDBFaVDYB/RqiR4I=
2+
ergo.services/actor/health v0.0.0-20260310144405-e07a5309c0ae/go.mod h1:F9FgeaY9cLRsSFMmGpr/bdJExrgDIhopKFtEhRX2AHM=
3+
ergo.services/actor/metrics v0.2.2-0.20260310144405-e07a5309c0ae h1:YJ3/UOrxQc1WzEeeTs5a4G2B1245hLYLYa4KO0kumpw=
4+
ergo.services/actor/metrics v0.2.2-0.20260310144405-e07a5309c0ae/go.mod h1:TnNbDtcyOMkR8vOECjSeIbVScFlkallRDgfyjiv0Wwo=
5+
ergo.services/application/mcp v0.0.0-20260310144611-bcf83e97036a h1:VtBIdu0LSaZ+DpJ6G5cjt81hFPG2LgHwsVVdc6SOBUI=
6+
ergo.services/application/mcp v0.0.0-20260310144611-bcf83e97036a/go.mod h1:mnJIT+WDSoRazS1H71Za0PygxJWWlKZS9NzGtKjCm80=
77
ergo.services/application/observer v0.1.0 h1:igJgDMgHNkY6BFBgrcsyEeJinZfRI1FRMGJcPy0wsv0=
88
ergo.services/application/observer v0.1.0/go.mod h1:EK27V3/Ts/mN4RQanOuYJCZkrSqL4QHDvg1v2BEOd/g=
9-
ergo.services/application/radar v0.0.0-20260305212329-10430de2a5c3 h1:IRMQcDelExhLHpDLaTrcfb3klHOLEe0HGkP41DzxLcI=
10-
ergo.services/application/radar v0.0.0-20260305212329-10430de2a5c3/go.mod h1:t9aufv/9Lk5ox7z+kQeLG3UU2A8ELzVmeKG9MeVNkWg=
11-
ergo.services/ergo v1.999.321-0.20260305211829-909f6f11d916 h1:RjQnfp7dsOeyIQvb7zNAMj7+0RbMuXUF1oPil9+Lmrg=
12-
ergo.services/ergo v1.999.321-0.20260305211829-909f6f11d916/go.mod h1:bLQ6PoO6Mz/8gVuzvPv3xfMfo1P9w6rZV1WnMXMeMdg=
9+
ergo.services/application/radar v0.0.0-20260310144611-bcf83e97036a h1:3rttIqyNiEORWVgge4w/Tg8Dn07z+hznGn6SnO/d8EE=
10+
ergo.services/application/radar v0.0.0-20260310144611-bcf83e97036a/go.mod h1:8aeiEkyEdWqZFbLZLHNYJ+XCjgi8SejXoVRQbbKCPYc=
11+
ergo.services/ergo v1.999.321-0.20260310144222-83254b31cb81 h1:Il3icnygAhhjYLrj5L9VC21d/4Dj1Xdvg4HIzF2ussc=
12+
ergo.services/ergo v1.999.321-0.20260310144222-83254b31cb81/go.mod h1:bLQ6PoO6Mz/8gVuzvPv3xfMfo1P9w6rZV1WnMXMeMdg=
1313
ergo.services/logger/colored v0.1.0 h1:jbibOaIVZnL+mUsEeyXzzjMaNFsNDcTd+8wdL6cPwu8=
1414
ergo.services/logger/colored v0.1.0/go.mod h1:OEqUiNzSrn3EMKGQuilmKWL0+DEx4Lts8QIkk5lbQoM=
1515
ergo.services/meta/websocket v0.1.0 h1:GXP3X8lLKOPBBHuqZP0Cg1zgg2Az0TLaEWjGbGihi/M=

0 commit comments

Comments
 (0)