Repository navigation
Conversation
fsspec license metadata is not detectable by licensecheck, causing CI failures. Add to ignore list since it's BSD-3-Clause licensed.
Implements issue #102: Base component for external communication via pub/sub message broker infrastructure. New base classes: - MessageDataReader: Abstract base for reading data from message brokers with connection management, reconnection with exponential backoff, retry logic, message acknowledgment, and chunked/buffered reading. - MessageDataWriter: Abstract base for writing data to message brokers with connection management, reconnection, retry logic, and chunked/buffered writing. Concrete implementations: - GCPPubSubDataReader/Writer: Google Cloud PubSub - AWSSQSDataReader/AWSSNSDataWriter: AWS SQS/SNS - KafkaDataReader/Writer: Apache Kafka Also includes: - Message broker exceptions (ConnectionError, TransientError, PermanentError) - Settings for GCP PubSub, AWS, and Kafka - Optional dependencies in pyproject.toml - Proposal document with design rationale - Comprehensive unit tests (72 new tests)
|
Benchmark comparison for |
- Fix ruff lint errors: import sorting, unused imports, S110 noqa comments - Fix ruff format errors in gcp_pubsub_io.py - Fix mypy overlap errors: remove duplicate fields from ArgsDict TypedDicts - Fix mypy multiple values error: use kwargs.setdefault instead of pop - Remove untracked test data files causing lint failures
|
Benchmark comparison for |
Increase connection establishment sleep in _ZMQPipelineConnectorProxy from 0.1s to 0.5s to allow the proxy subprocess's SUB socket subscription to propagate to XPUB before the sender starts publishing (ZMQ slow joiner problem). Also mark the test as flaky with 3 reruns following the existing pattern used elsewhere in the repo. Fixes: test_process_with_components_run[RayProcess-zmq_connector_cls-zmq_pubsub_proxy=True-10-2.0]
|
Benchmark comparison for |
Codecov Report❌ Patch coverage is 📢 Thoughts on this report? Let us know! |
- Fix invalid-raise errors in message_reader.py and message_writer.py by initializing last_exception with a non-None default instead of Optional[Exception] - Remove PublisherClient.close() call in gcp_pubsub_io.py (method does not exist on the client); just set reference to None for GC - Update test to match new disconnect behavior
|
Benchmark comparison for |
The design proposal is not tracked in the repo. Added to .gitignore to prevent accidental re-commit.
|
Benchmark comparison for |
Resolve conflict in tests/integration/test_process_with_components_run.py: keep main's native pytest parametrization (ConnectorCase + indirect + ids) replacing pytest-cases, retaining the flaky marker for the Ray + ZMQ proxy case. Regenerate uv.lock for merged pyproject extras. Verified: ruff check, ruff format, ty check all pass; new ZMQ backend/ channel/proxy unit tests (21) and message data tests (72) pass; merged integration file collects cleanly (25 tests).
|
Benchmark comparison for |
Blocking issues: - An empty broker poll no longer ends the stream. `_receive` returning [] means "nothing yet", so `step()` waits and polls again like `WebsocketReader`; only a subclass raising `NoMoreDataException` closes the IO (B1). - A batch is acknowledged only after its last record is consumed, instead of after the first one. Previously a crash mid-batch lost the unread messages (B2). - `KafkaDataReader._ack` commits the offsets following the processed records, per partition, instead of the consumer position, which could mark unprocessed records as done (B3). - GCP reader and writer dispatch every synchronous PubSub call (client construction, pull, acknowledge, publish, future.result, stop, close) to a worker thread, so they no longer block the shared event loop (B4, B5). Design and consolidation: - `MessageDataReader`/`MessageDataWriter` now extend `DataReader`/`DataWriter`, reusing the buffer, binding and batching logic instead of copying it (W14). - Retry, backoff and reconnect live in one place: `plugboard.utils.retry` with a frozen `RetryPolicy` value object shared by both bases (W13, L1, L2). - A failing reconnect no longer escapes the retry loop; it costs an attempt and the broker error is raised once attempts run out (W12). - Row-to-message encoding is shared as `iter_records`/`encode_records`/ `encode_records_bytes`; the old per-row deque indexing was quadratic in buffer size and duplicated across three writers (W16, W19). - Connection access is serialised with an `asyncio.Lock` across ack and reconnect (W15). - Kafka writer submits all records then flushes once, instead of one round trip per record (W18). SQS deletes ack batches of up to 10 with `delete_message_batch` and reports partial failures rather than pretending success (L3). - The GCP publisher is stopped on disconnect, so a reconnect no longer leaks a channel and its threads (W17). - `MessageBrokerPermanentError`/`TransientError` are now actually raised and honoured: each broker maps its SDK's failures onto the hierarchy, and permanent errors skip the retries that previously burnt against a guaranteed failure (W1). Configuration and packaging: - The GCP/AWS/Kafka settings are read as constructor fallbacks via `resolve_argument`, so `GCP_PUBSUB_PROJECT_ID`, `AWS_REGION` and `KAFKA_BOOTSTRAP_SERVERS` work (W2). - `aws-messaging` now uses `aiobotocore` instead of `aioboto3`, which no longer forces the shared botocore/boto3 stack backwards (W5). - Broker SDKs are declared in the `test` dependency group, which also resolves the three `ty` unresolved-import findings (W10 root cause). - The six concrete broker classes are exported from `plugboard.library` (W3). - Removed the misleading broker-specific ArgsDicts that documented attributes they did not declare (W20). Tests: 126 tests, up from 72, all using teardown-aware monkeypatch fixtures instead of global `sys.modules` mocks, asserting broker calls rather than private attribute identity (W10, W11). New coverage for retry exhaustion and the delay cap, empty-batch handling, error-branch mapping, offset commits, reconnect failure, and the acknowledgment timing that previously lost records (W7, W8, W9). Record encoding is parametrised once rather than copied across three modules (L4). Docs: new Message Data usage page (extras, worked example, delivery semantics), broker settings in the configuration page, and the new extras in the README (W4). Test isolation: the S3 file tests pinned their bucket region so they no longer depend on the developer's ambient AWS profile, which the newer botocore resolves over `AWS_REGION`.
The connection delay and the flaky marker on the Ray + ZMQ integration test fix an unrelated ZMQ connector issue; they now move with it to their own change.
|
Benchmark comparison for |
…ure at teardown Closes the two remaining uncovered branches from the error-handling paths the review asked about: a reconnect whose disconnect raises must still connect and let the retry proceed, and a send still in flight when destroy runs must be reported rather than swallowed. plugboard/utils/retry.py is now at 100% branch coverage.
Review responseAddressed in
Deferred (W6): real-infrastructure integration tests, pending the AWS account and GCP project plus the Terraform in |
|
Benchmark comparison for |
Summary
Implements #102: feat: Base component for external communication (#102).
Adds
MessageDataReader/MessageDataWriterfor reading and writing records through a pub/sub broker - the broker equivalent ofFileReader/FileWriter, where each message carries one record andfield_namesbecome the component's outputs (reader) or inputs (writer). Both extend the existingDataReader/DataWriterbases, adding what a broker needs and a finite source does not: a long-lived connection, reconnection with backoff, and acknowledgment.Three implementations: Google Cloud PubSub, AWS SQS/SNS, Apache Kafka.
Integration tests against real infrastructure are deliberately deferred - no AWS/GCP accounts or broker resources exist yet. Terraform for the required PubSub topic/subscription and SNS topic/SQS queue lives in
plugboard-infra(terraform/ci-data-message-testing), anddocs/running-message-data-tests.mdthere describes wiring its outputs into this suite. All 128 tests here are mock-based, so no credentials are needed.Changes
Base classes
MessageDataReader/MessageDataWriterextendDataReader/DataWriter(buffer, input binding and batching are inherited, not copied), adding_connect,_disconnect,_receive/_send,_convert, and_ack(reader).plugboard/utils/retry.py, with a frozenRetryPolicyvalue object used by both bases.MessageBrokerConnectionError/TransientError/PermanentError) are raised and honoured: each implementation maps its SDK's failures onto them, and permanent errors skip retries that cannot help.GCP_PUBSUB_PROJECT_ID,AWS_REGION,KAFKA_BOOTSTRAP_SERVERS); explicit arguments win.plugboard.library.Implementations
GCPPubSubDataReader/GCPPubSubDataWriter(google-cloud-pubsub) - every synchronous PubSub call is dispatched to a worker thread, so a pull or publish never stalls the shared event loop.AWSSQSDataReader/AWSSNSDataWriter(aiobotocore) - long polling,delete_message_batchfor acks, and the 10-message SQS cap is logged rather than applied silently.KafkaDataReader/KafkaDataWriter(aiokafka) - commits the offsets of processed records per partition; sends a batch then flushes once, so producer batching is used.Packaging: optional extras
gcp-pubsub,aws-messaging,kafka; the broker SDKs are also declared in thetestdependency group.Docs: new Message Data usage page (extras, worked example, delivery semantics), broker settings in the configuration page, new extras in the README.
Behaviour worth knowing
WebsocketReader. A subclass signals a genuinely exhausted source by raisingNoMoreDataException, which closes the IO as other data readers do.Testing
monkeypatchfixtures rather than globalsys.modulesmocks, asserting broker calls and arguments rather than private attribute identity.utils/retry.py100%,message_reader.py95%,message_writer.py94%,gcp_pubsub_io.py94%,aws_messaging_io.py95%,kafka_io.py96% - the remaining misses are abstract-method bodies.ruff check,ruff format --check,tyoverplugboard/,plugboard-schemas/andtests/all pass;mkdocs buildis clean.test_process_stop_event[LocalProcess-ZMQConnector-...]), which passes on repeat runs in isolation and is addressed in fix: avoid ZMQ slow-joiner loss on pipeline connect #297.Related
plugboard-infra/terraform/ci-data-message-testing- infrastructure for the deferred real-broker integration tests.