105 lines
3.9 KiB
Go
105 lines
3.9 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"io"
|
|
"log"
|
|
"testing"
|
|
"time"
|
|
|
|
amqp "github.com/rabbitmq/amqp091-go"
|
|
)
|
|
|
|
type saidTo struct{ acked, nacked, requeued bool }
|
|
|
|
func (a *saidTo) Ack(uint64, bool) error { a.acked = true; return nil }
|
|
func (a *saidTo) Nack(_ uint64, _ bool, requeue bool) error {
|
|
a.nacked, a.requeued = true, requeue
|
|
return nil
|
|
}
|
|
func (a *saidTo) Reject(uint64, bool) error { return nil }
|
|
func (a *saidTo) unsettled() bool { return !a.acked && !a.nacked }
|
|
|
|
type heardWith struct{ err error }
|
|
|
|
func (h heardWith) Heard(context.Context, Report) error { return h.err }
|
|
|
|
// switchable answers with whatever it is set to — the store away, then back.
|
|
type switchable struct{ err error }
|
|
|
|
func (h *switchable) Heard(context.Context, Report) error { return h.err }
|
|
|
|
var tag uint64
|
|
|
|
func aReport(t *testing.T, to *saidTo, node, declared string) amqp.Delivery {
|
|
t.Helper()
|
|
body, err := json.Marshal(Report{Node: node, Declared: declared, Applied: []string{"store"}})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
tag++
|
|
return amqp.Delivery{Acknowledger: to, RoutingKey: KeyReport, Body: body, DeliveryTag: tag}
|
|
}
|
|
|
|
func quiet() *log.Logger { return log.New(io.Discard, "", 0) }
|
|
|
|
// A report the store could not take right now is held, unsettled, and recorded when the store is
|
|
// back; one the store answered no to is acknowledged; one recorded is acknowledged (issue 082, 083).
|
|
func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) {
|
|
store := &switchable{err: errors.Join(ErrTryAgain, errors.New("starting up"))}
|
|
s := &Server{listener: store, log: quiet()}
|
|
held := &saidTo{}
|
|
s.handleReport(context.Background(), aReport(t, held, "anchor", "d1"))
|
|
if !held.unsettled() || len(s.parked) != 1 {
|
|
t.Fatalf("a report the store could not take was not held: %+v, %d held", held, len(s.parked))
|
|
}
|
|
store.err = nil
|
|
s.retryHeld(context.Background())
|
|
if !held.acked || len(s.parked) != 0 {
|
|
t.Fatalf("a held report was not recorded once the store was back: %+v, %d held", held, len(s.parked))
|
|
}
|
|
|
|
refused := &saidTo{}
|
|
s = &Server{listener: heardWith{err: errors.New("a report named no node")}, log: quiet()}
|
|
s.handleReport(context.Background(), aReport(t, refused, "anchor", "d1"))
|
|
if !refused.acked || refused.nacked {
|
|
t.Fatalf("a report the store answered no to was not acknowledged: %+v", refused)
|
|
}
|
|
}
|
|
|
|
// A newer report from the same node supersedes one of its reports still held: recorded after the
|
|
// newer, the older would overwrite what the node is doing now.
|
|
func TestANewerReportSupersedesAHeldOneFromTheSameNode(t *testing.T) {
|
|
s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet()}
|
|
older, newer, other := &saidTo{}, &saidTo{}, &saidTo{}
|
|
s.handleReport(context.Background(), aReport(t, older, "anchor", "d1"))
|
|
s.handleReport(context.Background(), aReport(t, other, "laptop", "d7"))
|
|
s.handleReport(context.Background(), aReport(t, newer, "anchor", "d2"))
|
|
if !older.acked {
|
|
t.Fatalf("the older report was not set aside by the newer: %+v", older)
|
|
}
|
|
if !newer.unsettled() || !other.unsettled() || len(s.parked) != 2 {
|
|
t.Fatalf("the newer report and another node's were not both held: newer %+v other %+v, %d held",
|
|
newer, other, len(s.parked))
|
|
}
|
|
}
|
|
|
|
// A store that has not come back within the bound is not restarting: the report is let go, loudly,
|
|
// rather than held for ever.
|
|
func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
|
|
s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("connection refused"))},
|
|
log: quiet(), giveUp: time.Millisecond}
|
|
held := &saidTo{}
|
|
s.handleReport(context.Background(), aReport(t, held, "anchor", "d1"))
|
|
if !held.unsettled() {
|
|
t.Fatalf("the first failure was not held: %+v", held)
|
|
}
|
|
time.Sleep(5 * time.Millisecond)
|
|
s.retryHeld(context.Background())
|
|
if !held.acked || len(s.parked) != 0 {
|
|
t.Fatalf("a report past the bound was not let go: %+v, %d held", held, len(s.parked))
|
|
}
|
|
}
|