Coverage for python/lsst/ctrl/bps/tests/gw_test_utils.py: 92%

265 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-25 15:17 -0700

1# This file is part of ctrl_bps. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

5# (https://www.lsst.org). 

6# See the COPYRIGHT file at the top-level directory of this distribution 

7# for details of code ownership. 

8# 

9# This software is dual licensed under the GNU General Public License and also 

10# under a 3-clause BSD license. Recipients may choose which of these licenses 

11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt, 

12# respectively. If you choose the GPL option then the following text applies 

13# (but note that there is still no warranty even if you opt for BSD instead): 

14# 

15# This program is free software: you can redistribute it and/or modify 

16# it under the terms of the GNU General Public License as published by 

17# the Free Software Foundation, either version 3 of the License, or 

18# (at your option) any later version. 

19# 

20# This program is distributed in the hope that it will be useful, 

21# but WITHOUT ANY WARRANTY; without even the implied warranty of 

22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the 

23# GNU General Public License for more details. 

24# 

25# You should have received a copy of the GNU General Public License 

26# along with this program. If not, see <https://www.gnu.org/licenses/>. 

27"""GenericWorkflow-related utilities to support ctrl_bps testing.""" 

28 

29__all__ = [ 

30 "make_3_label_workflow", 

31 "make_3_label_workflow_groups_sort", 

32 "make_3_label_workflow_noop_sort", 

33 "make_5_label_workflow", 

34 "make_5_label_workflow_2_groups", 

35 "make_5_label_workflow_middle_groups", 

36 "make_lazy_workflow", 

37] 

38 

39import logging 

40from collections import Counter 

41from typing import cast 

42 

43from lsst.ctrl.bps import ( 

44 GenericWorkflow, 

45 GenericWorkflowExec, 

46 GenericWorkflowGroup, 

47 GenericWorkflowJob, 

48 GenericWorkflowLazyGroup, 

49 GenericWorkflowNodeType, 

50 GenericWorkflowNoopJob, 

51) 

52 

53_LOG = logging.getLogger(__name__) 

54 

55 

56def make_3_label_workflow(workflow_name: str, final: bool) -> GenericWorkflow: 

57 """Create a simple 3 label test workflow. 

58 

59 Parameters 

60 ---------- 

61 workflow_name : `str` 

62 Name of the test workflow. 

63 final : `bool` 

64 Whether to add a final job. 

65 

66 Returns 

67 ------- 

68 gwf : `lsst.ctrl.bps.GenericWorkflow` 

69 The test workflow. 

70 """ 

71 gwexec = GenericWorkflowExec("exec1", "/usr/bin/uptime", False) 

72 gwf = GenericWorkflow(workflow_name) 

73 job = GenericWorkflowJob("pipetaskInit", label="pipetaskInit", executable=gwexec) 

74 gwf.add_job(job) 

75 for visit, vgroup in [ 

76 (10001, "2024-06-26T07:28:26.289"), 

77 (10002, "2024-06-26T07:29:06.969"), 

78 (301, "2024-06-26T07:27:45.775"), 

79 ]: # 301 is to ensure numeric sorting 

80 for detector in [10, 11]: 

81 prev_name = "pipetaskInit" 

82 for label in ["label1", "label2", "label3"]: 

83 name = f"{label}_{visit}_{detector}" 

84 job = GenericWorkflowJob( 

85 name, 

86 label=label, 

87 executable=gwexec, 

88 quanta_counts=Counter({label: 1}), 

89 tags={"visit": visit, "detector": detector, "group": vgroup}, 

90 ) 

91 gwf.add_job(job, [prev_name], None) 

92 prev_name = name 

93 

94 if final: 

95 gwexec = GenericWorkflowExec("finalJob.bash", "finalJob.bash", True) 

96 job = GenericWorkflowJob("finalJob", label="finalJob", executable=gwexec) 

97 gwf.add_final(job) 

98 

99 return gwf 

100 

101 

102def make_3_label_workflow_noop_sort(workflow_name: str, final: bool) -> GenericWorkflow: 

103 """Create a test workflow that has noop jobs. 

104 

105 Parameters 

106 ---------- 

107 workflow_name : `str` 

108 Name of the test workflow. 

109 final : `bool` 

110 Whether to add a final job. 

111 

112 Returns 

113 ------- 

114 gwf : `lsst.ctrl.bps.GenericWorkflow` 

115 The test workflow. 

116 """ 

117 gwexec = GenericWorkflowExec("exec1", "/usr/bin/uptime", False) 

118 gwf = GenericWorkflow(workflow_name) 

119 job = GenericWorkflowJob("pipetaskInit", label="pipetaskInit", executable=gwexec) 

120 gwf.add_job(job) 

121 prev_noop: GenericWorkflowNoopJob | None = None 

122 for visit in sorted([10001, 10002, 301]): # 301 is to ensure numeric sorting 

123 if visit != 10002: 

124 noop_job = GenericWorkflowNoopJob(f"noop_order1_{visit}", "order1") 

125 gwf.add_job(noop_job) 

126 for detector in [10, 11]: 

127 prev_name = "pipetaskInit" 

128 for label in ["label1", "label2", "label3"]: 

129 name = f"{label}_{visit}_{detector}" 

130 job = GenericWorkflowJob( 

131 name, label=label, executable=gwexec, tags={"visit": visit, "detector": detector} 

132 ) 

133 gwf.add_job(job, [prev_name], None) 

134 if label == "label1" and prev_noop: 

135 gwf.add_job_relationships([prev_noop.name], [name]) 

136 if label == "label2" and visit != 10002: 

137 gwf.add_job_relationships([name], [noop_job.name]) 

138 prev_name = name 

139 prev_noop = noop_job 

140 

141 if final: 141 ↛ 145line 141 didn't jump to line 145 because the condition on line 141 was always true

142 gwexec = GenericWorkflowExec("finalJob.bash", "finalJob.bash", True) 

143 job = GenericWorkflowJob("finalJob", label="finalJob", executable=gwexec) 

144 gwf.add_final(job) 

145 return gwf 

146 

147 

148def make_3_label_workflow_groups_sort(workflow_name: str, final: bool) -> GenericWorkflow: 

149 """Create a test workflow that has job groups. 

150 

151 Parameters 

152 ---------- 

153 workflow_name : `str` 

154 Name of the test workflow. 

155 final : `bool` 

156 Whether to add a final job. 

157 

158 Returns 

159 ------- 

160 gwf : `lsst.ctrl.bps.GenericWorkflow` 

161 The test workflow. 

162 """ 

163 gwexec = GenericWorkflowExec("exec1", "/usr/bin/uptime", False) 

164 gwf = GenericWorkflow(workflow_name) 

165 job = GenericWorkflowJob("pipetaskInit", label="pipetaskInit", executable=gwexec) 

166 gwf.add_job(job) 

167 prev_group: GenericWorkflowGroup | None = None 

168 for visit in sorted([10001, 10002, 301]): # 301 is to ensure numeric sorting 

169 job_group = GenericWorkflowGroup(f"group_order1_{visit}", "order1") 

170 for detector in [10, 11]: 

171 prev_name: str | None = None 

172 for label in ["label1", "label2"]: 

173 name = f"{label}_{visit}_{detector}" 

174 job = GenericWorkflowJob( 

175 name, label=label, executable=gwexec, tags={"visit": visit, "detector": detector} 

176 ) 

177 job_group.add_job(job) 

178 if prev_name: 

179 job_group.add_job_relationships(prev_name, name) 

180 prev_name = name 

181 gwf.add_job(job_group, ["pipetaskInit"], None) 

182 if prev_group: 

183 gwf.add_job_relationships([prev_group.name], [job_group.name]) 

184 

185 prev_group = job_group 

186 for visit in sorted([10001, 10002, 301]): # 301 is to ensure numeric sorting 

187 for detector in [10, 11]: 

188 for label in ["label3"]: 

189 name = f"{label}_{visit}_{detector}" 

190 job = GenericWorkflowJob( 

191 name, label=label, executable=gwexec, tags={"visit": visit, "detector": detector} 

192 ) 

193 gwf.add_job(job, [f"group_order1_{visit}"], None) 

194 

195 if final: 195 ↛ 200line 195 didn't jump to line 200 because the condition on line 195 was always true

196 gwexec = GenericWorkflowExec("finalJob.bash", "finalJob.bash", True) 

197 job = GenericWorkflowJob("finalJob", label="finalJob", executable=gwexec) 

198 gwf.add_final(job) 

199 

200 return gwf 

201 

202 

203# 301 is to ensure numeric sorting 

204DEFAULT_DIMS = [ 

205 (10001, 10), 

206 (10001, 11), 

207 (10001, 20), 

208 (10002, 10), 

209 (10002, 11), 

210 (10002, 20), 

211 (301, 10), 

212 (301, 11), 

213 (301, 20), 

214] 

215DIM_MAPPING = {301: "gval1", 10001: "gval2", 10002: "gval3"} 

216 

217UNEVEN_LABEL_DIMS = { 

218 "T1": [(10002, 11), (10002, 20)], 

219 "T2": [(10001, 11), (10001, 20), (10002, 10), (10002, 11), (10002, 20)], 

220 "T2b": [(301, 11), (301, 20), (10001, 11), (10001, 20), (10002, 10), (10002, 11), (10002, 20)], 

221 "T3": DEFAULT_DIMS, 

222 "T4": DEFAULT_DIMS, 

223} 

224 

225EVEN_LABEL_DIMS = { 

226 "T1": DEFAULT_DIMS, 

227 "T2": DEFAULT_DIMS, 

228 "T2b": DEFAULT_DIMS, 

229 "T3": DEFAULT_DIMS, 

230 "T4": DEFAULT_DIMS, 

231} 

232 

233 

234def make_5_label_workflow( 

235 workflow_name: str, final: bool, uneven: bool = False, equiv_dims: bool = False 

236) -> GenericWorkflow: 

237 """Create a simple 3 label test workflow. 

238 

239 Parameters 

240 ---------- 

241 workflow_name : `str` 

242 Name of the test workflow. 

243 final : `bool` 

244 Whether to add a final job. 

245 uneven : `bool`, optional 

246 Whether some of the jobs for initial tasks are 

247 not included as if finished in previous run. 

248 equiv_dims : `bool`, optional 

249 Whether first label jobs have a different but equivalent 

250 dim (like group and visit in AP pipeline). 

251 

252 Returns 

253 ------- 

254 gwf : `lsst.ctrl.bps.GenericWorkflow` 

255 The test workflow. 

256 """ 

257 gwexec = GenericWorkflowExec("exec1", "/usr/bin/uptime", False) 

258 gwf = GenericWorkflow(workflow_name) 

259 job = GenericWorkflowJob("pipetaskInit", label="pipetaskInit", executable=gwexec) 

260 gwf.add_job(job) 

261 if uneven: 

262 label_dims = UNEVEN_LABEL_DIMS 

263 else: 

264 label_dims = EVEN_LABEL_DIMS 

265 

266 prev_label = "pipetaskInit" 

267 for label in sorted(label_dims): 

268 for dim1, dim2 in label_dims[label]: 

269 tags: dict[str, str | int] = {"detector": dim2} 

270 # if want to test with equivalent dims (e.g., group and visit) 

271 if equiv_dims and label == "T1": 

272 tags["group"] = DIM_MAPPING[dim1] 

273 name = f"{label}_{DIM_MAPPING[dim1]}_{dim2}" 

274 else: 

275 tags["visit"] = dim1 

276 name = f"{label}_{dim1}_{dim2}" 

277 

278 job = GenericWorkflowJob( 

279 name, label=label, executable=gwexec, quanta_counts=Counter({label: 1}), tags=tags 

280 ) 

281 parents = [] 

282 if label == "T1": 

283 parents = ["pipetaskInit"] 

284 elif (dim1, dim2) in label_dims[prev_label]: 

285 if equiv_dims and label == "T2": 

286 prev_name = f"{prev_label}_{DIM_MAPPING[dim1]}_{dim2}" 

287 else: 

288 prev_name = f"{prev_label}_{dim1}_{dim2}" 

289 parents = [prev_name] 

290 else: 

291 parents = ["pipetaskInit"] 

292 

293 gwf.add_job(job, parents, None) 

294 

295 if label != "T2b": # nothing is a descenant of T2b 

296 prev_label = label 

297 

298 if final: 298 ↛ 303line 298 didn't jump to line 303 because the condition on line 298 was always true

299 gwexec = GenericWorkflowExec("finalJob.bash", "finalJob.bash", True) 

300 job = GenericWorkflowJob("finalJob", label="finalJob", executable=gwexec) 

301 gwf.add_final(job) 

302 

303 return gwf 

304 

305 

306def make_5_label_workflow_2_groups( 

307 workflow_name: str, final: bool, uneven: bool = False, equiv_dims: bool = False, blocking: bool = False 

308) -> GenericWorkflow: 

309 """Create a simple 3 label test workflow. 

310 

311 Parameters 

312 ---------- 

313 workflow_name : `str` 

314 Name of the test workflow. 

315 final : `bool` 

316 Whether to add a final job. 

317 uneven : `bool`, optional 

318 Whether some of the jobs for initial tasks are 

319 not included as if finished in previous run. 

320 equiv_dims : `bool`, optional 

321 Whether first label jobs have a different but equivalent 

322 dim (like group and visit in AP pipeline). 

323 blocking : `bool`, optional 

324 Value to use in group nodes. 

325 

326 Returns 

327 ------- 

328 gwf : `lsst.ctrl.bps.GenericWorkflow` 

329 The test workflow. 

330 """ 

331 gwf_orig = make_5_label_workflow("sink_uneven", final, uneven, equiv_dims) 

332 

333 if uneven: 

334 label_dims = UNEVEN_LABEL_DIMS 

335 else: 

336 label_dims = EVEN_LABEL_DIMS 

337 

338 # make job lists 

339 job_lists: dict[str, list[str]] = {} 

340 group_labels: dict[str, str] = {} 

341 

342 group_label = "order1" 

343 for dim1, dim2 in label_dims["T1"]: 

344 if equiv_dims: 344 ↛ 347line 344 didn't jump to line 347 because the condition on line 344 was always true

345 job_name = f"T1_{DIM_MAPPING[dim1]}_{dim2}" 

346 else: 

347 job_name = f"T1_{dim1}_{dim2}" 

348 group_name = f"group_{group_label}_{dim1}" 

349 group_labels[group_name] = group_label 

350 job_lists.setdefault(group_name, []).append(job_name) 

351 

352 for dim1, dim2 in label_dims["T2"]: 

353 job_name = f"T2_{dim1}_{dim2}" 

354 group_name = f"group_{group_label}_{dim1}" 

355 group_labels[group_name] = group_label 

356 job_lists.setdefault(group_name, []).append(job_name) 

357 

358 group_label = "order2" 

359 for label in ["T3", "T4"]: 

360 for dim1, dim2 in label_dims[label]: 

361 job_name = f"{label}_{dim1}_{dim2}" 

362 group_name = f"group_{group_label}_{dim1}" 

363 group_labels[group_name] = group_label 

364 job_lists.setdefault(group_name, []).append(job_name) 

365 

366 # make groups of jobs 

367 groups = {} 

368 for group_name, job_names in job_lists.items(): 

369 if job_names: 369 ↛ 368line 369 didn't jump to line 368 because the condition on line 369 was always true

370 group = GenericWorkflowGroup(group_name, group_labels[group_name], blocking=blocking) 

371 # Add all jobs first then add edges 

372 for job_name in job_names: 

373 group.add_job(gwf_orig.get_job(job_name)) 

374 

375 for name in job_names: 

376 edges = [(name, p) for p in gwf_orig.predecessors(name) if p in job_names] 

377 group.add_edges_from(edges) 

378 groups[group_name] = group 

379 

380 gwf = GenericWorkflow(workflow_name) 

381 

382 # add main workflow nodes 

383 gwf.add_job(gwf_orig.get_job("pipetaskInit")) 

384 for dim1, dim2 in label_dims["T2b"]: 

385 job_name = f"T2b_{dim1}_{dim2}" 

386 gwf.add_job(gwf_orig.get_job(job_name)) 

387 

388 for group in groups.values(): 

389 gwf.add_job(group) 

390 

391 # add main workflow edges 

392 edges = [ 

393 ("pipetaskInit", "group_order1_10001"), 

394 ("pipetaskInit", "group_order1_10002"), 

395 ("group_order1_10001", "group_order2_10001"), 

396 ("group_order1_10002", "group_order2_10002"), 

397 ("group_order1_10001", "T2b_10001_11"), 

398 ("group_order1_10001", "T2b_10001_20"), 

399 ("group_order1_10002", "T2b_10002_10"), 

400 ("group_order1_10002", "T2b_10002_11"), 

401 ("group_order1_10002", "T2b_10002_20"), 

402 # group order dependencies 

403 ("group_order1_10001", "group_order1_10002"), 

404 ("group_order2_301", "group_order2_10001"), 

405 ("group_order2_10001", "group_order2_10002"), 

406 ] 

407 

408 if uneven: 

409 edges.extend( 

410 [ 

411 ("pipetaskInit", "T2b_301_11"), 

412 ("pipetaskInit", "T2b_301_20"), 

413 ("pipetaskInit", "group_order2_301"), 

414 ("pipetaskInit", "group_order2_10001"), 

415 ] 

416 ) 

417 else: 

418 edges.extend( 

419 [ 

420 ("pipetaskInit", "group_order1_301"), 

421 ("group_order1_301", "group_order1_10001"), 

422 ("group_order1_301", "group_order2_301"), 

423 ("group_order1_301", "T2b_301_10"), 

424 ("group_order1_301", "T2b_301_11"), 

425 ("group_order1_301", "T2b_301_20"), 

426 ("group_order1_10001", "T2b_10001_10"), 

427 ] 

428 ) 

429 gwf.add_edges_from(edges) 

430 

431 if final: 431 ↛ 435line 431 didn't jump to line 435 because the condition on line 431 was always true

432 job = cast(GenericWorkflowJob, gwf_orig.get_final()) 

433 gwf.add_final(job) 

434 

435 return gwf 

436 

437 

438def make_5_label_workflow_middle_groups( 

439 workflow_name: str, final: bool, uneven: bool = False, equiv_dims: bool = False, blocking: bool = False 

440) -> GenericWorkflow: 

441 """Create a test workflow with a group in middle of workflow 

442 (T2, T2b, and T3). 

443 

444 Parameters 

445 ---------- 

446 workflow_name : `str` 

447 Name of the test workflow. 

448 final : `bool` 

449 Whether to add a final job. 

450 uneven : `bool`, optional 

451 Whether some of the jobs for initial tasks are 

452 not included as if finished in previous run. 

453 equiv_dims : `bool`, optional 

454 Whether first label jobs have a different but equivalent 

455 dim (like group and visit in AP pipeline). 

456 blocking : `bool`, optional 

457 Value to use in group nodes. 

458 

459 Returns 

460 ------- 

461 gwf : `lsst.ctrl.bps.GenericWorkflow` 

462 The test workflow. 

463 """ 

464 gwf_orig = make_5_label_workflow(workflow_name, final, uneven, equiv_dims) 

465 

466 if uneven: 

467 label_dims = UNEVEN_LABEL_DIMS 

468 else: 

469 label_dims = EVEN_LABEL_DIMS 

470 

471 # make job lists 

472 job_lists: dict[str, list[str]] = {} 

473 group_labels: dict[str, str] = {} 

474 

475 group_label = "mid" 

476 for label in ["T2", "T2b", "T3"]: 

477 for dim1, dim2 in label_dims[label]: 

478 job_name = f"{label}_{dim1}_{dim2}" 

479 group_name = f"group_{group_label}_{dim1}" 

480 group_labels[group_name] = group_label 

481 job_lists.setdefault(group_name, []).append(job_name) 

482 

483 # make groups of jobs 

484 groups = {} 

485 for group_name, job_names in job_lists.items(): 

486 if job_names: 486 ↛ 485line 486 didn't jump to line 485 because the condition on line 486 was always true

487 group = GenericWorkflowGroup(group_name, group_labels[group_name], blocking=blocking) 

488 # Add all jobs first then add edges 

489 for job_name in job_names: 

490 group.add_job(gwf_orig.get_job(job_name)) 

491 

492 for name in job_names: 

493 edges = [(name, p) for p in gwf_orig.predecessors(name) if p in job_names] 

494 group.add_edges_from(edges) 

495 

496 groups[group_name] = group 

497 

498 gwf = GenericWorkflow(workflow_name) 

499 

500 # add main workflow nodes 

501 gwf.add_job(gwf_orig.get_job("pipetaskInit")) 

502 for label in ["T1", "T4"]: 

503 for dim1, dim2 in label_dims[label]: 

504 if equiv_dims and label == "T1": 

505 job_name = f"T1_{DIM_MAPPING[dim1]}_{dim2}" 

506 else: 

507 job_name = f"{label}_{dim1}_{dim2}" 

508 gwf.add_job(gwf_orig.get_job(job_name)) 

509 

510 for group in groups.values(): 

511 gwf.add_job(group) 

512 

513 # add main workflow edges 

514 edges = [ 

515 ("group_mid_301", "T4_301_10"), 

516 ("group_mid_301", "T4_301_11"), 

517 ("group_mid_301", "T4_301_20"), 

518 ("group_mid_10001", "T4_10001_10"), 

519 ("group_mid_10001", "T4_10001_11"), 

520 ("group_mid_10001", "T4_10001_20"), 

521 ("group_mid_10002", "T4_10002_10"), 

522 ("group_mid_10002", "T4_10002_11"), 

523 ("group_mid_10002", "T4_10002_20"), 

524 # group order dependencies 

525 ("group_mid_301", "group_mid_10001"), 

526 ("group_mid_10001", "group_mid_10002"), 

527 ] 

528 

529 if uneven: 

530 if equiv_dims: 530 ↛ 540line 530 didn't jump to line 540 because the condition on line 530 was always true

531 edges.extend( 

532 [ 

533 ("pipetaskInit", "T1_gval3_11"), 

534 ("pipetaskInit", "T1_gval3_20"), 

535 ("T1_gval3_11", "group_mid_10002"), 

536 ("T1_gval3_20", "group_mid_10002"), 

537 ] 

538 ) 

539 else: 

540 edges.extend( 

541 [ 

542 ("pipetaskInit", "T1_10002_11"), 

543 ("pipetaskInit", "T1_10002_20"), 

544 ("T1_10002_11", "group_mid_10002"), 

545 ("T1_10002_20", "group_mid_10002"), 

546 ] 

547 ) 

548 

549 # Because in orig workflow, pipetaskInit has edge to T2(10002, 10), 

550 # there will be an "extra" edge from pipetaskInit to group_mid_10002. 

551 edges.extend( 

552 [ 

553 ("pipetaskInit", "group_mid_301"), 

554 ("pipetaskInit", "group_mid_10001"), 

555 ("pipetaskInit", "group_mid_10002"), 

556 ] 

557 ) 

558 else: 

559 dim1s = [301, 10001, 10002] 

560 for dim1 in dim1s: 

561 T1_dim1: str | int = dim1 

562 if equiv_dims: 562 ↛ 565line 562 didn't jump to line 565 because the condition on line 562 was always true

563 T1_dim1 = DIM_MAPPING[dim1] 

564 

565 edges.extend( 

566 [ 

567 ("pipetaskInit", f"T1_{T1_dim1}_10"), 

568 ("pipetaskInit", f"T1_{T1_dim1}_11"), 

569 ("pipetaskInit", f"T1_{T1_dim1}_20"), 

570 (f"T1_{T1_dim1}_10", f"group_mid_{dim1}"), 

571 (f"T1_{T1_dim1}_11", f"group_mid_{dim1}"), 

572 (f"T1_{T1_dim1}_20", f"group_mid_{dim1}"), 

573 ] 

574 ) 

575 gwf.add_edges_from(edges) 

576 

577 if final: 577 ↛ 581line 577 didn't jump to line 581 because the condition on line 577 was always true

578 job = cast(GenericWorkflowJob, gwf_orig.get_final()) 

579 gwf.add_final(job) 

580 

581 return gwf 

582 

583 

584def compare_generic_workflows(gwf1: GenericWorkflow, gwf2: GenericWorkflow) -> bool: 

585 """Compare two workflows printing log messages where not equal. 

586 

587 Parameters 

588 ---------- 

589 gwf1 : `lsst.ctrl.bps.GenericWorkflow` 

590 First workflow to compare. 

591 gwf2 : `lsst.ctrl.bps.GenericWorkflow` 

592 Second workflow to compare. 

593 

594 Returns 

595 ------- 

596 equal : bool 

597 Whether the two workflows are the same. 

598 """ 

599 equal = True 

600 

601 # check edges 

602 edges1 = set(gwf1.edges) 

603 edges2 = set(gwf2.edges) 

604 only_in_first = edges1 - edges2 

605 only_in_second = edges2 - edges1 

606 if only_in_first: 606 ↛ 607line 606 didn't jump to line 607 because the condition on line 606 was never true

607 _LOG.debug("Edges only in %s, but not in %s: %s", gwf1.name, gwf2.name, only_in_first) 

608 equal = False 

609 if only_in_second: 609 ↛ 610line 609 didn't jump to line 610 because the condition on line 609 was never true

610 _LOG.debug("Edges only in %s, but not in %s: %s", gwf2.name, gwf1.name, only_in_second) 

611 equal = False 

612 

613 # check nodes 

614 names1 = set(gwf1.nodes) 

615 names2 = set(gwf2.nodes) 

616 only_in_first = names1 - names2 

617 only_in_second = names2 - names1 

618 if only_in_first: 618 ↛ 619line 618 didn't jump to line 619 because the condition on line 618 was never true

619 _LOG.debug("Jobs only in %s, but not in %s: %s", gwf1.name, gwf2.name, only_in_first) 

620 equal = False 

621 if only_in_second: 621 ↛ 622line 621 didn't jump to line 622 because the condition on line 621 was never true

622 _LOG.debug("Jobs only in %s, but not in %s: %s", gwf2.name, gwf1.name, only_in_second) 

623 equal = False 

624 

625 # check node values 

626 for name in names1 & names2: 

627 job1 = gwf1.get_job(name) 

628 job2 = gwf2.get_job(name) 

629 

630 if job1.node_type != job2.node_type: 630 ↛ 631line 630 didn't jump to line 631 because the condition on line 630 was never true

631 _LOG.debug( 

632 "Group jobs` node_type not equal %s=%s, %s=%s", 

633 job1.name, 

634 job1.blocking, 

635 job2.name, 

636 job2.blocking, 

637 ) 

638 equal = False 

639 elif job1.node_type == GenericWorkflowNodeType.GROUP: 

640 if job1.blocking != job2.blocking: 640 ↛ 641line 640 didn't jump to line 641 because the condition on line 640 was never true

641 _LOG.debug( 

642 "Group jobs` blocking not equal %s=%s, %s=%s", 

643 job1.name, 

644 job1.blocking, 

645 job2.name, 

646 job2.blocking, 

647 ) 

648 equal = False 

649 

650 # compare workflows 

651 equal = equal or compare_generic_workflows(job1, job2) 

652 

653 # check final 

654 final1 = gwf1.get_final() 

655 final2 = gwf2.get_final() 

656 if final1 != final2: 656 ↛ 657line 656 didn't jump to line 657 because the condition on line 656 was never true

657 _LOG.debug("Final jobs are not equal: %s vs %s", final1, final2) 

658 equal = False 

659 

660 return equal 

661 

662 

663def make_lazy_workflow(workflow_name: str, final: bool) -> GenericWorkflow: # pragma: no cover 

664 """Create a simple workflow with a lazy workflow node for WMS tests. 

665 

666 Parameters 

667 ---------- 

668 workflow_name : `str` 

669 Name of the test workflow. 

670 final : `bool` 

671 Whether to add a final job. 

672 

673 Returns 

674 ------- 

675 gwf : `lsst.ctrl.bps.GenericWorkflow` 

676 The test workflow. 

677 """ 

678 gwf = GenericWorkflow(workflow_name) 

679 

680 # job 1 

681 gwexec1 = GenericWorkflowExec("exec1", "my_exec1.sh", False) 

682 job1 = GenericWorkflowJob("job1", "label1", executable=gwexec1) 

683 gwf.add_job(job1, None) 

684 

685 # lazy workflow job 

686 gwexec2 = GenericWorkflowExec("exec2", "${CTRL_BPS_DIR}/python/lsst/ctrl/bps/_make_workflow.sh", False) 

687 job2 = GenericWorkflowLazyGroup("lazy2", "label2", executable=gwexec2) 

688 gwf.add_job(job2, [job1.name]) 

689 

690 # Job after 

691 gwexec3 = GenericWorkflowExec("exec3", "my_exec3.sh", False) 

692 job3 = GenericWorkflowJob("job3", "label3", executable=gwexec3) 

693 gwf.add_job(job3, [job2.name]) 

694 

695 if final: 

696 gwexec = GenericWorkflowExec("finalJob.bash", "finalJob.bash", True) 

697 job = GenericWorkflowJob("finalJob", label="finalJob", executable=gwexec) 

698 gwf.add_final(job) 

699 

700 return gwf