MCPcopy Create free account
hub / github.com/RedisGraph/RedisGraph / consumeStream

Method consumeStream

tests/flow/test_graph_info.py:77–114  ·  view source on GitHub ↗
(self, stream, drop=True, n_items=1)

Source from the content-addressed store, hash-verified

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

Callers 4

test03_long_queryMethod · 0.95
test05_rename_graphMethod · 0.95

Calls 2

LoggedQueryClass · 0.85
typeMethod · 0.80

Tested by

no test coverage detected