An enrolling node asks again while the mesh cannot answer, and proves it holds its key (issue 083) #19

Merged
jschoubben merged 3 commits from multiple-fixes into main 2026-09-22 12:42:57 +00:00
2 changed files with 113 additions and 17 deletions
Showing only changes of commit c8bbdb2da5 - Show all commits
+79 -17
View File
@@ -50,9 +50,15 @@ type EnrolRequest struct {
// EnrolReply is what the mesh says back. // EnrolReply is what the mesh says back.
type EnrolReply struct { type EnrolReply struct {
Accepted bool `json:"accepted"` Accepted bool `json:"accepted"`
Node string `json:"node,omitempty"`
Queue string `json:"queue,omitempty"` // TryAgain is the mesh saying it cannot answer right now — its store is restarting, or the
// token is held for a moment by another enrolment — and that nothing was spent. The same
// request is asked again (novox/hq issue 083).
TryAgain bool `json:"try_again,omitempty"`
Node string `json:"node,omitempty"`
Queue string `json:"queue,omitempty"`
// What this node keeps so it can come back on its own. Without these a restart would need a // What this node keeps so it can come back on its own. Without these a restart would need a
// person with a new token, which would make disconnection a crisis rather than the ordinary // person with a new token, which would make disconnection a crisis rather than the ordinary
@@ -68,6 +74,33 @@ type EnrolReply struct {
// ErrRefused is what a node gets when the mesh will not have it. // ErrRefused is what a node gets when the mesh will not have it.
var ErrRefused = errors.New("the mesh refused this enrolment") var ErrRefused = errors.New("the mesh refused this enrolment")
// EnrolPatience is how long a node keeps asking while the mesh says "try again" — as long as the
// mesh holds a token for the one enrolment presenting it, so a node asking the whole time is never
// held off by its own earlier attempt.
const EnrolPatience = 2 * time.Minute
// AskAgainAfter is the pause between asks while the mesh says "try again".
const AskAgainAfter = 3 * time.Second
// ErrNotNow is a mesh that said "try again" for longer than this node would keep asking.
var ErrNotNow = errors.New("the mesh could not answer this enrolment")
// answered decides what one reply means: done, ask again, or stop with an error. Separate from the
// broker so it can be held to that by a test.
func answered(reply EnrolReply, asking time.Duration) (again bool, err error) {
switch {
case reply.Accepted:
return false, nil
case reply.TryAgain && asking < EnrolPatience:
return true, nil
case reply.TryAgain:
return false, fmt.Errorf("%w for %s: %s. The token was not spent — run enrol again with it",
ErrNotNow, EnrolPatience, reply.Refusal)
default:
return false, fmt.Errorf("%w: %s", ErrRefused, reply.Refusal)
}
}
// Enrol presents this node's key and its one-time secret, and waits to be told it is known. // Enrol presents this node's key and its one-time secret, and waits to be told it is known.
// //
// The broker has already authenticated this connection: the account was created when the token // The broker has already authenticated this connection: the account was created when the token
@@ -128,18 +161,29 @@ func Enrol(ctx context.Context, address, pin, node, secret string, public []byte
return EnrolReply{}, err return EnrolReply{}, err
} }
correlation := fmt.Sprintf("%s-%d", node, time.Now().UnixNano()) // Asked, and asked again with the same request while the mesh says "try again": the keys
publish, cancel := context.WithTimeout(ctx, timeout) // this node generated are the ones it keeps, so the same request is the same enrolment, and
defer cancel() // the mesh holds the token for it (novox/hq issue 083).
if err := channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false, ask := func() (string, error) {
amqp.Publishing{ correlation := fmt.Sprintf("%s-%d", node, time.Now().UnixNano())
ContentType: "application/json", publish, cancel := context.WithTimeout(ctx, timeout)
CorrelationId: correlation, defer cancel()
ReplyTo: queue.Name, if err := channel.PublishWithContext(publish, Exchange, KeyEnrol, false, false,
Body: body, amqp.Publishing{
}); err != nil { ContentType: "application/json",
return EnrolReply{}, fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err) CorrelationId: correlation,
ReplyTo: queue.Name,
Body: body,
}); err != nil {
return "", fmt.Errorf("cannot publish to the %s exchange: %w", Exchange, err)
}
return correlation, nil
} }
correlation, err := ask()
if err != nil {
return EnrolReply{}, err
}
began := time.Now()
// Waited for rather than assumed. A published message that nothing answers means the control // Waited for rather than assumed. A published message that nothing answers means the control
// plane is not running, and a node that carried on regardless would believe it had joined a // plane is not running, and a node that carried on regardless would believe it had joined a
@@ -170,10 +214,28 @@ func Enrol(ctx context.Context, address, pin, node, secret string, public []byte
if err := json.Unmarshal(delivery.Body, &reply); err != nil { if err := json.Unmarshal(delivery.Body, &reply); err != nil {
return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err) return EnrolReply{}, fmt.Errorf("the mesh's answer could not be read: %w", err)
} }
if !reply.Accepted { again, err := answered(reply, time.Since(began))
return reply, fmt.Errorf("%w: %s", ErrRefused, reply.Refusal) if err != nil {
return reply, err
} }
return reply, nil if !again {
return reply, nil
}
select {
case <-ctx.Done():
return EnrolReply{}, ctx.Err()
case <-time.After(AskAgainAfter):
}
if correlation, err = ask(); err != nil {
return EnrolReply{}, err
}
if !deadline.Stop() {
select {
case <-deadline.C:
default:
}
}
deadline.Reset(timeout)
} }
} }
} }
+34
View File
@@ -0,0 +1,34 @@
package link
import (
"errors"
"strings"
"testing"
"time"
)
// What one reply means to a node enrolling (novox/hq issue 083): accepted is done; "try again" is
// asked again, until the node's patience runs out and it says the token was not spent; anything
// else is a refusal.
func TestAnEnrollingNodeAsksAgainWhileTheMeshSaysNotNow(t *testing.T) {
if again, err := answered(EnrolReply{Accepted: true}, 0); again || err != nil {
t.Fatalf("an accepted enrolment was not done: again=%v err=%v", again, err)
}
if again, err := answered(EnrolReply{TryAgain: true}, time.Second); !again || err != nil {
t.Fatalf("a mesh saying not now was not asked again: again=%v err=%v", again, err)
}
again, err := answered(EnrolReply{TryAgain: true, Refusal: "restarting"}, EnrolPatience)
if again || !errors.Is(err, ErrNotNow) {
t.Fatalf("a mesh saying not now past the node's patience did not stop it: again=%v err=%v", again, err)
}
if !errorsMention(err, "not spent") {
t.Errorf("giving up did not say the token can be used again: %v", err)
}
if again, err := answered(EnrolReply{Refusal: "that token cannot be used"}, 0); again || !errors.Is(err, ErrRefused) {
t.Fatalf("a refusal was not a refusal: again=%v err=%v", again, err)
}
}
func errorsMention(err error, what string) bool {
return err != nil && strings.Contains(err.Error(), what)
}