Skip to content

native_fastsurfer

tit.pre.native_fastsurfer

Project-scoped mailbox client for the optional macOS FastSurfer worker.

run_native_fastsurfer

run_native_fastsurfer(project_dir: str, subject_id: str, input_path: str | Path, *, threads: int, logger: Logger, stop_event: Event) -> bool

Run via an enabled host worker; return False only when it is not enabled.

Only project-relative input and a bounded thread count cross the mailbox. Cancellation and host disconnection fail the job instead of retrying on CPU.

Source code in tit/pre/native_fastsurfer.py
def run_native_fastsurfer(
    project_dir: str,
    subject_id: str,
    input_path: str | Path,
    *,
    threads: int,
    logger: logging.Logger,
    stop_event: threading.Event,
) -> bool:
    """Run via an enabled host worker; return False only when it is not enabled.

    Only project-relative input and a bounded thread count cross the mailbox.
    Cancellation and host disconnection fail the job instead of retrying on CPU.
    """
    root = Path(project_dir).resolve()
    mailbox = _inside(root, root / "code/ti-toolbox/native-fastsurfer")
    session = _session(root, mailbox, logger=logger)
    if session is None:
        return False
    if not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9_-]{0,127}", subject_id):
        raise PreprocessError("Invalid native FastSurfer subject identifier.")
    if (
        isinstance(threads, bool)
        or not isinstance(threads, int)
        or not 1 <= threads <= 1024
    ):
        raise PreprocessError("Native FastSurfer threads must be between 1 and 1024.")
    candidate = Path(input_path)
    source = _inside(root, candidate if candidate.is_absolute() else root / candidate)
    if not source.is_file() or not source.name.lower().endswith((".nii", ".nii.gz")):
        raise PreprocessError(
            "Native FastSurfer requires an existing project NIfTI input."
        )
    if stop_event.is_set():
        raise PreprocessCancelled(
            "Pre-processing cancelled before native command start."
        )
    request_id = str(uuid.uuid4())
    requests = _inside(root, mailbox / "requests")
    requests.mkdir(parents=True, exist_ok=True)
    directory = _inside(root, requests / request_id)
    directory.mkdir()
    request = {
        "version": 1,
        "session": session,
        "id": request_id,
        "subject_id": subject_id,
        "threads": threads,
        "input_path": source.relative_to(root).as_posix(),
    }
    completed = False
    offset = 0
    heartbeat_at = 0.0
    try:
        temporary = directory / "request.tmp"
        temporary.write_text(json.dumps(request))
        (directory / "heartbeat.json").write_text(
            json.dumps({"lastSeen": time.time() * 1000})
        )
        heartbeat_at = time.monotonic()
        temporary.replace(directory / "request.json")
        logger.info("Running native FastSurfer on Apple GPU (CPU view aggregation).")
        while True:
            if time.monotonic() - heartbeat_at >= 2 or heartbeat_at == 0:
                heartbeat = _inside(root, directory / "heartbeat.tmp")
                heartbeat.write_text(json.dumps({"lastSeen": time.time() * 1000}))
                heartbeat.replace(_inside(root, directory / "heartbeat.json"))
                heartbeat_at = time.monotonic()
            if stop_event.is_set():
                raise PreprocessCancelled("Native FastSurfer cancelled.")
            log_path = _inside(root, directory / "stdout.log")
            if log_path.is_file():
                with log_path.open("rb") as stream:
                    stream.seek(offset)
                    chunk = stream.read(262144)
                    offset = stream.tell()
                if chunk:
                    logger.info("%s", chunk.decode("utf-8", errors="replace").rstrip())
            result_path = _inside(root, directory / "result.json")
            if result_path.exists():
                result = _read_json(result_path)
                if result.get("ok") is not True:
                    raise PreprocessError(
                        f"Native FastSurfer failed: {str(result.get('error', 'invalid result'))[:2000]}"
                    )
                completed = True
                return True
            if _session(root, mailbox) != session:
                raise PreprocessError(
                    "Native FastSurfer host session ended or changed."
                )
            stop_event.wait(_POLL_SECONDS)
    finally:
        if not completed:
            _inside(root, directory / "cancel").touch()