A host reaches either bus, and a genesis template that raises the mesh on the new one #33
@@ -1,9 +1,13 @@
|
||||
module github.com/novox/mesh-host
|
||||
|
||||
go 1.25.0
|
||||
go 1.26.0
|
||||
|
||||
require (
|
||||
github.com/klauspost/compress v1.20.0 // indirect
|
||||
github.com/nats-io/nats.go v1.54.0 // indirect
|
||||
github.com/nats-io/nkeys v0.4.16 // indirect
|
||||
github.com/nats-io/nuid v1.0.1 // indirect
|
||||
github.com/rabbitmq/amqp091-go v1.14.0 // indirect
|
||||
golang.org/x/crypto v0.55.0 // indirect
|
||||
golang.org/x/sys v0.47.0 // indirect
|
||||
golang.org/x/crypto v0.57.0 // indirect
|
||||
golang.org/x/sys v0.48.0 // indirect
|
||||
)
|
||||
|
||||
@@ -1,6 +1,18 @@
|
||||
github.com/klauspost/compress v1.20.0 h1:a3C1ke2ohxFymNlb2HWAHjDeKCI90scRskErZkR0ezA=
|
||||
github.com/klauspost/compress v1.20.0/go.mod h1:LUdAzn7YLVvxLpc7y3V1m40wESHTgc1422pwwBSKYuI=
|
||||
github.com/nats-io/nats.go v1.54.0 h1:vsXoOxjHp/GmPUN+EcI7uOf/uB+iAP+kEsAFNQN0yzA=
|
||||
github.com/nats-io/nats.go v1.54.0/go.mod h1:y+DZoD1oBOYfZTU681eTUiUjI0vbqYGixNVFHcjHJ0k=
|
||||
github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg=
|
||||
github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs=
|
||||
github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw=
|
||||
github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c=
|
||||
github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek=
|
||||
github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o=
|
||||
golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M=
|
||||
golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis=
|
||||
golang.org/x/crypto v0.57.0 h1:3ZVCjf8Ggz7zneR/EHRVx68Ctf+2pmIMP2UFhh9cC6M=
|
||||
golang.org/x/crypto v0.57.0/go.mod h1:Fdz0i5U6CoizGwLda9DttjSk6qlZo25zYNtR+ycvuZA=
|
||||
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/sys v0.48.0 h1:bbX/i/6MgT9BVLM9RT1thmxL04yeTAhbEz4SyadbXoo=
|
||||
golang.org/x/sys v0.48.0/go.mod h1:hNLxWAXmnKAxqDtdwIYC4bM9oQPEecfsnNMuSxOs3og=
|
||||
|
||||
@@ -0,0 +1,88 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
)
|
||||
|
||||
// Bus is what a host needs of the mesh's bus, in the mesh's own words.
|
||||
//
|
||||
// A host says exactly two things unprompted: what it applied, and that it is here. They are not
|
||||
// the same kind of statement and the difference is the whole of this interface — one must arrive
|
||||
// and one must not be insisted on.
|
||||
//
|
||||
// **The host still imports nothing of the mesh's own** (novox/hq ADR 0005): this is its own
|
||||
// interface over its own client libraries, not a contract shared with the controller. The two
|
||||
// agree because a conformance fixture holds them to one envelope, which is the only kind of
|
||||
// agreement that survives being in different repositories.
|
||||
type Bus interface {
|
||||
// Report says what this node applied. **It must arrive.** A report that fails leaves the
|
||||
// mesh believing the node never answered while the node believes it did, and the two go on
|
||||
// disagreeing with nothing anywhere saying so — the shape of fault this project keeps
|
||||
// finding. Returns false when it could not be delivered, so the caller can say so.
|
||||
Report(ctx context.Context, node string, body []byte) error
|
||||
|
||||
// Alive says this node is here, and nothing else. **Losing one is nothing**: the next is a
|
||||
// minute away and the mesh reads a gap rather than counting arrivals. Insisting on delivery
|
||||
// would turn a harmless miss into a logged failure every minute.
|
||||
Alive(ctx context.Context, node string, body []byte) error
|
||||
}
|
||||
|
||||
// --- The bus the mesh runs on today -----------------------------------------------------------
|
||||
|
||||
// OverCurrent is the bus as a channel, until the rollout.
|
||||
type OverCurrent struct{ Channel *amqp.Channel }
|
||||
|
||||
func (b OverCurrent) Report(ctx context.Context, node string, body []byte) error {
|
||||
// Mandatory: an unroutable report comes back rather than disappearing.
|
||||
return b.Channel.PublishWithContext(ctx, Exchange, KeyReport, true, false,
|
||||
amqp.Publishing{ContentType: "application/json", Body: body})
|
||||
}
|
||||
|
||||
func (b OverCurrent) Alive(ctx context.Context, node string, body []byte) error {
|
||||
return b.Channel.PublishWithContext(ctx, Exchange, KeyAlive, false, false,
|
||||
amqp.Publishing{ContentType: "application/json", Body: body})
|
||||
}
|
||||
|
||||
// --- NATS ---------------------------------------------------------------------------------
|
||||
|
||||
// OverNATS is the bus as a connection. A report goes through JetStream because it must survive
|
||||
// the controller's store restarting; a heartbeat does not, because it must not.
|
||||
type OverNATS struct {
|
||||
Conn *nats.Conn
|
||||
JS nats.JetStreamContext
|
||||
}
|
||||
|
||||
// ReportSubject and AliveSubject are this node's own, and no other node's: a host's account may
|
||||
// publish `mesh.control.<its own node>.>` and nothing wider, so the subject is the authority on
|
||||
// which node a report is about.
|
||||
func ReportSubject(node string) string { return "mesh.control." + node + ".report" }
|
||||
func AliveSubject(node string) string { return "mesh.control." + node + ".alive" }
|
||||
|
||||
func (b OverNATS) Report(ctx context.Context, node string, body []byte) error {
|
||||
// Into the CONTROL stream and awaited: this is the message the store-window guarantee is
|
||||
// about (novox/hq ADR 0083). The controller naks with a delay while its store is away and
|
||||
// the message is redelivered; a publish the bus never accepted must fail here rather than
|
||||
// be assumed.
|
||||
if _, err := b.JS.Publish(ReportSubject(node), body, nats.Context(ctx)); err != nil {
|
||||
return fmt.Errorf("reporting: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b OverNATS) Alive(ctx context.Context, node string, body []byte) error {
|
||||
// Core, deliberately: a heartbeat in a stream is the mesh's least valuable message competing
|
||||
// for retention with its most valuable, and a lost one is the next one.
|
||||
if err := b.Conn.Publish(AliveSubject(node), body); err != nil {
|
||||
return err
|
||||
}
|
||||
// Flushed rather than fired and forgotten, so "could not tell the mesh" means the write
|
||||
// failed rather than that nobody has looked yet.
|
||||
flush, cancel := context.WithTimeout(ctx, 2*time.Second)
|
||||
defer cancel()
|
||||
return b.Conn.FlushWithContext(flush)
|
||||
}
|
||||
+22
-12
@@ -247,7 +247,7 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
// report, and reports are rare where this is constant.
|
||||
beat := time.NewTicker(AliveEvery)
|
||||
defer beat.Stop()
|
||||
publishAlive(ctx, channel, m, say, timeout)
|
||||
publishAlive(ctx, OverCurrent{Channel: channel}, m, say, timeout)
|
||||
|
||||
closed := conn.NotifyClose(make(chan *amqp.Error, 1))
|
||||
|
||||
@@ -268,11 +268,11 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
case <-ctx.Done():
|
||||
return nil
|
||||
case <-beat.C:
|
||||
publishAlive(ctx, channel, m, say, timeout)
|
||||
publishAlive(ctx, OverCurrent{Channel: channel}, m, say, timeout)
|
||||
case unasked := <-outbox:
|
||||
// Said without having been asked: a reconcile found what an adopted node holds, or
|
||||
// its firewall, changed since it last said.
|
||||
published := publishReport(ctx, channel, m, unasked.Report, say, timeout)
|
||||
published := publishReport(ctx, OverCurrent{Channel: channel}, m, unasked.Report, say, timeout)
|
||||
if unasked.Done != nil {
|
||||
unasked.Done(published)
|
||||
}
|
||||
@@ -287,7 +287,7 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
delivery, superseded := newest(deliveries, delivery, drainWindow)
|
||||
for _, old := range superseded {
|
||||
say("set aside a declaration: a newer one arrived with it")
|
||||
publishReport(ctx, channel, m, Report{Node: m.Node, Declared: declaredIn(old.Body),
|
||||
publishReport(ctx, OverCurrent{Channel: channel}, m, Report{Node: m.Node, Declared: declaredIn(old.Body),
|
||||
Superseded: declaredIn(delivery.Body)}, say, timeout)
|
||||
_ = old.Ack(false)
|
||||
}
|
||||
@@ -300,7 +300,7 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
default:
|
||||
say(fmt.Sprintf("applied %d resource(s)", len(report.Applied)))
|
||||
}
|
||||
publishReport(ctx, channel, m, report, say, timeout)
|
||||
publishReport(ctx, OverCurrent{Channel: channel}, m, report, say, timeout)
|
||||
// Acknowledged after the report is published. A node that dies between applying and
|
||||
// reporting leaves the declaration on the broker and applies it again on return,
|
||||
// which is safe because applying is reconciliation — it converges rather than
|
||||
@@ -321,6 +321,18 @@ const (
|
||||
// newest takes what is already waiting behind `first` and returns the last of them to apply, and
|
||||
// the rest to set aside. It waits `window` for a straggler after each arrival and no longer: a
|
||||
// declaration in flight from the mesh arrives within that; one that does not is the next push.
|
||||
//
|
||||
// **Its job narrows once declarations are state rather than messages, and does not disappear.**
|
||||
// On the bus being built, a declaration is last-per-subject (novox/hq design 29 §4), so a node
|
||||
// that was away receives exactly the current one instead of a queue of superseded ones — the
|
||||
// catch-up half of what this does is then the stream's. And a stream sequence orders them
|
||||
// definitively, where this window only infers order from arrival time, which is the wire-level
|
||||
// answer to novox/hq issue 107.
|
||||
//
|
||||
// What remains is the live case: three pushes in quick succession to a *connected* node are
|
||||
// three deliveries, whatever the stream later retains. So this is narrowed at the rollout, not
|
||||
// deleted — and saying which half goes is worth more than a note that it "can probably be
|
||||
// removed", which is how a load-bearing window gets deleted by somebody in a hurry.
|
||||
func newest(deliveries <-chan amqp.Delivery, first amqp.Delivery, window time.Duration) (amqp.Delivery, []amqp.Delivery) {
|
||||
latest := first
|
||||
var superseded []amqp.Delivery
|
||||
@@ -404,13 +416,13 @@ func Publish(ctx context.Context, m Membership, report Report, timeout time.Dura
|
||||
}
|
||||
defer channel.Close()
|
||||
var said string
|
||||
if !publishReport(ctx, channel, m, report, func(s string) { said = s }, timeout) {
|
||||
if !publishReport(ctx, OverCurrent{Channel: channel}, m, report, func(s string) { said = s }, timeout) {
|
||||
return errors.New(said)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report,
|
||||
func publishReport(ctx context.Context, bus Bus, m Membership, report Report,
|
||||
say Announce, timeout time.Duration) bool {
|
||||
report.Node = m.Node
|
||||
body, err := json.Marshal(report)
|
||||
@@ -424,8 +436,7 @@ func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, rep
|
||||
// Said rather than swallowed. A report that fails to publish leaves the mesh believing this
|
||||
// node never answered, while the node believes it did — and the two would go on disagreeing
|
||||
// with nothing anywhere saying so. That shape of fault is the one this project keeps finding.
|
||||
if err := channel.PublishWithContext(publish, Exchange, KeyReport, true, false,
|
||||
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
|
||||
if err := bus.Report(publish, m.Node, body); err != nil {
|
||||
say(fmt.Sprintf("applied, and could not tell the mesh: %v", err))
|
||||
return false
|
||||
}
|
||||
@@ -433,7 +444,7 @@ func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, rep
|
||||
}
|
||||
|
||||
// publishAlive says this node is here, and nothing else.
|
||||
func publishAlive(ctx context.Context, channel *amqp.Channel, m Membership, say Announce,
|
||||
func publishAlive(ctx context.Context, bus Bus, m Membership, say Announce,
|
||||
timeout time.Duration) {
|
||||
body, err := json.Marshal(Alive{Node: m.Node})
|
||||
if err != nil {
|
||||
@@ -444,8 +455,7 @@ func publishAlive(ctx context.Context, channel *amqp.Channel, m Membership, say
|
||||
// Not mandatory, unlike a report. Losing one is nothing: the next is a minute away, and the
|
||||
// mesh is reading a gap rather than counting arrivals. Insisting on delivery would turn a
|
||||
// harmless miss into a logged failure every minute.
|
||||
if err := channel.PublishWithContext(publish, Exchange, KeyAlive, false, false,
|
||||
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
|
||||
if err := bus.Alive(publish, m.Node, body); err != nil {
|
||||
say("could not tell the mesh this node is here: " + err.Error())
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user