|
|
@ -1,6 +1,6 @@
|
|
|
|
// Package exegroup 管理1组执行单元的启动和终止;
|
|
|
|
// Package exegroup 按照顺序启动和终止任务;
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// 使用举例:
|
|
|
|
// Examples:
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// 新建 [Group]:
|
|
|
|
// 新建 [Group]:
|
|
|
|
//
|
|
|
|
//
|
|
|
@ -12,54 +12,48 @@
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// 新增1个 [Actor], 通过 ctx 控制 [Actor] 退出:
|
|
|
|
// 新增1个 [Actor], 通过 ctx 控制 [Actor] 退出:
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// g.New().WithName("do nothing").WithGo(func(ctx context.Context) error {
|
|
|
|
// g.New().WithName("do nothing").WithStartFunc(func(ctx context.Context) error {
|
|
|
|
// <-ctx.Done()
|
|
|
|
// <-ctx.Done()
|
|
|
|
// return ctx.Err()
|
|
|
|
// return ctx.Err()
|
|
|
|
// })
|
|
|
|
// })
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// 新增1个 [Actor], 通过 stopFunc 控制 [Actor] 退出:
|
|
|
|
|
|
|
|
//
|
|
|
|
|
|
|
|
// server := &http.Server{Addr: ":80"}
|
|
|
|
|
|
|
|
// inShutdown := &atomic.Bool{}
|
|
|
|
|
|
|
|
// c := make(chan error, 1)
|
|
|
|
|
|
|
|
// goFunc := func(_ context.Context) error {
|
|
|
|
|
|
|
|
// err := server.ListenAndServe()
|
|
|
|
|
|
|
|
// if inShutdown.Load() {
|
|
|
|
|
|
|
|
// err = <-c
|
|
|
|
|
|
|
|
// }
|
|
|
|
|
|
|
|
// return err
|
|
|
|
|
|
|
|
// }
|
|
|
|
|
|
|
|
// stopFunc := func(ctx context.Context) {
|
|
|
|
|
|
|
|
// inShutdown.Store(true)
|
|
|
|
|
|
|
|
// c <- server.Shutdown(ctx)
|
|
|
|
|
|
|
|
// }
|
|
|
|
|
|
|
|
// g.New().WithName("server").WithGoStop(goFunc, stopFunc)
|
|
|
|
|
|
|
|
//
|
|
|
|
|
|
|
|
// 启动 [Group]:
|
|
|
|
// 启动 [Group]:
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// g.Run()
|
|
|
|
// g.Run(ctx)
|
|
|
|
package exegroup
|
|
|
|
package exegroup
|
|
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"context"
|
|
|
|
|
|
|
|
"errors"
|
|
|
|
"fmt"
|
|
|
|
"fmt"
|
|
|
|
"time"
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
"git.blauwelle.com/go/crate/log"
|
|
|
|
)
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
// Group 管理1组 [Actor], 每个 [Actor] 在1个goroutine 中运行;
|
|
|
|
// Group 管理1组 [Actor], 每个 [Actor] 在1个goroutine 中运行;
|
|
|
|
// Group 的执行过程参考 [Group.Run];
|
|
|
|
// Group 的执行过程参考 [Group.Run];
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// Group 的配置包含:
|
|
|
|
// Group 的配置:
|
|
|
|
// - [WithConcurrentStop] 并发执行 [Actor] 的终止函数, 默认不并发执行;
|
|
|
|
|
|
|
|
// - [WithStopTimeout] 指定 [Group.Run] 从进入终止过程到返回的最长时间, 默认不限制最长时间;
|
|
|
|
// - [WithStopTimeout] 指定 [Group.Run] 从进入终止过程到返回的最长时间, 默认不限制最长时间;
|
|
|
|
type Group struct {
|
|
|
|
type Group struct {
|
|
|
|
|
|
|
|
errorsChan chan actorError
|
|
|
|
actors []*Actor
|
|
|
|
actors []*Actor
|
|
|
|
|
|
|
|
syncActors []*Actor
|
|
|
|
|
|
|
|
asyncActors []*Actor
|
|
|
|
cfg config
|
|
|
|
cfg config
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
type actorError struct {
|
|
|
|
|
|
|
|
actor *Actor
|
|
|
|
|
|
|
|
err error
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// New 创建 [Group];
|
|
|
|
// New 创建 [Group];
|
|
|
|
func New(opts ...Option) *Group {
|
|
|
|
func New(opts ...Option) *Group {
|
|
|
|
cfg := config{}
|
|
|
|
cfg := config{
|
|
|
|
|
|
|
|
stopTimeout: -1,
|
|
|
|
|
|
|
|
}
|
|
|
|
for _, opt := range opts {
|
|
|
|
for _, opt := range opts {
|
|
|
|
opt.apply(&cfg)
|
|
|
|
opt.apply(&cfg)
|
|
|
|
}
|
|
|
|
}
|
|
|
@ -71,88 +65,139 @@ func New(opts ...Option) *Group {
|
|
|
|
// Default 创建包含信号处理 [Actor] 的 [Group];
|
|
|
|
// Default 创建包含信号处理 [Actor] 的 [Group];
|
|
|
|
func Default(opts ...Option) *Group {
|
|
|
|
func Default(opts ...Option) *Group {
|
|
|
|
g := New(opts...)
|
|
|
|
g := New(opts...)
|
|
|
|
g.New().WithName("signal").WithGo(HandleSignal())
|
|
|
|
g.NewAsync().WithName("signal").WithStartFunc(HandleSignal())
|
|
|
|
return g
|
|
|
|
return g
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// New 创建并添加 [Actor],
|
|
|
|
func (g *Group) newActor() *Actor {
|
|
|
|
// 对 Actor 的配置通过链式调用完成;
|
|
|
|
actor := &Actor{
|
|
|
|
func (g *Group) New() *Actor {
|
|
|
|
group: g,
|
|
|
|
actor := new(Actor).WithName(fmt.Sprintf("actor-%03d", len(g.actors)+1))
|
|
|
|
closeAfterStop: make(chan struct{}),
|
|
|
|
|
|
|
|
closeAfterTimeout: make(chan struct{}),
|
|
|
|
|
|
|
|
name: fmt.Sprintf("actor-%d", len(g.actors)),
|
|
|
|
|
|
|
|
stopTimeout: -1,
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
actor.startCtx, actor.startCancel = context.WithCancel(context.Background())
|
|
|
|
g.actors = append(g.actors, actor)
|
|
|
|
g.actors = append(g.actors, actor)
|
|
|
|
return actor
|
|
|
|
return actor
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Run 启动所有 [Actor]s 并等待 [Actor]s 执行完成;
|
|
|
|
// NewAsync 创建并添加异步退出的 [Actor], 对 Actor 的配置通过链式调用完成;
|
|
|
|
|
|
|
|
func (g *Group) NewAsync() *Actor {
|
|
|
|
|
|
|
|
actor := g.newActor()
|
|
|
|
|
|
|
|
g.asyncActors = append(g.asyncActors, actor)
|
|
|
|
|
|
|
|
return actor
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// NewSync 创建并添加同步退出的 [Actor], 对 Actor 的配置通过链式调用完成;
|
|
|
|
|
|
|
|
func (g *Group) NewSync() *Actor {
|
|
|
|
|
|
|
|
actor := g.newActor()
|
|
|
|
|
|
|
|
g.syncActors = append(g.syncActors, actor)
|
|
|
|
|
|
|
|
return actor
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Run 启动所有 [Actor] 并等待 [Actor] 执行完成;
|
|
|
|
// 当没有 Actor 时执行 Run 会 panic;
|
|
|
|
// 当没有 Actor 时执行 Run 会 panic;
|
|
|
|
//
|
|
|
|
//
|
|
|
|
// Run 包含两个阶段:
|
|
|
|
// Run 包含两个阶段:
|
|
|
|
// 1. 启动 [Actor] 并等待, Actor 的运行函数返回错误时进入下1个阶段;
|
|
|
|
// 1. 启动所有 [Actor];
|
|
|
|
// 2. 执行 [Actor] 的终止函数并返回;
|
|
|
|
// 2. 等待 ctx 或第1个Actor返回, 终止所有Actor.
|
|
|
|
func (g *Group) Run(ctx context.Context) error {
|
|
|
|
func (g *Group) Run(ctx context.Context) error {
|
|
|
|
g.validate()
|
|
|
|
if err := g.validate(ctx); err != nil {
|
|
|
|
ctx, cancel := context.WithCancel(ctx)
|
|
|
|
return err
|
|
|
|
c := make(chan error, len(g.actors))
|
|
|
|
}
|
|
|
|
g.start(ctx, c)
|
|
|
|
g.errorsChan = make(chan actorError, len(g.actors))
|
|
|
|
return g.wait(c, cancel)
|
|
|
|
g.start(ctx)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
return g.stop(ctx, g.wait(ctx))
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func (g *Group) validate() {
|
|
|
|
func (g *Group) validate(ctx context.Context) error {
|
|
|
|
if len(g.actors) == 0 {
|
|
|
|
if len(g.actors) == 0 {
|
|
|
|
panic("no actor")
|
|
|
|
err := errors.New("no actor")
|
|
|
|
|
|
|
|
log.Error(ctx, err.Error())
|
|
|
|
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
for _, actor := range g.actors {
|
|
|
|
for i, actor := range g.actors {
|
|
|
|
if actor.goFunc == nil {
|
|
|
|
if actor.startFunc == nil {
|
|
|
|
panic(actor.name + " has nil goFunc")
|
|
|
|
err := fmt.Errorf("group-actor %d: %s has nil startFunc", i, actor.name)
|
|
|
|
|
|
|
|
log.Error(ctx, err.Error())
|
|
|
|
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func (g *Group) start(ctx context.Context, c chan error) {
|
|
|
|
// start 启动所有的 [Actor] 后返回.
|
|
|
|
|
|
|
|
func (g *Group) start(ctx context.Context) {
|
|
|
|
for _, actor := range g.actors {
|
|
|
|
for _, actor := range g.actors {
|
|
|
|
go func(actor *Actor) {
|
|
|
|
go actor.start()
|
|
|
|
var err error
|
|
|
|
log.Tracef(ctx, "group is starting group-actor %s", actor.name)
|
|
|
|
defer func() {
|
|
|
|
|
|
|
|
if v := recover(); v != nil {
|
|
|
|
|
|
|
|
err = fmt.Errorf("%v", v)
|
|
|
|
|
|
|
|
}
|
|
|
|
}
|
|
|
|
c <- err
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
err = actor.goFunc(ctx)
|
|
|
|
// wait 等待 ctx 或 [Actor] 返回
|
|
|
|
}(actor)
|
|
|
|
func (g *Group) wait(ctx context.Context) error {
|
|
|
|
|
|
|
|
log.Tracef(ctx, "group is waiting")
|
|
|
|
|
|
|
|
select {
|
|
|
|
|
|
|
|
case ae := <-g.errorsChan:
|
|
|
|
|
|
|
|
log.Tracef(ctx, "group is waiting: group-actor %s return", ae.actor.name)
|
|
|
|
|
|
|
|
if ae.err != nil {
|
|
|
|
|
|
|
|
return fmt.Errorf("group-actor %s return with %w", ae.actor.name, ae.err)
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
|
|
|
case <-ctx.Done():
|
|
|
|
|
|
|
|
log.Tracef(ctx, "group is waiting: %s", ctx.Err().Error())
|
|
|
|
|
|
|
|
return ctx.Err()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
func (g *Group) wait(c chan error, cancel context.CancelFunc) error {
|
|
|
|
// stop 等待 [Actor] 停止
|
|
|
|
err := <-c
|
|
|
|
func (g *Group) stop(ctx context.Context, err error) error {
|
|
|
|
cancel()
|
|
|
|
stopCtx := context.Background()
|
|
|
|
ctx := context.Background()
|
|
|
|
if g.cfg.stopTimeout >= 0 {
|
|
|
|
if g.cfg.stopTimeout > 0 {
|
|
|
|
log.Tracef(ctx, "set group stop timeout to %s", g.cfg.stopTimeout)
|
|
|
|
ctx, cancel = context.WithTimeout(context.Background(), g.cfg.stopTimeout)
|
|
|
|
var cancel func()
|
|
|
|
|
|
|
|
stopCtx, cancel = context.WithTimeout(stopCtx, g.cfg.stopTimeout)
|
|
|
|
defer cancel()
|
|
|
|
defer cancel()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
for i := len(g.actors) - 1; i >= 0; i-- {
|
|
|
|
go func() {
|
|
|
|
if g.actors[i].stopFunc != nil {
|
|
|
|
// 终止异步 [Actor]
|
|
|
|
if g.cfg.concurrentStop {
|
|
|
|
for _, actor := range g.asyncActors {
|
|
|
|
go g.actors[i].stopFunc(ctx)
|
|
|
|
log.Tracef(ctx, "group is stopping group-actor %s", actor.name)
|
|
|
|
} else {
|
|
|
|
actor.sendStopSignal()
|
|
|
|
g.actors[i].stopFunc(ctx)
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// 终止同步 [Actor] 并等待 [Actor] 退出
|
|
|
|
|
|
|
|
for i := len(g.syncActors) - 1; i >= 0; i-- {
|
|
|
|
|
|
|
|
actor := g.syncActors[i]
|
|
|
|
|
|
|
|
log.Tracef(ctx, "group is stopping group-actor %s", actor.name)
|
|
|
|
|
|
|
|
actor.sendStopSignal()
|
|
|
|
|
|
|
|
actor.wait()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
for i := 1; i < len(g.actors); i++ {
|
|
|
|
}()
|
|
|
|
|
|
|
|
// 等待所有 [Actor] 退出, 超时立即退出
|
|
|
|
|
|
|
|
for _, actor := range g.actors {
|
|
|
|
select {
|
|
|
|
select {
|
|
|
|
case <-c:
|
|
|
|
case <-actor.closeAfterStop:
|
|
|
|
case <-ctx.Done():
|
|
|
|
log.Tracef(ctx, "group is stopping group-actor %s: stopped", actor.name)
|
|
|
|
|
|
|
|
case <-actor.closeAfterTimeout:
|
|
|
|
|
|
|
|
log.Errorf(ctx, "group is stopping group-actor %s: timeout", actor.name)
|
|
|
|
|
|
|
|
case <-stopCtx.Done():
|
|
|
|
|
|
|
|
if err == nil {
|
|
|
|
|
|
|
|
err = stopCtx.Err()
|
|
|
|
|
|
|
|
} else {
|
|
|
|
|
|
|
|
err = fmt.Errorf("group is stopping group-actor %s: global timeout: %w", actor.name, err)
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
log.Error(ctx, err.Error())
|
|
|
|
return err
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
log.Tracef(ctx, "group return: %v", err)
|
|
|
|
return err
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
type config struct {
|
|
|
|
type config struct {
|
|
|
|
concurrentStop bool
|
|
|
|
stopTimeout time.Duration // <0表示没有超时时间
|
|
|
|
stopTimeout time.Duration
|
|
|
|
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Option 修改 [Group] 的配置;
|
|
|
|
// Option 修改 [Group] 的配置;
|
|
|
@ -166,14 +211,8 @@ func (fn optionFunc) apply(cfg *config) {
|
|
|
|
fn(cfg)
|
|
|
|
fn(cfg)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// WithConcurrentStop 并发执行 [Actor] 的终止函数;
|
|
|
|
// WithStopTimeout 指定 [Group.Run] 从进入终止过程到返回的最长时间.
|
|
|
|
func WithConcurrentStop() Option {
|
|
|
|
// <0表示没有超时时间.
|
|
|
|
return optionFunc(func(cfg *config) {
|
|
|
|
|
|
|
|
cfg.concurrentStop = true
|
|
|
|
|
|
|
|
})
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// WithStopTimeout 指定 [Group.Run] 从进入终止过程到返回的最长时间;
|
|
|
|
|
|
|
|
func WithStopTimeout(d time.Duration) Option {
|
|
|
|
func WithStopTimeout(d time.Duration) Option {
|
|
|
|
return optionFunc(func(cfg *config) {
|
|
|
|
return optionFunc(func(cfg *config) {
|
|
|
|
cfg.stopTimeout = d
|
|
|
|
cfg.stopTimeout = d
|
|
|
|