refactor: reduce cross-file duplication (round 2) - #183
Merged
Merged
Conversation
The file backends store each job as a JSON file and rebuild the Job from that dict on pop. The seven-field mapping (id, topic-or-default, data, priority, attempts, error) was spelt out at three pop sites — FileQueue's claim-by-id, LiteBackend.pop() and LiteBackend.pop_batch(). It now lives once as Job.from_data(queue, job_data, default_topic). Behaviour-preserving: same fields, same defaults, same topic fallback. The reclaim path in LiteBackend keeps its own construction — it overrides attempts and error and defaults data, so it is not the same mapping. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
RateLimiter.apply() already holds the enforcement path — resolve the IP, check it, set the x-ratelimit headers, and on refusal add retry-after and the 429. That same body was copied into two staticmethod entry points: RateLimiter.before_rate_limit (shared instance) and RateLimiterMiddleware.before_rate_limit (the middleware's own instance). Both now just resolve their limiter and delegate to apply(). The response shape, header set and 429 wording are unchanged — the enforcement lives in one place. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
The Kafka and RabbitMQ backends both keep dead jobs on a <topic>.dead_letter topic and adapt an external connector to the same dequeue/enqueue interface. Their dead_letters() and retry_job() were byte-identical drain-filter-requeue loops — pure over that interface, carrying no broker specifics. They now live once in TopicDeadLetterMixin, which both backends inherit. The broker-specific methods (push, pop, fail, retry, reject, purge, clear) stay in each backend unchanged — those differ genuinely (Kafka re-produces past a committed offset; RabbitMQ tracks jobs and re-publishes with an attempt count) and are not touched. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
The lazy relationship descriptor and the imperative has_many() both page through ALL children in _LAZY_PAGE_SIZE blocks so a parent with more than one page never loses the tail (REL-EAGER-UNBOUNDED). The page loop was written out in both places; it now lives once as fields.fetch_all_pages(db, sql, params, start_offset). Behaviour-preserving: same page size, same stop-on-short-block, same start offset. model.py already imported from fields, so the helper adds no new coupling. Both callers still wrap the returned rows in their own model. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
Both connection types carry room membership — the raw-socket WebSocketConnection and the ASGI _AsgiWebSocketConnection. join_room, leave_room, broadcast_to_room and the rooms property only touch this connection's own _rooms set and delegate to _manager, so they do not depend on the transport. They now live once in RoomMemberMixin, which both connections inherit; server.py already imports from the websocket package. Transport-specific code (framing, send, ping, close, the ASGI vs stream-reader wiring) stays in each class unchanged. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
…apter SQLite, MySQL and PostgreSQL each appended their own pagination clause with the same if/else: unpaginated returns the SQL and params untouched; paginated puts LIMIT/OFFSET on a new line (so a trailing -- comment cannot swallow it) with two bound values. Only the placeholder differed — ? on SQLite, %s on the servers — which each adapter already declares as PARAM_MARKER. The append now lives once as DatabaseAdapter._paginate_sql(). Each adapter still decides `paginated` for itself (SQLite and PostgreSQL exclude writes; MySQL does not), so no engine's pagination rule changes. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
RateLimiter.apply() and the two before_rate_limit staticmethods (the class-based middleware entry points) had no coverage — the production dispatch path uses check()/apply_headers() directly, so nothing exercised apply() or the wrappers. The dedup that routed both wrappers through apply() could have regressed with the suite still green. Adds real-object tests (real Response, plain request holder, no doubles) that drive apply(), RateLimiter.before_rate_limit and RateLimiterMiddleware.before_rate_limit over the limit and assert the 429, the retry-after header and the rate-limit headers. Mutation-proven: gutting apply()'s 429 fails 3, and cutting either wrapper's delegation fails its test. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
dead_letters() on the Kafka/RabbitMQ backends was already gated by the deterministic dead-letter test, but retry_job() — the other method now shared through TopicDeadLetterMixin — had only a surface "it resolves" check. Adds two live-Kafka regressions: retry(id) reports True and clears the dead-letter topic; retry(unknown-id) reports False and leaves the real dead letter in place. Re-consumption of the re-produced record is a Kafka offset concern outside this method's contract, so it is not asserted. Mutation-proven on live Kafka: forcing retry_job to return False fails the revive test; skipping its found-is-None guard fails the unknown-id test. Co-Authored-By: Claude Opus 4.8 <[email protected]> Co-Authored-By: Tina4 <[email protected]>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Reduce cross-file duplication (round 2)
Genuine cross-file dedup only.
tina4 metricsbefore → after:Every cross-file duplicate is gone. The 13 that remain are same-file blocks (response.py, frond/engine.py, public/js/frond.js, router.py, auth, orm/model.py, mongodb_handler, connection.py, messenger) — out of scope for this pass and left untouched.
Deduped — shared logic moved to a cohesive home
queue/__init__.py↔queue/lite_backend.py(Job built from a persisted dict, ×3)Job.from_data(queue, job_data, default_topic)inqueue/job.pycore/middleware.py↔core/rate_limiter.py(429 enforcement, ×3)before_rate_limitwrappers delegate toRateLimiter.apply()queue/kafka_backend.py↔queue/rabbitmq_backend.py(dead_letters()+retry_job(), byte-identical)TopicDeadLetterMixininqueue/topic_dead_letter.pyorm/fields.py↔orm/model.py(uncapped has_many paging loop)fields.fetch_all_pages(db, sql, params, start_offset)core/server.py↔websocket/__init__.py(room membership on two connection types)RoomMemberMixininwebsocket/__init__.pydatabase/mysql.py↔database/postgres.py(+sqlite.py) (LIMIT/OFFSET clause)DatabaseAdapter._paginate_sql()All six are behaviour-preserving: same JSON, same defaults, same wire output, same status codes.
Look-alikes deliberately LEFT (not force-merged)
push,pop,fail,retry,reject,purge,cleardiffer genuinely (Kafka re-produces past a committed offset; RabbitMQ tracks jobs and re-publishes with an attempt count). Only the two byte-identical, broker-agnostic methods were shared. The refusal docstrings stay per-backend.LiteBackendreservation-reclaim Job construction — looks likeJob.from_databut overridesattemptsanderrorand defaultsdata; not the same mapping, left alone.paginateddecision — SQLite/PostgreSQL exclude writes, MySQL does not; only the clause-append was shared, the decision stays per-engine.Tests — real engines, no mocks, mutation-proven
Characterization + regression, each shared helper proven to gate by mutation:
Job.from_dataerrorfield → queue suite redRateLimiter.apply+ wrappersTestRateLimiterEnforcement(realResponse)TopicDeadLetterMixin.dead_lettersreturn []→ 2 Kafka tests failTopicDeadLetterMixin.retry_jobreturn False→ revive test fails; skip found-guard → unknown-id test failsfetch_all_pagesRoomMemberMixinjoin_room→ 4 fail_paginate_sqlpaginated→ 2 failTwo new tests lock previously-untested surface: the class-based rate-limit middleware entry points, and
retry_job()on the topic backends (only surface-checked before).carbonah lint tina4_python: 62 findings before and after — identical set, only line numbers shifted by the removed/added lines. No new findings.Verification
Every touched subsystem was run green on real engines on the lab (no mocks): the queue backends against live Kafka, RabbitMQ and MongoDB; the ORM relationships, the rate limiter, the WebSocket rooms, and the database adapters against real SQLite, MySQL and PostgreSQL. Each shared helper is mutation-proven to gate (table above).
One caveat, and it is not this change: a whole-suite run on the shared lab hangs in the Firebird driver's commit under concurrent load. This reproduces identically on the untouched
v3baseline, and the Firebird tests — including the pagination path through the base adapter this PR edits — pass fast in isolation on both baseline and this branch (60 passed). So the hang is lab load contention, not the dedup. This PR's own GitHub CI is the authoritative full run and the merge gate.🤖 Generated with Claude Code