Source code for nornir_buildmanager.build

"""

*Note*: Certain arguments support regular expressions.  See the python :py:mod:`re` module for instructions on how to construct appropriate regular expressions.

.. argparse::
   :module: nornir_buildmanager.build
   :func: BuildParserRoot
   :prog: nornir_build volumepath

"""

import argparse
import logging
import os
import sys

import matplotlib

import nornir_buildmanager.volumemanager.volumemanager
import nornir_buildmanager.pipelinemanager as pipelinemanager
import nornir_imageregistration
from nornir_imageregistration.headless import is_headless

# Nornir build must use a backend that does not allocate windows in the GUI.
# Otherwise bugs will appear in multi-threaded environments.
#
# Importing this module used to switch the process to QtAgg unconditionally, undoing the
# Agg that nornir_imageregistration selects for headless runs. Any later plt.show()
# then blocked on a Qt event loop waiting for a window nobody could close. Under pytest
# that wedged whole sessions: collection imports every test module, so one test that
# imports nornir_buildmanager made the entire run GUI-backed, and a later plotting test
# hung with no output. The same files passed individually because nothing had pulled the
# GUI backend in.
if is_headless():
    matplotlib.use('Agg')
elif 'DEBUG' not in os.environ:
    try:
        matplotlib.use('QtAgg')
    except ImportError:
        matplotlib.use('Agg')  # Fallback to non-interactive backend

import matplotlib.pyplot as plt

plt.ioff()

import nornir_buildmanager
from nornir_buildmanager import *
from nornir_shared.misc import SetupLogging, lowpriority
from nornir_shared.tasktimer import TaskTimer
import pkgutil

import nornir_shared.prettyoutput as prettyoutput

CommandParserDict = {}


def _AddParserRootArguments(parser: argparse.ArgumentParser):
    """Add global flags shared by all build subcommands."""
    parser.add_argument('-debug',
                        action='store_true',
                        required=False,
                        default=False,
                        help='If true any exceptions raised by pipelines are not handled.',
                        dest='debug')

    parser.add_argument('-lowpriority', '-lp',
                        action='store_true',
                        required=False,
                        default=False,
                        help='Run the build with lower priority.  The machine may be more responsive at the expense of much slower builds. 3x-5x slower in tests.',
                        dest='lowpriority')

    parser.add_argument('-verbose',
                        action='store_true',
                        required=False,
                        default=False,
                        help='Log per-iterate pipeline progress (Iterate/Mapping dumps). '
                             'Same as NORNIR_LOG_PIPELINE_PROGRESS=1. Independent of -debug.',
                        dest='verbose')

    parser.add_argument('-computational_library',
                        choices=['cupy', 'numpy', 'detect'],
                        required=False,
                        default='detect',
                        help='If not specified, cupy will be used if a nVidia GPU is present.  Otherwise, force cupy (GPU) or numpy (CPU) use.',
                        dest='computational_library')

    parser.add_argument('-no-delete',
                        action='store_true',
                        required=False,
                        default=False,
                        help='Skip all destructive Clean() and file removals; log what would have been deleted instead.  Stages still run and missing outputs are still created.',
                        dest='no_delete')


#     parser.add_argument('-recover',
#                         action='store_true',
#                         required=False,
#                         default=False,
#                         help='Used to recover missing meta-data.  This searches child directories for VolumeData.xml files and re-links them to the parent element in volume path.  This command does not recurse and does not need to be run on the top-level volume directory.',
#                         dest='verbose')

def _AddRecoverNotesParser(root_parser: argparse.ArgumentParser, subparsers):
    recover_parser = subparsers.add_parser('RecoverNotes',
                                           help='Recover or update notes from *.txt files. With ImportDir, scan an IDoc-style raw-data tree and copy changed notes into matching volume sections.', )
    recover_parser.set_defaults(func=call_recover_import_meta_data, parser=root_parser)

    recover_parser.add_argument('volumepath',
                                action='store',
                                type=str,
                                help='The path to the volume')

    recover_parser.add_argument('ImportDir',
                                nargs='?',
                                default=None,
                                help='Optional raw-data directory in IDoc import layout. When provided, notes.txt files from section folders are copied into matching volume sections.')

    recover_parser.add_argument('-save',
                                action='store_true',
                                required=False,
                                default=False,
                                help='Set this flag to save the VolumeData.xml files with the located linked elements included.',
                                dest='save_restoration')


def _AddRecoverParser(root_parser: argparse.ArgumentParser, subparsers):
    recover_parser = subparsers.add_parser('RecoverLinks',
                                           help='Used to recover missing meta-data.  This searches child directories for VolumeData.xml files and re-links them to the parent element in volume path.  This command does not recurse and does not need to be run on the top-level volume directory.', )
    recover_parser.set_defaults(func=call_recover_links, parser=root_parser)

    recover_parser.add_argument('volumepath',
                                action='store',
                                type=str,
                                help='The path to the volume')

    recover_parser.add_argument('-recurse',
                                action='store_true',
                                required=False,
                                default=False,
                                help='Set this flag to include sub-directories',
                                dest='recurse')

    recover_parser.add_argument('-save',
                                action='store_true',
                                required=False,
                                default=False,
                                help='Set this flag to save the VolumeData.xml files with the located linked elements included.',
                                dest='save_restoration')


def _AddXMLRepairParser(root_parser: argparse.ArgumentParser, subparsers):
    recover_parser = subparsers.add_parser('RepairXML',
                                           help='Fixes an issue where XML files have extra characters after the closing tag.  Should only apply to data before Dec 2023', )
    recover_parser.set_defaults(func=call_repair_xml, parser=root_parser)

    recover_parser.add_argument('volumepath',
                                action='store',
                                type=str,
                                help='The path to the volume')


def _GetPipelineXMLPath() -> str:
    if __spec__ is None:
        return os.path.join(os.path.dirname(__file__), 'config', 'Pipelines.xml')
    else:
        return pkgutil.get_data(__name__, os.path.join('config', 'Pipelines.xml'))  # type: ignore[return-value]


[docs] def BuildParserRoot() -> argparse.ArgumentParser: # conflict_handler = 'resolve' replaces old arguments with new if both use the same option flag parser = argparse.ArgumentParser('Buildscript', conflict_handler='resolve', description='Options available to all build commands. Specific pipelines extend this argument list.', epilog='Examples:\n' ' nornir-build ImportIDoc /data/volume /data/idoc\n' ' nornir-build -debug -computational_library cupy Mosaic /data/volume -Sections 1-10\n' ' nornir-build help Mosaic') _AddParserRootArguments(parser) # Create subparsers for commands pipeline_subparsers = parser.add_subparsers(title='Commands', dest='command') # Add a special help command that doesn't require volumepath help_parser = pipeline_subparsers.add_parser('help', help='Show help for a specific command') help_parser.add_argument('command_name', nargs='?', help='Name of the command to show help for') help_parser.set_defaults(func=print_help, parser=parser) _AddRecoverParser(parser, pipeline_subparsers) _AddRecoverNotesParser(parser, pipeline_subparsers) _AddXMLRepairParser(parser, pipeline_subparsers) _AddPipelineParsers(pipeline_subparsers) return parser
def _AddPipelineParsers(subparsers: argparse._SubParsersAction): PipelineXML = _GetPipelineXMLPath() # Load the element tree once and pass it to the later functions so we aren't parsing the XML text in the loop PipelineTree = pipelinemanager.PipelineManager.LoadPipelineXML(PipelineXML) for pipeline_name in pipelinemanager.PipelineManager.ListPipelines(PipelineTree): pipeline = pipelinemanager.PipelineManager.Load(PipelineTree, pipeline_name) pipeline_parser = subparsers.add_parser(pipeline_name, help=pipeline.Help, epilog=pipeline.Epilog) # type: ignore[union-attr] # Add volumepath as first positional argument pipeline_parser.add_argument('volumepath', action='store', type=str, help='The path to the volume') pipeline.GetArgParser(pipeline_parser, IncludeGlobals=True) # type: ignore[union-attr] pipeline_parser.set_defaults(func=call_pipeline, PipelineXmlFile=_GetPipelineXMLPath(), PipelineName=pipeline_name) CommandParserDict[pipeline_name] = pipeline_parser
[docs] def call_repair_xml(args): """Repair malformed VolumeData XML files with trailing content. This migration-style repair targets older metadata where characters were written after the closing XML tag (primarily pre-Dec 2023 data sets). """ volumeObj = nornir_buildmanager.volumemanager.volumemanager.VolumeManager.Load(args.volumepath) nornir_buildmanager.operations.migration.RepairCroppedXMLFilesInElement(volumeObj) # type: ignore[arg-type]
[docs] def call_recover_import_meta_data(args): """Recover notes metadata from text files under the volume or an import directory.""" volumeObj = nornir_buildmanager.volumemanager.volumemanager.VolumeManager.Load(args.volumepath) if volumeObj is None: prettyoutput.Log(f"Volume not found: {args.volumepath}") return import_dir = getattr(args, 'ImportDir', None) if import_dir: notesAdded = nornir_buildmanager.importers.shared.RecoverNotesFromImportDir(volumeObj, import_dir, None) # type: ignore[union-attr] recurse = True scan_path = import_dir else: notesAdded = nornir_buildmanager.importers.shared.TryAddNotes(volumeObj, volumeObj.FullPath, None) # type: ignore[union-attr] recurse = False scan_path = volumeObj.FullPath # type: ignore[union-attr] if not notesAdded: prettyoutput.Log(f"No notes recovered from {scan_path}.") return if args.save_restoration: volumeObj.Save(recurse=recurse) # type: ignore[union-attr] prettyoutput.Log("Recovered notes file saved.") else: prettyoutput.Log("Save flag not set, recovered notes, but not saved.")
[docs] def call_pipeline(args): pipelinemanager.PipelineManager.RunPipeline(PipelineXmlFile=args.PipelineXmlFile, PipelineName=args.PipelineName, args=args)
def _run_pipeline_segment(args: argparse.Namespace, volume_tree=None, flush_at_boundary: bool = False): """Run one pipeline segment, optionally reusing an in-memory volume tree.""" tree = pipelinemanager.PipelineManager.RunPipeline( PipelineXmlFile=args.PipelineXmlFile, PipelineName=args.PipelineName, args=args, volume_tree=volume_tree, ) if flush_at_boundary and tree is not None: nornir_buildmanager.volumemanager.volumemanager.VolumeManager.Save(tree) return tree def _GetFromNamespace(ns, attribname, default=None): if attribname in ns: return getattr(ns, attribname) else: return default def _pipeline_name_from_args(args: argparse.Namespace) -> str | None: """Return PipelineName or utility command name from parsed CLI args.""" pipeline = getattr(args, 'PipelineName', None) if not pipeline: pipeline = getattr(args, 'command', None) return pipeline def _chain_pipeline_display_name(name: str, index: int, total: int) -> str: """Sidebar/meta title for a chain segment (``Name (k/n)`` when chained).""" if total > 1: return f"{name} ({index + 1}/{total})" return name # Pipelines whose -Output directory is the primary product root; include it in # the dashboard header so parallel ExportAnnotationCrops runs are distinguishable. _OUTPUT_IN_HEADER_PIPELINES = frozenset({ "ExportAnnotationCrops", "RepairAnnotationOverlays", "ScoreAnnotationCrops", }) def _dashboard_volumepath(args: argparse.Namespace, pipeline: str | None = None) -> str | None: """Volume path for MQTT/dashboard header, with -Output when it is the crop root.""" volumepath = getattr(args, "volumepath", None) if not volumepath: return None name = pipeline if pipeline is not None else _pipeline_name_from_args(args) # Chain display names look like "ExportAnnotationCrops (2/3)". base = name.split(" (", 1)[0] if name else "" output = getattr(args, "OutputPath", None) if base in _OUTPUT_IN_HEADER_PIPELINES and isinstance(output, str) and output: return f"{volumepath} → {output}" return volumepath def _publish_early_run_meta_from_args(args: argparse.Namespace, pipeline: str | None = None) -> None: """Publish retained dashboard meta as soon as CLI args are known. Uses ``PipelineName`` when present (pipeline commands); otherwise ``command`` (utilities such as RecoverLinks). Skips when volumepath is missing. Optional *pipeline* overrides the displayed name (e.g. ``Prune (1/6)``). For AnnotationCrops pipelines, ``volumepath`` includes ``-Output`` so the dashboard header distinguishes parallel export targets. """ if pipeline is None: pipeline = _pipeline_name_from_args(args) if not pipeline: return volumepath = _dashboard_volumepath(args, pipeline) if not volumepath: return prettyoutput.publish_early_run_meta( pipeline=pipeline, volumepath=volumepath, compute=os.environ.get('NORNIR_COMPUTATIONAL_LIBRARY'), ) def _publish_chain_segment_meta(args: argparse.Namespace, *, index: int, total: int) -> None: """Publish running meta for a later ``--then`` segment without resetting ``start_ts``.""" name = _pipeline_name_from_args(args) if not name: return display = _chain_pipeline_display_name(name, index, total) volumepath = _dashboard_volumepath(args, display) if not volumepath: return prettyoutput.publish_run_meta( pipeline=display, volumepath=volumepath, status="running", compute=os.environ.get('NORNIR_COMPUTATIONAL_LIBRARY'), ) def _publish_chain_progress(index: int, total: int) -> None: """Publish the depth-0 ``Pipelines`` track for a ``--then`` chain.""" prettyoutput.publish_run_event( "iterate_progress", track_id="chain", label="Pipelines", current=index, total=total, depth=0, ) def _publish_pipeline_segment_complete(name: str) -> None: """Publish a sticky completed depth-0 track for a finished chain segment.""" prettyoutput.publish_run_event( "iterate_progress", track_id=f"pipeline:{name}", label=name, current=1, total=1, fraction=1.0, depth=0, ) def _publish_run_completion(succeeded: bool) -> None: """Publish retained final run status for the dashboard.""" import time prettyoutput.publish_run_meta( status="completed" if succeeded else "failed", end_ts=time.time(), )
[docs] def InitLogging(buildArgs): """Initialize persistent logging for the current command invocation. Logging writes through ``nornir_shared.misc.SetupLogging``. When a ``volumepath`` is present, file logs go under ``<volumepath>/logs``. ``-debug`` sets the level to DEBUG; otherwise WARN is used. """ # nornir_shared.Misc.RunWithProfiler('Execute()', "C:/Temp/profile.pr") buildArgs = _ReorderArgs(list(buildArgs)) parser = BuildParserRoot() (args, extraargs) = parser.parse_known_args(buildArgs) if 'volumepath' in args and args.volumepath: log_dir = os.path.join(args.volumepath, 'logs') if _GetFromNamespace(args, 'debug', False): SetupLogging(OutputPath=log_dir, Level=logging.DEBUG) else: SetupLogging(OutputPath=log_dir, Level=logging.WARN) else: SetupLogging(Level=logging.WARN)
[docs] def init_computational_library(args: argparse.Namespace): """Select CPU/GPU computation backend and export process environment. Records the backend the build asked for in ``NORNIR_COMPUTATIONAL_LIBRARY_REQUESTED`` and updates ``nornir_imageregistration``'s active backend, which sets ``NORNIR_COMPUTATIONAL_LIBRARY`` to the backend actually in effect. The two can diverge: requesting ``cupy`` on a host whose CuPy runtime probe fails leaves requested=cupy / effective=numpy. Keeping both lets a build log show that a GPU run silently fell back to the CPU. """ if args.computational_library == 'detect': if nornir_imageregistration.HasCupy(): args.computational_library = 'cupy' else: args.computational_library = 'numpy' else: args.computational_library = args.computational_library.lower() os.environ['NORNIR_COMPUTATIONAL_LIBRARY_REQUESTED'] = args.computational_library os.environ['NORNIR_COMPUTATIONAL_LIBRARY'] = args.computational_library nornir_imageregistration.SetActiveComputationLib( nornir_imageregistration.ComputationLib.cupy if args.computational_library == 'cupy' else nornir_imageregistration.ComputationLib.numpy)
def _GetValidCommands() -> list[str]: """Get list of all valid commands/pipelines.""" commands = ['help', 'RecoverLinks', 'RecoverNotes', 'RepairXML'] # Add pipeline commands from XML PipelineXML = _GetPipelineXMLPath() PipelineTree = pipelinemanager.PipelineManager.LoadPipelineXML(PipelineXML) commands.extend(pipelinemanager.PipelineManager.ListPipelines(PipelineTree)) return commands # Flags defined on the root parser only (see _AddParserRootArguments). Used to recognize # [volumepath, <root flags...>, <command>, ...] test/harness argv and normalize to # [<root flags...>, <command>, volumepath, ...] before subparser dispatch. _ROOT_FLAGS_NO_VALUE = frozenset({'-debug', '-verbose', '-lowpriority', '-lp', '-no-delete'}) _ROOT_FLAGS_WITH_VALUE = frozenset({'-computational_library'}) def _segment_is_root_only_flags(segment: list[str]) -> bool: """True if *segment* is a sequence of root-parser flags (and values for value-taking flags).""" j = 0 while j < len(segment): t = segment[j] if t in _ROOT_FLAGS_NO_VALUE: j += 1 continue if t in _ROOT_FLAGS_WITH_VALUE: if j + 1 >= len(segment): return False j += 2 continue return False return True def _leading_root_flag_segment_length(args: list[str]) -> int: """Return the length of a leading argv prefix consumed by root-parser flags.""" j = 0 while j < len(args): token = args[j] if token in _ROOT_FLAGS_NO_VALUE: j += 1 continue if token in _ROOT_FLAGS_WITH_VALUE: if j + 1 >= len(args): break j += 2 continue break return j def _ReorderArgs(args: list[str]) -> list[str]: """Reorder argv so root flags and subcommand precede volumepath (argparse subparser layout). Accepts legacy test orderings: - ``[volumepath, command, ...]`` (swap first two when command is second) - ``[volumepath, -debug, ..., command, ...]`` (move root flags + command before volumepath) - ``[-debug, ..., volumepath, command, ...]`` (launch.json / flags-before-path order) Leaves canonical ``[flags..., command, volumepath, ...]`` unchanged. """ if not args: return args valid_commands = frozenset(_GetValidCommands()) # Command-first: nothing to do if args[0] in valid_commands: return args # [volumepath, command, ...] — do not swap when args[0] is a flag (e.g. -debug ImportPMG vol) if len(args) > 1 and not args[0].startswith('-') and args[1] in valid_commands: return [args[1], args[0]] + args[2:] # [volumepath, <root flags...>, command, <tail>] if not args[0].startswith('-'): volumepath = args[0] for i in range(2, len(args)): if args[i] not in valid_commands: continue middle = args[1:i] if _segment_is_root_only_flags(middle): return middle + [args[i], volumepath] + args[i + 1 :] # [<root flags...>, volumepath, command, <tail>] — e.g. launch.json TEM configs prefix_len = _leading_root_flag_segment_length(args) if prefix_len < len(args): rest = args[prefix_len:] if len(rest) > 1 and not rest[0].startswith('-') and rest[1] in valid_commands: return args[:prefix_len] + [rest[1], rest[0]] + rest[2:] return args def _SplitChainSegments(buildArgs: list[str]) -> list[list[str]]: """Split argv on ``--then`` into per-pipeline segments.""" segments: list[list[str]] = [] current: list[str] = [] for token in buildArgs: if token == '--then': if not current: raise ValueError('--then cannot precede the first pipeline segment') segments.append(current) current = [] else: current.append(token) if not current: raise ValueError('no pipeline segment after final --then') segments.append(current) return segments def _AppendTimingOutput(volumepath: str, timer: TaskTimer) -> None: """Append a timing record for one pipeline invocation to Timing.txt.""" out_str = str(timer) prettyoutput.Log(out_str) time_text_full_path = os.path.join(volumepath, 'Timing.txt') try: with open(time_text_full_path, 'a') as output_file: output_file.writelines(out_str) except OSError: prettyoutput.Log('Could not write %s' % time_text_full_path)
[docs] def ExecuteChain(buildArgs: list[str]) -> None: """Run multiple pipelines sequentially in one process, separated by ``--then``.""" segments = _SplitChainSegments(buildArgs) first_segment = _ReorderArgs(segments[0]) InitLogging(first_segment) parser = BuildParserRoot() first_args = parser.parse_args(first_segment) if getattr(first_args, 'command', None) == 'help': first_args.func(first_args) return if not hasattr(first_args, 'volumepath') or not first_args.volumepath: parser.error("the following arguments are required: volumepath") if first_args.lowpriority: lowpriority() print("Warning, using low priority flag. This can make builds much slower") if hasattr(first_args, 'computational_library'): init_computational_library(first_args) root_flags = first_segment[:_leading_root_flag_segment_length(first_segment)] volumepath = first_args.volumepath valid_commands = frozenset(_GetValidCommands()) volume_tree = None succeeded = False num_segments = len(segments) try: for index, segment in enumerate(segments): if index > 0 and segment[0] not in valid_commands: parser.error(f"unknown pipeline in chain segment: {segment[0]}") if index == 0: segment_argv = first_segment else: segment_argv = _ReorderArgs(root_flags + [segment[0], volumepath] + segment[1:]) args = parser.parse_args(segment_argv) cmd_name = _pipeline_name_from_args(args) or getattr(args, 'PipelineName', None) display_name = _chain_pipeline_display_name(cmd_name, index, num_segments) if index == 0: _publish_early_run_meta_from_args(args, pipeline=display_name) else: _publish_chain_segment_meta(args, index=index, total=num_segments) if num_segments > 1: _publish_chain_progress(index, num_segments) timer = TaskTimer() try: timer.Start(cmd_name) volume_tree = _run_pipeline_segment(args, volume_tree=volume_tree, flush_at_boundary=True) if num_segments > 1: _publish_pipeline_segment_complete(cmd_name) _publish_chain_progress(index + 1, num_segments) finally: timer.End(cmd_name) _AppendTimingOutput(volumepath, timer) succeeded = True finally: _publish_run_completion(succeeded)
[docs] def Execute(buildArgs=None): """Run the nornir-build command line entrypoint. The argv parser accepts canonical subcommand ordering and legacy test ordering that starts with ``volumepath``. """ # Spend more time on each thread before switching # sys.setswitchinterval(500) if buildArgs is None: buildArgs = sys.argv[1:] if '--then' in buildArgs: ExecuteChain(buildArgs) return # Reorder arguments to support both command-first and volumepath-first patterns buildArgs = _ReorderArgs(buildArgs) # Change the temp directory if nornir specifies an alternate in the environment variables if 'NORNIR_TEMP_DIR' in os.environ: temp_dir = os.environ['NORNIR_TEMP_DIR'] os.makedirs(temp_dir, exist_ok=True) os.environ['TEMP'] = temp_dir os.environ['TMP'] = temp_dir os.environ['TMPDIR'] = temp_dir InitLogging(buildArgs) Timer = TaskTimer() parser = BuildParserRoot() args = parser.parse_args(buildArgs) # Help command does not require volumepath or timing output. if getattr(args, 'command', None) == 'help': args.func(args) return # For all other commands, require volumepath if not hasattr(args, 'volumepath') or not args.volumepath: parser.error("the following arguments are required: volumepath") if args.lowpriority: lowpriority() print("Warning, using low priority flag. This can make builds much slower") if hasattr(args, 'computational_library'): init_computational_library(args) _publish_early_run_meta_from_args(args) cmd_name = None if hasattr(args, 'PipelineName'): cmd_name = args.PipelineName elif len(buildArgs) >= 2: cmd_name = buildArgs[1] succeeded = False try: if cmd_name is not None: Timer.Start(cmd_name) args.func(args) succeeded = True finally: if cmd_name is not None: Timer.End(cmd_name) _AppendTimingOutput(args.volumepath, Timer) _publish_run_completion(succeeded)
if __name__ == '__main__': Execute()