Skip to content

Re-raise and stop interrupted queries in the SQL cursors, sharing the Spark handling - #853

Open
laughingman7743 wants to merge 8 commits into
masterfrom
fix/840-sql-cursor-interrupt
Open

laughingman7743 wants to merge 8 commits into
masterfrom
fix/840-sql-cursor-interrupt

Conversation

@laughingman7743

@laughingman7743 laughingman7743 commented Sep 26, 2026 •

Copy link
Copy Markdown
Member

WHAT

With kill_on_interrupt enabled (the default), the SQL cursors now handle an interrupt the way the Spark cursors do since #833 and #861.
The handling lives in BaseCursor, where the synchronous Spark cursors share it, and AioBaseCursor has the asyncio counterparts.

Polling phase (_poll())

  • After an interrupt, the cursor calls _cancel_and_wait(): it requests cancellation (_cancel()) and waits for a terminal state (_poll_until_terminal()).
  • It then re-raises the original KeyboardInterrupt / asyncio.CancelledError, whatever the terminal state is. Before, the SQL cursors returned the final execution.
  • A failure to cancel or wait becomes the interrupt's __cause__.
  • The final execution is not stored for SQL, and no result set is built; query_id keeps the ID.

Start phase (_start_execution(), used by _execute())

  • Before, an interrupt during StartQueryExecution aborted the request (sync) or left it running in a worker thread (asyncio), so a query could start and keep running with no ID on the cursor.
  • With kill_on_interrupt, the request now runs on a short-lived helper thread (sync) or is shielded from task cancellation (asyncio).
  • On an interrupt, the cursor waits for the request, records the ID through _set_interrupted_execution_id(), calls _cancel_and_wait(), and re-raises.
  • A request that has not begun when the interrupt is handled is abandoned and never sent: the helper thread has not started it (sync), or the caller task was cancelled before the shielded start task ran (asyncio, checked with Task.cancelling()).

Sharing with Spark, building on the #881 helpers. No class hierarchy changes, and no new name-mangled methods (still 0).

  • SparkBaseCursor._poll() is removed; BaseCursor._poll() is the same code.
  • The helper-thread block in SparkBaseCursor._calculate(), _wait_for_calculation_start(), and _INTERRUPT_CHECK_INTERVAL move to BaseCursor._start_execution() / _wait_for_start() / pyathena.common._INTERRUPT_CHECK_INTERVAL.
  • Spark keeps only what differs: _poll_until_terminal(), _cancel(), _cancel_and_wait() (stores the calculation execution), and _set_interrupted_execution_id() (stores calculation_id).
  • WithResultSet overrides _set_interrupted_execution_id() to set query_id, for both WithFetch and WithAsyncFetch.
  • AioSparkCursor keeps its own asyncio _poll() / _calculate() handling from Cancel a Spark calculation interrupted while it is being started #861. It does not inherit AioBaseCursor, and the maintainer chose not to add module-level helpers for it. Its _calculate() gets the same abandon check as AioBaseCursor._start_execution(), so the check is written in both places.
  • The helper thread is now named pyathena-start (was pyathena-spark-start).
  • The "Query canceled by user." warning of the sync Spark cursors is now logged by the pyathena.common logger.

Other changes

  • SparkCursor.execute() / AioSparkCursor.execute() clear calculation_id and the previous calculation execution before starting. Before, a failed execution (for example, a failed cancel request after an interrupt) left the earlier calculation's state and outputs next to the new calculation_id.
  • executemany() stops at the interrupted execution and re-raises, with rowcount == -1 and query_id kept. Before, an interrupted execution that ended SUCCEEDED let the loop continue.
  • The thread-pool Async* cursors get the start-phase handling, because their execute() starts the query on the caller's thread. They have no query_id property, so the stopped query's ID is not available. They poll on worker threads, which never receive KeyboardInterrupt.
  • No ClientRequestToken is generated for SQL, unlike for Spark in Cancel a Spark calculation interrupted while it is being started #861. botocore auto-generates it for StartQueryExecution (idempotencyToken in the service model; not for StartCalculationExecution). PyAthena's own retries default to throttling errors, which do not start a query.
  • dbt-athena: its legacy PyAthena cursor calls _execute(), so it gets the start-phase handling, and overrides _poll(), so its polling is unchanged. Its current connection manager uses boto3 directly.
  • Docs: "Query cancellation on interrupt" in docs/usage.md and "Task cancellation" in docs/aio.md.

Behavior change (release note):

  • With the default kill_on_interrupt=True, an interrupt during a SQL cursor's execute() now propagates as KeyboardInterrupt / asyncio.CancelledError. Before, it surfaced as OperationalError (query ended CANCELLED/FAILED) or was lost (query ended SUCCEEDED). Callers that caught OperationalError after Ctrl-C or task cancellation need to handle the interrupt instead. asyncio.wait_for() now raises TimeoutError.
  • An interrupt while a query is being started now waits for StartQueryExecution and stops the query it started. With kill_on_interrupt, each StartQueryExecution of a synchronous cursor runs on a short-lived daemon thread.
  • After a failed execute() on SparkCursor / AioSparkCursor, state and the other calculation properties no longer describe the previous calculation.
  • AioSparkCursor: a task cancelled before its start task began no longer sends StartCalculationExecution. Before, the request was sent and then stopped.

WHY

Closes #840.
Swallowing the interrupt lost Ctrl-C when the query finished first. In asyncio, task.cancel() did not end with a cancelled task, and a timeout from asyncio.wait_for() surfaced as the query's OperationalError, or as a normal return, instead of TimeoutError.
After #861 (#841) fixed the start phase for Spark, the maintainer asked this PR to cover the SQL start phase and to share the handling with Spark.
This revision is rebased onto the #879 / #880 refactors (#881, #883, #909, #912, #914) and uses their single-underscore helpers.

TEST

Tested commit: c6357d0, rebased onto 9ef49d3 (#918 enables mypy explicit-override). After the rebase, SparkBaseCursor._cancel_and_wait and _set_interrupted_execution_id override BaseCursor and carry @override; without it mypy reports explicit-override (checked). WithResultSet._set_interrupted_execution_id stays unmarked like the mixin's other members, because the mixin has no base class.
Earlier @override changes (#917 rebase, a32a993): the asyncio overrides _start_query_execution, _start_execution, _cancel_and_wait. just lint (mypy with explicit-override) passes, and the offline checks below were rerun on this commit with the same results. Code and tests are otherwise as in 1ae69b4; the six later commits change docs only.
AWS CI on a32a993 (run 36896451089) passed: 2176 passed in the PyAthena suite (Spark included), and 589 passed in each SQLAlchemy suite. CI on this commit runs once Ready.

  • just format, just lint: passed (ruff, format check, mypy, cfn-lint, license headers). just docs lint: 0 errors.
  • Offline, no AWS: uv run --env-file .env pytest --noconftest -p no:xdist .... --noconftest skips the session setup that creates AWS resources.
    • SQL (tests/pyathena/test_cursor.py, tests/pyathena/aio/test_cursor.py, -k "interrupt or starting or before_request or on_poll_invoked or on_poll_none or legacy_kwargs_passthrough"): 27 passed.
    • Spark (tests/pyathena/spark, tests/pyathena/aio/spark): 114 passed and 2 skipped. 14 errors come only from tests that need the AWS fixtures (spark_cursor, async_spark_cursor, aio_spark_cursor).
    • The Cancel a Spark calculation interrupted while it is being started #861 Spark tests pass unchanged, except that the patch targets moved to pyathena.common.wait / pyathena.common.threading.Thread, the wait-interrupting helper moved to tests/pyathena/util.py, and the thread name changed.
    • A combined SQL and Spark selection on Python 3.11.11, the supported floor: 140 passed.
  • New SQL tests:
    • Polling phase: the interrupt or cancellation is re-raised only after RUNNING → terminal is polled (asserted through on_poll), for CANCELLED and SUCCEEDED. Cancel and wait failures become __cause__. asyncio.wait_for() raises TimeoutError. Without kill_on_interrupt, no stop is requested.
    • Start phase (sync): an interrupt injected while StartQueryExecution is blocked. One request is sent, the stop uses the returned ID, query_id is set, and no result set is built. Start and cancel failures become __cause__. Without kill_on_interrupt, no helper thread is created. AsyncCursor.execute() also stops the query.
    • Start phase (asyncio): task.cancel() while the start is blocked keeps the task pending until the request finishes, then stops the query, and task.cancelled() is true. wait_for() raises TimeoutError. Failures become __cause__. Without kill_on_interrupt, the cancellation propagates at once.
    • The asyncio start-phase timeout tests (SQL and Spark) release the request only after an asyncio.sleep() that ends after the timeout, so the order is fixed by the event loop's timer order.
    • test_execute_cancelled_before_request_is_sent (SQL and Spark asyncio) cancels the task right after execute() has scheduled the start task. No request is sent, no stop is requested, and the ID stays None. Both fail on 4480d3c. The asyncio window and timeout tests passed 30 of 30 repeated runs.
  • Spark reuse: the Define best-effort Spark calculation cancellation #833 poll-phase failure tests (sync and asyncio) start with a previous calculation on the cursor and assert that it was cleared.
  • Regression check, with pyathena/ restored to origin/master (49481d9):
    • 19 of the 22 SQL interrupt tests fail. The 3 that pass cover the unchanged kill_on_interrupt=False paths.
    • All 4 Spark reuse tests fail.
  • Real SIGINT (probe outside the repository, mocked client) while Cursor.execute() was blocked in StartQueryExecution, on 3.13.1 and 3.11.11: one start request, a stop with the returned ID, query_id set, and KeyboardInterrupt re-raised without a cause.
  • Live check from the earlier revision: StopQueryExecution on a SUCCEEDED query returns HTTP 200, and the state stays SUCCEEDED.
  • The AWS suites run in CI once the PR is Ready.

🤖 Generated with Claude Code

Comment thread pyathena/common.py
return query_execution
self.__poll(query_id)
except Exception as e:
raise interrupt from e

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 1: implementation behavior. Result: CLEAN

Scope: c01c56f73c7dbf973fe73b52b6093f0a32952321..2de91e768196dc061251f7609ae0d9b9be2905cd (full diff: pyathena/common.py, pyathena/aio/common.py, both test files, docs/usage.md, docs/aio.md).

Covered:

  • Callers of the shared _poll(): execute() of Cursor/DictCursor, pandas/arrow/polars/s3fs (sync), and AioCursor plus the aio pandas/arrow/polars/s3fs cursors. The thread-pool Async* cursors call _poll() from executor threads (pyathena/async_cursor.py:153, pyathena/pandas/async_cursor.py:131, etc.), which never receive KeyboardInterrupt, so they are unaffected. The Spark cursors override _poll(). SQLAlchemy does not call _poll() directly.
  • Cursor state after an interrupt: _reset_state() already cleared result_set and rowcount, query_id is set before polling, and no result set is built on the interrupt path, so nothing is left open.
  • executemany(): both pyathena/result_set.py:1064 and pyathena/aio/common.py:651 catch BaseException, close, and re-raise, so an interrupt stops the loop with rowcount == -1 and query_id kept. Before this change, an interrupted execution that ended SUCCEEDED let the loop continue with the next parameters.
  • Exception flow: the bare raise after the inner try/except Exception re-raises the outer interrupt. A second interrupt or cancellation during the wait is a BaseException, so it propagates instead of becoming a cause (same as Define best-effort Spark calculation cancellation #833).
  • Tests: the new tests fail on the original _poll() for the defect itself. The SUCCEEDED cases fail with DID NOT RAISE, because the fake result-set class lets the old code return normally.

No actionable findings. Round two will add the executemany() consequence to the PR description.

assert cursor.query_id == "query_id"
assert cursor.result_set is None

async def test_execute_kill_on_interrupt_timeout(self):

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 1: test determinism

The first status request blocks on an event that is never set, so the 0.01 s timeout always fires while polling.
The re-poll after the stop request returns at once, so asyncio.wait_for() sees CancelledError and raises asyncio.TimeoutError on both the 3.10 wait_for implementation and the timeout()-based one from 3.12.
asyncio.TimeoutError is used instead of the builtin TimeoutError, because the two are distinct on 3.10.

Comment thread docs/usage.md

With `kill_on_interrupt` enabled, which is the default, a `KeyboardInterrupt` while `execute()` waits for the query
requests cancellation, waits until the query reaches a terminal state, and then propagates.
Cancellation is a best-effort request, so the query can still end as `SUCCEEDED` or `FAILED`.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 2: claims, callers, and operations. Result: FINDINGS (PR description only, repaired)

Scope: c01c56f73c7dbf973fe73b52b6093f0a32952321..2de91e768196dc061251f7609ae0d9b9be2905cd, all claims in the PR description, commit message, _poll() docstrings, docs/usage.md, and docs/aio.md.

Claims checked:

  • Best-effort stop can still end SUCCEEDED: measured with one SELECT 1 on the CI account. StopQueryExecution on an already SUCCEEDED query returns HTTP 200 and the state stays SUCCEEDED. So the race re-raises the interrupt without a cause, and this sentence and the _poll() docstrings hold.
  • "If the cancellation request fails, ... as its cause": covered by the failure[cancel] tests on both bases.
  • kill_on_interrupt=False keeps the query running: no stop request is made (test_execute_without_kill_on_interrupt).
  • Timeouts: the commit message claim holds. On the original code, asyncio.wait_for() surfaces OperationalError (the new timeout test fails this way), and CPython's Timeout.__aexit__ converts only CancelledError into TimeoutError.
  • Existing callers: dbt-athena's current connection manager uses boto3 directly, and its legacy PyAthena cursor overrides _poll() (dbt-athena/src/dbt/adapters/athena/connections_legacy.py:167), so it is unaffected. The kill_on_interrupt parameter descriptions in pyathena/connection.py:228 and the cursor docstrings remain accurate.

Findings, repaired in the PR description:

  1. The WHY claimed that TaskGroup could not handle the cancellation. A TaskGroup still raises its ExceptionGroup when a sibling fails, so this is narrowed to the verified asyncio.wait_for() effect.
  2. The description omitted the executemany() consequence: it now stops at the interrupted execution, where an interrupted execution that ended SUCCEEDED used to let the loop continue. This is added, along with the dbt-athena compatibility note and the live stop measurement.

Deferred (pre-existing, out of scope): the thread-pool Async* cursor docstrings (e.g. pyathena/arrow/async_cursor.py:91) say kill_on_interrupt cancels on keyboard interrupt, but their polling runs in executor threads, which never receive KeyboardInterrupt. This PR does not change that path.

Comment thread pyathena/common.py
if not self._kill_on_interrupt:
raise
_logger.warning("Query canceled by user.")
try:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed): Codex CLI 0.157.0, model gpt-6-sol, reasoning effort high, session 01a0dcbd-22a4-7b63-b2a6-f2802f6e8c27. Result: FINDINGS

Scope: c01c56f73c7dbf973fe73b52b6093f0a32952321..2de91e768196dc061251f7609ae0d9b9be2905cd. The reviewer ran in a --sandbox read-only detached snapshot of the head with no .env. The prompt contained only the diff range, a file list, and the review questions; it had no PR number, description, commit message, or prior findings. Static review: the reviewer ran no tests. The snapshot and the PR worktree were unchanged afterwards.

Covered (reviewer's words): "synchronous and native asyncio execute() and executemany() paths through the default, dict, pandas, Arrow, Polars, and S3FS cursors; cursor state and exception chaining; the thread-backed async and Spark overrides; the new tests and both documentation examples."

[P2] Pre-existing, exposed by the new contract: "A second KeyboardInterrupt or task cancellation during the stop request or follow-up poll escapes the handler because both are BaseException subclasses, outside except Exception. With the query still running, execute() exits before observing a terminal state. ... The behavior predates the diff, while the new documentation states the wait without this qualification." (also pyathena/aio/common.py:155)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: documented; behavior kept (pre-existing). Verified: a second KeyboardInterrupt or task cancellation raised during _cancel() or the follow-up poll is a BaseException, so except Exception does not catch it and it propagates with the first interrupt as its context. This predates the diff, matches the Spark contract from #833, and gives users a way to stop waiting on a query that does not stop. da16975 documents it: docs/usage.md says "A second KeyboardInterrupt during that wait propagates without waiting for the terminal state.", and docs/aio.md has the equivalent sentence for a repeated task cancellation.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI, model gpt-6-sol, reasoning effort high, session 01a0dcc3-f3ef-7833-a415-86e914ba8fa2. Scope: 2de91e768196dc061251f7609ae0d9b9be2905cd..da169757662dda486e41349789ce7a77e67ddc91 (same base c01c56f73c7dbf973fe73b52b6093f0a32952321), read-only snapshot, static review.

Covered surfaces: the four files changed in 2de91e7..da16975, traced through pyathena/aio/cursor.py, pyathena/aio/common.py, pyathena/aio/util.py, pyathena/cursor.py, and pyathena/common.py. This was a static, read-only review; I did not run tests.

FINDINGS

  • P2 — docs/aio.md:142, docs/usage.md:511: The new qualification covers a second interruption while polling for the terminal state, but says nothing about one during the stop request. Both handlers call _cancel before that follow-up poll, and their except Exception blocks do not catch a second CancelledError or KeyboardInterrupt. For example, a second task cancellation while the stop request is awaiting its worker thread exits the handler without confirming that the stop request ran or that the query became terminal. The documentation still implies the first interruption completes those steps unless interrupted “during that wait.”

Prior items

  1. Resolved. docs/aio.md:145 limits the timeout claim to polling and explains that cancellation during startup can leave query_id unset while the worker-thread request continues.
  2. Not an actual defect. In this test, _execute is an immediately completing AsyncMock; the first status request then enters an Event.wait() and suspends. On Python 3.10, wait_for() schedules the child task before its timeout can cancel it; on 3.11–3.14, it awaits the coroutine inside the timeout context. The event loop therefore reaches that blocked status request before delivering the 10 ms timeout. A separate barrier is unnecessary for this fixture.
  3. Resolved. Both tests now return RUNNING before a terminal state and assert that the polling callback saw both states before the original interruption propagated.
  4. Unresolved for the stop-request window; resolved for follow-up polling. The source behavior described in the finding remains possible, and the added wording qualifies only “that wait.”

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in f528d50. The repeated-interrupt note now covers the cancellation request as well as the follow-up wait, in both docs/usage.md and docs/aio.md, and says the query can keep running.

Independent follow-up review (relayed): Codex CLI, model gpt-6-sol, reasoning effort high, session 01a0dcc8-cb13-7712-8373-34cdcfc3f7e1. Scope: da169757662dda486e41349789ce7a77e67ddc91..f528d50bc242ccacb06631081a7f21bda31d70a2, read-only snapshot, static review.

Covered surfaces: docs/aio.md “Task cancellation” and docs/usage.md “Query cancellation on interrupt,” checked against the five named source files. The new wording covers a second interruption during both the cancellation request and the follow-up poll.

CLEAN. No actionable inaccuracies found in either section. This was a read-only source review; no tests were run.

Comment thread docs/aio.md Outdated
If the cancellation request fails, `asyncio.CancelledError` is raised with the error as its cause.
With `kill_on_interrupt=False`, `asyncio.CancelledError` is raised immediately and the query keeps running.

A timeout from `asyncio.wait_for()` therefore cancels the query and raises `asyncio.TimeoutError`:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed, Codex gpt-6-sol): [P2] introduced

"The wait_for() example says a timeout cancels the query. If the timeout occurs while start_query_execution is running in a worker thread, execute() has not assigned query_id; cancellation cannot stop the request, and Athena may start a query after the task exits. The example can print None while that query keeps running."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in da16975. Verified: aio StartQueryExecution runs through asyncio.to_thread() (pyathena/aio/util.py:42), and query_id is assigned only when _execute() returns. The docs now limit the cancellation to a timeout that expires while execute() waits for the query, and state that a timeout during the start leaves query_id as None while the pending start request can still start the query. execute() has no await between the query_id assignment and _poll(), so there is no third window. The example prints query_id as a value that may be None.

kill_on_interrupt=True, final_state=AthenaQueryExecution.STATE_CANCELLED
)
with pytest.raises(asyncio.TimeoutError):
await asyncio.wait_for(cursor.execute("SELECT 1"), timeout=0.01)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed, Codex gpt-6-sol): [P2] introduced

"The 10 ms timeout test has no barrier confirming that the task reached polling. Under a scheduling delay, it can time out before a query ID is assigned; cancel is then never awaited and the assertion fails intermittently."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: rejected with evidence. No await in execute() yields to the event loop before the first status request blocks. _execute is an AsyncMock, which completes without suspending, and _call_on_start_query_execution is synchronous, so query_id is always assigned before the loop can run the timeout callback. On 3.12+, wait_for() awaits the coroutine inside the calling task under timeouts.timeout(), so the timer can fire only at the first suspension, the blocked poll. On 3.10/3.11, wait_for() wraps the coroutine in a task whose first step call_soon places in _ready before _run_once moves any expired timer into _ready, so that step runs first even after a scheduling delay. The test passed on 3.13.1 and 3.10.16 locally.

def test_execute_kill_on_interrupt(self, final_state):
"""An interrupt cancels the query, waits for it, and is re-raised (no AWS)."""
cursor, cancel = _offline_cursor(kill_on_interrupt=True)
cursor._get_query_execution = MagicMock(

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed, Codex gpt-6-sol): [P3] introduced

"The follow-up status is immediately terminal in the new tests. Given a query that stays RUNNING after the stop request, an implementation that polls once and raises before terminal state could still satisfy these tests. They exercise the prior implementation's failing path, but do not assert the promised wait through a nonterminal state." (also tests/pyathena/aio/test_cursor.py:197)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in da16975. Both test_execute_kill_on_interrupt tests now report RUNNING after the stop request, then the terminal state. They assert through the on_poll hook that the polled states were [RUNNING, final_state] before the interrupt or cancellation was re-raised. Mutation check: replacing the follow-up __poll() with a single _get_query_execution() call in both _poll() implementations makes all 4 of these tests fail. With the real implementation, 13 targeted tests pass on 3.13.1, and 11 interrupt tests pass on 3.10.16.

@laughingman7743
laughingman7743 marked this pull request as ready for review September 26, 2026 08:34
@laughingman7743
laughingman7743 force-pushed the fix/840-sql-cursor-interrupt branch from f528d50 to ef17dec Compare September 28, 2026 06:26
@laughingman7743
laughingman7743 marked this pull request as draft September 28, 2026 06:47
@laughingman7743 laughingman7743 changed the title Re-raise the interrupt after kill_on_interrupt cancellation Re-raise and stop interrupted queries in the SQL cursors, sharing the Spark handling Sep 28, 2026
Comment thread pyathena/common.py Outdated
_INTERRUPT_CHECK_INTERVAL = 0.1


def _start_interruptibly(

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 1 (expanded scope: start phase + shared helpers): implementation behavior. Result: FINDINGS (1, repaired)

Scope: full diff 659676c07e2c09397b7cbc5740cc10bfe6fe41cb..d66c6994034894313f6b68f23f9f5d421ac13632 (rebased onto #861). Because the contract expanded, this is a full pass, not a range-diff.

Covered:

  • _execute() callers: only the cursors' execute(). That is Cursor, the pandas/arrow/polars/s3fs sync cursors, their Async* thread-pool variants (start on the caller's thread), and the five Aio* cursors. SQLAlchemy calls execute() only. dbt-athena legacy calls _execute() and gets the start-phase handling.
  • Shared helpers: the sync helper keeps Cancel a Spark calculation interrupted while it is being started #861's abandon check (Future.set_running_or_notify_cancel() / cancel()) and its interruptible wait. Name mangling inside the lambda resolves per class (_BaseCursor__start_query_execution, _AioBaseCursor__..., _SparkBaseCursor__...). The bare raise after the inner except Exception re-raises the outer interrupt.
  • Spark: _poll/_calculate now pass their own poll/stop callables. __stop_started_calculation sets calculation_id before the cancel, as before, so a cancel failure still leaves the ID on the cursor. All Cancel a Spark calculation interrupted while it is being started #861 tests pass with only their patch targets moved.
  • SQL state: _set_interrupted_query_id() is a no-op on BaseCursor (AsyncCursor has no query_id) and sets query_id in WithFetch/WithAsyncFetch. It is called before the cancel. On a start failure, query_id stays None.
  • Poll-phase aio tests: they keep a mocked _execute, so the earlier determinism argument for the 10 ms wait_for test still holds. The start-phase tests use a separate helper with the real _execute().

Finding (repaired in d66c699): docs/aio.md said query_id is None only if the timeout expires before the query is started, but a start request that fails after the timeout also leaves it None. Reworded to "only if no query was started".

Out of scope, same as the Spark cursors since #861: in asyncio, a start task that has not begun when the cancellation arrives still sends its request, which is then stopped. The sync helper abandons it.

Comment thread docs/usage.md
The `query_id` property returns that query's ID.
If the request has not been sent yet when the interrupt is handled, it is never sent.
`AsyncCursor` and its variants also stop a query whose start is interrupted in `execute()`.
They wait for queries on worker threads, which do not receive `KeyboardInterrupt`.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 2 (expanded scope): claims, callers, and operations. Result: CLEAN

Scope: 659676c07e2c09397b7cbc5740cc10bfe6fe41cb..d66c6994034894313f6b68f23f9f5d421ac13632, the rewritten PR description, docstrings, and docs/usage.md / docs/aio.md.

Claims checked:

  • No SQL token generation needed: botocore's Athena model marks StartQueryExecution.ClientRequestToken as idempotencyToken (auto-generated per call) and StartCalculationExecution's not. RetryConfig defaults to THROTTLING_ERROR_CODES, which do not start a query.
  • "AsyncCursor and its variants also stop a query whose start is interrupted": pyathena/async_cursor.py:225 and the pandas/arrow/polars/s3fs async cursors call self._execute() on the caller's thread. Polling runs in self._executor.
  • "If the request has not been sent yet ..., it is never sent": this is the sync abandon path, covered by the Cancel a Spark calculation interrupted while it is being started #861 tests test_calculate_interrupted_before_request_is_sent[False/True], which now exercise the shared helper.
  • Real signal: a SIGINT via os.kill while Cursor.execute() was blocked in a mocked StartQueryExecution produced one start request, a stop with the returned ID, query_id set, and KeyboardInterrupt without a cause. Checked on 3.13.1 and 3.10.16.
  • Operational: one short-lived daemon thread per StartQueryExecution with kill_on_interrupt, and no extra AWS calls on the normal path. The start still goes through retry_api_call with the same config.
  • Evidence scope: local results are offline only (3.13.1 and 3.10.16). The AWS suites were not run locally for this revision. An accidental local run of all of tests/pyathena with --noconftest sent some queries from tests that call connect() directly; its missing-schema failures are not used as evidence, and the PR description says so.

No corrections were needed beyond the round 1 repair.

Comment thread tests/pyathena/aio/test_cursor.py Outdated
async def test_execute_timeout_while_starting(self):
"""A timeout during the start request stops the query and raises TimeoutError (no AWS)."""
cursor, cancel, started, release = _starting_cursor()
timer = threading.Timer(0.2, release.set)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed): Codex CLI 0.157.0, reported model gpt-6-astra, reasoning effort high, session 01a0e6c8-41ad-7423-a641-afcd8c29f7fc. Result: FINDINGS

Scope: 659676c07e2c09397b7cbc5740cc10bfe6fe41cb..d66c6994034894313f6b68f23f9f5d421ac13632 (full expanded scope). The reviewer ran in a --sandbox read-only detached snapshot with no .env. The prompt contained no PR number, description, commit message, or prior findings. Static review. The snapshot and the PR worktree were unchanged afterwards.

Covered (reviewer's words): "the full diff and the sync, thread-pool async, and native asyncio cursor families, including pandas/Arrow/Polars/S3FS and Spark. Traced start/poll/cancel handling, exception identity/chaining, repeated interruption, helper lifecycle, cursor state, legacy methods and overrides, Python ≥3.10 compatibility, tests, and documentation." It found "No additional defect ... in the shared-helper refactor or legacy method signatures."

1. [P2] introduced: "The release timer starts before wait_for() establishes its timeout. If the test thread is descheduled for over 200 ms after timer.start(), the request is released before execution begins. The mocked query then reaches CANCELLED and raises OperationalError, failing the expected TimeoutError assertion despite correct implementation."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in 10f7a83, together with the same pattern in the Spark test from #861 (tests/pyathena/aio/spark/test_cursor.py). Both tests now start wait_for() as a task, wait for the start request to begin, and then await asyncio.sleep(0.1) before releasing the request. The sleep begins after wait_for() scheduled its 0.05 s deadline, so its timer is due later. The event loop therefore runs the timeout callback first, and the cancellation reaches the task before the test resumes and releases the request. threading.Timer is gone. Both tests passed 30 of 30 repeated runs.

Comment thread docs/aio.md Outdated
With `kill_on_interrupt=False`, `asyncio.CancelledError` is raised immediately and the query keeps running.

A timeout from `asyncio.wait_for()` therefore cancels the query and raises `asyncio.TimeoutError`.
`query_id` is `None` only if no query was started, for example when the timeout expires while looking up a cached result.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed, Codex gpt-6-astra): 2. [P3] introduced

"With kill_on_interrupt=False, a timeout during StartQueryExecution cancels the await while its underlying thread continues. Athena can start the query, but the cursor retains query_id=None. Repeated cancellation during the protected start wait can also abandon ID recovery. Document None as an unavailable ID, rather than evidence that no query exists."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in 10f7a83. Verified: with kill_on_interrupt=False, or after a second cancellation during the shielded start, a query can start while query_id stays None. The sentence is now one-directional: "query_id is None if the timeout expires before the start request is sent, for example while looking up a cached result." The preceding paragraph already says that with kill_on_interrupt=False the query keeps running, and that another cancellation during the waits leaves it running.

Comment thread pyathena/spark/common.py Outdated
)

future: Future[str] = Future()
def __stop_started_calculation(self, calculation_id: str) -> None:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed, Codex gpt-6-astra): 3. [P2] pre-existing, preserved by the Spark refactor

Anchored here because the lines are outside the diff: pyathena/spark/cursor.py:147, pyathena/aio/spark/cursor.py:351.

"Execute calculation A successfully, then start B on the same cursor. Interrupt/cancel B while polling and make _cancel() fail. The exception propagates, but calculation_id identifies B while calculation_execution, state, and output accessors still describe A. Neither execution path clears the previous calculation object."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in 10f7a83 (pre-existing; folded in as a contained Spark consistency fix). Verified: SparkCursor.execute() and AioSparkCursor.execute() overwrote _calculation_id without clearing _calculation_execution. Both now set them to None before _calculate(), like the SQL cursors' _reset_state(). The #833 poll-phase failure tests (sync and asyncio) now start with a previous calculation on the cursor and assert that calculation_id is the new one and calculation_execution is None. Without the reset, all 4 fail. Side effect, checked: cancel() during a new start now raises ProgrammingError instead of stopping the previous, finished calculation. The live test_cancel tests wait for an ID that is neither None nor the previous one, so they are unaffected. This is recorded as a release-note item in the PR description.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI, reported model gpt-6-astra, reasoning effort high, session 01a0e710-57c6-7d73-b6b3-e9bb82e625e4. Scope: d66c6994034894313f6b68f23f9f5d421ac13632..10f7a832ebbac10440706167d1e370f90b23c2ab (same base 659676c07e2c09397b7cbc5740cc10bfe6fe41cb), read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.

Covered surfaces: all six changed files; Spark start/poll/cancellation paths; calculation_id, calculation_execution, state, and output accessors; repository callers; documentation; and local CPython asyncio sources for 3.10–3.14.

CLEAN — no actionable defects found in d66c699..10f7a832.

  1. Timeout-test race: resolved. Both rewritten tests wait for the start request to begin, then release it only after the event-loop sleep.

    • Python 3.10–3.11: wait_for() registers its timeout before execution starts. Its deadline precedes the subsequently registered sleep deadline. Even if both timers become overdue together, the timeout callback queues the wait_for() continuation first; that continuation requests cancellation before the test resumes and releases the worker.
    • Python 3.12–3.14: wait_for() uses the timeout context manager, whose earlier timer directly cancels the executing task before the sleep continuation releases the worker.

    The pending-task assertion, subsequent TimeoutError, cancellation-call assertion, and retained-ID assertion collectively exercise the intended behavior. The existing 10-second worker watchdog remains a finite scheduling limit; the original independent-thread release race is removed.

  2. query_id documentation: resolved. docs/aio.md:148 removes the exclusivity claim. It gives the pre-start/cache-lookup case without asserting that None proves no query started, allowing the documented disabled-cleanup and repeated-cancellation cases.

  3. Stale Spark calculation: resolved. pyathena/spark/cursor.py:148 and pyathena/aio/spark/cursor.py:352 clear both fields before starting another calculation. Start failures leave no previous ID/result; failures after obtaining the new ID retain that ID without the previous execution. Successful interruption cleanup still stores the current terminal execution.

    Compatibility is consistent: state and execution-derived properties already support None; cancel() rejects an unknown ID and targets the current calculation once known. The strengthened tests seed previous state and verify its removal for both cancellation-request and cleanup-wait failures. They directly protect the stale-execution regression, though they do not independently test clearing a previous ID when startup fails.

Static review only: no tests, builds, writes, or network access. HEAD remained 10f7a832; the worktree remained clean.


Author note: the remark that no test separately covers clearing a previous ID when the start fails is not taken up. The reset runs unconditionally before _calculate() (pyathena/spark/cursor.py:148, pyathena/aio/spark/cursor.py:352), and the new assertions already fail without it.

Comment thread pyathena/common.py
self._cancel(query_id)
self._poll_until_terminal(query_id)

def _start_execution(self, start: Callable[[], str]) -> str:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 1 (rebased onto the #879/#880 refactors): implementation behavior. Result: CLEAN

Scope: full diff 49481d90..4480d3cc7d885ccfa1f1eff3b2e3366271482bf4 (base = origin/master tip after #914). This is a re-implementation on the #881 single-underscore helpers rather than a replay of the old commits, so it is a full pass.

Covered:

  • Sync SQL (Cursor, DictCursor, arrow/pandas/polars/s3fs) and the thread-pool Async* cursors: _execute() → _start_execution(), and _poll() → _poll_until_terminal() / _cancel_and_wait(). WithResultSet._set_interrupted_execution_id() comes before BaseCursor in the WithFetch / WithAsyncFetch MRO, so it sets query_id. AsyncCursor keeps the BaseCursor no-op.
  • Sync Spark (SparkCursor, AsyncSparkCursor): the removed SparkBaseCursor._poll was line-for-line the new BaseCursor._poll. _calculate() → _start_execution() keeps Cancel a Spark calculation interrupted while it is being started #861's abandon check and interruptible wait, now BaseCursor._wait_for_start. The overrides _cancel_and_wait (stores the execution) and _set_interrupted_execution_id (stores calculation_id) keep the old state updates. AsyncSparkCursor polls in executor threads as before.
  • asyncio: AioBaseCursor._execute() → _start_execution() (shield), and _poll() / _cancel_and_wait(). AioSparkCursor is untouched apart from the execute() reset (maintainer's choice). It inherits the sync _set_interrupted_execution_id, but its own _calculate sets _calculation_id directly as before.
  • Structure: no class hierarchy or mixin-role change. Name-mangled def __x count: 0 on master and 0 on the branch (measured).
  • Tests: on master's pyathena/, 19 of the 22 new SQL interrupt tests fail; the 3 that pass are kill_on_interrupt=False paths. All 4 Spark reuse tests fail. All tests pass on 3.13.1 and 3.11.11, the floor.

Comment thread pyathena/spark/common.py
return self._start_execution(lambda: self._start_calculation_execution(request))

future: Future[str] = Future()
def _set_interrupted_execution_id(self, execution_id: str) -> None:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round 2: claims, callers, and operations. Result: FINDINGS (PR description only, repaired)

Scope: same range, plus the rewritten PR description, docstrings, and docs/usage.md / docs/aio.md.

Checked:

Finding (repaired in the description): the release-note bullet said every synchronous StartQueryExecution runs on a daemon thread. That holds only with kill_on_interrupt, and the bullet is now qualified.

Limits: offline evidence only. The AWS suites run once Ready.

Comment thread pyathena/aio/common.py Outdated
if not self._kill_on_interrupt:
return await start

task = asyncio.ensure_future(start)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed): Codex CLI, reported model gpt-6-sol, reasoning effort high, session 01a0f83a-cae6-7733-a9aa-9c2b45cc1f39. Result: FINDINGS

Scope: 49481d900e9e0098e3edd562f1a858ce91ec0573..4480d3cc7d885ccfa1f1eff3b2e3366271482bf4. The reviewer ran in a --sandbox read-only detached snapshot with no .env. The prompt contained no PR number, description, commit message, or prior findings. Static review. The snapshot and the PR worktree were unchanged afterwards.

Covered (reviewer's words): "the start, poll, cancellation, and state paths through the sync SQL, thread-pool, native asyncio, and Spark cursors, including their specialized variants, tests, and documentation." It found "no additional regression in Spark's moved sync polling handler or in method resolution for the inspected cursor classes."

1. [P2] introduced: "Cancellation can start a query after the caller has cancelled. If the outer task is cancelled after ensure_future(start) schedules the child but before that child runs, the handler awaits the child, which then sends StartQueryExecution. The timeout is delayed while a new query starts and is cancelled best effort. This also contradicts the query_id is None claim in docs/aio.md:148."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in 1ae69b4 (code + tests) and 1bcd22c (docs). Verified: the window is real. test_execute_cancelled_before_request_is_sent fails on 4480d3c with the request sent.

A first attempt, a started flag checked in the cancellation handler, does not close the window. task.cancel() only schedules the caller's wakeup, so the start task can take its first step before the handler runs. It was reverted.

The fix: _start_execution() records asyncio.current_task().cancelling() on entry. The start task returns None without calling start() if the caller's count has grown since. The handler then re-raises the original cancellation without a stop request. start is now a callable that returns the coroutine, so an abandoned request leaves no un-awaited coroutine. Task.cancelling() exists from Python 3.11, the supported floor.

The test cancels right after await asyncio.sleep(0) lets execute() schedule the start task: no request, no stop, query_id is None. It passed 30 of 30 repeated runs. The docs now say the request is never sent in that case, which makes the query_id is None sentence in docs/aio.md hold. (Maintainer chose to fix this in this PR.)

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI, reported model gpt-6-sol, reasoning effort high, session 01a0f845-314c-72c2-85c3-6bfb34edf6f9. Scope: 4480d3cc7d885ccfa1f1eff3b2e3366271482bf4..fdb6c2e0938d30d8a3dde843b0b85a6feb3a815c, read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.

Covered surfaces: the new commits in both asyncio start paths, their cursor state and cancellation handling, the new tests, and the three changed documentation pages. This was a read-only source review; no tests were run.

FINDINGS

  • P3 — docs/aio.md:149: query_id is not always None when a timeout expires before the AWS request is sent. If the thread pool is busy, the start task can pass its cancellation guard and queue the call with asyncio.to_thread(). A timeout can then occur before a worker sends it. The cursor waits for that start, records its ID, and requests cancellation. The documentation should distinguish “the start task has not begun” from “the AWS call has not been sent.”

Previous findings: Item 1’s scheduling window is fixed, but its documentation claim remains too broad as described above, so it is only partially resolved. Item 2 is resolved: docs/usage.md now says that AsyncCursor has no query_id property. Item 3 is resolved: Spark uses the same pre-start cancellation guard.

The new guard creates the start coroutine only after checking Task.cancelling(), so the reviewed pre-start path does not leave an un-awaited start coroutine. The inspected wait_for() and repeated-cancellation paths preserve the existing cancellation propagation and exception-cause behavior.


Disposition: fixed in 96ad462. Verified: once the start task has passed the guard, a busy thread pool can delay the actual AWS call past a timeout, and the request is then still sent and stopped. docs/aio.md now reads: "query_id is None if the timeout expires before execute() begins the start request, for example while looking up a cached result." This matches the "never sent" sentences, which also refer to execute() beginning the request. No other doc sentence ties behavior to the request being sent.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI, reported model gpt-6-sol, reasoning effort high, session 01a0f848-3107-7460-8411-5bbdc483975b. Scope: fdb6c2e0938d30d8a3dde843b0b85a6feb3a815c..96ad4624b46542e44742add80c405d2777ffcc0b, read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.

Covered surfaces: the one-sentence diff, every sentence in “Task cancellation” in docs/aio.md and “Query cancellation on interrupt” in docs/usage.md, and the start-phase sentences in docs/spark.md, traced through the requested cursor implementations and the direct AioCursor.execute() caller.

FINDINGS

The new sentence resolves the previously reported thread-pool timing case: once the async start task passes its cancellation guard, cancellation waits for the start result and retains the query ID.

  • Low — docs/aio.md:149: The revised condition is still too broad when the cache supplies a previous query ID. If cache lookup finds a reusable result and wait_for() times out during subsequent polling, execute() never begins a start request, yet query_id contains the reused ID. The cache-hit path is in AioBaseCursor._execute(); AioCursor.execute() assigns that ID before polling.

No other actionable documentation mismatch was found in the requested sections. This was a read-only source review; no tests or network checks were run.


Disposition: fixed in 07c22e3. Verified: on a cache hit, AioBaseCursor._execute() returns the reused ID without a start request, and AioCursor.execute() assigns it before polling. The sentence now names the only window that leaves the ID unset: "If the timeout expires while execute() is still looking up a cached result, no query is started and query_id is None." _reset_state() clears query_id at the start of execute(), and _execute() has no other await before the start phase.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI, reported model gpt-6-sol, reasoning effort high, session 01a0f84a-76a9-70d2-a2cf-4f002804208b. Scope: 96ad4624b46542e44742add80c405d2777ffcc0b..07c22e37bdb0dbf61eb92511ecf2e9c57e3d4a68, read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.

Covered surfaces: the one-sentence diff; “Task cancellation” in docs/aio.md; “Query cancellation on interrupt” in docs/usage.md; and the start-phase sentences in docs/spark.md, traced through the specified source files and the cursor entry points needed to check query_id.

FINDINGS. The new sentence at docs/aio.md:149 resolves the previous cache-hit finding: during lookup, query_id remains None; after a cache hit, polling uses the reused ID.

  • P2 — docs/aio.md:148: The sentence says a wait_for() timeout cancels the query. If the timeout occurs during cache lookup, no query has started, as the next sentence correctly explains. With kill_on_interrupt=False, a timeout during polling also leaves the query running.
  • P3 — docs/aio.md:146, docs/usage.md:522: Both say the query keeps running with kill_on_interrupt=False. If cancellation or interruption occurs during cache lookup, execution ends before a start request and there is no running query.

I found no further actionable mismatch in the reviewed Spark start-phase sentences. This was source inspection only; no tests were run.


Disposition: both fixed in 486a0e0.

  • docs/aio.md: "With kill_on_interrupt enabled, a timeout from asyncio.wait_for() therefore cancels a query that has started and raises asyncio.TimeoutError."
  • docs/aio.md / docs/usage.md: "With kill_on_interrupt=False, ... propagates/is raised immediately, and a query that has already started keeps running."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI, reported model gpt-6-sol, reasoning effort high, session 01a0f84c-eefe-7873-8d02-2d4ef850e963. Scope: 07c22e37bdb0dbf61eb92511ecf2e9c57e3d4a68..486a0e0f81bc3d8a4dfb2a28db9dfaca826ca6ae, read-only snapshot, static review. The prompt asked for every sentence of the three sections in one pass. The snapshot and the PR worktree were unchanged afterwards.

Covered surfaces: every sentence in the Task cancellation subsection of docs/aio.md, the Query cancellation on interrupt section of docs/usage.md, and the cancellation paragraphs of docs/spark.md, checked against the listed source files. This was a read-only source review; no tests were run.

Prior findings: Both reported counterexamples are addressed. The wait_for() claim now specifies kill_on_interrupt and accounts for a timeout during cache lookup; both kill_on_interrupt=False claims now refer only to a query that has started.

FINDINGS

  • Medium — docs/aio.md:148: The timeout claim still covers more of execute() than the cancellation handler does. After polling observes SUCCEEDED, AioCursor.execute() awaits result-set creation (pyathena/aio/cursor.py:143). If wait_for() expires during that await, it raises TimeoutError, but execution never enters the query cancellation path in _poll() (pyathena/aio/common.py:243). The query has started, yet the timeout does not cancel it.

Disposition: fixed in 69cacb3. Verified: after SUCCEEDED, AioCursor.execute() awaits result-set creation outside _poll(), so a timeout there stops nothing; the query has already finished. The sentence now covers only the phases that are handled: "With kill_on_interrupt enabled, a timeout from asyncio.wait_for() while execute() starts or waits for the query therefore cancels it and raises asyncio.TimeoutError." The other sentences were already limited to "while execute() waits for / is still starting the query".

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up review (relayed): Codex CLI, reported model gpt-6-sol, reasoning effort high, session 01a0f84f-1b28-7ec2-84ec-f4bc242c650c. Scope: 486a0e0f81bc3d8a4dfb2a28db9dfaca826ca6ae..69cacb31eb5a1311ded732d023cfc29af9fdec88, plus every sentence of the three cancellation sections. Read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.

Covered surfaces: Every sentence in “Task cancellation” in docs/aio.md, “Query cancellation on interrupt” in docs/usage.md, and the cancellation paragraphs in docs/spark.md, checked against the nine requested source files.

The previous finding is resolved: docs/aio.md:148 now limits the timeout claim to query start and polling, excluding result-set creation after SUCCEEDED.

FINDINGS

  • P2 — docs/spark.md:259: The bullet promises that a running Spark job stops within seconds and ends CANCELED. In a reachable case, a stop request sent just after the calculation starts has no effect and the calculation ends COMPLETED, as the same section states. The cursor sends StopCalculationExecution but does not enforce the promised outcome.

Disposition: the previous item is resolved, and no finding remains in this PR's diff. The new item (docs/spark.md:259) is pre-existing text from #833 (71a8168f) that this PR does not change. It sits under "Athena cancels the calculation on a best-effort basis", and the same section already says that "A cancellation request sent right after a calculation starts can occasionally have no effect". It is left out of scope and can be tightened separately if wanted.

Comment thread pyathena/common.py
wait((future,), timeout=_INTERRUPT_CHECK_INTERVAL)
return future.result()

def _set_interrupted_execution_id(self, execution_id: str) -> None: # noqa: B027

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed, Codex gpt-6-sol): 2. [P2] introduced

"An interrupted AsyncCursor start obtains an execution ID, but its inherited _set_interrupted_execution_id() does nothing. If cancellation then fails, execute() raises before returning its (query_id, future) tuple, leaving the caller without the ID needed to retry cancellation or inspect the query. The unqualified ID-retention claim in docs/usage.md:509 does not hold for AsyncCursor and its variants."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in the docs (1bcd22c); no code change. Verified: AsyncCursor and its variants have no query_id property, and execute() raises before returning (query_id, future). Adding an attribute would change the class design, which is out of scope here. docs/usage.md now says: "AsyncCursor and its variants also stop a query whose start is interrupted in execute(), but they have no query_id property, so the ID of that query is not available." The "query_id property keeps/returns" sentences describe the cursors that have the property. The PR description says the same.

Returns:
Self reference for method chaining.
"""
# A failure below must not leave the previous calculation on the cursor.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed, Codex gpt-6-sol): 3. [P2] pre-existing

Anchored here because the line is outside the diff: pyathena/aio/spark/cursor.py:136.

"Native asyncio Spark already has the same shielded-start window: cancellation before its scheduled start task runs can still initiate a calculation. This path was unchanged by the reviewed commit."

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Disposition: fixed in 1ae69b4 (pre-existing since #861; the maintainer chose to fix it in this PR). AioSparkCursor._calculate() gets the same Task.cancelling() check as AioBaseCursor._start_execution(). The class does not inherit AioBaseCursor, and the maintainer chose not to add module-level helpers, so the check is written in both places. New test tests/pyathena/aio/spark/test_cursor.py::test_execute_cancelled_before_request_is_sent: no StartCalculationExecution, no stop, calculation_id is None. It fails on 4480d3c. All #861 asyncio Spark tests still pass. docs/spark.md gets the matching sentence.

Comment thread pyathena/aio/common.py
return await start()

caller = asyncio.current_task()
cancel_requests = caller.cancelling() if caller else 0

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review of the repairs (rounds 1 and 2): range 4480d3cc7d885ccfa1f1eff3b2e3366271482bf4..fdb6c2e0938d30d8a3dde843b0b85a6feb3a815c. Result: FINDINGS (docs wording, repaired in fdb6c2e)

Round 1, behavior:

  • The cancellation count is taken on entry, so a cancellation that was already pending, or swallowed without uncancel(), does not abandon a later start.
  • asyncio.wait_for() on 3.11 runs the coroutine in its own task. The timeout cancels that task, which is the caller here. The start-phase timeout tests pass on 3.11.11.
  • A second cancellation while the handler awaits a start task that has not run cancels that task before its first step, so start() is never called and no coroutine is left un-awaited.
  • If the start task has begun, the path is unchanged: wait, record, stop, re-raise.
  • AioSparkCursor has the identical change.
  • Checks: SQL 27 passed, Spark 114 passed, 3.11.11 140 passed. The window, timeout, and cancelled-while-starting tests passed 30 of 30 repeated runs. Both window tests fail on 4480d3c.

Round 2, claims:

  • The new doc sentence said a request "not sent yet when the cancellation is handled" is never sent. The code abandons it only if the start task, or the sync helper thread, has not begun the request. A begun request whose HTTP call has not left yet is still sent.
  • The sync sentence in docs/usage.md had the same imprecision.
  • All three now say the request is never sent if the interrupt or cancellation comes before execute() begins it.

@laughingman7743
laughingman7743 force-pushed the fix/840-sql-cursor-interrupt branch from 69cacb3 to 2f4dcc4 Compare October 1, 2026 16:42
@laughingman7743
laughingman7743 marked this pull request as ready for review October 1, 2026 16:46
@laughingman7743
laughingman7743 marked this pull request as draft October 1, 2026 16:57
@laughingman7743
laughingman7743 force-pushed the fix/840-sql-cursor-interrupt branch from 2f4dcc4 to a32a993 Compare October 1, 2026 16:58
Comment thread pyathena/aio/common.py
raise DatabaseError(*e.args) from e
return cast(str, response.get("QueryExecutionId"))

@override

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rebase onto 39f22dc (#917 @override) and independent follow-up (relayed).

The series ecb19c15..2f4dcc43ec34ad63294925d75850eef7578eed05 was rebased to 39f22dc1..a32a9930c46943c2572f94627946088a1b145d45. The only conflict was the @override that #917 put on _get_query_execution. Following #917's convention, the new asyncio overrides _start_query_execution, _start_execution, and _cancel_and_wait also got @override, keeping their # type: ignore[override]. just lint passes. mypy --warn-unused-ignores reports only the three existing pyathena/sqlalchemy/ ignores. Offline SQL (27) and Spark (114) tests were rerun with the same results. AWS CI on the pre-rebase head 2f4dcc4 passed (run 36894523702).

Codex CLI, reported model gpt-6-astra, reasoning effort max (as reported), session 01a0f867-80a7-7a02-85f6-a436a679bf45. Scope: git range-diff ecb19c15..2f4dcc43 39f22dc1..a32a9930, read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.

CLEAN — no actionable findings.

Covered surfaces:

  • All eight commits via the requested range-diff. The only substantive rebase adjustments are three @override additions: _start_query_execution, _start_execution, and _cancel_and_wait.
  • All four named AioBaseCursor methods and the affected AioSparkCursor methods. Their markers target existing base methods and follow Mark the asyncio overrides of sync methods with @override #917’s decorator/type-ignore convention.
  • Mark the asyncio overrides of sync methods with @override #917’s imports, decorators, type-ignore removals, and pyathena.util.override shim. All are preserved, including in both overlapping asyncio files.
  • Behavior preservation: method bodies are unchanged from the reviewed head; the shim returns functions unchanged. The series’ remaining source files, documentation, and tests are byte-identical.

Static review only; no builds, tests, writes, or network access. HEAD remains a32a9930, and the worktree is clean.

@laughingman7743
laughingman7743 marked this pull request as ready for review October 1, 2026 17:02
@laughingman7743
laughingman7743 marked this pull request as draft October 2, 2026 00:23
laughingman7743 and others added 5 commits October 2, 2026 09:23
With kill_on_interrupt enabled, the SQL cursors returned the final
execution after cancelling an interrupted query instead of re-raising
the KeyboardInterrupt or asyncio.CancelledError, and an interrupt while
StartQueryExecution was in flight left the started query running without
an ID on the cursor.

Move the Spark cursor's interrupt handling into BaseCursor (_poll,
_cancel_and_wait, _start_execution, _wait_for_start) so that the SQL and
synchronous Spark cursors share it, and give AioBaseCursor the asyncio
counterparts. Cursors record the ID of a started execution through
_set_interrupted_execution_id(), which WithResultSet and SparkBaseCursor
override. The Spark cursors also clear the previous calculation before
starting a new one.

Closes #840

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The asyncio cursors sent StartQueryExecution or StartCalculationExecution
and then stopped the execution when the task was cancelled after
execute() had scheduled the shielded start task but before that task
began. Skip the request when the caller has been cancelled since, as the
synchronous cursors already do, so that it is never sent.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
laughingman7743 and others added 3 commits October 2, 2026 09:23
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…d queries

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@laughingman7743
laughingman7743 force-pushed the fix/840-sql-cursor-interrupt branch from a32a993 to c6357d0 Compare October 2, 2026 00:25
Comment thread pyathena/spark/common.py
@@ -336,38 +329,6 @@ def _poll_until_terminal(
time.sleep(self._poll_interval)

@override

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rebase onto 9ef49d3 (#918 explicit-override) and independent follow-up (relayed).

The series 39f22dc1..a32a9930c46943c2572f94627946088a1b145d45 was rebased to 9ef49d31..c6357d07e5c0e0508f46f535bcbfd9d596ba2fe8.

The conflicts were in pyathena/spark/common.py, where the series removes SparkBaseCursor._poll and _wait_for_calculation_start. Both removals are kept, and #918's @override on _cancel is preserved.

_cancel_and_wait and _set_interrupted_execution_id now override BaseCursor, so they carry @override. Removing the marker from _cancel_and_wait makes mypy report explicit-override (checked, then restored).

just lint and just docs lint pass. The offline SQL (27) and Spark (114) tests were rerun with the same results.

Codex CLI, reported model gpt-6-astra, reasoning effort max (as reported), session 01a0fa00-6299-7391-9054-92a2e537200c. Scope: git range-diff 39f22dc1..a32a9930 9ef49d31..c6357d07, read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.

Covered surfaces:

  • All eight commits via the supplied range-diff: seven unchanged; the first differs only in Spark annotation integration and context.
  • Mark every checkable override with @override #918’s imports, decorators, and mypy explicit-override setting remain intact.
  • Spark removes only the reviewed series’ code. _cancel_and_wait and _set_interrupted_execution_id correctly carry @override; surviving overrides retain their annotations.
  • New AioBaseCursor overrides are decorated. The base-less WithResultSet._set_interrupted_execution_id correctly remains undecorated.
  • No behavioral drift found: differences between reviewed and rebased heads are annotation-related; the decorator is a runtime no-op. Tests and documentation are identical.

CLEAN — no actionable findings.

Static inspection only; no builds, tests, network access, or writes. HEAD remains c6357d07, with a clean worktree.

@laughingman7743
laughingman7743 marked this pull request as ready for review October 2, 2026 00:30

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Cursors swallow the interrupt after kill_on_interrupt cancellation

1 participant