Nothing the control queue carries is lost while the store restarts (issue 083) #43
@@ -28,13 +28,15 @@ type following struct{ open *stores }
|
|||||||
func (f following) Upgraded(ctx context.Context, u link.Upgraded) error {
|
func (f following) Upgraded(ctx context.Context, u link.Upgraded) error {
|
||||||
inv := f.open.inventory
|
inv := f.open.inventory
|
||||||
|
|
||||||
|
// The store read first, and an outage there said as one, so the announcement is held and asked
|
||||||
|
// again (novox/hq issue 083). Only here: a push that fails further down is not asked again.
|
||||||
decision, err := inv.UpgradeOf(ctx, u.Module)
|
decision, err := inv.UpgradeOf(ctx, u.Module)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return notNow(err)
|
||||||
}
|
}
|
||||||
on, err := inv.Running(ctx, u.Module)
|
on, err := inv.Running(ctx, u.Module)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return notNow(err)
|
||||||
}
|
}
|
||||||
if len(on) == 0 {
|
if len(on) == 0 {
|
||||||
fmt.Printf("%s moved to %s; no machine runs it\n", u.Module, shortCommit(u.Commit))
|
fmt.Printf("%s moved to %s; no machine runs it\n", u.Module, shortCommit(u.Commit))
|
||||||
@@ -194,3 +196,12 @@ func (f following) Announceable(ctx context.Context) ([]link.Announcement, error
|
|||||||
}
|
}
|
||||||
return out, nil
|
return out, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// notNow marks a store that could not be read right now, so the announcement is held rather than
|
||||||
|
// lost; anything else is returned as it was.
|
||||||
|
func notNow(err error) error {
|
||||||
|
if inventory.Unreachable(err) {
|
||||||
|
return fmt.Errorf("%w: %w", link.ErrTryAgain, err)
|
||||||
|
}
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|||||||
@@ -383,3 +383,88 @@ func TestBeingHeardFromDoesNotChangeWhatANodeOwns(t *testing.T) {
|
|||||||
t.Errorf("after a bare word that the node is here, the mesh believes it owns %v", owned)
|
t.Errorf("after a bare word that the node is here, the mesh believes it owns %v", owned)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// **A token is claimed, then spent** (novox/hq issue 083). A presenter may claim it again — an
|
||||||
|
// enrolment interrupted by a restarting store asks again with the same key — while another is held
|
||||||
|
// off until the lease lapses; and it is spent only by the one holding the claim.
|
||||||
|
func TestATokenIsClaimedByOnePresenterAndSpentOnlyByIt(t *testing.T) {
|
||||||
|
inv := fresh(t)
|
||||||
|
ctx := t.Context()
|
||||||
|
if _, err := inv.AddNode(ctx, "laptop"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
issued, err := inv.IssueToken(ctx, "laptop", time.Hour)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
node, err := inv.Claim(ctx, issued.Secret, "key-a", false)
|
||||||
|
if err != nil || node.Name != "laptop" {
|
||||||
|
t.Fatalf("a fresh token was not claimed for its node: %v %q", err, node.Name)
|
||||||
|
}
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-a", false); err != nil {
|
||||||
|
t.Fatalf("the presenter holding the claim could not claim again after an interruption: %v", err)
|
||||||
|
}
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-b", false); !errors.Is(err, ErrTokenInUse) {
|
||||||
|
t.Fatalf("a second presenter was not held off while the claim is live: %v", err)
|
||||||
|
}
|
||||||
|
if err := inv.Spend(ctx, issued.Secret, "key-b"); !errors.Is(err, ErrTokenRefused) {
|
||||||
|
t.Fatalf("a presenter not holding the claim spent the token: %v", err)
|
||||||
|
}
|
||||||
|
if err := inv.Spend(ctx, issued.Secret, "key-a"); err != nil {
|
||||||
|
t.Fatalf("the presenter holding the claim could not spend it: %v", err)
|
||||||
|
}
|
||||||
|
// Spent: nobody else may claim it, ever.
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-b", false); !errors.Is(err, ErrTokenRefused) {
|
||||||
|
t.Fatalf("a spent token was claimed by another presenter: %v", err)
|
||||||
|
}
|
||||||
|
// Nor the presenter that spent it, without proof it holds the key: a public key is no secret.
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-a", false); !errors.Is(err, ErrTokenRefused) {
|
||||||
|
t.Fatalf("a spent token was claimed again with no proof of the key: %v", err)
|
||||||
|
}
|
||||||
|
// With it, and inside the lease, it may, and spend it again: its spend reached the store and
|
||||||
|
// the answer did not reach the node, which asked again — refusing it would lock out a machine
|
||||||
|
// the mesh holds as enrolled.
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-a", true); err != nil {
|
||||||
|
t.Fatalf("the proven presenter whose answer was lost after its spend was refused: %v", err)
|
||||||
|
}
|
||||||
|
if err := inv.Spend(ctx, issued.Secret, "key-a"); err != nil {
|
||||||
|
t.Fatalf("spending again by the same presenter failed: %v", err)
|
||||||
|
}
|
||||||
|
// And not once the lease is over: then a spent token is spent to everyone, proof or not.
|
||||||
|
if _, err := inv.store.Pool().Exec(ctx,
|
||||||
|
`update enrolment_token set claimed_until = now() - interval '1 second' where secret = $1`,
|
||||||
|
hashSecret(issued.Secret)); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-a", true); !errors.Is(err, ErrTokenRefused) {
|
||||||
|
t.Fatalf("a spent token was claimed again after its lease: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A claim lapses: a host that gave up and was started over, with keys of its own, is not held off
|
||||||
|
// for longer than the lease.
|
||||||
|
func TestAClaimThatLapsedCanBeTakenByAnotherPresenter(t *testing.T) {
|
||||||
|
inv := fresh(t)
|
||||||
|
ctx := t.Context()
|
||||||
|
if _, err := inv.AddNode(ctx, "laptop"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
issued, err := inv.IssueToken(ctx, "laptop", time.Hour)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-a", false); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := inv.store.Pool().Exec(ctx,
|
||||||
|
`update enrolment_token set claimed_until = now() - interval '1 second' where secret = $1`,
|
||||||
|
hashSecret(issued.Secret)); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "key-b", false); err != nil {
|
||||||
|
t.Fatalf("a lapsed claim held off a new presenter: %v", err)
|
||||||
|
}
|
||||||
|
if err := inv.Spend(ctx, issued.Secret, "key-a"); !errors.Is(err, ErrTokenRefused) {
|
||||||
|
t.Fatalf("the presenter whose claim lapsed could still spend the token: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,11 @@
|
|||||||
|
-- A token is claimed by the enrolment presenting it, and spent only when that enrolment has
|
||||||
|
-- written everything it needs (novox/hq 04-ISSUES/083).
|
||||||
|
--
|
||||||
|
-- Spending came first and the node's keys after, in another database: a store that went away
|
||||||
|
-- between the two left a spent token and a node with no key, and the host — which makes new keys
|
||||||
|
-- on every attempt — could not try again. The claim holds the token for one presenter for a short
|
||||||
|
-- lease, so an attempt that failed part-way can be made again by the same presenter, and a second
|
||||||
|
-- presenter cannot interleave with the first.
|
||||||
|
|
||||||
|
alter table enrolment_token add column claimed_by text;
|
||||||
|
alter table enrolment_token add column claimed_until timestamptz;
|
||||||
@@ -203,7 +203,75 @@ func (i *Inventory) IssueToken(ctx context.Context, nodeName string, validFor ti
|
|||||||
// token that had expired.
|
// token that had expired.
|
||||||
var ErrTokenRefused = errors.New("that token cannot be used")
|
var ErrTokenRefused = errors.New("that token cannot be used")
|
||||||
|
|
||||||
// Redeem spends a token and reports which node it was for.
|
// ClaimLease is how long a claimed token is held for the one presenter that claimed it. Long
|
||||||
|
// enough for an enrolment to be tried again through a store restart; short enough that a host
|
||||||
|
// which gave up and was started over, with keys of its own, is not kept waiting long.
|
||||||
|
const ClaimLease = 2 * time.Minute
|
||||||
|
|
||||||
|
// ErrTokenInUse is a token another presenter holds a claim on right now. Not a refusal: the claim
|
||||||
|
// lapses, and asking again after it is the answer.
|
||||||
|
var ErrTokenInUse = errors.New("the token is being used by another enrolment")
|
||||||
|
|
||||||
|
// Claim takes a token for one presenter — `by`, which names the key presenting it — for the length
|
||||||
|
// of a lease, and says which node it enrols. The same presenter may claim it again, as may anyone
|
||||||
|
// once the lease has lapsed; nothing is spent until Spend (novox/hq 04-ISSUES/083).
|
||||||
|
//
|
||||||
|
// A token this same presenter already spent may be claimed again when `again` says so — the
|
||||||
|
// caller has proof the presenter holds the key's private half — and only while its claim's lease
|
||||||
|
// is live: its spend reached the store and the answer did not reach the node, which asked again.
|
||||||
|
// Refusing it then would lock out a machine the mesh holds as enrolled. Without the proof a spent
|
||||||
|
// token stays spent to everyone, as ADR 0004 says.
|
||||||
|
func (i *Inventory) Claim(ctx context.Context, secret, by string, again bool) (Node, error) {
|
||||||
|
var id string
|
||||||
|
err := i.store.Pool().QueryRow(ctx,
|
||||||
|
`update enrolment_token set claimed_by = $2,
|
||||||
|
claimed_until = case when redeemed is null then now() + $3::interval else claimed_until end
|
||||||
|
where secret = $1 and expires > now()
|
||||||
|
and ((redeemed is null and (claimed_by is null or claimed_by = $2 or claimed_until < now()))
|
||||||
|
or (redeemed is not null and $4 and claimed_by = $2 and claimed_until > now()))
|
||||||
|
returning node`, hashSecret(secret), by, ClaimLease.String(), again).Scan(&id)
|
||||||
|
if errors.Is(err, pgx.ErrNoRows) {
|
||||||
|
// Unusable, or held by someone else — told apart, because the second passes.
|
||||||
|
var held bool
|
||||||
|
probe := i.store.Pool().QueryRow(ctx,
|
||||||
|
`select true from enrolment_token
|
||||||
|
where secret = $1 and redeemed is null and expires > now()`, hashSecret(secret)).Scan(&held)
|
||||||
|
switch {
|
||||||
|
case probe == nil && held:
|
||||||
|
return Node{}, ErrTokenInUse
|
||||||
|
case probe != nil && !errors.Is(probe, pgx.ErrNoRows):
|
||||||
|
// The store went away between the two questions: "not now", not a refusal.
|
||||||
|
return Node{}, probe
|
||||||
|
}
|
||||||
|
return Node{}, ErrTokenRefused
|
||||||
|
}
|
||||||
|
if err != nil {
|
||||||
|
return Node{}, err
|
||||||
|
}
|
||||||
|
var n Node
|
||||||
|
err = i.store.Pool().QueryRow(ctx,
|
||||||
|
`select id, name, created from node where id = $1`, id).Scan(&n.ID, &n.Name, &n.Created)
|
||||||
|
return n, err
|
||||||
|
}
|
||||||
|
|
||||||
|
// Spend makes a claimed token used, only for the presenter holding the claim. The last write to the
|
||||||
|
// store in an enrolment, so a token is spent exactly when the node it enrolled is complete. Spent
|
||||||
|
// again by the same presenter is not an error: an answer lost after the first spend.
|
||||||
|
func (i *Inventory) Spend(ctx context.Context, secret, by string) error {
|
||||||
|
tag, err := i.store.Pool().Exec(ctx,
|
||||||
|
`update enrolment_token set redeemed = coalesce(redeemed, now())
|
||||||
|
where secret = $1 and claimed_by = $2`, hashSecret(secret), by)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if tag.RowsAffected() == 0 {
|
||||||
|
return ErrTokenRefused
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Redeem spends a token in one step and reports which node it was for. Enrolment claims and then
|
||||||
|
// spends (Claim, Spend); this is the one-step form, kept for what spends a token outright.
|
||||||
//
|
//
|
||||||
// It does not issue an identity. What a node presents afterwards to prove it is that node is not
|
// It does not issue an identity. What a node presents afterwards to prove it is that node is not
|
||||||
// decided anywhere (novox/hq ADR 0004 names the property, not the mechanism), and guessing at it
|
// decided anywhere (novox/hq ADR 0004 names the property, not the mechanism), and guessing at it
|
||||||
|
|||||||
@@ -0,0 +1,113 @@
|
|||||||
|
package link_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"crypto/ed25519"
|
||||||
|
"crypto/rand"
|
||||||
|
"crypto/sha256"
|
||||||
|
"encoding/hex"
|
||||||
|
"errors"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
|
"github.com/novox/mesh-controller/internal/link"
|
||||||
|
)
|
||||||
|
|
||||||
|
// A token another enrolment holds for the moment is "not now", not a refusal: the node asks again
|
||||||
|
// with the same request, and nothing is spent (novox/hq issue 083).
|
||||||
|
func TestAnEnrolmentMetByAHeldTokenIsAskedToTryAgain(t *testing.T) {
|
||||||
|
inv := inventory.ForTest(t)
|
||||||
|
ctx := t.Context()
|
||||||
|
if _, err := inv.AddNode(ctx, "laptop"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
issued, err := inv.IssueToken(ctx, "laptop", time.Hour)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, "another enrolment's key", false); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
public, _, err := ed25519.GenerateKey(rand.Reader)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
_, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{
|
||||||
|
Node: "laptop", Secret: issued.Secret, PublicKey: public})
|
||||||
|
if !errors.Is(err, link.ErrTryAgain) {
|
||||||
|
t.Fatalf("an enrolment met by a held token was not asked to try again: %v", err)
|
||||||
|
}
|
||||||
|
// And a token that cannot be used at all is still refused outright.
|
||||||
|
_, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{
|
||||||
|
Node: "laptop", Secret: "not-a-token", PublicKey: public})
|
||||||
|
if err == nil || errors.Is(err, link.ErrTryAgain) {
|
||||||
|
t.Fatalf("a token that cannot be used was not refused outright: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// **A spent token cannot be replayed with a node's public key** (novox/hq issue 083, on review). A
|
||||||
|
// public key is no secret: finishing an enrolment whose token that key spent takes proof the
|
||||||
|
// presenter holds its private half, and a proof made with another key — or for another request —
|
||||||
|
// is refused outright.
|
||||||
|
func TestASpentTokenCannotBeReplayedWithAPublicKeyAlone(t *testing.T) {
|
||||||
|
inv := inventory.ForTest(t)
|
||||||
|
ctx := t.Context()
|
||||||
|
if _, err := inv.AddNode(ctx, "laptop"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
issued, err := inv.IssueToken(ctx, "laptop", time.Hour)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
victim, victimPrivate, err := ed25519.GenerateKey(rand.Reader)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
// The victim's enrolment spent the token.
|
||||||
|
by := claimantOf(victim)
|
||||||
|
if _, err := inv.Claim(ctx, issued.Secret, by, false); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := inv.Spend(ctx, issued.Secret, by); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Someone with the leaked token and the victim's public key, and no proof.
|
||||||
|
_, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{
|
||||||
|
Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the attacker's"})
|
||||||
|
if err == nil || errors.Is(err, link.ErrTryAgain) {
|
||||||
|
t.Fatalf("a spent token was taken again with a public key alone: %v", err)
|
||||||
|
}
|
||||||
|
// With a proof made by another key.
|
||||||
|
_, forger, err := ed25519.GenerateKey(rand.Reader)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
forged := ed25519.Sign(forger, link.EnrolProof(issued.Secret, victim, "", "the attacker's", ""))
|
||||||
|
_, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{
|
||||||
|
Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the attacker's", Proof: forged})
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("a spent token was taken again with a proof made by another key")
|
||||||
|
}
|
||||||
|
// With the victim's own proof, but for another request — keys swapped in.
|
||||||
|
theirs := ed25519.Sign(victimPrivate, link.EnrolProof(issued.Secret, victim, "", "the victim's", ""))
|
||||||
|
_, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{
|
||||||
|
Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the attacker's", Proof: theirs})
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("a spent token was taken again with a proof made for another request")
|
||||||
|
}
|
||||||
|
// And the victim's own, proven request is not honoured when the broker redelivered it: the
|
||||||
|
// first delivery may already have been answered.
|
||||||
|
_, err = link.Enrolment{Inventory: inv}.Enrol(ctx, link.EnrolRequest{
|
||||||
|
Node: "laptop", Secret: issued.Secret, PublicKey: victim, SealingKey: "the victim's",
|
||||||
|
Proof: theirs, Redelivered: true})
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("a redelivered request finished an enrolment already spent")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// claimantOf is the claimant the enrolment derives from a key, recomputed here.
|
||||||
|
func claimantOf(public ed25519.PublicKey) string {
|
||||||
|
sum := sha256.Sum256(public)
|
||||||
|
return hex.EncodeToString(sum[:])
|
||||||
|
}
|
||||||
@@ -0,0 +1,13 @@
|
|||||||
|
package link
|
||||||
|
|
||||||
|
import "testing"
|
||||||
|
|
||||||
|
// The bytes a node signs when it enrols, built here to check its signature. The host builds the
|
||||||
|
// same bytes in another repository; this known answer is repeated in its test, so the two cannot
|
||||||
|
// drift apart without one of them failing (novox/hq issue 083).
|
||||||
|
func TestWhatAnEnrollingNodeSignsIsFixed(t *testing.T) {
|
||||||
|
got := string(EnrolProof("s", []byte{1, 2, 3}, "o", "e", "v"))
|
||||||
|
if want := "novox-mesh-enrol\x00s\x00AQID\x00o\x00e\x00v"; got != want {
|
||||||
|
t.Fatalf("the enrolment proof's message changed: %q, want %q", got, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
+74
-35
@@ -4,7 +4,9 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"crypto/ed25519"
|
"crypto/ed25519"
|
||||||
"crypto/rand"
|
"crypto/rand"
|
||||||
|
"crypto/sha256"
|
||||||
"encoding/base64"
|
"encoding/base64"
|
||||||
|
"encoding/hex"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
@@ -28,32 +30,57 @@ type Enrolment struct {
|
|||||||
Broker broker.Broker
|
Broker broker.Broker
|
||||||
}
|
}
|
||||||
|
|
||||||
// Enrol spends the token and records what the node presented.
|
// Enrol records what the node presented and spends the token.
|
||||||
//
|
//
|
||||||
// Order matters and it is the order things become irreversible. The token is spent first, in a
|
// Order matters and it is the order things become irreversible (novox/hq issue 083). The token is
|
||||||
// single statement that both finds and marks it, so two machines racing on one secret produce one
|
// claimed first, in a single statement that both finds it and holds it for this presenter's key, so
|
||||||
// winner. Only then is a key recorded — because recording a key for a node whose token turned out
|
// two machines racing on one secret produce one holder. Then everything the node presented is
|
||||||
// to be spent would leave the mesh believing a machine that never had the right to join.
|
// written — each write an overwrite, so an attempt interrupted by the store going away can be made
|
||||||
func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply, error) {
|
// again by the same presenter. Then the token is spent. Last, the token's secret stops being the
|
||||||
|
// node's broker password: done after the spend, because a password replaced by an attempt that
|
||||||
|
// then failed would be one nobody holds, and the node could not even log in to ask again.
|
||||||
|
func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (reply EnrolReply, err error) {
|
||||||
secret, public, profile := request.Secret, ed25519.PublicKey(request.PublicKey), request.Profile
|
secret, public, profile := request.Secret, ed25519.PublicKey(request.PublicKey), request.Profile
|
||||||
|
|
||||||
|
// A store that could not be asked right now, or a token another presenter holds for the
|
||||||
|
// moment, is "not now": the node asks again with the same request (novox/hq issue 083).
|
||||||
|
defer func() {
|
||||||
|
if inventory.Unreachable(err) || errors.Is(err, inventory.ErrTokenInUse) ||
|
||||||
|
errors.Is(err, context.Canceled) {
|
||||||
|
err = fmt.Errorf("%w: %w", ErrTryAgain, err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
if len(public) != ed25519.PublicKeySize {
|
if len(public) != ed25519.PublicKeySize {
|
||||||
return EnrolReply{}, fmt.Errorf("a node presented a %d-byte key, and an identity is %d",
|
return EnrolReply{}, fmt.Errorf("a node presented a %d-byte key, and an identity is %d",
|
||||||
len(public), ed25519.PublicKeySize)
|
len(public), ed25519.PublicKeySize)
|
||||||
}
|
}
|
||||||
|
|
||||||
node, err := e.Inventory.Redeem(ctx, secret)
|
// Claimed, not spent: the token is held for this presenter while the node is written, and
|
||||||
|
// spent only as the last write. The node's identity lives in another database than the
|
||||||
|
// token, so the two cannot be one transaction; a failure between them used to leave a spent
|
||||||
|
// token and a node with no key, which the host — making new keys on every attempt — could not
|
||||||
|
// recover from. Every write below overwrites, so an attempt made again is safe.
|
||||||
|
// A proof that does not verify is refused outright: it was made with another key, or for
|
||||||
|
// another request. One that verifies lets this presenter finish an enrolment whose token it
|
||||||
|
// already spent — never a request the broker handed over a second time, which may already
|
||||||
|
// have been answered.
|
||||||
|
proven := false
|
||||||
|
if len(request.Proof) > 0 {
|
||||||
|
if !ed25519.Verify(public, EnrolProof(secret, public, request.OverlayKey, request.SealingKey,
|
||||||
|
request.ServingKey), request.Proof) {
|
||||||
|
return EnrolReply{}, errors.New("the enrolment's proof does not match the key it presents")
|
||||||
|
}
|
||||||
|
proven = true
|
||||||
|
}
|
||||||
|
by := claimant(public)
|
||||||
|
node, err := e.Inventory.Claim(ctx, secret, by, proven && !request.Redelivered)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return EnrolReply{}, err
|
return EnrolReply{}, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// From here the token is gone whatever happens next, so anything that fails leaves a node
|
|
||||||
// record with no live key — which is visible and fixable with a new token, where a spent
|
|
||||||
// token believed to be unspent is neither.
|
|
||||||
if _, err := e.Identity.RecordNodeKey(ctx, node.ID, public); err != nil {
|
if _, err := e.Identity.RecordNodeKey(ctx, node.ID, public); err != nil {
|
||||||
return EnrolReply{}, fmt.Errorf(
|
return EnrolReply{}, fmt.Errorf("%s's key could not be recorded: %w", node.Name, err)
|
||||||
"the token was spent and the key could not be recorded, so %s has no identity and "+
|
|
||||||
"needs a new token: %w", node.Name, err)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
key, err := e.Identity.Active(ctx)
|
key, err := e.Identity.Active(ctx)
|
||||||
@@ -61,7 +88,7 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply,
|
|||||||
return EnrolReply{}, err
|
return EnrolReply{}, err
|
||||||
}
|
}
|
||||||
|
|
||||||
reply := EnrolReply{
|
reply = EnrolReply{
|
||||||
Accepted: true,
|
Accepted: true,
|
||||||
Node: node.Name,
|
Node: node.Name,
|
||||||
Queue: QueueFor(node.Name),
|
Queue: QueueFor(node.Name),
|
||||||
@@ -70,22 +97,6 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply,
|
|||||||
Signer: key.Public,
|
Signer: key.Public,
|
||||||
}
|
}
|
||||||
|
|
||||||
// The token's secret was the broker password up to this moment, which is what let this
|
|
||||||
// connection exist at all. It is replaced now, so the one-time thing stays one-time and the
|
|
||||||
// credential the node keeps for years is not the one that was pasted into a terminal.
|
|
||||||
if e.Management != nil {
|
|
||||||
password, err := freshPassword()
|
|
||||||
if err != nil {
|
|
||||||
return EnrolReply{}, err
|
|
||||||
}
|
|
||||||
if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil {
|
|
||||||
return EnrolReply{}, fmt.Errorf(
|
|
||||||
"the token was spent and %s's broker password could not be replaced: %w",
|
|
||||||
node.Name, err)
|
|
||||||
}
|
|
||||||
reply.Password = password
|
|
||||||
}
|
|
||||||
|
|
||||||
// Recorded before the profile because the overlay is the first declaration this node will
|
// Recorded before the profile because the overlay is the first declaration this node will
|
||||||
// receive, and without this key the mesh cannot compose one. A node enrolled with no overlay
|
// receive, and without this key the mesh cannot compose one. A node enrolled with no overlay
|
||||||
// key is a node the graph skips — an ordinary in-between state, and one worth leaving as
|
// key is a node the graph skips — an ordinary in-between state, and one worth leaving as
|
||||||
@@ -98,21 +109,43 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply,
|
|||||||
if request.ServingKey != "" {
|
if request.ServingKey != "" {
|
||||||
if err := e.Identity.RecordServingKey(ctx, node.ID, request.ServingKey); err != nil {
|
if err := e.Identity.RecordServingKey(ctx, node.ID, request.ServingKey); err != nil {
|
||||||
return EnrolReply{}, fmt.Errorf(
|
return EnrolReply{}, fmt.Errorf(
|
||||||
"the token was spent and %s's serving key could not be recorded: %w",
|
"%s's serving key could not be recorded: %w", node.Name, err)
|
||||||
node.Name, err)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if request.SealingKey != "" {
|
if request.SealingKey != "" {
|
||||||
if err := e.Inventory.RecordSealingKey(ctx, node.ID, request.SealingKey); err != nil {
|
if err := e.Inventory.RecordSealingKey(ctx, node.ID, request.SealingKey); err != nil {
|
||||||
return EnrolReply{}, fmt.Errorf(
|
return EnrolReply{}, fmt.Errorf(
|
||||||
"the token was spent and %s's sealing key could not be recorded: %w",
|
"%s's sealing key could not be recorded: %w", node.Name, err)
|
||||||
node.Name, err)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if request.OverlayKey != "" {
|
if request.OverlayKey != "" {
|
||||||
if err := e.Inventory.RecordOverlayKey(ctx, node.ID, request.OverlayKey); err != nil {
|
if err := e.Inventory.RecordOverlayKey(ctx, node.ID, request.OverlayKey); err != nil {
|
||||||
return EnrolReply{}, fmt.Errorf(
|
return EnrolReply{}, fmt.Errorf(
|
||||||
"the token was spent and %s's overlay key could not be recorded: %w", node.Name, err)
|
"%s's overlay key could not be recorded: %w", node.Name, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Spent once the node is complete in the store.
|
||||||
|
if err := e.Inventory.Spend(ctx, secret, by); err != nil {
|
||||||
|
return EnrolReply{}, fmt.Errorf("%s was written and its token could not be spent: %w", node.Name, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The token's secret was the broker password up to this moment, which is what let this
|
||||||
|
// connection exist at all. It is replaced now, so the one-time thing stays one-time and the
|
||||||
|
// credential the node keeps for years is not the one that was pasted into a terminal. After
|
||||||
|
// the spend and not before: a replaced password on an attempt that failed would be held by
|
||||||
|
// nobody. If the broker will not take it now, the enrolment still stands — the node keeps
|
||||||
|
// the token's secret as its password, which it is told, and which is said here.
|
||||||
|
if e.Management != nil {
|
||||||
|
password, err := freshPassword()
|
||||||
|
if err != nil {
|
||||||
|
return EnrolReply{}, err
|
||||||
|
}
|
||||||
|
if err := e.Management.CreateNodeAccount(ctx, node.Name, password); err != nil {
|
||||||
|
log.Printf("%s is enrolled and its broker password could not be replaced, so it keeps "+
|
||||||
|
"the token's secret as its password: %v", node.Name, err)
|
||||||
|
} else {
|
||||||
|
reply.Password = password
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -125,6 +158,12 @@ func (e Enrolment) Enrol(ctx context.Context, request EnrolRequest) (EnrolReply,
|
|||||||
return reply, nil
|
return reply, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// claimant names the key presenting a token, so a claim can be held for it alone.
|
||||||
|
func claimant(public ed25519.PublicKey) string {
|
||||||
|
sum := sha256.Sum256(public)
|
||||||
|
return hex.EncodeToString(sum[:])
|
||||||
|
}
|
||||||
|
|
||||||
// freshPassword is the node's own broker credential from enrolment onward.
|
// freshPassword is the node's own broker credential from enrolment onward.
|
||||||
func freshPassword() (string, error) {
|
func freshPassword() (string, error) {
|
||||||
raw := make([]byte, 32)
|
raw := make([]byte, 32)
|
||||||
|
|||||||
@@ -6,6 +6,8 @@
|
|||||||
// traffic between them, each receiving half of what it expects. That has happened here before.
|
// traffic between them, each receiving half of what it expects. That has happened here before.
|
||||||
package link
|
package link
|
||||||
|
|
||||||
|
import "encoding/base64"
|
||||||
|
|
||||||
// Exchange is where nodes publish everything they have to say.
|
// Exchange is where nodes publish everything they have to say.
|
||||||
const Exchange = "mesh"
|
const Exchange = "mesh"
|
||||||
|
|
||||||
@@ -58,6 +60,18 @@ type EnrolRequest struct {
|
|||||||
// Profile is what this machine can be asked to do. The control plane cannot decide what a
|
// Profile is what this machine can be asked to do. The control plane cannot decide what a
|
||||||
// node should run without it, so it arrives with enrolment rather than being asked for after.
|
// node should run without it, so it arrives with enrolment rather than being asked for after.
|
||||||
Profile map[string]any `json:"profile,omitempty"`
|
Profile map[string]any `json:"profile,omitempty"`
|
||||||
|
|
||||||
|
// Proof is the node's identity key signing EnrolProof over this request: that the presenter
|
||||||
|
// holds the private half of PublicKey, not only knows the public one. Required to finish an
|
||||||
|
// enrolment whose token this key already spent — the case of an answer lost after the spend —
|
||||||
|
// because a public key is no secret, and without it anyone holding a leaked token and a
|
||||||
|
// node's public key could replay the spent token (novox/hq issue 083, on review).
|
||||||
|
Proof []byte `json:"proof,omitempty"`
|
||||||
|
|
||||||
|
// Redelivered is set by the control plane, never sent: the broker handed this request over a
|
||||||
|
// second time. Such a request does not finish an enrolment already spent — the first time may
|
||||||
|
// have answered, and the node holds what it was told.
|
||||||
|
Redelivered bool `json:"-"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// Signed is a declaration and the signature over it.
|
// Signed is a declaration and the signature over it.
|
||||||
@@ -109,6 +123,11 @@ type EnrolReply struct {
|
|||||||
// Accepted says whether the node is now known.
|
// Accepted says whether the node is now known.
|
||||||
Accepted bool `json:"accepted"`
|
Accepted bool `json:"accepted"`
|
||||||
|
|
||||||
|
// TryAgain says the mesh cannot answer right now — its store is restarting, or the token is
|
||||||
|
// held for a moment by another enrolment — and the node should ask again with the same
|
||||||
|
// request. Nothing was spent (novox/hq issue 083).
|
||||||
|
TryAgain bool `json:"try_again,omitempty"`
|
||||||
|
|
||||||
// Node is the name the mesh has for this machine, which settles any disagreement: the token
|
// Node is the name the mesh has for this machine, which settles any disagreement: the token
|
||||||
// was issued for a node record, and that record's name wins over what the machine called
|
// was issued for a node record, and that record's name wins over what the machine called
|
||||||
// itself.
|
// itself.
|
||||||
@@ -132,3 +151,10 @@ type EnrolReply struct {
|
|||||||
// token can fail — unknown, spent, expired — so that guessing learns nothing.
|
// token can fail — unknown, spent, expired — so that guessing learns nothing.
|
||||||
Refusal string `json:"refusal,omitempty"`
|
Refusal string `json:"refusal,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// EnrolProof is what a node signs with its identity key when it enrols: the token and every key it
|
||||||
|
// presents, so a proof cannot be moved to another request.
|
||||||
|
func EnrolProof(secret string, public []byte, overlay, sealing, serving string) []byte {
|
||||||
|
return []byte("novox-mesh-enrol\x00" + secret + "\x00" + base64.StdEncoding.EncodeToString(public) +
|
||||||
|
"\x00" + overlay + "\x00" + sealing + "\x00" + serving)
|
||||||
|
}
|
||||||
|
|||||||
@@ -20,85 +20,85 @@ func (a *saidTo) Nack(_ uint64, _ bool, requeue bool) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
func (a *saidTo) Reject(uint64, bool) error { 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 }
|
type heardWith struct{ err error }
|
||||||
|
|
||||||
func (h heardWith) Heard(context.Context, Report) error { return h.err }
|
func (h heardWith) Heard(context.Context, Report) error { return h.err }
|
||||||
|
|
||||||
func aReport(t *testing.T, to *saidTo) amqp.Delivery {
|
// 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()
|
t.Helper()
|
||||||
body, err := json.Marshal(Report{Node: "anchor", Applied: []string{"store"}})
|
body, err := json.Marshal(Report{Node: node, Declared: declared, Applied: []string{"store"}})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
return amqp.Delivery{Acknowledger: to, RoutingKey: KeyReport, Body: body}
|
tag++
|
||||||
|
return amqp.Delivery{Acknowledger: to, RoutingKey: KeyReport, Body: body, DeliveryTag: tag}
|
||||||
}
|
}
|
||||||
|
|
||||||
// A report the store could not take right now goes back to the broker to be asked again; one the
|
func quiet() *log.Logger { return log.New(io.Discard, "", 0) }
|
||||||
// store answered no to is acknowledged, or it would come back for ever (novox/hq issue 082).
|
|
||||||
func TestAReportTheStoreCouldNotTakeIsKeptAndOneItRefusedIsNot(t *testing.T) {
|
|
||||||
quiet := log.New(io.Discard, "", 0)
|
|
||||||
|
|
||||||
notNow := &saidTo{}
|
// A report the store could not take right now is held, unsettled, and recorded when the store is
|
||||||
s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet, again: 1}
|
// back; one the store answered no to is acknowledged; one recorded is acknowledged (issue 082, 083).
|
||||||
s.handleReport(context.Background(), aReport(t, notNow))
|
func TestAReportTheStoreCouldNotTakeIsHeldAndOneItRefusedIsNot(t *testing.T) {
|
||||||
if notNow.acked || !notNow.nacked || !notNow.requeued {
|
store := &switchable{err: errors.Join(ErrTryAgain, errors.New("starting up"))}
|
||||||
t.Fatalf("a report the store could not take yet was not handed back to be asked again: %+v", notNow)
|
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{}
|
refused := &saidTo{}
|
||||||
s = &Server{listener: heardWith{err: errors.New("a report named no node")}, log: quiet, again: 1}
|
s = &Server{listener: heardWith{err: errors.New("a report named no node")}, log: quiet()}
|
||||||
s.handleReport(context.Background(), aReport(t, refused))
|
s.handleReport(context.Background(), aReport(t, refused, "anchor", "d1"))
|
||||||
if !refused.acked || refused.nacked {
|
if !refused.acked || refused.nacked {
|
||||||
t.Fatalf("a report the store answered no to was not acknowledged, so it would spin: %+v", refused)
|
t.Fatalf("a report the store answered no to was not acknowledged: %+v", refused)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
recorded := &saidTo{}
|
// A newer report from the same node supersedes one of its reports still held: recorded after the
|
||||||
s = &Server{listener: heardWith{}, log: quiet}
|
// newer, the older would overwrite what the node is doing now.
|
||||||
s.handleReport(context.Background(), aReport(t, recorded))
|
func TestANewerReportSupersedesAHeldOneFromTheSameNode(t *testing.T) {
|
||||||
if !recorded.acked || recorded.nacked {
|
s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))}, log: quiet()}
|
||||||
t.Fatalf("a recorded report was not acknowledged: %+v", recorded)
|
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,
|
// A store that has not come back within the bound is not restarting: the report is let go, loudly,
|
||||||
// rather than holding every enrolment and report behind it for ever (novox/hq issue 082, review).
|
// rather than held for ever.
|
||||||
func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
|
func TestAReportIsLetGoOnceTheStoreHasBeenGoneTooLong(t *testing.T) {
|
||||||
quiet := log.New(io.Discard, "", 0)
|
|
||||||
s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("connection refused"))},
|
s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("connection refused"))},
|
||||||
log: quiet, again: 1, giveUp: time.Millisecond}
|
log: quiet(), giveUp: time.Millisecond}
|
||||||
|
held := &saidTo{}
|
||||||
first := &saidTo{}
|
s.handleReport(context.Background(), aReport(t, held, "anchor", "d1"))
|
||||||
s.handleReport(context.Background(), aReport(t, first))
|
if !held.unsettled() {
|
||||||
if !first.requeued {
|
t.Fatalf("the first failure was not held: %+v", held)
|
||||||
t.Fatalf("the first failure was not handed back: %+v", first)
|
|
||||||
}
|
}
|
||||||
time.Sleep(5 * time.Millisecond)
|
time.Sleep(5 * time.Millisecond)
|
||||||
later := &saidTo{}
|
s.retryHeld(context.Background())
|
||||||
s.handleReport(context.Background(), aReport(t, later))
|
if !held.acked || len(s.parked) != 0 {
|
||||||
if !later.acked || later.nacked {
|
t.Fatalf("a report past the bound was not let go: %+v, %d held", held, len(s.parked))
|
||||||
t.Fatalf("a report past the bound was not let go, so it would hold the queue for ever: %+v", later)
|
|
||||||
}
|
|
||||||
// Let go, and forgotten: the same report arriving fresh starts a new clock.
|
|
||||||
again := &saidTo{}
|
|
||||||
s.handleReport(context.Background(), aReport(t, again))
|
|
||||||
if !again.requeued {
|
|
||||||
t.Fatalf("a report let go was not forgotten, so its next arrival is given no chance: %+v", again)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// Shutting down does not wait out the pause.
|
|
||||||
func TestAReportBeingTriedAgainDoesNotHoldUpShutdown(t *testing.T) {
|
|
||||||
quiet := log.New(io.Discard, "", 0)
|
|
||||||
s := &Server{listener: heardWith{err: errors.Join(ErrTryAgain, errors.New("starting up"))},
|
|
||||||
log: quiet, again: time.Hour}
|
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
|
||||||
cancel()
|
|
||||||
done := make(chan struct{})
|
|
||||||
go func() { s.handleReport(ctx, aReport(t, &saidTo{})); close(done) }()
|
|
||||||
select {
|
|
||||||
case <-done:
|
|
||||||
case <-time.After(5 * time.Second):
|
|
||||||
t.Fatal("a cancelled context still waited out the pause")
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+217
-65
@@ -2,10 +2,13 @@ package link
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"crypto/sha256"
|
||||||
|
"encoding/hex"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"github.com/novox/mesh-controller/internal/envfile"
|
"github.com/novox/mesh-controller/internal/envfile"
|
||||||
|
"github.com/novox/mesh-controller/internal/inventory"
|
||||||
"log"
|
"log"
|
||||||
"os"
|
"os"
|
||||||
"time"
|
"time"
|
||||||
@@ -52,29 +55,44 @@ type Server struct {
|
|||||||
log *log.Logger
|
log *log.Logger
|
||||||
upgrader Upgrader
|
upgrader Upgrader
|
||||||
replayer Replayer
|
replayer Replayer
|
||||||
// again is how long a report the store could not take waits before it is handed back to
|
// Messages the store could not take right now, held unacknowledged and tried again on a
|
||||||
// the broker; zero means TryAgainAfter. giveUp is how long one report is kept trying before
|
// ticker, by subject (novox/hq issues 082, 083). again is the ticker's interval, zero meaning
|
||||||
// it is let go; zero means GiveUpAfter. waiting is when each report still trying first failed.
|
// TryAgainAfter; giveUp is how long one is kept, zero meaning GiveUpAfter.
|
||||||
again time.Duration
|
again time.Duration
|
||||||
giveUp time.Duration
|
giveUp time.Duration
|
||||||
waiting map[string]time.Time
|
parked map[string]*held
|
||||||
|
}
|
||||||
|
|
||||||
|
// held is one message the store could not take, kept to be tried again.
|
||||||
|
type held struct {
|
||||||
|
delivery amqp.Delivery
|
||||||
|
retry func(context.Context, amqp.Delivery)
|
||||||
|
what string
|
||||||
|
first time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
// ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and
|
// ErrTryAgain marks a listener's failure as "not now": what it was given is worth keeping and
|
||||||
// asking again, as when the store is restarting (novox/hq issue 082).
|
// asking again, as when the store is restarting (novox/hq issue 082).
|
||||||
var ErrTryAgain = errors.New("not now, try again")
|
var ErrTryAgain = errors.New("not now, try again")
|
||||||
|
|
||||||
// TryAgainAfter is the pause before a report the store could not take goes back to the broker.
|
// TryAgainAfter is how often messages the store could not take are tried again. A store comes
|
||||||
// The consumer takes one message at a time, so without it a restarting store would be asked in
|
// back in seconds, and a report a few seconds late is still current.
|
||||||
// a tight loop; a store comes back in seconds, and a report a few seconds late is still current.
|
|
||||||
const TryAgainAfter = 2 * time.Second
|
const TryAgainAfter = 2 * time.Second
|
||||||
|
|
||||||
// GiveUpAfter bounds how long one report holds the queue. The consumer takes one message at a
|
// GiveUpAfter bounds how long one message is kept trying. A store that has not come back in this
|
||||||
// time, so a report being tried again holds every enrolment, build result and other report
|
// long is not restarting, and the message is let go with a line saying it was lost.
|
||||||
// behind it; a store that has not come back in this long is not restarting, and holding the
|
|
||||||
// mesh's control queue for it would turn one lost report into a mesh that answers nothing.
|
|
||||||
const GiveUpAfter = 2 * time.Minute
|
const GiveUpAfter = 2 * time.Minute
|
||||||
|
|
||||||
|
// Prefetch is how many messages the broker hands the control plane before it has settled them.
|
||||||
|
// More than one because a message the store could not take is held, unsettled, while the loop goes
|
||||||
|
// on answering others — an enrolment above all, which a host is waiting on (novox/hq issue 083).
|
||||||
|
// Bounded, because what is held is also what the broker has not kept on its own disk as pending.
|
||||||
|
const Prefetch = 64
|
||||||
|
|
||||||
|
// PrefetchHeadroom is how much of the prefetch is never held, so the loop always has messages to
|
||||||
|
// answer — an enrolment above all — while others wait for the store.
|
||||||
|
const PrefetchHeadroom = 8
|
||||||
|
|
||||||
// Records tells the server where to keep build results.
|
// Records tells the server where to keep build results.
|
||||||
//
|
//
|
||||||
// Set after Connect rather than passed to it, because a control plane that only publishes — the
|
// Set after Connect rather than passed to it, because a control plane that only publishes — the
|
||||||
@@ -91,7 +109,8 @@ type Upgrader interface {
|
|||||||
// Upgraded is told which module moved and between which commits. An error is logged and the
|
// Upgraded is told which module moved and between which commits. An error is logged and the
|
||||||
// message is not requeued: an upgrade the control plane could not act on is not one it will
|
// message is not requeued: an upgrade the control plane could not act on is not one it will
|
||||||
// act on by being handed the same message again, and a poison message on a durable queue
|
// act on by being handed the same message again, and a poison message on a durable queue
|
||||||
// would stop every upgrade behind it.
|
// would stop every upgrade behind it — except the store unreachable for the moment, which is
|
||||||
|
// asked again for a bounded time (novox/hq issue 083).
|
||||||
Upgraded(ctx context.Context, u Upgraded) error
|
Upgraded(ctx context.Context, u Upgraded) error
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -202,10 +221,11 @@ func (s *Server) Close() {
|
|||||||
// receive half of what it expects — a fault this project has already had, between a module's
|
// receive half of what it expects — a fault this project has already had, between a module's
|
||||||
// daemon and its capability server.
|
// daemon and its capability server.
|
||||||
func (s *Server) Serve(ctx context.Context) error {
|
func (s *Server) Serve(ctx context.Context) error {
|
||||||
// Prefetch of one. The control plane writes to a database per message, and a burst of
|
// A bounded prefetch rather than one. The loop still takes messages one at a time; what the
|
||||||
// enrolments delivered all at once would be held in memory rather than left on the broker,
|
// prefetch buys is that a message the store could not take can be held while the loop goes on
|
||||||
// which is the one place they survive a restart.
|
// to the next, instead of every enrolment waiting behind it (novox/hq issue 083). Anything held
|
||||||
if err := s.channel.Qos(1, 0, false); err != nil {
|
// goes back to the broker if the control plane stops, because nothing held is acknowledged.
|
||||||
|
if err := s.channel.Qos(Prefetch, 0, false); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -246,10 +266,22 @@ func (s *Server) Serve(ctx context.Context) error {
|
|||||||
s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded)
|
s.log.Printf("consuming %s, bound to %s/%s", UpgradeQueue, EventsExchange, KeyModuleUpgraded)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
again := s.again
|
||||||
|
if again == 0 {
|
||||||
|
again = TryAgainAfter
|
||||||
|
}
|
||||||
|
ticker := time.NewTicker(again)
|
||||||
|
defer ticker.Stop()
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return nil
|
return nil
|
||||||
|
case <-ticker.C:
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
s.retryHeld(ctx)
|
||||||
case delivery, ok := <-catchups:
|
case delivery, ok := <-catchups:
|
||||||
if !ok {
|
if !ok {
|
||||||
if catchups != nil {
|
if catchups != nil {
|
||||||
@@ -330,29 +362,19 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
_ = delivery.Reject(false)
|
_ = delivery.Reject(false)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
// A node's newer report supersedes one of its older reports still held: the older is its
|
||||||
|
// past, and recorded after the newer it would overwrite what the node is doing now.
|
||||||
|
subject := "report " + report.Node
|
||||||
|
s.supersede(subject, delivery)
|
||||||
if s.listener != nil {
|
if s.listener != nil {
|
||||||
if err := s.listener.Heard(context.Background(), report); err != nil {
|
if err := s.listener.Heard(context.Background(), report); err != nil {
|
||||||
if errors.Is(err, ErrTryAgain) && s.keepTrying(report) {
|
// Held, not acknowledged, while the store cannot take it: the node reports an apply
|
||||||
// Kept, not acknowledged. The node reports an apply once, and a report lost here
|
// once, and a report lost here is a node the mesh never hears from again — the store
|
||||||
// is a node the mesh never hears from again: the store restarting under the
|
// restarting under the adoption that node just applied lost exactly that (issue 082).
|
||||||
// adoption that node just applied lost exactly that (novox/hq issue 082).
|
what := fmt.Sprintf("%s's report of declaration %s", report.Node, report.Declared)
|
||||||
again := s.again
|
if s.tryLater(ctx, delivery, subject, what, err, s.handle) {
|
||||||
if again == 0 {
|
|
||||||
again = TryAgainAfter
|
|
||||||
}
|
|
||||||
s.log.Printf("could not record %s's report yet, and will again in %s: %v", report.Node, again, err)
|
|
||||||
select {
|
|
||||||
case <-ctx.Done():
|
|
||||||
case <-time.After(again):
|
|
||||||
}
|
|
||||||
_ = delivery.Nack(false, true)
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if errors.Is(err, ErrTryAgain) {
|
|
||||||
s.log.Printf("LOST %s's report of declaration %s: the store has not come back in %s, "+
|
|
||||||
"and the control queue cannot wait longer — the node is current but the mesh will "+
|
|
||||||
"read it as unanswered until its next push: %v", report.Node, report.Declared, s.giveUpAfter(), err)
|
|
||||||
}
|
|
||||||
// Said rather than swallowed. A report the mesh heard and failed to write down is a
|
// Said rather than swallowed. A report the mesh heard and failed to write down is a
|
||||||
// node whose recovery copy is silently older than it looks.
|
// node whose recovery copy is silently older than it looks.
|
||||||
s.log.Printf("could not record %s's report: %v", report.Node, err)
|
s.log.Printf("could not record %s's report: %v", report.Node, err)
|
||||||
@@ -367,26 +389,96 @@ func (s *Server) handleReport(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
default:
|
default:
|
||||||
s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied))
|
s.log.Printf("%s applied %d resource(s)", report.Node, len(report.Applied))
|
||||||
}
|
}
|
||||||
s.stopTrying(report)
|
s.settled(subject, delivery)
|
||||||
_ = delivery.Ack(false)
|
_ = delivery.Ack(false)
|
||||||
}
|
}
|
||||||
|
|
||||||
// keepTrying says whether a report the store could not take is still within the time it may hold
|
// tryLater holds a message the store could not take right now, to be tried again on the ticker,
|
||||||
// the queue, starting that clock on its first failure.
|
// and says whether it did (novox/hq issues 082, 083).
|
||||||
func (s *Server) keepTrying(r Report) bool {
|
//
|
||||||
if s.waiting == nil {
|
// "Right now" is the store unreachable or restarting — ErrTryAgain from a listener, or an error the
|
||||||
s.waiting = map[string]time.Time{}
|
// inventory reads as an outage. Anything else is an answer, and is left to the caller to settle.
|
||||||
}
|
// Held means unacknowledged and set aside: the loop goes on to the next message, so an enrolment a
|
||||||
key := r.Node + " " + r.Declared
|
// host is waiting on is answered while a report waits for the store. One message is held at most
|
||||||
first, seen := s.waiting[key]
|
// giveUpAfter; past it, it is let go with a line saying it was lost, and the caller settles it.
|
||||||
if !seen {
|
func (s *Server) tryLater(ctx context.Context, delivery amqp.Delivery, subject, what string, err error,
|
||||||
s.waiting[key] = time.Now()
|
retry func(context.Context, amqp.Delivery)) bool {
|
||||||
|
// Shutting down: nothing is settled. Unsettled, the broker hands the message to whatever
|
||||||
|
// consumes next — a cancelled context is not an answer about the message (issue 083, review).
|
||||||
|
if ctx.Err() != nil {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
return time.Since(first) < s.giveUpAfter()
|
if !errors.Is(err, ErrTryAgain) && !inventory.Unreachable(err) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
if s.parked == nil {
|
||||||
|
s.parked = map[string]*held{}
|
||||||
|
}
|
||||||
|
h, ok := s.parked[subject]
|
||||||
|
if !ok || h.delivery.DeliveryTag != delivery.DeliveryTag {
|
||||||
|
// Held no further than the prefetch leaves room: past it, the broker would hand the loop
|
||||||
|
// nothing new — enrolments included — until something held was let go.
|
||||||
|
if !ok && len(s.parked) >= Prefetch-PrefetchHeadroom {
|
||||||
|
s.log.Printf("LOST %s: %d messages are already held for the store, and holding more "+
|
||||||
|
"would stop the queue: %v", what, len(s.parked), err)
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
h = &held{delivery: delivery, retry: retry, what: what, first: time.Now()}
|
||||||
|
s.parked[subject] = h
|
||||||
|
s.log.Printf("could not keep %s yet; holding it to try again: %v", what, err)
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
if time.Since(h.first) >= s.giveUpAfter() {
|
||||||
|
delete(s.parked, subject)
|
||||||
|
s.log.Printf("LOST %s: the store has not come back in %s: %v", what, s.giveUpAfter(), err)
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Server) stopTrying(r Report) { delete(s.waiting, r.Node+" "+r.Declared) }
|
// supersede drops a message held for a subject when a newer one for it arrives: the older is
|
||||||
|
// acknowledged, because acting on it after the newer would undo the newer.
|
||||||
|
func (s *Server) supersede(subject string, newer amqp.Delivery) {
|
||||||
|
h, ok := s.parked[subject]
|
||||||
|
if !ok || h.delivery.DeliveryTag == newer.DeliveryTag {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
delete(s.parked, subject)
|
||||||
|
s.log.Printf("set aside %s: a newer one arrived", h.what)
|
||||||
|
_ = h.delivery.Ack(false)
|
||||||
|
}
|
||||||
|
|
||||||
|
// settled forgets a message once it has been handled either way.
|
||||||
|
func (s *Server) settled(subject string, delivery amqp.Delivery) {
|
||||||
|
if h, ok := s.parked[subject]; ok && h.delivery.DeliveryTag == delivery.DeliveryTag {
|
||||||
|
delete(s.parked, subject)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// digest names a message by its content.
|
||||||
|
func digest(body []byte) string {
|
||||||
|
sum := sha256.Sum256(body)
|
||||||
|
return hex.EncodeToString(sum[:])
|
||||||
|
}
|
||||||
|
|
||||||
|
// retryHeld tries every held message again. Each handler holds it again, settles it, or lets it
|
||||||
|
// go past the bound.
|
||||||
|
func (s *Server) retryHeld(ctx context.Context) {
|
||||||
|
for _, h := range s.snapshot() {
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
h.retry(ctx, h.delivery)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) snapshot() []*held {
|
||||||
|
out := make([]*held, 0, len(s.parked))
|
||||||
|
for _, h := range s.parked {
|
||||||
|
out = append(out, h)
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
func (s *Server) giveUpAfter() time.Duration {
|
func (s *Server) giveUpAfter() time.Duration {
|
||||||
if s.giveUp == 0 {
|
if s.giveUp == 0 {
|
||||||
@@ -402,12 +494,27 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
if err := json.Unmarshal(delivery.Body, &request); err != nil {
|
if err := json.Unmarshal(delivery.Body, &request); err != nil {
|
||||||
s.log.Printf("an enrolment request could not be read: %v", err)
|
s.log.Printf("an enrolment request could not be read: %v", err)
|
||||||
} else {
|
} else {
|
||||||
|
request.Redelivered = delivery.Redelivered
|
||||||
accepted, err := s.enroller.Enrol(ctx, request)
|
accepted, err := s.enroller.Enrol(ctx, request)
|
||||||
if err != nil {
|
switch {
|
||||||
|
case errors.Is(err, ErrTryAgain):
|
||||||
|
// Not a refusal: nothing was spent, and the same request asked again will be
|
||||||
|
// answered. Replied at once rather than held, so the node — which is waiting on
|
||||||
|
// this answer — decides when to ask, and the queue behind it moves (issue 083).
|
||||||
|
reply = EnrolReply{TryAgain: true, Refusal: "the mesh cannot answer right now; ask again"}
|
||||||
|
s.log.Printf("asked %q to enrol again shortly: %v", request.Node, err)
|
||||||
|
case err != nil && request.Redelivered:
|
||||||
|
// Said as what it most likely is: the broker handed this request over again after
|
||||||
|
// the control plane stopped mid-answer, and an enrolment already spent is not
|
||||||
|
// finished a second time. The node may need a new token.
|
||||||
|
s.log.Printf("refusing a redelivered enrolment for %q — it may have finished before "+
|
||||||
|
"the control plane stopped, and if the node did not get its answer it needs a new "+
|
||||||
|
"token: %v", request.Node, err)
|
||||||
|
case err != nil:
|
||||||
// Logged in full here, where an operator can see it; sent back as one refusal, so
|
// Logged in full here, where an operator can see it; sent back as one refusal, so
|
||||||
// that somebody guessing learns nothing from which reason came back.
|
// that somebody guessing learns nothing from which reason came back.
|
||||||
s.log.Printf("refusing enrolment for %q: %v", request.Node, err)
|
s.log.Printf("refusing enrolment for %q: %v", request.Node, err)
|
||||||
} else {
|
default:
|
||||||
reply = accepted
|
reply = accepted
|
||||||
s.log.Printf("enrolled %s", accepted.Node)
|
s.log.Printf("enrolled %s", accepted.Node)
|
||||||
}
|
}
|
||||||
@@ -416,9 +523,8 @@ func (s *Server) handleEnrol(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
s.reply(ctx, delivery, reply)
|
s.reply(ctx, delivery, reply)
|
||||||
|
|
||||||
// Acknowledged after the reply is sent, so a control plane that dies mid-answer leaves the
|
// Acknowledged after the reply is sent, so a control plane that dies mid-answer leaves the
|
||||||
// request on the broker rather than having consumed it silently. Enrolment is idempotent
|
// request on the broker rather than having consumed it silently. Asked again by the same
|
||||||
// only in the sense that the token is spent — a redelivery gets the refusal, which is
|
// presenter, an enrolment finishes: the token is held for its key and spent last (issue 083).
|
||||||
// correct and visible, where a lost request is neither.
|
|
||||||
_ = delivery.Ack(false)
|
_ = delivery.Ack(false)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -464,11 +570,21 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
_ = delivery.Reject(false)
|
_ = delivery.Reject(false)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
// Each build result its own subject: none supersedes another, and recording one twice is
|
||||||
|
// harmless — the build is kept by its id.
|
||||||
|
subject := "build " + digest(delivery.Body)
|
||||||
|
s.supersede(subject, delivery)
|
||||||
if err := s.recorder.Built(ctx, result); err != nil {
|
if err := s.recorder.Built(ctx, result); err != nil {
|
||||||
|
// Held while the store cannot take it: a build result lost here is never announced (083).
|
||||||
|
if s.tryLater(ctx, delivery, subject, fmt.Sprintf("a build result from %s", result.On), err, s.handle) {
|
||||||
|
return
|
||||||
|
}
|
||||||
s.log.Printf("cannot keep a build result from %s: %v", result.On, err)
|
s.log.Printf("cannot keep a build result from %s: %v", result.On, err)
|
||||||
|
s.settled(subject, delivery)
|
||||||
_ = delivery.Reject(false)
|
_ = delivery.Reject(false)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
s.settled(subject, delivery)
|
||||||
switch {
|
switch {
|
||||||
case result.Failed != "":
|
case result.Failed != "":
|
||||||
s.log.Printf("%s could not build %s", result.On, result.Repository)
|
s.log.Printf("%s could not build %s", result.On, result.Repository)
|
||||||
@@ -478,26 +594,33 @@ func (s *Server) handleBuilt(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
_ = delivery.Ack(false)
|
_ = delivery.Ack(false)
|
||||||
}
|
}
|
||||||
|
|
||||||
// upgraded hands one announcement to whatever is following them.
|
|
||||||
//
|
|
||||||
// **Acknowledged whatever happens.** A failure here is the control plane being unable to act on an
|
|
||||||
// upgrade — a machine that cannot be resolved, a broker that will not take a declaration — and
|
|
||||||
// none of those get better by being handed the same message again. Requeuing would put a poison
|
|
||||||
// message at the head of a durable queue and stop every upgrade behind it, which turns one module
|
|
||||||
// nobody can push into a mesh that stops following its own catalogue.
|
|
||||||
// catchingUp answers a catalogue that has just started and may have missed builds.
|
// catchingUp answers a catalogue that has just started and may have missed builds.
|
||||||
//
|
//
|
||||||
// Acknowledged before the work, deliberately: a replay that fails is not one that succeeds by
|
// Acknowledged after the work. A replay that fails for a reason other than the store is not one
|
||||||
// being handed the same request again, and the catalogue asks every time it starts. Requeueing a
|
// that succeeds by being handed the same request again, so that is acknowledged and said; but a
|
||||||
// poison request would stop every later catch-up behind it.
|
// store that could not be read right now is held and asked again, bounded, rather than the
|
||||||
|
// request lost until the catalogue next restarts (issue 083). One request stands for all: a newer
|
||||||
|
// one supersedes one still held.
|
||||||
func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) {
|
func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) {
|
||||||
defer func() { _ = delivery.Ack(false) }()
|
const subject = "catch-up"
|
||||||
|
s.supersede(subject, delivery)
|
||||||
|
holding := false
|
||||||
|
defer func() {
|
||||||
|
if !holding {
|
||||||
|
s.settled(subject, delivery)
|
||||||
|
_ = delivery.Ack(false)
|
||||||
|
}
|
||||||
|
}()
|
||||||
if s.replayer == nil {
|
if s.replayer == nil {
|
||||||
s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay")
|
s.log.Printf("a catalogue asked to catch up and this control plane has nothing to replay")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
announcements, err := s.replayer.Announceable(ctx)
|
announcements, err := s.replayer.Announceable(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
if s.tryLater(ctx, delivery, subject, "a catalogue's request to catch up", err, s.catchingUp) {
|
||||||
|
holding = true
|
||||||
|
return
|
||||||
|
}
|
||||||
s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err)
|
s.log.Printf("a catalogue asked to catch up and the mesh could not read its builds: %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -516,9 +639,28 @@ func (s *Server) catchingUp(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent)
|
s.log.Printf("a catalogue asked to catch up; re-announced %d build(s)", sent)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// upgraded hands one announcement to whatever is following them.
|
||||||
|
//
|
||||||
|
// **Acknowledged whatever happens, but one thing.** A failure here is usually the control plane
|
||||||
|
// being unable to act on an upgrade — a machine that cannot be resolved, a broker that will not
|
||||||
|
// take a declaration — and none of those get better by being handed the same message again. The
|
||||||
|
// one exception is the upgrader saying the store could not be read for the moment (ErrTryAgain):
|
||||||
|
// that is held and asked again, bounded (novox/hq issue 083). Only the upgrader's word counts
|
||||||
|
// here, not an error that merely looks like an outage — a push that timed out on the second
|
||||||
|
// machine is not asked again, or the first would be pushed every few seconds for two minutes.
|
||||||
|
// A newer move of the same module supersedes one still held.
|
||||||
func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) {
|
func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) {
|
||||||
defer func() { _ = delivery.Ack(false) }()
|
|
||||||
var u Upgraded
|
var u Upgraded
|
||||||
|
_ = json.Unmarshal(delivery.Body, &u)
|
||||||
|
subject := "upgrade " + u.Module
|
||||||
|
s.supersede(subject, delivery)
|
||||||
|
holding := false
|
||||||
|
defer func() {
|
||||||
|
if !holding {
|
||||||
|
s.settled(subject, delivery)
|
||||||
|
_ = delivery.Ack(false)
|
||||||
|
}
|
||||||
|
}()
|
||||||
if err := json.Unmarshal(delivery.Body, &u); err != nil {
|
if err := json.Unmarshal(delivery.Body, &u); err != nil {
|
||||||
s.log.Printf("an upgrade announcement could not be read: %v", err)
|
s.log.Printf("an upgrade announcement could not be read: %v", err)
|
||||||
return
|
return
|
||||||
@@ -528,6 +670,16 @@ func (s *Server) upgraded(ctx context.Context, delivery amqp.Delivery) {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
if err := s.upgrader.Upgraded(ctx, u); err != nil {
|
if err := s.upgrader.Upgraded(ctx, u); err != nil {
|
||||||
|
// Shutting down is not an answer about the announcement: left for the broker.
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
holding = true
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if errors.Is(err, ErrTryAgain) &&
|
||||||
|
s.tryLater(ctx, delivery, subject, fmt.Sprintf("%s's move to %s", u.Module, short(u.Commit)), err, s.upgraded) {
|
||||||
|
holding = true
|
||||||
|
return
|
||||||
|
}
|
||||||
s.log.Printf("%s moved to %s and the mesh could not act on it: %v",
|
s.log.Printf("%s moved to %s and the mesh could not act on it: %v",
|
||||||
u.Module, short(u.Commit), err)
|
u.Module, short(u.Commit), err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,181 @@
|
|||||||
|
package link
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"log"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5/pgconn"
|
||||||
|
amqp "github.com/rabbitmq/amqp091-go"
|
||||||
|
)
|
||||||
|
|
||||||
|
// What a store restarting under an adoption answers with (novox/hq issues 082, 083).
|
||||||
|
var restarting = &pgconn.PgError{Code: "57P03", Message: "the database system is starting up"}
|
||||||
|
|
||||||
|
type recordsWith struct{ err error }
|
||||||
|
|
||||||
|
func (r recordsWith) Built(context.Context, BuildResult) error { return r.err }
|
||||||
|
|
||||||
|
type upgradesWith struct{ err error }
|
||||||
|
|
||||||
|
func (u upgradesWith) Upgraded(context.Context, Upgraded) error { return u.err }
|
||||||
|
|
||||||
|
type replaysWith struct{ err error }
|
||||||
|
|
||||||
|
func (r replaysWith) Announceable(context.Context) ([]Announcement, error) { return nil, r.err }
|
||||||
|
|
||||||
|
type settledAs struct{ acked, nacked, requeued, rejected bool }
|
||||||
|
|
||||||
|
func (a *settledAs) Ack(uint64, bool) error { a.acked = true; return nil }
|
||||||
|
func (a *settledAs) Nack(_ uint64, _ bool, requeue bool) error {
|
||||||
|
a.nacked, a.requeued = true, requeue
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
func (a *settledAs) Reject(uint64, bool) error { a.rejected = true; return nil }
|
||||||
|
|
||||||
|
func a(t *testing.T, to *settledAs, key string, v any) amqp.Delivery {
|
||||||
|
t.Helper()
|
||||||
|
body, err := json.Marshal(v)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
tag++
|
||||||
|
return amqp.Delivery{Acknowledger: to, RoutingKey: key, Body: body, DeliveryTag: tag}
|
||||||
|
}
|
||||||
|
|
||||||
|
func quietServer() *Server { return &Server{log: log.New(io.Discard, "", 0)} }
|
||||||
|
|
||||||
|
func (a *settledAs) held() bool { return !a.acked && !a.nacked && !a.rejected }
|
||||||
|
|
||||||
|
// A build result the store could not take right now is handed back; one it refused is rejected,
|
||||||
|
// as before; one it kept is acknowledged.
|
||||||
|
func TestABuildResultWaitsOutARestartingStore(t *testing.T) {
|
||||||
|
built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}
|
||||||
|
for _, c := range []struct {
|
||||||
|
what string
|
||||||
|
err error
|
||||||
|
want func(*settledAs) bool
|
||||||
|
}{
|
||||||
|
{"restarting", restarting, func(s *settledAs) bool { return s.held() }},
|
||||||
|
{"refused", errors.New("no such module"), func(s *settledAs) bool { return s.rejected && !s.nacked }},
|
||||||
|
{"kept", nil, func(s *settledAs) bool { return s.acked && !s.nacked }},
|
||||||
|
} {
|
||||||
|
s := quietServer()
|
||||||
|
s.recorder = recordsWith{err: c.err}
|
||||||
|
to := &settledAs{}
|
||||||
|
s.handleBuilt(context.Background(), a(t, to, KeyBuilt, built))
|
||||||
|
if !c.want(to) {
|
||||||
|
t.Errorf("%s: a build result was settled as %+v", c.what, *to)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// An upgrade announcement arriving while the store restarts is asked again; any other failure is
|
||||||
|
// acknowledged, so it cannot stop every upgrade behind it.
|
||||||
|
func TestAnUpgradeWaitsOutARestartingStoreAndNothingElse(t *testing.T) {
|
||||||
|
moved := Upgraded{Module: "gitea", Commit: "abcdef0123"}
|
||||||
|
for _, c := range []struct {
|
||||||
|
what string
|
||||||
|
err error
|
||||||
|
want func(*settledAs) bool
|
||||||
|
}{
|
||||||
|
{"the store away, said by the upgrader", errors.Join(ErrTryAgain, restarting), func(s *settledAs) bool { return s.held() }},
|
||||||
|
{"a push that timed out", context.DeadlineExceeded, func(s *settledAs) bool { return s.acked && !s.nacked }},
|
||||||
|
{"cannot act", errors.New("anchor cannot be resolved"), func(s *settledAs) bool { return s.acked && !s.nacked }},
|
||||||
|
{"acted", nil, func(s *settledAs) bool { return s.acked && !s.nacked }},
|
||||||
|
} {
|
||||||
|
s := quietServer()
|
||||||
|
s.upgrader = upgradesWith{err: c.err}
|
||||||
|
to := &settledAs{}
|
||||||
|
s.upgraded(context.Background(), a(t, to, "upgraded", moved))
|
||||||
|
if !c.want(to) {
|
||||||
|
t.Errorf("%s: an upgrade was settled as %+v", c.what, *to)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// A catalogue's request to catch up is acknowledged after the work, and asked again while the
|
||||||
|
// store cannot be read — not lost until the catalogue next restarts.
|
||||||
|
func TestACatchUpWaitsOutARestartingStore(t *testing.T) {
|
||||||
|
for _, c := range []struct {
|
||||||
|
what string
|
||||||
|
err error
|
||||||
|
want func(*settledAs) bool
|
||||||
|
}{
|
||||||
|
{"restarting", restarting, func(s *settledAs) bool { return s.held() }},
|
||||||
|
{"unreadable", errors.New("a build row is malformed"), func(s *settledAs) bool { return s.acked && !s.nacked }},
|
||||||
|
{"nothing to replay", nil, func(s *settledAs) bool { return s.acked && !s.nacked }},
|
||||||
|
} {
|
||||||
|
s := quietServer()
|
||||||
|
s.replayer = replaysWith{err: c.err}
|
||||||
|
to := &settledAs{}
|
||||||
|
s.catchingUp(context.Background(), a(t, to, "catch-up", map[string]string{}))
|
||||||
|
if !c.want(to) {
|
||||||
|
t.Errorf("%s: a catch-up request was settled as %+v", c.what, *to)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Shutting down is not an answer about a message: one handled with a cancelled context is left
|
||||||
|
// unsettled, for the broker to hand to whatever consumes next (issue 083, review).
|
||||||
|
func TestAMessageHandledDuringShutdownIsLeftForTheBroker(t *testing.T) {
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
s := quietServer()
|
||||||
|
s.recorder = recordsWith{err: context.Canceled}
|
||||||
|
to := &settledAs{}
|
||||||
|
s.handleBuilt(ctx, a(t, to, KeyBuilt, BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}))
|
||||||
|
if !to.held() {
|
||||||
|
t.Fatalf("a build result handled during shutdown was settled, and so lost: %+v", *to)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Two identical build results: the newer sets the older aside rather than leaving it unsettled
|
||||||
|
// for ever, holding a place in the prefetch.
|
||||||
|
func TestAnIdenticalBuildResultSetsTheHeldOneAside(t *testing.T) {
|
||||||
|
s := quietServer()
|
||||||
|
s.recorder = recordsWith{err: restarting}
|
||||||
|
built := BuildResult{On: "anchor", Repository: "/r", Commit: "abc"}
|
||||||
|
first, second := &settledAs{}, &settledAs{}
|
||||||
|
s.handleBuilt(context.Background(), a(t, first, KeyBuilt, built))
|
||||||
|
s.handleBuilt(context.Background(), a(t, second, KeyBuilt, built))
|
||||||
|
if !first.acked || !second.held() || len(s.parked) != 1 {
|
||||||
|
t.Fatalf("an identical build result did not set the held one aside: first %+v second %+v, %d held",
|
||||||
|
*first, *second, len(s.parked))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// What is held stops short of the prefetch, so the loop always has room to answer an enrolment.
|
||||||
|
func TestWhatIsHeldLeavesRoomInThePrefetch(t *testing.T) {
|
||||||
|
s := quietServer()
|
||||||
|
s.recorder = recordsWith{err: restarting}
|
||||||
|
var last *settledAs
|
||||||
|
for i := 0; i < Prefetch; i++ {
|
||||||
|
last = &settledAs{}
|
||||||
|
s.handleBuilt(context.Background(), a(t, last, KeyBuilt, BuildResult{On: "anchor", Commit: fmt.Sprint(i)}))
|
||||||
|
}
|
||||||
|
if len(s.parked) != Prefetch-PrefetchHeadroom {
|
||||||
|
t.Fatalf("%d messages were held; the ceiling is %d", len(s.parked), Prefetch-PrefetchHeadroom)
|
||||||
|
}
|
||||||
|
if last.held() {
|
||||||
|
t.Fatalf("a message past the ceiling was held: %+v", *last)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// An upgrade handled during shutdown is left for the broker too — the upgrader's error is the
|
||||||
|
// cancelled context, which is no answer about the announcement.
|
||||||
|
func TestAnUpgradeHandledDuringShutdownIsLeftForTheBroker(t *testing.T) {
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
s := quietServer()
|
||||||
|
s.upgrader = upgradesWith{err: context.Canceled}
|
||||||
|
to := &settledAs{}
|
||||||
|
s.upgraded(ctx, a(t, to, "upgraded", Upgraded{Module: "gitea", Commit: "abcdef0123"}))
|
||||||
|
if !to.held() {
|
||||||
|
t.Fatalf("an upgrade handled during shutdown was settled, and so lost: %+v", *to)
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user