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 }