135 lines
3.7 KiB
Go
135 lines
3.7 KiB
Go
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()
|
||
}
|
||
}
|
||
|
||
// LabelTaskPool 预标注任务并发池:逐张图片调 RF-DETR(IO 等待为主),
|
||
// 并发度来自 config.yml labelTask.poolSize,缺失或非法时回退 consts 默认值。
|
||
// 池内任务禁止提交本池(防 worker 饿死死锁);DB 写仍走 Serial 单写者。
|
||
type LabelTaskPool struct {
|
||
pool *grpool.Pool
|
||
}
|
||
|
||
var (
|
||
labelTaskPoolOnce sync.Once
|
||
labelTaskPool *LabelTaskPool
|
||
)
|
||
|
||
// LabelTaskPoolInstance 进程级预标注池单例(懒初始化,读取配置)。
|
||
func LabelTaskPoolInstance() *LabelTaskPool {
|
||
labelTaskPoolOnce.Do(func() {
|
||
ctx := context.Background()
|
||
size := g.Cfg().MustGet(ctx, "labelTask.poolSize", consts.LabelPoolDefaultSize).Int()
|
||
if size <= 0 {
|
||
size = consts.LabelPoolDefaultSize
|
||
}
|
||
labelTaskPool = &LabelTaskPool{pool: grpool.New(size, size)}
|
||
})
|
||
return labelTaskPool
|
||
}
|
||
|
||
// Submit 提交单张图片的预标注任务并等待完成,返回任务的 error。
|
||
func (p *LabelTaskPool) 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()
|
||
}
|
||
}
|
||
|
||
// GenTaskPool 文生图任务并发池:逐张调 local-ai 生成图片(IO 等待为主),
|
||
// 并发度来自 config.yml imageGen.poolSize(z-image 显存独占,默认 1),
|
||
// 缺失或非法时回退 consts 默认值。池内任务禁止提交本池(防 worker 饿死死锁)。
|
||
type GenTaskPool struct {
|
||
pool *grpool.Pool
|
||
}
|
||
|
||
var (
|
||
genTaskPoolOnce sync.Once
|
||
genTaskPool *GenTaskPool
|
||
)
|
||
|
||
// GenTaskPoolInstance 进程级文生图池单例(懒初始化,读取配置)。
|
||
func GenTaskPoolInstance() *GenTaskPool {
|
||
genTaskPoolOnce.Do(func() {
|
||
ctx := context.Background()
|
||
size := g.Cfg().MustGet(ctx, "imageGen.poolSize", consts.GenPoolDefaultSize).Int()
|
||
if size <= 0 {
|
||
size = consts.GenPoolDefaultSize
|
||
}
|
||
genTaskPool = &GenTaskPool{pool: grpool.New(size, size)}
|
||
})
|
||
return genTaskPool
|
||
}
|
||
|
||
// Submit 提交单张图片的生成任务并等待完成,返回任务的 error。
|
||
func (p *GenTaskPool) 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()
|
||
}
|
||
}
|