-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathorchestrator.py
executable file
·66 lines (56 loc) · 2.03 KB
/
orchestrator.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
#!/usr/bin/env python3
import argparse
import logging
from multiprocessing import pool
import sys
from respdiff import cli, sendrecv
from respdiff.database import LMDB, MetaDatabase
def main():
cli.setup_logging()
parser = argparse.ArgumentParser(
description="read queries from LMDB, send them in parallel to servers "
"listed in configuration file, and record answers into LMDB"
)
cli.add_arg_envdir(parser)
cli.add_arg_config(parser)
parser.add_argument(
"--ignore-timeout",
action="store_true",
help="continue despite consecutive timeouts from resolvers",
)
args = parser.parse_args()
sendrecv.module_init(args)
with LMDB(args.envdir) as lmdb:
meta = MetaDatabase(lmdb, args.cfg["servers"]["names"], create=True)
meta.write_version()
meta.write_start_time()
lmdb.open_db(LMDB.QUERIES)
adb = lmdb.open_db(LMDB.ANSWERS, create=True, check_notexists=True)
qstream = lmdb.key_value_stream(LMDB.QUERIES)
txn = lmdb.env.begin(adb, write=True)
try:
# process queries in parallel
with pool.Pool(
processes=args.cfg["sendrecv"]["jobs"], initializer=sendrecv.worker_init
) as p:
i = 0
for qkey, blob in p.imap_unordered(
sendrecv.worker_perform_query, qstream, chunksize=100
):
i += 1
if i % 10000 == 0:
logging.info("Received {:d} answers".format(i))
txn.put(qkey, blob)
except KeyboardInterrupt:
logging.info("SIGINT received, exiting...")
sys.exit(130)
except RuntimeError as err:
logging.error(err)
sys.exit(1)
finally:
# attempt to preserve data if something went wrong (or not)
logging.debug("Comitting LMDB transaction...")
txn.commit()
meta.write_end_time()
if __name__ == "__main__":
main()