Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions docs/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
# 1.2.8 (2026-07-27)

- `asynctaskext`/`busext` 新增 Prometheus 处理耗时/QPS 埋点(`asynctask_task_duration_seconds`/`bus_task_duration_seconds`),config `<NS>monitor_enable` 开关可选开启,默认关闭零开销;与 Python coast 库同名同 label 同 buckets,可跨语言合并查询

# 1.2.7 (2026-06-16)

- 移除 shanbay/go-redis fork 的 replace 指令,改用官方 `github.com/redis/go-redis/extra/redisotel/v9` v9.18.0 及 `rediscmd/v9` v9.18.0
Expand Down
26 changes: 26 additions & 0 deletions docs/ext_amqp_cn.md
Original file line number Diff line number Diff line change
Expand Up @@ -208,3 +208,29 @@ func TestHelloHandler_Run(t *testing.T) {
}
}
```

## 监控指标

在 config.yaml 里加一行即可开启处理耗时/QPS 埋点(默认关闭,零开销):

```yaml
bus_monitor_enable: true
```

开启后会记录 Prometheus Histogram 指标 `bus_task_duration_seconds`(label:`task_name`、`queue`、`status`;两者都取消息的 routing key,因为 busext 目前没有独立于 routing key 的"任务名"概念)。`status` 是真实的处理结果:`success` / `failure`(`handler.Run()` 返回 error)/ `parse_error`(payload 解析失败)/ `invalid_message`(headers、content-type、encoding、未注册 routing key 等消息本身有问题)。gobay 本身**不会**自建 `/metrics` HTTP server,需要使用方应用自己起一个(如 `promhttp.Handler()`),照搬 `cachext` 的既有惯例。

这个指标和 Python coast 库(`coast/celery.py`)里的 `bus_task_duration_seconds` 是**同一个指标**(同名、同 label、同 buckets),可以在同一个 Prometheus/Grafana 查询里合并 Go 和 Python 两边的数据,不需要额外处理。

常用 PromQL:

```promql
# QPS
sum(rate(bus_task_duration_seconds_count[1m])) by (task_name)

# 失败率
sum(rate(bus_task_duration_seconds_count{status="failure"}[5m])) by (task_name)
/ sum(rate(bus_task_duration_seconds_count[5m])) by (task_name)

# P95 耗时
histogram_quantile(0.95, sum(rate(bus_task_duration_seconds_bucket[5m])) by (le, task_name))
```
26 changes: 26 additions & 0 deletions docs/ext_asynctask_cn.md
Original file line number Diff line number Diff line change
Expand Up @@ -116,3 +116,29 @@ func SomeAsyncTask(ctx context.Context, relatedItemIDToProcess uint64) error {
```

## 测试代码只需测试 SomeAsyncTask 函数本身即可

## 监控指标

在 config.yaml 里加一行即可开启处理耗时/QPS 埋点(默认关闭,零开销):

```yaml
asynctask_monitor_enable: true
```

开启后会记录 Prometheus Histogram 指标 `asynctask_task_duration_seconds`(label:`task_name`、`queue`、`status`;`status` 固定为 `"unknown"`——machinery 的全局 hook 拿不到每次任务调用的成功/失败信息,只能保证记录到耗时和 QPS)。gobay 本身**不会**自建 `/metrics` HTTP server,需要使用方应用自己起一个(如 `promhttp.Handler()`),照搬 `cachext` 的既有惯例。

这个指标和 Python coast 库(`coast/celery.py`)里的 `asynctask_task_duration_seconds` 是**同一个指标**(同名、同 label、同 buckets),可以在同一个 Prometheus/Grafana 查询里合并 Go 和 Python 两边的数据,不需要额外处理。

常用 PromQL:

```promql
# QPS
sum(rate(asynctask_task_duration_seconds_count[1m])) by (task_name)

# 平均耗时
sum(rate(asynctask_task_duration_seconds_sum[5m])) by (task_name)
/ sum(rate(asynctask_task_duration_seconds_count[5m])) by (task_name)

# P95 耗时
histogram_quantile(0.95, sum(rate(asynctask_task_duration_seconds_bucket[5m])) by (le, task_name))
```
44 changes: 44 additions & 0 deletions extensions/asynctaskext/asynctaskext.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import (
"github.com/google/uuid"
"github.com/mitchellh/mapstructure"
"github.com/shanbay/gobay"
"github.com/shanbay/gobay/observability"
)

const (
Expand All @@ -39,6 +40,10 @@ type AsyncTaskExt struct {
lock sync.Mutex
healthCheckCompleteChan chan string
healthHandlerRegistered bool

monitorEnabled bool
taskStartTimes map[string]time.Time
taskStartTimesM sync.Mutex
}

func (t *AsyncTaskExt) Object() interface{} {
Expand Down Expand Up @@ -70,6 +75,12 @@ func (t *AsyncTaskExt) Init(app *gobay.Application) error {
return err
}
t.server = server

if config.GetBool("monitor_enable") {
t.monitorEnabled = true
t.taskStartTimes = make(map[string]time.Time)
}

return t.registerHealthCheck()
}

Expand Down Expand Up @@ -102,6 +113,13 @@ func (t *AsyncTaskExt) StartWorker(queue string, concurrency int, enableHealthCh
worker.Queue = queue
t.workers = append(t.workers, worker)

if t.monitorEnabled {
worker.SetPreTaskHandler(t.recordTaskStart)
worker.SetPostTaskHandler(func(sig *tasks.Signature) {
t.recordTaskDuration(sig, queue)
})
}

// run health check http server
if enableHealthCheck && !t.healthHandlerRegistered {
t.healthHandlerRegistered = true
Expand Down Expand Up @@ -146,6 +164,32 @@ func (t *AsyncTaskExt) genConsumerTag(queue string) string {
return fmt.Sprintf("%s@%s", queue, hostName)
}

// recordTaskStart 记录任务开始处理的时间点,按 signature.UUID 索引。
func (t *AsyncTaskExt) recordTaskStart(sig *tasks.Signature) {
t.taskStartTimesM.Lock()
defer t.taskStartTimesM.Unlock()
t.taskStartTimes[sig.UUID] = time.Now()
}

// recordTaskDuration 从 taskStartTimes 弹出起点算 duration,Observe 到
// observability.AsyncTaskDurationSeconds。status 固定写死 "unknown"——machinery
// 的 SetPostTaskHandler 全局 hook 拿不到该次调用的 error,无法在并发 worker 下
// 精确归因到具体任务,为避免引入 reflect.MakeFunc 包装 handler 的复杂度和风险,
// 明确放弃精确 status(详见 plan §1.3)。
func (t *AsyncTaskExt) recordTaskDuration(sig *tasks.Signature, queue string) {
t.taskStartTimesM.Lock()
start, ok := t.taskStartTimes[sig.UUID]
if ok {
delete(t.taskStartTimes, sig.UUID)
}
t.taskStartTimesM.Unlock()
if !ok {
return
}
observability.AsyncTaskDurationSeconds.WithLabelValues(sig.Name, queue, "unknown").
Observe(time.Since(start).Seconds())
}

func (t *AsyncTaskExt) registerHealthCheck() error {
return t.server.RegisterTask(healthCheckTaskName, func(healthCheckUUID string) error {
select {
Expand Down
217 changes: 217 additions & 0 deletions extensions/asynctaskext/asynctaskext_test.go
Original file line number Diff line number Diff line change
@@ -1,14 +1,34 @@
/*
# Partition table (ECP + BVA) — step 01-asynctask-bus-metrics (asynctaskext side)
#
# 参数/状态 | 等价类 | 类型 | 代表值/场景 | 期望输出 | 对应契约条目
# monitorEnabled | 关闭(默认/未设置) | 有效 | one_asynctask_ 在 "testing" env 下未设置 monitor_enable | 任务正常执行不受影响,不 panic,asynctask_task_duration_seconds 不含该 task_name 的记录 | 数据层 API 锚点 1
# monitorEnabled | 开启 | 有效 | three_asynctask_monitor_enable: true | Init 不 panic,monitorEnabled == true | 数据层 API 锚点 2/3/9
# 任务执行结果 (task outcome)| 成功 | 有效 | TaskAddThree 正常返回 | 指标 _count 增加、_sum > 0,status 仍为 "unknown" | 数据层 API 锚点 2、行为契约"status 固定 unknown"
# 任务执行结果 (task outcome)| 返回 error | 有效(错误路径) | TaskFailThree 返回 error | 指标依然被 Observe(不因失败跳过),status 仍为 "unknown" | 数据层 API 锚点 3
# 同类型多实例 (multi-instance)| 两个实例都开启 monitor_enable | 边界/并发场景 | three_asynctask_ + four_asynctask_ 同时 Init | 不 panic(验证包级单例、无重复 promauto 注册) | 数据层 API 锚点 9、MUST "包级单例"
#
# 说明:本 step 的 asynctaskext 侧只有 2 个独立参数(monitorEnabled ×2、task outcome ×2),
# 未达到 pairwise 触发阈值(≥3 参数 × 每个 ≥2 取值),因此未调用 pairwise 脚本,
# 采用穷举的 2×2 组合(关闭+成功 已由锚点1覆盖;开启+成功、开启+失败 由锚点2/3覆盖)。
*/
package asynctaskext

import (
"context"
"errors"
"io"
"log"
"net/http"
"strconv"
"strings"
"sync"
"testing"
"time"

"github.com/RichardKnop/machinery/v1/backends/result"
"github.com/RichardKnop/machinery/v1/tasks"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/stretchr/testify/assert"

"github.com/shanbay/gobay"
Expand Down Expand Up @@ -55,6 +75,81 @@ func TaskSubWithContext(ctx context.Context, arg1, arg2 int64) (int64, error) {
return arg1 - arg2, nil
}

// TaskAddThree / TaskFailThree are used exclusively by the monitor tests below,
// with names that never collide with "add"/"sub"/"subCtx" registered on taskOne,
// so metric assertions can't accidentally match unrelated task executions.
func TaskAddThree(args ...int64) (int64, error) {
sum := int64(0)
for _, arg := range args {
sum += arg
}
return sum, nil
}

func TaskFailThree(arg int64) (int64, error) {
return 0, errors.New("intentional failure for asynctaskext monitor test")
}

var (
taskThree AsyncTaskExt
taskFour AsyncTaskExt

asyncTaskMetricsOnce sync.Once
)

const asyncTaskMetricsAddr = "127.0.0.1:2113"

// startAsyncTaskMetricsServer exposes prometheus.DefaultRegisterer via
// promhttp.Handler(), mirroring extensions/cachext's TestCacheExt_Cached_Monitor
// pattern (config-gated instrumentation + no self-built /metrics server in
// production code, only in the test harness).
func startAsyncTaskMetricsServer() {
asyncTaskMetricsOnce.Do(func() {
go func() {
http.Handle("/metrics", promhttp.Handler())
if err := http.ListenAndServe(asyncTaskMetricsAddr, nil); err != nil {
log.Fatalf("error when start prometheus server: %v\n", err)
}
}()
time.Sleep(200 * time.Millisecond)
})
}

func fetchAsyncTaskMetrics(t *testing.T) string {
resp, err := http.Get("http://" + asyncTaskMetricsAddr + "/metrics")
if err != nil {
t.Fatal(err)
}
defer resp.Body.Close()
data, err := io.ReadAll(resp.Body)
if err != nil {
t.Fatal(err)
}
return string(data)
}

// extractMetricValue looks up "<family>{<labels>} <value>" in the exposition
// text and parses the numeric value, so tests can assert e.g. "_sum > 0"
// instead of only checking for substring presence.
func extractMetricValue(t *testing.T, metrics, family, labels string) float64 {
prefix := family + "{" + labels + "} "
idx := strings.Index(metrics, prefix)
if idx == -1 {
t.Fatalf("metric not found: %s", prefix)
}
rest := metrics[idx+len(prefix):]
end := strings.IndexByte(rest, '\n')
if end == -1 {
end = len(rest)
}
valueStr := strings.TrimSpace(rest[:end])
value, err := strconv.ParseFloat(valueStr, 64)
if err != nil {
t.Fatalf("failed to parse metric value %q: %v", valueStr, err)
}
return value
}

func TestPushConsume(t *testing.T) {
if err := taskOne.RegisterWorkerHandlers(map[string]interface{}{
"add": TaskAdd, "sub": TaskSub, "subCtx": TaskSubWithContext,
Expand Down Expand Up @@ -184,3 +279,125 @@ func TestMultiTaskExtStartWorker(t *testing.T) {
})
})
}

// TestAsyncTaskExt_Monitor_Disabled 覆盖锚点 1:monitor_enable 默认关闭时,
// monitorEnabled 保持 false,任务正常执行不受影响、不 panic,且不产生任何
// asynctask_task_duration_seconds 记录。
func TestAsyncTaskExt_Monitor_Disabled(t *testing.T) {
startAsyncTaskMetricsServer()

assert.False(t, taskOne.monitorEnabled)

sign := &tasks.Signature{
Name: "add",
Args: []tasks.Arg{
{Type: "int64", Value: 5},
{Type: "int64", Value: 6},
},
}
asyncResult, err := taskOne.SendTask(sign)
if err != nil {
t.Fatal(err)
}
if _, err := asyncResult.Get(5 * time.Millisecond); err != nil {
t.Fatal(err)
}
time.Sleep(200 * time.Millisecond)

data := fetchAsyncTaskMetrics(t)
assert.NotContains(t, data,
`asynctask_task_duration_seconds_count{queue="gobay.task.one",status="unknown",task_name="add"}`)
}

// TestAsyncTaskExt_Monitor covers 锚点 2/3/9:monitor_enable=true 下,
// 两个同类型实例(three_asynctask_/four_asynctask_)各自 Init 不 panic,
// 成功任务和失败任务都会被 Observe,status 始终固定 "unknown"。
func TestAsyncTaskExt_Monitor(t *testing.T) {
startAsyncTaskMetricsServer()

taskThree = AsyncTaskExt{NS: "three_asynctask_"}
taskFour = AsyncTaskExt{NS: "four_asynctask_"}

t.Run("9: 同类型两个实例都开启 monitor_enable 各自 Init 不 panic", func(t *testing.T) {
assert.NotPanics(t, func() {
app, err := gobay.CreateApp(
"../../testdata",
"asynctaskmonitored",
map[gobay.Key]gobay.Extension{
"taskthree": &taskThree,
"taskfour": &taskFour,
},
)
if err != nil {
t.Fatal(err)
}
if err := app.Init(); err != nil {
t.Fatal(err)
}
})
assert.True(t, taskThree.monitorEnabled)
assert.True(t, taskFour.monitorEnabled)
})

if err := taskThree.RegisterWorkerHandlers(map[string]interface{}{
"addThree": TaskAddThree,
"failThree": TaskFailThree,
}); err != nil {
t.Fatal(err)
}
go func() {
// use default queue "gobay.task.three"
if err := taskThree.StartWorker("", 1, false); err != nil {
t.Error(err)
}
}()
time.Sleep(500 * time.Millisecond) // make sure the worker is started

t.Run("2: monitor_enable=true 下执行已注册任务成功 -> _count 增加、_sum > 0", func(t *testing.T) {
sign := &tasks.Signature{
Name: "addThree",
Args: []tasks.Arg{
{Type: "int64", Value: 3},
{Type: "int64", Value: 4},
},
}
asyncResult, err := taskThree.SendTask(sign)
if err != nil {
t.Fatal(err)
}
if results, err := asyncResult.Get(5 * time.Millisecond); err != nil {
t.Fatal(err)
} else if res, ok := results[0].Interface().(int64); !ok || res != 7 {
t.Fatalf("unexpected task result: %v", results)
}
time.Sleep(300 * time.Millisecond)

data := fetchAsyncTaskMetrics(t)
labels := `queue="gobay.task.three",status="unknown",task_name="addThree"`
assert.Contains(t, data, `asynctask_task_duration_seconds_count{`+labels+`} 1`)
sum := extractMetricValue(t, data, "asynctask_task_duration_seconds_sum", labels)
assert.Greater(t, sum, 0.0)
})

t.Run("3: monitor_enable=true 下执行返回 error 的任务仍被 Observe,status 仍为 unknown", func(t *testing.T) {
sign := &tasks.Signature{
Name: "failThree",
Args: []tasks.Arg{
{Type: "int64", Value: 1},
},
}
asyncResult, err := taskThree.SendTask(sign)
if err != nil {
t.Fatal(err)
}
// TaskFailThree itself returns an error; we only care whether the
// metric was recorded regardless of the task outcome, matching the
// "不因失败而跳过" 行为契约.
_, _ = asyncResult.Get(5 * time.Millisecond)
time.Sleep(300 * time.Millisecond)

data := fetchAsyncTaskMetrics(t)
labels := `queue="gobay.task.three",status="unknown",task_name="failThree"`
assert.Contains(t, data, `asynctask_task_duration_seconds_count{`+labels+`} 1`)
})
}
Loading
Loading