openrvdas.logger.transforms.xml_aggregator_transform

No module-level documentation available.
 1#!/usr/bin/env python3
 2
 3import logging
 4
 5from threading import Lock
 6from xml.sax.handler import ContentHandler
 7from xml.sax import make_parser
 8
 9from logger.transforms.transform import Transform  # noqa: E402
10
11
12################################################################################
13class XMLAggregatorTransform(Transform):
14    """Aggregate passed lines of XML until a complete XML record whose
15    outermost element matches 'tag' has been seen, then pass it on as a
16    single record."""
17    ############################
18
19    def __init__(self, tag, **kwargs):
20        """
21        'tag' should be the identity of the top-level XML element that
22        we're expecting to read, e.g. 'OSU_DAS_Record'.
23        """
24        super().__init__(**kwargs)  # processes 'quiet' and type hints
25        self.tag = tag
26
27        # Only let one thread touch buffer at a time. Of course, if we're
28        # getting interleaved lines from different XML records here, we're
29        # screwed anyway.
30        self.buffer_lock = Lock()
31        self.buffer = ''
32
33        self.handler = _XMLHandler(tag=tag)
34        self.parser = make_parser(['xml.sax.IncrementalParser'])
35        self.parser.setContentHandler(self.handler)
36
37    ############################
38    def transform(self, record: str):
39        """Aggregate, returning None until we're done, then return record."""
40
41        # See if it's something we can process, and if not, try digesting
42        if not self.can_process_record(record):  # inherited from BaseModule()
43            return self.digest_record(record)  # inherited from BaseModule()
44
45        with self.buffer_lock:
46            # Feed record to the incremental parser
47            self.buffer += record + '\n'
48            self.parser.feed(record)
49            logging.debug('transform() got line: %s', record)
50
51            # If the record completes and XML record, it will be added to the
52            # queue in self.handler.items() - pop it off and return
53            if self.handler.complete():
54                xml_record = self.buffer
55                self.buffer = ''
56                self.parser.close()
57                self.parser.reset()
58                self.handler.reset()
59                logging.debug('transform() got closing tag: %s', xml_record)
60                return xml_record
61
62        # Otherwise go home emptyhanded
63        return None
64
65
66################################################################################
67class _XMLHandler(ContentHandler):
68    """
69    Helper class here - _XMLHandler.endElement() will get called on
70    closing tags, so we can detect when we've got a closing tag for
71    whatever XML element we're after. We're omitting startElement,
72    and characters methods to store data on a stack during processing.
73    """
74    ############################
75
76    def __init__(self, tag):
77        super().__init__()
78        self.tag = tag
79        self.item_list = []
80        self.element_complete = False
81
82    ############################
83    def complete(self):
84        return self.element_complete
85
86    ############################
87    def reset(self):
88        self.element_complete = False
89
90    ############################
91    def endElement(self, name):
92        if name == self.tag:
93            # Create item from stored data on stack
94            self.element_complete = True
class XMLAggregatorTransform(logger.transforms.transform.Transform):
14class XMLAggregatorTransform(Transform):
15    """Aggregate passed lines of XML until a complete XML record whose
16    outermost element matches 'tag' has been seen, then pass it on as a
17    single record."""
18    ############################
19
20    def __init__(self, tag, **kwargs):
21        """
22        'tag' should be the identity of the top-level XML element that
23        we're expecting to read, e.g. 'OSU_DAS_Record'.
24        """
25        super().__init__(**kwargs)  # processes 'quiet' and type hints
26        self.tag = tag
27
28        # Only let one thread touch buffer at a time. Of course, if we're
29        # getting interleaved lines from different XML records here, we're
30        # screwed anyway.
31        self.buffer_lock = Lock()
32        self.buffer = ''
33
34        self.handler = _XMLHandler(tag=tag)
35        self.parser = make_parser(['xml.sax.IncrementalParser'])
36        self.parser.setContentHandler(self.handler)
37
38    ############################
39    def transform(self, record: str):
40        """Aggregate, returning None until we're done, then return record."""
41
42        # See if it's something we can process, and if not, try digesting
43        if not self.can_process_record(record):  # inherited from BaseModule()
44            return self.digest_record(record)  # inherited from BaseModule()
45
46        with self.buffer_lock:
47            # Feed record to the incremental parser
48            self.buffer += record + '\n'
49            self.parser.feed(record)
50            logging.debug('transform() got line: %s', record)
51
52            # If the record completes and XML record, it will be added to the
53            # queue in self.handler.items() - pop it off and return
54            if self.handler.complete():
55                xml_record = self.buffer
56                self.buffer = ''
57                self.parser.close()
58                self.parser.reset()
59                self.handler.reset()
60                logging.debug('transform() got closing tag: %s', xml_record)
61                return xml_record
62
63        # Otherwise go home emptyhanded
64        return None

Aggregate passed lines of XML until a complete XML record whose outermost element matches 'tag' has been seen, then pass it on as a single record.

XMLAggregatorTransform(tag, **kwargs)
20    def __init__(self, tag, **kwargs):
21        """
22        'tag' should be the identity of the top-level XML element that
23        we're expecting to read, e.g. 'OSU_DAS_Record'.
24        """
25        super().__init__(**kwargs)  # processes 'quiet' and type hints
26        self.tag = tag
27
28        # Only let one thread touch buffer at a time. Of course, if we're
29        # getting interleaved lines from different XML records here, we're
30        # screwed anyway.
31        self.buffer_lock = Lock()
32        self.buffer = ''
33
34        self.handler = _XMLHandler(tag=tag)
35        self.parser = make_parser(['xml.sax.IncrementalParser'])
36        self.parser.setContentHandler(self.handler)

'tag' should be the identity of the top-level XML element that we're expecting to read, e.g. 'OSU_DAS_Record'.

tag
buffer_lock
buffer
handler
parser
def transform(self, record: str):
39    def transform(self, record: str):
40        """Aggregate, returning None until we're done, then return record."""
41
42        # See if it's something we can process, and if not, try digesting
43        if not self.can_process_record(record):  # inherited from BaseModule()
44            return self.digest_record(record)  # inherited from BaseModule()
45
46        with self.buffer_lock:
47            # Feed record to the incremental parser
48            self.buffer += record + '\n'
49            self.parser.feed(record)
50            logging.debug('transform() got line: %s', record)
51
52            # If the record completes and XML record, it will be added to the
53            # queue in self.handler.items() - pop it off and return
54            if self.handler.complete():
55                xml_record = self.buffer
56                self.buffer = ''
57                self.parser.close()
58                self.parser.reset()
59                self.handler.reset()
60                logging.debug('transform() got closing tag: %s', xml_record)
61                return xml_record
62
63        # Otherwise go home emptyhanded
64        return None

Aggregate, returning None until we're done, then return record.