openrvdas.logger.transforms.true_winds_transform

Compute true winds by processing and aggregating vessel course/speed/heading and relative wind dir/speed records.

There are plenty of challenges with computing a universally-accepted true wind value. Even with the correct algorithm (not a given), unless the vessel nav and anemometer values have identical timestamps, there's the question of how one integrates/interpolates/extrapolates values with different timestamps.

The update_on_fields dict allows specifying which fields will trigger a new computation. This allows making simplifying assumptions, e.g. that vessel course/speed/heading is less variable than wind dir/speed, so we will only produce updated results when we receive new anemometer records. The max_field_age dict allows specifying a time beyond which we are unwilling to trust any computations.

A more robust approach would be to wait until we got the next vessel record and interpolate the course/speed/heading values between the two vessel records (or, conversely, output when we got a vessel record, using an interpolation of the preceding and following anemometer readings?).

  1#!/usr/bin/env python3
  2"""Compute true winds by processing and aggregating vessel
  3course/speed/heading and relative wind dir/speed records.
  4
  5There are plenty of challenges with computing a universally-accepted
  6true wind value. Even with the correct algorithm (not a given), unless
  7the vessel nav and anemometer values have identical timestamps,
  8there's the question of how one integrates/interpolates/extrapolates
  9values with different timestamps.
 10
 11The update_on_fields dict allows specifying which fields will trigger
 12a new computation. This allows making simplifying assumptions, e.g. that
 13vessel course/speed/heading is less variable than wind dir/speed, so we
 14will only produce updated results when we receive new anemometer records.
 15The max_field_age dict allows specifying a time beyond which we are unwilling
 16to trust any computations.
 17
 18A more robust approach would be to wait until we got the next vessel
 19record and interpolate the course/speed/heading values between the two
 20vessel records (or, conversely, output when we got a vessel record,
 21using an interpolation of the preceding and following anemometer
 22readings?).
 23
 24"""
 25
 26import logging
 27import time
 28
 29from pprint import pformat
 30from typing import Union
 31
 32from logger.utils.das_record import DASRecord, to_das_record_list  # noqa: E402
 33from logger.utils.truewinds.truew import truew  # noqa: E402
 34from logger.transforms.derived_data_transform import DerivedDataTransform  # noqa: E402
 35
 36
 37################################################################################
 38class TrueWindsTransform(DerivedDataTransform):
 39    """Transform that computes true winds from vessel
 40    course/speed/heading and anemometer relative wind speed/dir.
 41    """
 42
 43    def __init__(self,
 44                 course_field, speed_field, heading_field,
 45                 wind_dir_field, wind_speed_field,
 46                 true_dir_name,
 47                 true_speed_name,
 48                 apparent_dir_name,
 49                 update_on_fields=[],
 50                 max_field_age={},
 51                 zero_line_reference=0,
 52                 convert_wind_factor=1,
 53                 convert_speed_factor=1,
 54                 data_id=None,
 55                 metadata_interval=None,
 56                 **kwargs):
 57        """
 58        ```
 59        course_field
 60        speed_field
 61        heading_field
 62        wind_dir_field
 63        wind_speed_field
 64                 Field names from which we should take values for
 65                 course, speed over ground, heading, relative wind speed
 66                 and relative wind direction.
 67
 68        true_dir_name
 69        true_speed_name
 70        apparent_dir_name
 71                 Names that should be given to transform output values.
 72
 73        update_on_fields
 74                 If non-empty, a list of fields, any of whose arrival should
 75                 trigger an output record. If None, generate output when any
 76                 field is updated.
 77
 78        max_field_age
 79                 If non-empty, a dict of field_name:seconds, specifying that
 80                 no output is to be produces if the age of any of the specified
 81                 names is older than the specified number of seconds.
 82
 83        zero_line_reference
 84                 Angle between bow and zero line on anemometer, referenced
 85                 to ship.
 86
 87        convert_wind_factor
 88        convert_speed_factor
 89                 Wind speed and vessel speed may be in different units; e.g
 90                 wind speed in meters/sec and vessel speed in knots. Multiply
 91                 the respective values we get by the respective factors.
 92                 Typically, only one of these will be not equal to 1; e.g. to
 93                 output true winds as meters/sec, we'll leave convert_wind_factor
 94                 as 1 and specify convert_speed_factor=0.5144
 95
 96        metadata_interval - how many seconds between when we attach field metadata
 97                     to a record we send out.
 98
 99        data_id  Optional name that will be attached to the resulting DASRecord
100        ```
101        """
102        super().__init__(**kwargs)  # processes 'quiet' and type hints
103
104        self.course_field = course_field
105        self.speed_field = speed_field
106        self.heading_field = heading_field
107        self.wind_dir_field = wind_dir_field
108        self.wind_speed_field = wind_speed_field
109
110        self.true_dir_name = true_dir_name
111        self.true_speed_name = true_speed_name
112        self.apparent_dir_name = apparent_dir_name
113
114        self.update_on_fields = update_on_fields
115        self.max_field_age = max_field_age
116        self.field_age = {}
117
118        self.zero_line_reference = zero_line_reference
119
120        self.convert_wind_factor = convert_wind_factor
121        self.convert_speed_factor = convert_speed_factor
122
123        self.metadata_interval = metadata_interval
124        self.last_metadata_send = 0
125
126        self.data_id = data_id
127
128        # TODO: It may make sense for us to cache most recent values so
129        # that, for example, we can take single DASRecords in the
130        # transform() method and use the most recent values we've seen
131        # from previous calls.
132        self.course_val = None
133        self.speed_val = None
134        self.heading_val = None
135        self.wind_dir_val = None
136        self.wind_speed_val = None
137
138        self.course_val_time = 0
139        self.speed_val_time = 0
140        self.heading_val_time = 0
141        self.wind_dir_val_time = 0
142        self.wind_speed_val_time = 0
143
144    ############################
145    def fields(self):
146        """Which fields are we interested in to produce transformed data?"""
147        return [self.course_field, self.speed_field, self.heading_field,
148                self.wind_dir_field, self.wind_speed_field]
149
150    ############################
151    def _metadata(self):
152        """Return a dict of metadata for our derived fields."""
153
154        metadata_fields = {
155            self.true_dir_name: {
156                'description': 'Derived true wind direction from %s, %s, %s, %s, %s'
157                % (self.course_field, self.speed_field, self.heading_field,
158                   self.wind_dir_field, self.wind_speed_field),
159                'units': 'degrees',
160                'device': 'TrueWindTransform',
161                'device_type': 'DerivedTrueWindTransform',
162                'device_type_field': self.true_dir_name
163            },
164            self.true_speed_name: {
165                'description': 'Derived true wind speed from %s, %s, %s, %s, %s'
166                % (self.course_field, self.speed_field, self.heading_field,
167                   self.wind_dir_field, self.wind_speed_field),
168                'units': 'depends on conversion used for %s (%g) and %s (%g)'
169                % (self.speed_field, self.convert_speed_factor,
170                   self.wind_speed_field, self.convert_wind_factor),
171                'device': 'TrueWindTransform',
172                'device_type': 'DerivedTrueWindTransform',
173                'device_type_field': self.true_speed_name
174            },
175            self.apparent_dir_name: {
176                'description': 'Derived apparent wind speed from %s, %s, %s, %s, %s'
177                % (self.course_field, self.speed_field, self.heading_field,
178                   self.wind_dir_field, self.wind_speed_field),
179                'units': 'degrees',
180                'device': 'TrueWindTransform',
181                'device_type': 'DerivedTrueWindTransform',
182                'device_type_field': self.apparent_dir_name
183            }
184        }
185        return metadata_fields
186
187    ############################
188    def transform(self, record: Union[dict, DASRecord]):
189        """Incorporate any useable fields in this record, and if it gives
190        us a new true wind value, return the results."""
191
192        # See if it's something we can process, and if not, try digesting
193        if not self.can_process_record(record):  # inherited from BaseModule()
194            return self.digest_record(record)  # inherited from BaseModule()
195
196        results = []
197        for das_record in to_das_record_list(record):
198            # If they haven't specified specific fields we should wait for
199            # before updates, plan to emit an update after every new record
200            # we process. Otherwise, assume we're not going to update unless
201            # we see one of the named fields.
202            if not self.update_on_fields:
203                update = True
204            else:
205                update = False
206
207            timestamp = das_record.timestamp
208            if not timestamp:
209                logging.info('DASRecord is missing timestamp - skipping')
210                continue
211
212            # Get latest values for any of our fields
213            fields = das_record.fields
214            if self.course_field in fields:
215                if timestamp >= self.course_val_time:
216                    self.course_val = fields.get(self.course_field)
217                    self.course_val_time = timestamp
218                    if self.course_field in self.update_on_fields:
219                        update = True
220
221            if self.speed_field in fields:
222                if timestamp >= self.speed_val_time:
223                    self.speed_val = fields.get(self.speed_field)
224                    self.speed_val *= self.convert_speed_factor
225                    self.speed_val_time = timestamp
226                    if self.speed_field in self.update_on_fields:
227                        update = True
228
229            if self.heading_field in fields:
230                if timestamp >= self.heading_val_time:
231                    self.heading_val = fields.get(self.heading_field)
232                    self.heading_val_time = timestamp
233                    if self.heading_field in self.update_on_fields:
234                        update = True
235
236            if self.wind_dir_field in fields:
237                if timestamp >= self.wind_dir_val_time:
238                    self.wind_dir_val = fields.get(self.wind_dir_field)
239                    self.wind_dir_val_time = timestamp
240                    if self.wind_dir_field in self.update_on_fields:
241                        update = True
242
243            if self.wind_speed_field in fields:
244                if timestamp >= self.wind_speed_val_time:
245                    self.wind_speed_val = fields.get(self.wind_speed_field)
246                    self.wind_speed_val *= self.convert_wind_factor
247                    self.wind_speed_val_time = timestamp
248                    if self.wind_speed_field in self.update_on_fields:
249                        update = True
250
251            # Check if needed all values are present, and none are too old to use
252            if self._values_too_old(timestamp):
253                continue
254
255            # If we've not seen anything that updates fields that would
256            # trigger a new true winds value, skip rest of computation.
257            if not update:
258                logging.debug('No update needed')
259                continue
260
261            logging.debug('Computing new true winds')
262            (true_dir, true_speed, apparent_dir) = truew(crse=self.course_val,
263                                                         cspd=self.speed_val,
264                                                         hd=self.heading_val,
265                                                         wdir=self.wind_dir_val,
266                                                         zlr=self.zero_line_reference,
267                                                         wspd=self.wind_speed_val)
268
269            logging.debug('Got true winds: dir: %s, speed: %s, apparent_dir: %s',
270                          true_dir, true_speed, apparent_dir)
271            if None in (true_dir, true_speed, apparent_dir):
272                logging.info('Got invalid true winds')
273                continue
274
275            # If here, we've got a valid new true wind result
276            true_wind_fields = {self.true_dir_name: true_dir,
277                                self.true_speed_name: true_speed,
278                                self.apparent_dir_name: apparent_dir}
279
280            # Add in metadata if so specified and it's been long enough since
281            # we last sent it.
282            now = time.time()
283            if self.metadata_interval and \
284               now - self.metadata_interval > self.last_metadata_send:
285                metadata = {'fields': self._metadata()}
286                self.last_metadata_send = now
287                logging.debug('Emitting metadata: %s', pformat(metadata))
288            else:
289                metadata = None
290
291            results.append(DASRecord(timestamp=timestamp, fields=true_wind_fields,
292                                     metadata=metadata, data_id=self.data_id))
293
294        logging.debug('Sending %d true wind results.', len(results))
295        return results
296
297    ############################
298    def _values_too_old(self, timestamp):
299        """Return true if any values are missing or too old to use."""
300
301        if None in (self.course_val, self.speed_val, self.heading_val,
302                    self.wind_dir_val, self.wind_speed_val):
303            logging.debug('Not all required values for true winds are present: '
304                          'time: %s: %s: %s, %s: %s, %s: %s, %s: %s, %s: %s',
305                          timestamp,
306                          self.course_field, self.course_val,
307                          self.speed_field, self.speed_val,
308                          self.heading_field, self.heading_val,
309                          self.wind_dir_field, self.wind_dir_val,
310                          self.wind_speed_field, self.wind_speed_val)
311            return True
312
313        course_max_age = self.max_field_age.get(self.course_field)
314        if (course_max_age and timestamp - self.course_val_time > course_max_age):
315            logging.debug('course_field too old - max age %g, age %g',
316                          course_max_age, timestamp - self.course_val_time)
317            return True
318
319        speed_max_age = self.max_field_age.get(self.speed_field)
320        if speed_max_age:
321            if timestamp - self.speed_val_time > speed_max_age:
322                logging.debug('speed_field too old - max age %g, age %g',
323                              speed_max_age, timestamp - self.speed_val_time)
324                return True
325
326        heading_max_age = self.max_field_age.get(self.heading_field)
327        if heading_max_age:
328            if timestamp - self.heading_val_time > heading_max_age:
329                logging.debug('heading_field too old - max age %g, age %g',
330                              heading_max_age, timestamp - self.heading_val_time)
331                return True
332
333        wind_dir_max_age = self.max_field_age.get(self.wind_dir_field)
334        if wind_dir_max_age:
335            if timestamp - self.wind_dir_val_time > wind_dir_max_age:
336                logging.debug('wind_dir_field too old - max age %g, age %g',
337                              wind_dir_max_age, timestamp - self.wind_dir_val_time)
338                return True
339
340        wind_speed_max_age = self.max_field_age.get(self.wind_speed_field)
341        if wind_speed_max_age:
342            if timestamp - self.wind_speed_val_time > wind_speed_max_age:
343                logging.debug('wind_speed_field too old - max age %g, age %g',
344                              wind_speed_max_age, timestamp - self.wind_speed_val_time)
345                return True
346
347        # Everything is present, and nothing's too old...
348        return False
class TrueWindsTransform(logger.transforms.derived_data_transform.DerivedDataTransform):
 39class TrueWindsTransform(DerivedDataTransform):
 40    """Transform that computes true winds from vessel
 41    course/speed/heading and anemometer relative wind speed/dir.
 42    """
 43
 44    def __init__(self,
 45                 course_field, speed_field, heading_field,
 46                 wind_dir_field, wind_speed_field,
 47                 true_dir_name,
 48                 true_speed_name,
 49                 apparent_dir_name,
 50                 update_on_fields=[],
 51                 max_field_age={},
 52                 zero_line_reference=0,
 53                 convert_wind_factor=1,
 54                 convert_speed_factor=1,
 55                 data_id=None,
 56                 metadata_interval=None,
 57                 **kwargs):
 58        """
 59        ```
 60        course_field
 61        speed_field
 62        heading_field
 63        wind_dir_field
 64        wind_speed_field
 65                 Field names from which we should take values for
 66                 course, speed over ground, heading, relative wind speed
 67                 and relative wind direction.
 68
 69        true_dir_name
 70        true_speed_name
 71        apparent_dir_name
 72                 Names that should be given to transform output values.
 73
 74        update_on_fields
 75                 If non-empty, a list of fields, any of whose arrival should
 76                 trigger an output record. If None, generate output when any
 77                 field is updated.
 78
 79        max_field_age
 80                 If non-empty, a dict of field_name:seconds, specifying that
 81                 no output is to be produces if the age of any of the specified
 82                 names is older than the specified number of seconds.
 83
 84        zero_line_reference
 85                 Angle between bow and zero line on anemometer, referenced
 86                 to ship.
 87
 88        convert_wind_factor
 89        convert_speed_factor
 90                 Wind speed and vessel speed may be in different units; e.g
 91                 wind speed in meters/sec and vessel speed in knots. Multiply
 92                 the respective values we get by the respective factors.
 93                 Typically, only one of these will be not equal to 1; e.g. to
 94                 output true winds as meters/sec, we'll leave convert_wind_factor
 95                 as 1 and specify convert_speed_factor=0.5144
 96
 97        metadata_interval - how many seconds between when we attach field metadata
 98                     to a record we send out.
 99
100        data_id  Optional name that will be attached to the resulting DASRecord
101        ```
102        """
103        super().__init__(**kwargs)  # processes 'quiet' and type hints
104
105        self.course_field = course_field
106        self.speed_field = speed_field
107        self.heading_field = heading_field
108        self.wind_dir_field = wind_dir_field
109        self.wind_speed_field = wind_speed_field
110
111        self.true_dir_name = true_dir_name
112        self.true_speed_name = true_speed_name
113        self.apparent_dir_name = apparent_dir_name
114
115        self.update_on_fields = update_on_fields
116        self.max_field_age = max_field_age
117        self.field_age = {}
118
119        self.zero_line_reference = zero_line_reference
120
121        self.convert_wind_factor = convert_wind_factor
122        self.convert_speed_factor = convert_speed_factor
123
124        self.metadata_interval = metadata_interval
125        self.last_metadata_send = 0
126
127        self.data_id = data_id
128
129        # TODO: It may make sense for us to cache most recent values so
130        # that, for example, we can take single DASRecords in the
131        # transform() method and use the most recent values we've seen
132        # from previous calls.
133        self.course_val = None
134        self.speed_val = None
135        self.heading_val = None
136        self.wind_dir_val = None
137        self.wind_speed_val = None
138
139        self.course_val_time = 0
140        self.speed_val_time = 0
141        self.heading_val_time = 0
142        self.wind_dir_val_time = 0
143        self.wind_speed_val_time = 0
144
145    ############################
146    def fields(self):
147        """Which fields are we interested in to produce transformed data?"""
148        return [self.course_field, self.speed_field, self.heading_field,
149                self.wind_dir_field, self.wind_speed_field]
150
151    ############################
152    def _metadata(self):
153        """Return a dict of metadata for our derived fields."""
154
155        metadata_fields = {
156            self.true_dir_name: {
157                'description': 'Derived true wind direction from %s, %s, %s, %s, %s'
158                % (self.course_field, self.speed_field, self.heading_field,
159                   self.wind_dir_field, self.wind_speed_field),
160                'units': 'degrees',
161                'device': 'TrueWindTransform',
162                'device_type': 'DerivedTrueWindTransform',
163                'device_type_field': self.true_dir_name
164            },
165            self.true_speed_name: {
166                'description': 'Derived true wind speed from %s, %s, %s, %s, %s'
167                % (self.course_field, self.speed_field, self.heading_field,
168                   self.wind_dir_field, self.wind_speed_field),
169                'units': 'depends on conversion used for %s (%g) and %s (%g)'
170                % (self.speed_field, self.convert_speed_factor,
171                   self.wind_speed_field, self.convert_wind_factor),
172                'device': 'TrueWindTransform',
173                'device_type': 'DerivedTrueWindTransform',
174                'device_type_field': self.true_speed_name
175            },
176            self.apparent_dir_name: {
177                'description': 'Derived apparent wind speed from %s, %s, %s, %s, %s'
178                % (self.course_field, self.speed_field, self.heading_field,
179                   self.wind_dir_field, self.wind_speed_field),
180                'units': 'degrees',
181                'device': 'TrueWindTransform',
182                'device_type': 'DerivedTrueWindTransform',
183                'device_type_field': self.apparent_dir_name
184            }
185        }
186        return metadata_fields
187
188    ############################
189    def transform(self, record: Union[dict, DASRecord]):
190        """Incorporate any useable fields in this record, and if it gives
191        us a new true wind value, return the results."""
192
193        # See if it's something we can process, and if not, try digesting
194        if not self.can_process_record(record):  # inherited from BaseModule()
195            return self.digest_record(record)  # inherited from BaseModule()
196
197        results = []
198        for das_record in to_das_record_list(record):
199            # If they haven't specified specific fields we should wait for
200            # before updates, plan to emit an update after every new record
201            # we process. Otherwise, assume we're not going to update unless
202            # we see one of the named fields.
203            if not self.update_on_fields:
204                update = True
205            else:
206                update = False
207
208            timestamp = das_record.timestamp
209            if not timestamp:
210                logging.info('DASRecord is missing timestamp - skipping')
211                continue
212
213            # Get latest values for any of our fields
214            fields = das_record.fields
215            if self.course_field in fields:
216                if timestamp >= self.course_val_time:
217                    self.course_val = fields.get(self.course_field)
218                    self.course_val_time = timestamp
219                    if self.course_field in self.update_on_fields:
220                        update = True
221
222            if self.speed_field in fields:
223                if timestamp >= self.speed_val_time:
224                    self.speed_val = fields.get(self.speed_field)
225                    self.speed_val *= self.convert_speed_factor
226                    self.speed_val_time = timestamp
227                    if self.speed_field in self.update_on_fields:
228                        update = True
229
230            if self.heading_field in fields:
231                if timestamp >= self.heading_val_time:
232                    self.heading_val = fields.get(self.heading_field)
233                    self.heading_val_time = timestamp
234                    if self.heading_field in self.update_on_fields:
235                        update = True
236
237            if self.wind_dir_field in fields:
238                if timestamp >= self.wind_dir_val_time:
239                    self.wind_dir_val = fields.get(self.wind_dir_field)
240                    self.wind_dir_val_time = timestamp
241                    if self.wind_dir_field in self.update_on_fields:
242                        update = True
243
244            if self.wind_speed_field in fields:
245                if timestamp >= self.wind_speed_val_time:
246                    self.wind_speed_val = fields.get(self.wind_speed_field)
247                    self.wind_speed_val *= self.convert_wind_factor
248                    self.wind_speed_val_time = timestamp
249                    if self.wind_speed_field in self.update_on_fields:
250                        update = True
251
252            # Check if needed all values are present, and none are too old to use
253            if self._values_too_old(timestamp):
254                continue
255
256            # If we've not seen anything that updates fields that would
257            # trigger a new true winds value, skip rest of computation.
258            if not update:
259                logging.debug('No update needed')
260                continue
261
262            logging.debug('Computing new true winds')
263            (true_dir, true_speed, apparent_dir) = truew(crse=self.course_val,
264                                                         cspd=self.speed_val,
265                                                         hd=self.heading_val,
266                                                         wdir=self.wind_dir_val,
267                                                         zlr=self.zero_line_reference,
268                                                         wspd=self.wind_speed_val)
269
270            logging.debug('Got true winds: dir: %s, speed: %s, apparent_dir: %s',
271                          true_dir, true_speed, apparent_dir)
272            if None in (true_dir, true_speed, apparent_dir):
273                logging.info('Got invalid true winds')
274                continue
275
276            # If here, we've got a valid new true wind result
277            true_wind_fields = {self.true_dir_name: true_dir,
278                                self.true_speed_name: true_speed,
279                                self.apparent_dir_name: apparent_dir}
280
281            # Add in metadata if so specified and it's been long enough since
282            # we last sent it.
283            now = time.time()
284            if self.metadata_interval and \
285               now - self.metadata_interval > self.last_metadata_send:
286                metadata = {'fields': self._metadata()}
287                self.last_metadata_send = now
288                logging.debug('Emitting metadata: %s', pformat(metadata))
289            else:
290                metadata = None
291
292            results.append(DASRecord(timestamp=timestamp, fields=true_wind_fields,
293                                     metadata=metadata, data_id=self.data_id))
294
295        logging.debug('Sending %d true wind results.', len(results))
296        return results
297
298    ############################
299    def _values_too_old(self, timestamp):
300        """Return true if any values are missing or too old to use."""
301
302        if None in (self.course_val, self.speed_val, self.heading_val,
303                    self.wind_dir_val, self.wind_speed_val):
304            logging.debug('Not all required values for true winds are present: '
305                          'time: %s: %s: %s, %s: %s, %s: %s, %s: %s, %s: %s',
306                          timestamp,
307                          self.course_field, self.course_val,
308                          self.speed_field, self.speed_val,
309                          self.heading_field, self.heading_val,
310                          self.wind_dir_field, self.wind_dir_val,
311                          self.wind_speed_field, self.wind_speed_val)
312            return True
313
314        course_max_age = self.max_field_age.get(self.course_field)
315        if (course_max_age and timestamp - self.course_val_time > course_max_age):
316            logging.debug('course_field too old - max age %g, age %g',
317                          course_max_age, timestamp - self.course_val_time)
318            return True
319
320        speed_max_age = self.max_field_age.get(self.speed_field)
321        if speed_max_age:
322            if timestamp - self.speed_val_time > speed_max_age:
323                logging.debug('speed_field too old - max age %g, age %g',
324                              speed_max_age, timestamp - self.speed_val_time)
325                return True
326
327        heading_max_age = self.max_field_age.get(self.heading_field)
328        if heading_max_age:
329            if timestamp - self.heading_val_time > heading_max_age:
330                logging.debug('heading_field too old - max age %g, age %g',
331                              heading_max_age, timestamp - self.heading_val_time)
332                return True
333
334        wind_dir_max_age = self.max_field_age.get(self.wind_dir_field)
335        if wind_dir_max_age:
336            if timestamp - self.wind_dir_val_time > wind_dir_max_age:
337                logging.debug('wind_dir_field too old - max age %g, age %g',
338                              wind_dir_max_age, timestamp - self.wind_dir_val_time)
339                return True
340
341        wind_speed_max_age = self.max_field_age.get(self.wind_speed_field)
342        if wind_speed_max_age:
343            if timestamp - self.wind_speed_val_time > wind_speed_max_age:
344                logging.debug('wind_speed_field too old - max age %g, age %g',
345                              wind_speed_max_age, timestamp - self.wind_speed_val_time)
346                return True
347
348        # Everything is present, and nothing's too old...
349        return False

Transform that computes true winds from vessel course/speed/heading and anemometer relative wind speed/dir.

TrueWindsTransform( course_field, speed_field, heading_field, wind_dir_field, wind_speed_field, true_dir_name, true_speed_name, apparent_dir_name, update_on_fields=[], max_field_age={}, zero_line_reference=0, convert_wind_factor=1, convert_speed_factor=1, data_id=None, metadata_interval=None, **kwargs)
 44    def __init__(self,
 45                 course_field, speed_field, heading_field,
 46                 wind_dir_field, wind_speed_field,
 47                 true_dir_name,
 48                 true_speed_name,
 49                 apparent_dir_name,
 50                 update_on_fields=[],
 51                 max_field_age={},
 52                 zero_line_reference=0,
 53                 convert_wind_factor=1,
 54                 convert_speed_factor=1,
 55                 data_id=None,
 56                 metadata_interval=None,
 57                 **kwargs):
 58        """
 59        ```
 60        course_field
 61        speed_field
 62        heading_field
 63        wind_dir_field
 64        wind_speed_field
 65                 Field names from which we should take values for
 66                 course, speed over ground, heading, relative wind speed
 67                 and relative wind direction.
 68
 69        true_dir_name
 70        true_speed_name
 71        apparent_dir_name
 72                 Names that should be given to transform output values.
 73
 74        update_on_fields
 75                 If non-empty, a list of fields, any of whose arrival should
 76                 trigger an output record. If None, generate output when any
 77                 field is updated.
 78
 79        max_field_age
 80                 If non-empty, a dict of field_name:seconds, specifying that
 81                 no output is to be produces if the age of any of the specified
 82                 names is older than the specified number of seconds.
 83
 84        zero_line_reference
 85                 Angle between bow and zero line on anemometer, referenced
 86                 to ship.
 87
 88        convert_wind_factor
 89        convert_speed_factor
 90                 Wind speed and vessel speed may be in different units; e.g
 91                 wind speed in meters/sec and vessel speed in knots. Multiply
 92                 the respective values we get by the respective factors.
 93                 Typically, only one of these will be not equal to 1; e.g. to
 94                 output true winds as meters/sec, we'll leave convert_wind_factor
 95                 as 1 and specify convert_speed_factor=0.5144
 96
 97        metadata_interval - how many seconds between when we attach field metadata
 98                     to a record we send out.
 99
100        data_id  Optional name that will be attached to the resulting DASRecord
101        ```
102        """
103        super().__init__(**kwargs)  # processes 'quiet' and type hints
104
105        self.course_field = course_field
106        self.speed_field = speed_field
107        self.heading_field = heading_field
108        self.wind_dir_field = wind_dir_field
109        self.wind_speed_field = wind_speed_field
110
111        self.true_dir_name = true_dir_name
112        self.true_speed_name = true_speed_name
113        self.apparent_dir_name = apparent_dir_name
114
115        self.update_on_fields = update_on_fields
116        self.max_field_age = max_field_age
117        self.field_age = {}
118
119        self.zero_line_reference = zero_line_reference
120
121        self.convert_wind_factor = convert_wind_factor
122        self.convert_speed_factor = convert_speed_factor
123
124        self.metadata_interval = metadata_interval
125        self.last_metadata_send = 0
126
127        self.data_id = data_id
128
129        # TODO: It may make sense for us to cache most recent values so
130        # that, for example, we can take single DASRecords in the
131        # transform() method and use the most recent values we've seen
132        # from previous calls.
133        self.course_val = None
134        self.speed_val = None
135        self.heading_val = None
136        self.wind_dir_val = None
137        self.wind_speed_val = None
138
139        self.course_val_time = 0
140        self.speed_val_time = 0
141        self.heading_val_time = 0
142        self.wind_dir_val_time = 0
143        self.wind_speed_val_time = 0
course_field
speed_field
heading_field
wind_dir_field
wind_speed_field
         Field names from which we should take values for
         course, speed over ground, heading, relative wind speed
         and relative wind direction.

true_dir_name
true_speed_name
apparent_dir_name
         Names that should be given to transform output values.

update_on_fields
         If non-empty, a list of fields, any of whose arrival should
         trigger an output record. If None, generate output when any
         field is updated.

max_field_age
         If non-empty, a dict of field_name:seconds, specifying that
         no output is to be produces if the age of any of the specified
         names is older than the specified number of seconds.

zero_line_reference
         Angle between bow and zero line on anemometer, referenced
         to ship.

convert_wind_factor
convert_speed_factor
         Wind speed and vessel speed may be in different units; e.g
         wind speed in meters/sec and vessel speed in knots. Multiply
         the respective values we get by the respective factors.
         Typically, only one of these will be not equal to 1; e.g. to
         output true winds as meters/sec, we'll leave convert_wind_factor
         as 1 and specify convert_speed_factor=0.5144

metadata_interval - how many seconds between when we attach field metadata
             to a record we send out.

data_id  Optional name that will be attached to the resulting DASRecord
course_field
speed_field
heading_field
wind_dir_field
wind_speed_field
true_dir_name
true_speed_name
apparent_dir_name
update_on_fields
max_field_age
field_age
zero_line_reference
convert_wind_factor
convert_speed_factor
metadata_interval
last_metadata_send
data_id
course_val
speed_val
heading_val
wind_dir_val
wind_speed_val
course_val_time
speed_val_time
heading_val_time
wind_dir_val_time
wind_speed_val_time
def fields(self):
146    def fields(self):
147        """Which fields are we interested in to produce transformed data?"""
148        return [self.course_field, self.speed_field, self.heading_field,
149                self.wind_dir_field, self.wind_speed_field]

Which fields are we interested in to produce transformed data?

def transform(self, record: Union[dict, logger.utils.das_record.DASRecord]):
189    def transform(self, record: Union[dict, DASRecord]):
190        """Incorporate any useable fields in this record, and if it gives
191        us a new true wind value, return the results."""
192
193        # See if it's something we can process, and if not, try digesting
194        if not self.can_process_record(record):  # inherited from BaseModule()
195            return self.digest_record(record)  # inherited from BaseModule()
196
197        results = []
198        for das_record in to_das_record_list(record):
199            # If they haven't specified specific fields we should wait for
200            # before updates, plan to emit an update after every new record
201            # we process. Otherwise, assume we're not going to update unless
202            # we see one of the named fields.
203            if not self.update_on_fields:
204                update = True
205            else:
206                update = False
207
208            timestamp = das_record.timestamp
209            if not timestamp:
210                logging.info('DASRecord is missing timestamp - skipping')
211                continue
212
213            # Get latest values for any of our fields
214            fields = das_record.fields
215            if self.course_field in fields:
216                if timestamp >= self.course_val_time:
217                    self.course_val = fields.get(self.course_field)
218                    self.course_val_time = timestamp
219                    if self.course_field in self.update_on_fields:
220                        update = True
221
222            if self.speed_field in fields:
223                if timestamp >= self.speed_val_time:
224                    self.speed_val = fields.get(self.speed_field)
225                    self.speed_val *= self.convert_speed_factor
226                    self.speed_val_time = timestamp
227                    if self.speed_field in self.update_on_fields:
228                        update = True
229
230            if self.heading_field in fields:
231                if timestamp >= self.heading_val_time:
232                    self.heading_val = fields.get(self.heading_field)
233                    self.heading_val_time = timestamp
234                    if self.heading_field in self.update_on_fields:
235                        update = True
236
237            if self.wind_dir_field in fields:
238                if timestamp >= self.wind_dir_val_time:
239                    self.wind_dir_val = fields.get(self.wind_dir_field)
240                    self.wind_dir_val_time = timestamp
241                    if self.wind_dir_field in self.update_on_fields:
242                        update = True
243
244            if self.wind_speed_field in fields:
245                if timestamp >= self.wind_speed_val_time:
246                    self.wind_speed_val = fields.get(self.wind_speed_field)
247                    self.wind_speed_val *= self.convert_wind_factor
248                    self.wind_speed_val_time = timestamp
249                    if self.wind_speed_field in self.update_on_fields:
250                        update = True
251
252            # Check if needed all values are present, and none are too old to use
253            if self._values_too_old(timestamp):
254                continue
255
256            # If we've not seen anything that updates fields that would
257            # trigger a new true winds value, skip rest of computation.
258            if not update:
259                logging.debug('No update needed')
260                continue
261
262            logging.debug('Computing new true winds')
263            (true_dir, true_speed, apparent_dir) = truew(crse=self.course_val,
264                                                         cspd=self.speed_val,
265                                                         hd=self.heading_val,
266                                                         wdir=self.wind_dir_val,
267                                                         zlr=self.zero_line_reference,
268                                                         wspd=self.wind_speed_val)
269
270            logging.debug('Got true winds: dir: %s, speed: %s, apparent_dir: %s',
271                          true_dir, true_speed, apparent_dir)
272            if None in (true_dir, true_speed, apparent_dir):
273                logging.info('Got invalid true winds')
274                continue
275
276            # If here, we've got a valid new true wind result
277            true_wind_fields = {self.true_dir_name: true_dir,
278                                self.true_speed_name: true_speed,
279                                self.apparent_dir_name: apparent_dir}
280
281            # Add in metadata if so specified and it's been long enough since
282            # we last sent it.
283            now = time.time()
284            if self.metadata_interval and \
285               now - self.metadata_interval > self.last_metadata_send:
286                metadata = {'fields': self._metadata()}
287                self.last_metadata_send = now
288                logging.debug('Emitting metadata: %s', pformat(metadata))
289            else:
290                metadata = None
291
292            results.append(DASRecord(timestamp=timestamp, fields=true_wind_fields,
293                                     metadata=metadata, data_id=self.data_id))
294
295        logging.debug('Sending %d true wind results.', len(results))
296        return results

Incorporate any useable fields in this record, and if it gives us a new true wind value, return the results.