Skip to content

Commit 127aae1

Browse files
committed
fix(load): a failed status check must not discard a running job
Review found a regression this PR introduced, reachable from the opposite direction to the one it fixes. `_poll_job` aborted on the first bad `GET /v1/jobs/{id}`. That exception reaches `ManagedDatabaseClient._request_with_retry`, which classifies a 502/503 as transient and re-invokes THE WHOLE OPERATION for every non-append mode -- so a single blipped status check discards a healthy running job and submits a second load against a table whose write lock the first one still holds. 409 RESOURCE_LOCKED, which is also transient, and round it goes. Worse than the behaviour it replaces, not better: before, a lost request took with it a load that was genuinely being torn down. Now the first job is neither going away nor releasing anything, and a load polls for up to `_LOAD_JOB_TIMEOUT_S`, so there are hundreds of status checks where there used to be one request -- every one a fresh chance to trigger the collision. Exposure up, collision from incidental to certain. So the poll tolerates a failed status check and keeps waiting, bounded by consecutive failures rather than by the first one. An isolated 502 is noise; a run of them means the API is gone and there is nothing to wait for. The poll's own deadline bounds the total wait either way. THE JOB ID IS NOW RETURNED. The docstring claimed the id removed `append`'s "did it land?" ambiguity, which was not true as implemented: it was a local, discarded on every failure path, and never reached the caller. `LoadManagedTableResult` carries it now, for the same reason `CreateIndexResult` does. `append` stays non-retryable -- knowing the id makes the question answerable, it does not make a blind re-submission safe, and that call belongs to the caller. The docstring says that instead of the stronger thing. `result` is read with `getattr(..., "actual_instance", ...)`, matching the index path rather than silently disagreeing with it: the strict read raised a bare AttributeError if the model ever arrived unwrapped. Did NOT add a strict isinstance check on the inline response. The index path can assert its inline type because it names it positively first; here the inline shape is read duck-typed and callers and tests rely on that, so asserting it narrows an interface this change has no business narrowing. Only the job branch is type-checked. Import order fixed for ruff's I rule; `ruff check` clean. 149 passed.
1 parent 115efd6 commit 127aae1

3 files changed

Lines changed: 130 additions & 7 deletions

File tree

‎hotdata_framework/client.py‎

Lines changed: 37 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -24,8 +24,8 @@
2424
from hotdata.models.database_default_table_decl import DatabaseDefaultTableDecl
2525
from hotdata.models.index_info_response import IndexInfoResponse
2626
from hotdata.models.job_status_response import JobStatusResponse
27-
from hotdata.models.load_managed_table_response import LoadManagedTableResponse
2827
from hotdata.models.load_managed_table_request import LoadManagedTableRequest
28+
from hotdata.models.load_managed_table_response import LoadManagedTableResponse
2929
from hotdata.models.query_request import QueryRequest
3030
from hotdata.models.query_response import QueryResponse
3131
from hotdata.models.submit_job_response import SubmitJobResponse
@@ -95,6 +95,9 @@
9595
# first place. Bounded rather than unbounded so a wedged job still surfaces.
9696
_LOAD_JOB_TIMEOUT_S = 3600.0
9797

98+
# Consecutive failed status checks before a poll gives up.
99+
_JOB_POLL_MAX_CONSECUTIVE_ERRORS = 5
100+
98101

99102
@dataclass(frozen=True)
100103
class ResultSummary:
@@ -461,14 +464,21 @@ def load_managed_table(
461464
)
462465
except ApiException as e:
463466
raise RuntimeError(api_error_message(e)) from e
467+
# Only the job branch is type-checked. The index path can also assert its
468+
# inline type because it names it positively first; here the inline shape is
469+
# read duck-typed, which callers and tests already rely on, so asserting it
470+
# would narrow an interface this change has no business narrowing.
471+
job_id: str | None = None
464472
if isinstance(loaded, SubmitJobResponse):
465-
loaded = self._load_response_from_job(loaded.id)
473+
job_id = loaded.id
474+
loaded = self._load_response_from_job(job_id)
466475
return LoadManagedTableResult(
467476
connection_id=loaded.connection_id,
468477
schema_name=loaded.schema_name,
469478
table_name=loaded.table_name,
470479
row_count=loaded.row_count,
471480
full_name=f"{db.id}.{loaded.schema_name}.{loaded.table_name}",
481+
job_id=job_id,
472482
)
473483

474484
def add_managed_table(
@@ -911,11 +921,27 @@ def _poll_job(
911921
jobs = self._jobs_api()
912922
deadline = time.monotonic() + timeout_s
913923
last: JobStatusResponse | None = None
924+
consecutive_errors = 0
914925
while time.monotonic() < deadline:
915926
try:
916927
last = jobs.get_job(job_id)
928+
consecutive_errors = 0
917929
except ApiException as e:
918-
raise RuntimeError(api_error_message(e)) from e
930+
# A failed STATUS CHECK is not a failed job. Aborting here throws
931+
# away work that is still running, and the caller's retry then
932+
# re-submits it while the original still holds its resources -- so
933+
# one blip becomes a collision with the job it just abandoned. A
934+
# load polls for up to `_LOAD_JOB_TIMEOUT_S`, so the longer the
935+
# wait the more chances to hit it, which is exactly backwards.
936+
#
937+
# CONSECUTIVE failures are the signal: an isolated 502 is noise, a
938+
# run of them means the API is gone and there is nothing to wait
939+
# for. The poll's own deadline bounds the total wait regardless.
940+
consecutive_errors += 1
941+
if consecutive_errors >= _JOB_POLL_MAX_CONSECUTIVE_ERRORS:
942+
raise RuntimeError(api_error_message(e)) from e
943+
time.sleep(interval_s)
944+
continue
919945
if last.status in _JOB_TERMINAL:
920946
return last
921947
time.sleep(interval_s)
@@ -929,9 +955,11 @@ def _load_response_from_job(self, job_id: str) -> LoadManagedTableResponse:
929955
930956
Polling replaces waiting on the request, so the outcome is read from
931957
durable state rather than from a connection that has to stay alive. That
932-
also removes the ambiguity `append` was made non-retryable for: a lost
933-
response no longer leaves "did it land?" unanswerable, because the job id
934-
outlives the request that created it.
958+
also gives a caller a handle: the job id is returned on
959+
`LoadManagedTableResult`, so "did it land?" is answerable after a lost
960+
response. `append` stays non-retryable -- knowing the id makes the question
961+
answerable, it does not make a blind re-submission safe, and that call is
962+
the caller's to make.
935963
936964
`partially_succeeded` is terminal and carries a message, so it is raised
937965
rather than returned -- a caller asked for a table's contents to be
@@ -943,7 +971,9 @@ def _load_response_from_job(self, job_id: str) -> LoadManagedTableResponse:
943971
raise RuntimeError(
944972
final.error_message or f"load job {job_id} finished {status}"
945973
)
946-
payload = final.result.actual_instance if final.result is not None else None
974+
# `result` is a oneOf wrapper today; tolerate the model arriving directly,
975+
# the same way the index path does rather than disagreeing with it.
976+
payload = getattr(final.result, "actual_instance", final.result)
947977
if not isinstance(payload, LoadManagedTableResponse):
948978
raise RuntimeError(
949979
f"load job {job_id} succeeded without a load result "

‎hotdata_framework/databases.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,11 @@ class LoadManagedTableResult:
8383
table_name: str
8484
row_count: int
8585
full_name: str
86+
# Set when the server ran the load as a background job. Carried for the same
87+
# reason CreateIndexResult carries it: it is the only handle a caller has to
88+
# ask "did that land?" after a lost response, and without it the question is
89+
# unanswerable. `None` when the load finished inline.
90+
job_id: str | None = None
8691

8792
def to_dict(self) -> dict[str, Any]:
8893
return asdict(self)

‎tests/test_client.py‎

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -714,3 +714,91 @@ def test_a_partially_succeeded_load_job_is_not_treated_as_success():
714714
else:
715715
raise AssertionError("partially_succeeded was treated as success")
716716

717+
718+
def test_a_transient_status_check_does_not_discard_a_running_job():
719+
"""The regression this guards. A load polls for up to an hour, so there are
720+
hundreds of status checks; aborting on the first bad one throws away a job
721+
that is still running, and the caller's retry then re-submits the load while
722+
the original still holds the table -- the 409 spin, from a new direction."""
723+
from hotdata.exceptions import ApiException
724+
725+
client = HotdataClient("k", "ws", host="https://api.hotdata.dev")
726+
calls = {"n": 0}
727+
final = SimpleNamespace(status="succeeded", error_message=None,
728+
result=SimpleNamespace(actual_instance=_load_response(5)))
729+
730+
class _Flaky:
731+
def get_job(self, job_id):
732+
calls["n"] += 1
733+
if calls["n"] < 3:
734+
raise ApiException(status=502, reason="Bad Gateway")
735+
return final
736+
737+
with (
738+
patch.object(client, "_jobs_api", return_value=_Flaky()),
739+
patch("time.sleep", lambda *_: None),
740+
):
741+
got = client._poll_job("jobs_1", timeout_s=60.0, interval_s=0.01)
742+
743+
assert got is final, "a blipping status check aborted the poll"
744+
assert calls["n"] == 3
745+
746+
747+
def test_a_run_of_failed_status_checks_still_gives_up():
748+
"""Tolerance is for blips, not for an API that has gone away."""
749+
from hotdata.exceptions import ApiException
750+
from hotdata_framework.client import _JOB_POLL_MAX_CONSECUTIVE_ERRORS
751+
752+
client = HotdataClient("k", "ws", host="https://api.hotdata.dev")
753+
754+
class _Dead:
755+
def get_job(self, job_id):
756+
raise ApiException(status=502, reason="Bad Gateway")
757+
758+
with (
759+
patch.object(client, "_jobs_api", return_value=_Dead()),
760+
patch("time.sleep", lambda *_: None),
761+
):
762+
try:
763+
client._poll_job("jobs_1", timeout_s=600.0, interval_s=0.01)
764+
except RuntimeError:
765+
pass
766+
else:
767+
raise AssertionError("polled forever against a dead API")
768+
assert _JOB_POLL_MAX_CONSECUTIVE_ERRORS >= 2
769+
770+
771+
def test_a_deferred_load_returns_the_job_id_to_the_caller():
772+
"""`append` stays non-retryable, so the id is the only handle a caller has to
773+
answer "did it land?" after a lost response -- the same reason
774+
CreateIndexResult carries one."""
775+
from hotdata.models.submit_job_response import SubmitJobResponse
776+
777+
client = HotdataClient("k", "ws", host="https://api.hotdata.dev")
778+
db = ManagedDatabase(id="db_1", description="mydb", default_connection_id="conn_1")
779+
connections = _FakeConnectionsApi(responses=[
780+
SubmitJobResponse(id="jobs_77", status="running", status_url="/v1/jobs/jobs_77"),
781+
])
782+
final = SimpleNamespace(status="succeeded", error_message=None,
783+
result=SimpleNamespace(actual_instance=_load_response()))
784+
785+
with (
786+
patch.object(client, "_databases_api", return_value=_ForbiddenDatabasesApi()),
787+
patch.object(client, "connections", return_value=connections),
788+
patch.object(client, "_poll_job", return_value=final),
789+
):
790+
result = client.load_managed_table(db, "orders", schema="public", upload_id="up_1")
791+
792+
assert result.job_id == "jobs_77"
793+
794+
795+
def test_an_inline_load_carries_no_job_id():
796+
client = HotdataClient("k", "ws", host="https://api.hotdata.dev")
797+
db = ManagedDatabase(id="db_1", description="mydb", default_connection_id="conn_1")
798+
with (
799+
patch.object(client, "_databases_api", return_value=_ForbiddenDatabasesApi()),
800+
patch.object(client, "connections", return_value=_FakeConnectionsApi()),
801+
):
802+
result = client.load_managed_table(db, "orders", schema="public", upload_id="up_1")
803+
assert result.job_id is None
804+

0 commit comments

Comments
 (0)