The controller acted only on what its events consumer handed it, so a merge the bus skipped left modules behind with nothing said. The stream is now read back every five minutes on a single-filter consumer, and any merge that would still move a module after ten minutes is said and acted on.
80 lines
2.8 KiB
Go
80 lines
2.8 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
)
|
|
|
|
// AnnouncedMerge is one merge the forge announced, as the events stream holds it: what it said, and
|
|
// when the bus took it.
|
|
type AnnouncedMerge struct {
|
|
SourceMoved
|
|
// At is when the bus took the announcement — the stream's own time, not the forge's.
|
|
At time.Time
|
|
// Seq is its place in the events stream.
|
|
Seq uint64
|
|
}
|
|
|
|
// MergedSubject is where the forge's merges land: the controller's own follow of them.
|
|
var MergedSubject = broker.ControllerFollows[3]
|
|
|
|
// readQuiet is how long a read of the stream waits for one more message before it takes the stream
|
|
// as read to its end. The stream answers at once when it holds something; this is only the wait at
|
|
// the end.
|
|
const readQuiet = 2 * time.Second
|
|
|
|
// AnnouncedMerges is every merge the forge announced since a moment, oldest first, read back from
|
|
// the events stream (novox/hq issue 266).
|
|
//
|
|
// **Read on a consumer of its own, filtered on the one subject.** The controller's durable consumer
|
|
// carries several filters, and the bus server the mesh ran when this was written (2.10) skips a
|
|
// message now and then on a consumer with more than one filter: it moves its delivered pointer past
|
|
// the message without ever handing it over, so the controller never hears of it and nothing says so.
|
|
// A consumer filtered on a single subject reads the stream another way and was not seen to skip. This
|
|
// one is ordered, ephemeral and acknowledges nothing, so reading it changes nothing on the bus.
|
|
func (s *Server) AnnouncedMerges(ctx context.Context, since time.Time) ([]AnnouncedMerge, error) {
|
|
if s.js == nil {
|
|
return nil, errors.New("this control plane is not on the bus, so it cannot read what the forge announced")
|
|
}
|
|
sub, err := s.js.Context().SubscribeSync(MergedSubject, nats.OrderedConsumer(), nats.StartTime(since))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("reading the forge's merges from the events stream: %w", err)
|
|
}
|
|
defer func() { _ = sub.Unsubscribe() }()
|
|
|
|
var out []AnnouncedMerge
|
|
for {
|
|
wait, cancel := context.WithTimeout(ctx, readQuiet)
|
|
msg, err := sub.NextMsgWithContext(wait)
|
|
cancel()
|
|
if err != nil {
|
|
if ctx.Err() != nil {
|
|
return nil, ctx.Err()
|
|
}
|
|
// Nothing more within the quiet wait: the stream has been read to its end.
|
|
break
|
|
}
|
|
meta, err := msg.Metadata()
|
|
if err != nil {
|
|
continue
|
|
}
|
|
var moved SourceMoved
|
|
if json.Unmarshal(msg.Data, &moved) == nil && moved.Commit != "" && moved.Repo != "" {
|
|
out = append(out, AnnouncedMerge{SourceMoved: moved, At: meta.Timestamp, Seq: meta.Sequence.Stream})
|
|
}
|
|
if meta.NumPending == 0 {
|
|
break
|
|
}
|
|
}
|
|
sort.Slice(out, func(i, j int) bool { return out[i].Seq < out[j].Seq })
|
|
return out, nil
|
|
}
|