191 lines
6.7 KiB
TypeScript
191 lines
6.7 KiB
TypeScript
export const useUserEvents = () => {
|
|
const eventSource = ref<EventSource | null>(null);
|
|
const isConnected = useState('userEvents:connected', () => false);
|
|
const lastEventTimestamp = useState<number>('userEvents:lastTimestamp', () => 0);
|
|
|
|
const reconnectAttempts = ref(0);
|
|
const connectionStartTime = ref<number>(0);
|
|
const pendingReconnect = ref<NodeJS.Timeout | null>(null);
|
|
|
|
const BASE_DELAY = 1000;
|
|
const MAX_DELAY = 30000;
|
|
const MAX_RECONNECT_ATTEMPTS = 10;
|
|
|
|
const cleanup = () => {
|
|
if (pendingReconnect.value) {
|
|
clearTimeout(pendingReconnect.value);
|
|
}
|
|
if (eventSource.value) {
|
|
eventSource.value.close();
|
|
eventSource.value = null;
|
|
}
|
|
isConnected.value = false;
|
|
};
|
|
|
|
const getReconnectDelay = (attempt: number): number => {
|
|
const exponentialDelay = Math.min(
|
|
BASE_DELAY * Math.pow(2, attempt),
|
|
MAX_DELAY
|
|
);
|
|
const jitter = exponentialDelay * Math.random() * 0.25;
|
|
return exponentialDelay + jitter;
|
|
};
|
|
|
|
const scheduleReconnect = () => {
|
|
if (reconnectAttempts.value >= MAX_RECONNECT_ATTEMPTS) {
|
|
console.error('Max reconnection attempts reached');
|
|
return;
|
|
}
|
|
|
|
const delay = getReconnectDelay(reconnectAttempts.value);
|
|
console.log(`Scheduling reconnect in ${delay}ms (attempt ${reconnectAttempts.value + 1})`);
|
|
|
|
pendingReconnect.value = setTimeout(() => {
|
|
reconnectAttempts.value++;
|
|
connect();
|
|
}, delay);
|
|
};
|
|
|
|
const connect = async () => {
|
|
cleanup();
|
|
|
|
connectionStartTime.value = Date.now();
|
|
reconnectAttempts.value = 0;
|
|
|
|
const source = new EventSource('/api/events');
|
|
eventSource.value = source;
|
|
|
|
source.onopen = () => {
|
|
console.log('User events connected');
|
|
isConnected.value = true;
|
|
reconnectAttempts.value = 0;
|
|
};
|
|
|
|
source.onerror = () => {
|
|
isConnected.value = false;
|
|
|
|
if (source.readyState === EventSource.CLOSED) {
|
|
console.log('EventSource closed, scheduling reconnect');
|
|
scheduleReconnect();
|
|
}
|
|
};
|
|
|
|
source.addEventListener('message', async (event) => {
|
|
if (event.data === '') return;
|
|
|
|
try {
|
|
const data = JSON.parse(event.data);
|
|
|
|
if (data.timestamp && data.timestamp > lastEventTimestamp.value) {
|
|
lastEventTimestamp.value = data.timestamp;
|
|
}
|
|
|
|
handleUserEvent(data);
|
|
} catch (e) {
|
|
console.error('Failed to parse user event:', e);
|
|
}
|
|
});
|
|
};
|
|
|
|
const handleUserEvent = async (data: { entity: string; op: string; payload: any; timestamp?: number }) => {
|
|
const { agents } = await useAgents();
|
|
|
|
switch (data.entity) {
|
|
case 'topics': {
|
|
switch (data.op) {
|
|
case 'create': {
|
|
const existing = agents.value.find(a => a.id === data.payload.agentId);
|
|
if (existing) {
|
|
const topicExists = existing.topics.some(t => t.id === data.payload.id);
|
|
if (!topicExists) {
|
|
agents.value = agents.value.map(agent => {
|
|
return agent.id === data.payload.agentId
|
|
? { ...agent, topics: [data.payload, ...agent.topics] }
|
|
: agent;
|
|
});
|
|
}
|
|
}
|
|
break;
|
|
}
|
|
case 'update': {
|
|
const topics = agents.value.flatMap(agent => agent.topics);
|
|
const topic = topics.find(t => t.id === data.payload.topicId);
|
|
if (topic) {
|
|
agents.value = agents.value.map(agent => {
|
|
return agent.id === topic.agentId
|
|
? { ...agent, topics: agent.topics.map(t => t.id === topic.id ? { ...t, ...data.payload } : t) }
|
|
: agent;
|
|
});
|
|
}
|
|
break;
|
|
}
|
|
case 'delete': {
|
|
const topics = agents.value.flatMap(agent => agent.topics);
|
|
const topic = topics.find(t => t.id === data.payload.topicId);
|
|
if (topic) {
|
|
agents.value = agents.value.map(agent => {
|
|
return agent.id === topic.agentId
|
|
? { ...agent, topics: agent.topics.filter(t => t.id !== topic.id) }
|
|
: agent;
|
|
});
|
|
}
|
|
break;
|
|
}
|
|
}
|
|
break;
|
|
}
|
|
case 'agents': {
|
|
switch (data.op) {
|
|
case 'create': {
|
|
const existing = agents.value.find(a => a.id === data.payload.id);
|
|
if (existing) {
|
|
agents.value = agents.value.map(agent =>
|
|
agent.id === data.payload.id ? { ...agent, ...data.payload } : agent
|
|
);
|
|
} else {
|
|
agents.value = [...agents.value, data.payload];
|
|
}
|
|
break;
|
|
}
|
|
case 'update': {
|
|
agents.value = agents.value.map(agent =>
|
|
agent.id === data.payload.id ? { ...agent, ...data.payload } : agent
|
|
);
|
|
break;
|
|
}
|
|
case 'delete': {
|
|
agents.value = agents.value.filter(agent => agent.id !== data.payload.id);
|
|
break;
|
|
}
|
|
}
|
|
break;
|
|
}
|
|
}
|
|
};
|
|
|
|
const reconcile = async () => {
|
|
console.log('Reconciling user state...');
|
|
const { data } = await useFetch<AgentWithTopics[]>('/api/agents');
|
|
if (data.value) {
|
|
const { agents } = await useAgents();
|
|
agents.value = data.value;
|
|
}
|
|
};
|
|
|
|
onMounted(() => {
|
|
connect();
|
|
});
|
|
|
|
onBeforeUnmount(() => {
|
|
cleanup();
|
|
});
|
|
|
|
return {
|
|
eventSource,
|
|
isConnected,
|
|
connect,
|
|
cleanup,
|
|
reconcile,
|
|
};
|
|
};
|