fix: add a lock check on starter project initialization (#8631)
* Adds a lock check on starter project initialization * [autofix.ci] apply automated fixes * ruff fix for exception * [autofix.ci] apply automated fixes * Adds a lock check on starter project initialization * [autofix.ci] apply automated fixes * fix: Improve logging for lock acquisition failures in starter project creation * [autofix.ci] apply automated fixes --------- Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com> Co-authored-by: Gabriel Luiz Freitas Almeida <gabriel@langflow.org>
This commit is contained in:
parent
c05327d773
commit
360c2395fb
4 changed files with 411 additions and 423 deletions
|
|
@ -124,7 +124,9 @@ dependencies = [
|
||||||
"cleanlab-tlm>=1.1.2",
|
"cleanlab-tlm>=1.1.2",
|
||||||
'gassist>=0.0.1; sys_platform == "win32"',
|
'gassist>=0.0.1; sys_platform == "win32"',
|
||||||
"twelvelabs>=0.4.7",
|
"twelvelabs>=0.4.7",
|
||||||
|
"filelock>=3.18.0",
|
||||||
"docling>=2.36.1",
|
"docling>=2.36.1",
|
||||||
|
"filelock>=3.18.0"
|
||||||
]
|
]
|
||||||
|
|
||||||
[dependency-groups]
|
[dependency-groups]
|
||||||
|
|
|
||||||
|
|
@ -1006,6 +1006,7 @@
|
||||||
"group_outputs": false,
|
"group_outputs": false,
|
||||||
"method": "retrieve_messages_dataframe",
|
"method": "retrieve_messages_dataframe",
|
||||||
"name": "dataframe",
|
"name": "dataframe",
|
||||||
|
"selected": null,
|
||||||
"tool_mode": true,
|
"tool_mode": true,
|
||||||
"types": [
|
"types": [
|
||||||
"DataFrame"
|
"DataFrame"
|
||||||
|
|
|
||||||
|
|
@ -154,10 +154,31 @@ def get_lifespan(*, fix_migration=False, version=None):
|
||||||
all_types_dict = await get_and_cache_all_types_dict(get_settings_service())
|
all_types_dict = await get_and_cache_all_types_dict(get_settings_service())
|
||||||
logger.debug(f"Types cached in {asyncio.get_event_loop().time() - current_time:.2f}s")
|
logger.debug(f"Types cached in {asyncio.get_event_loop().time() - current_time:.2f}s")
|
||||||
|
|
||||||
|
# Use file-based lock to prevent multiple workers from creating duplicate starter projects concurrently.
|
||||||
|
# Note that it's still possible that one worker may complete this task, release the lock,
|
||||||
|
# then another worker pick it up, but the operation is idempotent so worst case it duplicates
|
||||||
|
# the initialization work.
|
||||||
current_time = asyncio.get_event_loop().time()
|
current_time = asyncio.get_event_loop().time()
|
||||||
logger.debug("Creating/updating starter projects")
|
logger.debug("Creating/updating starter projects")
|
||||||
|
import tempfile
|
||||||
|
|
||||||
|
from filelock import FileLock
|
||||||
|
|
||||||
|
lock_file = Path(tempfile.gettempdir()) / "langflow_starter_projects.lock"
|
||||||
|
lock = FileLock(lock_file, timeout=1)
|
||||||
|
try:
|
||||||
|
with lock:
|
||||||
await create_or_update_starter_projects(all_types_dict)
|
await create_or_update_starter_projects(all_types_dict)
|
||||||
logger.debug(f"Starter projects updated in {asyncio.get_event_loop().time() - current_time:.2f}s")
|
logger.debug(
|
||||||
|
f"Starter projects created/updated in {asyncio.get_event_loop().time() - current_time:.2f}s"
|
||||||
|
)
|
||||||
|
except TimeoutError:
|
||||||
|
# Another process has the lock
|
||||||
|
logger.debug("Another worker is creating starter projects, skipping")
|
||||||
|
except Exception as e: # noqa: BLE001
|
||||||
|
logger.warning(
|
||||||
|
f"Failed to acquire lock for starter projects: {e}. Starter projects may not be created or updated."
|
||||||
|
)
|
||||||
|
|
||||||
current_time = asyncio.get_event_loop().time()
|
current_time = asyncio.get_event_loop().time()
|
||||||
logger.debug("Starting telemetry service")
|
logger.debug("Starting telemetry service")
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue