|
|
|
|
@@ -27,7 +27,10 @@ import {
|
|
|
|
|
import type { AgentLifecycleStatus } from "@server/shared/agent-lifecycle";
|
|
|
|
|
import type { DaemonClient } from "@server/client/daemon-client";
|
|
|
|
|
import { File } from "expo-file-system";
|
|
|
|
|
import { useHostRuntimeSession } from "@/runtime/host-runtime";
|
|
|
|
|
import {
|
|
|
|
|
getHostRuntimeStore,
|
|
|
|
|
useHostRuntimeSession,
|
|
|
|
|
} from "@/runtime/host-runtime";
|
|
|
|
|
import {
|
|
|
|
|
useSessionStore,
|
|
|
|
|
type Agent,
|
|
|
|
|
@@ -67,6 +70,7 @@ export type {
|
|
|
|
|
} from "@/stores/session-store";
|
|
|
|
|
|
|
|
|
|
const HISTORY_STALE_AFTER_MS = 60_000;
|
|
|
|
|
const AUTHORITATIVE_REVALIDATION_DEBOUNCE_MS = 300;
|
|
|
|
|
|
|
|
|
|
const findLatestAssistantMessageText = (items: StreamItem[]): string | null => {
|
|
|
|
|
for (let i = items.length - 1; i >= 0; i -= 1) {
|
|
|
|
|
@@ -286,20 +290,6 @@ function SessionProviderInternal({
|
|
|
|
|
(state) => state.sessions[serverId]?.agents
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
const handleAppResumed = useCallback(
|
|
|
|
|
(awayMs: number) => {
|
|
|
|
|
if (awayMs < HISTORY_STALE_AFTER_MS) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
bumpHistorySyncGeneration(serverId);
|
|
|
|
|
},
|
|
|
|
|
[bumpHistorySyncGeneration, serverId]
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
// Client activity tracking (heartbeat, push token registration)
|
|
|
|
|
useClientActivity({ client, focusedAgentId, onAppResumed: handleAppResumed });
|
|
|
|
|
usePushTokenRegistration({ client, serverId });
|
|
|
|
|
|
|
|
|
|
// State for voice detection flags (will be set by RealtimeContext)
|
|
|
|
|
const isDetectingRef = useRef(false);
|
|
|
|
|
const isSpeakingRef = useRef(false);
|
|
|
|
|
@@ -326,6 +316,12 @@ function SessionProviderInternal({
|
|
|
|
|
);
|
|
|
|
|
const attentionNotifiedRef = useRef<Map<string, number>>(new Map());
|
|
|
|
|
const appStateRef = useRef(AppState.currentState);
|
|
|
|
|
const revalidationTimerRef = useRef<ReturnType<typeof setTimeout> | null>(
|
|
|
|
|
null
|
|
|
|
|
);
|
|
|
|
|
const revalidationInFlightRef = useRef<Promise<void> | null>(null);
|
|
|
|
|
const revalidationQueuedRef = useRef(false);
|
|
|
|
|
const wasConnectedRef = useRef(isConnected);
|
|
|
|
|
|
|
|
|
|
useEffect(() => {
|
|
|
|
|
const subscription = AppState.addEventListener("change", (nextState) => {
|
|
|
|
|
@@ -349,6 +345,233 @@ function SessionProviderInternal({
|
|
|
|
|
previousAgentStatusRef.current = nextStatuses;
|
|
|
|
|
}, [sessionAgents]);
|
|
|
|
|
|
|
|
|
|
const hydrateWorkspaces = useCallback(
|
|
|
|
|
async (options?: { subscribe?: boolean; isCancelled?: () => boolean }) => {
|
|
|
|
|
if (!client || !isConnected) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const workspaces = new Map<string, WorkspaceDescriptor>();
|
|
|
|
|
let cursor: string | null = null;
|
|
|
|
|
let includeSubscribe = options?.subscribe ?? false;
|
|
|
|
|
|
|
|
|
|
while (true) {
|
|
|
|
|
const payload = await client.fetchWorkspaces({
|
|
|
|
|
sort: [{ key: "activity_at", direction: "desc" }],
|
|
|
|
|
...(includeSubscribe ? { subscribe: {} } : {}),
|
|
|
|
|
page: cursor ? { limit: 200, cursor } : { limit: 200 },
|
|
|
|
|
});
|
|
|
|
|
if (options?.isCancelled?.()) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for (const entry of payload.entries) {
|
|
|
|
|
const workspace = normalizeWorkspaceDescriptor(entry);
|
|
|
|
|
workspaces.set(workspace.id, workspace);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (!payload.pageInfo.hasMore || !payload.pageInfo.nextCursor) {
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
cursor = payload.pageInfo.nextCursor;
|
|
|
|
|
includeSubscribe = false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (options?.isCancelled?.()) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
setWorkspaces(serverId, workspaces);
|
|
|
|
|
setHasHydratedWorkspaces(serverId, true);
|
|
|
|
|
},
|
|
|
|
|
[client, isConnected, serverId, setHasHydratedWorkspaces, setWorkspaces]
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
const applyAuthoritativeAgentSnapshot = useCallback(
|
|
|
|
|
(agent: Agent) => {
|
|
|
|
|
setAgents(serverId, (prev) => {
|
|
|
|
|
const current = prev.get(agent.id);
|
|
|
|
|
if (current && agent.updatedAt.getTime() < current.updatedAt.getTime()) {
|
|
|
|
|
return prev;
|
|
|
|
|
}
|
|
|
|
|
const next = new Map(prev);
|
|
|
|
|
next.set(agent.id, agent);
|
|
|
|
|
return next;
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
if (agent.archivedAt) {
|
|
|
|
|
clearArchiveAgentPending({
|
|
|
|
|
queryClient,
|
|
|
|
|
serverId,
|
|
|
|
|
agentId: agent.id,
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
setAgentLastActivity(agent.id, agent.lastActivityAt);
|
|
|
|
|
|
|
|
|
|
setPendingPermissions(serverId, (prev) => {
|
|
|
|
|
const existingKeysForAgent: string[] = [];
|
|
|
|
|
for (const [key, pending] of prev.entries()) {
|
|
|
|
|
if (pending.agentId === agent.id) {
|
|
|
|
|
existingKeysForAgent.push(key);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const nextEntries = agent.pendingPermissions.map((request) => ({
|
|
|
|
|
key: derivePendingPermissionKey(agent.id, request),
|
|
|
|
|
agentId: agent.id,
|
|
|
|
|
request,
|
|
|
|
|
}));
|
|
|
|
|
|
|
|
|
|
let changed = existingKeysForAgent.length !== nextEntries.length;
|
|
|
|
|
if (!changed) {
|
|
|
|
|
const existingKeySet = new Set(existingKeysForAgent);
|
|
|
|
|
for (const entry of nextEntries) {
|
|
|
|
|
const existing = prev.get(entry.key);
|
|
|
|
|
if (!existingKeySet.has(entry.key) || !existing) {
|
|
|
|
|
changed = true;
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const currentRequest = existing.request;
|
|
|
|
|
if (
|
|
|
|
|
currentRequest.id !== entry.request.id ||
|
|
|
|
|
currentRequest.kind !== entry.request.kind ||
|
|
|
|
|
currentRequest.name !== entry.request.name ||
|
|
|
|
|
currentRequest.title !== entry.request.title ||
|
|
|
|
|
currentRequest.description !== entry.request.description
|
|
|
|
|
) {
|
|
|
|
|
changed = true;
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (!changed) {
|
|
|
|
|
return prev;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const next = new Map(prev);
|
|
|
|
|
for (const key of existingKeysForAgent) {
|
|
|
|
|
next.delete(key);
|
|
|
|
|
}
|
|
|
|
|
for (const entry of nextEntries) {
|
|
|
|
|
next.set(entry.key, entry);
|
|
|
|
|
}
|
|
|
|
|
return next;
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
const prevStatus = previousAgentStatusRef.current.get(agent.id);
|
|
|
|
|
if (prevStatus === "running" && agent.status !== "running") {
|
|
|
|
|
const session = useSessionStore.getState().sessions[serverId];
|
|
|
|
|
const queue = session?.queuedMessages.get(agent.id);
|
|
|
|
|
if (queue && queue.length > 0) {
|
|
|
|
|
const [next, ...rest] = queue;
|
|
|
|
|
console.log(
|
|
|
|
|
"[Session] Flushing queued message for agent:",
|
|
|
|
|
agent.id,
|
|
|
|
|
next.text
|
|
|
|
|
);
|
|
|
|
|
if (sendAgentMessageRef.current) {
|
|
|
|
|
void sendAgentMessageRef.current(agent.id, next.text, next.images);
|
|
|
|
|
}
|
|
|
|
|
setQueuedMessages(serverId, (prev) => {
|
|
|
|
|
const updated = new Map(prev);
|
|
|
|
|
updated.set(agent.id, rest);
|
|
|
|
|
return updated;
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
previousAgentStatusRef.current.set(agent.id, agent.status);
|
|
|
|
|
},
|
|
|
|
|
[
|
|
|
|
|
queryClient,
|
|
|
|
|
serverId,
|
|
|
|
|
setAgentLastActivity,
|
|
|
|
|
setAgents,
|
|
|
|
|
setPendingPermissions,
|
|
|
|
|
setQueuedMessages,
|
|
|
|
|
]
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
const runAuthoritativeRevalidation = useCallback(async () => {
|
|
|
|
|
await Promise.all([
|
|
|
|
|
getHostRuntimeStore().refreshAgentDirectory({ serverId }),
|
|
|
|
|
hydrateWorkspaces(),
|
|
|
|
|
]);
|
|
|
|
|
}, [hydrateWorkspaces, serverId]);
|
|
|
|
|
|
|
|
|
|
const flushAuthoritativeRevalidation = useCallback(() => {
|
|
|
|
|
if (!client || !isConnected) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
if (revalidationInFlightRef.current) {
|
|
|
|
|
revalidationQueuedRef.current = true;
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const run = runAuthoritativeRevalidation()
|
|
|
|
|
.catch((error) => {
|
|
|
|
|
console.error("[Session] authoritative revalidation failed", {
|
|
|
|
|
serverId,
|
|
|
|
|
error,
|
|
|
|
|
});
|
|
|
|
|
})
|
|
|
|
|
.finally(() => {
|
|
|
|
|
if (revalidationInFlightRef.current === run) {
|
|
|
|
|
revalidationInFlightRef.current = null;
|
|
|
|
|
}
|
|
|
|
|
if (!revalidationQueuedRef.current) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
revalidationQueuedRef.current = false;
|
|
|
|
|
if (revalidationTimerRef.current) {
|
|
|
|
|
clearTimeout(revalidationTimerRef.current);
|
|
|
|
|
}
|
|
|
|
|
revalidationTimerRef.current = setTimeout(() => {
|
|
|
|
|
revalidationTimerRef.current = null;
|
|
|
|
|
flushAuthoritativeRevalidation();
|
|
|
|
|
}, AUTHORITATIVE_REVALIDATION_DEBOUNCE_MS);
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
revalidationInFlightRef.current = run;
|
|
|
|
|
}, [client, isConnected, runAuthoritativeRevalidation, serverId]);
|
|
|
|
|
|
|
|
|
|
const scheduleAuthoritativeRevalidation = useCallback(() => {
|
|
|
|
|
if (!client || !isConnected) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
revalidationQueuedRef.current = true;
|
|
|
|
|
if (revalidationTimerRef.current) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
revalidationTimerRef.current = setTimeout(() => {
|
|
|
|
|
revalidationTimerRef.current = null;
|
|
|
|
|
if (!revalidationQueuedRef.current) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
revalidationQueuedRef.current = false;
|
|
|
|
|
flushAuthoritativeRevalidation();
|
|
|
|
|
}, AUTHORITATIVE_REVALIDATION_DEBOUNCE_MS);
|
|
|
|
|
}, [client, flushAuthoritativeRevalidation, isConnected]);
|
|
|
|
|
|
|
|
|
|
const handleAppResumed = useCallback(
|
|
|
|
|
(awayMs: number) => {
|
|
|
|
|
scheduleAuthoritativeRevalidation();
|
|
|
|
|
if (awayMs < HISTORY_STALE_AFTER_MS) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
bumpHistorySyncGeneration(serverId);
|
|
|
|
|
},
|
|
|
|
|
[bumpHistorySyncGeneration, scheduleAuthoritativeRevalidation, serverId]
|
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
// Client activity tracking (heartbeat, push token registration)
|
|
|
|
|
useClientActivity({ client, focusedAgentId, onAppResumed: handleAppResumed });
|
|
|
|
|
usePushTokenRegistration({ client, serverId });
|
|
|
|
|
|
|
|
|
|
const notifyAgentAttention = useCallback(
|
|
|
|
|
(params: {
|
|
|
|
|
agentId: string;
|
|
|
|
|
@@ -450,38 +673,11 @@ function SessionProviderInternal({
|
|
|
|
|
|
|
|
|
|
let cancelled = false;
|
|
|
|
|
void (async () => {
|
|
|
|
|
const workspaces = new Map<string, WorkspaceDescriptor>();
|
|
|
|
|
let cursor: string | null = null;
|
|
|
|
|
let includeSubscribe = true;
|
|
|
|
|
|
|
|
|
|
try {
|
|
|
|
|
while (true) {
|
|
|
|
|
const payload = await client.fetchWorkspaces({
|
|
|
|
|
sort: [{ key: "activity_at", direction: "desc" }],
|
|
|
|
|
...(includeSubscribe ? { subscribe: {} } : {}),
|
|
|
|
|
page: cursor ? { limit: 200, cursor } : { limit: 200 },
|
|
|
|
|
});
|
|
|
|
|
if (cancelled) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for (const entry of payload.entries) {
|
|
|
|
|
const workspace = normalizeWorkspaceDescriptor(entry);
|
|
|
|
|
workspaces.set(workspace.id, workspace);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (!payload.pageInfo.hasMore || !payload.pageInfo.nextCursor) {
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
cursor = payload.pageInfo.nextCursor;
|
|
|
|
|
includeSubscribe = false;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (cancelled) {
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
setWorkspaces(serverId, workspaces);
|
|
|
|
|
setHasHydratedWorkspaces(serverId, true);
|
|
|
|
|
await hydrateWorkspaces({
|
|
|
|
|
subscribe: true,
|
|
|
|
|
isCancelled: () => cancelled,
|
|
|
|
|
});
|
|
|
|
|
} catch (error) {
|
|
|
|
|
console.error("[Session] Failed to hydrate workspaces:", error);
|
|
|
|
|
}
|
|
|
|
|
@@ -490,7 +686,7 @@ function SessionProviderInternal({
|
|
|
|
|
return () => {
|
|
|
|
|
cancelled = true;
|
|
|
|
|
};
|
|
|
|
|
}, [client, isConnected, serverId, setHasHydratedWorkspaces, setWorkspaces]);
|
|
|
|
|
}, [client, hydrateWorkspaces, isConnected]);
|
|
|
|
|
|
|
|
|
|
const applyAgentUpdatePayload = useCallback(
|
|
|
|
|
(update: AgentUpdatePayload) => {
|
|
|
|
|
@@ -553,111 +749,13 @@ function SessionProviderInternal({
|
|
|
|
|
}),
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
setAgents(serverId, (prev) => {
|
|
|
|
|
const current = prev.get(agent.id);
|
|
|
|
|
if (current && agent.updatedAt.getTime() < current.updatedAt.getTime()) {
|
|
|
|
|
return prev;
|
|
|
|
|
}
|
|
|
|
|
const next = new Map(prev);
|
|
|
|
|
next.set(agent.id, agent);
|
|
|
|
|
return next;
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
if (agent.archivedAt) {
|
|
|
|
|
clearArchiveAgentPending({
|
|
|
|
|
queryClient,
|
|
|
|
|
serverId,
|
|
|
|
|
agentId: agent.id,
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Update agentLastActivity slice (top-level)
|
|
|
|
|
setAgentLastActivity(agent.id, agent.lastActivityAt);
|
|
|
|
|
|
|
|
|
|
setPendingPermissions(serverId, (prev) => {
|
|
|
|
|
const existingKeysForAgent: string[] = [];
|
|
|
|
|
for (const [key, pending] of prev.entries()) {
|
|
|
|
|
if (pending.agentId === agent.id) {
|
|
|
|
|
existingKeysForAgent.push(key);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const nextEntries = agent.pendingPermissions.map((request) => ({
|
|
|
|
|
key: derivePendingPermissionKey(agent.id, request),
|
|
|
|
|
agentId: agent.id,
|
|
|
|
|
request,
|
|
|
|
|
}));
|
|
|
|
|
|
|
|
|
|
let changed = existingKeysForAgent.length !== nextEntries.length;
|
|
|
|
|
if (!changed) {
|
|
|
|
|
const existingKeySet = new Set(existingKeysForAgent);
|
|
|
|
|
for (const entry of nextEntries) {
|
|
|
|
|
const existing = prev.get(entry.key);
|
|
|
|
|
if (!existingKeySet.has(entry.key) || !existing) {
|
|
|
|
|
changed = true;
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const currentRequest = existing.request;
|
|
|
|
|
if (
|
|
|
|
|
currentRequest.id !== entry.request.id ||
|
|
|
|
|
currentRequest.kind !== entry.request.kind ||
|
|
|
|
|
currentRequest.name !== entry.request.name ||
|
|
|
|
|
currentRequest.title !== entry.request.title ||
|
|
|
|
|
currentRequest.description !== entry.request.description
|
|
|
|
|
) {
|
|
|
|
|
changed = true;
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if (!changed) {
|
|
|
|
|
return prev;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
const next = new Map(prev);
|
|
|
|
|
for (const key of existingKeysForAgent) {
|
|
|
|
|
next.delete(key);
|
|
|
|
|
}
|
|
|
|
|
for (const entry of nextEntries) {
|
|
|
|
|
next.set(entry.key, entry);
|
|
|
|
|
}
|
|
|
|
|
return next;
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
// Flush queued messages when agent transitions from running to not running
|
|
|
|
|
const prevStatus = previousAgentStatusRef.current.get(agent.id);
|
|
|
|
|
if (prevStatus === "running" && agent.status !== "running") {
|
|
|
|
|
const session = useSessionStore.getState().sessions[serverId];
|
|
|
|
|
const queue = session?.queuedMessages.get(agent.id);
|
|
|
|
|
if (queue && queue.length > 0) {
|
|
|
|
|
const [next, ...rest] = queue;
|
|
|
|
|
console.log(
|
|
|
|
|
"[Session] Flushing queued message for agent:",
|
|
|
|
|
agent.id,
|
|
|
|
|
next.text
|
|
|
|
|
);
|
|
|
|
|
if (sendAgentMessageRef.current) {
|
|
|
|
|
void sendAgentMessageRef.current(agent.id, next.text, next.images);
|
|
|
|
|
}
|
|
|
|
|
setQueuedMessages(serverId, (prev) => {
|
|
|
|
|
const updated = new Map(prev);
|
|
|
|
|
updated.set(agent.id, rest);
|
|
|
|
|
return updated;
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
previousAgentStatusRef.current.set(agent.id, agent.status);
|
|
|
|
|
applyAuthoritativeAgentSnapshot(agent);
|
|
|
|
|
},
|
|
|
|
|
[
|
|
|
|
|
queryClient,
|
|
|
|
|
applyAuthoritativeAgentSnapshot,
|
|
|
|
|
serverId,
|
|
|
|
|
setAgents,
|
|
|
|
|
setAgentLastActivity,
|
|
|
|
|
setPendingPermissions,
|
|
|
|
|
setQueuedMessages,
|
|
|
|
|
setAgentTimelineCursor,
|
|
|
|
|
]
|
|
|
|
|
);
|
|
|
|
|
@@ -701,6 +799,14 @@ function SessionProviderInternal({
|
|
|
|
|
const currentTail = session?.agentStreamTail.get(agentId) ?? [];
|
|
|
|
|
const currentHead = session?.agentStreamHead.get(agentId) ?? [];
|
|
|
|
|
|
|
|
|
|
if (payload.agent) {
|
|
|
|
|
const normalized = normalizeAgentSnapshot(payload.agent, serverId);
|
|
|
|
|
applyAuthoritativeAgentSnapshot({
|
|
|
|
|
...normalized,
|
|
|
|
|
projectPlacement: session?.agents.get(agentId)?.projectPlacement ?? null,
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Call pure reducer
|
|
|
|
|
const result = processTimelineResponse({
|
|
|
|
|
payload,
|
|
|
|
|
@@ -810,6 +916,7 @@ function SessionProviderInternal({
|
|
|
|
|
}
|
|
|
|
|
},
|
|
|
|
|
[
|
|
|
|
|
applyAuthoritativeAgentSnapshot,
|
|
|
|
|
applyAgentUpdatePayload,
|
|
|
|
|
clearAgentStreamHead,
|
|
|
|
|
markAgentHistorySynchronized,
|
|
|
|
|
@@ -828,6 +935,22 @@ function SessionProviderInternal({
|
|
|
|
|
clearPendingAgentUpdates(serverId);
|
|
|
|
|
}, [isConnected, serverId]);
|
|
|
|
|
|
|
|
|
|
useEffect(() => {
|
|
|
|
|
const wasConnected = wasConnectedRef.current;
|
|
|
|
|
wasConnectedRef.current = isConnected;
|
|
|
|
|
if (!wasConnected && isConnected) {
|
|
|
|
|
scheduleAuthoritativeRevalidation();
|
|
|
|
|
}
|
|
|
|
|
}, [isConnected, scheduleAuthoritativeRevalidation]);
|
|
|
|
|
|
|
|
|
|
useEffect(() => {
|
|
|
|
|
return () => {
|
|
|
|
|
if (revalidationTimerRef.current) {
|
|
|
|
|
clearTimeout(revalidationTimerRef.current);
|
|
|
|
|
}
|
|
|
|
|
};
|
|
|
|
|
}, []);
|
|
|
|
|
|
|
|
|
|
// Daemon message handlers - directly update Zustand store
|
|
|
|
|
useEffect(() => {
|
|
|
|
|
const unsubAgentUpdate = client.on("agent_update", (message) => {
|
|
|
|
|
|