A refused membership does not stop the controller #190
@@ -708,7 +708,11 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, na
|
|||||||
if !ok {
|
if !ok {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
issued := 0
|
// The declarations are sent and recorded by now; a membership that cannot be issued is said
|
||||||
|
// and does not unsay them. Every runtime without one serves the shape it derives (ADR 0160), so
|
||||||
|
// the push stands, the first failure is named once, and the next push tries again.
|
||||||
|
issued, failed := 0, 0
|
||||||
|
var first error
|
||||||
for _, node := range names {
|
for _, node := range names {
|
||||||
for _, d := range records.Assigned[node] {
|
for _, d := range records.Assigned[node] {
|
||||||
body, err := json.Marshal(broker.MembershipFor(node, d, where))
|
body, err := json.Marshal(broker.MembershipFor(node, d, where))
|
||||||
@@ -716,7 +720,11 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, na
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := bus.PublishMembership(ctx, node, d.Module, body); err != nil {
|
if err := bus.PublishMembership(ctx, node, d.Module, body); err != nil {
|
||||||
return err
|
if first == nil {
|
||||||
|
first = err
|
||||||
|
}
|
||||||
|
failed++
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
issued++
|
issued++
|
||||||
}
|
}
|
||||||
@@ -724,6 +732,10 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, na
|
|||||||
if issued > 0 {
|
if issued > 0 {
|
||||||
fmt.Printf(" issued %d membership(s)\n", issued)
|
fmt.Printf(" issued %d membership(s)\n", issued)
|
||||||
}
|
}
|
||||||
|
if failed > 0 {
|
||||||
|
fmt.Printf(" %d membership(s) could not be issued; the first: %v — the machines keep what "+
|
||||||
|
"they derive until the next push\n", failed, first)
|
||||||
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -148,7 +148,15 @@ func (b OverNATS) PublishSeatEvent(ctx context.Context, seat, event string, body
|
|||||||
|
|
||||||
// PublishMembership issues one assignment what it serves and reaches (novox/hq ADR 0160), last per
|
// PublishMembership issues one assignment what it serves and reaches (novox/hq ADR 0160), last per
|
||||||
// subject, so the runtime that connects later reads the current one and one that is running follows.
|
// subject, so the runtime that connects later reads the current one and one that is running follows.
|
||||||
|
// MembershipWait bounds how long issuing one membership may take. A publish the server refuses is
|
||||||
|
// never acknowledged, and a stream publish waits for its acknowledgement for as long as its
|
||||||
|
// context lives: on 2026-10-01 the daemon's own context was that long, and one refused membership
|
||||||
|
// held the controller's receive loop for good (novox/hq issue 185).
|
||||||
|
const MembershipWait = 10 * time.Second
|
||||||
|
|
||||||
func (b OverNATS) PublishMembership(ctx context.Context, node, module string, body []byte) error {
|
func (b OverNATS) PublishMembership(ctx context.Context, node, module string, body []byte) error {
|
||||||
|
ctx, cancel := context.WithTimeout(ctx, MembershipWait)
|
||||||
|
defer cancel()
|
||||||
_, err := b.JS.Publish(broker.MembershipSubject(node, module), body, nats.Context(ctx))
|
_, err := b.JS.Publish(broker.MembershipSubject(node, module), body, nats.Context(ctx))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("issuing %s on %s its membership: %w", module, node, err)
|
return fmt.Errorf("issuing %s on %s its membership: %w", module, node, err)
|
||||||
|
|||||||
Reference in New Issue
Block a user