@@ -43,6 +43,14 @@ class WorkflowCatalogError(Exception):
4343 """Base error for workflow catalog operations."""
4444
4545
46+ class WorkflowCatalogFetchError (WorkflowCatalogError ):
47+ """A configured workflow catalog could not be fetched."""
48+
49+
50+ class WorkflowCatalogValidationError (WorkflowCatalogError ):
51+ """A workflow catalog supplied malformed metadata."""
52+
53+
4654class WorkflowValidationError (WorkflowCatalogError ):
4755 """Validation error for catalog config or workflow data."""
4856
@@ -517,18 +525,36 @@ def _fetch_single_catalog(
517525 """Fetch a single catalog, using cache when possible."""
518526 cache_file , meta_file = self ._get_cache_paths (entry .url )
519527
528+ def validate_payload (data : Any ) -> dict [str , Any ]:
529+ if not isinstance (data , dict ):
530+ raise WorkflowCatalogValidationError (
531+ f"Catalog from { entry .url } is not a valid JSON object."
532+ )
533+ if "workflows" in data and not isinstance (data ["workflows" ], (dict , list )):
534+ raise WorkflowCatalogValidationError (
535+ f"Catalog from { entry .url } has malformed workflows metadata."
536+ )
537+ return data
538+
520539 if not force_refresh and self ._is_url_cache_valid (entry .url ):
521540 try :
522541 with open (cache_file , encoding = "utf-8" ) as f :
523542 cached = json .load (f )
524- if isinstance (cached , dict ):
525- return cached
526- except (UnicodeDecodeError , json .JSONDecodeError , OSError ):
543+ return validate_payload (cached )
544+ except (
545+ UnicodeDecodeError ,
546+ json .JSONDecodeError ,
547+ OSError ,
548+ WorkflowCatalogValidationError ,
549+ ):
527550 # Ignore invalid/unreadable cache and fall back to fetching from source.
528551 pass
529552
530553 # Fetch from URL — validate scheme before opening and after redirects
554+ from http .client import HTTPException
531555 from urllib .parse import urlparse
556+
557+ from specify_cli .authentication .http import RedirectPolicyError
532558 from specify_cli .authentication .http import open_url as _open_url
533559
534560 def _validate_catalog_url (url : str ) -> None :
@@ -542,18 +568,18 @@ def _validate_catalog_url(url: str) -> None:
542568 hostname = parsed .hostname
543569 _ = parsed .port
544570 except (TypeError , ValueError ):
545- raise WorkflowCatalogError (
571+ raise WorkflowCatalogValidationError (
546572 f"Refusing to fetch catalog from malformed URL: { url } "
547573 ) from None
548574 is_localhost = hostname in ("localhost" , "127.0.0.1" , "::1" )
549575 if parsed .scheme != "https" and not (
550576 parsed .scheme == "http" and is_localhost
551577 ):
552- raise WorkflowCatalogError (
578+ raise WorkflowCatalogValidationError (
553579 f"Refusing to fetch catalog from non-HTTPS URL: { url } "
554580 )
555581 if not hostname :
556- raise WorkflowCatalogError (
582+ raise WorkflowCatalogValidationError (
557583 f"Refusing to fetch catalog from URL with no hostname: { url } "
558584 )
559585
@@ -578,29 +604,39 @@ def _validate_redirect(_old_url: str, new_url: str) -> None:
578604 read_response_limited (
579605 resp ,
580606 max_bytes = _max_json_catalog_bytes (),
581- error_type = WorkflowCatalogError ,
607+ error_type = WorkflowCatalogValidationError ,
582608 label = "workflow catalog" ,
583609 ).decode ("utf-8" )
584610 )
585- except Exception as exc :
611+ except (
612+ WorkflowCatalogValidationError ,
613+ RedirectPolicyError ,
614+ UnicodeError ,
615+ json .JSONDecodeError ,
616+ ) as exc :
617+ raise WorkflowCatalogValidationError (
618+ f"Invalid workflow catalog from { entry .url } : { exc } "
619+ ) from exc
620+ except (OSError , HTTPException ) as exc :
586621 # Fall back to cache if available
587622 if cache_file .exists ():
588623 try :
589624 with open (cache_file , encoding = "utf-8" ) as f :
590625 cached = json .load (f )
591- if isinstance (cached , dict ):
592- return cached
593- except (json .JSONDecodeError , ValueError , OSError ):
626+ return validate_payload (cached )
627+ except (
628+ json .JSONDecodeError ,
629+ ValueError ,
630+ OSError ,
631+ WorkflowCatalogValidationError ,
632+ ):
594633 # Stale-cache read failed; let the original fetch error propagate.
595634 pass
596- raise WorkflowCatalogError (
635+ raise WorkflowCatalogFetchError (
597636 f"Failed to fetch catalog from { entry .url } : { exc } "
598637 ) from exc
599638
600- if not isinstance (data , dict ):
601- raise WorkflowCatalogError (
602- f"Catalog from { entry .url } is not a valid JSON object."
603- )
639+ data = validate_payload (data )
604640
605641 # Write cache
606642 try :
@@ -615,26 +651,44 @@ def _validate_redirect(_old_url: str, new_url: str) -> None:
615651 return data
616652
617653 def _get_merged_workflows (
618- self , force_refresh : bool = False
654+ self , force_refresh : bool = False , * , workflow_id : str | None = None
619655 ) -> dict [str , dict [str , Any ]]:
620- """Merge workflows from all active catalogs (lower priority number wins) ."""
656+ """Merge for search, or resolve one ID from the first valid winning source ."""
621657 catalogs = self .get_active_catalogs ()
622658 merged : dict [str , dict [str , Any ]] = {}
623659 fetch_errors = 0
660+ validation_error : WorkflowCatalogValidationError | None = None
624661
625- # Process later/higher-numbered entries first so earlier/lower-numbered
626- # entries overwrite them on workflow ID conflicts.
627- for entry in reversed (catalogs ):
662+ # Search uses overwrite order; exact-ID lookup visits the highest
663+ # priority source first and stops at its matching entry.
664+ sources = catalogs if workflow_id is not None else reversed (catalogs )
665+ for entry in sources :
628666 try :
629667 data = self ._fetch_single_catalog (entry , force_refresh )
630- except WorkflowCatalogError :
668+ except WorkflowCatalogError as exc :
669+ if workflow_id is not None and isinstance (
670+ exc , WorkflowCatalogValidationError
671+ ):
672+ raise
673+ if (
674+ isinstance (exc , WorkflowCatalogValidationError )
675+ and validation_error is None
676+ ):
677+ validation_error = exc
631678 fetch_errors += 1
632679 continue
633680 workflows = data .get ("workflows" , {})
634681 # Handle both dict and list formats
635682 if isinstance (workflows , dict ):
636683 for wf_id , wf_data in workflows .items ():
684+ if workflow_id is not None and wf_id != workflow_id :
685+ continue
637686 if not isinstance (wf_data , dict ):
687+ if workflow_id is not None :
688+ raise WorkflowCatalogValidationError (
689+ f"Invalid workflow catalog entry for '{ wf_id } ' "
690+ f"from { entry .url } : expected a JSON object"
691+ )
638692 continue
639693 wf_data ["_catalog_name" ] = entry .name
640694 wf_data ["_install_allowed" ] = entry .install_allowed
@@ -645,11 +699,17 @@ def _get_merged_workflows(
645699 continue
646700 wf_id = wf_data .get ("id" , "" )
647701 if wf_id :
702+ if workflow_id is not None and wf_id != workflow_id :
703+ continue
648704 wf_data ["_catalog_name" ] = entry .name
649705 wf_data ["_install_allowed" ] = entry .install_allowed
650706 merged [wf_id ] = wf_data
707+ if workflow_id is not None and workflow_id in merged :
708+ return merged
651709 if fetch_errors == len (catalogs ) and catalogs :
652- raise WorkflowCatalogError (
710+ if validation_error is not None :
711+ raise validation_error
712+ raise WorkflowCatalogFetchError (
653713 "All configured catalogs failed to fetch."
654714 )
655715 return merged
@@ -698,7 +758,7 @@ def get_workflow_info(
698758 """Get the current or an exact advertised release from the winning source."""
699759 from ._versions import select_release
700760
701- merged = self ._get_merged_workflows ()
761+ merged = self ._get_merged_workflows (workflow_id = workflow_id )
702762 wf = merged .get (workflow_id )
703763 if wf is None :
704764 return None
@@ -716,7 +776,7 @@ def get_workflow_version_details(
716776 """Return advertised versions and whether their source allows installation."""
717777 from ._versions import available_versions
718778
719- merged = self ._get_merged_workflows ()
779+ merged = self ._get_merged_workflows (workflow_id = workflow_id )
720780 wf = merged .get (workflow_id )
721781 if wf is None :
722782 return None
0 commit comments