Skip to content

Commit c8cc598

Browse files
committed
refactor(examples): better graceful shutdown
1 parent 670c61b commit c8cc598

8 files changed

Lines changed: 128 additions & 171 deletions

File tree

examples/confluent/cmd/consumer/main.go

Lines changed: 30 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -2,32 +2,28 @@ package main
22

33
import (
44
"context"
5+
"errors"
56
"fmt"
67
"net/http"
78
"os"
89
"os/signal"
910
"syscall"
11+
"time"
1012

1113
"github.com/alebabai/go-kafka"
1214
adapter "github.com/alebabai/go-kafka/adapter/confluent"
1315
ckafka "github.com/confluentinc/confluent-kafka-go/v2/kafka"
1416
"github.com/go-kit/log"
1517
"github.com/go-kit/log/level"
18+
"golang.org/x/sync/errgroup"
1619

1720
"github.com/alebabai/go-kit-kafka/v2/examples/common"
1821
"github.com/alebabai/go-kit-kafka/v2/examples/common/consumer"
1922
)
2023

2124
func main() {
22-
var (
23-
ctx context.Context
24-
cancel context.CancelFunc
25-
)
26-
{
27-
ctx = context.Background()
28-
ctx, cancel = context.WithCancel(ctx)
29-
defer cancel()
30-
}
25+
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
26+
defer stop()
3127

3228
var logger log.Logger
3329
{
@@ -90,49 +86,44 @@ func main() {
9086

9187
_ = logger.Log("msg", "initialize http server")
9288

93-
var httpServer *http.Server
94-
{
95-
httpServer = &http.Server{
96-
Addr: consumer.HTTPServerAddr,
97-
Handler: consumer.NewHTTPHandler(consumer.MakeListEventsEndpoint(svc)),
98-
}
99-
100-
defer func() {
101-
if err := httpServer.Shutdown(ctx); err != nil {
102-
fatal(logger, fmt.Errorf("failed to shutdown http server: %w", err))
103-
}
104-
}()
89+
httpServer := &http.Server{
90+
Addr: consumer.HTTPServerAddr,
91+
Handler: consumer.NewHTTPHandler(consumer.MakeListEventsEndpoint(svc)),
10592
}
10693

107-
errc := make(chan error, 1)
94+
g, ctx := errgroup.WithContext(ctx)
10895

109-
go func() {
96+
g.Go(func() error {
11097
if err := consumerListener.Listen(ctx, -1); err != nil {
111-
errc <- err
98+
return fmt.Errorf("kafka listener: %w", err)
11299
}
113-
}()
100+
return nil
101+
})
114102

115-
go func() {
116-
if err := httpServer.ListenAndServe(); err != nil {
117-
errc <- err
103+
g.Go(func() error {
104+
if err := httpServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
105+
return fmt.Errorf("http server: %w", err)
118106
}
119-
}()
120-
121-
sigc := make(chan os.Signal, 1)
107+
return nil
108+
})
122109

123-
go func() {
124-
sigc := make(chan os.Signal, 1)
125-
signal.Notify(sigc, syscall.SIGINT, syscall.SIGTERM)
126-
}()
110+
g.Go(func() error {
111+
<-ctx.Done()
112+
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
113+
defer cancel()
114+
if err := httpServer.Shutdown(shutdownCtx); err != nil {
115+
return fmt.Errorf("http server shutdown: %w", err)
116+
}
117+
return nil
118+
})
127119

128120
_ = logger.Log("msg", "application started")
129121

130-
select {
131-
case sig := <-sigc:
132-
_ = logger.Log("msg", "application stopped", "exit", sig)
133-
case err := <-errc:
122+
if err := g.Wait(); err != nil {
134123
fatal(logger, err)
135124
}
125+
126+
_ = logger.Log("msg", "application stopped")
136127
}
137128

138129
func fatal(logger log.Logger, err error) {

examples/confluent/cmd/producer/main.go

Lines changed: 29 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package main
22

33
import (
44
"context"
5+
"errors"
56
"fmt"
67
"net/http"
78
"os"
@@ -13,26 +14,16 @@ import (
1314
"github.com/go-kit/kit/endpoint"
1415
"github.com/go-kit/log"
1516
"github.com/go-kit/log/level"
17+
"golang.org/x/sync/errgroup"
1618

1719
"github.com/alebabai/go-kit-kafka/v2/examples/common"
1820
"github.com/alebabai/go-kit-kafka/v2/examples/common/producer"
19-
2021
"github.com/alebabai/go-kit-kafka/v2/examples/confluent/pkg/kafka/adapter"
2122
)
2223

2324
func main() {
24-
var (
25-
ctx context.Context
26-
cancel context.CancelFunc
27-
)
28-
{
29-
ctx = context.Background()
30-
ctx, cancel = context.WithCancel(ctx)
31-
defer func() {
32-
cancel()
33-
time.Sleep(3 * time.Second)
34-
}()
35-
}
25+
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
26+
defer stop()
3627

3728
var logger log.Logger
3829
{
@@ -78,46 +69,41 @@ func main() {
7869

7970
_ = logger.Log("msg", "initializing http server")
8071

81-
var httpServer *http.Server
82-
{
83-
httpServer = &http.Server{
84-
Addr: producer.HTTPServerAddr,
85-
Handler: producer.NewHTTPHandler(
86-
producerMiddleware(
87-
producer.MakeGenerateEventEndpoint(svc),
88-
),
72+
httpServer := &http.Server{
73+
Addr: producer.HTTPServerAddr,
74+
Handler: producer.NewHTTPHandler(
75+
producerMiddleware(
76+
producer.MakeGenerateEventEndpoint(svc),
8977
),
90-
}
91-
92-
defer func() {
93-
if err := httpServer.Shutdown(ctx); err != nil {
94-
fatal(logger, fmt.Errorf("failed to shutdown http server: %w", err))
95-
}
96-
}()
78+
),
9779
}
9880

99-
errc := make(chan error, 1)
81+
g, ctx := errgroup.WithContext(ctx)
10082

101-
go func() {
102-
if err := httpServer.ListenAndServe(); err != nil {
103-
errc <- err
83+
g.Go(func() error {
84+
if err := httpServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
85+
return fmt.Errorf("http server: %w", err)
10486
}
105-
}()
106-
107-
sigc := make(chan os.Signal, 1)
108-
109-
go func() {
110-
signal.Notify(sigc, syscall.SIGINT, syscall.SIGTERM)
111-
}()
87+
return nil
88+
})
89+
90+
g.Go(func() error {
91+
<-ctx.Done()
92+
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
93+
defer cancel()
94+
if err := httpServer.Shutdown(shutdownCtx); err != nil {
95+
return fmt.Errorf("http server shutdown: %w", err)
96+
}
97+
return nil
98+
})
11299

113100
_ = logger.Log("msg", "application started")
114101

115-
select {
116-
case sig := <-sigc:
117-
_ = logger.Log("msg", "application stopped", "exit", sig)
118-
case err := <-errc:
102+
if err := g.Wait(); err != nil {
119103
fatal(logger, err)
120104
}
105+
106+
_ = logger.Log("msg", "application stopped")
121107
}
122108

123109
func fatal(logger log.Logger, err error) {

examples/confluent/go.mod

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
module github.com/alebabai/go-kit-kafka/v2/examples/confluent
22

3-
go 1.22
3+
go 1.25.0
44

55
require (
66
github.com/alebabai/go-kafka v0.3.2
@@ -9,6 +9,7 @@ require (
99
github.com/confluentinc/confluent-kafka-go/v2 v2.4.0
1010
github.com/go-kit/kit v0.13.0
1111
github.com/go-kit/log v0.2.1
12+
golang.org/x/sync v0.21.0
1213
)
1314

1415
require (

examples/confluent/go.sum

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -326,8 +326,8 @@ golang.org/x/net v0.22.0 h1:9sGLhx7iRIHEiX0oAJ3MRZMUCElJgy7Br1nO+AMN3Tc=
326326
golang.org/x/net v0.22.0/go.mod h1:JKghWKKOSdJwpW2GEx0Ja7fmaKnMsbu+MWVZTokSYmg=
327327
golang.org/x/oauth2 v0.16.0 h1:aDkGMBSYxElaoP81NpoUoz2oo2R2wHdZpGToUxfyQrQ=
328328
golang.org/x/oauth2 v0.16.0/go.mod h1:hqZ+0LWXsiVoZpeld6jVt06P3adbS2Uu911W1SsJv2o=
329-
golang.org/x/sync v0.6.0 h1:5BMeUDZ7vkXGfEr1x9B4bRcTH4lpkTkpdh0T/J+qjbQ=
330-
golang.org/x/sync v0.6.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
329+
golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
330+
golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
331331
golang.org/x/sys v0.18.0 h1:DBdB3niSjOA/O0blCZBqDefyWNYveAYMNF1Wum0DYQ4=
332332
golang.org/x/sys v0.18.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
333333
golang.org/x/term v0.18.0 h1:FcHjZXDMxI8mM3nwhX9HlKop4C0YQvCVCdwYl2wOtE8=

examples/sarama/cmd/consumer/main.go

Lines changed: 32 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package main
22

33
import (
44
"context"
5+
"errors"
56
"fmt"
67
"net/http"
78
"os"
@@ -14,21 +15,15 @@ import (
1415
adapter "github.com/alebabai/go-kafka/adapter/sarama"
1516
"github.com/go-kit/log"
1617
"github.com/go-kit/log/level"
18+
"golang.org/x/sync/errgroup"
1719

1820
"github.com/alebabai/go-kit-kafka/v2/examples/common"
1921
"github.com/alebabai/go-kit-kafka/v2/examples/common/consumer"
2022
)
2123

2224
func main() {
23-
var (
24-
ctx context.Context
25-
cancel context.CancelFunc
26-
)
27-
{
28-
ctx = context.Background()
29-
ctx, cancel = context.WithCancel(ctx)
30-
defer cancel()
31-
}
25+
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
26+
defer stop()
3227

3328
var logger log.Logger
3429
{
@@ -101,55 +96,51 @@ func main() {
10196

10297
_ = logger.Log("msg", "initializing http server")
10398

104-
var httpServer *http.Server
105-
{
106-
httpServer = &http.Server{
107-
Addr: consumer.HTTPServerAddr,
108-
Handler: consumer.NewHTTPHandler(consumer.MakeListEventsEndpoint(svc)),
109-
}
110-
111-
defer func() {
112-
if err := httpServer.Shutdown(ctx); err != nil {
113-
fatal(logger, fmt.Errorf("failed to shutdown http server: %w", err))
114-
}
115-
}()
99+
httpServer := &http.Server{
100+
Addr: consumer.HTTPServerAddr,
101+
Handler: consumer.NewHTTPHandler(consumer.MakeListEventsEndpoint(svc)),
116102
}
117103

118-
errc := make(chan error, 1)
104+
g, ctx := errgroup.WithContext(ctx)
119105

120-
go func() {
106+
g.Go(func() error {
121107
if err := consumerGroupListener.Listen(ctx, common.KafkaTopic); err != nil {
122-
errc <- err
108+
return fmt.Errorf("kafka listener: %w", err)
123109
}
124-
}()
110+
return nil
111+
})
125112

126-
go func() {
113+
g.Go(func() error {
127114
for err := range consumerGroupListener.Errors() {
128-
time.Sleep(2 * time.Second) // debounce
129115
level.Error(logger).Log("err", err)
130116
}
131-
}()
117+
return nil
118+
})
132119

133-
go func() {
134-
if err := httpServer.ListenAndServe(); err != nil {
135-
errc <- err
120+
g.Go(func() error {
121+
if err := httpServer.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
122+
return fmt.Errorf("http server: %w", err)
136123
}
137-
}()
138-
139-
sigc := make(chan os.Signal, 1)
124+
return nil
125+
})
140126

141-
go func() {
142-
signal.Notify(sigc, syscall.SIGINT, syscall.SIGTERM)
143-
}()
127+
g.Go(func() error {
128+
<-ctx.Done()
129+
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
130+
defer cancel()
131+
if err := httpServer.Shutdown(shutdownCtx); err != nil {
132+
return fmt.Errorf("http server shutdown: %w", err)
133+
}
134+
return nil
135+
})
144136

145137
_ = logger.Log("msg", "application started")
146138

147-
select {
148-
case sig := <-sigc:
149-
_ = logger.Log("msg", "application stopped", "exit", sig)
150-
case err := <-errc:
139+
if err := g.Wait(); err != nil {
151140
fatal(logger, err)
152141
}
142+
143+
_ = logger.Log("msg", "application stopped")
153144
}
154145

155146
func fatal(logger log.Logger, err error) {

0 commit comments

Comments
 (0)