fix: show worker added message

This commit is contained in:
jialin
2026-06-05 10:20:55 +08:00
committed by jialin
parent faf6246fe6
commit 907187d53a
8 changed files with 118 additions and 28 deletions
@@ -1,6 +1,7 @@
import { workerAddedCountAtom } from '@/atoms/clusters';
import useSetChunkRequest from '@/hooks/use-chunk-request';
import useUpdateChunkedList from '@/hooks/use-update-chunk-list';
import useQueryWorkerList from '@/pages/resources/services/use-query-worker-list';
import { useAtom } from 'jotai';
import _ from 'lodash';
import qs from 'query-string';
@@ -13,7 +14,31 @@ export default function useAddWorkerMessage() {
const [addedCount, setAddedCount] = useState(0);
const timerRef = useRef<any>(null);
const triggerAtRef = useRef<number>(0);
const existingIdsRef = useRef<Set<string | number>>(new Set());
const snapshotReceivedRef = useRef<boolean>(false);
const [, setWorkerAddedCount] = useAtom(workerAddedCountAtom);
// fetchData auto-cancels any in-flight request on each call and on unmount,
// so a stale seed from a previous open/cluster can't leak in
const { fetchData: fetchWorkerList, cancelRequest: cancelWorkerListRequest } =
useQueryWorkerList();
const isNewWorker = (item: any) => {
console.log(
'isNewWorker item:',
item,
existingIdsRef.current,
snapshotReceivedRef.current
);
if (item?.id == null) {
return false;
}
if (existingIdsRef.current.has(item.id)) {
return false;
}
existingIdsRef.current.add(item.id);
// before the full snapshot has been received, every item is pre-existing
return snapshotReceivedRef.current;
};
const updateAddedCount = (count: number) => {
setAddedCount(count);
@@ -30,8 +55,9 @@ export default function useAddWorkerMessage() {
const { updateChunkedList } = useUpdateChunkedList({
events: ['CREATE', 'INSERT'],
dataList: [],
triggerAt: triggerAtRef,
isNewItem: isNewWorker,
onCreate: (newItems: any) => {
console.log('onCreate newItems:', newItems, triggerAtRef.current);
if (triggerAtRef.current) {
newItemsRef.current = newItemsRef.current.concat(newItems);
showAddWorkerMessage();
@@ -49,6 +75,27 @@ export default function useAddWorkerMessage() {
_.each(list, (data: any) => {
updateChunkedList(data);
});
// fallback: if seeding the baseline failed, the first watch chunk marks the
// snapshot as received (only reliable when there is at least one worker)
snapshotReceivedRef.current = true;
};
// Seed the baseline set of already-existing worker ids via REST before the
// watch starts. Relying on the watch's first chunk fails when a cluster has
// zero workers (no CREATE event is sent, so the snapshot flag never flips and
// genuinely new workers get misclassified as pre-existing).
const seedExistingWorkers = async (params: Record<string, any>) => {
try {
const items = await fetchWorkerList({ ...params, page: -1 } as any);
(items || []).forEach((item: any) => {
if (item?.id != null) {
existingIdsRef.current.add(item.id);
}
});
snapshotReceivedRef.current = true;
} catch (error) {
// ignore: fall back to the first watch chunk (see updateHandler)
}
};
const resetAddedCount = () => {
@@ -56,11 +103,18 @@ export default function useAddWorkerMessage() {
chunkRequestRef.current?.current?.cancel?.();
newItemsRef.current = [];
triggerAtRef.current = 0;
existingIdsRef.current = new Set();
snapshotReceivedRef.current = false;
// cancel any in-flight seed request
cancelWorkerListRequest();
clearTimeout(timerRef.current);
};
const createModelsChunkRequest = async (params = {}) => {
resetAddedCount();
// seed the baseline before watching so new workers are detected even when
// the cluster currently has zero workers
await seedExistingWorkers(params);
try {
chunkRequestRef.current = setChunkRequest({
url: `${WORKERS_API}?${qs.stringify(_.pickBy(params, (val: any) => !!val))}`,