fix(frontend): add useRef hook to manage eventSource in ChatMessage component

feat(frontend): add support for process.env.PORT environment variable in server.ts
feat(frontend): add updateFlowPool function to NewChatView component
feat(frontend): add buildId parameter to addDataToFlowPool function in flowStore
feat(frontend): add stream_url property to ChatOutputType in flow types
This commit is contained in:
anovazzi1 2024-02-27 18:29:57 -03:00
commit df07cf413b
7 changed files with 64 additions and 60 deletions

View file

@ -1,5 +1,5 @@
import Convert from "ansi-to-html"; import Convert from "ansi-to-html";
import { useEffect, useMemo, useState } from "react"; import { useEffect, useMemo, useState,useRef } from "react";
import Markdown from "react-markdown"; import Markdown from "react-markdown";
import rehypeMathjax from "rehype-mathjax"; import rehypeMathjax from "rehype-mathjax";
import remarkGfm from "remark-gfm"; import remarkGfm from "remark-gfm";
@ -12,6 +12,7 @@ import IconComponent from "../../../components/genericIconComponent";
import { chatMessagePropsType } from "../../../types/components"; import { chatMessagePropsType } from "../../../types/components";
import { classNames } from "../../../utils/utils"; import { classNames } from "../../../utils/utils";
import FileCard from "../fileComponent"; import FileCard from "../fileComponent";
import useFlowStore from "../../../stores/flowStore";
export default function ChatMessage({ export default function ChatMessage({
chat, chat,
@ -29,6 +30,9 @@ export default function ChatMessage({
const chatMessageString = chat.message ? chat.message.toString() : ""; const chatMessageString = chat.message ? chat.message.toString() : "";
const [chatMessage, setChatMessage] = useState(chatMessageString); const [chatMessage, setChatMessage] = useState(chatMessageString);
const [isStreaming, setIsStreaming] = useState(false); const [isStreaming, setIsStreaming] = useState(false);
const eventSource = useRef<EventSource | undefined>(undefined);
const updateFlowPool = useFlowStore((state) => state.updateFlowPool);
// The idea now is that chat.stream_url MAY be a URL if we should stream the output of the chat // The idea now is that chat.stream_url MAY be a URL if we should stream the output of the chat
// probably the message is empty when we have a stream_url // probably the message is empty when we have a stream_url
@ -36,49 +40,48 @@ export default function ChatMessage({
const streamChunks = (url: string) => { const streamChunks = (url: string) => {
setIsStreaming(true); // Streaming starts setIsStreaming(true); // Streaming starts
return new Promise<boolean>((resolve, reject) => { return new Promise<boolean>((resolve, reject) => {
const eventSource = new EventSource(url); eventSource.current = new EventSource(url);
eventSource.onmessage = (event) => { eventSource.current.onmessage = (event) => {
let parsedData = JSON.parse(event.data); let parsedData = JSON.parse(event.data);
if (parsedData.chunk) { if (parsedData.chunk) {
setChatMessage((prev) => prev + parsedData.chunk); setChatMessage((prev) => prev + parsedData.chunk);
} }
}; };
eventSource.onerror = (event) => { eventSource.current.onerror = (event) => {
setIsStreaming(false);
eventSource.current?.close();
setStreamUrl(undefined);
reject(new Error("Streaming failed")); reject(new Error("Streaming failed"));
setIsStreaming(false);
eventSource.close();
}; };
eventSource.addEventListener("close", (event) => { eventSource.current.addEventListener("close", (event) => {
setStreamUrl(null); // Update state to reflect the stream is closed setStreamUrl(undefined); // Update state to reflect the stream is closed
resolve(true); eventSource.current?.close();
setIsStreaming(false); setIsStreaming(false);
eventSource.close(); resolve(true);
}); });
}); });
}; };
useEffect(() => { useEffect(() => {
if (streamUrl && chat.message === "") { console.log(streamUrl)
if (streamUrl&& !isStreaming) {
streamChunks(streamUrl) streamChunks(streamUrl)
.then(() => { .then(() => {
if (updateChat) { if (updateChat) {
updateChat(chat, chatMessage, streamUrl); console.log("rodou")
updateChat(chat, chatMessage);
} }
}) })
.catch((error) => { .catch((error) => {
console.error(error); console.error(error);
}); });
} }
}, [streamUrl]); return () => {
eventSource.current?.close();
useEffect(() => {
// This effect is specifically for calling updateChat after streaming ends
if (!isStreaming && streamUrl) {
if (updateChat) {
updateChat(chat, chatMessage, streamUrl);
}
} }
}, [isStreaming]); }, [streamUrl,chatMessage]);
useEffect(() => { useEffect(() => {
const element = document.getElementById("last-chat-message"); const element = document.getElementById("last-chat-message");

View file

@ -34,6 +34,7 @@ export default function NewChatView({
const inputIds = inputs.map((obj) => obj.id); const inputIds = inputs.map((obj) => obj.id);
const outputIds = outputs.map((obj) => obj.id); const outputIds = outputs.map((obj) => obj.id);
const outputTypes = outputs.map((obj) => obj.type); const outputTypes = outputs.map((obj) => obj.type);
const updateFlowPool = useFlowStore((state)=>state.updateFlowPool)
useEffect(() => { useEffect(() => {
if (!outputTypes.includes("ChatOutput")) { if (!outputTypes.includes("ChatOutput")) {
@ -67,14 +68,12 @@ export default function NewChatView({
const { sender, message, sender_name, stream_url } = output.data const { sender, message, sender_name, stream_url } = output.data
.artifacts as ChatOutputType; .artifacts as ChatOutputType;
const componentId = output.id + index;
const is_ai = sender === "Machine" || sender === null; const is_ai = sender === "Machine" || sender === null;
return { return {
isSend: !is_ai, isSend: !is_ai,
message: message, message: message,
sender_name, sender_name,
id: componentId, componentId: output.id,
stream_url: stream_url, stream_url: stream_url,
}; };
} catch (e) { } catch (e) {
@ -83,7 +82,7 @@ export default function NewChatView({
isSend: false, isSend: false,
message: "Error parsing message", message: "Error parsing message",
sender_name: "Error", sender_name: "Error",
id: output.id + index, componentId: output.id,
}; };
} }
}); });
@ -120,27 +119,25 @@ export default function NewChatView({
function updateChat( function updateChat(
chat: ChatMessageType, chat: ChatMessageType,
message: string, message: string,
stream_url: string | null stream_url?: string
) { ) {
if (message === "") return; if (message === "") return;
console.log(`updateChat: ${message}`); chat.message = message;
console.log("chatHistory:", chatHistory); console.log(message)
chat.message = message;
chat.stream_url = stream_url;
// chat is one of the chatHistory // chat is one of the chatHistory
setChatHistory((oldChatHistory) => { updateFlowPool(chat.componentId,{message,sender_name:chat.sender_name??"Bot",sender:"Machine"})
const index = oldChatHistory.findIndex((ch) => ch.id === chat.id); // setChatHistory((oldChatHistory) => {
// const index = oldChatHistory.findIndex((ch) => ch.id === chat.id);
if (index === -1) return oldChatHistory; // if (index === -1) return oldChatHistory;
let newChatHistory = _.cloneDeep(oldChatHistory); // let newChatHistory = _.cloneDeep(oldChatHistory);
newChatHistory = [ // newChatHistory = [
...newChatHistory.slice(0, index), // ...newChatHistory.slice(0, index),
chat, // chat,
...newChatHistory.slice(index + 1), // ...newChatHistory.slice(index + 1),
]; // ];
console.log("newChatHistory:", newChatHistory); // console.log("newChatHistory:", newChatHistory);
return newChatHistory; // return newChatHistory;
}); // });
} }
return ( return (
@ -167,7 +164,7 @@ export default function NewChatView({
lockChat={lockChat} lockChat={lockChat}
chat={chat} chat={chat}
lastMessage={chatHistory.length - 1 === index ? true : false} lastMessage={chatHistory.length - 1 === index ? true : false}
key={`${chat.id}-${index}`} key={`${chat.componentId}-${index}`}
updateChat={updateChat} updateChat={updateChat}
/> />
)) ))

View file

@ -51,7 +51,7 @@ const useFlowStore = create<FlowStoreType>((set, get) => ({
setFlowPool: (flowPool) => { setFlowPool: (flowPool) => {
set({ flowPool }); set({ flowPool });
}, },
addDataToFlowPool: (data: any, nodeId: string) => { addDataToFlowPool: (data: FlowPoolObjectType, nodeId: string) => {
let newFlowPool = cloneDeep({ ...get().flowPool }); let newFlowPool = cloneDeep({ ...get().flowPool });
if (!newFlowPool[nodeId]) newFlowPool[nodeId] = [data]; if (!newFlowPool[nodeId]) newFlowPool[nodeId] = [data];
else { else {
@ -416,12 +416,13 @@ const useFlowStore = create<FlowStoreType>((set, get) => ({
} }
function handleBuildUpdate( function handleBuildUpdate(
vertexBuildData: VertexBuildTypeAPI, vertexBuildData: VertexBuildTypeAPI,
status: BuildStatus status: BuildStatus,
buildId:string
) { ) {
if (vertexBuildData && vertexBuildData.inactive_vertices) { if (vertexBuildData && vertexBuildData.inactive_vertices) {
get().removeFromVerticesBuild(vertexBuildData.inactive_vertices); get().removeFromVerticesBuild(vertexBuildData.inactive_vertices);
} }
get().addDataToFlowPool(vertexBuildData, vertexBuildData.id); get().addDataToFlowPool({...vertexBuildData,buildId}, vertexBuildData.id);
useFlowStore.getState().updateBuildStatus([vertexBuildData.id], status); useFlowStore.getState().updateBuildStatus([vertexBuildData.id], status);
} }
await updateFlowInDatabase({ await updateFlowInDatabase({

View file

@ -9,7 +9,7 @@ export type ChatMessageType = {
files?: Array<{ data: string; type: string; data_type: string }>; files?: Array<{ data: string; type: string; data_type: string }>;
prompt?: string; prompt?: string;
chatKey?: string; chatKey?: string;
id?: string; componentId: string;
stream_url?: string | null; stream_url?: string | null;
sender_name?: string; sender_name?: string;
}; };

View file

@ -527,7 +527,7 @@ export type chatMessagePropsType = {
updateChat: ( updateChat: (
chat: ChatMessageType, chat: ChatMessageType,
message: string, message: string,
stream_url: string stream_url?: string
) => void; ) => void;
}; };
@ -632,9 +632,9 @@ export type validationStatusType = {
id: string; id: string;
data: object | any; data: object | any;
params: string; params: string;
progress: number; progress?: number;
valid: boolean; valid: boolean;
duration: string; duration?: string;
}; };
export type ApiKey = { export type ApiKey = {

View file

@ -18,6 +18,7 @@ export type ChatOutputType = {
message: string; message: string;
sender: string; sender: string;
sender_name: string; sender_name: string;
stream_url?: string;
}; };
export type FlowPoolObjectType = { export type FlowPoolObjectType = {
@ -25,9 +26,10 @@ export type FlowPoolObjectType = {
valid: boolean; valid: boolean;
params: any; params: any;
data: { artifacts: any | ChatOutputType | chatInputType; results: any | ChatOutputType | chatInputType }; data: { artifacts: any | ChatOutputType | chatInputType; results: any | ChatOutputType | chatInputType };
duration: string; duration?: string;
progress: number; progress?: number;
id: string; id: string;
buildId: string;
}; };
export type FlowPoolType = { export type FlowPoolType = {
@ -40,7 +42,7 @@ export type FlowStoreType = {
outputs: Array<{ type: string; id: string }>; outputs: Array<{ type: string; id: string }>;
hasIO: boolean; hasIO: boolean;
setFlowPool: (flowPool: FlowPoolType) => void; setFlowPool: (flowPool: FlowPoolType) => void;
addDataToFlowPool: (data: any, nodeId: string) => void; addDataToFlowPool: (data: FlowPoolObjectType, nodeId: string) => void;
CleanFlowPool: () => void; CleanFlowPool: () => void;
isBuilding: boolean; isBuilding: boolean;
isPending: boolean; isPending: boolean;

View file

@ -9,7 +9,7 @@ type BuildVerticesParams = {
flowId: string; // Assuming FlowType is the type for your flow flowId: string; // Assuming FlowType is the type for your flow
nodeId?: string | null; // Assuming nodeId is of type string, and it's optional nodeId?: string | null; // Assuming nodeId is of type string, and it's optional
onGetOrderSuccess?: () => void; onGetOrderSuccess?: () => void;
onBuildUpdate?: (data: VertexBuildTypeAPI, status: BuildStatus) => void; // Replace any with the actual type if it's not any onBuildUpdate?: (data: VertexBuildTypeAPI, status: BuildStatus,buildId:string) => void; // Replace any with the actual type if it's not any
onBuildComplete?: (allNodesValid: boolean) => void; onBuildComplete?: (allNodesValid: boolean) => void;
onBuildError?: (title, list, idList: string[]) => void; onBuildError?: (title, list, idList: string[]) => void;
onBuildStart?: (idList: string[]) => void; onBuildStart?: (idList: string[]) => void;
@ -48,7 +48,7 @@ export async function buildVertices({
let orderResponse; let orderResponse;
try { try {
orderResponse = await getVerticesOrder(flowId, nodeId); orderResponse = await getVerticesOrder(flowId, nodeId);
} catch (error) { } catch (error:any) {
console.log(error); console.log(error);
setErrorData({ setErrorData({
title: "Oops! Looks like you missed something", title: "Oops! Looks like you missed something",
@ -59,6 +59,7 @@ export async function buildVertices({
} }
if (onGetOrderSuccess) onGetOrderSuccess(); if (onGetOrderSuccess) onGetOrderSuccess();
let verticesOrder: Array<Array<string>> = orderResponse.data.ids; let verticesOrder: Array<Array<string>> = orderResponse.data.ids;
const runId = orderResponse.data.run_id;
let vertices_layers: Array<Array<string>> = []; let vertices_layers: Array<Array<string>> = [];
let stop = false; let stop = false;
if (validateNodes) { if (validateNodes) {
@ -102,14 +103,14 @@ export async function buildVertices({
onBuildUpdate onBuildUpdate
) { ) {
// If it is, skip building and set the state to inactive // If it is, skip building and set the state to inactive
onBuildUpdate(getInactiveVertexData(id), BuildStatus.INACTIVE); onBuildUpdate(getInactiveVertexData(id), BuildStatus.INACTIVE,runId);
buildResults.push(false); buildResults.push(false);
continue; continue;
} }
await buildVertex({ await buildVertex({
flowId, flowId,
id, id,
onBuildUpdate, onBuildUpdate:(data: VertexBuildTypeAPI, status: BuildStatus) => {if(onBuildUpdate) onBuildUpdate(data, status,runId)},
onBuildError, onBuildError,
verticesIds, verticesIds,
buildResults, buildResults,