"""Brain-5D command-line simulation entry point with dashboard integration. This module runs the simulation and optionally starts the dashboard in the main thread (to handle signals). The simulation runs in a daemon thread so that the dashboard can control it via the OperatorBridge. CLI semantics: 1. Dashboard mode (default): python -m src.main --config configs/poc_config.yaml -> starts RuntimeController in IDLE -> does NOT automatically advance ticks -> operator controls runtime through dashboard/API 2. Headless mode: python -m src.main --config configs/poc_structural_live.yaml --no-dashboard --ticks 500 -> runs exactly N ticks through RuntimeController.run_ticks(N) 3. --ticks in dashboard mode only overrides the config value; the controller still starts IDLE. Usage: python -m src.main --config configs/poc_config.yaml python -m src.main --config configs/poc_config.yaml --no-dashboard python -m src.main --config configs/poc_config.yaml --no-dashboard --ticks 500 python -m src.main --config configs/poc_config.yaml --observe --benchmark """ from __future__ import annotations import argparse import random import statistics import sys import threading from dataclasses import asdict from pathlib import Path from typing import Any, Literal, cast from src.config.loader import load_config # ================================================================ # Canonical Runtime Controller # ================================================================ from src.controller.runtime import PostTickHook from src.controller.runtime import RuntimeController as _RuntimeController from src.core import Brain5DConfig, NeuralNetwork from src.core.spatial_index import ( coords_to_linear, linear_to_5d, make_boundary_coord, unpack_coords, ) from src.diagnostics.propagation import PropagationAnalyzer from src.diagnostics.stimulus import StimulusEngine, StimulusResult from src.diagnostics.topology_health import TopologyHealth from src.experience import build_experience_subsystem from src.homeostasis import HomeostasisEngine from src.learning.learning_engine import LearningEngine # ================================================================ # Snapshot Writer # ================================================================ from src.storage import B5DSnapshotWriter from src.telemetry.history import History from src.telemetry.probes import ProbeManager from src.telemetry.spike_history import SpikeHistory from src.utils.run_artifacts import RunArtifacts from src.version import BRAIN5D_VERSION_DISPLAY # ================================================================ # Dashboard Integration – with None‑fallback # ================================================================ _dashboard_available = False _OperatorBridge: type | None = None _serve_dashboard: Any = None _DashboardStateStore: type | None = None try: from src.dashboard.operator_bridge import OperatorBridge as _OperatorBridge from src.dashboard.server import serve_dashboard as _serve_dashboard from src.dashboard.state import DashboardStateStore as _DashboardStateStore _dashboard_available = True except ImportError as e: print(f"⚠️ Dashboard not available: {e}") # ================================================================ # Helper Functions # ================================================================ def sample_positions_excluding_poc( total: int, reserved: set[int], n: int, rng: random.Random, ) -> list[int]: """Sample free linear positions while preserving reserved PoC cells.""" available = [idx for idx in range(total) if idx not in reserved] if n > len(available): raise ValueError( f"Not enough unreserved positions: need {n}, have {len(available)}" ) return rng.sample(available, n) def build_network(config_dict: dict[str, Any]) -> tuple[NeuralNetwork, random.Random]: """Build the network from configuration using the new Brain5DConfig.""" # Convert dict to Brain5DConfig config = Brain5DConfig.from_dict(config_dict) rng = random.Random(int(config_dict.get("seed", 42))) # Create network network = NeuralNetwork(config, rng) dims = config.dimensions # Extract topology information from the original config dict topology = config_dict.get("topology", {}) input_dim = topology.get("input", {}).get("dimension", "x") input_coord = topology.get("input", {}).get("coordinate", 0) output_dim = topology.get("output", {}).get("dimension", "x") output_coord = topology.get("output", {}).get("coordinate", dims[0] - 1) # Add reserved neurons (input, output, diagnostic) input_coord_5d = make_boundary_coord(dims, input_dim, input_coord) output_coord_5d = make_boundary_coord(dims, output_dim, output_coord) diag_coord = tuple( config_dict.get("diagnostics", {}).get("target_coord", (0, 0, 0, 0, 0)) ) reserved_coords = {input_coord_5d, output_coord_5d, diag_coord} reserved_indices = {coords_to_linear(coord, dims) for coord in reserved_coords} total_positions = 1 for d in dims: total_positions *= d initial_neurons = int(config_dict.get("initial_neurons", 5000)) chosen = sample_positions_excluding_poc( total_positions, reserved_indices, initial_neurons - len(reserved_coords), rng, ) for idx in chosen: network.add_neuron(linear_to_5d(idx, dims)) for coord in sorted(reserved_coords): network.add_neuron(coord) # Set input/output cells network.set_input_output_cells( input_dim, input_coord, output_dim, output_coord, ) # Initialize random connections conn_per_neuron = config.network.initial_connections_per_neuron radius = config.network.neighbour_radius network.initialize_random_connections(conn_per_neuron, radius) return network, rng def setup_learning( network: NeuralNetwork, config_dict: dict[str, Any], ) -> LearningEngine | None: """Set up and attach the learning engine if enabled.""" try: learning = LearningEngine(network, config_dict) if learning.enabled: learning.attach() print("✅ Learning engine attached") return learning except Exception as e: print(f"⚠️ Learning engine setup failed: {e}") return None def setup_homeostasis( network: NeuralNetwork, config_dict: dict[str, Any], ) -> HomeostasisEngine | None: """Set up and attach the homeostasis engine if enabled.""" try: homeostasis = HomeostasisEngine(network, config_dict) if homeostasis.enabled: homeostasis.attach() print("✅ Homeostasis engine attached") return homeostasis except Exception as e: print(f"⚠️ Homeostasis engine setup failed: {e}") return None def setup_observatory( network: NeuralNetwork, config_dict: dict[str, Any], spike_history: SpikeHistory, history: History, probes: ProbeManager, ) -> Any | None: """Set up the observatory if visualization is enabled.""" vis = config_dict.get("visualization", {}) if not vis.get("enabled", False): return None try: from src.visualization.observatory import Observatory observatory = Observatory(network, config_dict, spike_history, history, probes) print("✅ Observatory ready") return observatory except Exception as e: print(f"⚠️ Observatory setup failed: {e}") return None # ================================================================ # Main # ================================================================ def _configure_utf8_streams() -> None: """Ensure stdout/stderr can emit Unicode even on cp1252 consoles.""" for stream_name in ("stdout", "stderr"): stream = getattr(sys, stream_name) if stream is not None and hasattr(stream, "reconfigure"): try: stream.reconfigure(encoding="utf-8", errors="replace") except Exception: pass def main() -> int: _configure_utf8_streams() parser = argparse.ArgumentParser( description="Brain 5D v0.5 - homeostatic self-regulation with dashboard" ) parser.add_argument("--config", default="configs/poc_config.yaml") parser.add_argument("--observe", action="store_true") parser.add_argument("--benchmark", action="store_true") parser.add_argument("--no-dashboard", action="store_true") parser.add_argument("--no-learning", action="store_true") parser.add_argument("--no-homeostasis", action="store_true") parser.add_argument( "--dashboard-host", default="127.0.0.1", help="Dashboard bind host (default: 127.0.0.1; use 0.0.0.0 for LAN)", ) parser.add_argument( "--dashboard-port", type=int, default=8765, help="Dashboard HTTP server port (default: 8765)", ) parser.add_argument( "--ticks", type=int, default=None, help="Override config ticks. In dashboard mode (default) this only updates config;\n the controller starts IDLE and must be advanced via API.\n Use --no-dashboard for automatic execution.", ) args = parser.parse_args() # Load configuration config_dict: dict[str, Any] = cast(dict[str, Any], load_config(args.config)) if args.observe: config_dict["visualization"] = config_dict.get("visualization", {}) config_dict["visualization"]["enabled"] = True if args.ticks is not None: config_dict["ticks"] = args.ticks # Compute config SHA-256 for provenance import hashlib try: _effective_config_path = Path(args.config).resolve() _config_bytes = _effective_config_path.read_bytes() config_dict["_sha256"] = hashlib.sha256(_config_bytes).hexdigest() config_dict["_path"] = str(_effective_config_path) except Exception: config_dict["_sha256"] = "" config_dict["_path"] = str(args.config) print(f"🚀 Brain 5D - v{BRAIN5D_VERSION_DISPLAY} with dashboard") print(f"📄 Config: {args.config} (sha256={config_dict['_sha256'][:16]}...)") # --- Build network --- network, rng = build_network(config_dict) # --- Setup engines --- learning = None if args.no_learning else setup_learning(network, config_dict) homeostasis = ( None if args.no_homeostasis else setup_homeostasis(network, config_dict) ) # --- Topology health --- health = cast(dict[str, Any], TopologyHealth(network).analyze()) # type: ignore[reportUnknownMemberType] # --- Telemetry --- telemetry_cfg = config_dict.get("telemetry", {}) history = History(int(telemetry_cfg.get("history_ticks", 10000))) spike_history = SpikeHistory(int(telemetry_cfg.get("spike_history_ticks", 1000))) probes = ProbeManager(network, config_dict) # --- Stimulus --- stimulus = StimulusEngine(config_dict, rng) # --- Propagation --- propagation = PropagationAnalyzer(network.output_cells) # --- Diagnostic probe --- diag_target = tuple( config_dict.get("diagnostics", {}).get("target_coord", (0, 0, 0, 0, 0)) ) diag_id = next( (nid for nid in network.neurons if unpack_coords(nid) == diag_target), None, ) if diag_id is not None: probes.add_probe(diag_id) # --- Observatory --- observatory = setup_observatory( network, config_dict, spike_history, history, probes ) # ================================================================ # Canonical RuntimeController Setup # ================================================================ dashboard_stop_event = threading.Event() operator_bridge = None state_store = None controller = None # Shared telemetry holders (filled by hooks, read by dashboard publisher) _storage_telemetry: dict[str, Any] = {"available": False} _self_org_stats: dict[str, Any] = {"available": False} def _self_org_stats_func() -> dict[str, Any]: return _self_org_stats # Telemetry and logging state (shared via closure) core_times: list[float] = [] warmup = 100 # Snapshot path _snapshot_dir = Path("artifacts") _snapshot_dir.mkdir(parents=True, exist_ok=True) _snapshot_path = _snapshot_dir / "latest.b5d" _snapshot_temp = _snapshot_dir / "latest.b5d.tmp" # Snapshot writer _snapshot_writer = B5DSnapshotWriter(restart_capable=True) def _write_snapshot() -> None: """Write a .b5d snapshot atomically (temp file → validate → rename).""" try: # Collect metadata git_commit = "unknown" git_dirty = True try: import subprocess result = subprocess.run( ["git", "rev-parse", "HEAD"], capture_output=True, text=True, timeout=5, cwd=Path(__file__).resolve().parents[1], ) if result.returncode == 0: git_commit = result.stdout.strip() dirty_result = subprocess.run( ["git", "status", "--porcelain"], capture_output=True, text=True, timeout=5, cwd=Path(__file__).resolve().parents[1], ) git_dirty = bool(dirty_result.stdout.strip()) except Exception: pass from src.storage.b5d import JSONMapping metadata: JSONMapping = { "type": "brain5d-snapshot", "version": BRAIN5D_VERSION_DISPLAY, "tick": network.current_tick, "neuron_count": len(network.neurons), "synapse_count": network.synapse_count, "dimensions": list(network.dimensions), "seed": config_dict.get("seed", 42), "git_commit": git_commit, "git_dirty": git_dirty, "config": { "path": str(args.config), "sha256": config_dict.get("_sha256", ""), }, } # Write to temp file from src.storage.b5d import NetworkSnapshotLike _snapshot_writer.write( str(_snapshot_temp), cast("NetworkSnapshotLike", network), metadata=metadata, ) # Validate written file (reader must be closed before rename on Windows) from src.storage.b5d import B5DReader reader = B5DReader(str(_snapshot_temp)) try: reader.validate_invariants(full_scan=False) finally: reader.close() # Atomic rename _snapshot_temp.replace(_snapshot_path) # Preserve immutable historical snapshot with timestamp import time as _time tick = network.current_tick ts = _time.strftime("%Y%m%d_%H%M%S") historical = _snapshot_dir / f"snapshot_t{tick}_{ts}.b5d" try: import shutil shutil.copy2(_snapshot_path, historical) except Exception: pass except Exception as exc: print(f"⚠️ Snapshot write failed: {type(exc).__name__}: {exc}") # Clean up temp file on failure try: _snapshot_temp.unlink(missing_ok=True) except Exception: pass # ================================================================ # Dashboard State Publishing Helper # ================================================================ def _publish_dashboard_state( state_store: Any, network: NeuralNetwork, result: Any, learning: Any, homeostasis: Any, storage_telemetry: dict[str, Any], self_org_stats: dict[str, Any], status: str, ) -> None: """Build and publish a DashboardSnapshot from live runtime data.""" from src.dashboard.health_builder import enrich_snapshot from src.dashboard.models import ( DashboardSnapshot, HomeostasisMetrics, LearningMetrics, NetworkMetrics, SelfOrganizationMetrics, SpikeMetrics, StorageMetrics, SystemMetrics, ) learning_stats = learning.stats if learning is not None else None homeo_stats = homeostasis.stats if homeostasis is not None else None # Storage telemetry: only populate if available storage_available = storage_telemetry.get("available", False) storage = StorageMetrics(available=storage_available) if storage_available: storage = StorageMetrics( available=True, queue_depth=storage_telemetry.get("queue_depth"), queue_capacity=storage_telemetry.get("queue_capacity"), batches_enqueued=storage_telemetry.get("batches_enqueued"), batches_written=storage_telemetry.get("batches_written"), deltas_written=storage_telemetry.get("deltas_written"), bytes_written=storage_telemetry.get("bytes_written"), dropped_batches=storage_telemetry.get("dropped_batches"), write_latency_ms=storage_telemetry.get("write_latency_ms"), commit_latency_ms=storage_telemetry.get("commit_latency_ms"), journal_size_bytes=storage_telemetry.get("journal_size_bytes"), worker_failed=storage_telemetry.get("worker_failed"), ) # Self-organization metrics: only populate if available so_available = self_org_stats.get("available", False) self_org = SelfOrganizationMetrics(available=so_available) if so_available: self_org = SelfOrganizationMetrics( available=True, neurons_created=self_org_stats.get("neurons_created"), neurons_removed=self_org_stats.get("neurons_removed"), synapses_created=self_org_stats.get("synapses_created"), synapses_pruned=self_org_stats.get("synapses_pruned"), ) snapshot = DashboardSnapshot( status=status, version=BRAIN5D_VERSION_DISPLAY, system=SystemMetrics( tick=result.tick, neurons=len(network.neurons), synapses=getattr(result, "total_synapses", network.synapse_count), spikes_total=result.total_spikes, spikes_last_tick=getattr(result, "spikes_this_tick", 0), core_step_ms=getattr(result, "core_step_ms", 0.0), mean_energy=getattr(result, "mean_energy", 0.0), ), learning=LearningMetrics( stdp_updates=getattr(learning_stats, "stdp_weight_updates", 0), reward_updates=getattr(learning_stats, "reward_weight_updates", 0), rewards_received=getattr(learning_stats, "rewards_received", 0), rewards_applied=getattr(learning_stats, "rewards_applied", 0), pending_rewards=getattr(learning_stats, "pending_rewards", 0), update_ms=getattr(learning_stats, "last_update_ms", 0.0), engine_attached=learning is not None, stdp_enabled=bool(config_dict.get("stdp", {}).get("enabled", False)), eligibility_enabled=bool( config_dict.get("eligibility", {}).get("enabled", False) ), reward_enabled=bool( config_dict.get("reward", {}).get("enabled", False) ), ), storage=storage, self_organization=self_org, homeostasis=HomeostasisMetrics( enabled=getattr(homeo_stats, "enabled", False), target_rate_hz=getattr(homeo_stats, "target_rate_hz", 0.0), actual_rate_hz=getattr(homeo_stats, "mean_rate_hz", 0.0), rate_error_hz=getattr(homeo_stats, "mean_rate_error_hz", 0.0), mean_rate_hz=getattr(homeo_stats, "mean_rate_hz", 0.0), mean_rate_error_hz=getattr(homeo_stats, "mean_rate_error_hz", 0.0), mean_threshold_adaptation=getattr( homeo_stats, "mean_threshold_adaptation", 0.0 ), target_energy=getattr(homeo_stats, "target_energy", 0.0), mean_energy=getattr(homeo_stats, "mean_energy", 0.0), mean_energy_error=getattr(homeo_stats, "mean_energy_error", 0.0), active_neurons=getattr(homeo_stats, "active_neurons", 0), updates=getattr(homeo_stats, "updates", 0), ), spikes=SpikeMetrics( total_spikes=result.total_spikes, active_neurons=len(getattr(result, "spike_ids", ())), mean_firing_rate_hz=0.0, burst_index=0.0, synchrony=0.0, spike_count_last_tick=getattr(result, "spikes_this_tick", 0), ), network=NetworkMetrics( tick=result.tick, neuron_count=len(network.neurons), synapse_count=getattr(result, "total_synapses", network.synapse_count), active_neurons=len(getattr(result, "spike_ids", ())), silent_neurons=len(network.neurons) - len(getattr(result, "spike_ids", ())), mean_firing_rate_hz=0.0, burst_index=0.0, synchrony=0.0, mean_energy=getattr(result, "mean_energy", 0.0), mean_threshold_adaptation=0.0, e_i_ratio=0.0, clustering_coefficient=0.0, mean_path_length=0.0, ), runtime={ "config_path": str(config_dict.get("_path", "")), "config_sha256": str(config_dict.get("_sha256", "")), "state_publish_interval_ticks": int( config_dict.get("dashboard", {}).get( "state_publish_interval_ticks", 10 ) ), }, ) state_store.publish(enrich_snapshot(snapshot, config_dict)) try: controller = _RuntimeController( network=network, homeostasis=None, # Homeostasis is attached as post-step hook batch_size=10, loop_delay_ms=0.0, telemetry_interval_ticks=10, snapshot_callback=_write_snapshot, ) print("✅ Canonical RuntimeController created (idle)") experience = build_experience_subsystem(config_dict, network, learning) if experience is not None: experience.attach_runtime(controller) print("✅ ExperienceEngine attached via runtime hooks") # Shared state for stimulus result (set by pre-hook, read by post-hook) _last_stim: list[StimulusResult | None] = [None] # Register pre-tick hook for stimulus def _pre_tick(tick: int) -> None: if dashboard_stop_event.is_set(): return _last_stim[0] = stimulus.apply(network, tick) # type: ignore[reportUnknownMemberType] controller.add_pre_hook(_pre_tick) # Register post-tick hook for telemetry, artifacts, logging, dashboard def _on_tick(tick: int, result: Any) -> None: if dashboard_stop_event.is_set(): return stim_result = _last_stim[0] # Reward reward_cfg = config_dict.get("reward", {}) if ( learning and learning.params.reward_enabled and reward_cfg.get("reward_source", "external") == "output_spike" and result.output_spike_ids ): learning.set_reward( float(reward_cfg.get("output_spike_value", 1.0)), result.tick, ) # Telemetry history.append_from_stepresult(result) # type: ignore[reportUnknownMemberType] spike_history.append(result.tick, result.spike_ids) if stim_result is not None: propagation.observe(stim_result, result) # Artifacts metric = history.get_all()[-1] # type: ignore[reportUnknownVariableType] artifacts.log_metrics(metric) # type: ignore[reportUnknownArgumentType] artifacts.log_spikes(result.tick, result.spike_ids) if stim_result is not None: artifacts.log_stimulus(stim_result) # Benchmark if args.benchmark and result.tick >= warmup: core_times.append(result.core_step_ms) # Console logging log_cfg = config_dict.get("logging", {}) log_interval = log_cfg.get("interval_ticks", 100) if (result.tick + 1) % log_interval == 0: print( f"Tick {result.tick + 1:4d} | spikes={result.spikes_this_tick:4d} " f"| total={result.total_spikes:6d} " f"| queue={result.queued_events:5d} " f"| {result.core_step_ms:.3f} ms" ) # Snapshot enrichment traverses all dashboard components/parameters; # publish at a bounded cadence instead of taxing every simulation tick. _dashboard_cfg = config_dict.get("dashboard", {}) _publish_interval = max( 1, int(_dashboard_cfg.get("state_publish_interval_ticks", 10)) ) if state_store is not None and (result.tick + 1) % _publish_interval == 0: try: _publish_dashboard_state( state_store=state_store, network=network, result=result, learning=learning, homeostasis=homeostasis, storage_telemetry=_storage_telemetry, self_org_stats=_self_org_stats_func(), status="running", ) except Exception: pass # Dashboard publishing must never break simulation controller.add_hook(_on_tick) # --- Observatory --- vis_cfg = config_dict.get("visualization", {}) refresh_interval = vis_cfg.get("refresh_interval_ticks", 100) if observatory: controller.add_hook( lambda t, r=None: ( observatory.draw() if (t + 1) % refresh_interval == 0 else None ) ) except Exception as e: print(f"⚠️ RuntimeController setup failed: {e}") controller = None # ================================================================ # Self-Organization Setup (optional, verdrahtet Coordinator + Plasticity + Journal) # ================================================================ _self_org_coordinator = None _self_org_plasticity = None _self_org_approval_policy = None _structural_journal = None so_cfg = config_dict.get("self_organization", {}) if so_cfg.get("enabled", False) and controller is not None: try: from src.self_organization.composition import compose_structural_subsystem # Config-authoritative plasticity limits: derive from YAML _allow_neurogenesis = bool(so_cfg.get("neurogenesis_enabled", True)) _allow_neuron_pruning = bool(so_cfg.get("pruning_enabled", False)) _allow_synapse_sprouting = bool(so_cfg.get("sprouting_enabled", False)) _allow_synapse_pruning = bool(so_cfg.get("synapse_pruning_enabled", False)) _max_changes = int(so_cfg.get("neurogenesis_max_per_cycle", 1)) _structural_journal_path = _snapshot_dir / "structural.journal" _composed = compose_structural_subsystem( network, _structural_journal_path, coordinator_enabled=True, coordinator_dry_run=False, max_changes_per_tick=_max_changes, allow_neurogenesis=_allow_neurogenesis, allow_neuron_pruning=_allow_neuron_pruning, allow_synapse_sprouting=_allow_synapse_sprouting, allow_synapse_pruning=_allow_synapse_pruning, ) _self_org_plasticity = _composed["plasticity"] _self_org_coordinator = _composed["coordinator"] _structural_journal = _composed["journal"] _manipulator = _composed["manipulator"] # Build the canonical approval policy from config so the dashboard # can read the real self-organization gate state. from src.self_organization.approval import ( ProposalApprovalPolicy, StructuralPlasticityConfig, ) _self_org_approval_policy = ProposalApprovalPolicy( StructuralPlasticityConfig( enabled=True, dry_run=False, auto_approval=bool(so_cfg.get("auto_approval", False)), auto_approval_threshold=float( so_cfg.get("auto_approval_threshold", 0.8) ), max_changes_per_tick=int(so_cfg.get("max_changes_per_tick", 5)), max_neuron_additions_per_tick=int( so_cfg.get("neurogenesis_max_per_cycle", 1) ), max_neuron_removals_per_tick=int( so_cfg.get("max_neuron_removals_per_tick", 0) ), max_synapse_additions_per_tick=int( so_cfg.get("sprouting_max_out_degree", 5) ), max_synapse_removals_per_tick=int( so_cfg.get("max_synapse_removals_per_tick", 5) ), min_neurons=int(so_cfg.get("min_neurons", 100)), max_neurons=int(so_cfg.get("max_neurons", 100_000)), allow_neuron_pruning=_allow_neuron_pruning, allow_synapse_pruning=_allow_synapse_pruning, cooldown_ticks=int(so_cfg.get("cooldown_ticks", 100)), ) ) # The legacy SelfOrganizationEngine ran_cycle() mutates the # network directly through its own manipulator, bypassing the # canonical Coordinator -> Approval -> PlasticityEngine path. # For Alpha.5, structural mutation must flow exclusively through # the canonical path. The legacy engine is NOT attached. # It remains available for proposal-generation research only. # Update self-org telemetry _self_org_stats.update( available=True, neurons_created=0, neurons_removed=0, synapses_created=0, synapses_pruned=0, ) print("✅ SelfOrganizationCoordinator + PlasticityEngine + Journal created") print( " (canonical path only; legacy SelfOrganizationEngine NOT attached)" ) # Attach the SelfOrganizationRuntimeAdapter as a post-tick hook. # This feeds real HomeostasisSignals through the policy and into # the coordinator — closing the production signal->policy->coordinator # path. The adapter does NOT mutate the network. # # CONFIG-AUTHORITATIVE: interval_ticks and policy_config are # derived from the YAML self_organization section. Hardcoded # defaults are never used when production config exists. if homeostasis is not None: try: from src.self_organization.policy import ( SelfOrganizationPolicyConfig, ) from src.self_organization.runtime_adapter import ( SelfOrganizationRuntimeAdapter, ) _so_interval = int(so_cfg.get("interval_ticks", 100)) _so_policy_config = SelfOrganizationPolicyConfig.from_config( config_dict ) _so_adapter = SelfOrganizationRuntimeAdapter( homeostasis_engine=homeostasis, coordinator=_self_org_coordinator, interval_ticks=_so_interval, policy_config=_so_policy_config, ) controller.add_hook(_so_adapter) print( f" ✅ SelfOrganizationRuntimeAdapter attached (interval={_so_interval}, config-authoritative)" ) except Exception as adapter_err: print(f" ⚠️ SelfOrganizationRuntimeAdapter failed: {adapter_err}") except Exception as e: print(f"⚠️ Self-organization setup failed: {e}") # ================================================================ # Runtime Delta Persistence (AsyncStorageSession) # ------------------------------------------------ # CONFIG-AUTHORITATIVE: This subsystem is ONLY started when the # authoritative config enables it. poc_config.yaml has # storage.enabled = false # storage.runtime.enabled = false # so the dashboard must show "disabled by config" for Delta Storage, # NOT fake zeros from a silently-active worker. # # The four persistence systems are deliberately separated: # 1. Snapshot Service -> always on (dashboard inspection) # 2. Runtime Delta Persistence -> gated by storage.runtime.enabled # 3. Structural Journal -> gated by self_organization.enabled # 4. Runtime Checkpoint -> gated by storage.checkpoint.enabled # ================================================================ _storage_cfg_raw = config_dict.get("storage", {}) _storage_cfg: dict[str, Any] = ( cast("dict[str, Any]", _storage_cfg_raw) if isinstance(_storage_cfg_raw, dict) else {} ) _storage_runtime_cfg_raw = _storage_cfg.get("runtime", {}) _storage_runtime_cfg: dict[str, Any] = ( cast("dict[str, Any]", _storage_runtime_cfg_raw) if isinstance(_storage_runtime_cfg_raw, dict) else {} ) _journal_cfg_raw = _storage_cfg.get("journal", {}) _journal_cfg: dict[str, Any] = ( cast("dict[str, Any]", _journal_cfg_raw) if isinstance(_journal_cfg_raw, dict) else {} ) _storage_runtime_enabled = bool(_storage_cfg.get("enabled", False)) and bool( _storage_runtime_cfg.get("enabled", False) ) if _storage_runtime_enabled: try: from src.storage.async_runtime import ( AsyncStorageConfig, AsyncStorageSession, ) from src.storage.delta_journal import DeltaJournal, JournalCorruptionError from src.storage.runtime import StorageRuntimeConfig _journal_path = _snapshot_dir / "latest.b5d.journal" _storage_runtime_config = StorageRuntimeConfig( snapshot_path=_snapshot_path, journal_path=_journal_path, commit_interval_ticks=int( _journal_cfg.get("commit_interval_ticks", 10) ), capture_policy=cast( "Literal['full_change_scan', 'dirty_tracking']", _storage_runtime_cfg.get("capture_policy", "full_change_scan"), ), ) _delta_journal = DeltaJournal(str(_journal_path)) try: _delta_journal.open() except JournalCorruptionError as _journal_err: # A corrupt journal from a previous run would prevent startup. # Alpha.5 uses the journal as runtime delta persistence; losing # the tail is acceptable because the canonical .b5d snapshot is # the source of truth. Rename the corrupt file and start fresh. _corrupt_path = _journal_path.with_suffix(".b5d.journal.corrupt") try: _corrupt_path.unlink(missing_ok=True) _journal_path.rename(_corrupt_path) except Exception: pass _delta_journal = DeltaJournal(str(_journal_path)) _delta_journal.open() print( f"⚠️ Reset corrupt delta journal ({_journal_err}); " f"old file preserved at {_corrupt_path.name}" ) if _delta_journal.last_tick > network.current_tick: import time as _time previous_last_tick = _delta_journal.last_tick _delta_journal.close() _restart_archive = _journal_path.with_name( f"{_journal_path.name}.restart-{previous_last_tick}-" f"{_time.time_ns()}" ) _journal_path.replace(_restart_archive) _delta_journal = DeltaJournal(str(_journal_path)) _delta_journal.open() print( "ℹ️ Archived delta journal from a later runtime " f"(last_tick={previous_last_tick}, " f"current_tick={network.current_tick}) as " f"{_restart_archive.name}" ) _delta_journal.close() from src.storage.runtime import RuntimeNetworkLike _async_config = AsyncStorageConfig() _async_storage = AsyncStorageSession( network=cast("RuntimeNetworkLike", network), runtime_config=_storage_runtime_config, async_config=_async_config, ) _async_storage.start() # Periodisch Telemetrie auslesen def _update_storage_telemetry() -> None: try: tel = _async_storage.telemetry _storage_telemetry.update( available=True, queue_depth=tel.queue_depth, queue_capacity=tel.queue_capacity, batches_enqueued=tel.batches_enqueued, batches_written=tel.batches_written, deltas_written=tel.deltas_written, bytes_written=tel.bytes_written, dropped_batches=tel.dropped_batches, write_latency_ms=tel.write_latency_ms, commit_latency_ms=tel.commit_latency_ms, journal_size_bytes=( _delta_journal.path.stat().st_size if _delta_journal.path.exists() else 0 ), worker_failed=tel.worker_failed, ) except Exception: pass # Storage-Telemetrie im post-tick hook aktualisieren if controller is not None: controller.add_hook(lambda _tick, _result: _update_storage_telemetry()) print("✅ AsyncStorageSession attached with telemetry (config-enabled)") except Exception as e: print(f"⚠️ Storage telemetry setup failed: {type(e).__name__}: {e}") else: # Storage is disabled by config — keep telemetry explicitly unavailable # so the dashboard renders "disabled by config" instead of fake zeros. _storage_telemetry.update(available=False) print( "ℹ️ Runtime Delta Persistence: disabled by config (storage.runtime.enabled=false)" ) # ================================================================ # OperatorBridge & Dashboard Setup # ================================================================ if not args.no_dashboard and _dashboard_available and controller is not None: try: # Create TelemetryFrameStore for atomic live visualization # Read telemetry config from config_dict, fall back to defaults _dashboard_cfg: dict[str, Any] = config_dict.get("dashboard", {}) # type: ignore[type-arg] _live_telemetry_cfg: dict[str, Any] = _dashboard_cfg.get("live_telemetry", {}) # type: ignore[type-arg] _lt_enabled = bool(_live_telemetry_cfg.get("enabled", True)) _lt_capture = int(_live_telemetry_cfg.get("capture_interval_ticks", 5)) _lt_window = int(_live_telemetry_cfg.get("activity_window_ticks", 20)) _sim_cfg: dict[str, Any] = config_dict.get("simulation", {}) # type: ignore[type-arg] _sim_dt_ms = float(_sim_cfg.get("dt_ms", 1.0)) _telemetry_store: Any = None if _lt_enabled: from src.dashboard.live_projection import ( TelemetryFrameStore, make_telemetry_hook, ) _telemetry_store = TelemetryFrameStore( capture_interval_ticks=_lt_capture, activity_window_ticks=_lt_window, ) _telemetry_store.set_dt_ms(_sim_dt_ms) # Prime Tick-0 frame so dashboard can respond immediately _telemetry_store.prime(controller.network) # Register post-tick hook via safe wrapper (routes errors to error buffer) from src.dashboard.live_projection import ( NetworkAccess as _NetworkAccess, ) _hook: PostTickHook = make_telemetry_hook( _telemetry_store, cast("_NetworkAccess", controller.network) ) controller.add_hook(_hook) print( f"✅ Live telemetry enabled (capture={_lt_capture}, window={_lt_window}, dt_ms={_sim_dt_ms})" ) else: print("⚠️ Live telemetry disabled by config") _OperatorBridge_cls = cast(type, _OperatorBridge) operator_bridge = _OperatorBridge_cls( controller=controller, coordinator=_self_org_coordinator, plasticity=_self_org_plasticity, approval_policy=_self_org_approval_policy, telemetry_store=_telemetry_store, ) # Attach the runtime config so the dashboard gate builder can # distinguish "disabled by config" from "config enabled but # component missing" (ERROR). operator_bridge.config_dict = config_dict # type: ignore[attr-defined] print("✅ OperatorBridge created with canonical RuntimeController") _DashboardStateStore_cls = cast(type, _DashboardStateStore) state_store = _DashboardStateStore_cls() from src.dashboard.state import set_dashboard_config set_dashboard_config(config_dict) # Initialen Dashboard-Snapshot bei Tick 0 publizieren # (bevor der erste Tick ausgeführt wird, damit das Dashboard # echte Netzwerkdaten anzeigt und nicht Nullen) from src.core.network import StepResult _initial_result = StepResult( tick=0, spike_ids=(), output_spike_ids=(), spikes_this_tick=0, total_spikes=0, delivered_events=0, queued_events=0, external_injection_count=0, external_total_current=0.0, synaptic_current_targets=0, mean_v=0.0, min_v=0.0, max_v=0.0, mean_energy=0.0, core_step_ms=0.0, neuron_activity={}, total_synapses=network.synapse_count, ) _publish_dashboard_state( state_store=state_store, network=network, result=_initial_result, learning=learning, homeostasis=homeostasis, storage_telemetry=_storage_telemetry, self_org_stats=_self_org_stats_func(), status="idle", ) print("✅ Initial dashboard state published (Tick 0)") except Exception as e: print(f"⚠️ Dashboard setup failed: {type(e).__name__}: {e}") operator_bridge = None state_store = None # ================================================================ # Artifacts (shared by hooks) # ================================================================ artifacts_ctx = RunArtifacts(config_dict) artifacts = artifacts_ctx.__enter__() artifacts.save_topology(health) print( f"🧠 Neurons={len(network.neurons)} Synapses={network.synapse_count} " f"Input={len(network.input_cells)} Output={len(network.output_cells)} " f"Learning={'on' if learning and learning.enabled else 'off'} " f"Homeostasis={'on' if homeostasis and homeostasis.enabled else 'off'}" f"Controller={'idle' if controller else 'none'}" ) # ================================================================ # Dashboard im Hauptthread starten (falls aktiviert) # ================================================================ if ( not args.no_dashboard and _dashboard_available and operator_bridge is not None and state_store is not None ): try: # Write initial snapshot so the heatmap source has real data print("💾 Writing initial .b5d snapshot...") _write_snapshot() docs_root = Path("docs") if Path("docs").exists() else None research_root = Path("research") if Path("research").exists() else None _dashboard_host = args.dashboard_host _dashboard_port = args.dashboard_port print( f"🧠 Starting Brain-5D dashboard on http://{_dashboard_host}:{_dashboard_port}" ) print("⏸️ Simulation starts in idle state. Use dashboard controls to run.") if _serve_dashboard is not None: _serve_dashboard(host=_dashboard_host, port=_dashboard_port, state=state_store, snapshot_path=_snapshot_path, structural_bridge=operator_bridge, docs_root=docs_root, research_root=research_root, chat_settings=cast(dict[str, Any], config_dict.get("research_chat", {}))) # type: ignore[reportOptionalCall, call-arg, operator] except KeyboardInterrupt: print("\n⏹️ Dashboard interrupted, stopping simulation...") finally: if controller is not None: controller.stop() else: # Kein Dashboard – starte Simulation mit konfigurierten Ticks total_ticks = int(config_dict.get("ticks", 1000)) print(f"▶️ Running {total_ticks} ticks (no dashboard)...") if controller is not None: controller.run_ticks(total_ticks) else: print("⚠️ Controller not available, skipping simulation") # Write final snapshot print("💾 Writing final .b5d snapshot...") _write_snapshot() # ================================================================ # Final summary # ================================================================ report = propagation.get_report() summary: dict[str, Any] = { "seed": config_dict.get("seed", 42), "ticks": network.current_tick, "final_neurons": len(network.neurons), "final_synapses": network.synapse_count, "total_spikes": network.total_spikes, "topology": health, "propagation": asdict(report), } if learning and learning.enabled: summary["learning"] = asdict(learning.stats) if homeostasis and homeostasis.enabled: summary["homeostasis"] = asdict(homeostasis.stats) if core_times: ordered = sorted(core_times) p95 = ordered[max(0, int(len(ordered) * 0.95) - 1)] summary["benchmark"] = { "mean_ms": statistics.mean(core_times), "median_ms": statistics.median(core_times), "p95_ms": p95, } print("📊 Benchmark:", summary["benchmark"]) artifacts.save_summary(summary) artifacts_ctx.__exit__(None, None, None) # Final dashboard state publication if state_store is not None: try: from src.dashboard.models import ( DashboardSnapshot, HomeostasisMetrics, LearningMetrics, NetworkMetrics, SpikeMetrics, SystemMetrics, ) learning_stats = learning.stats if learning is not None else None homeo_stats = homeostasis.stats if homeostasis is not None else None final_snapshot = DashboardSnapshot( status="completed", version=BRAIN5D_VERSION_DISPLAY, system=SystemMetrics( tick=network.current_tick, neurons=len(network.neurons), synapses=network.synapse_count, spikes_total=network.total_spikes, ), learning=LearningMetrics( stdp_updates=getattr(learning_stats, "stdp_weight_updates", 0), reward_updates=getattr(learning_stats, "reward_weight_updates", 0), rewards_received=getattr(learning_stats, "rewards_received", 0), rewards_applied=getattr(learning_stats, "rewards_applied", 0), pending_rewards=getattr(learning_stats, "pending_rewards", 0), ), homeostasis=HomeostasisMetrics( enabled=getattr(homeo_stats, "enabled", False), target_rate_hz=getattr(homeo_stats, "target_rate_hz", 0.0), actual_rate_hz=getattr(homeo_stats, "mean_rate_hz", 0.0), rate_error_hz=getattr(homeo_stats, "mean_rate_error_hz", 0.0), mean_rate_hz=getattr(homeo_stats, "mean_rate_hz", 0.0), mean_rate_error_hz=getattr(homeo_stats, "mean_rate_error_hz", 0.0), mean_threshold_adaptation=getattr( homeo_stats, "mean_threshold_adaptation", 0.0 ), target_energy=getattr(homeo_stats, "target_energy", 0.0), mean_energy=getattr(homeo_stats, "mean_energy", 0.0), mean_energy_error=getattr(homeo_stats, "mean_energy_error", 0.0), active_neurons=getattr(homeo_stats, "active_neurons", 0), updates=getattr(homeo_stats, "updates", 0), ), spikes=SpikeMetrics( total_spikes=network.total_spikes, spike_count_last_tick=0, ), network=NetworkMetrics( tick=network.current_tick, neuron_count=len(network.neurons), synapse_count=network.synapse_count, ), ) state_store.publish(final_snapshot) except Exception: pass print("\n📈 Propagation:", report) if learning and learning.enabled: print("📚 Learning:", learning.stats) if homeostasis and homeostasis.enabled: print("⚖️ Homeostasis:", homeostasis.stats) if observatory: print("🔭 Observatory running — close window to exit") observatory.block_until_closed() return 0 if __name__ == "__main__": raise SystemExit(main())