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..>` 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) }