package link import ( "context" "encoding/json" "io" "log" "testing" "time" amqp "github.com/rabbitmq/amqp091-go" ) // The harness for the consume side on the bus the mesh runs on today. // // Messages arrive through the seam, so what these tests exercise is the controller's decision // about a message and this transport's way of keeping one — which is what the seam separated. A // fake acknowledger stands in for the bus, because what is asserted is how a message was settled // and that needs no server. // settled is how the bus was told to settle one message. type settled struct{ acked, nacked, requeued, rejected bool } func (a *settled) Ack(uint64, bool) error { a.acked = true; return nil } func (a *settled) Nack(_ uint64, _ bool, requeue bool) error { a.nacked, a.requeued = true, requeue return nil } func (a *settled) Reject(uint64, bool) error { a.rejected = true; return nil } // 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 } var tag uint64 func quiet() *log.Logger { return log.New(io.Discard, "", 0) } // 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, *currentInbound) { in := ¤tInbound{held: map[uint64]*holding{}} return &Server{inbound: in, bus: OverCurrent{}, log: quiet()}, in } // sends is one message arriving over this transport, as the controller reads it. func (c *currentInbound) 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 ¤tControl{kind: kind, on: c, delivery: amqp.Delivery{ Acknowledger: to, Body: body, DeliveryTag: tag, }} } // dueNow brings every held message forward, so a test need not wait out the backoff a real store // restart would be given (RedeliverAfter). func (c *currentInbound) dueNow() { for _, h := range c.held { h.due = time.Now().Add(-time.Second) } } // retries hands every held message back to the controller, the way the ticker does. func (c *currentInbound) retries(ctx context.Context, s *Server) { c.dueNow() c.retryHeld(ctx, s.act) }