Skip to content

Commit b855b5e

Browse files
authored
Merge pull request #382 from OpenMS/claude/fix-duplicate-config-Ju3LV
Centralize queue timeout settings in QueueManager
2 parents 36d2af7 + a681540 commit b855b5e

2 files changed

Lines changed: 26 additions & 14 deletions

File tree

src/workflow/QueueManager.py

Lines changed: 26 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -55,25 +55,34 @@ class QueueManager:
5555
def __init__(self):
5656
self._redis = None
5757
self._queue = None
58-
self._is_online = self._check_online_mode()
5958
self._init_attempted = False
6059

60+
settings = self._load_settings()
61+
self._is_online = self._check_online_mode(settings)
62+
63+
queue_settings = settings.get("queue_settings", {})
64+
self._default_timeout = queue_settings.get("default_timeout", 7200)
65+
self._default_result_ttl = queue_settings.get("result_ttl", 86400)
66+
6167
if self._is_online:
6268
self._init_redis()
6369

64-
def _check_online_mode(self) -> bool:
70+
@staticmethod
71+
def _load_settings() -> dict:
72+
"""Load settings.json once; return empty dict on failure."""
73+
try:
74+
with open("settings.json", "r") as f:
75+
return json.load(f)
76+
except Exception:
77+
return {}
78+
79+
def _check_online_mode(self, settings: dict) -> bool:
6580
"""Check if running in online mode"""
6681
# Check environment variable first (set in Docker)
6782
if os.environ.get("REDIS_URL"):
6883
return True
6984

70-
# Fallback: check settings file
71-
try:
72-
with open("settings.json", "r") as f:
73-
settings = json.load(f)
74-
return settings.get("online_deployment", False)
75-
except Exception:
76-
return False
85+
return settings.get("online_deployment", False)
7786

7887
def _init_redis(self) -> None:
7988
"""Initialize Redis connection and queue"""
@@ -108,8 +117,8 @@ def submit_job(
108117
args: tuple = (),
109118
kwargs: dict = None,
110119
job_id: Optional[str] = None,
111-
timeout: int = 7200, # 2 hour default
112-
result_ttl: int = 86400, # 24 hours
120+
timeout: Optional[int] = None,
121+
result_ttl: Optional[int] = None,
113122
description: str = ""
114123
) -> Optional[str]:
115124
"""
@@ -120,8 +129,8 @@ def submit_job(
120129
args: Positional arguments for the function
121130
kwargs: Keyword arguments for the function
122131
job_id: Optional custom job ID (defaults to UUID)
123-
timeout: Job timeout in seconds
124-
result_ttl: How long to keep results
132+
timeout: Job timeout in seconds (defaults to settings.json queue_settings.default_timeout)
133+
result_ttl: How long to keep results (defaults to settings.json queue_settings.result_ttl)
125134
description: Human-readable job description
126135
127136
Returns:
@@ -131,6 +140,10 @@ def submit_job(
131140
return None
132141

133142
kwargs = kwargs or {}
143+
if timeout is None:
144+
timeout = self._default_timeout
145+
if result_ttl is None:
146+
result_ttl = self._default_result_ttl
134147

135148
try:
136149
job = self._queue.enqueue(

src/workflow/WorkflowManager.py

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,6 @@ def _start_workflow_queued(self) -> None:
7272
"workflow_module": self.__class__.__module__,
7373
},
7474
job_id=job_id,
75-
timeout=7200, # 2 hour timeout
7675
description=f"Workflow: {self.name}"
7776
)
7877

0 commit comments

Comments
 (0)