| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 309c379 commit ad50bc7
2 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -570,7 +570,7 @@ def from_proto(data_source: DataSourceProto): | |||
| 570 | 570 | owner=data_source.owner, | |
| 571 | 571 | batch_source=( | |
| 572 | 572 | DataSource.from_proto(data_source.batch_source) | |
| 573 | - if data_source.batch_source | ||
| 573 | + if data_source.HasField("batch_source") | ||
| 574 | 574 | else None | |
| 575 | 575 | ), | |
| 576 | 576 | ) | |
@@ -757,7 +757,7 @@ def from_proto(data_source: DataSourceProto): | |||
| 757 | 757 | owner=data_source.owner, | |
| 758 | 758 | batch_source=( | |
| 759 | 759 | DataSource.from_proto(data_source.batch_source) | |
| 760 | - if data_source.batch_source | ||
| 760 | + if data_source.HasField("batch_source") | ||
| 761 | 761 | else None | |
| 762 | 762 | ), | |
| 763 | 763 | ) | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -379,3 +379,41 @@ def kafka_source(watermark): | |||
| 379 | 379 | KafkaSource.from_proto(unset_proto).kafka_options.watermark_delay_threshold | |
| 380 | 380 | is None | |
| 381 | 381 | ) | |
| 382 | + | ||
| 383 | + | ||
| 384 | + @pytest.mark.parametrize( | ||
| 385 | + "stream_source", | ||
| 386 | + [ | ||
| 387 | + KafkaSource( | ||
| 388 | + name="test_source", | ||
| 389 | + kafka_bootstrap_servers="test_servers", | ||
| 390 | + message_format=ProtoFormat("class_path"), | ||
| 391 | + topic="test_topic", | ||
| 392 | + timestamp_field="event_timestamp", | ||
| 393 | + ), | ||
| 394 | + KinesisSource( | ||
| 395 | + name="test_source", | ||
| 396 | + region="test_region", | ||
| 397 | + record_format=ProtoFormat("class_path"), | ||
| 398 | + stream_name="test_stream", | ||
| 399 | + timestamp_field="event_timestamp", | ||
| 400 | + ), | ||
| 401 | + ], | ||
| 402 | + ids=["kafka", "kinesis"], | ||
| 403 | + ) | ||
| 404 | + def test_stream_source_without_batch_source_round_trips(stream_source): | ||
| 405 | + """``batch_source`` is optional on Kafka and Kinesis sources. | ||
| 406 | + | ||
| 407 | + ``from_proto`` checked ``if data_source.batch_source``, but an unset proto | ||
| 408 | + sub-message is still truthy, so the empty message was parsed as a data | ||
| 409 | + source and failed with "Could not identify the source type being added." | ||
| 410 | + That made these sources impossible to read back from the registry. | ||
| 411 | + """ | ||
| 412 | + proto = stream_source.to_proto() | ||
| 413 | + assert not proto.HasField("batch_source") | ||
| 414 | + | ||
| 415 | + restored = DataSource.from_proto(proto) | ||
| 416 | + | ||
| 417 | + assert isinstance(restored, type(stream_source)) | ||
| 418 | + assert restored.batch_source is None | ||
| 419 | + assert restored.name == stream_source.name | ||
| Back | FazBrowse Home | New Git URL |
0 commit comments