Added workflow test flag to server and updated workflow setup to use it
This commit is contained in:
+15
-13
@@ -59,6 +59,7 @@ workflow_redis_manager = None
|
||||
workflow_runner = None
|
||||
|
||||
USE_SIMULATED_WORKFLOW = os.getenv("WORKFLOW_SIMULATION", "0") == "1"
|
||||
WORKFLOW_TEST = False
|
||||
|
||||
_all_pgroups_cache: dict[str, tuple[list[str], float]] = {}
|
||||
_ALL_PGROUPS_TTL_S = 60.0 # adjust TTL as needed
|
||||
@@ -91,19 +92,20 @@ async def lifespan(application: FastAPI):
|
||||
daq = AareDAQ(cfg, bl)
|
||||
|
||||
# ── Workflow system ──
|
||||
# workflow_redis_manager = WorkflowRedisManager(
|
||||
# client=cfg._BeamlineConfig__client, # Reuse this worker's Redis connection
|
||||
# beamline=bl.value,
|
||||
# )
|
||||
# workflow_runner = PersistentWorkflowRunner(
|
||||
# redis_manager=workflow_redis_manager,
|
||||
# registry=STATE_REGISTRY,
|
||||
# handlers=SIMULATED_HANDLER_REGISTRY if USE_SIMULATED_WORKFLOW else HANDLER_REGISTRY,
|
||||
# )
|
||||
# if USE_SIMULATED_WORKFLOW:
|
||||
# logger.warning("⚠️ Workflow system running in SIMULATION mode - no actual DAQ operations")
|
||||
#
|
||||
# set_workflow_dependencies(workflow_redis_manager, workflow_runner, cfg)
|
||||
if WORKFLOW_TEST:
|
||||
workflow_redis_manager = WorkflowRedisManager(
|
||||
client=cfg._BeamlineConfig__client, # Reuse this worker's Redis connection
|
||||
beamline=bl.value,
|
||||
)
|
||||
workflow_runner = PersistentWorkflowRunner(
|
||||
redis_manager=workflow_redis_manager,
|
||||
registry=STATE_REGISTRY,
|
||||
handlers=SIMULATED_HANDLER_REGISTRY if USE_SIMULATED_WORKFLOW else HANDLER_REGISTRY,
|
||||
)
|
||||
if USE_SIMULATED_WORKFLOW:
|
||||
logger.warning("⚠️ Workflow system running in SIMULATION mode - no actual DAQ operations")
|
||||
|
||||
set_workflow_dependencies(workflow_redis_manager, workflow_runner, cfg)
|
||||
|
||||
# ── Initial TELL sync ──
|
||||
try:
|
||||
|
||||
Reference in New Issue
Block a user