diff --git a/tconnectsync/sync/tandemsource/autoupdate.py b/tconnectsync/sync/tandemsource/autoupdate.py index 4085ee2..c3b1aec 100644 --- a/tconnectsync/sync/tandemsource/autoupdate.py +++ b/tconnectsync/sync/tandemsource/autoupdate.py @@ -28,7 +28,7 @@ class TandemSourceAutoupdate: def process(self, tconnect, nightscout, time_start, time_end, pretend, features=None): if features is None: features = DEFAULT_FEATURES - + # Read from android api, find exact interval to cut down on API calls # Refresh API token. If failure, die, have wrapper script re-run. @@ -72,7 +72,7 @@ class TandemSourceAutoupdate: if pretend: logger.info('Would update now if not in pretend mode') else: - added = ProcessTimeRange(tconnect, nightscout, tconnectDevice['tconnectDeviceId'], pretend, features=features).process(time_start, time_end) + added = ProcessTimeRange(tconnect, nightscout, tconnectDevice, pretend, features=features).process(time_start, time_end) logger.info('Added %d items from ProcessTimeRange' % added) # Track the time it took to find a new event between runs, @@ -127,12 +127,12 @@ class TandemSourceAutoupdate: # If it's been 3 loops since the last time we found new data, # then we're not in sync with the rate at which pump data is being - # uploaded, so + # uploaded, so if len(self.time_diffs_between_attempts) >= 3: # The pump hasn't sent us data that, based on previous cadence, we were expecting - logger.warning(AutoupdateNoIndexChangeWarning("Sleeping %d seconds after unexpected no index change based on previous cadence. (New data might be delayed.)" % + logger.warning(AutoupdateNoIndexChangeWarning("Sleeping %d seconds after unexpected no index change based on previous cadence. (New data might be delayed.)" % int(self.secret.AUTOUPDATE_UNEXPECTED_NO_INDEX_SLEEP_SECONDS))) - + logger.debug("Last event time: %s, time diffs between attempts: %s" % (self.last_event_time, self.time_diffs_between_attempts)) time.sleep(self.secret.AUTOUPDATE_UNEXPECTED_NO_INDEX_SLEEP_SECONDS) @@ -154,7 +154,7 @@ class TandemSourceAutoupdate: if len(self.time_diffs_between_updates) > 10: self.time_diffs_between_updates = self.time_diffs_between_updates[1:] - # If we have less than 3 data points, + # If we have less than 3 data points, if len(self.time_diffs_between_updates) > 2: sleep_secs = sum(self.time_diffs_between_updates) / len(self.time_diffs_between_updates) diff --git a/tconnectsync/sync/tandemsource/process.py b/tconnectsync/sync/tandemsource/process.py index 60e2941..e53fd6d 100644 --- a/tconnectsync/sync/tandemsource/process.py +++ b/tconnectsync/sync/tandemsource/process.py @@ -1,19 +1,33 @@ import logging +import collections from ...features import DEFAULT_FEATURES from ...eventparser.generic import Events, decode_raw_events, EVENT_LEN from ...eventparser import events as eventtypes +from ...domain.tandemsource.event_class import EventClass +from .process_basal import ProcessBasal +from .process_basal_suspension import ProcessBasalSuspension +from .process_basal_resume import ProcessBasalResume +from .process_alarm import ProcessAlarm logger = logging.getLogger(__name__) class ProcessTimeRange: - def __init__(self, tconnect, nightscout, tconnect_device_id, pretend, features=DEFAULT_FEATURES): + def __init__(self, tconnect, nightscout, tconnectDevice, pretend, features=DEFAULT_FEATURES): self.tconnect = tconnect self.nightscout = nightscout - self.tconnect_device_id = tconnect_device_id + self.tconnect_device_id = tconnectDevice['tconnectDeviceId'] + self.max_date_with_events = tconnectDevice['maxDateWithEvents'] self.pretend = pretend self.features = features - + + event_classes = { + EventClass.BASAL.name: ProcessBasal, + EventClass.BASAL_SUSPENSION.name: ProcessBasalSuspension, + EventClass.BASAL_RESUME.name: ProcessBasalResume, + EventClass.ALARM.name: ProcessAlarm + } + def process(self, time_start, time_end): logger.info(f"ProcessTimeRange {time_start=} {time_end=} {self.tconnect_device_id=} {self.features=}") @@ -22,11 +36,32 @@ class ProcessTimeRange: logger.info(f"Read {len(pump_events_decoded)=} (est. {len(pump_events_decoded)/EVENT_LEN} events)") events = Events(pump_events_decoded) - processed_count = 0 + events_first_time = None + events_last_time = None + for_eventclass = collections.defaultdict(list) for event in events: - if isinstance(event, eventtypes.LidBasalRateChange): - logger.debug(f"Found {event=}") - + if not events_first_time: + events_first_time = event.eventTimestamp + if not events_last_time: + events_last_time = event.eventTimestamp + events_first_time = min(events_first_time, event.eventTimestamp) + events_last_time = max(events_last_time, event.eventTimestamp) + + clazz = EventClass.for_event(event) + if clazz: + processed_count += 1 + for_eventclass[clazz.name].append(event) + + count_by_eventclass = {k: len(v) for k,v in for_eventclass.items()} + logger.info(f"Found events: {count_by_eventclass}") + + processed_count = 0 + for clazz, events in for_eventclass.items(): + if clazz in self.event_classes.keys(): + c = self.event_classes[clazz](self.tconnect, self.nightscout, self.tconnect_device_id, self.pretend, self.features) + ns_entries = c.process(events, events_first_time, events_last_time) + processed_count += c.write(ns_entries) + return processed_count - \ No newline at end of file + diff --git a/tconnectsync/sync/tandemsource/process_alarm.py b/tconnectsync/sync/tandemsource/process_alarm.py new file mode 100644 index 0000000..ec97e95 --- /dev/null +++ b/tconnectsync/sync/tandemsource/process_alarm.py @@ -0,0 +1,67 @@ +import logging +import arrow + +from ...features import DEFAULT_FEATURES +from ...eventparser.generic import Events, decode_raw_events, EVENT_LEN +from ...eventparser.utils import bitmask_to_list +from ...eventparser import events as eventtypes +from ...domain.tandemsource.event_class import EventClass +from ...parser.nightscout import ( + ALARM_EVENTTYPE, + NightscoutEntry +) + +logger = logging.getLogger(__name__) + +class ProcessAlarm: + def __init__(self, tconnect, nightscout, tconnect_device_id, pretend, features=DEFAULT_FEATURES): + self.tconnect = tconnect + self.nightscout = nightscout + self.tconnect_device_id = tconnect_device_id + self.pretend = pretend + self.features = features + + def process(self, events, time_start, time_end): + logger.debug("ProcessAlarm: querying for last uploaded alarm") + last_upload = self.nightscout.last_uploaded_entry(ALARM_EVENTTYPE, time_start=time_start, time_end=time_end) + last_upload_time = None + if last_upload: + last_upload_time = arrow.get(last_upload["created_at"]) + logger.info("Last Nightscout alarm upload: %s" % last_upload_time) + + ns_entries = [] + for event in sorted(events, key=lambda x: x.eventTimestamp): + if last_upload_time and arrow.get(event.eventTimestamp) < last_upload_time: + if self.pretend: + logger.info("Skipping Alarm event before last upload time: %s (time range: %s - %s)" % (event, time_start, time_end)) + continue + + ns_entries.append(self.alarm_to_nsentry(event)) + + + return ns_entries + + def write(self, ns_entries): + count = 0 + for entry in ns_entries: + if self.pretend: + logger.info("Would upload to Nightscout: %s" % entry) + else: + logger.info("Uploading to Nightscout: %s" % entry) + self.nightscout.upload_entry(entry) + count += 1 + + return count + + + def alarm_to_nsentry(self, event): + if type(event) == eventtypes.LidAlarmActivated: + return NightscoutEntry.alarm( + created_at = event.eventTimestamp, + reason = event.alarmid + ) + elif type(event) == eventtypes.LidMalfunctionActivated: + return NightscoutEntry.alarm( + created_at = event.eventTimestamp, + reason = "Malfunction" + ) diff --git a/tconnectsync/sync/tandemsource/process_basal.py b/tconnectsync/sync/tandemsource/process_basal.py new file mode 100644 index 0000000..58dbc94 --- /dev/null +++ b/tconnectsync/sync/tandemsource/process_basal.py @@ -0,0 +1,79 @@ +import logging +import arrow + +from ...features import DEFAULT_FEATURES +from ...eventparser.generic import Events, decode_raw_events, EVENT_LEN +from ...eventparser.utils import bitmask_to_list +from ...eventparser import events as eventtypes +from ...domain.tandemsource.event_class import EventClass +from ...parser.nightscout import ( + BASAL_EVENTTYPE, + NightscoutEntry +) + +logger = logging.getLogger(__name__) + +class ProcessBasal: + def __init__(self, tconnect, nightscout, tconnect_device_id, pretend, features=DEFAULT_FEATURES): + self.tconnect = tconnect + self.nightscout = nightscout + self.tconnect_device_id = tconnect_device_id + self.pretend = pretend + self.features = features + + def process(self, events, time_start, time_end): + logger.debug("ProcessBasal: querying for last uploaded entry") + last_upload = self.nightscout.last_uploaded_entry(BASAL_EVENTTYPE, time_start=time_start, time_end=time_end) + last_upload_time = None + if last_upload: + last_upload_time = arrow.get(last_upload["created_at"]) + logger.info("Last Nightscout basal upload: %s" % last_upload_time) + + with_duration = [] + for event in sorted(events, key=lambda x: x.eventTimestamp): + if last_upload_time and arrow.get(event.eventTimestamp) < last_upload_time: + if self.pretend: + logger.info("Skipping basal event before last upload time: %s (time range: %s - %s)" % (event, time_start, time_end)) + continue + + with_duration.append([event.eventTimestamp, None, event]) + + for i in range(len(with_duration)-1): + with_duration[i][1] = with_duration[i+1][0] - with_duration[i][0] + + with_duration[-1][1] = time_end - with_duration[-1][0] + + ns_entries = [] + for item in with_duration: + ns_entries.append(self.basal_to_nsentry(*item)) + + return ns_entries + + def write(self, ns_entries): + count = 0 + for entry in ns_entries: + if self.pretend: + logger.info("Would upload to Nightscout: %s" % entry) + else: + logger.info("Uploading to Nightscout: %s" % entry) + self.nightscout.upload_entry(entry) + count += 1 + + return count + + + def basal_to_nsentry(self, start, duration, event): + if type(event) == eventtypes.LidBasalRateChange: + return NightscoutEntry.basal( + value = event.commandedbasalrate, + duration_mins = duration.seconds / 60, + created_at = start, + reason = ', '.join(bitmask_to_list(event.changetype)) + ) + if type(event) == eventtypes.LidBasalDelivery: + return NightscoutEntry.basal( + value = event.commandedRate, + duration_mins = duration.seconds / 60, + created_at = start, + reason = ', '.join(bitmask_to_list(event.commandedRateSource)) + ) diff --git a/tconnectsync/sync/tandemsource/process_basal_resume.py b/tconnectsync/sync/tandemsource/process_basal_resume.py new file mode 100644 index 0000000..5141348 --- /dev/null +++ b/tconnectsync/sync/tandemsource/process_basal_resume.py @@ -0,0 +1,61 @@ +import logging +import arrow + +from ...features import DEFAULT_FEATURES +from ...eventparser.generic import Events, decode_raw_events, EVENT_LEN +from ...eventparser.utils import bitmask_to_list +from ...eventparser import events as eventtypes +from ...domain.tandemsource.event_class import EventClass +from ...parser.nightscout import ( + BASALRESUME_EVENTTYPE, + NightscoutEntry +) + +logger = logging.getLogger(__name__) + +class ProcessBasalResume: + def __init__(self, tconnect, nightscout, tconnect_device_id, pretend, features=DEFAULT_FEATURES): + self.tconnect = tconnect + self.nightscout = nightscout + self.tconnect_device_id = tconnect_device_id + self.pretend = pretend + self.features = features + + def process(self, events, time_start, time_end): + logger.debug("ProcessBasalResume: querying for last uploaded resume-suspension") + last_upload = self.nightscout.last_uploaded_entry(BASALRESUME_EVENTTYPE, time_start=time_start, time_end=time_end) + last_upload_time = None + if last_upload: + last_upload_time = arrow.get(last_upload["created_at"]) + logger.info("Last Nightscout BasalResume upload: %s" % last_upload_time) + + ns_entries = [] + for event in sorted(events, key=lambda x: x.eventTimestamp): + if last_upload_time and arrow.get(event.eventTimestamp) < last_upload_time: + if self.pretend: + logger.info("Skipping BasalResume event before last upload time: %s (time range: %s - %s)" % (event, time_start, time_end)) + continue + + ns_entries.append(self.resume_to_nsentry(event)) + + + return ns_entries + + def write(self, ns_entries): + count = 0 + for entry in ns_entries: + if self.pretend: + logger.info("Would upload to Nightscout: %s" % entry) + else: + logger.info("Uploading to Nightscout: %s" % entry) + self.nightscout.upload_entry(entry) + count += 1 + + return count + + + def resume_to_nsentry(self, event): + if type(event) == eventtypes.LidPumpingResumed: + return NightscoutEntry.basalresume( + created_at = event.eventTimestamp + ) diff --git a/tconnectsync/sync/tandemsource/process_basal_suspension.py b/tconnectsync/sync/tandemsource/process_basal_suspension.py new file mode 100644 index 0000000..2ec2e2d --- /dev/null +++ b/tconnectsync/sync/tandemsource/process_basal_suspension.py @@ -0,0 +1,62 @@ +import logging +import arrow + +from ...features import DEFAULT_FEATURES +from ...eventparser.generic import Events, decode_raw_events, EVENT_LEN +from ...eventparser.utils import bitmask_to_list +from ...eventparser import events as eventtypes +from ...domain.tandemsource.event_class import EventClass +from ...parser.nightscout import ( + BASALSUSPENSION_EVENTTYPE, + NightscoutEntry +) + +logger = logging.getLogger(__name__) + +class ProcessBasalSuspension: + def __init__(self, tconnect, nightscout, tconnect_device_id, pretend, features=DEFAULT_FEATURES): + self.tconnect = tconnect + self.nightscout = nightscout + self.tconnect_device_id = tconnect_device_id + self.pretend = pretend + self.features = features + + def process(self, events, time_start, time_end): + logger.debug("ProcessBasalSuspension: querying for last uploaded suspension") + last_upload = self.nightscout.last_uploaded_entry(BASALSUSPENSION_EVENTTYPE, time_start=time_start, time_end=time_end) + last_upload_time = None + if last_upload: + last_upload_time = arrow.get(last_upload["created_at"]) + logger.info("Last Nightscout basalsuspension upload: %s" % last_upload_time) + + ns_entries = [] + for event in sorted(events, key=lambda x: x.eventTimestamp): + if last_upload_time and arrow.get(event.eventTimestamp) < last_upload_time: + if self.pretend: + logger.info("Skipping basalsuspension event before last upload time: %s (time range: %s - %s)" % (event, time_start, time_end)) + continue + + ns_entries.append(self.suspension_to_nsentry(event)) + + + return ns_entries + + def write(self, ns_entries): + count = 0 + for entry in ns_entries: + if self.pretend: + logger.info("Would upload to Nightscout: %s" % entry) + else: + logger.info("Uploading to Nightscout: %s" % entry) + self.nightscout.upload_entry(entry) + count += 1 + + return count + + + def suspension_to_nsentry(self, event): + if type(event) == eventtypes.LidPumpingSuspended: + return NightscoutEntry.basalsuspension( + created_at = event.eventTimestamp, + reason = ', '.join(bitmask_to_list(event.suspendreason)) + )