""" hub.py — the traffic cop and the brain. Everything (your hand, the room, a satellite, bitcoin) fires OSC at this hub on 127.0.0.1:9000. It doesn't care where a signal came from — a fingertip and a plane over Tokyo are the same kind of thing here. The hub: 1. receives every OSC message into a thread-safe raw state, 2. on a fixed control-rate clock, normalizes + One-Euro-filters each signal to a clean 0..1 (events become decaying impulses), 3. runs the modulation matrix (matrix.py) to turn signals into destinations, 4. emits results as OSC /out/ (for Ableton / TouchDesigner), optional MIDI CC/notes, and a live WebSocket feed for the browser visualizer. Run: python hub.py (uses config.json) python hub.py --config myset.json Ports: OSC in 9000 | OSC out 9001 | WebSocket 8765 | HTTP viz 8088 """ from __future__ import annotations import argparse import asyncio import json import math import os import queue import threading import time from functools import partial from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer import websockets from pythonosc.dispatcher import Dispatcher from pythonosc.osc_server import ThreadingOSCUDPServer from pythonosc.udp_client import SimpleUDPClient from normalize import AdaptiveNormalizer, OneEuroFilter, ImpulseEnvelope from matrix import ModMatrix from transform import TweakBank HERE = os.path.dirname(os.path.abspath(__file__)) START = time.monotonic() # sidereal orbital periods in Earth-years — a group assigned a planet cycles at # this rate relative to `orbit_base_seconds` (one Earth orbit = base seconds). PLANET_PERIODS = { "mercury": 0.2408, "venus": 0.6152, "earth": 1.0, "mars": 1.8808, "jupiter": 11.862, "saturn": 29.457, "uranus": 84.02, "neptune": 164.8, "pluto": 248.0, } def now() -> float: return time.monotonic() - START # --------------------------------------------------------------------------- # Per-signal processing chain # --------------------------------------------------------------------------- class Signal: def __init__(self, name: str, spec: dict): self.name = name self.type = spec.get("type", "continuous") self.label = spec.get("label", name) nrm = spec.get("norm", {}) mode = nrm.get("mode", "minmax") fixed = None if "lo" in nrm and "hi" in nrm: fixed = (nrm["lo"], nrm["hi"]) mode = "fixed" self.norm = AdaptiveNormalizer(nrm.get("window", 600), mode, fixed) flt = spec.get("filter", {}) self.filter = OneEuroFilter(flt.get("min_cutoff", 1.2), flt.get("beta", 0.02)) self.env = ImpulseEnvelope(spec.get("attack", 0.01), spec.get("decay", 0.6)) self.raw = 0.0 self.value = 0.0 def update_continuous(self, raw: float, t: float) -> float: self.raw = raw n = self.norm(raw) self.value = self.filter(n, t) return self.value def trigger(self, magnitude: float, t: float): self.raw = magnitude self.env.trigger(self.norm(magnitude), t) def event_value(self, t: float) -> float: self.value = self.env.value(t) return self.value # --------------------------------------------------------------------------- # The hub # --------------------------------------------------------------------------- class Hub: def __init__(self, cfg: dict): self.cfg = cfg self.specs = cfg.get("signals", {}) self.matrix = ModMatrix.from_config(cfg) self.tweaks = TweakBank(cfg) self.rate = cfg.get("control_rate_hz", 60) self.signals: dict[str, Signal] = {} self._raw_lock = threading.Lock() self._raw_inbox: dict[str, tuple[float, float]] = {} # key -> (val,t) self._event_q: "queue.Queue[tuple[str,float,float]]" = queue.Queue() self.out = SimpleUDPClient(cfg.get("osc_out_host", "127.0.0.1"), cfg.get("osc_out_port", 9001)) self.warp = cfg.get("warp", {}) self._jog_clients: dict[int, SimpleUDPClient] = {} self.orbit_base = cfg.get("orbit_base_seconds", 60) self.midi = _open_midi(cfg.get("midi", {})) self.clients: set = set() self.recent_events: list[dict] = [] self._last_notes: dict[str, int] = {} self.route_meta = [ {"source": r.get("source"), "dest": r.get("dest"), "label": r.get("label", "")} for r in cfg.get("routes", []) ] # ---- OSC ingest (runs in the OSC server thread) -------------------- def _osc_handler(self, address: str, *args): if not args: return try: val = float(args[0]) except (TypeError, ValueError): return # /gs/a/b -> "a.b" ; /gs/quake/event -> event on "quake.event" parts = address.strip("/").split("/") if parts and parts[0] == "gs": parts = parts[1:] key = ".".join(parts) t = now() if key.endswith(".event"): self._event_q.put((key, val, t)) else: with self._raw_lock: self._raw_inbox[key] = (val, t) def _ensure_signal(self, key: str) -> Signal: sig = self.signals.get(key) if sig is None: spec = self.specs.get(key, {}) if key.endswith(".event") and "type" not in spec: spec = {**spec, "type": "event"} sig = Signal(key, spec) self.signals[key] = sig return sig # ---- control-rate loop (runs in asyncio) -------------------------- async def control_loop(self): period = 1.0 / self.rate while True: t = now() # drain events -> triggers while True: try: key, mag, et = self._event_q.get_nowait() except queue.Empty: break sig = self._ensure_signal(key) sig.trigger(mag, et) ev = {"src": key.rsplit(".", 1)[0], "mag": round(mag, 3), "t": et} self.recent_events.append(ev) self.recent_events = self.recent_events[-40:] # snapshot raw inbox with self._raw_lock: inbox = dict(self._raw_inbox) sources: dict[str, float] = {} wire_sources: dict[str, dict] = {} for key, (raw, rt) in inbox.items(): sig = self._ensure_signal(key) v = sig.update_continuous(raw, t) sources[key] = v wire_sources[key] = {"raw": round(raw, 4), "norm": round(v, 4), "label": sig.label, "muted": self.tweaks.is_muted(key)} # event envelopes decay every tick for key, sig in self.signals.items(): if sig.type == "event": v = sig.event_value(t) sources[key] = v wire_sources[key] = {"raw": round(sig.raw, 4), "norm": round(v, 4), "label": sig.label, "event": True, "muted": self.tweaks.is_muted(key)} # planetary orbits — each orbit-group's gain waxes/wanes at its # planet's real relative rate (Mercury whips, Neptune barely drifts) for g, planet in self.tweaks.group_orbit.items(): orb = PLANET_PERIODS.get(planet, 1.0) * self.orbit_base phase = (t / orb) % 1.0 self.tweaks.set_orbit_lfo( g, 0.25 + 0.75 * (0.5 + 0.5 * math.sin(2 * math.pi * phase))) # scrub any time-warped replay workers from live controls self._forward_jog(sources) # performance layer: live controls -> tweaks -> groups -> matrix self.tweaks.apply_controls(sources) tweaked = self.tweaks.transform(sources) dests, active = self.matrix.evaluate(tweaked) self._emit(dests) await self._broadcast(t, wire_sources, dests, active, self.tweaks.macros()) await asyncio.sleep(period) # ---- time-warp jog forwarding ------------------------------------- def _forward_jog(self, sources: dict[str, float]): for w in self.warp.values(): port = w.get("port") if not port: continue cl = self._jog_clients.get(port) if cl is None: cl = SimpleUDPClient("127.0.0.1", port) self._jog_clients[port] = cl rs, ns = w.get("rate_src"), w.get("nudge_src") if rs and rs in sources: cl.send_message("/jog/rate", float(sources[rs])) if ns and ns in sources: cl.send_message("/jog/nudge", float(sources[ns])) # ---- outputs ------------------------------------------------------- def _emit(self, dests: dict[str, float]): midi_map = self.cfg.get("midi", {}).get("map", {}) for dest, val in dests.items(): self.out.send_message(f"/out/{dest}", float(val)) m = midi_map.get(dest) if m and self.midi: if m.get("note") and val >= 40: # quantized note number note = int(val) if self._last_notes.get(dest) != note: if self._last_notes.get(dest): self.midi.send_message([0x80 | m.get("channel", 0), self._last_notes[dest], 0]) self.midi.send_message([0x90 | m.get("channel", 0), note, 100]) self._last_notes[dest] = note elif "cc" in m: cc = int(_clamp01(val) * 127) self.midi.send_message([0xB0 | m.get("channel", 0), m["cc"], cc]) async def _broadcast(self, t, sources, dests, active, macros=None): if not self.clients: return routes = [] for i, r in enumerate(self.matrix.routes): routes.append({"source": r.get("source"), "dest": r.get("dest"), "label": r.get("label", ""), "amount": round(r.get("amount", 0), 3), "active": round(active.get(i, 0.0), 4)}) payload = json.dumps({ "t": round(t, 3), "sources": sources, "dests": {k: (round(v, 4) if isinstance(v, float) else v) for k, v in dests.items()}, "routes": routes, "macros": macros or {}, "events": self.recent_events[-12:], }) websockets.broadcast(self.clients, payload) async def ws_handler(self, ws): self.clients.add(ws) try: await ws.send(json.dumps({ "hello": "godstrument", "routes": [{"source": r.get("source"), "dest": r.get("dest"), "label": r.get("label", ""), "amount": round(r.get("amount", 0), 3)} for r in self.matrix.routes], "groups": self.cfg.get("groups", {}), "midi": self.cfg.get("midi", {}).get("map", {})})) async for msg in ws: # browser -> hub commands try: self._handle_cmd(json.loads(msg)) except Exception: # noqa pass finally: self.clients.discard(ws) def _handle_cmd(self, m: dict): cmd = m.get("cmd") if cmd == "mute_toggle": self.tweaks.toggle(m.get("target"), "mute") elif cmd == "toggle": self.tweaks.toggle(m.get("target"), m.get("param", "mute")) elif cmd == "set": self.tweaks.set_param(m.get("target"), m.get("param", "gain"), float(m.get("value", 0.0))) elif cmd == "add_route": self.matrix.add_route(m.get("source"), m.get("dest"), float(m.get("amount", 0.5))) elif cmd == "remove_route": self.matrix.remove_route(m.get("source"), m.get("dest")) elif cmd == "set_route": self.matrix.set_route_amount(m.get("source"), m.get("dest"), float(m.get("amount", 0.5))) elif cmd == "make_group": self.tweaks.add_group(m.get("name"), m.get("members", [])) # ---- run ----------------------------------------------------------- async def run(self): # OSC server thread disp = Dispatcher() disp.set_default_handler(self._osc_handler) osc_srv = ThreadingOSCUDPServer( (self.cfg.get("osc_in_host", "127.0.0.1"), self.cfg.get("osc_in_port", 9000)), disp) threading.Thread(target=osc_srv.serve_forever, daemon=True).start() # HTTP static server for the visualizer _serve_http(os.path.join(HERE, "viz"), self.cfg.get("http_port", 8080)) ws_port = self.cfg.get("ws_port", 8765) print("╔══════════════════════════════════════════════╗") print("║ GODSTRUMENT hub online ") print(f"║ OSC in udp/{self.cfg.get('osc_in_port',9000)} " f"OSC out udp/{self.cfg.get('osc_out_port',9001)}") print(f"║ WebSocket ws/{ws_port}") print(f"║ Visualizer -> http://localhost:" f"{self.cfg.get('http_port',8080)}") print(f"║ MIDI out: {'connected' if self.midi else 'off (optional)'}") print(f"║ {len(self.route_meta)} routes patched, " f"clock {self.rate}Hz") print("╚══════════════════════════════════════════════╝") async with websockets.serve(self.ws_handler, "127.0.0.1", ws_port): await self.control_loop() def _open_midi(midi_cfg: dict): if not midi_cfg.get("enabled"): return None try: import rtmidi out = rtmidi.MidiOut() ports = out.get_ports() target = midi_cfg.get("port_name", "IAC") idx = next((i for i, p in enumerate(ports) if target.lower() in p.lower()), None) if idx is None and ports: idx = 0 if idx is None: out.open_virtual_port("Godstrument") else: out.open_port(idx) print(f" MIDI -> {ports[idx]}") return out except Exception as e: # noqa print(f" MIDI unavailable ({e}); continuing without it.") return None def _serve_http(directory: str, port: int): handler = partial(SimpleHTTPRequestHandler, directory=directory) httpd = ThreadingHTTPServer(("127.0.0.1", port), handler) httpd.RequestHandlerClass.log_message = lambda *a, **k: None threading.Thread(target=httpd.serve_forever, daemon=True).start() def _clamp01(x: float) -> float: return 0.0 if x < 0.0 else 1.0 if x > 1.0 else x def main(): ap = argparse.ArgumentParser() ap.add_argument("--config", default=os.path.join(HERE, "config.json")) args = ap.parse_args() with open(args.config) as f: cfg = json.load(f) hub = Hub(cfg) try: asyncio.run(hub.run()) except KeyboardInterrupt: print("\n godstrument sleeps.") if __name__ == "__main__": main()