import asyncio from tasks import ATask from ws.message import MsgIo from utils.value import ChangedFloat from components import APid, AStirrer from components.plant import APlant from components.sud import Sud, SudState class SudTask(ATask): def __init__(self, sud: Sud, tc: APid, stirrer: AStirrer, pot: APlant, dt, interval, msg_handler: MsgIo, forecast_estimator=None): ATask.__init__(self, interval) self.sud = sud self.tc = tc self.stirrer = stirrer self.pot = pot # Simulated seconds per tick, vs. interval's wall-clock seconds per # tick - same dt/interval split Pot/TempController/Stirrer already # use internally. Their warp-induced speedup falls out naturally # from ticking real physics at simulated dt; Sud has no physics of # its own, so hold_remaining must be ticked by dt explicitly to get # the same speedup instead of running in real time. self.dt = dt self.msg_handler = msg_handler # Predicts a schedule's actual duration by simulating it with the # same plant/controller machinery this server uses for real - see # components/sud_forecast.py. Optional only so tests/demos that # build a SudTask without one still work; the server always passes # one. self.forecast_estimator = forecast_estimator msg_handler.set_recv_handler(self.recv) def apply_plant_params(self, step): """Keeps the real plant's and the controller's internal model's lumped (M, C) in sync with the step's grain_mass/water_mass, since those vary over the course of a brew (malt going in, water boiling off) - mirrors demo_sud.py's apply_plant_params().""" params = self.sud.derive_plant_params(step.get('grain_mass', 0), step.get('water_mass', 0)) self.pot.set_thermal_params(params['M'], params['C']) if hasattr(self.tc, 'set_model_params'): self.tc.set_model_params(params['M'], params['C']) def apply_stirrer(self, phase): stirrer_cfg = phase.get('stirrer', {}) speed = stirrer_cfg.get('speed', 0) interval_time = stirrer_cfg.get('interval_time', 0) on_ratio = stirrer_cfg.get('on_ratio', 1.0) if interval_time > 0: self.stirrer.set_cycle_time(interval_time) self.stirrer.set_duty_cycle(on_ratio) else: self.stirrer.set_cycle_time(1.0) self.stirrer.set_duty_cycle(1.0 if speed > 0 else 0.0) self.stirrer.set_speed(speed) def on_step_changed(self, step): ramp = step.get('ramp') if step else None hold = step.get('hold') if step else None # Every step ramps to 'temperature' first, then optionally holds - # self.sud.state tells us which phase is currently active; it's # already up to date by the time this callback fires, both on a # full step transition and on the ramp->hold phase switch within # one step (components/sud.py's temp_reached()). ramping = self.sud.state == SudState.RAMPING phase = ramp if ramping else hold if step is not None: self.apply_plant_params(step) if ramping and step['temperature'] is not None: self.tc.set_theta_soll(step['temperature']) self.tc.set_heatrate_soll(ramp['rate']) self.apply_stirrer(phase) asyncio.create_task(self.send({'Step': { 'Index': self.sud.index, 'Type': 'ramp' if ramping else 'hold' if hold is not None else None, 'Descr': step.get('descr') if step else None, 'Temp': step.get('temperature') if step else None, 'Rate': ramp.get('rate') if ramp else None, 'Duration': hold.get('duration') if (hold is not None and not ramping) else None, 'WaitForUser': step.get('user_wait_for_continue', False) if step else None, }})) asyncio.create_task(self.send_remaining_forecast()) def on_state_changed(self, value): asyncio.create_task(self.send({'State': str(value)})) if value in (SudState.DONE, SudState.IDLE): self.stirrer.set_duty_cycle(1.0) self.stirrer.set_speed(0) # A finished/stopped run no longer owns the controller - hand # control back to manual mode (off by default there too). self.tc.set_enabled(False) else: # Any other state (RAMPING/HOLDING/WAIT_USER/PAUSED) means a # run is in progress and needs the controller actively driving # the heater. self.tc.set_enabled(True) def on_user_message_changed(self, value): asyncio.create_task(self.send({'UserMessage': value})) def on_hold_remaining_changed(self, value): asyncio.create_task(self.send({'HoldRemaining': value})) def on_elapsed_changed(self, value): asyncio.create_task(self.send({'Elapsed': value})) async def send_forecast(self, doc): if self.forecast_estimator is None: return # Runs the simulation in a worker thread - it's CPU-bound and can # take a couple hundred ms for a long schedule, which would # otherwise stall every other task (heater, sensor, ...) for that # whole window. loop = asyncio.get_event_loop() t, theta = await loop.run_in_executor(None, self.forecast_estimator.estimate, doc) await self.send({'Forecast': {'T': t, 'Theta': theta}}) def remaining_schedule(self): """The not-yet-done part of the running schedule - the current step's already-finished phase is dropped, and its still-to-run phase is left for SudForecastEstimator to size from the live start_theta/hold_remaining it's seeded with, rather than the step's nominal start. Mirrors client/brewpi_gui.py's Window._remaining_schedule(), which does the same thing client-side for its own (less accurate) naive fallback.""" schedule = self.sud.schedule index = self.sud.index if not schedule or not (0 <= index < len(schedule)): return [] current = schedule[index] rest = schedule[index + 1:] if self.sud.state == SudState.WAIT_USER: return rest if self.sud.state == SudState.HOLDING: # Keep the rest of the current (already fully-resolved) step - # grain_mass/water_mass/etc. - SudForecastEstimator re-resolves # each step against an empty default.step, so a bare # {'hold': {...}} here would leave those None instead of # inherited, breaking derive_plant_params(). The ramp phase is # already done, so drop 'temperature' (not 'ramp' - every step # always carries a 'ramp' dict now, dropping it would leave # 'rate' missing the moment this synthetic step is re-resolved): # without a 'temperature' of its own, the rebuilt step pushes # no new target, start_theta already seeds the simulation at # the target anyway, so it resolves out of its (now synthetic) # ramp phase in a tick or two regardless. Override the hold's # duration to what's actually left. remaining_minutes = max(self.sud.hold_remaining, 0.0) / 60.0 synthetic = {k: v for k, v in current.items() if k != 'temperature'} synthetic['hold'] = {**current.get('hold', {}), 'duration': remaining_minutes} return [synthetic] + rest return [current] + rest async def send_remaining_forecast(self): """Like send_forecast(), but for the remaining steps only, seeded from the live current temperature - the dynamic forecast's projected-remainder line (GUI's SudForecastPlot.line_projected) uses this instead of its own naive abs(delta)/rate estimate whenever it's available.""" if self.forecast_estimator is None: return remaining = self.remaining_schedule() if not remaining: return doc = { 'Name': self.sud.name, 'Description': self.sud.description, 'pot_mass': self.sud.pot_mass, 'pot_material': self.sud.pot_material, 'steps': remaining, } start_theta = self.tc.get_theta_ist() loop = asyncio.get_event_loop() t, theta = await loop.run_in_executor(None, self.forecast_estimator.estimate, doc, start_theta) await self.send({'RemainingForecast': {'T': t, 'Theta': theta}}) async def recv(self, data): for pair in data.items(): if 'Start' in pair[0]: self.sud.start() elif 'Confirm' in pair[0]: self.sud.confirm() elif 'Pause' in pair[0]: self.sud.pause() elif 'Stop' in pair[0]: self.sud.stop() elif 'Save' in pair[0]: doc = self.sud.save() await self.send({'Json': doc}) await self.send_forecast(doc) elif 'Load' in pair[0]: if self.sud.load(pair[1]): await self.send({'Name': self.sud.name, 'Description': self.sud.description}) await self.send({'Json': pair[1]}) await self.send_forecast(pair[1]) else: # Sud.load() refuses while a run is in progress (state # not IDLE/DONE) - tell the client why instead of # silently dropping the request. The client has no # business pre-emptively guessing this itself from its # own (replicated, laggy) view of the state. # # Immediately cleared back to None - unlike every other # field here, this is a one-shot event, not state. The # dispatcher has no concept of "don't persist this into # global_state" (see ws/user.py's update()), so without # clearing it, any client connecting later - even one # that never touched Load - would get this stale error # replayed on connect, with nothing it just did to # explain why. await self.send({'Error': 'Cannot load a new schedule while a run is in progress - stop it first.'}) await self.send({'Error': None}) async def send(self, data): await self.msg_handler.send(data) async def on_process(self): print("{}: Started with interval {} s".format(self.msg_handler.get_key(), self.interval)) self.sud.set_on_changed('step', self.on_step_changed) self.sud.set_on_changed('state', self.on_state_changed) self.sud.set_on_changed('user_message', self.on_user_message_changed) self.sud.set_on_changed('hold_remaining', ChangedFloat(self.on_hold_remaining_changed, prec=0).set) self.sud.set_on_changed('elapsed', ChangedFloat(self.on_elapsed_changed, prec=1).set) asyncio.create_task(self.send({'Name': self.sud.name, 'Description': self.sud.description})) while True: if self.sud.state == SudState.RAMPING and self.tc.is_holding(): self.sud.temp_reached() self.sud.tick(self.dt) await asyncio.sleep(self.interval)