Skip to content

feat(realtime): add ordered fetch worker - #84

Merged
swkeever merged 3 commits into
mainfrom
skeever/realtime-fetch-queue
Sep 3, 2026
Merged

swkeever merged 3 commits into
mainfrom
skeever/realtime-fetch-queue

Conversation

@swkeever

@swkeever swkeever commented Sep 3, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

  • add a bounded FIFO worker for realtime Postgres row fetches
  • preserve lightweight fallbacks in fetch outcomes
  • continue ordered processing after ordinary fetch failures
  • surface worker failure to both immediate and backpressured producers
  • make shutdown cancellation-safe

Scope

Internal worker only. Public channel configuration, session invalidation, and subscription lifecycle integration remain separate slices.

Verification

  • OpenAPI generated-client check
  • Ruff check and format check
  • mypy and Pyright
  • 388 unit tests
  • contract feature dry run
  • package build
  • Hosting SDK contract unit and shell suites

@swkeever

swkeever commented Sep 3, 2026

Copy link
Copy Markdown
Collaborator Author

@codex review

@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Sep 3, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-09-03T16:05:57.197101Z ca2cf88 Manual request
ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: c120c57c77

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment on lines +72 to +74
if self._task is None:
self._task = asyncio.create_task(self._run())
await self._queue.put(job)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Reject enqueues after the worker task has failed

When _deliver raises, or _run exits for any other exception, _task remains a completed non-None task, so subsequent calls never start another consumer. Up to queue_limit jobs then appear to enqueue successfully, while the next enqueue—or close() while the queue is full—blocks forever because nothing drains the queue; detect task.done() and surface the worker failure before accepting more jobs.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Fixed in 9c0e30b. Enqueue now surfaces a completed worker task failure before accepting another job; the regression verifies the original delivery exception is returned.

Comment on lines +81 to +83
self._closed = True
if task is not None:
await self._queue.put(_STOP_WORKER)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Make stop-marker insertion cancellation-safe

If close() is cancelled while the bounded queue is full, cancellation can occur in await self._queue.put(...) after _closed has already become true but before the stop marker is inserted. The worker then drains the accepted jobs and waits indefinitely, and every later close() skips insertion because _closed is already true and consequently waits forever on that worker task.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Fixed in 9c0e30b. Stop-marker insertion now runs in a retained task shielded from caller cancellation, so a later close resumes the same shutdown; the regression cancels a full-queue close and verifies all accepted work drains.

@swkeever

swkeever commented Sep 3, 2026

Copy link
Copy Markdown
Collaborator Author

@codex review

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 9c0e30bab2

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

self._raise_worker_failure()
if self._task is None:
self._task = asyncio.create_task(self._run())
await self._queue.put(job)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Stop waiting on a full queue when the worker exits

When the queue is full and an enqueue is already suspended here, _deliver can raise after the pre-put _raise_worker_failure() check has run. The consumer then exits without freeing the queued slot, so this await never completes and continues holding _state_lock; close() and later enqueues consequently block forever as well. The new pre-put check is fresh evidence because it cannot observe a failure that occurs during backpressure; wait for both the queue insertion and worker task, and propagate the worker failure if it finishes first.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Fixed in ca2cf88. A blocked enqueue now races its queue insertion against the worker task, cancels the pending put when the worker exits, and propagates the worker failure. The regression reproduces the full-queue delivery-failure race.

@swkeever

swkeever commented Sep 3, 2026

Copy link
Copy Markdown
Collaborator Author

@codex review

@chatgpt-codex-connector

Copy link
Copy Markdown

Codex Review: Didn't find any major issues. Chef's kiss.

Reviewed commit: ca2cf88d5c

ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

@swkeever
swkeever marked this pull request as ready for review September 3, 2026 16:07
@swkeever
swkeever merged commit bf8e94a into main Sep 3, 2026
2 checks passed
@swkeever
swkeever deleted the skeever/realtime-fetch-queue branch September 3, 2026 16:07
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.

1 participant