The task container must periodically update 'lastUpdatedAt' field in the progress file as a heartbeat that worker container monitors. The task will fail if there's no change in 'lastUpdatedAt' over 5 min. Parameters: :param progress_file_path: progress file path
(progress_file_path: str)
| 196 | |
| 197 | |
| 198 | def heartbeat(progress_file_path: str): |
| 199 | """ |
| 200 | The task container must periodically update 'lastUpdatedAt' field |
| 201 | in the progress file as a heartbeat that worker container monitors. |
| 202 | The task will fail if there's no change in 'lastUpdatedAt' over 5 min. |
| 203 | |
| 204 | Parameters: |
| 205 | :param progress_file_path: progress file path |
| 206 | """ |
| 207 | while TASK_RUNNING: |
| 208 | ts = datetime.datetime.now(datetime.timezone.utc).strftime( |
| 209 | "%Y-%m-%dT%H:%M:%S.%fZ" |
| 210 | ) |
| 211 | if not os.path.isfile(progress_file_path): |
| 212 | logger.info("Creating progress file") |
| 213 | progress = Progress( |
| 214 | taskId=os.getenv("NVCT_TASK_ID", ""), |
| 215 | percentComplete=0, |
| 216 | name="", |
| 217 | lastUpdatedAt=ts, |
| 218 | metadata=dict(), |
| 219 | ) |
| 220 | with open(progress_file_path, "w") as f: |
| 221 | f.write(progress.output_to_json()) |
| 222 | |
| 223 | else: |
| 224 | progress = parse_progress_file(progress_file_path) |
| 225 | progress.lastUpdatedAt = ts |
| 226 | |
| 227 | temp_file = progress_file_path + ".tmp" |
| 228 | with open(temp_file, "w") as f: |
| 229 | f.write(progress.output_to_json()) |
| 230 | |
| 231 | try: |
| 232 | os.rename(temp_file, progress_file_path) |
| 233 | except Exception as e: |
| 234 | logger.exception(f"Failed to write to progress file: {e}") |
| 235 | raise e |
| 236 | |
| 237 | logger.info(f"Updated timestamp in progress file {ts}") |
| 238 | time.sleep(60 + random.randint(1, 10)) |
| 239 | |
| 240 | |
| 241 | def wait_for_secrets_file(file_path, max_retries=30, retry_delay=5.0): |
nothing calls this directly
no test coverage detected