get map stream -> instrument from servicemanager for feeder

This commit is contained in:
2025-05-05 08:32:34 +02:00
parent d683306beb
commit 1ed579f1cb
3 changed files with 10 additions and 3 deletions

View File

@ -25,9 +25,14 @@ def main(dbname=None):
db = SEHistory(dbname, access='write')
db.enable_write_access()
host = socket.gethostname().split('.')[0]
fm = FrappyManager()
fm.get_info()
host = socket.gethostname().split('.')[0]
# create map to get instrument from internal stream uri
insmap = db.instrument_by_stream
for ins, ports in fm.info.items():
for p in ports.values():
insmap[f'{host}:{p}'] = ins
cfginfo = {}
for ins, procs in fm.get_procs(cfginfo=cfginfo).items():
for service in procs:

View File

@ -35,6 +35,7 @@ class SEHistory(InfluxDBWrapper):
parser.optionxform = str
parser.read([Path('~/.config/sehistory').expanduser()])
section = parser[dbname] if dbname else parser[parser.sections()[0]]
self.instrument_by_stream = {}
super().__init__(*(section[k] for k in ('uri', 'token', 'org', 'bucket')), access=access)
def curves(self, start=None, stop=None, measurement=('*.value', '*.target'), field='float',
@ -407,4 +408,6 @@ class SEHistory(InfluxDBWrapper):
self.flush()
def add_stream(self, value, tags, key, ts):
if value == '': # unknown instrument
value = self.instrument or self.instrument_by_stream.get(key, '0')
self.set_instrument(key, value, ts, **tags)

View File

@ -254,9 +254,8 @@ class EventStream:
# note: a stream with buffered content might not be ready to emit any event, because
# of filtering
def __init__(self, *udp, instrument=None, **streams):
def __init__(self, *udp, **streams):
self.streams = streams
self.instrument = instrument
self.udp = {v.socket.fileno(): v for v in udp}
def wait_ready(self, timeout):