Эх сурвалжийг харах

metadata_manager: refactor to a provider based versions manager

Shiv Tyagi 2 сар өмнө
parent
commit
839c33eae6

+ 31 - 13
metadata_manager/__init__.py

@@ -1,21 +1,39 @@
-from .versions_fetcher import (
-    VersionsFetcher,
+from .ap_src_meta_fetcher import APSourceMetadataFetcher
+from .firmware_server import ManifestJSON, ManifestFetchError, ReleaseRecord
+from .firmware_server.client import ManifestClient
+from .firmware_server.index import ManifestIndex
+from .vehicles_manager import DEFAULT_VEHICLES, Vehicle, VehiclesManager
+from .versions_manager import (
+    DEFAULT_WHITELISTED_FORK_REMOTES,
+    ForkRemoteSpec,
+    ManifestJsonVersionsProvider,
     RemoteInfo,
-)
-
-from .ap_src_meta_fetcher import (
-    APSourceMetadataFetcher,
-)
-
-from .vehicles_manager import (
-    VehiclesManager,
-    Vehicle,
+    RemotesJsonVersionsProvider,
+    VersionInfo,
+    VersionsManager,
+    VersionsProvider,
+    WhitelistedForkTagVersionsProvider,
+    build_default_providers,
 )
 
 __all__ = [
     "APSourceMetadataFetcher",
-    "VersionsFetcher",
+    "DEFAULT_VEHICLES",
+    "DEFAULT_WHITELISTED_FORK_REMOTES",
+    "ForkRemoteSpec",
+    "ManifestClient",
+    "ManifestFetchError",
+    "ManifestIndex",
+    "ManifestJSON",
+    "ManifestJsonVersionsProvider",
+    "ReleaseRecord",
     "RemoteInfo",
-    "VehiclesManager",
+    "RemotesJsonVersionsProvider",
     "Vehicle",
+    "VersionInfo",
+    "VersionsManager",
+    "VersionsProvider",
+    "VehiclesManager",
+    "WhitelistedForkTagVersionsProvider",
+    "build_default_providers",
 ]

+ 1 - 1
metadata_manager/ap_src_meta_fetcher.py

@@ -505,7 +505,7 @@ class APSourceMetadataFetcher:
 
         # Heli builds are stored under a separate folder
         artifacts_subdir = board_id
-        if vehicle_id == "Heli":
+        if vehicle_id == "heli":
             artifacts_subdir += "-heli"
 
         features_txt_url = f"{artifacts_url}/{artifacts_subdir}/features.txt"

+ 9 - 0
metadata_manager/firmware_server/__init__.py

@@ -0,0 +1,9 @@
+from .manifest import ManifestJSON
+from .models import ReleaseRecord
+from .exceptions import ManifestFetchError
+
+__all__ = [
+    "ManifestJSON",
+    "ManifestFetchError",
+    "ReleaseRecord",
+]

+ 149 - 0
metadata_manager/firmware_server/client.py

@@ -0,0 +1,149 @@
+import json
+import logging
+import os
+from dataclasses import dataclass
+from datetime import datetime, timezone
+from pathlib import Path
+from typing import Optional
+
+import requests
+
+from .exceptions import ManifestFetchError
+
+
+@dataclass
+class _CacheMeta:
+    etag: Optional[str] = None
+    last_modified: Optional[str] = None
+    fetched_at: Optional[str] = None
+
+    def to_dict(self) -> dict:
+        return {
+            "etag": self.etag,
+            "last_modified": self.last_modified,
+            "fetched_at": self.fetched_at,
+        }
+
+    @classmethod
+    def from_dict(cls, data: dict) -> "_CacheMeta":
+        return cls(
+            etag=data.get("etag"),
+            last_modified=data.get("last_modified"),
+            fetched_at=data.get("fetched_at"),
+        )
+
+
+class ManifestClient:
+    """Fetch and cache the ArduPilot firmware manifest.json file."""
+
+    def __init__(
+        self,
+        url: str,
+        cache_dir: str,
+        timeout: int = 120,
+        user_agent: str = "CustomBuild/1.0",
+    ):
+        self.url = url
+        self.cache_dir = Path(cache_dir)
+        self.cache_path = self.cache_dir / "manifest.json"
+        self.meta_path = self.cache_dir / "manifest.json.meta"
+        self.timeout = timeout
+        self.user_agent = user_agent
+        self.logger = logging.getLogger(__name__)
+
+    def fetch_raw(self) -> bytes:
+        headers = {"User-Agent": self.user_agent}
+        meta = self._read_meta() if self._has_cache() else _CacheMeta()
+
+        if meta.etag:
+            headers["If-None-Match"] = meta.etag
+        if meta.last_modified:
+            headers["If-Modified-Since"] = meta.last_modified
+
+        try:
+            response = requests.get(
+                self.url,
+                headers=headers,
+                timeout=self.timeout,
+            )
+        except requests.RequestException as exc:
+            return self._fallback_or_raise(exc)
+
+        if response.status_code == 304:
+            self.logger.info("Manifest not modified (304), using cache")
+            self._touch_fetched_at(self._now_iso())
+            return self._read_cache_bytes()
+
+        if response.status_code != 200:
+            return self._fallback_or_raise(
+                ManifestFetchError(
+                    f"Manifest fetch failed with status {response.status_code}"
+                )
+            )
+
+        raw = response.content
+        self._write_cache(
+            raw,
+            _CacheMeta(
+                etag=response.headers.get("ETag"),
+                last_modified=response.headers.get("Last-Modified"),
+                fetched_at=self._now_iso(),
+            ),
+        )
+        self.logger.info("Downloaded manifest (%d bytes)", len(raw))
+        return raw
+
+    def fetch(self) -> dict:
+        return json.loads(self.fetch_raw().decode("utf-8"))
+
+    def _has_cache(self) -> bool:
+        return self.cache_path.is_file()
+
+    def _read_cache_bytes(self) -> bytes:
+        return self.cache_path.read_bytes()
+
+    def _read_meta(self) -> _CacheMeta:
+        if not self.meta_path.is_file():
+            return _CacheMeta()
+        return _CacheMeta.from_dict(
+            json.loads(self.meta_path.read_text(encoding="utf-8"))
+        )
+
+    def _write_cache(self, raw: bytes, meta: _CacheMeta) -> None:
+        self._atomic_write_bytes(self.cache_path, raw)
+        self._atomic_write_text(
+            self.meta_path,
+            json.dumps(meta.to_dict(), indent=2),
+        )
+
+    def _touch_fetched_at(self, fetched_at: str) -> None:
+        meta = self._read_meta()
+        meta.fetched_at = fetched_at
+        self._atomic_write_text(
+            self.meta_path,
+            json.dumps(meta.to_dict(), indent=2),
+        )
+
+    def _atomic_write_bytes(self, path: Path, raw: bytes) -> None:
+        path.parent.mkdir(parents=True, exist_ok=True)
+        tmp_path = path.with_name(f"{path.name}.tmp")
+        tmp_path.write_bytes(raw)
+        os.replace(tmp_path, path)
+
+    def _atomic_write_text(self, path: Path, text: str) -> None:
+        path.parent.mkdir(parents=True, exist_ok=True)
+        tmp_path = path.with_name(f"{path.name}.tmp")
+        tmp_path.write_text(text, encoding="utf-8")
+        os.replace(tmp_path, path)
+
+    def _fallback_or_raise(self, exc: Exception) -> bytes:
+        if self._has_cache():
+            self.logger.warning(
+                "Manifest fetch failed (%s), using stale cache", exc
+            )
+            return self._read_cache_bytes()
+        raise ManifestFetchError(str(exc)) from exc
+
+    @staticmethod
+    def _now_iso() -> str:
+        return datetime.now(timezone.utc).isoformat()

+ 2 - 0
metadata_manager/firmware_server/exceptions.py

@@ -0,0 +1,2 @@
+class ManifestFetchError(Exception):
+    """Raised when the firmware manifest cannot be fetched and no cache exists."""

+ 207 - 0
metadata_manager/firmware_server/index.py

@@ -0,0 +1,207 @@
+import logging
+import re
+from collections import defaultdict
+from typing import Optional
+
+from packaging.version import InvalidVersion, Version
+
+from .models import ReleaseRecord
+
+# Minimum version for a vehicle to expose from manifest entries.
+MIN_VERSION_BY_VEHICLE_ID = {
+    "copter": "4.3",
+    "plane": "4.3",
+    "rover": "4.3",
+    "sub": "4.3",
+    "tracker": "4.3",
+    "blimp": "4.3",
+    "heli": "4.3",
+    "ap-periph": "1.8.1",
+}
+
+
+def vehicle_id_for_manifest_entry(entry: dict) -> Optional[str]:
+    match entry.get("vehicletype", ""):
+        case "Copter":
+            return "heli" if entry.get("mav-type") == "HELICOPTER" else "copter"
+        case "Plane":
+            return "plane"
+        case "Rover":
+            return "rover"
+        case "Sub":
+            return "sub"
+        case "Blimp":
+            return "blimp"
+        case "AntennaTracker":
+            return "tracker"
+        case "AP_Periph":
+            return "ap-periph"
+        case _:
+            return None
+
+
+def _should_skip_version(
+    vehicle_id: str, release_type: str, manifest_version: str
+) -> bool:
+    if release_type == "latest":
+        return False
+
+    min_version = MIN_VERSION_BY_VEHICLE_ID.get(vehicle_id)
+    if not min_version:
+        return False
+
+    try:
+        return Version(manifest_version) < Version(min_version)
+    except InvalidVersion:
+        return True
+
+
+def map_manifest_release_type(mav_firmware_version_type: str) -> Optional[str]:
+    if not mav_firmware_version_type:
+        return None
+    if mav_firmware_version_type.startswith("STABLE-"):
+        return "stable"
+    if mav_firmware_version_type == "OFFICIAL":
+        return "stable"
+    if mav_firmware_version_type == "BETA":
+        return "beta"
+    if mav_firmware_version_type == "DEV":
+        return "latest"
+    return mav_firmware_version_type.lower()
+
+
+def parse_artifacts_base_url(url: str) -> Optional[str]:
+    match = re.match(
+        r"(https://firmware\.ardupilot\.org/[^/]+/[^/]+)/",
+        url or "",
+    )
+    if match:
+        return match.group(1)
+    return None
+
+
+def _release_fields_from_entry(
+    entry: dict,
+) -> Optional[tuple[str, str, str, str, str]]:
+    vehicle_id = vehicle_id_for_manifest_entry(entry)
+    if vehicle_id is None:
+        return None
+
+    release_type = map_manifest_release_type(
+        entry.get("mav-firmware-version-type", "")
+    )
+    manifest_version = entry.get("mav-firmware-version")
+    if not release_type or not manifest_version:
+        return None
+
+    if _should_skip_version(vehicle_id, release_type, manifest_version):
+        return None
+
+    base_url = parse_artifacts_base_url(entry.get("url", ""))
+    if not base_url:
+        return None
+
+    git_sha = entry.get("git-sha")
+    if not git_sha:
+        return None
+
+    version_number = "NA" if release_type == "latest" else manifest_version
+    return vehicle_id, release_type, version_number, base_url, git_sha
+
+
+def _record_release_meta(
+    release_meta: dict[tuple, dict],
+    vehicle_id: str,
+    release_type: str,
+    version_number: str,
+    base_url: str,
+    git_sha: str,
+    logger: logging.Logger,
+) -> None:
+    key = (vehicle_id, release_type, version_number)
+    if key not in release_meta:
+        release_meta[key] = {
+            "vehicle_id": vehicle_id,
+            "release_type": release_type,
+            "version_number": version_number,
+            "ap_build_artifacts_url": base_url,
+            "git_sha": git_sha,
+        }
+        return
+
+    if release_meta[key]["git_sha"] != git_sha:
+        logger.debug(
+            "Conflicting git-sha for %s: keeping %s, ignoring %s",
+            key,
+            release_meta[key]["git_sha"][:8],
+            git_sha[:8],
+        )
+
+
+def _releases_from_meta(
+    release_meta: dict[tuple, dict],
+) -> dict[str, list[ReleaseRecord]]:
+    releases_by_vehicle: dict[str, list[ReleaseRecord]] = defaultdict(list)
+    for meta in release_meta.values():
+        releases_by_vehicle[meta["vehicle_id"]].append(
+            ReleaseRecord(
+                vehicle_id=meta["vehicle_id"],
+                release_type=meta["release_type"],
+                version_number=meta["version_number"],
+                commit_reference=meta["git_sha"],
+                ap_build_artifacts_url=meta["ap_build_artifacts_url"],
+            )
+        )
+
+    for vehicle_id in releases_by_vehicle:
+        releases_by_vehicle[vehicle_id].sort(
+            key=lambda r: (
+                0 if r.release_type == "latest" else 1,
+                r.release_type,
+                _version_sort_key(r.version_number),
+            )
+        )
+
+    return dict(releases_by_vehicle)
+
+
+class ManifestIndex:
+    """In-memory index built from a parsed manifest document."""
+
+    def __init__(self, releases_by_vehicle: dict[str, list[ReleaseRecord]]):
+        self.releases_by_vehicle = releases_by_vehicle
+
+    @classmethod
+    def build(cls, manifest: dict) -> "ManifestIndex":
+        logger = logging.getLogger(__name__)
+        release_meta: dict[tuple, dict] = {}
+
+        for entry in manifest.get("firmware") or []:
+            fields = _release_fields_from_entry(entry)
+            if fields is None:
+                continue
+            vehicle_id, release_type, version_number, base_url, git_sha = fields
+            _record_release_meta(
+                release_meta,
+                vehicle_id,
+                release_type,
+                version_number,
+                base_url,
+                git_sha,
+                logger,
+            )
+
+        return cls(_releases_from_meta(release_meta))
+
+    def get_releases(self, vehicle_id: str) -> list[ReleaseRecord]:
+        return list(self.releases_by_vehicle.get(vehicle_id, []))
+
+
+def _version_sort_key(version_number: str) -> tuple:
+    if version_number == "NA":
+        return (0,)
+    try:
+        parsed = Version(version_number)
+        return (1, parsed.major, parsed.minor, parsed.micro)
+    except InvalidVersion:
+        return (2, version_number)

+ 40 - 0
metadata_manager/firmware_server/manifest.py

@@ -0,0 +1,40 @@
+import logging
+from typing import Optional
+
+from .client import ManifestClient
+from .exceptions import ManifestFetchError
+from .index import ManifestIndex
+
+
+class ManifestJSON:
+    """Facade over cached manifest data for version and firmware-server lookups."""
+
+    def __init__(self, url: str, cache_dir: str):
+        self._client = ManifestClient(url=url, cache_dir=cache_dir)
+        self._index: Optional[ManifestIndex] = None
+        self.logger = logging.getLogger(__name__)
+
+    @property
+    def is_available(self) -> bool:
+        return self._index is not None
+
+    def refresh(self) -> None:
+        try:
+            manifest = self._client.fetch()
+            self._index = ManifestIndex.build(manifest)
+            self.logger.info(
+                "Manifest index built with releases for %d vehicles",
+                len(self._index.releases_by_vehicle),
+            )
+        except ManifestFetchError:
+            if self._index is not None:
+                self.logger.warning(
+                    "Manifest refresh failed, continuing with stale index"
+                )
+                return
+            raise
+
+    def get_releases(self, vehicle_id: str) -> list:
+        if not self.is_available:
+            return []
+        return self._index.get_releases(vehicle_id)

+ 10 - 0
metadata_manager/firmware_server/models.py

@@ -0,0 +1,10 @@
+from dataclasses import dataclass
+
+
+@dataclass(frozen=True)
+class ReleaseRecord:
+    vehicle_id: str
+    release_type: str
+    version_number: str
+    commit_reference: str
+    ap_build_artifacts_url: str

+ 3 - 3
metadata_manager/remotes.schema.json

@@ -21,9 +21,9 @@
             "type": "object",
             "description": "Vehicle object",
             "properties": {
-              "name": {
+              "id": {
                 "type": "string",
-                "description": "Name of vehicle"
+                "description": "Vehicle id"
               },
               "releases": {
                 "type": "array",
@@ -55,7 +55,7 @@
               }
             },
             "required": [
-              "name",
+              "id",
               "releases"
             ]
           }

+ 0 - 10
metadata_manager/vehicles_manager.py

@@ -3,13 +3,11 @@ class Vehicle:
                  id: str,
                  name: str,
                  ap_source_subdir: str,
-                 fw_server_vehicle_sdir: str,
                  waf_build_command: str,
                  ) -> None:
         self.id = id
         self.name = name
         self.ap_source_subdir = ap_source_subdir
-        self.fw_server_vehicle_sdir = fw_server_vehicle_sdir
         self.waf_build_command = waf_build_command
 
     def __eq__(self, other):
@@ -27,56 +25,48 @@ DEFAULT_VEHICLES = [
         id="copter",
         name="Copter",
         ap_source_subdir="ArduCopter",
-        fw_server_vehicle_sdir="Copter",
         waf_build_command="copter"
     ),
     Vehicle(
         id="plane",
         name="Plane",
         ap_source_subdir="ArduPlane",
-        fw_server_vehicle_sdir="Plane",
         waf_build_command="plane"
     ),
     Vehicle(
         id="rover",
         name="Rover",
         ap_source_subdir="Rover",
-        fw_server_vehicle_sdir="Rover",
         waf_build_command="rover"
     ),
     Vehicle(
         id="sub",
         name="Sub",
         ap_source_subdir="ArduSub",
-        fw_server_vehicle_sdir="Sub",
         waf_build_command="sub"
     ),
     Vehicle(
         id="heli",
         name="Heli",
         ap_source_subdir="ArduCopter",
-        fw_server_vehicle_sdir="Copter",
         waf_build_command="heli"
     ),
     Vehicle(
         id="blimp",
         name="Blimp",
         ap_source_subdir="Blimp",
-        fw_server_vehicle_sdir="Blimp",
         waf_build_command="blimp"
     ),
     Vehicle(
         id="tracker",
         name="Tracker",
         ap_source_subdir="AntennaTracker",
-        fw_server_vehicle_sdir="AntennaTracker",
         waf_build_command="antennatracker"
     ),
     Vehicle(
         id="ap-periph",
         name="AP_Periph",
         ap_source_subdir="Tools/AP_Periph",
-        fw_server_vehicle_sdir="AP_Periph",
         waf_build_command="AP_Periph"
     ),
 ]

+ 0 - 358
metadata_manager/versions_fetcher.py

@@ -1,358 +0,0 @@
-import logging
-import os
-import ap_git
-import json
-import jsonschema
-import hashlib
-from pathlib import Path
-from threading import Lock
-from utils import TaskRunner
-from .vehicles_manager import VehiclesManager as vehm
-
-
-class VersionInfo:
-    """
-    Class to wrap version info properties.
-    """
-    def __init__(self,
-                 remote_info: 'RemoteInfo',
-                 commit_ref: str,
-                 release_type: str,
-                 version_number: str,
-                 ap_build_artifacts_url) -> None:
-        self.remote_info = remote_info
-        self.commit_ref = commit_ref
-        self.release_type = release_type
-        self.version_number = version_number
-        self.ap_build_artifacts_url = ap_build_artifacts_url
-
-        # Generate version_id as remote-sanitized_commit_ref-hash
-        # Commit ref is sanitized for URL safety by replacing '/' with '-'
-        # Hash is used to ensure unique ID after sanitization of commit_ref,
-        # as different commit refs may become identical after replacing '/'
-        commit_ref_sanitized = commit_ref.replace('/', '-')
-        commit_ref_hash = hashlib.md5(commit_ref.encode()).hexdigest()[:8]
-        self.version_id = (
-            f"{remote_info.name}-{commit_ref_sanitized}-{commit_ref_hash}"
-        )
-
-
-class RemoteInfo:
-    """
-    Class to wrap remote info properties.
-    """
-    def __init__(self,
-                 name: str,
-                 url: str) -> None:
-        self.name = name
-        self.url = url
-
-    def to_dict(self):
-        return {
-            'name': self.name,
-            'url': self.url,
-        }
-
-
-class VersionsFetcher:
-    """
-    Class to fetch the version-to-build metadata from remotes.json
-    and provide methods to view the same
-    """
-
-    __singleton = None
-
-    def __init__(self, remotes_json_path: str,
-                 ap_repo: ap_git.GitRepo):
-        """
-        Initializes the VersionsFetcher instance
-        with a given remotes.json path.
-
-        Parameters:
-            remotes_json_path (str): Path to the remotes.json file.
-            ap_repo (GitRepo): ArduPilot local git repository. This local
-                               repository is shared between the VersionsFetcher
-                               and the APSourceMetadataFetcher.
-
-        Raises:
-            RuntimeError: If an instance of this class already exists,
-                          enforcing a singleton pattern.
-        """
-        if vehm.get_singleton() is None:
-            raise RuntimeError("VehiclesManager should be initialised first")
-
-        # Enforce singleton pattern by raising an error if
-        # an instance already exists.
-        if VersionsFetcher.__singleton:
-            raise RuntimeError("VersionsFetcher must be a singleton.")
-
-        self.logger = logging.getLogger(__name__)
-
-        self.__remotes_json_path = remotes_json_path
-        self.__ensure_remotes_json()
-        self.__access_lock_versions_metadata = Lock()
-        self.__versions_metadata = []
-        tasks = (
-            (self.fetch_ap_releases, 1200),
-            (self.fetch_whitelisted_tags, 1200),
-        )
-        self.__task__runner = TaskRunner(tasks=tasks)
-        self.repo = ap_repo
-        VersionsFetcher.__singleton = self
-
-    def start(self) -> None:
-        """
-        Start auto-fetch jobs.
-        """
-        self.logger.info(
-            "Starting VersionsFetcher background auto-fetch jobs."
-        )
-        self.__task__runner.start()
-
-    def stop(self) -> None:
-        """
-        Stop auto-fetch jobs.
-        """
-        self.logger.info(
-            "Stopping VersionsFetcher background auto-fetch jobs."
-        )
-        self.__task__runner.stop()
-
-    def get_all_remotes_info(self) -> list[RemoteInfo]:
-        """
-        Return the list of RemoteInfo objects constructed from the
-        information in the remotes.json file
-
-        Returns:
-            list: RemoteInfo objects for all remotes mentioned in remotes.json
-        """
-        return [
-            RemoteInfo(
-                name=remote.get('name', None),
-                url=remote.get('url', None)
-            )
-            for remote in self.__get_versions_metadata()
-        ]
-
-    def get_remote_info(self, remote_name: str) -> RemoteInfo:
-        """
-        Return the RemoteInfo for the given remote name, None otherwise.
-
-        Returns:
-            RemoteInfo: The remote information object.
-        """
-        return next(
-            (
-                remote for remote in self.get_all_remotes_info()
-                if remote.name == remote_name
-            ),
-            None
-        )
-
-    def get_versions_for_vehicle(self, vehicle_id: str) -> list[VersionInfo]:
-        """
-        Return the list of dictionaries containing the info about the
-        versions listed to be built for a particular vehicle.
-
-        Parameters:
-            vehicle_id (str): the vehicle ID to fetch versions list for
-
-        Returns:
-            list: VersionInfo objects for all versions allowed to be
-                  built for the said vehicle.
-
-        """
-        if vehicle_id is None:
-            raise ValueError("Vehicle ID is a required parameter.")
-
-        vehicle = vehm.get_singleton().get_vehicle_by_id(vehicle_id)
-        if vehicle is None:
-            raise ValueError(f"Invalid vehicle ID '{vehicle_id}'.")
-
-        vehicle_name = vehicle.name
-
-        versions_list = []
-        for remote in self.__get_versions_metadata():
-            remote_info = RemoteInfo(
-                name=remote.get('name', None),
-                url=remote.get('url', None)
-            )
-            for vehicle in remote['vehicles']:
-                if vehicle['name'] != vehicle_name:
-                    continue
-
-                for release in vehicle['releases']:
-                    versions_list.append(VersionInfo(
-                        remote_info=remote_info,
-                        commit_ref=release.get('commit_reference', None),
-                        release_type=release.get('release_type', None),
-                        version_number=release.get('version_number', None),
-                        ap_build_artifacts_url=release.get(
-                            'ap_build_artifacts_url',
-                            None
-                        )
-                    ))
-        return versions_list
-
-    def is_version_listed(self, vehicle_id: str, version_id: str) -> bool:
-        """
-        Check if a version with given properties mentioned in remotes.json
-
-        Parameters:
-            vehicle_id (str): ID of the vehicle for which version is listed
-            version_id (str): version ID
-
-        Returns:
-            bool: True if the said version is mentioned in remotes.json,
-                  False otherwise
-
-        """
-        if vehicle_id is None:
-            raise ValueError("vehicle_id is a required parameter.")
-
-        if version_id is None:
-            raise ValueError("version_id is a required parameter.")
-
-        return version_id in [
-            version_info.version_id
-            for version_info in
-            self.get_versions_for_vehicle(vehicle_id=vehicle_id)
-        ]
-
-    def get_version_info(self, vehicle_id: str,
-                         version_id: str) -> VersionInfo:
-        """
-        Find first version matching the given properties in remotes.json
-
-        Parameters:
-            vehicle_id (str): ID of the vehicle for which version is listed
-            version_id (str): version ID
-
-        Returns:
-            VersionInfo: Object for the version matching the properties,
-                         None if not found
-
-        """
-        return next(
-            (
-                version
-                for version in self.get_versions_for_vehicle(
-                    vehicle_id=vehicle_id
-                )
-                if version.version_id == version_id
-            ),
-            None
-        )
-
-    def reload_remotes_json(self) -> None:
-        """
-        Read remotes.json, validate its structure against the schema
-        and cache it in memory
-        """
-        # load file containing vehicles listed to be built for each
-        # remote along with the branches/tags/commits on which the
-        # firmware can be built
-        remotes_json_schema_path = os.path.join(
-            os.path.dirname(__file__),
-            'remotes.schema.json'
-        )
-        with open(self.__remotes_json_path, 'r') as f, \
-             open(remotes_json_schema_path, 'r') as s:
-            f_content = f.read()
-
-            # Early return if file is empty
-            if not f_content:
-                return
-            versions_metadata = json.loads(f_content)
-            schema = json.loads(s.read())
-            # validate schema
-            jsonschema.validate(instance=versions_metadata, schema=schema)
-            self.__set_versions_metadata(versions_metadata=versions_metadata)
-
-        # update git repo with latest remotes list
-        self.__sync_remotes_with_ap_repo()
-
-    def __ensure_remotes_json(self) -> None:
-        """
-        Ensures remotes.json exists and is a valid JSON file.
-        """
-        p = Path(self.__remotes_json_path)
-
-        if not p.exists():
-            # Ensure parent directory exists
-            Path.mkdir(p.parent, parents=True, exist_ok=True)
-
-            # write empty json list
-            with open(p, 'w') as f:
-                f.write('[]')
-
-    def __set_versions_metadata(self, versions_metadata: list) -> None:
-        """
-        Set versions metadata property with the one passed as parameter
-        This requires to acquire the access lock to avoid overwriting the
-        object while it is being read
-        """
-        if versions_metadata is None:
-            raise ValueError("versions_metadata is a required parameter. "
-                             "Cannot be None.")
-
-        with self.__access_lock_versions_metadata:
-            self.__versions_metadata = versions_metadata
-
-    def __get_versions_metadata(self) -> list:
-        """
-        Read versions metadata property
-        This requires to acquire the access lock to avoid reading the list
-        while it is being modified
-
-        Returns:
-            list: the versions metadata list
-        """
-        with self.__access_lock_versions_metadata:
-            return self.__versions_metadata
-
-    def __sync_remotes_with_ap_repo(self):
-        """
-        Update the remotes in ArduPilot local repository with the latest
-        remotes list.
-        """
-        remotes = tuple(
-            (remote.name, remote.url)
-            for remote in self.get_all_remotes_info()
-        )
-        self.repo.remote_add_bulk(remotes=remotes, force=True)
-
-    def fetch_ap_releases(self) -> None:
-        """
-        Execute the fetch_releases.py script to update remotes.json
-        with Ardupilot's official releases
-        """
-        from scripts import fetch_releases
-        fetch_releases.run(
-            base_dir=os.path.join(
-                os.path.dirname(self.__remotes_json_path),
-                '..',
-            ),
-            remote_name="ardupilot",
-        )
-        self.reload_remotes_json()
-        return
-
-    def fetch_whitelisted_tags(self) -> None:
-        """
-        Execute the fetch_whitelisted_tags.py script to update
-        remotes.json with tags from whitelisted repos
-        """
-        from scripts import fetch_whitelisted_tags
-        fetch_whitelisted_tags.run(
-            base_dir=os.path.join(
-                os.path.dirname(self.__remotes_json_path),
-                '..',
-            )
-        )
-        self.reload_remotes_json()
-        return
-
-    @staticmethod
-    def get_singleton():
-        return VersionsFetcher.__singleton

+ 23 - 0
metadata_manager/versions_manager/__init__.py

@@ -0,0 +1,23 @@
+from .manager import VersionsManager
+from .models import ForkRemoteSpec, RemoteInfo, VersionInfo
+from .providers import (
+    WhitelistedForkTagVersionsProvider,
+    ManifestJsonVersionsProvider,
+    RemotesJsonVersionsProvider,
+    VersionsProvider,
+    build_default_providers,
+    DEFAULT_WHITELISTED_FORK_REMOTES,
+)
+
+__all__ = [
+    "VersionsManager",
+    "ForkRemoteSpec",
+    "RemoteInfo",
+    "VersionInfo",
+    "VersionsProvider",
+    "ManifestJsonVersionsProvider",
+    "WhitelistedForkTagVersionsProvider",
+    "RemotesJsonVersionsProvider",
+    "build_default_providers",
+    "DEFAULT_WHITELISTED_FORK_REMOTES",
+]

+ 151 - 0
metadata_manager/versions_manager/manager.py

@@ -0,0 +1,151 @@
+import logging
+from pathlib import Path
+
+import ap_git
+from utils import TaskRunner
+
+from metadata_manager.versions_manager.models import RemoteInfo, VersionInfo
+from metadata_manager.versions_manager.providers import (
+    VersionsProvider,
+    build_default_providers,
+)
+from metadata_manager.vehicles_manager import VehiclesManager as vehm
+
+REMOTES_SCHEMA_PATH = Path(__file__).resolve().parent.parent / "remotes.schema.json"
+
+
+def _release_type_dedup_priority(release_type: str) -> int:
+    """Lower value wins when multiple releases share the same version_id."""
+    match release_type:
+        case "stable":
+            return 0
+        case "beta":
+            return 1
+        case "latest":
+            return 2
+        case "tag":
+            return 3
+        case _:
+            return 99
+
+
+class VersionsManager:
+    """
+    Version manager merges version metadata from different pluggable providers.
+    """
+
+    __singleton = None
+
+    def __init__(
+        self,
+        ap_repo: ap_git.GitRepo,
+        remotes_json_path: str,
+        manifest_json=None,
+        providers: list[VersionsProvider] | None = None,
+    ):
+        if vehm.get_singleton() is None:
+            raise RuntimeError("VehiclesManager should be initialised first")
+
+        if VersionsManager.__singleton:
+            raise RuntimeError("VersionsManager must be a singleton.")
+
+        self.logger = logging.getLogger(__name__)
+        if providers is None:
+            if manifest_json is None:
+                raise ValueError("manifest_json is required when providers is omitted")
+            providers = build_default_providers(
+                manifest_json=manifest_json,
+                remotes_json_path=remotes_json_path,
+                schema_path=str(REMOTES_SCHEMA_PATH),
+            )
+        self._providers = providers
+        self.repo = ap_repo
+        self._remotes_json_path = remotes_json_path
+        self.__task__runner = TaskRunner(tasks=((self.refresh_all, 1200),))
+        VersionsManager.__singleton = self
+
+    def start(self) -> None:
+        self.logger.info("Starting VersionsManager background refresh jobs.")
+        self.__task__runner.start()
+
+    def stop(self) -> None:
+        self.logger.info("Stopping VersionsManager background refresh jobs.")
+        self.__task__runner.stop()
+
+    def refresh_all(self) -> None:
+        for provider in self._providers:
+            try:
+                provider.refresh()
+            except Exception as exc:
+                self.logger.error(
+                    "Provider %s refresh failed: %s", provider.name, exc
+                )
+        self._sync_remotes_with_ap_repo()
+
+    def get_all_remotes_info(self) -> list[RemoteInfo]:
+        remotes = {}
+        for provider in self._providers:
+            for remote in provider.get_remotes():
+                remotes[remote.name] = remote
+        return list(remotes.values())
+
+    def get_remote_info(self, remote_name: str) -> RemoteInfo | None:
+        return next(
+            (remote for remote in self.get_all_remotes_info()
+             if remote.name == remote_name),
+            None,
+        )
+
+    def get_versions_for_vehicle(self, vehicle_id: str) -> list[VersionInfo]:
+        if vehicle_id is None:
+            raise ValueError("Vehicle ID is a required parameter.")
+
+        vehicle = vehm.get_singleton().get_vehicle_by_id(vehicle_id)
+        if vehicle is None:
+            raise ValueError(f"Invalid vehicle ID '{vehicle_id}'.")
+
+        by_version_id: dict[str, VersionInfo] = {}
+        for provider in self._providers:
+            for version in provider.get_versions(vehicle_id):
+                existing = by_version_id.get(version.version_id)
+                if existing is None:
+                    by_version_id[version.version_id] = version
+                    continue
+                if _release_type_dedup_priority(version.release_type) < (
+                    _release_type_dedup_priority(existing.release_type)
+                ):
+                    by_version_id[version.version_id] = version
+        return list(by_version_id.values())
+
+    def is_version_listed(self, vehicle_id: str, version_id: str) -> bool:
+        if vehicle_id is None:
+            raise ValueError("vehicle_id is a required parameter.")
+        if version_id is None:
+            raise ValueError("version_id is a required parameter.")
+
+        return version_id in [
+            version.version_id
+            for version in self.get_versions_for_vehicle(vehicle_id=vehicle_id)
+        ]
+
+    def get_version_info(self, vehicle_id: str, version_id: str) -> VersionInfo | None:
+        return next(
+            (
+                version
+                for version in self.get_versions_for_vehicle(vehicle_id=vehicle_id)
+                if version.version_id == version_id
+            ),
+            None,
+        )
+
+    def _sync_remotes_with_ap_repo(self) -> None:
+        remotes = tuple(
+            (remote.name, remote.url)
+            for remote in self.get_all_remotes_info()
+        )
+        if remotes:
+            self.repo.remote_add_bulk(remotes=remotes, force=True)
+
+    @staticmethod
+    def get_singleton():
+        return VersionsManager.__singleton

+ 70 - 0
metadata_manager/versions_manager/models.py

@@ -0,0 +1,70 @@
+import hashlib
+from dataclasses import dataclass
+
+from metadata_manager.firmware_server.models import ReleaseRecord
+
+
+@dataclass(frozen=True)
+class ForkRemoteSpec:
+    """GitHub owner and repo for a whitelisted fork."""
+
+    owner: str
+    repo: str
+
+    @property
+    def github_repo(self) -> str:
+        return f"{self.owner}/{self.repo}"
+
+    @property
+    def url(self) -> str:
+        return f"https://github.com/{self.github_repo}.git"
+
+
+@dataclass(frozen=True)
+class RemoteInfo:
+    name: str
+    url: str
+
+    def to_dict(self):
+        return {
+            "name": self.name,
+            "url": self.url,
+        }
+
+
+class VersionInfo:
+    """Version metadata exposed to the rest of the application."""
+
+    def __init__(
+        self,
+        remote_info: RemoteInfo,
+        commit_ref: str,
+        release_type: str,
+        version_number: str,
+        ap_build_artifacts_url,
+    ) -> None:
+        self.remote_info = remote_info
+        self.commit_ref = commit_ref
+        self.release_type = release_type
+        self.version_number = version_number
+        self.ap_build_artifacts_url = ap_build_artifacts_url
+
+        commit_ref_sanitized = commit_ref.replace("/", "-")
+        commit_ref_hash = hashlib.md5(commit_ref.encode()).hexdigest()[:8]
+        self.version_id = (
+            f"{remote_info.name}-{commit_ref_sanitized}-{commit_ref_hash}"
+        )
+
+    @classmethod
+    def from_release_record(
+        cls,
+        remote_info: RemoteInfo,
+        release: ReleaseRecord,
+    ) -> "VersionInfo":
+        return cls(
+            remote_info=remote_info,
+            commit_ref=release.commit_reference,
+            release_type=release.release_type,
+            version_number=release.version_number,
+            ap_build_artifacts_url=release.ap_build_artifacts_url,
+        )

+ 318 - 0
metadata_manager/versions_manager/providers.py

@@ -0,0 +1,318 @@
+import json
+import logging
+import os
+from abc import ABC, abstractmethod
+from pathlib import Path
+from threading import Lock
+from typing import Optional
+
+import jsonschema
+import requests
+
+from metadata_manager.vehicles_manager import DEFAULT_VEHICLES
+
+from metadata_manager.firmware_server import ManifestJSON
+from metadata_manager.firmware_server.models import ReleaseRecord
+
+from .models import ForkRemoteSpec, RemoteInfo, VersionInfo
+
+OFFICIAL_REMOTE_NAME = "ardupilot"
+OFFICIAL_REMOTE_URL = "https://github.com/ardupilot/ardupilot.git"
+FIRMWARE_SERVER_BASE = "https://firmware.ardupilot.org"
+
+# CBS vehicle id -> firmware.ardupilot.org top-level directory.
+_FIRMWARE_SERVER_DIR_BY_VEHICLE_ID = {
+    "copter": "Copter",
+    "plane": "Plane",
+    "rover": "Rover",
+    "sub": "Sub",
+    "heli": "Copter",
+    "blimp": "Blimp",
+    "tracker": "AntennaTracker",
+    "ap-periph": "AP_Periph",
+}
+
+
+def _firmware_server_dir(vehicle_id: str) -> str:
+    try:
+        return _FIRMWARE_SERVER_DIR_BY_VEHICLE_ID[vehicle_id]
+    except KeyError as exc:
+        raise ValueError(f"Unknown vehicle id: {vehicle_id}") from exc
+
+
+def _vehicle_id_for_tag_segment(segment: str) -> Optional[str]:
+    for vehicle in DEFAULT_VEHICLES:
+        if vehicle.name == segment or vehicle.id == segment:
+            return vehicle.id
+    return None
+
+
+DEFAULT_WHITELISTED_FORK_REMOTES = [
+    ForkRemoteSpec(owner="tridge", repo="ardupilot"),
+    ForkRemoteSpec(owner="peterbarker", repo="ardupilot"),
+    ForkRemoteSpec(owner="rmackay9", repo="rmackay9-ardupilot"),
+    ForkRemoteSpec(owner="shiv-tyagi", repo="ardupilot"),
+    ForkRemoteSpec(owner="andyp1per", repo="ardupilot"),
+]
+
+
+class VersionsProvider(ABC):
+    @property
+    @abstractmethod
+    def name(self) -> str:
+        ...
+
+    @abstractmethod
+    def get_versions(self, vehicle_id: str) -> list[VersionInfo]:
+        ...
+
+    @abstractmethod
+    def get_remotes(self) -> list[RemoteInfo]:
+        ...
+
+    def refresh(self) -> None:
+        """Optional hook to refresh provider data."""
+
+    @property
+    def is_available(self) -> bool:
+        return True
+
+
+class ManifestJsonVersionsProvider(VersionsProvider):
+    name = "official-ardupilot"
+
+    def __init__(self, manifest_json: ManifestJSON):
+        self._manifest_json = manifest_json
+        self._remote = RemoteInfo(OFFICIAL_REMOTE_NAME, OFFICIAL_REMOTE_URL)
+
+    @property
+    def is_available(self) -> bool:
+        return self._manifest_json.is_available
+
+    def refresh(self) -> None:
+        self._manifest_json.refresh()
+
+    def get_remotes(self) -> list[RemoteInfo]:
+        if not self.is_available:
+            return []
+        return [self._remote]
+
+    def get_versions(self, vehicle_id: str) -> list[VersionInfo]:
+        if not self.is_available:
+            return []
+
+        return [
+            VersionInfo.from_release_record(self._remote, release)
+            for release in self._manifest_json.get_releases(vehicle_id)
+        ]
+
+
+class WhitelistedForkTagVersionsProvider(VersionsProvider):
+    """
+    Discover custom-build tags from whitelisted GitHub forks.
+
+    Only tags with the ``custom-build/`` prefix are included. Vehicle-specific
+    tag segments may use a vehicle name (e.g. ``Copter``) or id (e.g. ``copter``).
+
+    Tag format examples:
+    - ``custom-build/my-feature`` — listed for all vehicles
+    - ``custom-build/Copter/my-feature`` — listed for Copter only
+    - ``custom-build/copter/my-feature`` — listed for Copter only
+    """
+
+    name = "whitelisted-fork-tags"
+
+    def __init__(
+        self,
+        fork_remotes: Optional[list[ForkRemoteSpec]] = None,
+    ):
+        self.logger = logging.getLogger(__name__)
+        self._fork_remotes = fork_remotes or DEFAULT_WHITELISTED_FORK_REMOTES
+        self._vehicle_ids = [v.id for v in DEFAULT_VEHICLES]
+        self._versions_by_remote_vehicle: dict[str, dict[str, list[ReleaseRecord]]] = {}
+        self._available = False
+
+    @property
+    def is_available(self) -> bool:
+        return self._available
+
+    def get_remotes(self) -> list[RemoteInfo]:
+        if not self.is_available:
+            return []
+        fork_by_owner = {spec.owner: spec for spec in self._fork_remotes}
+        return [
+            RemoteInfo(name=remote_name, url=fork_by_owner[remote_name].url)
+            for remote_name in self._versions_by_remote_vehicle
+            if remote_name != OFFICIAL_REMOTE_NAME and remote_name in fork_by_owner
+        ]
+
+    def refresh(self) -> None:
+        versions_map = {
+            spec.owner: {vehicle_id: [] for vehicle_id in self._vehicle_ids}
+            for spec in self._fork_remotes
+        }
+
+        for spec in self._fork_remotes:
+            if spec.owner == OFFICIAL_REMOTE_NAME:
+                continue
+            try:
+                tag_objs = self._fetch_tags_from_github(spec.github_repo)
+            except Exception as exc:
+                self.logger.warning(
+                    "Skipping remote %s (%s): %s",
+                    spec.owner,
+                    spec.github_repo,
+                    exc,
+                )
+                continue
+
+            for tag_info in tag_objs:
+                ref = tag_info["ref"].replace("refs/tags/", "")
+                parts = ref.split("/", 3)
+                if len(parts) <= 1 or parts[0] != "custom-build":
+                    continue
+
+                tag_vehicle_id = _vehicle_id_for_tag_segment(parts[1])
+                if tag_vehicle_id is not None:
+                    if len(parts) == 2:
+                        continue
+                    vehicles_for_tag = [tag_vehicle_id]
+                else:
+                    vehicles_for_tag = self._vehicle_ids
+
+                for vid in vehicles_for_tag:
+                    versions_map[spec.owner][vid].append(
+                        ReleaseRecord(
+                            vehicle_id=vid,
+                            release_type="tag",
+                            version_number=parts[-1],
+                            commit_reference=tag_info["object"]["sha"],
+                            ap_build_artifacts_url=(
+                                f"{FIRMWARE_SERVER_BASE}/"
+                                f"{_firmware_server_dir(vid)}/latest"
+                            ),
+                        )
+                    )
+
+        self._versions_by_remote_vehicle = versions_map
+        self._available = True
+
+    def get_versions(self, vehicle_id: str) -> list[VersionInfo]:
+        if not self.is_available:
+            return []
+
+        fork_by_owner = {spec.owner: spec for spec in self._fork_remotes}
+        versions = []
+        for remote_name, vehicles_map in self._versions_by_remote_vehicle.items():
+            if remote_name == OFFICIAL_REMOTE_NAME:
+                continue
+            spec = fork_by_owner.get(remote_name)
+            if spec is None:
+                continue
+            remote = RemoteInfo(name=spec.owner, url=spec.url)
+            for release in vehicles_map.get(vehicle_id, []):
+                versions.append(VersionInfo.from_release_record(remote, release))
+        return versions
+
+    def _fetch_tags_from_github(self, github_repo: str) -> list:
+        url = f"https://api.github.com/repos/{github_repo}/git/refs/tags"
+        headers = {
+            "X-GitHub-Api-Version": "2022-11-28",
+            "Accept": "application/vnd.github+json",
+        }
+        token = os.getenv("CBS_GITHUB_ACCESS_TOKEN")
+        if token:
+            headers["Authorization"] = f"Bearer {token}"
+
+        response = requests.get(url=url, headers=headers, timeout=60)
+        response.raise_for_status()
+        return response.json()
+
+
+class RemotesJsonVersionsProvider(VersionsProvider):
+    name = "remotes-json"
+
+    def __init__(self, remotes_json_path: str, schema_path: str):
+        self._remotes_json_path = remotes_json_path
+        self._schema_path = schema_path
+        self._lock = Lock()
+        self._metadata: list = []
+
+    @property
+    def is_available(self) -> bool:
+        return bool(self._metadata)
+
+    def refresh(self) -> None:
+        self.reload()
+
+    def reload(self) -> None:
+        path = Path(self._remotes_json_path)
+        if not path.is_file():
+            with self._lock:
+                self._metadata = []
+            return
+
+        content = path.read_text(encoding="utf-8")
+        if not content.strip():
+            with self._lock:
+                self._metadata = []
+            return
+
+        metadata = json.loads(content)
+        schema = json.loads(Path(self._schema_path).read_text(encoding="utf-8"))
+        jsonschema.validate(instance=metadata, schema=schema)
+
+        with self._lock:
+            self._metadata = metadata
+
+    def get_remotes(self) -> list[RemoteInfo]:
+        with self._lock:
+            metadata = list(self._metadata)
+        return [
+            RemoteInfo(name=remote.get("name"), url=remote.get("url"))
+            for remote in metadata
+            if remote.get("name") and remote.get("url")
+        ]
+
+    def get_versions(self, vehicle_id: str) -> list[VersionInfo]:
+        versions = []
+        with self._lock:
+            metadata = list(self._metadata)
+
+        for remote in metadata:
+            remote_info = RemoteInfo(
+                name=remote.get("name"),
+                url=remote.get("url"),
+            )
+            for remote_vehicle in remote.get("vehicles", []):
+                if remote_vehicle.get("id") != vehicle_id:
+                    continue
+                for release in remote_vehicle.get("releases", []):
+                    versions.append(
+                        VersionInfo(
+                            remote_info=remote_info,
+                            commit_ref=release.get("commit_reference"),
+                            release_type=release.get("release_type"),
+                            version_number=release.get("version_number"),
+                            ap_build_artifacts_url=release.get(
+                                "ap_build_artifacts_url"
+                            ),
+                        )
+                    )
+        return versions
+
+
+def build_default_providers(
+    manifest_json: ManifestJSON,
+    remotes_json_path: str,
+    schema_path: str,
+    fork_remotes: Optional[list[ForkRemoteSpec]] = None,
+) -> list[VersionsProvider]:
+    return [
+        ManifestJsonVersionsProvider(manifest_json),
+        WhitelistedForkTagVersionsProvider(fork_remotes=fork_remotes),
+        RemotesJsonVersionsProvider(
+            remotes_json_path=remotes_json_path,
+            schema_path=schema_path,
+        ),
+    ]