package common import ( "context" "sync" "github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/os/grpool" "observer-server/biz/consts" ) // CallbackPool 支付回调处理并发池:验签/落授权等 IO 任务在池内执行并返回结果, // 并发度来自 config.yml payment.poolSize,缺失或非法时回退 consts 默认值。 // 池内任务禁止再提交本池(防 worker 饿死死锁)。 type CallbackPool struct { pool *grpool.Pool } var ( callbackPoolOnce sync.Once callbackPool *CallbackPool ) // CallbackPoolInstance 进程级回调池单例(懒初始化,读取配置)。 func CallbackPoolInstance() *CallbackPool { callbackPoolOnce.Do(func() { ctx := context.Background() size := g.Cfg().MustGet(ctx, "payment.poolSize", consts.PaymentPoolDefaultSize).Int() if size <= 0 { size = consts.PaymentPoolDefaultSize } callbackPool = &CallbackPool{pool: grpool.New(size, size)} }) return callbackPool } // Submit 提交回调任务并等待执行完成,返回任务的 error。 func (p *CallbackPool) Submit(ctx context.Context, fn func(ctx context.Context) error) error { res := make(chan error, 1) if err := p.pool.Add(ctx, func(ctx context.Context) { res <- fn(ctx) }); err != nil { return err } select { case err := <-res: return err case <-ctx.Done(): return ctx.Err() } }