(params, requests)
| 934 | |
| 935 | |
| 936 | def execute_batch(params, requests): |
| 937 | # type: (KeeperParams, List[dict]) -> List[dict] |
| 938 | responses = [] |
| 939 | if not requests: |
| 940 | return responses |
| 941 | |
| 942 | throttle_delay = 10 |
| 943 | chunk_size = 999 |
| 944 | queue = requests.copy() |
| 945 | delay_next_batch = False |
| 946 | |
| 947 | while len(queue) > 0: |
| 948 | chunk = queue[:chunk_size] |
| 949 | queue = queue[chunk_size:] |
| 950 | |
| 951 | rq = { |
| 952 | 'command': 'execute', |
| 953 | 'requests': chunk |
| 954 | } |
| 955 | try: |
| 956 | if delay_next_batch: |
| 957 | time.sleep(throttle_delay) |
| 958 | rs = communicate(params, rq) |
| 959 | if 'results' in rs: |
| 960 | results = rs['results'] # type: list |
| 961 | if len(results) > 0: |
| 962 | error_rs = results[-1] |
| 963 | throttled = error_rs.get('result') != 'success' and error_rs.get('result_code') == 'throttled' |
| 964 | if throttled: |
| 965 | delay_next_batch = True |
| 966 | results.pop() |
| 967 | responses.extend(results) |
| 968 | |
| 969 | if len(results) < len(chunk): |
| 970 | queue = chunk[len(results):] + queue |
| 971 | except Exception as e: |
| 972 | logging.error(e) |
| 973 | |
| 974 | return responses |
| 975 | |
| 976 | |
| 977 | def update_record(params, record, **kwargs): |
no test coverage detected