|
3143 | 3143 | "\n", |
3144 | 3144 | " return date_ranges\n", |
3145 | 3145 | "\n", |
| 3146 | + " def _verify_running_lock(self, existing_df, workspace_id: str, dataflow_id: str):\n", |
| 3147 | + " \"\"\"\n", |
| 3148 | + " Query the Fabric API to verify whether the job reflected by a 'Running' tracking row\n", |
| 3149 | + " is genuinely active or is already finished (stale lock).\n", |
| 3150 | + "\n", |
| 3151 | + " Called only from the concurrency guard when the row is young enough to block and\n", |
| 3152 | + " force_run is False. Never raises — all errors are absorbed and returned as inconclusive.\n", |
| 3153 | + "\n", |
| 3154 | + " Returns a 2-tuple:\n", |
| 3155 | + " (\"confirmed_running\", api_status_str) — a live job is still executing\n", |
| 3156 | + " (\"confirmed_terminal\", api_status_str) — job finished; tracking row is stale\n", |
| 3157 | + " (\"inconclusive\", None) — API unavailable; caller falls back\n", |
| 3158 | + " to time-based stale check\n", |
| 3159 | + " \"\"\"\n", |
| 3160 | + " _TERMINAL = {\"Success\", \"Failed\", \"Cancelled\", \"Completed\", \"Error\", \"Succeeded\", \"Deduped\"}\n", |
| 3161 | + "\n", |
| 3162 | + " try:\n", |
| 3163 | + " row = existing_df.iloc[0]\n", |
| 3164 | + " job_instance_id = row.get('job_instance_id') or None\n", |
| 3165 | + " stored_cicd = row.get('is_cicd_dataflow')\n", |
| 3166 | + " # Coerce stored bit/int to bool if present\n", |
| 3167 | + " if stored_cicd is not None:\n", |
| 3168 | + " try:\n", |
| 3169 | + " stored_cicd = bool(int(stored_cicd))\n", |
| 3170 | + " except (TypeError, ValueError):\n", |
| 3171 | + " stored_cicd = None\n", |
| 3172 | + "\n", |
| 3173 | + " # ── Branch a: exact job-instance poll (most reliable) ────────────────────\n", |
| 3174 | + " if job_instance_id:\n", |
| 3175 | + " try:\n", |
| 3176 | + " endpoint = (\n", |
| 3177 | + " f\"/v1/workspaces/{workspace_id}/items/{dataflow_id}\"\n", |
| 3178 | + " f\"/jobs/instances/{job_instance_id}\"\n", |
| 3179 | + " )\n", |
| 3180 | + " resp = self.client.get(endpoint)\n", |
| 3181 | + " sc = getattr(resp, 'status_code', 200)\n", |
| 3182 | + "\n", |
| 3183 | + " if sc == 404:\n", |
| 3184 | + " # Job does not exist on Fabric — was never created or has been\n", |
| 3185 | + " # cleaned up. The tracking row's Running status is definitely stale.\n", |
| 3186 | + " self.logger.info(\n", |
| 3187 | + " f\"Lock verification: job {job_instance_id} not found (HTTP 404) \"\n", |
| 3188 | + " f\"for dataflow {dataflow_id} — lock is stale (job never existed).\"\n", |
| 3189 | + " )\n", |
| 3190 | + " return (\"confirmed_terminal\", \"Failed\")\n", |
| 3191 | + "\n", |
| 3192 | + " if sc == 429:\n", |
| 3193 | + " retry_after = int(resp.headers.get('Retry-After', 30))\n", |
| 3194 | + " self.logger.warning(\n", |
| 3195 | + " f\"Lock verification rate-limited (429). Sleeping {retry_after}s.\"\n", |
| 3196 | + " )\n", |
| 3197 | + " time.sleep(retry_after)\n", |
| 3198 | + " resp = self.client.get(endpoint)\n", |
| 3199 | + " sc = getattr(resp, 'status_code', 200)\n", |
| 3200 | + "\n", |
| 3201 | + " if sc in (200, 202):\n", |
| 3202 | + " api_status = resp.json().get(\"status\", \"Unknown\")\n", |
| 3203 | + " self.logger.info(\n", |
| 3204 | + " f\"Lock verification: job {job_instance_id} API status={api_status!r} \"\n", |
| 3205 | + " f\"for dataflow {dataflow_id}.\"\n", |
| 3206 | + " )\n", |
| 3207 | + " if api_status in _TERMINAL:\n", |
| 3208 | + " return (\"confirmed_terminal\", api_status)\n", |
| 3209 | + " return (\"confirmed_running\", api_status)\n", |
| 3210 | + "\n", |
| 3211 | + " # 403, 400, 5xx, etc. — cannot verify\n", |
| 3212 | + " self.logger.warning(\n", |
| 3213 | + " f\"Lock verification: exact-instance poll returned HTTP {sc} \"\n", |
| 3214 | + " f\"for job {job_instance_id} — inconclusive.\"\n", |
| 3215 | + " )\n", |
| 3216 | + " return (\"inconclusive\", None)\n", |
| 3217 | + "\n", |
| 3218 | + " except Exception as _e:\n", |
| 3219 | + " self.logger.warning(\n", |
| 3220 | + " f\"Lock verification: exact-instance poll raised {_e} \"\n", |
| 3221 | + " f\"for job {job_instance_id} — inconclusive.\"\n", |
| 3222 | + " )\n", |
| 3223 | + " return (\"inconclusive\", None)\n", |
| 3224 | + "\n", |
| 3225 | + " # ── Branch b: CI/CD dataflow, no job_instance_id — list endpoint ─────────\n", |
| 3226 | + " if stored_cicd is True:\n", |
| 3227 | + " try:\n", |
| 3228 | + " endpoint = f\"/v1/workspaces/{workspace_id}/items/{dataflow_id}/jobs/instances\"\n", |
| 3229 | + " resp = self.client.get(endpoint)\n", |
| 3230 | + " sc = getattr(resp, 'status_code', 200)\n", |
| 3231 | + "\n", |
| 3232 | + " if sc == 429:\n", |
| 3233 | + " retry_after = int(resp.headers.get('Retry-After', 30))\n", |
| 3234 | + " self.logger.warning(\n", |
| 3235 | + " f\"Lock verification rate-limited (429). Sleeping {retry_after}s.\"\n", |
| 3236 | + " )\n", |
| 3237 | + " time.sleep(retry_after)\n", |
| 3238 | + " resp = self.client.get(endpoint)\n", |
| 3239 | + " sc = getattr(resp, 'status_code', 200)\n", |
| 3240 | + "\n", |
| 3241 | + " if sc in (200, 202):\n", |
| 3242 | + " data = resp.json()\n", |
| 3243 | + " jobs = data.get('value', data) if isinstance(data, dict) else data\n", |
| 3244 | + " if isinstance(jobs, list) and len(jobs) > 0:\n", |
| 3245 | + " api_status = jobs[0].get(\"status\", \"Unknown\")\n", |
| 3246 | + " self.logger.info(\n", |
| 3247 | + " f\"Lock verification (CI/CD list): most-recent job status={api_status!r} \"\n", |
| 3248 | + " f\"for dataflow {dataflow_id}.\"\n", |
| 3249 | + " )\n", |
| 3250 | + " if api_status in _TERMINAL:\n", |
| 3251 | + " return (\"confirmed_terminal\", api_status)\n", |
| 3252 | + " return (\"confirmed_running\", api_status)\n", |
| 3253 | + " else:\n", |
| 3254 | + " # No jobs found — nothing is running\n", |
| 3255 | + " self.logger.info(\n", |
| 3256 | + " f\"Lock verification (CI/CD list): no job instances found \"\n", |
| 3257 | + " f\"for dataflow {dataflow_id} — lock is stale.\"\n", |
| 3258 | + " )\n", |
| 3259 | + " return (\"confirmed_terminal\", \"Failed\")\n", |
| 3260 | + "\n", |
| 3261 | + " self.logger.warning(\n", |
| 3262 | + " f\"Lock verification (CI/CD list): HTTP {sc} for dataflow {dataflow_id} \"\n", |
| 3263 | + " \"— inconclusive.\"\n", |
| 3264 | + " )\n", |
| 3265 | + " return (\"inconclusive\", None)\n", |
| 3266 | + "\n", |
| 3267 | + " except Exception as _e:\n", |
| 3268 | + " self.logger.warning(\n", |
| 3269 | + " f\"Lock verification (CI/CD list) raised {_e} for dataflow {dataflow_id} \"\n", |
| 3270 | + " \"— inconclusive.\"\n", |
| 3271 | + " )\n", |
| 3272 | + " return (\"inconclusive\", None)\n", |
| 3273 | + "\n", |
| 3274 | + " # ── Branch c: regular dataflow — transactions endpoint ────────────────────\n", |
| 3275 | + " if stored_cicd is False:\n", |
| 3276 | + " try:\n", |
| 3277 | + " endpoint = f\"/v1.0/myorg/groups/{workspace_id}/dataflows/{dataflow_id}/transactions\"\n", |
| 3278 | + " resp = self.pbi_client.get(endpoint)\n", |
| 3279 | + " sc = getattr(resp, 'status_code', 200)\n", |
| 3280 | + "\n", |
| 3281 | + " if sc == 429:\n", |
| 3282 | + " retry_after = int(resp.headers.get('Retry-After', 30))\n", |
| 3283 | + " self.logger.warning(\n", |
| 3284 | + " f\"Lock verification rate-limited (429). Sleeping {retry_after}s.\"\n", |
| 3285 | + " )\n", |
| 3286 | + " time.sleep(retry_after)\n", |
| 3287 | + " resp = self.pbi_client.get(endpoint)\n", |
| 3288 | + " sc = getattr(resp, 'status_code', 200)\n", |
| 3289 | + "\n", |
| 3290 | + " if sc in (200, 201):\n", |
| 3291 | + " data = resp.json()\n", |
| 3292 | + " txns = data.get('value', []) if isinstance(data, dict) else []\n", |
| 3293 | + " if txns:\n", |
| 3294 | + " api_status = txns[0].get(\"status\", \"Unknown\")\n", |
| 3295 | + " self.logger.info(\n", |
| 3296 | + " f\"Lock verification (transactions): most-recent txn status={api_status!r} \"\n", |
| 3297 | + " f\"for dataflow {dataflow_id}.\"\n", |
| 3298 | + " )\n", |
| 3299 | + " if api_status in _TERMINAL:\n", |
| 3300 | + " return (\"confirmed_terminal\", api_status)\n", |
| 3301 | + " return (\"confirmed_running\", api_status)\n", |
| 3302 | + " else:\n", |
| 3303 | + " # No transactions found — nothing is running\n", |
| 3304 | + " self.logger.info(\n", |
| 3305 | + " f\"Lock verification (transactions): no transactions found \"\n", |
| 3306 | + " f\"for dataflow {dataflow_id} — lock is stale.\"\n", |
| 3307 | + " )\n", |
| 3308 | + " return (\"confirmed_terminal\", \"Failed\")\n", |
| 3309 | + "\n", |
| 3310 | + " self.logger.warning(\n", |
| 3311 | + " f\"Lock verification (transactions): HTTP {sc} for dataflow {dataflow_id} \"\n", |
| 3312 | + " \"— inconclusive.\"\n", |
| 3313 | + " )\n", |
| 3314 | + " return (\"inconclusive\", None)\n", |
| 3315 | + "\n", |
| 3316 | + " except Exception as _e:\n", |
| 3317 | + " self.logger.warning(\n", |
| 3318 | + " f\"Lock verification (transactions) raised {_e} for dataflow {dataflow_id} \"\n", |
| 3319 | + " \"— inconclusive.\"\n", |
| 3320 | + " )\n", |
| 3321 | + " return (\"inconclusive\", None)\n", |
| 3322 | + "\n", |
| 3323 | + " # ── Branch d: dataflow type unknown, no job_instance_id ───────────────────\n", |
| 3324 | + " self.logger.warning(\n", |
| 3325 | + " f\"Lock verification skipped for dataflow {dataflow_id}: \"\n", |
| 3326 | + " \"is_cicd_dataflow not stored and no job_instance_id available — inconclusive.\"\n", |
| 3327 | + " )\n", |
| 3328 | + " return (\"inconclusive\", None)\n", |
| 3329 | + "\n", |
| 3330 | + " except Exception as _outer:\n", |
| 3331 | + " self.logger.warning(\n", |
| 3332 | + " f\"Lock verification raised unexpected error for dataflow {dataflow_id}: {_outer} \"\n", |
| 3333 | + " \"— inconclusive.\"\n", |
| 3334 | + " )\n", |
| 3335 | + " return (\"inconclusive\", None)\n", |
| 3336 | + "\n", |
3146 | 3337 | " def execute_incremental_refresh(self,\n", |
3147 | 3338 | " workspace_id: str,\n", |
3148 | 3339 | " dataflow_id: str,\n", |
|
3575 | 3766 | " )\n", |
3576 | 3767 | "\n", |
3577 | 3768 | " # Concurrency guard: refuse to start if another run is already active for this dataflow\n", |
3578 | | - " # unless the existing run is stale (older than 2x timeout) or force_run=True.\n", |
| 3769 | + " # unless the existing run is stale (older than 2x timeout), force_run=True, or the\n", |
| 3770 | + " # Fabric API confirms the job already finished (API-verified stale lock).\n", |
3579 | 3771 | " existing_df = self.get_incremental(dataflow_id, is_adhoc=False)\n", |
3580 | 3772 | " if existing_df is not None:\n", |
3581 | 3773 | " existing_status = existing_df.iloc[0].get('status', '')\n", |
|
3588 | 3780 | " stale_threshold = timedelta(minutes=timeout_minutes * 2)\n", |
3589 | 3781 | " if age < stale_threshold:\n", |
3590 | 3782 | " if not force_run:\n", |
3591 | | - " raise RuntimeError(\n", |
3592 | | - " f\"Dataflow {dataflow_id} appears to have an active run \"\n", |
3593 | | - " f\"(status=Running, last updated {age} ago). \"\n", |
3594 | | - " \"Wait for it to complete, or pass force_run=True to override.\"\n", |
| 3783 | + " # Before blocking, verify actual job state with the Fabric API.\n", |
| 3784 | + " # This automatically clears locks left by crashes/OOM kills that\n", |
| 3785 | + " # occurred before any dataflow job was triggered, without requiring\n", |
| 3786 | + " # operator intervention or force_run=True.\n", |
| 3787 | + " _lock_verdict, _api_status = self._verify_running_lock(\n", |
| 3788 | + " existing_df, workspace_id, dataflow_id\n", |
| 3789 | + " )\n", |
| 3790 | + "\n", |
| 3791 | + " if _lock_verdict == \"confirmed_terminal\":\n", |
| 3792 | + " # API confirmed the job is done — the tracking row is stale.\n", |
| 3793 | + " # Repair only status + update_time; preserve last_backup_table\n", |
| 3794 | + " # and job_instance_id so _recover_from_backup_if_needed can\n", |
| 3795 | + " # still find the backup on the next step if one exists.\n", |
| 3796 | + " _repair_status = _api_status or \"Failed\"\n", |
| 3797 | + " try:\n", |
| 3798 | + " self._execute_with_retry(\n", |
| 3799 | + " f\"UPDATE {self.incremental_table} \"\n", |
| 3800 | + " \"SET [status] = ?, [update_time] = ? \"\n", |
| 3801 | + " \"WHERE [dataflow_id] = ? \"\n", |
| 3802 | + " \" AND ([is_adhoc] IS NULL OR [is_adhoc] = 0) \"\n", |
| 3803 | + " \" AND [status] = 'Running'\",\n", |
| 3804 | + " (_repair_status,\n", |
| 3805 | + " datetime.utcnow().strftime('%Y-%m-%d %H:%M:%S.%f')[:-3],\n", |
| 3806 | + " dataflow_id)\n", |
| 3807 | + " )\n", |
| 3808 | + " self.logger.warning(\n", |
| 3809 | + " f\"Stale Running lock repaired for dataflow {dataflow_id}: \"\n", |
| 3810 | + " f\"API reported {_api_status!r}, tracking row updated to \"\n", |
| 3811 | + " f\"{_repair_status!r} (age={age}). Proceeding with execution.\"\n", |
| 3812 | + " )\n", |
| 3813 | + " # Refresh the snapshot so downstream branch-selection logic\n", |
| 3814 | + " # (last_status, last_refresh_data) sees the corrected status.\n", |
| 3815 | + " existing_df = self.get_incremental(dataflow_id, is_adhoc=False)\n", |
| 3816 | + " except Exception as _repair_err:\n", |
| 3817 | + " self.logger.error(\n", |
| 3818 | + " f\"API confirmed stale lock for dataflow {dataflow_id} \"\n", |
| 3819 | + " f\"but row repair failed: {_repair_err}. \"\n", |
| 3820 | + " \"Raising concurrency error to avoid unsafe state.\"\n", |
| 3821 | + " )\n", |
| 3822 | + " raise RuntimeError(\n", |
| 3823 | + " f\"Dataflow {dataflow_id} appears to have an active run \"\n", |
| 3824 | + " f\"(status=Running, last updated {age} ago). \"\n", |
| 3825 | + " \"Wait for it to complete, or pass force_run=True to override.\"\n", |
| 3826 | + " )\n", |
| 3827 | + "\n", |
| 3828 | + " elif _lock_verdict == \"confirmed_running\":\n", |
| 3829 | + " # API confirmed a live job is still executing — lock is legitimate.\n", |
| 3830 | + " raise RuntimeError(\n", |
| 3831 | + " f\"Dataflow {dataflow_id} has a confirmed active run \"\n", |
| 3832 | + " f\"(API status: {_api_status!r}, tracking row last updated {age} ago). \"\n", |
| 3833 | + " \"Wait for it to complete, or pass force_run=True to override.\"\n", |
| 3834 | + " )\n", |
| 3835 | + "\n", |
| 3836 | + " else:\n", |
| 3837 | + " # inconclusive — API unavailable or dataflow type unknown.\n", |
| 3838 | + " # Conservative fallback: trust the time-based stale check.\n", |
| 3839 | + " self.logger.warning(\n", |
| 3840 | + " f\"API verification inconclusive for dataflow {dataflow_id} \"\n", |
| 3841 | + " f\"(age={age}). Falling back to time-based stale check.\"\n", |
| 3842 | + " )\n", |
| 3843 | + " raise RuntimeError(\n", |
| 3844 | + " f\"Dataflow {dataflow_id} appears to have an active run \"\n", |
| 3845 | + " f\"(status=Running, last updated {age} ago). \"\n", |
| 3846 | + " \"Wait for it to complete, or pass force_run=True to override.\"\n", |
| 3847 | + " )\n", |
| 3848 | + "\n", |
| 3849 | + " else:\n", |
| 3850 | + " self.logger.warning(\n", |
| 3851 | + " f\"force_run=True: proceeding despite active-looking run \"\n", |
| 3852 | + " f\"(last updated {age} ago)\"\n", |
3595 | 3853 | " )\n", |
3596 | | - " self.logger.warning(\n", |
3597 | | - " f\"force_run=True: proceeding despite active-looking run \"\n", |
3598 | | - " f\"(last updated {age} ago)\"\n", |
3599 | | - " )\n", |
3600 | 3854 | "\n", |
3601 | 3855 | " # Flags for tracking CAS lock ownership and bucket execution entry.\n", |
3602 | 3856 | " # Used by the outer except block and the no-new-data branch to repair the\n", |
|
0 commit comments