@@ -9957,119 +9957,134 @@ func TestJetStreamClusterJszRaftLeaderReporting(t *testing.T) {
99579957}
99589958
99599959func 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
1007510090func TestJetStreamClusterMetaPeerRemoveResponseAfterQuorum (t * testing.T ) {
0 commit comments