Skip to content

Commit da4e6fe

Browse files
authored
Stateless drivers
* Make default devices file. * Move setup into tomato.daemon instead of cli. * Fix reload tests. * Implement DeviceFile model. * Remove int channels. * Modify test files channels int->str * Stop tomato.init from overwriting files. * Add DeviceFile to Daemon, remove devs * Touch up tests for devices. * Rework drivers on Daemon. * Ruff. * Fix typos.
1 parent b3e149f commit da4e6fe

24 files changed

Lines changed: 762 additions & 626 deletions

MANIFEST.in

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,2 @@
11
include versioneer.py
22
include src/tomato/_version.py
3-
include src/tomato/data/*

src/tomato/daemon/__init__.py

Lines changed: 14 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -6,19 +6,17 @@
66
77
"""
88

9-
import logging
109
import argparse
11-
from pathlib import Path
12-
from threading import Thread
13-
import toml
10+
import logging
1411
import time
15-
import zmq
16-
17-
from tomato.models import Reply, Daemon
1812
import tomato.daemon.cmd as cmd
19-
import tomato.daemon.job
2013
import tomato.daemon.driver
2114
import tomato.daemon.io as io
15+
import tomato.daemon.job
16+
import zmq
17+
from pathlib import Path
18+
from threading import Thread
19+
from tomato.models import Daemon, Reply
2220

2321
logger = logging.getLogger(__name__)
2422

@@ -30,7 +28,7 @@ def setup_logging(daemon: Daemon):
3028
"""
3129
logdir = Path(daemon.settings["logdir"])
3230
logdir.mkdir(parents=True, exist_ok=True)
33-
logfile = logdir / f"daemon_{daemon.port}.log"
31+
logfile = logdir / f"tomato_daemon_{daemon.port}.log"
3432
logging.basicConfig(
3533
level=daemon.verbosity,
3634
format="%(asctime)s - %(levelname)8s - %(name)-30s - %(message)s",
@@ -53,12 +51,13 @@ def tomato_daemon():
5351
parser.add_argument("--appdir", "-A", type=str, default=str(Path.cwd()))
5452

5553
args = parser.parse_args()
56-
settings = toml.load(Path(args.appdir) / "settings.toml")
5754

58-
daemon = Daemon(**vars(args), status="bootstrap", settings=settings)
55+
daemon = Daemon(**vars(args), status="bootstrap")
5956
setup_logging(daemon)
6057
logger.info("logging set up with verbosity %s", daemon.verbosity)
6158

59+
# TODO: setup should not be a thing really.
60+
cmd.setup(msg={}, daemon=daemon)
6261
logger.debug("attempting to restore daemon state")
6362
io.load(daemon)
6463

@@ -71,10 +70,10 @@ def tomato_daemon():
7170

7271
logger.debug("entering main loop")
7372
jmgr = Thread(target=tomato.daemon.job.manager, args=(daemon.port,), daemon=True)
74-
jmgr.do_run = True
73+
setattr(jmgr, "do_run", True)
7574
jmgr.start()
7675
dmgr = Thread(target=tomato.daemon.driver.manager, args=(daemon.port,), daemon=True)
77-
dmgr.do_run = True
76+
setattr(dmgr, "do_run", True)
7877
dmgr.start()
7978
t0 = time.process_time()
8079
while True:
@@ -93,9 +92,9 @@ def tomato_daemon():
9392
rep.send_pyobj(ret)
9493
if daemon.status == "stop":
9594
for mgr, label in [(jmgr, "job"), (dmgr, "driver")]:
96-
if mgr is not None and mgr.do_run:
95+
if mgr is not None and getattr(mgr, "do_run"):
9796
logger.debug("stopping %s manager thread", label)
98-
mgr.do_run = False
97+
setattr(mgr, "do_run", False)
9998
if jmgr is not None:
10099
jmgr.join(1e-3)
101100
if not jmgr.is_alive():

src/tomato/daemon/cmd.py

Lines changed: 97 additions & 62 deletions
Original file line numberDiff line numberDiff line change
@@ -12,21 +12,22 @@
1212
1313
"""
1414

15+
import logging
16+
import tomato.daemon.io as io
17+
import tomato.daemon.jobdb as jobdb
18+
import tomato.utils
19+
from pathlib import Path
20+
from pydantic import BaseModel
1521
from tomato.models import (
22+
Component,
1623
Daemon,
17-
Driver,
1824
Device,
19-
Reply,
20-
Pipeline,
2125
Job,
22-
Component,
26+
Pipeline,
27+
Reply,
28+
SpawnData,
2329
)
24-
from pydantic import BaseModel
2530
from typing import Any
26-
import logging
27-
28-
import tomato.daemon.io as io
29-
import tomato.daemon.jobdb as jobdb
3031

3132
logger = logging.getLogger(__name__)
3233

@@ -44,7 +45,7 @@ def merge_pipelines(
4445
if pip.jobid is not None:
4546
ret[pname] = pip
4647
else:
47-
if pip.devs == new[pname].devs:
48+
if pip.components == new[pname].components:
4849
ret[pname] = pip
4950
elif pip.jobid is None:
5051
ret[pname] = new[pname]
@@ -74,27 +75,68 @@ def stop(msg: dict, daemon: Daemon) -> Reply:
7475

7576
def setup(msg: dict, daemon: Daemon) -> Reply:
7677
logger = logging.getLogger(f"{__name__}.setup")
77-
logger.debug("%s", msg)
78+
79+
# TODO: Rework this!
80+
devicefile = tomato.utils.load_device_file(
81+
Path(daemon.settings["devices"]["config"]),
82+
logger,
83+
)
84+
logger.debug(f"{devicefile=}")
85+
devs = {dev["name"]: Device(**dev) for dev in devicefile["devices"]}
86+
logger.debug(f"{devs=}")
87+
pips, cmps = tomato.utils.get_pipelines(
88+
devs,
89+
devicefile["pipelines"],
90+
logger,
91+
)
92+
logger.debug(f"{pips=}")
93+
logger.debug(f"{cmps=}")
94+
7895
if daemon.status == "bootstrap":
79-
for key in ["drvs", "devs", "pips", "cmps"]:
80-
setattr(daemon, key, msg[key])
96+
for key, val in [
97+
# ("drvs", drvs),
98+
# ("devs", devs),
99+
("pips", pips),
100+
("cmps", cmps),
101+
]:
102+
setattr(daemon, key, val)
103+
104+
daemon.drivers = {
105+
key: SpawnData(name=key) for key in daemon.devicefile.drivers.keys()
106+
}
81107
logger.info("setup successful with pipelines: '%s'", daemon.pips.keys())
82108
daemon.status = "running"
83109
else:
110+
try:
111+
nd = Daemon(
112+
status=daemon.status,
113+
port=daemon.port,
114+
appdir=daemon.appdir,
115+
verbosity=daemon.verbosity,
116+
)
117+
except Exception as e:
118+
logger.critical("Error", exc_info=e)
119+
return Reply(
120+
success=False,
121+
msg="could not parse updated settings",
122+
)
123+
logger.debug(f"{nd=}")
124+
ndf = nd.devicefile
84125
# First, check that we're not touching anything associated with a running job
85126
check_components = set()
86127
check_devices = set()
87128
check_drivers = set()
88129
for dpip in daemon.pips.values():
130+
logger.debug(f"{dpip=}")
89131
if dpip.jobid is None:
90132
continue
91-
if dpip.name not in msg["pips"]:
133+
if dpip.name not in pips:
92134
return Reply(
93135
success=False,
94136
msg="reload would delete a running pipeline",
95137
data=dpip,
96138
)
97-
pip = msg["pips"][dpip.name]
139+
pip = pips[dpip.name]
98140
if pip.components != dpip.components:
99141
return Reply(
100142
success=False,
@@ -105,13 +147,13 @@ def setup(msg: dict, daemon: Daemon) -> Reply:
105147

106148
for cname in check_components:
107149
dcomp = daemon.cmps[cname]
108-
if cname not in msg["cmps"]:
150+
if cname not in cmps:
109151
return Reply(
110152
success=False,
111153
msg="reload would delete a component of a running pipeline",
112154
data=dcomp,
113155
)
114-
comp = msg["cmps"][cname]
156+
comp = cmps[cname]
115157
if (
116158
dcomp.name != comp.name
117159
or dcomp.driver != comp.driver
@@ -128,54 +170,37 @@ def setup(msg: dict, daemon: Daemon) -> Reply:
128170
check_devices.add(dcomp.device)
129171
check_drivers.add(dcomp.driver)
130172

131-
for dname in check_devices:
132-
ddev = daemon.devs[dname]
133-
if dname not in msg["devs"]:
134-
return Reply(
135-
success=False,
136-
msg="reload would delete a device of a component in a running pipeline",
137-
data=ddev,
138-
)
139-
dev = msg["devs"][dname]
140-
if (
141-
ddev.name != dev.name
142-
or ddev.driver != dev.driver
143-
or ddev.address != dev.address
144-
or ddev.pollrate != dev.pollrate
145-
or any(ch not in dev.channels for ch in ddev.channels)
146-
):
147-
return Reply(
148-
success=False,
149-
msg="reload would modify a device of a component in a running pipeline",
150-
data=ddev,
151-
)
152-
153173
for dname in check_drivers:
154-
ddrv = daemon.drvs[dname]
155-
if dname not in msg["drvs"]:
174+
if dname not in ndf.drivers:
156175
return Reply(
157176
success=False,
158177
msg="reload would delete a driver of a device in a running pipeline",
159-
data=ddev,
178+
data=daemon.drivers[dname],
160179
)
161-
drv = msg["drvs"][dname]
162-
if ddrv.name != drv.name or ddrv.settings != drv.settings:
180+
181+
if daemon.devicefile.drivers[dname].settings != ndf.drivers[dname].settings:
163182
return Reply(
164183
success=False,
165184
msg="reload would modify a driver of a device in a running pipeline",
166-
data=ddrv,
185+
data=daemon.devicefile.drivers[dname].settings,
167186
)
168187

169-
_api_reload(msg["drvs"], daemon.drvs, "driver", ["settings"])
170-
171-
_api_reload(msg["pips"], daemon.pips, "pipeline", ["components"])
188+
logger.critical("goint into api reload")
172189

190+
_api_reload(pips, daemon.pips, "pipeline", ["components"])
173191
attrlist = ["driver", "device", "address", "channel", "role"]
174-
_api_reload(msg["cmps"], daemon.cmps, "component", attrlist)
175-
176-
_api_reload(msg["devs"], daemon.devs, "device", ["channels", "pollrate"])
177-
192+
_api_reload(cmps, daemon.cmps, "component", attrlist)
193+
# Add new drivers, they will be spawned by driver.manager
194+
for dname in ndf.drivers.keys():
195+
if dname not in daemon.drivers:
196+
logger.info("adding new driver '%s'", dname)
197+
daemon.drivers[dname] = SpawnData(name=dname)
198+
199+
# We want to trigger re-parse of config files on daemon
200+
daemon.settings = nd.settings
201+
daemon.devicefile = ndf
178202
logger.info("reload successful with pipelines: '%s'", daemon.pips.keys())
203+
179204
return Reply(success=True, data=daemon)
180205

181206

@@ -201,7 +226,7 @@ def pipeline(msg: dict, daemon: Daemon) -> Reply:
201226
logger.debug("%s", msg)
202227
pip = msg["params"]
203228
if pip["name"] is None:
204-
logger.error()
229+
logger.error("no pipeline name supplied")
205230
return Reply(success=False, msg="no pipeline name supplied", data=msg)
206231
if pip["name"] not in daemon.pips:
207232
dest = Pipeline(**pip)
@@ -243,15 +268,6 @@ def get_jobs(msg: dict, daemon: Daemon) -> Reply:
243268
return Reply(success=True, msg=f"found {len(jobs)} jobs", data=jobs)
244269

245270

246-
def driver(msg: dict, daemon: Daemon) -> Reply:
247-
return _api(
248-
otype="driver",
249-
msg=msg,
250-
ddict=daemon.drvs,
251-
Cls=Driver,
252-
)
253-
254-
255271
def device(msg: dict, daemon: Daemon) -> Reply:
256272
return _api(
257273
otype="device",
@@ -293,3 +309,22 @@ def _api(otype: str, msg: dict, ddict: dict[str, Any], Cls: BaseModel) -> Reply:
293309
msg=f"{otype} {obj['name']!r} updated",
294310
data=ddict[obj["name"]],
295311
)
312+
313+
314+
def reload(msg: dict, daemon: Daemon, **kwargs: dict) -> Reply:
315+
# daemon.settings = toml.load(Path(daemon.appdir) / "settings.toml")
316+
return Reply(success=True, msg="daemon settings reloaded", data=daemon.settings)
317+
318+
319+
def driver_set(msg: dict, daemon: Daemon) -> Reply:
320+
params = msg.pop("params")
321+
name = params.pop("name")
322+
for k, v in params.items():
323+
setattr(daemon.drivers[name], k, v)
324+
return Reply(success=True, msg=f"updated driver {name}", data=daemon.drivers[name])
325+
326+
327+
def driver_del(msg: dict, daemon: Daemon) -> Reply:
328+
name = msg.pop("params").pop("name")
329+
del daemon.drivers[name]
330+
return Reply(success=True, msg=f"removed driver {name}")

0 commit comments

Comments
 (0)