This is an automated email from the ASF dual-hosted git repository.

alfredlu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 756f9cb  [INLONG-2066] Confirm the message correctly in Go SDK (#2067)
756f9cb is described below

commit 756f9cb591056383b3fdd4c2ce85b936fffc6466
Author: Zijie Lu <[email protected]>
AuthorDate: Tue Dec 28 15:22:20 2021 +0800

    [INLONG-2066] Confirm the message correctly in Go SDK (#2067)
    
    Signed-off-by: Zijie Lu <[email protected]>
---
 .../tubemq-client-go/client/consumer_impl.go       | 13 +++---
 .../tubemq-client-go/config/config.go              | 34 +++++++++++----
 .../tubemq-client-go/config/config_test.go         | 49 ++++++++++++++++++++++
 .../tubemq-client-go/tdmsg/td_msg_test.go          |  2 +-
 4 files changed, 84 insertions(+), 14 deletions(-)

diff --git 
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go 
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go
index bc26097..ca9cca0 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go
@@ -280,7 +280,8 @@ func (c *consumer) Confirm(confirmContext string, consumed 
bool) (*ConsumerResul
                return nil, errs.New(errs.RetErrConfirmTimeout, "Not found the 
partition by confirm_context!")
        }
 
-       rsp, err := c.sendConfirmReq2Broker(partition)
+       defer c.rmtDataCache.ReleasePartition(true, 
c.subInfo.IsFiltered(topic), confirmContext, consumed)
+       rsp, err := c.sendConfirmReq2Broker(partition, consumed)
        if err != nil {
                log.Infof("[CONSUMER]Confirm error %s", err.Error())
                return nil, err
@@ -290,7 +291,8 @@ func (c *consumer) Confirm(confirmContext string, consumed 
bool) (*ConsumerResul
                BrokerHost:   partition.GetBroker().GetHost(),
                PartitionID:  uint32(partition.GetPartitionID()),
                PartitionKey: partition.GetPartitionKey(),
-               CurrOffset:   util.InvalidValue,
+               CurrOffset:   rsp.GetCurrOffset(),
+               MaxOffset:    rsp.GetMaxOffset(),
        }
        cs := &ConsumerResult{
                TopicName: partition.GetTopic(),
@@ -299,13 +301,11 @@ func (c *consumer) Confirm(confirmContext string, 
consumed bool) (*ConsumerResul
        if !rsp.GetSuccess() {
                return cs, errs.New(rsp.GetErrCode(), rsp.GetErrMsg())
        }
-       currOffset := rsp.GetCurrOffset()
-       c.rmtDataCache.BookPartitionInfo(partitionKey, currOffset, 
util.InvalidValue)
-       err = c.rmtDataCache.ReleasePartition(true, 
c.subInfo.IsFiltered(topic), confirmContext, consumed)
+       c.rmtDataCache.BookPartitionInfo(partitionKey, rsp.GetCurrOffset(), 
rsp.GetMaxOffset())
        return cs, err
 }
 
-func (c *consumer) sendConfirmReq2Broker(partition *metadata.Partition) 
(*protocol.CommitOffsetResponseB2C, error) {
+func (c *consumer) sendConfirmReq2Broker(partition *metadata.Partition, 
consumed bool) (*protocol.CommitOffsetResponseB2C, error) {
        m := &metadata.Metadata{}
        node := &metadata.Node{}
        node.SetHost(util.GetLocalHost())
@@ -313,6 +313,7 @@ func (c *consumer) sendConfirmReq2Broker(partition 
*metadata.Partition) (*protoc
        m.SetNode(node)
        sub := &metadata.SubscribeInfo{}
        sub.SetGroup(c.config.Consumer.Group)
+       partition.SetLastConsumed(consumed)
        sub.SetPartition(partition)
        m.SetSubscribeInfo(sub)
 
diff --git 
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go 
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go
index bd54505..f419d1a 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go
@@ -32,13 +32,16 @@ import (
 )
 
 const (
-       MaxRPCTimeout      = 300000 * time.Millisecond
-       MinRPCTimeout      = 8000 * time.Millisecond
-       MaxSessionKeyLen   = 1024
-       MaxGroupLen        = 1024
-       MaxTopicLen        = 64
-       MaxFilterLen       = 256
-       MaxFilterItemCount = 500
+       MaxRPCTimeout              = 300000 * time.Millisecond
+       MinRPCTimeout              = 8000 * time.Millisecond
+       MaxSessionKeyLen           = 1024
+       MaxGroupLen                = 1024
+       MaxTopicLen                = 64
+       MaxFilterLen               = 256
+       MaxFilterItemCount         = 500
+       ConsumeFromFirstOffset     = -1
+       ConsumeFromLatestOffset    = 0
+       ConsumeFromMaxOffsetAlways = 1
 )
 
 // Config defines multiple configuration options.
@@ -196,6 +199,10 @@ func (c *Config) ValidateConsumer() error {
                }
        }
 
+       if err := c.validateConsumePosition(); err != nil {
+               return err
+       }
+
        err := c.validateMaster()
        if err != nil {
                return errs.New(errs.RetInvalidConfig, err.Error())
@@ -302,6 +309,13 @@ func (c *Config) validateGroup(group string) error {
        return nil
 }
 
+func (c *Config) validateConsumePosition() error {
+       if c.Consumer.ConsumePosition < ConsumeFromFirstOffset || 
c.Consumer.ConsumePosition > ConsumeFromMaxOffsetAlways {
+               return errs.New(errs.RetInvalidConfig, 
fmt.Sprintf("consumePosition should be only in (-1, 0, 1), while %d is passed", 
c.Consumer.ConsumePosition))
+       }
+       return nil
+}
+
 // ParseAddress parses the address to user-defined config.
 func ParseAddress(address string) (config *Config, err error) {
        c := NewDefaultConfig()
@@ -545,3 +559,9 @@ func WithBoundConsume(sessionKey string, sourceCount int, 
selectBig bool, partOf
                c.Consumer.PartitionOffset = partOffset
        }
 }
+
+func WithConsumePosition(consumePosition int) Option {
+       return func(c *Config) {
+               c.Consumer.ConsumePosition = consumePosition
+       }
+}
diff --git 
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go 
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go
index c284935..9d5e910 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go
@@ -67,6 +67,40 @@ func TestParseAddress(t *testing.T) {
        address = "127.0.0.1:9092,127.0.0.1:9093?filters=12312323&filters=1212"
        _, err = ParseAddress(address)
        assert.NotNil(t, err)
+
+       address = 
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=-1"
+       c, err = ParseAddress(address)
+       assert.Nil(t, err)
+       err = c.ValidateConsumer()
+       assert.Nil(t, err)
+
+       address = 
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=1"
+       c, err = ParseAddress(address)
+       assert.Nil(t, err)
+       err = c.ValidateConsumer()
+       assert.Nil(t, err)
+
+       address = 
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=0"
+       c, err = ParseAddress(address)
+       assert.Nil(t, err)
+       err = c.ValidateConsumer()
+       assert.Nil(t, err)
+
+       address = 
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=-2"
+       c, err = ParseAddress(address)
+       assert.Nil(t, err)
+       err = c.ValidateConsumer()
+       assert.NotNil(t, err)
+
+       address = 
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=2"
+       c, err = ParseAddress(address)
+       assert.Nil(t, err)
+       err = c.ValidateConsumer()
+       assert.NotNil(t, err)
+
+       address = 
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=a"
+       c, err = ParseAddress(address)
+       assert.NotNil(t, err)
 }
 
 func TestValidateGroup(t *testing.T) {
@@ -131,4 +165,19 @@ func TestValidateConsumer(t *testing.T) {
        partitionOffset = map[string]int64{"181895251:topic1:1": 0, 
"181895251:topic2:2": 10}
        WithBoundConsume("11", 0, true, partitionOffset)(c)
        assert.Nil(t, c.ValidateConsumer())
+
+       WithConsumePosition(2)(c)
+       assert.NotNil(t, c.ValidateConsumer())
+
+       WithConsumePosition(-2)(c)
+       assert.NotNil(t, c.ValidateConsumer())
+
+       WithConsumePosition(-1)(c)
+       assert.Nil(t, c.ValidateConsumer())
+
+       WithConsumePosition(0)(c)
+       assert.Nil(t, c.ValidateConsumer())
+
+       WithConsumePosition(1)(c)
+       assert.Nil(t, c.ValidateConsumer())
 }
diff --git 
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go 
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go
index 1e8408d..70005f6 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go
@@ -34,7 +34,7 @@ func TestTDMsgV4(t *testing.T) {
 }
 
 func TestTDMsgV1(t *testing.T) {
-       b := 
[]byte{15,1,0,0,1,125,71,98,161,138,0,0,0,1,0,206,100,116,61,49,54,51,55,53,56,48,49,56,53,57,56,55,38,109,115,103,85,85,73,68,61,100,98,52,97,51,101,51,100,45,50,100,101,55,45,52,99,102,102,45,56,54,97,101,45,98,53,55,52,55,101,57,49,98,51,101,53,38,99,110,116,61,49,38,109,116,61,112,98,38,78,111,100,101,73,80,61,49,49,46,49,53,52,46,50,48,57,46,49,55,57,38,114,116,61,49,54,51,55,53,56,48,49,56,53,57,57,52,38,109,61,57,38,116,105,100,61,116,95,115,110,103,95,103,100,116,95,117,110
 [...]
+       b := []byte{15, 1, 0, 0, 1, 125, 71, 98, 161, 138, 0, 0, 0, 1, 0, 206, 
100, 116, 61, 49, 54, 51, 55, 53, 56, 48, 49, 56, 53, 57, 56, 55, 38, 109, 115, 
103, 85, 85, 73, 68, 61, 100, 98, 52, 97, 51, 101, 51, 100, 45, 50, 100, 101, 
55, 45, 52, 99, 102, 102, 45, 56, 54, 97, 101, 45, 98, 53, 55, 52, 55, 101, 57, 
49, 98, 51, 101, 53, 38, 99, 110, 116, 61, 49, 38, 109, 116, 61, 112, 98, 38, 
78, 111, 100, 101, 73, 80, 61, 49, 49, 46, 49, 53, 52, 46, 50, 48, 57, 46, 49, 
55, 57, 38, 114, 116, 61, [...]
        tm, err := New(b)
        assert.Equal(t, uint64(1637580185994), tm.CreateTime)
        assert.Equal(t, int32(1), tm.Version)

Reply via email to