feat: diagnostics of workflow dispatcher

This commit is contained in:
Fu Diwei
2025-08-28 13:24:15 +08:00
parent 496341fd2e
commit 302dd4272a
8 changed files with 206 additions and 25 deletions
+6
View File
@@ -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"`
}
+15 -4
View File
@@ -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")
@@ -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")
+9
View File
@@ -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 {
+23
View File
@@ -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();
+7 -2
View File
@@ -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"
}
+7 -2
View File
@@ -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": "运行中"
}
+112 -17
View File
@@ -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 = () => {
<Divider />
<h2>{t("settings.diagnostics.workflow_dispatcher.title")}</h2>
<SettingsDiagnosticsWorkflowDispatcher className="md:max-w-160" />
<SettingsDiagnosticsWorkflowDispatcher />
</>
);
};
@@ -53,6 +54,8 @@ const SettingsDiagnosticsLogs = ({ className, style }: { className?: string; sty
},
{
refreshDeps: [page, pageSize],
debounceWait: 300,
debounceLeading: true,
onSuccess: (res) => {
if (page === 1) {
setListData([]);
@@ -112,6 +115,11 @@ const SettingsDiagnosticsLogs = ({ className, style }: { className?: string; sty
refreshData();
};
const handleRefreshClick = () => {
setPage(1);
refreshData();
};
const handleLoadMoreClick = () => {
setPage((prev) => prev + 1);
};
@@ -127,13 +135,11 @@ const SettingsDiagnosticsLogs = ({ className, style }: { className?: string; sty
</Button>
</div>
</Show>
<Show when={listData.length === 0}>
<Show when={listData.length === 0 && !loading}>
<Empty description={loadError ? getErrMsg(loadError) : t("common.text.nodata")} image={Empty.PRESENTED_IMAGE_SIMPLE}>
{loadError && (
<Button icon={<IconReload size="1.25em" />} type="primary" onClick={handleReloadClick}>
{t("common.button.reload")}
</Button>
)}
<Button icon={<IconReload size="1.25em" />} type="primary" onClick={handleReloadClick}>
{t("common.button.reload")}
</Button>
</Empty>
</Show>
<Show when={listData.length > 0}>
@@ -147,11 +153,19 @@ const SettingsDiagnosticsLogs = ({ className, style }: { className?: string; sty
);
})}
</div>
{hasMore && (
<a onClick={handleLoadMoreClick}>
<span className="text-xs">Load more</span>
<div className="flex w-full items-center">
<a onClick={handleRefreshClick}>
<span className="text-xs">{t("settings.diagnostics.logs.button.refresh.label")}</span>
</a>
)}
{hasMore && (
<>
<Divider type="vertical" />
<a onClick={handleLoadMoreClick}>
<span className="text-xs">{t("settings.diagnostics.logs.button.load_more.label")}</span>
</a>
</>
)}
</div>
</div>
</Show>
</div>
@@ -166,7 +180,9 @@ const SettingsDiagnosticsCrons = ({ className, style }: { className?: string; st
const [page, setPage] = useState(1);
const [pageSize, setPageSize] = useState(10);
type CronJob = Awaited<ReturnType<typeof listCronJobs>>["items"][number];
type CronJob = Awaited<ReturnType<typeof listCronJobs>>["items"][number] & {
nextTriggerTime: string;
};
const [listData, setListData] = useState<CronJob[]>([]);
const [listTotal, setListTotal] = useState(0);
@@ -180,7 +196,13 @@ const SettingsDiagnosticsCrons = ({ className, style }: { className?: string; st
const startIndex = (page - 1) * pageSize;
const endIndex = startIndex + pageSize;
return {
items: res.items.slice(startIndex, endIndex),
items: res.items
.slice(startIndex, endIndex)
.map((item) => ({
...item,
nextTriggerTime: dayjs(getNextCronExecutions(item.cron)[0]).format("YYYY-MM-DD HH:mm:ss"),
}))
.sort((a, b) => a.nextTriggerTime.localeCompare(b.nextTriggerTime)),
totalItems: res.items.length,
};
});
@@ -224,11 +246,12 @@ const SettingsDiagnosticsCrons = ({ className, style }: { className?: string; st
renderItem={(record) => (
<List.Item>
<Tooltip
className="block xl:hidden"
title={
<>
{t("settings.diagnostics.crons.job.next_trigger_time")}
{t("settings.diagnostics.crons.props.next_trigger_time")}
<br />
{dayjs(getNextCronExecutions(record.cron)[0]).format("YYYY-MM-DD HH:mm:ss")}
{record.nextTriggerTime}
</>
}
mouseEnterDelay={1}
@@ -243,6 +266,20 @@ const SettingsDiagnosticsCrons = ({ className, style }: { className?: string; st
</div>
</div>
</Tooltip>
<div className="hidden w-full items-center justify-between gap-4 overflow-hidden xl:flex">
<div className="flex-1 truncate">
<Typography.Text>{record.id}</Typography.Text>
</div>
<div className="flex items-center justify-end">
<Typography.Text type="secondary">{record.cron}</Typography.Text>
<Divider type="vertical" />
<Typography.Text type="secondary">
{t("settings.diagnostics.crons.props.next_trigger_time")}
{record.nextTriggerTime}
</Typography.Text>
</div>
</div>
</List.Item>
)}
/>
@@ -258,9 +295,67 @@ const SettingsDiagnosticsCrons = ({ className, style }: { className?: string; st
const SettingsDiagnosticsWorkflowDispatcher = ({ className, style }: { className?: string; style?: React.CSSProperties }) => {
const { t } = useTranslation();
type Statistics = Awaited<ReturnType<typeof getWorkflowStats>>["data"];
const [statistics, setStatistics] = useState<Statistics>();
const { loading } = useRequest(
() => {
return getWorkflowStats();
},
{
throttleWait: 1000,
throttleLeading: true,
pollingInterval: 3000,
pollingWhenHidden: true,
onSuccess: (res) => {
setStatistics(res.data);
},
}
);
return (
<div className={className} style={style}>
<div>TODO ...</div>
<div className="flex w-full flex-wrap items-stretch justify-center gap-4 sm:flex-nowrap">
<Card className="w-full sm:flex-1 md:w-1/3" loading={loading && !statistics}>
<Statistic title={t("settings.diagnostics.workflow_dispatcher.statistics.concurrency")} value={statistics?.concurrency ?? "-"} />
</Card>
<Tooltip
mouseEnterDelay={1}
placement="topLeft"
title={
statistics?.pendingRunIds?.length
? statistics?.pendingRunIds?.map((id) => (
<div key={id}>
<Tag>#{id}</Tag>
</div>
))
: null
}
>
<Card className="w-full sm:flex-1 md:w-1/3" loading={loading && !statistics}>
<Statistic title={t("settings.diagnostics.workflow_dispatcher.statistics.pending")} value={statistics?.pendingRunIds?.length ?? "-"} />
</Card>
</Tooltip>
<Tooltip
mouseEnterDelay={1}
placement="topLeft"
title={
statistics?.processingRunIds?.length
? statistics?.processingRunIds?.map((id) => (
<div key={id}>
<Tag>#{id}</Tag>
</div>
))
: null
}
>
<Card className="w-full sm:flex-1 md:w-1/3" loading={loading && !statistics}>
<Statistic title={t("settings.diagnostics.workflow_dispatcher.statistics.processing")} value={statistics?.processingRunIds?.length ?? "-"} />
</Card>
</Tooltip>
</div>
</div>
);
};