fix: Resolve event delivery configuration header passing between frontend and backend (#7514)
Co-authored-by: Gabriel Luiz Freitas Almeida <gabriel@langflow.org> Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
This commit is contained in:
parent
467ae9e34f
commit
59a7440045
12 changed files with 54 additions and 46 deletions
|
|
@ -13,6 +13,7 @@ from sqlmodel import select
|
||||||
from langflow.api.disconnect import DisconnectHandlerStreamingResponse
|
from langflow.api.disconnect import DisconnectHandlerStreamingResponse
|
||||||
from langflow.api.utils import (
|
from langflow.api.utils import (
|
||||||
CurrentActiveUser,
|
CurrentActiveUser,
|
||||||
|
EventDeliveryType,
|
||||||
build_graph_from_data,
|
build_graph_from_data,
|
||||||
build_graph_from_db,
|
build_graph_from_db,
|
||||||
format_elapsed_time,
|
format_elapsed_time,
|
||||||
|
|
@ -84,12 +85,12 @@ async def get_flow_events_response(
|
||||||
*,
|
*,
|
||||||
job_id: str,
|
job_id: str,
|
||||||
queue_service: JobQueueService,
|
queue_service: JobQueueService,
|
||||||
stream: bool = True,
|
event_delivery: EventDeliveryType,
|
||||||
):
|
):
|
||||||
"""Get events for a specific build job, either as a stream or single event."""
|
"""Get events for a specific build job, either as a stream or single event."""
|
||||||
try:
|
try:
|
||||||
main_queue, event_manager, event_task, _ = queue_service.get_queue_data(job_id)
|
main_queue, event_manager, event_task, _ = queue_service.get_queue_data(job_id)
|
||||||
if stream:
|
if event_delivery in (EventDeliveryType.STREAMING, EventDeliveryType.DIRECT):
|
||||||
if event_task is None:
|
if event_task is None:
|
||||||
logger.error(f"No event task found for job {job_id}")
|
logger.error(f"No event task found for job {job_id}")
|
||||||
raise HTTPException(status_code=404, detail="No event task found for job")
|
raise HTTPException(status_code=404, detail="No event task found for job")
|
||||||
|
|
|
||||||
|
|
@ -2,6 +2,7 @@ from __future__ import annotations
|
||||||
|
|
||||||
import uuid
|
import uuid
|
||||||
from datetime import timedelta
|
from datetime import timedelta
|
||||||
|
from enum import Enum
|
||||||
from typing import TYPE_CHECKING, Annotated, Any
|
from typing import TYPE_CHECKING, Annotated, Any
|
||||||
|
|
||||||
from fastapi import Depends, HTTPException, Query
|
from fastapi import Depends, HTTPException, Query
|
||||||
|
|
@ -34,6 +35,12 @@ CurrentActiveUser = Annotated[User, Depends(get_current_active_user)]
|
||||||
DbSession = Annotated[AsyncSession, Depends(get_session)]
|
DbSession = Annotated[AsyncSession, Depends(get_session)]
|
||||||
|
|
||||||
|
|
||||||
|
class EventDeliveryType(str, Enum):
|
||||||
|
STREAMING = "streaming"
|
||||||
|
DIRECT = "direct"
|
||||||
|
POLLING = "polling"
|
||||||
|
|
||||||
|
|
||||||
def has_api_terms(word: str):
|
def has_api_terms(word: str):
|
||||||
return "api" in word and ("key" in word or ("token" in word and "tokens" not in word))
|
return "api" in word and ("key" in word or ("token" in word and "tokens" not in word))
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -27,6 +27,7 @@ from langflow.api.limited_background_tasks import LimitVertexBuildBackgroundTask
|
||||||
from langflow.api.utils import (
|
from langflow.api.utils import (
|
||||||
CurrentActiveUser,
|
CurrentActiveUser,
|
||||||
DbSession,
|
DbSession,
|
||||||
|
EventDeliveryType,
|
||||||
build_and_cache_graph_from_data,
|
build_and_cache_graph_from_data,
|
||||||
build_graph_from_db,
|
build_graph_from_db,
|
||||||
format_elapsed_time,
|
format_elapsed_time,
|
||||||
|
|
@ -55,12 +56,10 @@ from langflow.services.deps import (
|
||||||
get_chat_service,
|
get_chat_service,
|
||||||
get_queue_service,
|
get_queue_service,
|
||||||
get_session,
|
get_session,
|
||||||
get_settings_service,
|
|
||||||
get_telemetry_service,
|
get_telemetry_service,
|
||||||
session_scope,
|
session_scope,
|
||||||
)
|
)
|
||||||
from langflow.services.job_queue.service import JobQueueNotFoundError, JobQueueService
|
from langflow.services.job_queue.service import JobQueueNotFoundError, JobQueueService
|
||||||
from langflow.services.settings.service import SettingsService
|
|
||||||
from langflow.services.telemetry.schema import ComponentPayload, PlaygroundPayload
|
from langflow.services.telemetry.schema import ComponentPayload, PlaygroundPayload
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
|
|
@ -154,7 +153,7 @@ async def build_flow(
|
||||||
current_user: CurrentActiveUser,
|
current_user: CurrentActiveUser,
|
||||||
queue_service: Annotated[JobQueueService, Depends(get_queue_service)],
|
queue_service: Annotated[JobQueueService, Depends(get_queue_service)],
|
||||||
flow_name: str | None = None,
|
flow_name: str | None = None,
|
||||||
settings_service: Annotated[SettingsService, Depends(get_settings_service)],
|
event_delivery: EventDeliveryType = EventDeliveryType.POLLING,
|
||||||
):
|
):
|
||||||
"""Build and process a flow, returning a job ID for event polling.
|
"""Build and process a flow, returning a job ID for event polling.
|
||||||
|
|
||||||
|
|
@ -174,6 +173,7 @@ async def build_flow(
|
||||||
queue_service: Queue service for job management
|
queue_service: Queue service for job management
|
||||||
flow_name: Optional name for the flow
|
flow_name: Optional name for the flow
|
||||||
settings_service: Settings service
|
settings_service: Settings service
|
||||||
|
event_delivery: Optional event delivery type - default is streaming
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
Dict with job_id that can be used to poll for build status
|
Dict with job_id that can be used to poll for build status
|
||||||
|
|
@ -197,12 +197,14 @@ async def build_flow(
|
||||||
queue_service=queue_service,
|
queue_service=queue_service,
|
||||||
flow_name=flow_name,
|
flow_name=flow_name,
|
||||||
)
|
)
|
||||||
if settings_service.settings.event_delivery != "direct":
|
|
||||||
|
# This is required to support FE tests - we need to be able to set the event delivery to direct
|
||||||
|
if event_delivery == EventDeliveryType.DIRECT:
|
||||||
return {"job_id": job_id}
|
return {"job_id": job_id}
|
||||||
return await get_flow_events_response(
|
return await get_flow_events_response(
|
||||||
job_id=job_id,
|
job_id=job_id,
|
||||||
queue_service=queue_service,
|
queue_service=queue_service,
|
||||||
stream=True,
|
event_delivery=event_delivery,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -211,13 +213,13 @@ async def get_build_events(
|
||||||
job_id: str,
|
job_id: str,
|
||||||
queue_service: Annotated[JobQueueService, Depends(get_queue_service)],
|
queue_service: Annotated[JobQueueService, Depends(get_queue_service)],
|
||||||
*,
|
*,
|
||||||
stream: bool = True,
|
event_delivery: EventDeliveryType = EventDeliveryType.STREAMING,
|
||||||
):
|
):
|
||||||
"""Get events for a specific build job."""
|
"""Get events for a specific build job."""
|
||||||
return await get_flow_events_response(
|
return await get_flow_events_response(
|
||||||
job_id=job_id,
|
job_id=job_id,
|
||||||
queue_service=queue_service,
|
queue_service=queue_service,
|
||||||
stream=stream,
|
event_delivery=event_delivery,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -87,7 +87,7 @@ export default function NodeStatus({
|
||||||
const isBuilding = useFlowStore((state) => state.isBuilding);
|
const isBuilding = useFlowStore((state) => state.isBuilding);
|
||||||
const setNode = useFlowStore((state) => state.setNode);
|
const setNode = useFlowStore((state) => state.setNode);
|
||||||
const version = useDarkStore((state) => state.version);
|
const version = useDarkStore((state) => state.version);
|
||||||
const config = useGetConfig();
|
const eventDeliveryConfig = useUtilityStore((state) => state.eventDelivery);
|
||||||
const setErrorData = useAlertStore((state) => state.setErrorData);
|
const setErrorData = useAlertStore((state) => state.setErrorData);
|
||||||
|
|
||||||
const postTemplateValue = usePostTemplateValue({
|
const postTemplateValue = usePostTemplateValue({
|
||||||
|
|
@ -96,10 +96,6 @@ export default function NodeStatus({
|
||||||
node: data.node,
|
node: data.node,
|
||||||
});
|
});
|
||||||
|
|
||||||
const shouldStreamEvents = () => {
|
|
||||||
return config.data?.event_delivery === EventDeliveryType.STREAMING;
|
|
||||||
};
|
|
||||||
|
|
||||||
// Start polling when connection is initiated
|
// Start polling when connection is initiated
|
||||||
const startPolling = () => {
|
const startPolling = () => {
|
||||||
window.open(connectionLink, "_blank");
|
window.open(connectionLink, "_blank");
|
||||||
|
|
@ -169,8 +165,7 @@ export default function NodeStatus({
|
||||||
setValidationStatus(null);
|
setValidationStatus(null);
|
||||||
buildFlow({
|
buildFlow({
|
||||||
stopNodeId: nodeId,
|
stopNodeId: nodeId,
|
||||||
stream: shouldStreamEvents(),
|
eventDelivery: eventDeliveryConfig,
|
||||||
eventDelivery: config.data?.event_delivery,
|
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -265,8 +260,7 @@ export default function NodeStatus({
|
||||||
if (buildStatus === BuildStatus.BUILDING || isBuilding) return;
|
if (buildStatus === BuildStatus.BUILDING || isBuilding) return;
|
||||||
buildFlow({
|
buildFlow({
|
||||||
stopNodeId: nodeId,
|
stopNodeId: nodeId,
|
||||||
stream: shouldStreamEvents(),
|
eventDelivery: eventDeliveryConfig,
|
||||||
eventDelivery: config.data?.event_delivery,
|
|
||||||
});
|
});
|
||||||
track("Flow Build - Clicked", { stopNodeId: nodeId });
|
track("Flow Build - Clicked", { stopNodeId: nodeId });
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,7 @@ import axios, { AxiosError, AxiosInstance, AxiosRequestConfig } from "axios";
|
||||||
import * as fetchIntercept from "fetch-intercept";
|
import * as fetchIntercept from "fetch-intercept";
|
||||||
import { useEffect } from "react";
|
import { useEffect } from "react";
|
||||||
import { Cookies } from "react-cookie";
|
import { Cookies } from "react-cookie";
|
||||||
import { BuildStatus } from "../../constants/enums";
|
import { BuildStatus, EventDeliveryType } from "../../constants/enums";
|
||||||
import useAlertStore from "../../stores/alertStore";
|
import useAlertStore from "../../stores/alertStore";
|
||||||
import useFlowStore from "../../stores/flowStore";
|
import useFlowStore from "../../stores/flowStore";
|
||||||
import { checkDuplicateRequestAndStoreRequest } from "./helpers/check-duplicate-requests";
|
import { checkDuplicateRequestAndStoreRequest } from "./helpers/check-duplicate-requests";
|
||||||
|
|
@ -263,6 +263,7 @@ export type StreamingRequestParams = {
|
||||||
onError?: (statusCode: number) => void;
|
onError?: (statusCode: number) => void;
|
||||||
onNetworkError?: (error: Error) => void;
|
onNetworkError?: (error: Error) => void;
|
||||||
buildController: AbortController;
|
buildController: AbortController;
|
||||||
|
eventDeliveryConfig?: EventDeliveryType;
|
||||||
};
|
};
|
||||||
|
|
||||||
async function performStreamingRequest({
|
async function performStreamingRequest({
|
||||||
|
|
|
||||||
|
|
@ -39,7 +39,7 @@ export const useGetConfig: useQueryFunctionType<undefined, ConfigResponse> = (
|
||||||
const setWebhookPollingInterval = useUtilityStore(
|
const setWebhookPollingInterval = useUtilityStore(
|
||||||
(state) => state.setWebhookPollingInterval,
|
(state) => state.setWebhookPollingInterval,
|
||||||
);
|
);
|
||||||
|
const setEventDelivery = useUtilityStore((state) => state.setEventDelivery);
|
||||||
const { query } = UseRequestProcessor();
|
const { query } = UseRequestProcessor();
|
||||||
|
|
||||||
const getConfigFn = async () => {
|
const getConfigFn = async () => {
|
||||||
|
|
@ -59,6 +59,7 @@ export const useGetConfig: useQueryFunctionType<undefined, ConfigResponse> = (
|
||||||
setWebhookPollingInterval(
|
setWebhookPollingInterval(
|
||||||
data.webhook_polling_interval ?? DEFAULT_POLLING_INTERVAL,
|
data.webhook_polling_interval ?? DEFAULT_POLLING_INTERVAL,
|
||||||
);
|
);
|
||||||
|
setEventDelivery(data.event_delivery ?? EventDeliveryType.POLLING);
|
||||||
}
|
}
|
||||||
return data;
|
return data;
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -155,11 +155,7 @@ export default function IOModal({
|
||||||
|
|
||||||
const chatValue = useUtilityStore((state) => state.chatValueStore);
|
const chatValue = useUtilityStore((state) => state.chatValueStore);
|
||||||
const setChatValue = useUtilityStore((state) => state.setChatValueStore);
|
const setChatValue = useUtilityStore((state) => state.setChatValueStore);
|
||||||
const config = useGetConfig();
|
const eventDeliveryConfig = useUtilityStore((state) => state.eventDelivery);
|
||||||
|
|
||||||
function shouldStreamEvents() {
|
|
||||||
return config.data?.event_delivery === EventDeliveryType.STREAMING;
|
|
||||||
}
|
|
||||||
|
|
||||||
const sendMessage = useCallback(
|
const sendMessage = useCallback(
|
||||||
async ({
|
async ({
|
||||||
|
|
@ -178,8 +174,7 @@ export default function IOModal({
|
||||||
files: files,
|
files: files,
|
||||||
silent: true,
|
silent: true,
|
||||||
session: sessionId,
|
session: sessionId,
|
||||||
stream: shouldStreamEvents(),
|
eventDelivery: eventDeliveryConfig,
|
||||||
eventDelivery: config.data?.event_delivery,
|
|
||||||
}).catch((err) => {
|
}).catch((err) => {
|
||||||
console.error(err);
|
console.error(err);
|
||||||
});
|
});
|
||||||
|
|
|
||||||
|
|
@ -835,7 +835,6 @@ const useFlowStore = create<FlowStoreType>((set, get) => ({
|
||||||
edges: get().edges || undefined,
|
edges: get().edges || undefined,
|
||||||
logBuilds: get().onFlowPage,
|
logBuilds: get().onFlowPage,
|
||||||
playgroundPage,
|
playgroundPage,
|
||||||
stream,
|
|
||||||
eventDelivery,
|
eventDelivery,
|
||||||
});
|
});
|
||||||
get().setIsBuilding(false);
|
get().setIsBuilding(false);
|
||||||
|
|
|
||||||
|
|
@ -1,3 +1,4 @@
|
||||||
|
import { EventDeliveryType } from "@/constants/enums";
|
||||||
import { Pagination, Tag } from "@/types/utils/types";
|
import { Pagination, Tag } from "@/types/utils/types";
|
||||||
import { UtilityStoreType } from "@/types/zustand/utility";
|
import { UtilityStoreType } from "@/types/zustand/utility";
|
||||||
import { create } from "zustand";
|
import { create } from "zustand";
|
||||||
|
|
@ -43,4 +44,7 @@ export const useUtilityStore = create<UtilityStoreType>((set, get) => ({
|
||||||
currentSessionId: "",
|
currentSessionId: "",
|
||||||
setCurrentSessionId: (sessionId: string) =>
|
setCurrentSessionId: (sessionId: string) =>
|
||||||
set({ currentSessionId: sessionId }),
|
set({ currentSessionId: sessionId }),
|
||||||
|
eventDelivery: EventDeliveryType.POLLING,
|
||||||
|
setEventDelivery: (eventDelivery: EventDeliveryType) =>
|
||||||
|
set({ eventDelivery }),
|
||||||
}));
|
}));
|
||||||
|
|
|
||||||
|
|
@ -149,7 +149,6 @@ export type FlowStoreType = {
|
||||||
files,
|
files,
|
||||||
silent,
|
silent,
|
||||||
session,
|
session,
|
||||||
stream,
|
|
||||||
eventDelivery,
|
eventDelivery,
|
||||||
}: {
|
}: {
|
||||||
startNodeId?: string;
|
startNodeId?: string;
|
||||||
|
|
@ -158,7 +157,6 @@ export type FlowStoreType = {
|
||||||
files?: string[];
|
files?: string[];
|
||||||
silent?: boolean;
|
silent?: boolean;
|
||||||
session?: string;
|
session?: string;
|
||||||
stream?: boolean;
|
|
||||||
eventDelivery?: EventDeliveryType;
|
eventDelivery?: EventDeliveryType;
|
||||||
}) => Promise<void>;
|
}) => Promise<void>;
|
||||||
getFlow: () => { nodes: Node[]; edges: EdgeType[]; viewport: Viewport };
|
getFlow: () => { nodes: Node[]; edges: EdgeType[]; viewport: Viewport };
|
||||||
|
|
|
||||||
|
|
@ -1,3 +1,4 @@
|
||||||
|
import { EventDeliveryType } from "@/constants/enums";
|
||||||
import { Pagination, Tag } from "@/types/utils/types";
|
import { Pagination, Tag } from "@/types/utils/types";
|
||||||
|
|
||||||
export type UtilityStoreType = {
|
export type UtilityStoreType = {
|
||||||
|
|
@ -25,4 +26,6 @@ export type UtilityStoreType = {
|
||||||
setCurrentSessionId: (sessionId: string) => void;
|
setCurrentSessionId: (sessionId: string) => void;
|
||||||
setClientId: (clientId: string) => void;
|
setClientId: (clientId: string) => void;
|
||||||
clientId: string;
|
clientId: string;
|
||||||
|
eventDelivery: EventDeliveryType;
|
||||||
|
setEventDelivery: (eventDelivery: EventDeliveryType) => void;
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -40,8 +40,7 @@ type BuildVerticesParams = {
|
||||||
logBuilds?: boolean;
|
logBuilds?: boolean;
|
||||||
session?: string;
|
session?: string;
|
||||||
playgroundPage?: boolean;
|
playgroundPage?: boolean;
|
||||||
stream?: boolean;
|
eventDelivery: EventDeliveryType;
|
||||||
eventDelivery?: EventDeliveryType;
|
|
||||||
};
|
};
|
||||||
|
|
||||||
function getInactiveVertexData(vertexId: string): VertexBuildTypeAPI {
|
function getInactiveVertexData(vertexId: string): VertexBuildTypeAPI {
|
||||||
|
|
@ -148,7 +147,10 @@ export async function buildFlowVerticesWithFallback(
|
||||||
e.message === POLLING_MESSAGES.STREAMING_NOT_SUPPORTED
|
e.message === POLLING_MESSAGES.STREAMING_NOT_SUPPORTED
|
||||||
) {
|
) {
|
||||||
// Fallback to polling
|
// Fallback to polling
|
||||||
return await buildFlowVertices({ ...params, stream: false });
|
return await buildFlowVertices({
|
||||||
|
...params,
|
||||||
|
eventDelivery: EventDeliveryType.POLLING,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
throw e;
|
throw e;
|
||||||
}
|
}
|
||||||
|
|
@ -176,13 +178,16 @@ async function pollBuildEvents(
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
let isDone = false;
|
let isDone = false;
|
||||||
while (!isDone) {
|
while (!isDone) {
|
||||||
const response = await fetch(`${url}?stream=false`, {
|
const response = await fetch(
|
||||||
method: "GET",
|
`${url}?event_delivery=${EventDeliveryType.POLLING}`,
|
||||||
headers: {
|
{
|
||||||
"Content-Type": "application/json",
|
method: "GET",
|
||||||
|
headers: {
|
||||||
|
"Content-Type": "application/json",
|
||||||
|
},
|
||||||
|
signal: abortController.signal, // Add abort signal to fetch
|
||||||
},
|
},
|
||||||
signal: abortController.signal, // Add abort signal to fetch
|
);
|
||||||
});
|
|
||||||
|
|
||||||
if (!response.ok) {
|
if (!response.ok) {
|
||||||
const errorData = await response.json().catch(() => ({}));
|
const errorData = await response.json().catch(() => ({}));
|
||||||
|
|
@ -241,7 +246,6 @@ export async function buildFlowVertices({
|
||||||
logBuilds,
|
logBuilds,
|
||||||
session,
|
session,
|
||||||
playgroundPage,
|
playgroundPage,
|
||||||
stream = true,
|
|
||||||
eventDelivery,
|
eventDelivery,
|
||||||
}: BuildVerticesParams) {
|
}: BuildVerticesParams) {
|
||||||
const inputs = {};
|
const inputs = {};
|
||||||
|
|
@ -260,10 +264,10 @@ export async function buildFlowVertices({
|
||||||
queryParams.append("log_builds", logBuilds.toString());
|
queryParams.append("log_builds", logBuilds.toString());
|
||||||
}
|
}
|
||||||
|
|
||||||
// Add stream parameter when using direct event delivery
|
queryParams.append(
|
||||||
if (eventDelivery === EventDeliveryType.DIRECT) {
|
"event_delivery",
|
||||||
queryParams.append("stream", "true");
|
eventDelivery ?? EventDeliveryType.POLLING,
|
||||||
}
|
);
|
||||||
|
|
||||||
if (queryParams.toString()) {
|
if (queryParams.toString()) {
|
||||||
buildUrl = `${buildUrl}?${queryParams.toString()}`;
|
buildUrl = `${buildUrl}?${queryParams.toString()}`;
|
||||||
|
|
@ -376,13 +380,12 @@ export async function buildFlowVertices({
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
useFlowStore.getState().setBuildController(buildController);
|
useFlowStore.getState().setBuildController(buildController);
|
||||||
|
|
||||||
// Then stream the events
|
// Then stream the events
|
||||||
const eventsUrl = `${BASE_URL_API}build/${job_id}/events`;
|
const eventsUrl = `${BASE_URL_API}build/${job_id}/events`;
|
||||||
const buildResults: Array<boolean> = [];
|
const buildResults: Array<boolean> = [];
|
||||||
const verticesStartTimeMs: Map<string, number> = new Map();
|
const verticesStartTimeMs: Map<string, number> = new Map();
|
||||||
|
|
||||||
if (stream) {
|
if (eventDelivery === EventDeliveryType.STREAMING) {
|
||||||
return performStreamingRequest({
|
return performStreamingRequest({
|
||||||
method: "GET",
|
method: "GET",
|
||||||
url: eventsUrl,
|
url: eventsUrl,
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue