88from uuid import UUID
99
1010import orjson
11- from fastapi import APIRouter , Depends , File , HTTPException , Request , UploadFile
11+ from fastapi import APIRouter , Depends , File , HTTPException , Request , UploadFile , status
1212from fastapi .encoders import jsonable_encoder
1313from fastapi_pagination import Page , Params
1414from fastapi_pagination .ext .sqlmodel import apaginate
15+ from lfx .log .logger import logger
1516from lfx .services .cache .utils import CACHE_MISS
1617from pydantic import ValidationError
1718from sqlmodel import and_ , col , select
5657from langflow .services .authorization .fetch import deny_to_404
5758from langflow .services .authorization .utils import _resolve_authz_domain
5859from langflow .services .cache .service import ThreadingInMemoryCache
60+ from langflow .services .database .lock_retry import (
61+ is_database_lock_error ,
62+ run_with_lock_retry ,
63+ )
5964from langflow .services .database .models .deployment .exceptions import (
6065 araise_if_deployment_guard_error_or_skip ,
6166)
7479# and FlowVersionError from the flow_version modules.
7580from langflow .services .database .models .folder .constants import DEFAULT_FOLDER_NAME
7681from langflow .services .database .models .folder .model import Folder
82+ from langflow .services .database .models .user .model import UserRead
7783from langflow .services .deps import get_settings_service , get_storage_service
7884from langflow .services .storage .service import StorageService
7985from langflow .utils .compression import compress_response
@@ -103,6 +109,9 @@ def _handle_unique_constraint_error(exc: Exception, *, status_code: int = 400) -
103109# build router
104110router = APIRouter (prefix = "/flows" , tags = ["Flows" ])
105111
112+ FLOW_UPDATE_FAILED = "Could not update the flow."
113+ FLOW_UPDATE_BUSY = "The database is busy. Please retry the request."
114+
106115
107116@router .post ("/" , response_model = FlowRead , status_code = 201 )
108117async def create_flow (
@@ -328,6 +337,7 @@ async def update_flow(
328337 storage_service : Annotated [StorageService , Depends (get_storage_service )],
329338):
330339 """Update a flow."""
340+ actor = UserRead .model_validate (current_user , from_attributes = True )
331341 try :
332342 # Destination check: if the payload moves the flow into a new
333343 # workspace/folder, the caller must also be authorized to write at the
@@ -339,7 +349,7 @@ async def update_flow(
339349 if target_workspace_id != db_flow .workspace_id or target_folder_id != db_flow .folder_id :
340350 try :
341351 await ensure_flow_permission (
342- current_user ,
352+ actor ,
343353 FlowAction .WRITE ,
344354 flow_id = flow_id ,
345355 flow_user_id = db_flow .user_id ,
@@ -357,7 +367,7 @@ async def update_flow(
357367
358368 async def operation () -> FlowRead :
359369 # Re-load inside each attempt so retry after nested rollback never uses an expired ORM instance.
360- db_flow_for_attempt = await _read_flow (session = session , flow_id = flow_id , user_id = current_user .id )
370+ db_flow_for_attempt = await _read_flow (session = session , flow_id = flow_id , user_id = actor .id )
361371 if not db_flow_for_attempt :
362372 raise HTTPException (status_code = 404 , detail = "Flow not found" )
363373 # TOCTOU: a concurrent PATCH could have moved this flow to a
@@ -367,7 +377,7 @@ async def operation() -> FlowRead:
367377 # stale check across a race.
368378 try :
369379 await ensure_flow_permission (
370- current_user ,
380+ actor ,
371381 FlowAction .WRITE ,
372382 flow_id = flow_id ,
373383 flow_user_id = db_flow_for_attempt .user_id ,
@@ -386,7 +396,7 @@ async def operation() -> FlowRead:
386396 ):
387397 try :
388398 await ensure_flow_permission (
389- current_user ,
399+ actor ,
390400 FlowAction .WRITE ,
391401 flow_id = flow_id ,
392402 flow_user_id = db_flow_for_attempt .user_id ,
@@ -399,26 +409,47 @@ async def operation() -> FlowRead:
399409 session = session ,
400410 db_flow = db_flow_for_attempt ,
401411 flow = flow ,
402- user_id = current_user .id ,
412+ user_id = actor .id ,
403413 storage_service = storage_service ,
404414 )
405415
406- if folder_id_will_change :
407- return await retry_flow_operation_on_deployment_guard (
408- db = session ,
409- user_id = current_user .id ,
410- flow_ids = [flow_id ],
411- operation = operation ,
412- )
413- return await operation ()
416+ async def update_attempt (_attempt : int ) -> FlowRead :
417+ if folder_id_will_change :
418+ return await retry_flow_operation_on_deployment_guard (
419+ db = session ,
420+ user_id = actor .id ,
421+ flow_ids = [flow_id ],
422+ operation = operation ,
423+ )
424+ return await operation ()
425+
426+ return await run_with_lock_retry (
427+ update_attempt ,
428+ session = session ,
429+ description = f"update_flow { flow_id } " ,
430+ )
414431 except HTTPException :
415432 raise
416433 except Exception as e :
417434 await araise_if_deployment_guard_error_or_skip (
418435 e ,
419436 log_message = f"op=update_flow flow_id={ flow_id } " ,
420437 )
421- raise _handle_unique_constraint_error (e ) from e
438+ if is_database_lock_error (e ):
439+ await logger .awarning ("op=update_flow flow_id=%s exhausted lock retries" , flow_id )
440+ raise HTTPException (
441+ status_code = status .HTTP_503_SERVICE_UNAVAILABLE ,
442+ detail = FLOW_UPDATE_BUSY ,
443+ headers = {"Retry-After" : "1" },
444+ ) from e
445+ handled_error = _handle_unique_constraint_error (e )
446+ if handled_error .status_code != status .HTTP_500_INTERNAL_SERVER_ERROR :
447+ raise handled_error from e
448+ await logger .aerror ("op=update_flow flow_id=%s failed with %s" , flow_id , type (e ).__name__ )
449+ raise HTTPException (
450+ status_code = status .HTTP_500_INTERNAL_SERVER_ERROR ,
451+ detail = FLOW_UPDATE_FAILED ,
452+ ) from e
422453
423454
424455@router .put ("/{flow_id}" , response_model = FlowRead )
0 commit comments