(self, stream, drop=True, n_items=1)
| 75 | self.graph = Graph(self.conn, GRAPH_ID) |
| 76 | |
| 77 | def consumeStream(self, stream, drop=True, n_items=1): |
| 78 | # wait for telemetry stream to be created |
| 79 | t = 'none' # type of stream_key |
| 80 | |
| 81 | while t == 'none': |
| 82 | t = self.conn.type(stream) |
| 83 | |
| 84 | self.env.assertEquals(t, "stream") |
| 85 | |
| 86 | # convert stream events to LoggedQueries |
| 87 | logged_queries = [] |
| 88 | streams = {stream: '0-0'} |
| 89 | |
| 90 | elapsed = 10 |
| 91 | while len(logged_queries) < n_items and elapsed > 0: |
| 92 | # read messages from the stream |
| 93 | messages = self.conn.xread(streams, block=0) |
| 94 | |
| 95 | if len(messages) > 0: |
| 96 | # process each message received |
| 97 | stream_messages = messages[0][1] |
| 98 | for message_id, message_payload in stream_messages: |
| 99 | logged_queries.append(LoggedQuery(message_payload)) |
| 100 | |
| 101 | # update stream last ID |
| 102 | streams[stream] = stream_messages[-1][0] |
| 103 | |
| 104 | time.sleep(0.2) |
| 105 | elapsed -= 0.2 |
| 106 | |
| 107 | # drop stream |
| 108 | if drop: |
| 109 | self.conn.delete(stream) |
| 110 | |
| 111 | # reverse order to match expected order of events |
| 112 | logged_queries.reverse() |
| 113 | |
| 114 | return logged_queries |
| 115 | |
| 116 | def assertLoggedQuery(self, logged_query, query, utilized_cache): |
| 117 | # validate event values |
no test coverage detected