diff --git a/internal/domain/dtos/workflow.go b/internal/domain/dtos/workflow.go
index ea4b1fe0a..6f6a77fdb 100644
--- a/internal/domain/dtos/workflow.go
+++ b/internal/domain/dtos/workflow.go
@@ -17,3 +17,9 @@ type WorkflowCancelRunReq struct {
}
type WorkflowCancelRunResp struct{}
+
+type WorkflowStatisticsResp struct {
+ Concurrency int `json:"concurrency"`
+ PendingRunIds []string `json:"pendingRunIds"`
+ ProcessingRunIds []string `json:"processingRunIds"`
+}
diff --git a/internal/rest/handlers/workflow.go b/internal/rest/handlers/workflow.go
index ff3df00c9..4b7aa42b9 100644
--- a/internal/rest/handlers/workflow.go
+++ b/internal/rest/handlers/workflow.go
@@ -13,6 +13,7 @@ import (
)
type workflowService interface {
+ GetStatistics(ctx context.Context) (*dtos.WorkflowStatisticsResp, error)
StartRun(ctx context.Context, req *dtos.WorkflowStartRunReq) (*dtos.WorkflowStartRunResp, error)
CancelRun(ctx context.Context, req *dtos.WorkflowCancelRunReq) (*dtos.WorkflowCancelRunResp, error)
Shutdown(ctx context.Context)
@@ -28,11 +29,21 @@ func NewWorkflowHandler(router *router.RouterGroup[*core.RequestEvent], service
}
group := router.Group("/workflows")
- group.POST("/{workflowId}/runs", handler.run)
- group.POST("/{workflowId}/runs/{runId}/cancel", handler.cancel)
+ group.GET("/stats", handler.getStatistics)
+ group.POST("/{workflowId}/runs", handler.startRun)
+ group.POST("/{workflowId}/runs/{runId}/cancel", handler.cancelRun)
}
-func (handler *WorkflowHandler) run(e *core.RequestEvent) error {
+func (handler *WorkflowHandler) getStatistics(e *core.RequestEvent) error {
+ res, err := handler.service.GetStatistics(e.Request.Context())
+ if err != nil {
+ return resp.Err(e, err)
+ }
+
+ return resp.Ok(e, res)
+}
+
+func (handler *WorkflowHandler) startRun(e *core.RequestEvent) error {
req := &dtos.WorkflowStartRunReq{}
req.WorkflowId = e.Request.PathValue("workflowId")
if err := e.BindBody(req); err != nil {
@@ -50,7 +61,7 @@ func (handler *WorkflowHandler) run(e *core.RequestEvent) error {
return resp.Ok(e, res)
}
-func (handler *WorkflowHandler) cancel(e *core.RequestEvent) error {
+func (handler *WorkflowHandler) cancelRun(e *core.RequestEvent) error {
req := &dtos.WorkflowCancelRunReq{}
req.WorkflowId = e.Request.PathValue("workflowId")
req.RunId = e.Request.PathValue("runId")
diff --git a/internal/workflow/dispatcher/dispatcher.go b/internal/workflow/dispatcher/dispatcher.go
index e5c8de825..0e047f2a5 100644
--- a/internal/workflow/dispatcher/dispatcher.go
+++ b/internal/workflow/dispatcher/dispatcher.go
@@ -35,12 +35,20 @@ func init() {
}
type WorkflowDispatcher interface {
+ GetStatistics() Statistics
+
Bootup(ctx context.Context) error
Shutdown(ctx context.Context) error
Start(ctx context.Context, runId string) error
Cancel(ctx context.Context, runId string) error
}
+type Statistics struct {
+ Concurrency int
+ PendingRunIds []string
+ ProcessingRunIds []string
+}
+
type workflowDispatcher struct {
booted bool
concurrency int
@@ -58,6 +66,25 @@ type workflowDispatcher struct {
var _ WorkflowDispatcher = (*workflowDispatcher)(nil)
+func (wd *workflowDispatcher) GetStatistics() Statistics {
+ wd.taskMtx.RLock()
+ defer wd.taskMtx.RUnlock()
+
+ stats := Statistics{
+ Concurrency: wd.concurrency,
+ PendingRunIds: make([]string, 0),
+ ProcessingRunIds: make([]string, 0),
+ }
+ for _, pendingRunId := range wd.pendingRunQueue {
+ stats.PendingRunIds = append(stats.PendingRunIds, pendingRunId)
+ }
+ for _, processingRunId := range wd.processingTasks {
+ stats.ProcessingRunIds = append(stats.ProcessingRunIds, processingRunId.RunId)
+ }
+
+ return stats
+}
+
func (wd *workflowDispatcher) Bootup(ctx context.Context) error {
if wd.booted {
return errors.New("could not re-bootup")
diff --git a/internal/workflow/service.go b/internal/workflow/service.go
index 2ce5506db..d16e1c174 100644
--- a/internal/workflow/service.go
+++ b/internal/workflow/service.go
@@ -75,6 +75,15 @@ func (s *WorkflowService) InitSchedule(ctx context.Context) error {
return nil
}
+func (s *WorkflowService) GetStatistics(ctx context.Context) (*dtos.WorkflowStatisticsResp, error) {
+ stats := s.dispatcher.GetStatistics()
+ return &dtos.WorkflowStatisticsResp{
+ Concurrency: stats.Concurrency,
+ PendingRunIds: stats.PendingRunIds,
+ ProcessingRunIds: stats.ProcessingRunIds,
+ }, nil
+}
+
func (s *WorkflowService) StartRun(ctx context.Context, req *dtos.WorkflowStartRunReq) (*dtos.WorkflowStartRunResp, error) {
workflow, err := s.workflowRepo.GetById(ctx, req.WorkflowId)
if err != nil {
diff --git a/ui/src/api/workflows.ts b/ui/src/api/workflows.ts
index 6a9d69f6d..81688c896 100644
--- a/ui/src/api/workflows.ts
+++ b/ui/src/api/workflows.ts
@@ -3,6 +3,29 @@ import { ClientResponseError } from "pocketbase";
import { WORKFLOW_TRIGGERS } from "@/domain/workflow";
import { getPocketBase } from "@/repository/_pocketbase";
+export const getStats = async () => {
+ const pb = getPocketBase();
+
+ const resp = await pb.send<
+ BaseResponse<{
+ concurrency: number;
+ pendingRunIds: string[];
+ processingRunIds: string[];
+ }>
+ >(`/api/workflows/stats`, {
+ method: "GET",
+ headers: {
+ "Content-Type": "application/json",
+ },
+ });
+
+ if (resp.code != 0) {
+ throw new ClientResponseError({ status: resp.code, response: resp, data: {} });
+ }
+
+ return resp;
+};
+
export const startRun = async (workflowId: string) => {
const pb = getPocketBase();
diff --git a/ui/src/i18n/locales/en/nls.settings.json b/ui/src/i18n/locales/en/nls.settings.json
index 78f0d83e5..642a11a0b 100644
--- a/ui/src/i18n/locales/en/nls.settings.json
+++ b/ui/src/i18n/locales/en/nls.settings.json
@@ -50,7 +50,12 @@
"settings.diagnostics.tab": "Diagnostics",
"settings.diagnostics.logs.title": "System logs",
+ "settings.diagnostics.logs.button.refresh.label": "Refresh",
+ "settings.diagnostics.logs.button.load_more.label": "Load more",
"settings.diagnostics.crons.title": "CRON jobs",
- "settings.diagnostics.crons.job.next_trigger_time": "Expected next execution time: ",
- "settings.diagnostics.workflow_dispatcher.title": "Workflow dispatcher"
+ "settings.diagnostics.crons.props.next_trigger_time": "Expected next execution time: ",
+ "settings.diagnostics.workflow_dispatcher.title": "Workflow dispatcher",
+ "settings.diagnostics.workflow_dispatcher.statistics.concurrency": "Concurrency",
+ "settings.diagnostics.workflow_dispatcher.statistics.pending": "Pending",
+ "settings.diagnostics.workflow_dispatcher.statistics.processing": "Processing"
}
diff --git a/ui/src/i18n/locales/zh/nls.settings.json b/ui/src/i18n/locales/zh/nls.settings.json
index e804d0b92..94d398492 100644
--- a/ui/src/i18n/locales/zh/nls.settings.json
+++ b/ui/src/i18n/locales/zh/nls.settings.json
@@ -50,7 +50,12 @@
"settings.diagnostics.tab": "系统诊断",
"settings.diagnostics.logs.title": "系统日志",
+ "settings.diagnostics.logs.button.refresh.label": "刷新日志",
+ "settings.diagnostics.logs.button.load_more.label": "加载更多",
"settings.diagnostics.crons.title": "后台任务",
- "settings.diagnostics.crons.job.next_trigger_time": "预计下次运行时间:",
- "settings.diagnostics.workflow_dispatcher.title": "工作流调度器"
+ "settings.diagnostics.crons.props.next_trigger_time": "预计下次运行时间:",
+ "settings.diagnostics.workflow_dispatcher.title": "工作流调度器",
+ "settings.diagnostics.workflow_dispatcher.statistics.concurrency": "最大并发",
+ "settings.diagnostics.workflow_dispatcher.statistics.pending": "等待运行",
+ "settings.diagnostics.workflow_dispatcher.statistics.processing": "运行中"
}
diff --git a/ui/src/pages/settings/SettingsDiagnostics.tsx b/ui/src/pages/settings/SettingsDiagnostics.tsx
index 7e77278e3..f90f27d79 100644
--- a/ui/src/pages/settings/SettingsDiagnostics.tsx
+++ b/ui/src/pages/settings/SettingsDiagnostics.tsx
@@ -2,9 +2,10 @@ import { useState } from "react";
import { useTranslation } from "react-i18next";
import { IconReload } from "@tabler/icons-react";
import { useRequest } from "ahooks";
-import { Button, Divider, Empty, List, Pagination, Tooltip, Typography } from "antd";
+import { Button, Card, Divider, Empty, List, Pagination, Statistic, Tag, Tooltip, Typography } from "antd";
import dayjs from "dayjs";
+import { getStats as getWorkflowStats } from "@/api/workflows";
import Show from "@/components/Show";
import { listCronJobs, listLogs } from "@/repository/system";
import { getNextCronExecutions } from "@/utils/cron";
@@ -27,7 +28,7 @@ const SettingsDiagnostics = () => {