Skip to content
Merged
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
161 changes: 144 additions & 17 deletions packages/core/realtime/use-realtime-sync-ws-instance.test.tsx
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
/**
* @vitest-environment jsdom
*/
import { QueryClient, QueryClientProvider, type InvalidateQueryFilters } from "@tanstack/react-query";
import { renderHook, waitFor } from "@testing-library/react";
import { QueryClient, QueryClientProvider, QueryObserver, type InvalidateQueryFilters, type QueryKey } from "@tanstack/react-query";
import { cleanup, renderHook, waitFor } from "@testing-library/react";
import type { ReactNode } from "react";
import { describe, expect, it, vi, beforeEach, afterEach } from "vitest";
import type { WSClient } from "../api/ws-client";
Expand All @@ -19,6 +19,8 @@ import {
} from "../workspace/pending-delete";
import { setApiInstance } from "../api";
import type { ApiClient } from "../api/client";
import type { IssueTableQuerySpec } from "../types";
import { getCurrentWsId } from "../platform/workspace-storage";
import { forgetLocalSearchIndex } from "../search-index/instance";
import { useRealtimeSync, type RealtimeSyncStores } from "./use-realtime-sync";

Expand All @@ -27,7 +29,7 @@ vi.mock("../search-index/instance", () => ({
}));

vi.mock("../platform/workspace-storage", () => ({
getCurrentWsId: () => "ws-1",
getCurrentWsId: vi.fn(() => "ws-1"),
getCurrentSlug: () => "test-ws",
// Draft stores are now loaded transitively (storage-cleanup → register-all-drafts)
// so their persist wiring must resolve against this mock.
Expand Down Expand Up @@ -379,33 +381,158 @@ describe("useRealtimeSync — Table server membership invalidation", () => {
let stores: RealtimeSyncStores;

beforeEach(() => {
qc = new QueryClient({ defaultOptions: { queries: { retry: false } } });
qc = new QueryClient({ defaultOptions: { queries: { retry: false, staleTime: Infinity } } });
stores = createStores();
});

afterEach(() => {
cleanup();
qc.clear();
vi.mocked(getCurrentWsId).mockReturnValue("ws-1");
vi.useRealTimers();
});

it("invalidates Table queries after a task lifecycle event", () => {
const spec: IssueTableQuerySpec = {
scope: { kind: "workspace" },
filters: {},
sort: { field: "position", direction: "asc" },
};
const workingFacet = (wsId = "ws-1") => issueKeys.tableFacets(wsId, {
query: spec, facets: [{ kind: "working_agents" }],
});

function mountRealtime() {
vi.useFakeTimers();
const ws = createMockWs();
const invalidate = vi.spyOn(qc, "invalidateQueries");
renderHook(() => useRealtimeSync(ws, stores), {
const hook = renderHook(() => useRealtimeSync(ws, stores), {
wrapper: createWrapper(qc),
});
const onAny = vi.mocked(ws.onAny).mock.calls[0]?.[0];
expect(onAny).toBeDefined();
const onAny = vi.mocked(ws.onAny).mock.calls[0]![0];
return { ...hook, emit: (type: string) => onAny({ type, payload: {} } as never) };
}

it("only invalidates working facets and queries whose membership depends on working state", async () => {
const { emit } = mountRealtime();
const untouched: QueryKey[] = [workingFacet("ws-2")];
const affected: QueryKey[] = [workingFacet(), workspaceWorkingAgentsKeys.list("ws-1", "issue")];
// Test both server membership forms, including an explicit empty id set,
// across rows, descriptors and ordinary facets such as status counts.
const filterCases: IssueTableQuerySpec["filters"][] = [{}, { working_only: false }, { working_only: true }, { working_issue_ids: [] }, { working_issue_ids: ["i1"] }];
for (const filters of filterCases) {
const query = { ...spec, filters };
const keys = [
issueKeys.tableRows("ws-1", query, { kind: "none" }, null, false, null),
issueKeys.tableGroups("ws-1", query, { kind: "status" }),
issueKeys.tableFacets("ws-1", { query, facets: [{ kind: "status" }] }),
];
(filters.working_only || filters.working_issue_ids ? affected : untouched).push(...keys);
}
for (const key of [...untouched, ...affected]) qc.setQueryData(key, []);
emit("task:completed");
await vi.advanceTimersByTimeAsync(1_000);
for (const key of affected) expect(qc.getQueryState(key)?.isInvalidated).toBe(true);
for (const key of untouched) expect(qc.getQueryState(key)?.isInvalidated).toBe(false);
});

onAny!({ type: "task:completed", payload: {} } as never);
vi.advanceTimersByTime(100);
it("coalesces a continuous lifecycle stream and eventually clears the final completed run", async () => {
const { emit } = mountRealtime();
let count = 1;
const queryFn = vi.fn(async () => count);
qc.setQueryData(workingFacet(), count);
const observer = new QueryObserver(qc, { queryKey: workingFacet(), queryFn });
const unsubscribe = observer.subscribe(() => {});
for (let i = 0; i < 20; i++) {
emit(i % 2 ? "agent:updated" : "task:started");
await vi.advanceTimersByTimeAsync(200);
}
expect(queryFn).toHaveBeenCalledTimes(4);
count = 0;
emit("task:completed");
await vi.advanceTimersByTimeAsync(1_000);
expect(queryFn).toHaveBeenCalledTimes(5);
expect(observer.getCurrentResult().data).toBe(0);
unsubscribe();
});

expect(invalidate).toHaveBeenCalledWith({
queryKey: issueKeys.tableAll("ws-1"),
});
expect(invalidate).toHaveBeenCalledWith({
queryKey: workspaceWorkingAgentsKeys.all("ws-1"),
});
it.each(["facet", "projection"])("restarts an in-flight first %s request after completion instead of accepting its stale response", async (kind) => {
const { emit } = mountRealtime();
const queryKey = kind === "facet" ? workingFacet() : workspaceWorkingAgentsKeys.list("ws-1", "issue");
let finishFirst!: (value: number) => void;
const queryFn = vi.fn(async () => 0).mockImplementationOnce(() => new Promise<number>((resolve) => { finishFirst = resolve; }));
const observer = new QueryObserver(qc, { queryKey, queryFn });
const unsubscribe = observer.subscribe(() => {});
emit("task:completed");
await vi.advanceTimersByTimeAsync(1_000);
expect(queryFn).toHaveBeenCalledTimes(2);
expect(observer.getCurrentResult().data).toBe(0);
finishFirst(1);
await vi.advanceTimersByTimeAsync(0);
expect(observer.getCurrentResult().data).toBe(0);
unsubscribe();
});

it("lets slow refreshes finish and performs a trailing refresh for events received in flight", async () => {
const { emit } = mountRealtime();
let finishSlow!: (value: number) => void;
const queryFn = vi.fn(async () => 0).mockImplementationOnce(() => new Promise<number>((resolve) => { finishSlow = resolve; }));
qc.setQueryData(workingFacet(), 1);
const observer = new QueryObserver(qc, { queryKey: workingFacet(), queryFn });
const unsubscribe = observer.subscribe(() => {});
emit("task:started");
await vi.advanceTimersByTimeAsync(1_000);
for (let i = 0; i < 10; i++) {
emit("task:completed");
await vi.advanceTimersByTimeAsync(200);
}
expect(queryFn).toHaveBeenCalledTimes(1);
finishSlow(1);
await vi.advanceTimersByTimeAsync(1_000);
expect(queryFn).toHaveBeenCalledTimes(2);
expect(observer.getCurrentResult().data).toBe(0);
unsubscribe();
});

it("ignores progress/messages and discards a queued refresh on unmount", async () => {
const { emit, unmount } = mountRealtime();
qc.setQueryData(workingFacet(), []);
qc.setQueryData(workspaceWorkingAgentsKeys.list("ws-1", "issue"), []);
emit("task:progress");
emit("task:message");
await vi.advanceTimersByTimeAsync(2_000);
expect(qc.getQueryState(workingFacet())?.isInvalidated).toBe(false);
expect(qc.getQueryState(workspaceWorkingAgentsKeys.list("ws-1", "issue"))?.isInvalidated).toBe(false);
emit("task:completed");
unmount();
await vi.advanceTimersByTimeAsync(1_000);
expect(qc.getQueryState(workingFacet())?.isInvalidated).toBe(false);
});

it("keeps a queued refresh scoped to the workspace that received it", async () => {
const { emit } = mountRealtime();
qc.setQueryData(workingFacet(), []);
qc.setQueryData(workingFacet("ws-2"), []);
emit("task:completed");
vi.mocked(getCurrentWsId).mockReturnValue("ws-2");
await vi.advanceTimersByTimeAsync(1_000);
expect(qc.getQueryState(workingFacet())?.isInvalidated).toBe(true);
expect(qc.getQueryState(workingFacet("ws-2"))?.isInvalidated).toBe(false);
});

it("does not schedule a trailing refresh after unmounting during a slow request", async () => {
const { emit, unmount } = mountRealtime();
let finish!: (value: number) => void;
const queryFn = vi.fn(() => new Promise<number>((resolve) => { finish = resolve; }));
qc.setQueryData(workingFacet(), 1);
const observer = new QueryObserver(qc, { queryKey: workingFacet(), queryFn });
const unsubscribe = observer.subscribe(() => {});
emit("task:started");
await vi.advanceTimersByTimeAsync(1_000);
emit("task:completed");
unmount();
finish(0);
await vi.advanceTimersByTimeAsync(2_000);
expect(queryFn).toHaveBeenCalledTimes(1);
unsubscribe();
});

it("invalidates Table queries after a property definition changes", () => {
Expand Down
77 changes: 69 additions & 8 deletions packages/core/realtime/use-realtime-sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

import { useEffect, useRef } from "react";
import { forgetLocalSearchIndex } from "../search-index/instance";
import { useQueryClient, type InfiniteData, type QueryClient } from "@tanstack/react-query";
import { useQueryClient, type InfiniteData, type QueryClient, type QueryFilters } from "@tanstack/react-query";
import type { WSClient } from "../api/ws-client";
import type { StoreApi, UseBoundStore } from "zustand";
import type { AuthState } from "../auth/store";
Expand Down Expand Up @@ -73,6 +73,8 @@ import {
} from "../chat/pending";
import { resolvePostAuthDestination, useHasOnboarded } from "../paths";
import type {
IssueTableFacetsRequest,
IssueTableQuerySpec,
MemberAddedPayload,
WorkspaceDeletedPayload,
WorkspaceUpdatedPayload,
Expand Down Expand Up @@ -698,6 +700,36 @@ function invalidateWorkspaceScopedQueries(qc: QueryClient): void {
qc.invalidateQueries({ queryKey: workspaceKeys.list() });
}

async function refreshWorkingAgentQueries(qc: QueryClient, wsId: string): Promise<void> {
const filters: QueryFilters[] = [
{ queryKey: workspaceWorkingAgentsKeys.all(wsId) },
{
queryKey: issueKeys.tableAll(wsId),
predicate: ({ queryKey }) => {
const kind = queryKey[3];
let spec: IssueTableQuerySpec | undefined;
if (kind === "facets") {
const request = queryKey[4] as IssueTableFacetsRequest;
if (request.facets.some((facet) => facet.kind === "working_agents")) return true;
spec = request.query;
} else if (kind === "rows" || kind === "groups") {
spec = queryKey[4] as IssueTableQuerySpec;
}
// An explicit empty id list is still a working filter: its next
// projection may contain newly started work. Other facets also
// depend on this membership, not only rows and group descriptors.
return spec?.filters.working_only === true ||
spec?.filters.working_issue_ids !== undefined;
},
},
];
// A lifecycle event during the FIRST fetch must not dedupe onto a response
// captured before that event. With staleTime=Infinity it could otherwise
// leave a completed run visible until another event happens to arrive.
await Promise.all(filters.map((filter) => qc.cancelQueries(filter)));
await Promise.all(filters.map((filter) => qc.invalidateQueries(filter)));
}

function invalidateSquadMemberStatusQueries(qc: QueryClient, wsId: string): void {
qc.invalidateQueries({
predicate: (query) => {
Expand Down Expand Up @@ -768,7 +800,6 @@ export function useRealtimeSync(
const wsId = getCurrentWsId();
if (wsId) {
qc.invalidateQueries({ queryKey: workspaceKeys.agents(wsId) });
qc.invalidateQueries({ queryKey: workspaceWorkingAgentsKeys.all(wsId) });
// Squad members status is derived per agent, so any agent
// change (status flip, archive, runtime swap) needs to refresh the
// per-squad members-status cache without refetching the static squad
Expand Down Expand Up @@ -916,12 +947,7 @@ export function useRealtimeSync(
const wsId = getCurrentWsId();
if (!wsId) return;
qc.invalidateQueries({ queryKey: agentTaskSnapshotKeys.list(wsId) });
qc.invalidateQueries({ queryKey: workspaceWorkingAgentsKeys.all(wsId) });
// The Table working-agent shortcut derives an assignee set from the
// projection above. Refresh its server-owned graph alongside that set
// so rows/groups/facets cannot remain on an old task transition while
// the projection refetches (global staleTime is Infinity).
qc.invalidateQueries({ queryKey: issueKeys.tableAll(wsId) });
// Working-agent projections have their own bounded refresh below.
// 30d activity series shares the same lifecycle signal — any task
// completion / failure shifts the histogram. (Dispatch alone
// doesn't change a completed_at-anchored series, but invalidating
Expand Down Expand Up @@ -965,6 +991,33 @@ export function useRealtimeSync(
},
};

const workingAgentTimers = new Map<string, ReturnType<typeof setTimeout>>();
const workingAgentInFlight = new Set<string>();
const workingAgentDirty = new Set<string>();
let disposed = false;
const scheduleWorkingAgentRefresh = (wsId: string) => {
// Flush one second after the FIRST event in a batch. Later events do
// not postpone it, so continuous activity cannot starve the projection.
// Capture the workspace now: navigation must not redirect a queued
// refresh into a different workspace's cache.
workingAgentDirty.add(wsId);
if (workingAgentTimers.has(wsId) || workingAgentInFlight.has(wsId)) return;
workingAgentTimers.set(wsId, setTimeout(async () => {
workingAgentTimers.delete(wsId);
workingAgentDirty.delete(wsId);
workingAgentInFlight.add(wsId);
try {
await refreshWorkingAgentQueries(qc, wsId);
} finally {
workingAgentInFlight.delete(wsId);
// Serialize slow requests: cancelling our own refresh every second
// would starve it under a continuous stream. Events received while
// it runs still get one trailing refresh, including the final stop.
if (!disposed && workingAgentDirty.has(wsId)) scheduleWorkingAgentRefresh(wsId);
}
}, 1_000));
};

const timers = new Map<string, ReturnType<typeof setTimeout>>();
const debouncedRefresh = (prefix: string, fn: () => void) => {
const existing = timers.get(prefix);
Expand Down Expand Up @@ -1007,6 +1060,10 @@ export function useRealtimeSync(
const unsubAny = ws.onAny((msg) => {
if (specificEvents.has(msg.type)) return;
const prefix = msg.type.split(":")[0] ?? "";
if (prefix === "agent" || (prefix === "task" && msg.type !== "task:progress")) {
const wsId = getCurrentWsId();
if (wsId) scheduleWorkingAgentRefresh(wsId);
}
const refresh = refreshMap[prefix];
if (refresh) debouncedRefresh(prefix, refresh);
});
Expand Down Expand Up @@ -1802,6 +1859,10 @@ export function useRealtimeSync(
if (aggregateRefreshTimer) clearTimeout(aggregateRefreshTimer);
timers.forEach(clearTimeout);
timers.clear();
disposed = true;
workingAgentTimers.forEach(clearTimeout);
workingAgentTimers.clear();
workingAgentDirty.clear();
};
}, [ws, qc, authStore, onToast]);

Expand Down
Loading
Loading