fixup panic
	
		
			
	
		
	
	
		
	
		
			Some checks failed
		
		
	
	
		
			
				
	
				codeql / analyze (go) (pull_request) Failing after 2m42s
				
			
		
			
				
	
				prbuild / test (pull_request) Failing after 1m29s
				
			
		
			
				
	
				prbuild / lint (pull_request) Failing after 2m37s
				
			
		
			
				
	
				autoapprove / autoapprove (pull_request) Failing after 1m24s
				
			
		
			
				
	
				automerge / automerge (pull_request) Failing after 4s
				
			
		
			
				
	
				dependabot-automerge / automerge (pull_request) Has been skipped
				
			
		
		
	
	
				
					
				
			
		
			Some checks failed
		
		
	
	codeql / analyze (go) (pull_request) Failing after 2m42s
				
			prbuild / test (pull_request) Failing after 1m29s
				
			prbuild / lint (pull_request) Failing after 2m37s
				
			autoapprove / autoapprove (pull_request) Failing after 1m24s
				
			automerge / automerge (pull_request) Failing after 4s
				
			dependabot-automerge / automerge (pull_request) Has been skipped
				
			Signed-off-by: Vasiliy Tolstov <v.tolstov@unistack.org>
This commit is contained in:
		
							
								
								
									
										10
									
								
								go.mod
									
									
									
									
									
								
							
							
						
						
									
										10
									
								
								go.mod
									
									
									
									
									
								
							| @@ -3,12 +3,12 @@ module go.unistack.org/micro-broker-kgo/v3 | ||||
| go 1.17 | ||||
|  | ||||
| require ( | ||||
| 	github.com/twmb/franz-go v1.11.5 | ||||
| 	github.com/twmb/franz-go/pkg/kmsg v1.3.0 | ||||
| 	go.unistack.org/micro/v3 v3.10.14 | ||||
| 	github.com/twmb/franz-go v1.12.1 | ||||
| 	github.com/twmb/franz-go/pkg/kmsg v1.4.0 | ||||
| 	go.unistack.org/micro/v3 v3.10.23 | ||||
| ) | ||||
|  | ||||
| require ( | ||||
| 	github.com/klauspost/compress v1.15.9 // indirect | ||||
| 	github.com/pierrec/lz4/v4 v4.1.15 // indirect | ||||
| 	github.com/klauspost/compress v1.16.0 // indirect | ||||
| 	github.com/pierrec/lz4/v4 v4.1.17 // indirect | ||||
| ) | ||||
|   | ||||
							
								
								
									
										15
									
								
								go.sum
									
									
									
									
									
								
							
							
						
						
									
										15
									
								
								go.sum
									
									
									
									
									
								
							| @@ -1,16 +1,12 @@ | ||||
| github.com/imdario/mergo v0.3.13/go.mod h1:4lJ1jqUDcsbIECGy0RUJAXNIhg+6ocWgb1ALK2O4oXg= | ||||
| github.com/klauspost/compress v1.15.9 h1:wKRjX6JRtDdrE9qwa4b/Cip7ACOshUI4smpCQanqjSY= | ||||
| github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= | ||||
| github.com/klauspost/compress v1.16.0 h1:iULayQNOReoYUe+1qtKOqw9CwJv3aNQu8ivo7lw1HU4= | ||||
| github.com/patrickmn/go-cache v2.1.0+incompatible/go.mod h1:3Qf8kWWT7OJRJbdiICTKqZju1ZixQ/KpMGzzAfe6+WQ= | ||||
| github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= | ||||
| github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= | ||||
| github.com/pierrec/lz4/v4 v4.1.17 h1:kV4Ip+/hUBC+8T6+2EgburRtkE9ef4nbY3f4dFhGjMc= | ||||
| github.com/silas/dag v0.0.0-20211117232152-9d50aa809f35/go.mod h1:7RTUFBdIRC9nZ7/3RyRNH1bdqIShrDejd1YbLwgPS+I= | ||||
| github.com/twmb/franz-go v1.11.5 h1:TTv5lVJd+87XkmP9dWN9Jgpf7IUUr7a7jee+byR8LBE= | ||||
| github.com/twmb/franz-go v1.11.5/go.mod h1:FvaHNlpT6woVYIl6LAuIeL7yHol1Fp6Gv2Dn21AvH78= | ||||
| github.com/twmb/franz-go/pkg/kmsg v1.3.0 h1:ouBETB7nTqRxiO5E8/pySoFZtVEW2VWw55z3/bsUzTw= | ||||
| github.com/twmb/franz-go/pkg/kmsg v1.3.0/go.mod h1:SxG/xJKhgPu25SamAq0rrucfp7lbzCpEXOC+vH/ELrY= | ||||
| go.unistack.org/micro/v3 v3.10.14 h1:7fgLpwGlCN67twhwtngJDEQvrMkUBDSA5vzZqxIDqNE= | ||||
| go.unistack.org/micro/v3 v3.10.14/go.mod h1:uMAc0U/x7dmtICCrblGf0ZLgYegu3VwQAquu+OFCw1Q= | ||||
| github.com/twmb/franz-go v1.12.1 h1:8lWT8q0spL40Nfw6eonJ8OoPGLvF9arvadRRmcSiu9Y= | ||||
| github.com/twmb/franz-go/pkg/kmsg v1.4.0 h1:tbp9hxU6m8qZhQTlpGiaIJOm4BXix5lsuEZ7K00dF0s= | ||||
| go.unistack.org/micro/v3 v3.10.23 h1:4BE7NwwyJbCWOfzjzztamBxJSgRHHW1uQtMGNDLHG3s= | ||||
| golang.org/x/crypto v0.0.0-20220817201139-bc19a97f63c8/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= | ||||
| golang.org/x/net v0.0.0-20211112202133-69e39bad7dc2/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= | ||||
| golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= | ||||
| @@ -20,5 +16,4 @@ golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9sn | ||||
| golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= | ||||
| golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= | ||||
| gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= | ||||
| gopkg.in/yaml.v3 v3.0.0/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= | ||||
| gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= | ||||
|   | ||||
							
								
								
									
										7
									
								
								kgo.go
									
									
									
									
									
								
							
							
						
						
									
										7
									
								
								kgo.go
									
									
									
									
									
								
							| @@ -197,8 +197,10 @@ func (k *Broker) Publish(ctx context.Context, topic string, msg *broker.Message, | ||||
|  | ||||
| func (k *Broker) publish(ctx context.Context, msgs []*broker.Message, opts ...broker.PublishOption) error { | ||||
| 	k.RLock() | ||||
| 	if !k.connected { | ||||
| 		k.RUnlock() | ||||
| 	ok := k.connected | ||||
| 	k.RUnlock() | ||||
|  | ||||
| 	if !ok { | ||||
| 		k.Lock() | ||||
| 		c, err := k.connect(ctx, k.kopts...) | ||||
| 		if err != nil { | ||||
| @@ -209,7 +211,6 @@ func (k *Broker) publish(ctx context.Context, msgs []*broker.Message, opts ...br | ||||
| 		k.connected = true | ||||
| 		k.Unlock() | ||||
| 	} | ||||
| 	k.RUnlock() | ||||
|  | ||||
| 	options := broker.NewPublishOptions(opts...) | ||||
| 	records := make([]*kgo.Record, 0, len(msgs)) | ||||
|   | ||||
		Reference in New Issue
	
	Block a user