##// END OF EJS Templates
wireprotov2peer: stream decoded responses...
wireprotov2peer: stream decoded responses Previously, wire protocol version 2 would buffer all response data. Only once all data was received did we CBOR decode it and resolve the future associated with the command. This was obviously not desirable. In future commits that introduce large response payloads, this caused significant memory bloat and slowed down client operations due to waiting on the server. This commit refactors the response handling code so that response data can be streamed. Command response objects now contain a buffered CBOR decoder. As new data arrives, it is fed into the decoder. Decoded objects are made available to the generator as they are decoded. Because there is a separate thread processing incoming frames and feeding data into the response object, there is the potential for race conditions when mutating response objects. So a lock has been added to guard access to critical state variables. Because the generator emitting decoded objects needs to wait on those objects to become available, we've added an Event for the generator to wait on so it doesn't busy loop. This does mean there is the potential for deadlocks. And I'm pretty sure they can occur in some scenarios. We already have a handful of TODOs around this. But I've added some more. Fixing this will likely require moving the background thread receiving frames into clienthandler. We likely would have done this anyway when implementing the client bits for the SSH transport. Test output changes because the initial CBOR map holding the overall response state is now always handled internally by the response object. Differential Revision: https://phab.mercurial-scm.org/D4474

File last commit:

r39288:9a81f126 default
r39597:d06834e0 default
Show More
catapipe.py
90 lines | 2.8 KiB | text/x-python | PythonLexer
Augie Fackler
contrib: new script to read events from a named pipe and emit catapult traces...
r39288 #!/usr/bin/env python3
#
# Copyright 2018 Google LLC.
#
# This software may be used and distributed according to the terms of the
# GNU General Public License version 2 or any later version.
"""Tool read primitive events from a pipe to produce a catapult trace.
For now the event stream supports
START $SESSIONID ...
and
END $SESSIONID ...
events. Everything after the SESSIONID (which must not contain spaces)
is used as a label for the event. Events are timestamped as of when
they arrive in this process and are then used to produce catapult
traces that can be loaded in Chrome's about:tracing utility. It's
important that the event stream *into* this process stay simple,
because we have to emit it from the shell scripts produced by
run-tests.py.
Typically you'll want to place the path to the named pipe in the
HGCATAPULTSERVERPIPE environment variable, which both run-tests and hg
understand.
"""
from __future__ import absolute_import, print_function
import argparse
import json
import os
Boris Feld
contrib: use a monotonic timer in catapipe...
r39550 import timeit
Augie Fackler
contrib: new script to read events from a named pipe and emit catapult traces...
r39288
_TYPEMAP = {
'START': 'B',
'END': 'E',
}
_threadmap = {}
Boris Feld
contrib: use a monotonic timer in catapipe...
r39550 # Timeit already contains the whole logic about which timer to use based on
# Python version and OS
timer = timeit.default_timer
Augie Fackler
contrib: new script to read events from a named pipe and emit catapult traces...
r39288 def main():
parser = argparse.ArgumentParser()
parser.add_argument('pipe', type=str, nargs=1,
help='Path of named pipe to create and listen on.')
parser.add_argument('output', default='trace.json', type=str, nargs='?',
Boris Feld
contrib: fix catapipe output argument documentation...
r39549 help='Path of json file to create where the traces '
'will be stored.')
Augie Fackler
contrib: new script to read events from a named pipe and emit catapult traces...
r39288 parser.add_argument('--debug', default=False, action='store_true',
help='Print useful debug messages')
args = parser.parse_args()
fn = args.pipe[0]
os.mkfifo(fn)
try:
with open(fn) as f, open(args.output, 'w') as out:
out.write('[\n')
Boris Feld
contrib: use a monotonic timer in catapipe...
r39550 start = timer()
Augie Fackler
contrib: new script to read events from a named pipe and emit catapult traces...
r39288 while True:
ev = f.readline().strip()
if not ev:
continue
Boris Feld
contrib: use a monotonic timer in catapipe...
r39550 now = timer()
Augie Fackler
contrib: new script to read events from a named pipe and emit catapult traces...
r39288 if args.debug:
print(ev)
verb, session, label = ev.split(' ', 2)
if session not in _threadmap:
_threadmap[session] = len(_threadmap)
pid = _threadmap[session]
ts_micros = (now - start).total_seconds() * 1000000
out.write(json.dumps(
{
"name": label,
"cat": "misc",
"ph": _TYPEMAP[verb],
"ts": ts_micros,
"pid": pid,
"tid": 1,
"args": {}
}))
out.write(',\n')
finally:
os.unlink(fn)
if __name__ == '__main__':
main()