Step 3.5's first half, mirroring the controller's. A host says exactly two things unprompted, and the difference between them is the whole interface: a report must arrive, and a heartbeat must not be insisted on. So a report goes through JetStream — it is the message the store-window guarantee is about — and a heartbeat stays on core, because a heartbeat in a stream is the mesh's least valuable message competing for retention with its most valuable. The host still imports nothing of the mesh's own (ADR 0005): this is its own interface over its own libraries. It agrees with the controller because a fixture holds both to one envelope, which is the only agreement that survives two repositories. Also recorded, where the next person reads it rather than in a plan: the "newest wins" window narrows at the rollout and does not disappear. Last- per-subject makes the catch-up half the stream's, and sequence orders them definitively — but three pushes to a connected node are still three deliveries. Saying which half goes is worth more than "can probably be removed", which is how a load-bearing window gets deleted in a hurry.
89 lines
4.0 KiB
Go
89 lines
4.0 KiB
Go
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)
|
|
}
|