Files
mesh-controller/internal/link/receive_fake_test.go
T
jochen 8d9d33ae85 Report a provider that keeps failing a consumer in status (hq ADR 0224)
The identity provider failed every consumer for a day and status called the
mesh well (hq issue 179). The controller now follows every provider's
provisioner.failing/recovered, keeps the newest failing word per provider,
machine and consumer (migration 0065), and status, its JSON and node show
name it until it recovers. Every module that receives contributions is
granted the two events, so no manifest can forget them.
2026-10-06 00:13:23 +02:00

97 lines
3.2 KiB
Go

package link
import (
"context"
"encoding/json"
"testing"
"time"
)
// A harness for the controller's decisions about a message, with no bus behind it.
//
// What these tests exercise is the server's verdict — taken, dropped, held for the store, let go
// once the store has been away too long — and the transport's part in that is only to keep a held
// message and hand it back. A fake that does exactly that stands in for the bus, so what is
// asserted is how a message was settled and the decision needs no server.
// settled is how the bus was told to settle one message.
type settled struct{ acked, nacked, requeued, rejected bool }
// unsettled is a message the controller has neither taken nor let go: it is held, and the bus will
// hand it to whatever consumes next if the controller stops.
func (a *settled) unsettled() bool { return !a.acked && !a.nacked && !a.rejected }
// fakeInbound keeps the messages the controller held, the way a stream would.
type fakeInbound struct{ held map[uint64]*fakeControl }
var tag uint64
// serving is a controller with nothing but a way of receiving, ready for a listener, a recorder,
// an upgrader or a replayer to be set on it.
func serving() (*Server, *fakeInbound) {
in := &fakeInbound{held: map[uint64]*fakeControl{}}
return &Server{inbound: in, log: quiet()}, in
}
func (c *fakeInbound) Also(string) error { return nil }
func (c *fakeInbound) Close() {}
func (c *fakeInbound) Receive(context.Context, func(context.Context, Control)) error {
return nil
}
// sends is one message arriving, as the controller reads it.
func (c *fakeInbound) sends(t *testing.T, to *settled, kind string, v any) Control {
t.Helper()
body, err := json.Marshal(v)
if err != nil {
t.Fatal(err)
}
tag++
return &fakeControl{kind: kind, body: body, tag: tag, to: to, on: c}
}
// retries hands every held message back to the controller, the way a stream redelivers.
func (c *fakeInbound) retries(ctx context.Context, s *Server) {
for tag, m := range c.held {
delete(c.held, tag)
m.redelivered = true
s.act(ctx, m)
}
}
type fakeControl struct {
kind string
subject string
body []byte
tag uint64
to *settled
on *fakeInbound
redelivered bool
since time.Time
}
func (m *fakeControl) Kind() string { return m.kind }
func (m *fakeControl) Body() []byte { return m.body }
func (m *fakeControl) Subject() string { return m.subject }
func (m *fakeControl) Redelivered() bool { return m.redelivered }
func (m *fakeControl) About(string) {}
func (m *fakeControl) Answer(context.Context, []byte) error { return nil }
func (m *fakeControl) Took() error { m.to.acked = true; return nil }
func (m *fakeControl) Drop() error { m.to.rejected = true; return nil }
// HeldFor is how long the store has been waited on for this message — zero on a first delivery.
func (m *fakeControl) HeldFor() time.Duration {
if m.since.IsZero() {
return 0
}
return time.Since(m.since)
}
func (m *fakeControl) Hold(time.Duration) error {
if m.since.IsZero() {
m.since = time.Now()
}
m.on.held[m.tag] = m
return nil
}