Bulk publish: per-message metadata in publish_events - #1228
Conversation
publish_events built each BulkPublishRequestEntry from entry_id, event and content_type only, so per-message metadata such as partitionKey could not be set. The runtime also merges request-level metadata into an entry only when that entry already carries metadata, so publish_metadata never reached the cloud event envelope either: ttlInSeconds set expiration on publish_event but not on publish_events. - Add BulkPublishEntry(event, metadata, content_type, entry_id), exported from dapr.clients, mirroring go-sdk PublishEventsEvent and java-sdk BulkPublishEntry. Plain bytes and str entries keep working. - Copy publish_metadata onto every entry. Entry keys override request keys, matching the runtime merge order. The request-level metadata is still sent, because the runtime reads rawPayload from it. - Share one entry builder between the sync and async clients. - Unit tests assert on the recorded BulkPublishRequest. The integration test proves ttlInSeconds reaches a plain entry through the expiration extension. Fixes dapr#1214 Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Nelson Parente <nelson_parente@live.com.pt>
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #1228 +/- ##
==========================================
- Coverage 83.91% 83.89% -0.03%
==========================================
Files 123 123
Lines 10260 10265 +5
==========================================
+ Hits 8610 8612 +2
- Misses 1650 1653 +3 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
- Port the plain-entry and type-default unit tests to the async client so both clients carry the same coverage of the shared entry builder. - Document that a BulkPublishEntry fixes its entry_id at construction, so one instance must not appear twice in the same data sequence. - Reuse _next_message in the existing streaming integration test instead of keeping two copies of the thread-and-future pattern. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Nelson Parente <nelson_parente@live.com.pt>
|
A few small things, none blocking:
|
sicoyle
left a comment
There was a problem hiding this comment.
one comment from me - thanks!!
| ValueError: event is not bytes or str. | ||
| """ | ||
| if not isinstance(event, (bytes, str)): | ||
| raise ValueError(f'invalid type for event {type(event)}') |
There was a problem hiding this comment.
| raise ValueError(f'invalid type for event {type(event)}') | |
| raise TypeError(f'invalid type for event {type(event)}') |
There was a problem hiding this comment.
Done in d54d25a. BulkPublishEntry now raises TypeError. The docstrings and the unit tests, sync and async, follow.
Resolves the conflict in tests/integration/test_pubsub.py. The next_message None loop from dapr#1234 moves into the shared _next_message helper, so the bulk metadata test gets it too. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Nelson Parente <nelson_parente@live.com.pt>
- BulkPublishEntry raises TypeError for an event that is not bytes or str. It raised ValueError before. - BulkPublishEntry defines __eq__ and __repr__. A failed entry is now readable in a log or in a test failure. - The async integration test mirrors the sync test. It publishes a plain entry and an entry with its own metadata, and asserts the expiration extension on both. Without the publish_metadata copy, the plain entry fails. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Nelson Parente <nelson_parente@live.com.pt>
|
Thanks Casper. All four points, in d54d25a:
Also merged upstream |
Description
publish_eventsbuilt eachBulkPublishRequestEntryfromentry_id,eventandcontent_typeonly, so per-message metadata such aspartitionKeyor a Kafka key could not be set. The runtime also merges request-level metadata into an entry only when that entry already carries metadata, sopublish_metadatanever reached the cloud event envelope either. On 1.18.4,ttlInSecondssetsexpirationwithpublish_eventbut not withpublish_events.This PR:
BulkPublishEntry(event, metadata=None, content_type=None, entry_id=None), exported fromdapr.clients. The shape mirrors go-sdkPublishEventsEventand java-sdkBulkPublishEntry. Plainbytesandstrentries keep working unchanged. An event of any other type raisesTypeError. The class defines__eq__and__repr__.publish_metadataonto every entry. Entry keys override request keys, which matches the runtime merge order. The request-level metadata is still sent because the runtime readsrawPayloadfrom it.pubsub-simpleexample to wrap one bulk event inBulkPublishEntry. The validated output is unchanged.Tests:
BulkPublishRequestthe fake sidecar recorded: entry metadata, request metadata copied onto every entry, entry-over-request precedence, caller entry IDs, content type precedence, type defaults, invalid event type.BulkPublishEntrywithttlInSecondsand assert that both received messages carry theexpirationextension. Redis pub/sub has no native TTL, soexpirationappears only when the metadata reaches the entry. Verified locally against runtime 1.17.2: with the copy disabled, both tests fail on the plain entry.CI note:
codecov/projectreports -0.07% with the patch at 100% coverage. The two clients lost their inline entry loops, which were covered lines, and the aio client's ratio shifts by 11 lines. No remaining line lost coverage.RELEASE NOTE: ADD Per-message metadata support in
publish_eventsbulk publish.Issue reference
Please reference the issue this PR will close: #1214
Checklist
Please make sure you've completed the relevant tasks for this PR, out of the following list:
pubsub-simpleexample here; the bulk publish page and the Python client page in Bulk publish: Python tab uses the SDK with per-entry metadata docs#5336🤖 Generated with Claude Code