Count an unasked report as said only once the broker has taken it (hq ADR 0100)
This commit is contained in:
+19
-5
@@ -76,7 +76,15 @@ func Hold(ctx context.Context, m Membership, apply Applier, say Announce, timeou
|
||||
// Outbox carries reports the node has to say without having been sent anything — what a
|
||||
// reconcile found changed on an adopted node (novox/hq ADR 0100). Published while the link is up;
|
||||
// a report made while it is down waits in the channel for the next one. Nil is allowed.
|
||||
type Outbox <-chan Report
|
||||
type Outbox <-chan Unasked
|
||||
|
||||
// Unasked is one such report, with the way to say whether it reached the mesh. Done is called
|
||||
// with true only when the broker took it — a node that marked a change said because it queued it
|
||||
// would never say it again, and the mesh would go on believing nothing changed.
|
||||
type Unasked struct {
|
||||
Report Report
|
||||
Done func(published bool)
|
||||
}
|
||||
|
||||
// HoldRoused is Hold, told when the machine has reason to think its link is stale, and handed
|
||||
// reports to publish between deliveries.
|
||||
@@ -261,10 +269,13 @@ func Run(ctx context.Context, m Membership, apply Applier, say Announce, timeout
|
||||
return nil
|
||||
case <-beat.C:
|
||||
publishAlive(ctx, channel, m, say, timeout)
|
||||
case report := <-outbox:
|
||||
case unasked := <-outbox:
|
||||
// Said without having been asked: a reconcile found what an adopted node holds, or
|
||||
// its firewall, changed since it last said.
|
||||
publishReport(ctx, channel, m, report, say, timeout)
|
||||
published := publishReport(ctx, channel, m, unasked.Report, say, timeout)
|
||||
if unasked.Done != nil {
|
||||
unasked.Done(published)
|
||||
}
|
||||
case reason := <-closed:
|
||||
return fmt.Errorf("the link closed: %v", reason)
|
||||
case delivery, ok := <-deliveries:
|
||||
@@ -365,13 +376,14 @@ func handleBody(ctx context.Context, m Membership, body []byte, apply Applier) R
|
||||
return apply(ctx, signed.Declaration, signed.Signature)
|
||||
}
|
||||
|
||||
// publishReport tells the mesh what this node did, and says whether the broker took it.
|
||||
func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, report Report,
|
||||
say Announce, timeout time.Duration) {
|
||||
say Announce, timeout time.Duration) bool {
|
||||
report.Node = m.Node
|
||||
body, err := json.Marshal(report)
|
||||
if err != nil {
|
||||
say("cannot encode this node's own report: " + err.Error())
|
||||
return
|
||||
return false
|
||||
}
|
||||
publish, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
@@ -382,7 +394,9 @@ func publishReport(ctx context.Context, channel *amqp.Channel, m Membership, rep
|
||||
if err := channel.PublishWithContext(publish, Exchange, KeyReport, true, false,
|
||||
amqp.Publishing{ContentType: "application/json", Body: body}); err != nil {
|
||||
say(fmt.Sprintf("applied, and could not tell the mesh: %v", err))
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// publishAlive says this node is here, and nothing else.
|
||||
|
||||
Reference in New Issue
Block a user