IQ.Pilot Release Commit @ 8d3c939

This commit is contained in:
IQ.Lvbs CI [bot]
2026-08-30 22:55:58 -05:00
parent c9629a3603
commit 90015b2835
57 changed files with 776 additions and 98 deletions

View File

@@ -0,0 +1,209 @@
"""
Copyright © IQ.Lvbs, apart of Project Teal Lvbs, All Rights Reserved, licensed under https://konn3kt.com/tos/
"""
from __future__ import annotations
import hashlib
import json
import os
from iqpilot.common.swaglog import cloudlog
import urllib.request
from pathlib import Path
from iqpilot.system.hardware.usb import egpu_dock_ready
USB_SYSFS_ROOT = "/sys/bus/usb/devices"
FIRMWARE_MIRROR = os.getenv("IQ_EGPU_FIRMWARE_MIRROR", "/data/firmware/tinygrad")
TINYGRAD_CACHE = "/data/.cache"
COMMA_LFS_BATCH_URL = "https://gitlab.com/commaai/openpilot-lfs.git/info/lfs/objects/batch"
DOWNLOAD_CHUNK = 4 * 1024 * 1024
def usbgpu_present(sysfs_root: str = USB_SYSFS_ROOT) -> bool:
return egpu_dock_ready(Path(sysfs_root))
def egpu_present_consented(params, sysfs_root: str = USB_SYSFS_ROOT) -> bool:
try:
if params is not None and params.get_bool("IQEgpuDisabled"):
return False
except Exception:
pass
return usbgpu_present(sysfs_root)
def egpu_selected(params, sysfs_root: str = USB_SYSFS_ROOT) -> bool:
try:
if params is not None and params.get_bool("IQEgpuDisabled"):
return False
if params is not None and params.get_bool("IQEgpuEnabled"):
return True
except Exception:
pass
return usbgpu_present(sysfs_root)
def resolve_backend(emac_enabled: bool, egpu_enabled: bool, egpu_present: bool = False) -> str | None:
if egpu_present:
return "egpu"
if emac_enabled:
return "emac"
if egpu_enabled:
return "egpu"
return None
def egpu_pkl_path(meta: dict) -> str:
from iqpilot.system.hardware.hw import Paths
return os.path.join(Paths.model_root(), f"egpu_{meta['key']}_{meta['sha256'][:8]}_amd_tinygrad.pkl")
def egpu_policy_pkl_path(meta: dict) -> str:
from iqpilot.system.hardware.hw import Paths
return os.path.join(Paths.model_root(), f"egpu_{meta['key']}_{meta['sha256'][:8]}_amd_policy.pkl")
def egpu_oob_pkl_path(meta: dict) -> str:
from iqpilot.system.hardware.hw import Paths
return os.path.join(Paths.model_root(), f"egpu_{meta['key']}_{meta['sha256'][:8]}_amd_policy_oob.pkl")
def onnx_cache_path(meta: dict) -> str:
from iqpilot.system.hardware.hw import Paths
return os.path.join(Paths.model_root(), f"{meta['model_name']}_{meta['sha256'][:8]}.onnx")
def _sha256_file(path: str) -> str:
digest = hashlib.sha256()
with open(path, "rb") as f:
while chunk := f.read(DOWNLOAD_CHUNK):
digest.update(chunk)
return digest.hexdigest()
def quarantine_artifact(path: str, why: str) -> None:
try:
if os.path.isfile(path):
os.replace(path, path + ".unusable")
except OSError:
try:
os.remove(path)
except OSError:
pass
def local_onnx(meta: dict) -> str | None:
path = onnx_cache_path(meta)
if not os.path.isfile(path):
return None
size = int(meta.get("download", {}).get("size", 0))
if size and os.path.getsize(path) != size:
quarantine_artifact(path, "onnx size mismatch")
return None
if _sha256_file(path) != meta["sha256"]:
quarantine_artifact(path, "onnx sha256 mismatch")
return None
return path
def resolve_download_url(download_url: str, sha256: str, size: int, timeout: float = 30.0) -> str:
if download_url.startswith("commalfs:"):
oid = download_url.split(":", 1)[1]
body = json.dumps({"operation": "download", "transfers": ["basic"],
"objects": [{"oid": oid, "size": size}]}).encode()
req = urllib.request.Request(COMMA_LFS_BATCH_URL, data=body, headers={
"Accept": "application/vnd.git-lfs+json", "Content-Type": "application/vnd.git-lfs+json"})
with urllib.request.urlopen(req, timeout=timeout) as r:
d = json.load(r)
return d["objects"][0]["actions"]["download"]["href"]
return download_url
def download_onnx(meta: dict, progress_cb=None) -> str:
from iqpilot.selfdrive.iqmodeld.egpu_model import download_descriptor
download_url, size = download_descriptor(meta)
if not download_url:
raise RuntimeError(f"model {meta['key']} has no download source; stage the onnx at {onnx_cache_path(meta)}")
path = onnx_cache_path(meta)
os.makedirs(os.path.dirname(path), exist_ok=True)
try:
from iqpilot.selfdrive.iqmodeld.model_bundle_downloader import download_hf_file
return download_hf_file(f"onnx/{meta['sha256']}.onnx", path, meta["sha256"], int(size or 0), progress_cb=progress_cb)
except Exception as e:
cloudlog.warning(f"onnx {meta['key']} unavailable from HF ({e}); falling back to {download_url.split(':', 1)[0]}")
url = resolve_download_url(download_url, meta["sha256"], size)
tmp = path + ".part"
digest = hashlib.sha256()
got = 0
with urllib.request.urlopen(url, timeout=60) as r, open(tmp, "wb") as f:
while chunk := r.read(DOWNLOAD_CHUNK):
f.write(chunk)
digest.update(chunk)
got += len(chunk)
if progress_cb is not None and size:
progress_cb(got / size)
if size and got != size:
os.remove(tmp)
raise RuntimeError(f"onnx download truncated: {got}/{size} bytes")
if digest.hexdigest() != meta["sha256"]:
os.remove(tmp)
raise RuntimeError(f"onnx sha256 mismatch for {meta['key']}")
os.replace(tmp, path)
return path
def download_precompiled(meta: dict, progress_cb=None, policy: bool = False, oob: bool = False) -> str | None:
field = "egpu_oob_artifact" if oob else "egpu_policy_artifact" if policy else "egpu_artifact"
art = meta.get(field)
if not art or not (art.get("objects") or art.get("hf_path")):
return None
from iqpilot.selfdrive.iqmodeld.model_bundle_downloader import download_hf_file, download_lfs_bundle
dest = egpu_oob_pkl_path(meta) if oob else egpu_policy_pkl_path(meta) if policy else egpu_pkl_path(meta)
if art.get("hf_path"):
try:
return download_hf_file(art["hf_path"], dest, art["sha256"], int(art.get("size", 0)), progress_cb=progress_cb)
except Exception as e:
cloudlog.warning(f"precompiled {meta['key']} unavailable from HF ({e}); trying LFS")
if not art.get("objects"):
raise
return download_lfs_bundle(art["objects"], dest, art["sha256"], int(art.get("size", 0)), progress_cb=progress_cb)
def patch_tinygrad_fetch_fw() -> None:
import pathlib
import zstandard
from tinygrad import helpers
if getattr(helpers.fetch_fw, "_iq_patched", False):
return
_orig = helpers.fetch_fw
def fetch_fw(path, name, sha256):
mirror = pathlib.Path(FIRMWARE_MIRROR) / path / name
if mirror.is_file():
blob = mirror.read_bytes()
if hashlib.sha256(blob).hexdigest() == sha256:
return blob
p = pathlib.Path(f"/lib/firmware/{path}/{name}.zst")
if p.is_file():
blob = zstandard.ZstdDecompressor().stream_reader(p.read_bytes()).read()
if hashlib.sha256(blob).hexdigest() == sha256:
return blob
blob = _orig(path, name, sha256)
# The dock's GPU firmware otherwise lives only in tinygrad's per-user download cache, which is
# a network fetch the first time a new HOME sees it; onroad the car is usually offline.
try:
mirror.parent.mkdir(parents=True, exist_ok=True)
tmp = mirror.with_suffix(mirror.suffix + ".part")
tmp.write_bytes(blob)
os.replace(tmp, mirror)
except OSError:
pass
return blob
fetch_fw._iq_patched = True
helpers.fetch_fw = fetch_fw

View File

@@ -0,0 +1,9 @@
"""
Copyright © IQ.Lvbs, apart of Project Teal Lvbs, All Rights Reserved, licensed under https://konn3kt.com/tos/
"""
from iqpilot._proprietary_loader import ProprietaryModuleMissing, load_private_module
try:
load_private_module(__name__, "iqpilot_private.models.egpu_model")
except ProprietaryModuleMissing:
from iqpilot.models_private_src.egpu_model import *

View File

@@ -0,0 +1,211 @@
"""
Copyright © IQ.Lvbs, apart of Project Teal Lvbs, All Rights Reserved, licensed under https://konn3kt.com/tos/
"""
from __future__ import annotations
import hashlib
import json
import os
MODELS_BASE_URLS = (
"https://git.konn3kt.com/teal/IQModels/raw/branch/main",
"https://gitlvb.teallvbs.xyz/teal/IQModels/raw/branch/main",
)
CHUNK = 4 * 1024 * 1024
HTTP_TIMEOUT_S = 60.0
STREAM_RETRIES = 6
def _requests_auth():
import importlib
for mod in ("iqpilot_private.models.git_auth", "iqpilot.models_private_src.git_auth",
"iqpilot.selfdrive.iqmodeld.models.git_auth"):
try:
return importlib.import_module(mod).get_requests_auth()
except Exception:
continue
return None
def _hf():
import importlib
for mod in ("iqpilot_private.models.git_auth", "iqpilot.selfdrive.iqmodeld.models.git_auth"):
try:
m = importlib.import_module(mod)
return m.get_hf_headers(), m.hf_resolve_url
except Exception:
continue
return None, None
def download_hf_file(hf_path: str, dst: str, sha256: str, size: int, progress_cb=None) -> str:
import requests
headers, resolve = _hf()
if resolve is None:
raise RuntimeError("no HF credentials available")
url = resolve(hf_path)
os.makedirs(os.path.dirname(dst), exist_ok=True)
tmp = dst + ".hfpart"
last_error: Exception | None = None
for _attempt in range(STREAM_RETRIES):
try:
have = os.path.getsize(tmp) if os.path.isfile(tmp) else 0
if size and have > size:
os.remove(tmp)
have = 0
if not size or have < size:
req_headers = dict(headers)
if have:
req_headers["Range"] = f"bytes={have}-"
with requests.get(url, headers=req_headers, stream=True, timeout=HTTP_TIMEOUT_S, allow_redirects=True) as r:
r.raise_for_status()
if have and r.status_code != 206:
have = 0
with open(tmp, "ab" if have else "wb") as f:
got = have
for chunk in r.iter_content(CHUNK):
f.write(chunk)
got += len(chunk)
if progress_cb is not None and size:
progress_cb(min(1.0, got / size))
digest = hashlib.sha256()
with open(tmp, "rb") as f:
for chunk in iter(lambda: f.read(CHUNK), b""):
digest.update(chunk)
if size and os.path.getsize(tmp) != size:
raise RuntimeError(f"size mismatch: {os.path.getsize(tmp)}/{size} bytes")
if sha256 and digest.hexdigest() != sha256:
os.remove(tmp)
raise RuntimeError("sha256 mismatch")
os.replace(tmp, dst)
return dst
except Exception as e:
last_error = e
raise RuntimeError(f"HF download failed: {last_error}")
def _lfs_endpoint(base_url: str) -> str:
return base_url.split("/raw/", 1)[0] + ".git/info/lfs"
def _resolve_oid(session, base_url: str, oid: str, size: int, auth):
import requests
batch = session.post(f"{_lfs_endpoint(base_url)}/objects/batch",
data=json.dumps({"operation": "download", "transfers": ["basic"],
"objects": [{"oid": oid, "size": size}]}),
headers={"Content-Type": "application/vnd.git-lfs+json",
"Accept": "application/vnd.git-lfs+json"},
auth=auth, timeout=HTTP_TIMEOUT_S)
batch.raise_for_status()
entry = batch.json()["objects"][0]
if "actions" not in entry:
raise requests.RequestException(f"LFS object unavailable: {entry.get('error', oid)}")
action = entry["actions"]["download"]
return action["href"], action.get("header", {})
def _part_path(dst: str, oid: str) -> str:
return os.path.join(dst + ".parts", oid)
def _part_complete(path: str, oid: str, size: int) -> bool:
if not os.path.isfile(path) or os.path.getsize(path) != size:
return False
digest = hashlib.sha256()
with open(path, "rb") as f:
for chunk in iter(lambda: f.read(CHUNK), b""):
digest.update(chunk)
return digest.hexdigest() == oid
def _fetch_part(session, base_url: str, obj: dict, path: str, auth, progress) -> None:
size = int(obj["size"])
have = os.path.getsize(path) if os.path.isfile(path) else 0
if have > size:
os.remove(path)
have = 0
href, headers = _resolve_oid(session, base_url, obj["oid"], size, auth)
obj_auth = None if headers.get("Authorization") else auth
# LFS parts are content-addressed (oid == sha256), so a half-written part can be resumed with a
# Range request and verified afterwards instead of being thrown away on every restart.
if have:
headers = {**headers, "Range": f"bytes={have}-"}
with session.get(href, headers=headers, stream=True, timeout=HTTP_TIMEOUT_S, auth=obj_auth) as r:
r.raise_for_status()
if have and r.status_code != 206:
have = 0
with open(path, "ab" if have else "wb") as f:
for chunk in r.iter_content(CHUNK):
f.write(chunk)
progress(len(chunk))
def download_lfs_bundle(objects: list, dst: str, sha256: str, size: int, progress_cb=None) -> str:
import requests
auth = _requests_auth()
session = requests.Session()
os.makedirs(dst + ".parts", exist_ok=True)
total = int(size) or sum(int(o["size"]) for o in objects)
done_bytes = sum(int(o["size"]) for o in objects if _part_complete(_part_path(dst, o["oid"]), o["oid"], int(o["size"])))
got = [done_bytes]
def progress(n: int) -> None:
got[0] += n
if progress_cb is not None and total:
progress_cb(min(1.0, got[0] / total))
last_error: Exception | None = None
for base_url in MODELS_BASE_URLS:
for _attempt in range(STREAM_RETRIES):
try:
for obj in objects:
path = _part_path(dst, obj["oid"])
if _part_complete(path, obj["oid"], int(obj["size"])):
continue
got[0] = done_bytes
_fetch_part(session, base_url, obj, path, auth, progress)
if not _part_complete(path, obj["oid"], int(obj["size"])):
if os.path.getsize(path) >= int(obj["size"]):
os.remove(path)
raise RuntimeError(f"part {obj['oid'][:12]} incomplete or failed verification")
done_bytes += int(obj["size"])
got[0] = done_bytes
break
except Exception as e:
last_error = e
else:
continue
break
else:
raise RuntimeError(f"model bundle download failed: {last_error}")
tmp = dst + ".part"
digest = hashlib.sha256()
with open(tmp, "wb") as out:
for obj in objects:
with open(_part_path(dst, obj["oid"]), "rb") as f:
for chunk in iter(lambda: f.read(CHUNK), b""):
out.write(chunk)
digest.update(chunk)
if total and os.path.getsize(tmp) != total:
os.remove(tmp)
raise RuntimeError(f"size mismatch: {os.path.getsize(tmp) if os.path.exists(tmp) else 0}/{total} bytes")
if sha256 and digest.hexdigest() != sha256:
os.remove(tmp)
for obj in objects:
try:
os.remove(_part_path(dst, obj["oid"]))
except OSError:
pass
raise RuntimeError("sha256 mismatch")
os.replace(tmp, dst)
for obj in objects:
try:
os.remove(_part_path(dst, obj["oid"]))
except OSError:
pass
try:
os.rmdir(dst + ".parts")
except OSError:
pass
return dst

View File

@@ -253,6 +253,15 @@ EVENTS: dict[int, dict[str, Alert | AlertCallbackType]] = {
"Ensure road ahead is clear"),
},
EventName.bigModelLoading: {
ET.NO_ENTRY: NoEntryAlert("Big Model Loading"),
},
EventName.bigModelFailed: {
ET.SOFT_DISABLE: soft_disable_alert("Big Model Failed"),
ET.PERMANENT: NormalPermanentAlert("Big Model Failed ", "Restart the car to retry,\nsmall model is still available", duration=20.),
},
EventName.lateralManeuver: {
ET.WARNING: longitudinal_maneuver_alert,
ET.PERMANENT: NormalPermanentAlert("Lateral Maneuver Mode"),
@@ -366,7 +375,7 @@ EVENTS: dict[int, dict[str, Alert | AlertCallbackType]] = {
"Pay Attention",
"",
AlertStatus.normal, AlertSize.small,
Priority.LOW, VisualAlert.none, AudibleAlert.none, .1),
Priority.MID, VisualAlert.none, AudibleAlert.none, .1),
},
EventName.promptDriverDistracted: {
@@ -892,7 +901,7 @@ if HARDWARE.get_device_type() == 'mici':
"Pay Attention",
"",
AlertStatus.normal, AlertSize.small,
Priority.LOW, VisualAlert.none, AudibleAlert.none, 2),
Priority.MID, VisualAlert.none, AudibleAlert.none, 2),
},
EventName.promptDriverDistracted: {
ET.PERMANENT: Alert(