Skip to content

Commit ea9680a

Browse files
Cherry-picks for 2.11.12 (#7776)
Includes the following: - #7772 - #7773 - #7769 - #7766 Signed-off-by: Neil Twigg <[email protected]>
2 parents 6f77800 + eb53e0d commit ea9680a

9 files changed

Lines changed: 344 additions & 127 deletions

server/consumer.go

Lines changed: 17 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -2304,7 +2304,8 @@ func (o *consumer) updateConfig(cfg *ConsumerConfig) error {
23042304

23052305
// Check for Subject Filters update.
23062306
newSubjects := gatherSubjectFilters(cfg.FilterSubject, cfg.FilterSubjects)
2307-
if !subjectSliceEqual(newSubjects, o.subjf.subjects()) {
2307+
updatedFilters := !subjectSliceEqual(newSubjects, o.subjf.subjects())
2308+
if updatedFilters {
23082309
newSubjf := make(subjectFilters, 0, len(newSubjects))
23092310
for _, newFilter := range newSubjects {
23102311
fs := &subjectFilter{
@@ -2339,13 +2340,6 @@ func (o *consumer) updateConfig(cfg *ConsumerConfig) error {
23392340
}
23402341
}
23412342

2342-
// Check if any filters were updated.
2343-
oldFilters := gatherSubjectFilters(o.cfg.FilterSubject, o.cfg.FilterSubjects)
2344-
newFilters := gatherSubjectFilters(cfg.FilterSubject, cfg.FilterSubjects)
2345-
slices.Sort(oldFilters)
2346-
slices.Sort(newFilters)
2347-
updatedFilters := !slices.Equal(oldFilters, newFilters)
2348-
23492343
// Record new config for others that do not need special handling.
23502344
// Allowed but considered no-op, [Description, SampleFrequency, MaxWaiting, HeadersOnly]
23512345
o.cfg = *cfg
@@ -6284,6 +6278,10 @@ func (o *consumer) checkStateForInterestStream(ss *StreamState) error {
62846278
retryAsflr = seq
62856279
}
62866280
} else if seq <= dflr {
6281+
// Store the first entry above our ack floor, so we don't need to look it up again on retryAsflr=0.
6282+
if retryAsflr == 0 {
6283+
retryAsflr = seq
6284+
}
62876285
// If we have pending, we will need to walk through to delivered in case we missed any of those acks as well.
62886286
if _, ok := state.Pending[seq]; !ok {
62896287
// The filters are already taken into account,
@@ -6295,8 +6293,18 @@ func (o *consumer) checkStateForInterestStream(ss *StreamState) error {
62956293
}
62966294
}
62976295
// If retry floor was not overwritten, set to ack floor+1, we don't need to account for any retries below it.
6296+
// However, our ack floor may be lower than the next message we can receive, so we correct it upward if needed.
62986297
if retryAsflr == 0 {
6299-
retryAsflr = asflr + 1
6298+
if filters != nil {
6299+
_, nseq, err = store.LoadNextMsgMulti(filters, asflr+1, &smv)
6300+
} else {
6301+
_, nseq, err = store.LoadNextMsg(filter, wc, asflr+1, &smv)
6302+
}
6303+
if err == nil {
6304+
retryAsflr = max(asflr+1, nseq)
6305+
} else if err == ErrStoreEOF {
6306+
retryAsflr = ss.LastSeq + 1
6307+
}
63006308
}
63016309

63026310
o.mu.Lock()

server/filestore.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8776,6 +8776,11 @@ func (fs *fileStore) compact(seq uint64) (uint64, error) {
87768776
fs.mu.Unlock()
87778777
return fs.purge(seq)
87788778
}
8779+
// Short-circuit if the store was already compacted past this point.
8780+
if fs.state.FirstSeq > seq {
8781+
fs.mu.Unlock()
8782+
return purged, nil
8783+
}
87798784
// We have to delete interior messages.
87808785
smb := fs.selectMsgBlock(seq)
87818786
if smb == nil {

server/jetstream_cluster_1_test.go

Lines changed: 112 additions & 97 deletions
Original file line numberDiff line numberDiff line change
@@ -9957,119 +9957,134 @@ func TestJetStreamClusterJszRaftLeaderReporting(t *testing.T) {
99579957
}
99589958

99599959
func TestJetStreamClusterNoInterestDesyncOnConsumerCreate(t *testing.T) {
9960-
c := createJetStreamClusterExplicit(t, "R3S", 3)
9961-
defer c.shutdown()
9960+
test := func(t *testing.T, twoConsumers bool) {
9961+
c := createJetStreamClusterExplicit(t, "R3S", 3)
9962+
defer c.shutdown()
99629963

9963-
nc, js := jsClientConnect(t, c.randomServer())
9964-
defer nc.Close()
9964+
nc, js := jsClientConnect(t, c.randomServer())
9965+
defer nc.Close()
99659966

9966-
_, err := js.AddStream(&nats.StreamConfig{
9967-
Name: "TEST",
9968-
Subjects: []string{"foo"},
9969-
Replicas: 3,
9970-
Retention: nats.InterestPolicy,
9971-
})
9972-
require_NoError(t, err)
9967+
_, err := js.AddStream(&nats.StreamConfig{
9968+
Name: "TEST",
9969+
Subjects: []string{"foo", "bar"},
9970+
Replicas: 3,
9971+
Retention: nats.InterestPolicy,
9972+
})
9973+
require_NoError(t, err)
99739974

9974-
// Pick a random server that will not know about the new consumer being created.
9975-
// If servers determine "no interest" individually, these servers will desync.
9976-
rs := c.randomNonLeader()
9977-
sjs := rs.getJetStream()
9978-
meta := sjs.getMetaGroup()
9979-
require_NoError(t, meta.PauseApply())
9975+
// Pick a random server that will not know about the new consumer being created.
9976+
// If servers determine "no interest" individually, these servers will desync.
9977+
rs := c.randomNonLeader()
9978+
sjs := rs.getJetStream()
9979+
meta := sjs.getMetaGroup()
9980+
require_NoError(t, meta.PauseApply())
99809981

9981-
sub, err := js.PullSubscribe(_EMPTY_, "DURABLE", nats.BindStream("TEST"))
9982-
require_NoError(t, err)
9983-
defer sub.Drain()
9982+
sub, err := js.PullSubscribe("foo", "DURABLE", nats.BindStream("TEST"))
9983+
require_NoError(t, err)
9984+
defer sub.Drain()
99849985

9985-
checkConsumersAssigned := func(expected int) {
9986-
t.Helper()
9987-
checkFor(t, 2*time.Second, 200*time.Millisecond, func() error {
9988-
var count int
9989-
for _, s := range c.servers {
9990-
_, _, jsa := s.globalAccount().getJetStreamFromAccount()
9991-
if jsa.consumerAssigned("TEST", "DURABLE") {
9992-
count++
9986+
checkConsumersAssigned := func(expected int) {
9987+
t.Helper()
9988+
checkFor(t, 2*time.Second, 200*time.Millisecond, func() error {
9989+
var count int
9990+
for _, s := range c.servers {
9991+
mset, err := s.globalAccount().lookupStream("TEST")
9992+
if err != nil {
9993+
return err
9994+
}
9995+
count += mset.numConsumers()
99939996
}
9994-
}
9995-
if count != expected {
9996-
return fmt.Errorf("expected %d, got %d", expected, count)
9997-
}
9998-
return nil
9999-
})
10000-
}
10001-
// Confirm only two servers know about the consumer.
10002-
checkConsumersAssigned(2)
10003-
c.waitOnConsumerLeader(globalAccountName, "TEST", "DURABLE")
9997+
if count != expected {
9998+
return fmt.Errorf("expected %d, got %d", expected, count)
9999+
}
10000+
return nil
10001+
})
10002+
}
10003+
// Confirm only two servers know about the consumer.
10004+
checkConsumersAssigned(2)
10005+
c.waitOnConsumerLeader(globalAccountName, "TEST", "DURABLE")
10006+
10007+
// Publish a single message. All servers will receive this, but only two will store it.
10008+
_, err = js.Publish("foo", nil)
10009+
require_NoError(t, err)
10010+
checkLastSeq := func(lseq uint64) {
10011+
t.Helper()
10012+
checkFor(t, 2*time.Second, 200*time.Millisecond, func() error {
10013+
for _, s := range c.servers {
10014+
mset, err := s.globalAccount().lookupStream("TEST")
10015+
if err != nil {
10016+
return err
10017+
}
10018+
if seq := mset.lastSeq(); seq != lseq {
10019+
return fmt.Errorf("expected %d, got %d", lseq, seq)
10020+
}
10021+
}
10022+
return nil
10023+
})
10024+
}
10025+
checkLastSeq(1)
10026+
10027+
// Resume the meta layer such that the consumer gets created on the remaining server.
10028+
meta.ResumeApply()
10029+
checkConsumersAssigned(3)
10030+
10031+
if twoConsumers {
10032+
_, err = js.AddConsumer("TEST", &nats.ConsumerConfig{Durable: "DURABLE2", FilterSubject: "bar"})
10033+
require_NoError(t, err)
10034+
checkConsumersAssigned(6)
10035+
}
10036+
10037+
// All servers will now store another published message.
10038+
_, err = js.Publish("foo", nil)
10039+
require_NoError(t, err)
10040+
checkLastSeq(2)
10041+
10042+
// Make sure the consumer leader is on the same server that didn't store the first message.
10043+
cl := c.consumerLeader(globalAccountName, "TEST", "DURABLE")
10044+
if cl != rs {
10045+
mset, err := cl.globalAccount().lookupStream("TEST")
10046+
require_NoError(t, err)
10047+
o := mset.lookupConsumer("DURABLE")
10048+
require_NotNil(t, o)
10049+
n := o.raftNode()
10050+
require_NoError(t, n.StepDown(rs.NodeName()))
10051+
c.waitOnConsumerLeader(globalAccountName, "TEST", "DURABLE")
10052+
cl = c.consumerLeader(globalAccountName, "TEST", "DURABLE")
10053+
require_Equal(t, cl, rs)
10054+
}
10055+
10056+
// Since the consumer leader is the same as the server that didn't store the first message,
10057+
// it can only receive and ack the second message.
10058+
msgs, err := sub.Fetch(1, nats.MaxWait(time.Second))
10059+
require_NoError(t, err)
10060+
require_Len(t, len(msgs), 1)
10061+
metadata, err := msgs[0].Metadata()
10062+
require_NoError(t, err)
10063+
require_Equal(t, metadata.Sequence.Stream, 2)
10064+
require_Equal(t, metadata.NumPending, 0)
10065+
require_NoError(t, msgs[0].AckSync())
10066+
10067+
if twoConsumers {
10068+
require_NoError(t, js.DeleteConsumer("TEST", "DURABLE"))
10069+
}
1000410070

10005-
// Publish a single message. All servers will receive this, but only two will store it.
10006-
_, err = js.Publish("foo", nil)
10007-
require_NoError(t, err)
10008-
checkLastSeq := func(lseq uint64) {
10009-
t.Helper()
1001010071
checkFor(t, 2*time.Second, 200*time.Millisecond, func() error {
10072+
// The servers will eventually be synced up again, but this relies on the interest state being checked.
1001110073
for _, s := range c.servers {
10074+
if s == rs {
10075+
continue
10076+
}
1001210077
mset, err := s.globalAccount().lookupStream("TEST")
1001310078
if err != nil {
1001410079
return err
1001510080
}
10016-
if seq := mset.lastSeq(); seq != lseq {
10017-
return fmt.Errorf("expected %d, got %d", lseq, seq)
10018-
}
10081+
mset.checkInterestState()
1001910082
}
10020-
return nil
10083+
return checkState(t, c, globalAccountName, "TEST")
1002110084
})
1002210085
}
10023-
checkLastSeq(1)
10024-
10025-
// Resume the meta layer such that the consumer gets created on the remaining server.
10026-
meta.ResumeApply()
10027-
checkConsumersAssigned(3)
10028-
10029-
// All servers will now store another published message.
10030-
_, err = js.Publish("foo", nil)
10031-
require_NoError(t, err)
10032-
checkLastSeq(2)
10033-
10034-
// Make sure the consumer leader is on the same server that didn't store the first message.
10035-
cl := c.consumerLeader(globalAccountName, "TEST", "DURABLE")
10036-
if cl != rs {
10037-
mset, err := cl.globalAccount().lookupStream("TEST")
10038-
require_NoError(t, err)
10039-
o := mset.lookupConsumer("DURABLE")
10040-
require_NotNil(t, o)
10041-
n := o.raftNode()
10042-
require_NoError(t, n.StepDown(rs.NodeName()))
10043-
c.waitOnConsumerLeader(globalAccountName, "TEST", "DURABLE")
10044-
cl = c.consumerLeader(globalAccountName, "TEST", "DURABLE")
10045-
require_Equal(t, cl, rs)
10046-
}
10047-
10048-
// Since the consumer leader is the same as the server that didn't store the first message,
10049-
// it can only receive and ack the second message.
10050-
msgs, err := sub.Fetch(1, nats.MaxWait(time.Second))
10051-
require_NoError(t, err)
10052-
require_Len(t, len(msgs), 1)
10053-
metadata, err := msgs[0].Metadata()
10054-
require_NoError(t, err)
10055-
require_Equal(t, metadata.Sequence.Stream, 2)
10056-
require_Equal(t, metadata.NumPending, 0)
10057-
require_NoError(t, msgs[0].AckSync())
10058-
10059-
checkFor(t, 2*time.Second, 200*time.Millisecond, func() error {
10060-
// The servers will eventually be synced up again, but this relies on the interest state being checked.
10061-
for _, s := range c.servers {
10062-
if s == rs {
10063-
continue
10064-
}
10065-
mset, err := s.globalAccount().lookupStream("TEST")
10066-
if err != nil {
10067-
return err
10068-
}
10069-
mset.checkInterestState()
10070-
}
10071-
return checkState(t, c, globalAccountName, "TEST")
10072-
})
10086+
t.Run("OneConsumer", func(t *testing.T) { test(t, false) })
10087+
t.Run("TwoConsumers", func(t *testing.T) { test(t, true) })
1007310088
}
1007410089

1007510090
func TestJetStreamClusterMetaPeerRemoveResponseAfterQuorum(t *testing.T) {

0 commit comments

Comments
 (0)