diff --git a/Makefile b/Makefile index b0127755..65e62b70 100644 --- a/Makefile +++ b/Makefile @@ -26,14 +26,10 @@ setup: poetry install --with dev --no-root proto: - poetry run python3 -m grpc_tools.protoc --pyi_out=pynumaflow/proto/sinker -I=pynumaflow/proto/sinker --python_out=pynumaflow/proto/sinker --grpc_python_out=pynumaflow/proto/sinker pynumaflow/proto/sinker/*.proto - poetry run python3 -m grpc_tools.protoc --pyi_out=pynumaflow/proto/mapper -I=pynumaflow/proto/mapper --python_out=pynumaflow/proto/mapper --grpc_python_out=pynumaflow/proto/mapper pynumaflow/proto/mapper/*.proto - poetry run python3 -m grpc_tools.protoc --pyi_out=pynumaflow/proto/reducer -I=pynumaflow/proto/reducer --python_out=pynumaflow/proto/reducer --grpc_python_out=pynumaflow/proto/reducer pynumaflow/proto/reducer/*.proto - poetry run python3 -m grpc_tools.protoc --pyi_out=pynumaflow/proto/sourcetransformer -I=pynumaflow/proto/sourcetransformer --python_out=pynumaflow/proto/sourcetransformer --grpc_python_out=pynumaflow/proto/sourcetransformer pynumaflow/proto/sourcetransformer/*.proto - poetry run python3 -m grpc_tools.protoc --pyi_out=pynumaflow/proto/sideinput -I=pynumaflow/proto/sideinput --python_out=pynumaflow/proto/sideinput --grpc_python_out=pynumaflow/proto/sideinput pynumaflow/proto/sideinput/*.proto - poetry run python3 -m grpc_tools.protoc --pyi_out=pynumaflow/proto/sourcer -I=pynumaflow/proto/sourcer --python_out=pynumaflow/proto/sourcer --grpc_python_out=pynumaflow/proto/sourcer pynumaflow/proto/sourcer/*.proto - poetry run python3 -m grpc_tools.protoc --pyi_out=pynumaflow/proto/accumulator -I=pynumaflow/proto/accumulator --python_out=pynumaflow/proto/accumulator --grpc_python_out=pynumaflow/proto/accumulator pynumaflow/proto/accumulator/*.proto - - - sed -i.bak -e 's/^\(import.*_pb2\)/from . \1/' pynumaflow/proto/*/*.py - rm pynumaflow/proto/*/*.py.bak + poetry run python3 -m grpc_tools.protoc -Ipynumaflow/proto/sinker=pynumaflow/proto/sinker --pyi_out=. --python_out=. --grpc_python_out=. pynumaflow/proto/sinker/*.proto + poetry run python3 -m grpc_tools.protoc -Ipynumaflow/proto/mapper=pynumaflow/proto/mapper --pyi_out=. --python_out=. --grpc_python_out=. pynumaflow/proto/mapper/*.proto + poetry run python3 -m grpc_tools.protoc -Ipynumaflow/proto/reducer=pynumaflow/proto/reducer --pyi_out=. --python_out=. --grpc_python_out=. pynumaflow/proto/reducer/*.proto + poetry run python3 -m grpc_tools.protoc -Ipynumaflow/proto/sourcetransformer=pynumaflow/proto/sourcetransformer --pyi_out=. --python_out=. --grpc_python_out=. pynumaflow/proto/sourcetransformer/*.proto + poetry run python3 -m grpc_tools.protoc -Ipynumaflow/proto/sideinput=pynumaflow/proto/sideinput --pyi_out=. --python_out=. --grpc_python_out=. pynumaflow/proto/sideinput/*.proto + poetry run python3 -m grpc_tools.protoc -Ipynumaflow/proto/sourcer=pynumaflow/proto/sourcer --pyi_out=. --python_out=. --grpc_python_out=. pynumaflow/proto/sourcer/*.proto + poetry run python3 -m grpc_tools.protoc -Ipynumaflow/proto/accumulator=pynumaflow/proto/accumulator --pyi_out=. --python_out=. --grpc_python_out=. pynumaflow/proto/accumulator/*.proto diff --git a/pynumaflow/proto/accumulator/accumulator_pb2.py b/pynumaflow/proto/accumulator/accumulator_pb2.py index 30a12ed2..f12f012f 100644 --- a/pynumaflow/proto/accumulator/accumulator_pb2.py +++ b/pynumaflow/proto/accumulator/accumulator_pb2.py @@ -1,7 +1,7 @@ # -*- coding: utf-8 -*- # Generated by the protocol buffer compiler. DO NOT EDIT! # NO CHECKED-IN PROTOBUF GENCODE -# source: accumulator.proto +# source: pynumaflow/proto/accumulator/accumulator.proto # Protobuf Python Version: 6.31.1 """Generated protocol buffer code.""" from google.protobuf import descriptor as _descriptor @@ -15,7 +15,7 @@ 31, 1, '', - 'accumulator.proto' + 'pynumaflow/proto/accumulator/accumulator.proto' ) # @@protoc_insertion_point(imports) @@ -26,32 +26,32 @@ from google.protobuf import timestamp_pb2 as google_dot_protobuf_dot_timestamp__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x11\x61\x63\x63umulator.proto\x12\x0e\x61\x63\x63umulator.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xf8\x01\n\x07Payload\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\n\n\x02id\x18\x05 \x01(\t\x12\x35\n\x07headers\x18\x06 \x03(\x0b\x32$.accumulator.v1.Payload.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"\xbe\x02\n\x12\x41\x63\x63umulatorRequest\x12(\n\x07payload\x18\x01 \x01(\x0b\x32\x17.accumulator.v1.Payload\x12\x45\n\toperation\x18\x02 \x01(\x0b\x32\x32.accumulator.v1.AccumulatorRequest.WindowOperation\x1a\xb6\x01\n\x0fWindowOperation\x12G\n\x05\x65vent\x18\x01 \x01(\x0e\x32\x38.accumulator.v1.AccumulatorRequest.WindowOperation.Event\x12\x30\n\x0bkeyedWindow\x18\x02 \x01(\x0b\x32\x1b.accumulator.v1.KeyedWindow\"(\n\x05\x45vent\x12\x08\n\x04OPEN\x10\x00\x12\t\n\x05\x43LOSE\x10\x01\x12\n\n\x06\x41PPEND\x10\x02\"}\n\x0bKeyedWindow\x12)\n\x05start\x18\x01 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\'\n\x03\x65nd\x18\x02 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04slot\x18\x03 \x01(\t\x12\x0c\n\x04keys\x18\x04 \x03(\t\"\x87\x01\n\x13\x41\x63\x63umulatorResponse\x12(\n\x07payload\x18\x01 \x01(\x0b\x32\x17.accumulator.v1.Payload\x12+\n\x06window\x18\x02 \x01(\x0b\x32\x1b.accumulator.v1.KeyedWindow\x12\x0c\n\x04tags\x18\x03 \x03(\t\x12\x0b\n\x03\x45OF\x18\x04 \x01(\x08\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\xac\x01\n\x0b\x41\x63\x63umulator\x12[\n\x0c\x41\x63\x63umulateFn\x12\".accumulator.v1.AccumulatorRequest\x1a#.accumulator.v1.AccumulatorResponse(\x01\x30\x01\x12@\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x1d.accumulator.v1.ReadyResponseBd\n#io.numaproj.numaflow.accumulator.v1Z=github.com/numaproj/numaflow-go/pkg/apis/proto/accumulator/v1b\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n.pynumaflow/proto/accumulator/accumulator.proto\x12\x0e\x61\x63\x63umulator.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xf8\x01\n\x07Payload\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\n\n\x02id\x18\x05 \x01(\t\x12\x35\n\x07headers\x18\x06 \x03(\x0b\x32$.accumulator.v1.Payload.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"\xbe\x02\n\x12\x41\x63\x63umulatorRequest\x12(\n\x07payload\x18\x01 \x01(\x0b\x32\x17.accumulator.v1.Payload\x12\x45\n\toperation\x18\x02 \x01(\x0b\x32\x32.accumulator.v1.AccumulatorRequest.WindowOperation\x1a\xb6\x01\n\x0fWindowOperation\x12G\n\x05\x65vent\x18\x01 \x01(\x0e\x32\x38.accumulator.v1.AccumulatorRequest.WindowOperation.Event\x12\x30\n\x0bkeyedWindow\x18\x02 \x01(\x0b\x32\x1b.accumulator.v1.KeyedWindow\"(\n\x05\x45vent\x12\x08\n\x04OPEN\x10\x00\x12\t\n\x05\x43LOSE\x10\x01\x12\n\n\x06\x41PPEND\x10\x02\"}\n\x0bKeyedWindow\x12)\n\x05start\x18\x01 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\'\n\x03\x65nd\x18\x02 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04slot\x18\x03 \x01(\t\x12\x0c\n\x04keys\x18\x04 \x03(\t\"\x87\x01\n\x13\x41\x63\x63umulatorResponse\x12(\n\x07payload\x18\x01 \x01(\x0b\x32\x17.accumulator.v1.Payload\x12+\n\x06window\x18\x02 \x01(\x0b\x32\x1b.accumulator.v1.KeyedWindow\x12\x0c\n\x04tags\x18\x03 \x03(\t\x12\x0b\n\x03\x45OF\x18\x04 \x01(\x08\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\xac\x01\n\x0b\x41\x63\x63umulator\x12[\n\x0c\x41\x63\x63umulateFn\x12\".accumulator.v1.AccumulatorRequest\x1a#.accumulator.v1.AccumulatorResponse(\x01\x30\x01\x12@\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x1d.accumulator.v1.ReadyResponseBd\n#io.numaproj.numaflow.accumulator.v1Z=github.com/numaproj/numaflow-go/pkg/apis/proto/accumulator/v1b\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) -_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'accumulator_pb2', _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'pynumaflow.proto.accumulator.accumulator_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: _globals['DESCRIPTOR']._loaded_options = None _globals['DESCRIPTOR']._serialized_options = b'\n#io.numaproj.numaflow.accumulator.v1Z=github.com/numaproj/numaflow-go/pkg/apis/proto/accumulator/v1' _globals['_PAYLOAD_HEADERSENTRY']._loaded_options = None _globals['_PAYLOAD_HEADERSENTRY']._serialized_options = b'8\001' - _globals['_PAYLOAD']._serialized_start=100 - _globals['_PAYLOAD']._serialized_end=348 - _globals['_PAYLOAD_HEADERSENTRY']._serialized_start=302 - _globals['_PAYLOAD_HEADERSENTRY']._serialized_end=348 - _globals['_ACCUMULATORREQUEST']._serialized_start=351 - _globals['_ACCUMULATORREQUEST']._serialized_end=669 - _globals['_ACCUMULATORREQUEST_WINDOWOPERATION']._serialized_start=487 - _globals['_ACCUMULATORREQUEST_WINDOWOPERATION']._serialized_end=669 - _globals['_ACCUMULATORREQUEST_WINDOWOPERATION_EVENT']._serialized_start=629 - _globals['_ACCUMULATORREQUEST_WINDOWOPERATION_EVENT']._serialized_end=669 - _globals['_KEYEDWINDOW']._serialized_start=671 - _globals['_KEYEDWINDOW']._serialized_end=796 - _globals['_ACCUMULATORRESPONSE']._serialized_start=799 - _globals['_ACCUMULATORRESPONSE']._serialized_end=934 - _globals['_READYRESPONSE']._serialized_start=936 - _globals['_READYRESPONSE']._serialized_end=966 - _globals['_ACCUMULATOR']._serialized_start=969 - _globals['_ACCUMULATOR']._serialized_end=1141 + _globals['_PAYLOAD']._serialized_start=129 + _globals['_PAYLOAD']._serialized_end=377 + _globals['_PAYLOAD_HEADERSENTRY']._serialized_start=331 + _globals['_PAYLOAD_HEADERSENTRY']._serialized_end=377 + _globals['_ACCUMULATORREQUEST']._serialized_start=380 + _globals['_ACCUMULATORREQUEST']._serialized_end=698 + _globals['_ACCUMULATORREQUEST_WINDOWOPERATION']._serialized_start=516 + _globals['_ACCUMULATORREQUEST_WINDOWOPERATION']._serialized_end=698 + _globals['_ACCUMULATORREQUEST_WINDOWOPERATION_EVENT']._serialized_start=658 + _globals['_ACCUMULATORREQUEST_WINDOWOPERATION_EVENT']._serialized_end=698 + _globals['_KEYEDWINDOW']._serialized_start=700 + _globals['_KEYEDWINDOW']._serialized_end=825 + _globals['_ACCUMULATORRESPONSE']._serialized_start=828 + _globals['_ACCUMULATORRESPONSE']._serialized_end=963 + _globals['_READYRESPONSE']._serialized_start=965 + _globals['_READYRESPONSE']._serialized_end=995 + _globals['_ACCUMULATOR']._serialized_start=998 + _globals['_ACCUMULATOR']._serialized_end=1170 # @@protoc_insertion_point(module_scope) diff --git a/pynumaflow/proto/accumulator/accumulator_pb2_grpc.py b/pynumaflow/proto/accumulator/accumulator_pb2_grpc.py index 232fc018..e8a6e86d 100644 --- a/pynumaflow/proto/accumulator/accumulator_pb2_grpc.py +++ b/pynumaflow/proto/accumulator/accumulator_pb2_grpc.py @@ -3,8 +3,8 @@ import grpc import warnings -from . import accumulator_pb2 as accumulator__pb2 from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 +from pynumaflow.proto.accumulator import accumulator_pb2 as pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2 GRPC_GENERATED_VERSION = '1.75.0' GRPC_VERSION = grpc.__version__ @@ -19,7 +19,7 @@ if _version_not_supported: raise RuntimeError( f'The grpc package installed is at version {GRPC_VERSION},' - + f' but the generated code in accumulator_pb2_grpc.py depends on' + + f' but the generated code in pynumaflow/proto/accumulator/accumulator_pb2_grpc.py depends on' + f' grpcio>={GRPC_GENERATED_VERSION}.' + f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}' + f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.' @@ -40,13 +40,13 @@ def __init__(self, channel): """ self.AccumulateFn = channel.stream_stream( '/accumulator.v1.Accumulator/AccumulateFn', - request_serializer=accumulator__pb2.AccumulatorRequest.SerializeToString, - response_deserializer=accumulator__pb2.AccumulatorResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.AccumulatorRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.AccumulatorResponse.FromString, _registered_method=True) self.IsReady = channel.unary_unary( '/accumulator.v1.Accumulator/IsReady', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=accumulator__pb2.ReadyResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.ReadyResponse.FromString, _registered_method=True) @@ -75,13 +75,13 @@ def add_AccumulatorServicer_to_server(servicer, server): rpc_method_handlers = { 'AccumulateFn': grpc.stream_stream_rpc_method_handler( servicer.AccumulateFn, - request_deserializer=accumulator__pb2.AccumulatorRequest.FromString, - response_serializer=accumulator__pb2.AccumulatorResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.AccumulatorRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.AccumulatorResponse.SerializeToString, ), 'IsReady': grpc.unary_unary_rpc_method_handler( servicer.IsReady, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=accumulator__pb2.ReadyResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.ReadyResponse.SerializeToString, ), } generic_handler = grpc.method_handlers_generic_handler( @@ -112,8 +112,8 @@ def AccumulateFn(request_iterator, request_iterator, target, '/accumulator.v1.Accumulator/AccumulateFn', - accumulator__pb2.AccumulatorRequest.SerializeToString, - accumulator__pb2.AccumulatorResponse.FromString, + pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.AccumulatorRequest.SerializeToString, + pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.AccumulatorResponse.FromString, options, channel_credentials, insecure, @@ -140,7 +140,7 @@ def IsReady(request, target, '/accumulator.v1.Accumulator/IsReady', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - accumulator__pb2.ReadyResponse.FromString, + pynumaflow_dot_proto_dot_accumulator_dot_accumulator__pb2.ReadyResponse.FromString, options, channel_credentials, insecure, diff --git a/pynumaflow/proto/mapper/map_pb2.py b/pynumaflow/proto/mapper/map_pb2.py index a2d52bfb..a9bb7332 100644 --- a/pynumaflow/proto/mapper/map_pb2.py +++ b/pynumaflow/proto/mapper/map_pb2.py @@ -1,7 +1,7 @@ # -*- coding: utf-8 -*- # Generated by the protocol buffer compiler. DO NOT EDIT! # NO CHECKED-IN PROTOBUF GENCODE -# source: map.proto +# source: pynumaflow/proto/mapper/map.proto # Protobuf Python Version: 6.31.1 """Generated protocol buffer code.""" from google.protobuf import descriptor as _descriptor @@ -15,7 +15,7 @@ 31, 1, '', - 'map.proto' + 'pynumaflow/proto/mapper/map.proto' ) # @@protoc_insertion_point(imports) @@ -26,32 +26,32 @@ from google.protobuf import timestamp_pb2 as google_dot_protobuf_dot_timestamp__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\tmap.proto\x12\x06map.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xac\x03\n\nMapRequest\x12+\n\x07request\x18\x01 \x01(\x0b\x32\x1a.map.v1.MapRequest.Request\x12\n\n\x02id\x18\x02 \x01(\t\x12)\n\thandshake\x18\x03 \x01(\x0b\x32\x11.map.v1.HandshakeH\x00\x88\x01\x01\x12/\n\x06status\x18\x04 \x01(\x0b\x32\x1a.map.v1.TransmissionStatusH\x01\x88\x01\x01\x1a\xef\x01\n\x07Request\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x38\n\x07headers\x18\x05 \x03(\x0b\x32\'.map.v1.MapRequest.Request.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x0c\n\n_handshakeB\t\n\x07_status\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"!\n\x12TransmissionStatus\x12\x0b\n\x03\x65ot\x18\x01 \x01(\x08\"\xf0\x01\n\x0bMapResponse\x12+\n\x07results\x18\x01 \x03(\x0b\x32\x1a.map.v1.MapResponse.Result\x12\n\n\x02id\x18\x02 \x01(\t\x12)\n\thandshake\x18\x03 \x01(\x0b\x32\x11.map.v1.HandshakeH\x00\x88\x01\x01\x12/\n\x06status\x18\x04 \x01(\x0b\x32\x1a.map.v1.TransmissionStatusH\x01\x88\x01\x01\x1a\x33\n\x06Result\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12\x0c\n\x04tags\x18\x03 \x03(\tB\x0c\n\n_handshakeB\t\n\x07_status\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32u\n\x03Map\x12\x34\n\x05MapFn\x12\x12.map.v1.MapRequest\x1a\x13.map.v1.MapResponse(\x01\x30\x01\x12\x38\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x15.map.v1.ReadyResponseB7Z5github.com/numaproj/numaflow-go/pkg/apis/proto/map/v1b\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n!pynumaflow/proto/mapper/map.proto\x12\x06map.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xac\x03\n\nMapRequest\x12+\n\x07request\x18\x01 \x01(\x0b\x32\x1a.map.v1.MapRequest.Request\x12\n\n\x02id\x18\x02 \x01(\t\x12)\n\thandshake\x18\x03 \x01(\x0b\x32\x11.map.v1.HandshakeH\x00\x88\x01\x01\x12/\n\x06status\x18\x04 \x01(\x0b\x32\x1a.map.v1.TransmissionStatusH\x01\x88\x01\x01\x1a\xef\x01\n\x07Request\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x38\n\x07headers\x18\x05 \x03(\x0b\x32\'.map.v1.MapRequest.Request.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x0c\n\n_handshakeB\t\n\x07_status\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"!\n\x12TransmissionStatus\x12\x0b\n\x03\x65ot\x18\x01 \x01(\x08\"\xf0\x01\n\x0bMapResponse\x12+\n\x07results\x18\x01 \x03(\x0b\x32\x1a.map.v1.MapResponse.Result\x12\n\n\x02id\x18\x02 \x01(\t\x12)\n\thandshake\x18\x03 \x01(\x0b\x32\x11.map.v1.HandshakeH\x00\x88\x01\x01\x12/\n\x06status\x18\x04 \x01(\x0b\x32\x1a.map.v1.TransmissionStatusH\x01\x88\x01\x01\x1a\x33\n\x06Result\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12\x0c\n\x04tags\x18\x03 \x03(\tB\x0c\n\n_handshakeB\t\n\x07_status\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32u\n\x03Map\x12\x34\n\x05MapFn\x12\x12.map.v1.MapRequest\x1a\x13.map.v1.MapResponse(\x01\x30\x01\x12\x38\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x15.map.v1.ReadyResponseB7Z5github.com/numaproj/numaflow-go/pkg/apis/proto/map/v1b\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) -_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'map_pb2', _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'pynumaflow.proto.mapper.map_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: _globals['DESCRIPTOR']._loaded_options = None _globals['DESCRIPTOR']._serialized_options = b'Z5github.com/numaproj/numaflow-go/pkg/apis/proto/map/v1' _globals['_MAPREQUEST_REQUEST_HEADERSENTRY']._loaded_options = None _globals['_MAPREQUEST_REQUEST_HEADERSENTRY']._serialized_options = b'8\001' - _globals['_MAPREQUEST']._serialized_start=84 - _globals['_MAPREQUEST']._serialized_end=512 - _globals['_MAPREQUEST_REQUEST']._serialized_start=248 - _globals['_MAPREQUEST_REQUEST']._serialized_end=487 - _globals['_MAPREQUEST_REQUEST_HEADERSENTRY']._serialized_start=441 - _globals['_MAPREQUEST_REQUEST_HEADERSENTRY']._serialized_end=487 - _globals['_HANDSHAKE']._serialized_start=514 - _globals['_HANDSHAKE']._serialized_end=538 - _globals['_TRANSMISSIONSTATUS']._serialized_start=540 - _globals['_TRANSMISSIONSTATUS']._serialized_end=573 - _globals['_MAPRESPONSE']._serialized_start=576 - _globals['_MAPRESPONSE']._serialized_end=816 - _globals['_MAPRESPONSE_RESULT']._serialized_start=740 - _globals['_MAPRESPONSE_RESULT']._serialized_end=791 - _globals['_READYRESPONSE']._serialized_start=818 - _globals['_READYRESPONSE']._serialized_end=848 - _globals['_MAP']._serialized_start=850 - _globals['_MAP']._serialized_end=967 + _globals['_MAPREQUEST']._serialized_start=108 + _globals['_MAPREQUEST']._serialized_end=536 + _globals['_MAPREQUEST_REQUEST']._serialized_start=272 + _globals['_MAPREQUEST_REQUEST']._serialized_end=511 + _globals['_MAPREQUEST_REQUEST_HEADERSENTRY']._serialized_start=465 + _globals['_MAPREQUEST_REQUEST_HEADERSENTRY']._serialized_end=511 + _globals['_HANDSHAKE']._serialized_start=538 + _globals['_HANDSHAKE']._serialized_end=562 + _globals['_TRANSMISSIONSTATUS']._serialized_start=564 + _globals['_TRANSMISSIONSTATUS']._serialized_end=597 + _globals['_MAPRESPONSE']._serialized_start=600 + _globals['_MAPRESPONSE']._serialized_end=840 + _globals['_MAPRESPONSE_RESULT']._serialized_start=764 + _globals['_MAPRESPONSE_RESULT']._serialized_end=815 + _globals['_READYRESPONSE']._serialized_start=842 + _globals['_READYRESPONSE']._serialized_end=872 + _globals['_MAP']._serialized_start=874 + _globals['_MAP']._serialized_end=991 # @@protoc_insertion_point(module_scope) diff --git a/pynumaflow/proto/mapper/map_pb2_grpc.py b/pynumaflow/proto/mapper/map_pb2_grpc.py index 446aaa9f..7325b904 100644 --- a/pynumaflow/proto/mapper/map_pb2_grpc.py +++ b/pynumaflow/proto/mapper/map_pb2_grpc.py @@ -4,7 +4,7 @@ import warnings from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -from . import map_pb2 as map__pb2 +from pynumaflow.proto.mapper import map_pb2 as pynumaflow_dot_proto_dot_mapper_dot_map__pb2 GRPC_GENERATED_VERSION = '1.75.0' GRPC_VERSION = grpc.__version__ @@ -19,7 +19,7 @@ if _version_not_supported: raise RuntimeError( f'The grpc package installed is at version {GRPC_VERSION},' - + f' but the generated code in map_pb2_grpc.py depends on' + + f' but the generated code in pynumaflow/proto/mapper/map_pb2_grpc.py depends on' + f' grpcio>={GRPC_GENERATED_VERSION}.' + f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}' + f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.' @@ -37,13 +37,13 @@ def __init__(self, channel): """ self.MapFn = channel.stream_stream( '/map.v1.Map/MapFn', - request_serializer=map__pb2.MapRequest.SerializeToString, - response_deserializer=map__pb2.MapResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_mapper_dot_map__pb2.MapRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_mapper_dot_map__pb2.MapResponse.FromString, _registered_method=True) self.IsReady = channel.unary_unary( '/map.v1.Map/IsReady', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=map__pb2.ReadyResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_mapper_dot_map__pb2.ReadyResponse.FromString, _registered_method=True) @@ -69,13 +69,13 @@ def add_MapServicer_to_server(servicer, server): rpc_method_handlers = { 'MapFn': grpc.stream_stream_rpc_method_handler( servicer.MapFn, - request_deserializer=map__pb2.MapRequest.FromString, - response_serializer=map__pb2.MapResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_mapper_dot_map__pb2.MapRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_mapper_dot_map__pb2.MapResponse.SerializeToString, ), 'IsReady': grpc.unary_unary_rpc_method_handler( servicer.IsReady, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=map__pb2.ReadyResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_mapper_dot_map__pb2.ReadyResponse.SerializeToString, ), } generic_handler = grpc.method_handlers_generic_handler( @@ -103,8 +103,8 @@ def MapFn(request_iterator, request_iterator, target, '/map.v1.Map/MapFn', - map__pb2.MapRequest.SerializeToString, - map__pb2.MapResponse.FromString, + pynumaflow_dot_proto_dot_mapper_dot_map__pb2.MapRequest.SerializeToString, + pynumaflow_dot_proto_dot_mapper_dot_map__pb2.MapResponse.FromString, options, channel_credentials, insecure, @@ -131,7 +131,7 @@ def IsReady(request, target, '/map.v1.Map/IsReady', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - map__pb2.ReadyResponse.FromString, + pynumaflow_dot_proto_dot_mapper_dot_map__pb2.ReadyResponse.FromString, options, channel_credentials, insecure, diff --git a/pynumaflow/proto/reducer/reduce_pb2.py b/pynumaflow/proto/reducer/reduce_pb2.py index 38e711db..9ed9fe19 100644 --- a/pynumaflow/proto/reducer/reduce_pb2.py +++ b/pynumaflow/proto/reducer/reduce_pb2.py @@ -1,7 +1,7 @@ # -*- coding: utf-8 -*- # Generated by the protocol buffer compiler. DO NOT EDIT! # NO CHECKED-IN PROTOBUF GENCODE -# source: reduce.proto +# source: pynumaflow/proto/reducer/reduce.proto # Protobuf Python Version: 6.31.1 """Generated protocol buffer code.""" from google.protobuf import descriptor as _descriptor @@ -15,7 +15,7 @@ 31, 1, '', - 'reduce.proto' + 'pynumaflow/proto/reducer/reduce.proto' ) # @@protoc_insertion_point(imports) @@ -26,33 +26,33 @@ from google.protobuf import timestamp_pb2 as google_dot_protobuf_dot_timestamp__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x0creduce.proto\x12\treduce.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\x98\x04\n\rReduceRequest\x12\x31\n\x07payload\x18\x01 \x01(\x0b\x32 .reduce.v1.ReduceRequest.Payload\x12;\n\toperation\x18\x02 \x01(\x0b\x32(.reduce.v1.ReduceRequest.WindowOperation\x1a\x9e\x01\n\x0fWindowOperation\x12=\n\x05\x65vent\x18\x01 \x01(\x0e\x32..reduce.v1.ReduceRequest.WindowOperation.Event\x12\"\n\x07windows\x18\x02 \x03(\x0b\x32\x11.reduce.v1.Window\"(\n\x05\x45vent\x12\x08\n\x04OPEN\x10\x00\x12\t\n\x05\x43LOSE\x10\x01\x12\n\n\x06\x41PPEND\x10\x04\x1a\xf5\x01\n\x07Payload\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12>\n\x07headers\x18\x05 \x03(\x0b\x32-.reduce.v1.ReduceRequest.Payload.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"j\n\x06Window\x12)\n\x05start\x18\x01 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\'\n\x03\x65nd\x18\x02 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04slot\x18\x03 \x01(\t\"\xa7\x01\n\x0eReduceResponse\x12\x30\n\x06result\x18\x01 \x01(\x0b\x32 .reduce.v1.ReduceResponse.Result\x12!\n\x06window\x18\x02 \x01(\x0b\x32\x11.reduce.v1.Window\x12\x0b\n\x03\x45OF\x18\x03 \x01(\x08\x1a\x33\n\x06Result\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12\x0c\n\x04tags\x18\x03 \x03(\t\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\x8a\x01\n\x06Reduce\x12\x43\n\x08ReduceFn\x12\x18.reduce.v1.ReduceRequest\x1a\x19.reduce.v1.ReduceResponse(\x01\x30\x01\x12;\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x18.reduce.v1.ReadyResponseb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n%pynumaflow/proto/reducer/reduce.proto\x12\treduce.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\x98\x04\n\rReduceRequest\x12\x31\n\x07payload\x18\x01 \x01(\x0b\x32 .reduce.v1.ReduceRequest.Payload\x12;\n\toperation\x18\x02 \x01(\x0b\x32(.reduce.v1.ReduceRequest.WindowOperation\x1a\x9e\x01\n\x0fWindowOperation\x12=\n\x05\x65vent\x18\x01 \x01(\x0e\x32..reduce.v1.ReduceRequest.WindowOperation.Event\x12\"\n\x07windows\x18\x02 \x03(\x0b\x32\x11.reduce.v1.Window\"(\n\x05\x45vent\x12\x08\n\x04OPEN\x10\x00\x12\t\n\x05\x43LOSE\x10\x01\x12\n\n\x06\x41PPEND\x10\x04\x1a\xf5\x01\n\x07Payload\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12>\n\x07headers\x18\x05 \x03(\x0b\x32-.reduce.v1.ReduceRequest.Payload.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\"j\n\x06Window\x12)\n\x05start\x18\x01 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\'\n\x03\x65nd\x18\x02 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04slot\x18\x03 \x01(\t\"\xa7\x01\n\x0eReduceResponse\x12\x30\n\x06result\x18\x01 \x01(\x0b\x32 .reduce.v1.ReduceResponse.Result\x12!\n\x06window\x18\x02 \x01(\x0b\x32\x11.reduce.v1.Window\x12\x0b\n\x03\x45OF\x18\x03 \x01(\x08\x1a\x33\n\x06Result\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12\x0c\n\x04tags\x18\x03 \x03(\t\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\x8a\x01\n\x06Reduce\x12\x43\n\x08ReduceFn\x12\x18.reduce.v1.ReduceRequest\x1a\x19.reduce.v1.ReduceResponse(\x01\x30\x01\x12;\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x18.reduce.v1.ReadyResponseb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) -_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'reduce_pb2', _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'pynumaflow.proto.reducer.reduce_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: DESCRIPTOR._loaded_options = None _globals['_REDUCEREQUEST_PAYLOAD_HEADERSENTRY']._loaded_options = None _globals['_REDUCEREQUEST_PAYLOAD_HEADERSENTRY']._serialized_options = b'8\001' - _globals['_REDUCEREQUEST']._serialized_start=90 - _globals['_REDUCEREQUEST']._serialized_end=626 - _globals['_REDUCEREQUEST_WINDOWOPERATION']._serialized_start=220 - _globals['_REDUCEREQUEST_WINDOWOPERATION']._serialized_end=378 - _globals['_REDUCEREQUEST_WINDOWOPERATION_EVENT']._serialized_start=338 - _globals['_REDUCEREQUEST_WINDOWOPERATION_EVENT']._serialized_end=378 - _globals['_REDUCEREQUEST_PAYLOAD']._serialized_start=381 - _globals['_REDUCEREQUEST_PAYLOAD']._serialized_end=626 - _globals['_REDUCEREQUEST_PAYLOAD_HEADERSENTRY']._serialized_start=580 - _globals['_REDUCEREQUEST_PAYLOAD_HEADERSENTRY']._serialized_end=626 - _globals['_WINDOW']._serialized_start=628 - _globals['_WINDOW']._serialized_end=734 - _globals['_REDUCERESPONSE']._serialized_start=737 - _globals['_REDUCERESPONSE']._serialized_end=904 - _globals['_REDUCERESPONSE_RESULT']._serialized_start=853 - _globals['_REDUCERESPONSE_RESULT']._serialized_end=904 - _globals['_READYRESPONSE']._serialized_start=906 - _globals['_READYRESPONSE']._serialized_end=936 - _globals['_REDUCE']._serialized_start=939 - _globals['_REDUCE']._serialized_end=1077 + _globals['_REDUCEREQUEST']._serialized_start=115 + _globals['_REDUCEREQUEST']._serialized_end=651 + _globals['_REDUCEREQUEST_WINDOWOPERATION']._serialized_start=245 + _globals['_REDUCEREQUEST_WINDOWOPERATION']._serialized_end=403 + _globals['_REDUCEREQUEST_WINDOWOPERATION_EVENT']._serialized_start=363 + _globals['_REDUCEREQUEST_WINDOWOPERATION_EVENT']._serialized_end=403 + _globals['_REDUCEREQUEST_PAYLOAD']._serialized_start=406 + _globals['_REDUCEREQUEST_PAYLOAD']._serialized_end=651 + _globals['_REDUCEREQUEST_PAYLOAD_HEADERSENTRY']._serialized_start=605 + _globals['_REDUCEREQUEST_PAYLOAD_HEADERSENTRY']._serialized_end=651 + _globals['_WINDOW']._serialized_start=653 + _globals['_WINDOW']._serialized_end=759 + _globals['_REDUCERESPONSE']._serialized_start=762 + _globals['_REDUCERESPONSE']._serialized_end=929 + _globals['_REDUCERESPONSE_RESULT']._serialized_start=878 + _globals['_REDUCERESPONSE_RESULT']._serialized_end=929 + _globals['_READYRESPONSE']._serialized_start=931 + _globals['_READYRESPONSE']._serialized_end=961 + _globals['_REDUCE']._serialized_start=964 + _globals['_REDUCE']._serialized_end=1102 # @@protoc_insertion_point(module_scope) diff --git a/pynumaflow/proto/reducer/reduce_pb2_grpc.py b/pynumaflow/proto/reducer/reduce_pb2_grpc.py index 7795282c..4a8dd390 100644 --- a/pynumaflow/proto/reducer/reduce_pb2_grpc.py +++ b/pynumaflow/proto/reducer/reduce_pb2_grpc.py @@ -4,7 +4,7 @@ import warnings from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -from . import reduce_pb2 as reduce__pb2 +from pynumaflow.proto.reducer import reduce_pb2 as pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2 GRPC_GENERATED_VERSION = '1.75.0' GRPC_VERSION = grpc.__version__ @@ -19,7 +19,7 @@ if _version_not_supported: raise RuntimeError( f'The grpc package installed is at version {GRPC_VERSION},' - + f' but the generated code in reduce_pb2_grpc.py depends on' + + f' but the generated code in pynumaflow/proto/reducer/reduce_pb2_grpc.py depends on' + f' grpcio>={GRPC_GENERATED_VERSION}.' + f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}' + f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.' @@ -37,13 +37,13 @@ def __init__(self, channel): """ self.ReduceFn = channel.stream_stream( '/reduce.v1.Reduce/ReduceFn', - request_serializer=reduce__pb2.ReduceRequest.SerializeToString, - response_deserializer=reduce__pb2.ReduceResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReduceRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReduceResponse.FromString, _registered_method=True) self.IsReady = channel.unary_unary( '/reduce.v1.Reduce/IsReady', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=reduce__pb2.ReadyResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReadyResponse.FromString, _registered_method=True) @@ -69,13 +69,13 @@ def add_ReduceServicer_to_server(servicer, server): rpc_method_handlers = { 'ReduceFn': grpc.stream_stream_rpc_method_handler( servicer.ReduceFn, - request_deserializer=reduce__pb2.ReduceRequest.FromString, - response_serializer=reduce__pb2.ReduceResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReduceRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReduceResponse.SerializeToString, ), 'IsReady': grpc.unary_unary_rpc_method_handler( servicer.IsReady, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=reduce__pb2.ReadyResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReadyResponse.SerializeToString, ), } generic_handler = grpc.method_handlers_generic_handler( @@ -103,8 +103,8 @@ def ReduceFn(request_iterator, request_iterator, target, '/reduce.v1.Reduce/ReduceFn', - reduce__pb2.ReduceRequest.SerializeToString, - reduce__pb2.ReduceResponse.FromString, + pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReduceRequest.SerializeToString, + pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReduceResponse.FromString, options, channel_credentials, insecure, @@ -131,7 +131,7 @@ def IsReady(request, target, '/reduce.v1.Reduce/IsReady', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - reduce__pb2.ReadyResponse.FromString, + pynumaflow_dot_proto_dot_reducer_dot_reduce__pb2.ReadyResponse.FromString, options, channel_credentials, insecure, diff --git a/pynumaflow/proto/sideinput/sideinput_pb2.py b/pynumaflow/proto/sideinput/sideinput_pb2.py index 47f2bbdd..af68192d 100644 --- a/pynumaflow/proto/sideinput/sideinput_pb2.py +++ b/pynumaflow/proto/sideinput/sideinput_pb2.py @@ -1,7 +1,7 @@ # -*- coding: utf-8 -*- # Generated by the protocol buffer compiler. DO NOT EDIT! # NO CHECKED-IN PROTOBUF GENCODE -# source: sideinput.proto +# source: pynumaflow/proto/sideinput/sideinput.proto # Protobuf Python Version: 6.31.1 """Generated protocol buffer code.""" from google.protobuf import descriptor as _descriptor @@ -15,7 +15,7 @@ 31, 1, '', - 'sideinput.proto' + 'pynumaflow/proto/sideinput/sideinput.proto' ) # @@protoc_insertion_point(imports) @@ -25,17 +25,17 @@ from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x0fsideinput.proto\x12\x0csideinput.v1\x1a\x1bgoogle/protobuf/empty.proto\"8\n\x11SideInputResponse\x12\r\n\x05value\x18\x01 \x01(\x0c\x12\x14\n\x0cno_broadcast\x18\x02 \x01(\x08\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\x99\x01\n\tSideInput\x12L\n\x11RetrieveSideInput\x12\x16.google.protobuf.Empty\x1a\x1f.sideinput.v1.SideInputResponse\x12>\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x1b.sideinput.v1.ReadyResponseb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n*pynumaflow/proto/sideinput/sideinput.proto\x12\x0csideinput.v1\x1a\x1bgoogle/protobuf/empty.proto\"8\n\x11SideInputResponse\x12\r\n\x05value\x18\x01 \x01(\x0c\x12\x14\n\x0cno_broadcast\x18\x02 \x01(\x08\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\x99\x01\n\tSideInput\x12L\n\x11RetrieveSideInput\x12\x16.google.protobuf.Empty\x1a\x1f.sideinput.v1.SideInputResponse\x12>\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x1b.sideinput.v1.ReadyResponseb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) -_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'sideinput_pb2', _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'pynumaflow.proto.sideinput.sideinput_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: DESCRIPTOR._loaded_options = None - _globals['_SIDEINPUTRESPONSE']._serialized_start=62 - _globals['_SIDEINPUTRESPONSE']._serialized_end=118 - _globals['_READYRESPONSE']._serialized_start=120 - _globals['_READYRESPONSE']._serialized_end=150 - _globals['_SIDEINPUT']._serialized_start=153 - _globals['_SIDEINPUT']._serialized_end=306 + _globals['_SIDEINPUTRESPONSE']._serialized_start=89 + _globals['_SIDEINPUTRESPONSE']._serialized_end=145 + _globals['_READYRESPONSE']._serialized_start=147 + _globals['_READYRESPONSE']._serialized_end=177 + _globals['_SIDEINPUT']._serialized_start=180 + _globals['_SIDEINPUT']._serialized_end=333 # @@protoc_insertion_point(module_scope) diff --git a/pynumaflow/proto/sideinput/sideinput_pb2_grpc.py b/pynumaflow/proto/sideinput/sideinput_pb2_grpc.py index b90412c4..2f1c8203 100644 --- a/pynumaflow/proto/sideinput/sideinput_pb2_grpc.py +++ b/pynumaflow/proto/sideinput/sideinput_pb2_grpc.py @@ -4,7 +4,7 @@ import warnings from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -from . import sideinput_pb2 as sideinput__pb2 +from pynumaflow.proto.sideinput import sideinput_pb2 as pynumaflow_dot_proto_dot_sideinput_dot_sideinput__pb2 GRPC_GENERATED_VERSION = '1.75.0' GRPC_VERSION = grpc.__version__ @@ -19,7 +19,7 @@ if _version_not_supported: raise RuntimeError( f'The grpc package installed is at version {GRPC_VERSION},' - + f' but the generated code in sideinput_pb2_grpc.py depends on' + + f' but the generated code in pynumaflow/proto/sideinput/sideinput_pb2_grpc.py depends on' + f' grpcio>={GRPC_GENERATED_VERSION}.' + f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}' + f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.' @@ -46,12 +46,12 @@ def __init__(self, channel): self.RetrieveSideInput = channel.unary_unary( '/sideinput.v1.SideInput/RetrieveSideInput', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=sideinput__pb2.SideInputResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_sideinput_dot_sideinput__pb2.SideInputResponse.FromString, _registered_method=True) self.IsReady = channel.unary_unary( '/sideinput.v1.SideInput/IsReady', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=sideinput__pb2.ReadyResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_sideinput_dot_sideinput__pb2.ReadyResponse.FromString, _registered_method=True) @@ -86,12 +86,12 @@ def add_SideInputServicer_to_server(servicer, server): 'RetrieveSideInput': grpc.unary_unary_rpc_method_handler( servicer.RetrieveSideInput, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=sideinput__pb2.SideInputResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_sideinput_dot_sideinput__pb2.SideInputResponse.SerializeToString, ), 'IsReady': grpc.unary_unary_rpc_method_handler( servicer.IsReady, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=sideinput__pb2.ReadyResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_sideinput_dot_sideinput__pb2.ReadyResponse.SerializeToString, ), } generic_handler = grpc.method_handlers_generic_handler( @@ -128,7 +128,7 @@ def RetrieveSideInput(request, target, '/sideinput.v1.SideInput/RetrieveSideInput', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - sideinput__pb2.SideInputResponse.FromString, + pynumaflow_dot_proto_dot_sideinput_dot_sideinput__pb2.SideInputResponse.FromString, options, channel_credentials, insecure, @@ -155,7 +155,7 @@ def IsReady(request, target, '/sideinput.v1.SideInput/IsReady', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - sideinput__pb2.ReadyResponse.FromString, + pynumaflow_dot_proto_dot_sideinput_dot_sideinput__pb2.ReadyResponse.FromString, options, channel_credentials, insecure, diff --git a/pynumaflow/proto/sinker/sink_pb2.py b/pynumaflow/proto/sinker/sink_pb2.py index c64f77c7..10462d7c 100644 --- a/pynumaflow/proto/sinker/sink_pb2.py +++ b/pynumaflow/proto/sinker/sink_pb2.py @@ -1,7 +1,7 @@ # -*- coding: utf-8 -*- # Generated by the protocol buffer compiler. DO NOT EDIT! # NO CHECKED-IN PROTOBUF GENCODE -# source: sink.proto +# source: pynumaflow/proto/sinker/sink.proto # Protobuf Python Version: 6.31.1 """Generated protocol buffer code.""" from google.protobuf import descriptor as _descriptor @@ -15,7 +15,7 @@ 31, 1, '', - 'sink.proto' + 'pynumaflow/proto/sinker/sink.proto' ) # @@protoc_insertion_point(imports) @@ -26,33 +26,33 @@ from google.protobuf import timestamp_pb2 as google_dot_protobuf_dot_timestamp__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\nsink.proto\x12\x07sink.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xa3\x03\n\x0bSinkRequest\x12-\n\x07request\x18\x01 \x01(\x0b\x32\x1c.sink.v1.SinkRequest.Request\x12+\n\x06status\x18\x02 \x01(\x0b\x32\x1b.sink.v1.TransmissionStatus\x12*\n\thandshake\x18\x03 \x01(\x0b\x32\x12.sink.v1.HandshakeH\x00\x88\x01\x01\x1a\xfd\x01\n\x07Request\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\n\n\x02id\x18\x05 \x01(\t\x12:\n\x07headers\x18\x06 \x03(\x0b\x32).sink.v1.SinkRequest.Request.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x0c\n\n_handshake\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\"!\n\x12TransmissionStatus\x12\x0b\n\x03\x65ot\x18\x01 \x01(\x08\"\xfc\x01\n\x0cSinkResponse\x12-\n\x07results\x18\x01 \x03(\x0b\x32\x1c.sink.v1.SinkResponse.Result\x12*\n\thandshake\x18\x02 \x01(\x0b\x32\x12.sink.v1.HandshakeH\x00\x88\x01\x01\x12\x30\n\x06status\x18\x03 \x01(\x0b\x32\x1b.sink.v1.TransmissionStatusH\x01\x88\x01\x01\x1a\x46\n\x06Result\x12\n\n\x02id\x18\x01 \x01(\t\x12\x1f\n\x06status\x18\x02 \x01(\x0e\x32\x0f.sink.v1.Status\x12\x0f\n\x07\x65rr_msg\x18\x03 \x01(\tB\x0c\n\n_handshakeB\t\n\x07_status*0\n\x06Status\x12\x0b\n\x07SUCCESS\x10\x00\x12\x0b\n\x07\x46\x41ILURE\x10\x01\x12\x0c\n\x08\x46\x41LLBACK\x10\x02\x32|\n\x04Sink\x12\x39\n\x06SinkFn\x12\x14.sink.v1.SinkRequest\x1a\x15.sink.v1.SinkResponse(\x01\x30\x01\x12\x39\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x16.sink.v1.ReadyResponseb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\"pynumaflow/proto/sinker/sink.proto\x12\x07sink.v1\x1a\x1bgoogle/protobuf/empty.proto\x1a\x1fgoogle/protobuf/timestamp.proto\"\xa3\x03\n\x0bSinkRequest\x12-\n\x07request\x18\x01 \x01(\x0b\x32\x1c.sink.v1.SinkRequest.Request\x12+\n\x06status\x18\x02 \x01(\x0b\x32\x1b.sink.v1.TransmissionStatus\x12*\n\thandshake\x18\x03 \x01(\x0b\x32\x12.sink.v1.HandshakeH\x00\x88\x01\x01\x1a\xfd\x01\n\x07Request\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\n\n\x02id\x18\x05 \x01(\t\x12:\n\x07headers\x18\x06 \x03(\x0b\x32).sink.v1.SinkRequest.Request.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x0c\n\n_handshake\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\"!\n\x12TransmissionStatus\x12\x0b\n\x03\x65ot\x18\x01 \x01(\x08\"\xfc\x01\n\x0cSinkResponse\x12-\n\x07results\x18\x01 \x03(\x0b\x32\x1c.sink.v1.SinkResponse.Result\x12*\n\thandshake\x18\x02 \x01(\x0b\x32\x12.sink.v1.HandshakeH\x00\x88\x01\x01\x12\x30\n\x06status\x18\x03 \x01(\x0b\x32\x1b.sink.v1.TransmissionStatusH\x01\x88\x01\x01\x1a\x46\n\x06Result\x12\n\n\x02id\x18\x01 \x01(\t\x12\x1f\n\x06status\x18\x02 \x01(\x0e\x32\x0f.sink.v1.Status\x12\x0f\n\x07\x65rr_msg\x18\x03 \x01(\tB\x0c\n\n_handshakeB\t\n\x07_status*0\n\x06Status\x12\x0b\n\x07SUCCESS\x10\x00\x12\x0b\n\x07\x46\x41ILURE\x10\x01\x12\x0c\n\x08\x46\x41LLBACK\x10\x02\x32|\n\x04Sink\x12\x39\n\x06SinkFn\x12\x14.sink.v1.SinkRequest\x1a\x15.sink.v1.SinkResponse(\x01\x30\x01\x12\x39\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x16.sink.v1.ReadyResponseb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) -_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'sink_pb2', _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'pynumaflow.proto.sinker.sink_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: DESCRIPTOR._loaded_options = None _globals['_SINKREQUEST_REQUEST_HEADERSENTRY']._loaded_options = None _globals['_SINKREQUEST_REQUEST_HEADERSENTRY']._serialized_options = b'8\001' - _globals['_STATUS']._serialized_start=855 - _globals['_STATUS']._serialized_end=903 - _globals['_SINKREQUEST']._serialized_start=86 - _globals['_SINKREQUEST']._serialized_end=505 - _globals['_SINKREQUEST_REQUEST']._serialized_start=238 - _globals['_SINKREQUEST_REQUEST']._serialized_end=491 - _globals['_SINKREQUEST_REQUEST_HEADERSENTRY']._serialized_start=445 - _globals['_SINKREQUEST_REQUEST_HEADERSENTRY']._serialized_end=491 - _globals['_HANDSHAKE']._serialized_start=507 - _globals['_HANDSHAKE']._serialized_end=531 - _globals['_READYRESPONSE']._serialized_start=533 - _globals['_READYRESPONSE']._serialized_end=563 - _globals['_TRANSMISSIONSTATUS']._serialized_start=565 - _globals['_TRANSMISSIONSTATUS']._serialized_end=598 - _globals['_SINKRESPONSE']._serialized_start=601 - _globals['_SINKRESPONSE']._serialized_end=853 - _globals['_SINKRESPONSE_RESULT']._serialized_start=758 - _globals['_SINKRESPONSE_RESULT']._serialized_end=828 - _globals['_SINK']._serialized_start=905 - _globals['_SINK']._serialized_end=1029 + _globals['_STATUS']._serialized_start=879 + _globals['_STATUS']._serialized_end=927 + _globals['_SINKREQUEST']._serialized_start=110 + _globals['_SINKREQUEST']._serialized_end=529 + _globals['_SINKREQUEST_REQUEST']._serialized_start=262 + _globals['_SINKREQUEST_REQUEST']._serialized_end=515 + _globals['_SINKREQUEST_REQUEST_HEADERSENTRY']._serialized_start=469 + _globals['_SINKREQUEST_REQUEST_HEADERSENTRY']._serialized_end=515 + _globals['_HANDSHAKE']._serialized_start=531 + _globals['_HANDSHAKE']._serialized_end=555 + _globals['_READYRESPONSE']._serialized_start=557 + _globals['_READYRESPONSE']._serialized_end=587 + _globals['_TRANSMISSIONSTATUS']._serialized_start=589 + _globals['_TRANSMISSIONSTATUS']._serialized_end=622 + _globals['_SINKRESPONSE']._serialized_start=625 + _globals['_SINKRESPONSE']._serialized_end=877 + _globals['_SINKRESPONSE_RESULT']._serialized_start=782 + _globals['_SINKRESPONSE_RESULT']._serialized_end=852 + _globals['_SINK']._serialized_start=929 + _globals['_SINK']._serialized_end=1053 # @@protoc_insertion_point(module_scope) diff --git a/pynumaflow/proto/sinker/sink_pb2_grpc.py b/pynumaflow/proto/sinker/sink_pb2_grpc.py index 33d09e2e..cecd2790 100644 --- a/pynumaflow/proto/sinker/sink_pb2_grpc.py +++ b/pynumaflow/proto/sinker/sink_pb2_grpc.py @@ -4,7 +4,7 @@ import warnings from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -from . import sink_pb2 as sink__pb2 +from pynumaflow.proto.sinker import sink_pb2 as pynumaflow_dot_proto_dot_sinker_dot_sink__pb2 GRPC_GENERATED_VERSION = '1.75.0' GRPC_VERSION = grpc.__version__ @@ -19,7 +19,7 @@ if _version_not_supported: raise RuntimeError( f'The grpc package installed is at version {GRPC_VERSION},' - + f' but the generated code in sink_pb2_grpc.py depends on' + + f' but the generated code in pynumaflow/proto/sinker/sink_pb2_grpc.py depends on' + f' grpcio>={GRPC_GENERATED_VERSION}.' + f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}' + f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.' @@ -37,13 +37,13 @@ def __init__(self, channel): """ self.SinkFn = channel.stream_stream( '/sink.v1.Sink/SinkFn', - request_serializer=sink__pb2.SinkRequest.SerializeToString, - response_deserializer=sink__pb2.SinkResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.SinkRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.SinkResponse.FromString, _registered_method=True) self.IsReady = channel.unary_unary( '/sink.v1.Sink/IsReady', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=sink__pb2.ReadyResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.ReadyResponse.FromString, _registered_method=True) @@ -69,13 +69,13 @@ def add_SinkServicer_to_server(servicer, server): rpc_method_handlers = { 'SinkFn': grpc.stream_stream_rpc_method_handler( servicer.SinkFn, - request_deserializer=sink__pb2.SinkRequest.FromString, - response_serializer=sink__pb2.SinkResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.SinkRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.SinkResponse.SerializeToString, ), 'IsReady': grpc.unary_unary_rpc_method_handler( servicer.IsReady, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=sink__pb2.ReadyResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.ReadyResponse.SerializeToString, ), } generic_handler = grpc.method_handlers_generic_handler( @@ -103,8 +103,8 @@ def SinkFn(request_iterator, request_iterator, target, '/sink.v1.Sink/SinkFn', - sink__pb2.SinkRequest.SerializeToString, - sink__pb2.SinkResponse.FromString, + pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.SinkRequest.SerializeToString, + pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.SinkResponse.FromString, options, channel_credentials, insecure, @@ -131,7 +131,7 @@ def IsReady(request, target, '/sink.v1.Sink/IsReady', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - sink__pb2.ReadyResponse.FromString, + pynumaflow_dot_proto_dot_sinker_dot_sink__pb2.ReadyResponse.FromString, options, channel_credentials, insecure, diff --git a/pynumaflow/proto/sourcer/source_pb2.py b/pynumaflow/proto/sourcer/source_pb2.py index f9645ce8..f85d827f 100644 --- a/pynumaflow/proto/sourcer/source_pb2.py +++ b/pynumaflow/proto/sourcer/source_pb2.py @@ -1,7 +1,7 @@ # -*- coding: utf-8 -*- # Generated by the protocol buffer compiler. DO NOT EDIT! # NO CHECKED-IN PROTOBUF GENCODE -# source: source.proto +# source: pynumaflow/proto/sourcer/source.proto # Protobuf Python Version: 6.31.1 """Generated protocol buffer code.""" from google.protobuf import descriptor as _descriptor @@ -15,7 +15,7 @@ 31, 1, '', - 'source.proto' + 'pynumaflow/proto/sourcer/source.proto' ) # @@protoc_insertion_point(imports) @@ -26,61 +26,61 @@ from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x0csource.proto\x12\tsource.v1\x1a\x1fgoogle/protobuf/timestamp.proto\x1a\x1bgoogle/protobuf/empty.proto\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"\xb1\x01\n\x0bReadRequest\x12/\n\x07request\x18\x01 \x01(\x0b\x32\x1e.source.v1.ReadRequest.Request\x12,\n\thandshake\x18\x02 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a\x35\n\x07Request\x12\x13\n\x0bnum_records\x18\x01 \x01(\x04\x12\x15\n\rtimeout_in_ms\x18\x02 \x01(\rB\x0c\n\n_handshake\"\x81\x05\n\x0cReadResponse\x12.\n\x06result\x18\x01 \x01(\x0b\x32\x1e.source.v1.ReadResponse.Result\x12.\n\x06status\x18\x02 \x01(\x0b\x32\x1e.source.v1.ReadResponse.Status\x12,\n\thandshake\x18\x03 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a\xe8\x01\n\x06Result\x12\x0f\n\x07payload\x18\x01 \x01(\x0c\x12!\n\x06offset\x18\x02 \x01(\x0b\x32\x11.source.v1.Offset\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04keys\x18\x04 \x03(\t\x12<\n\x07headers\x18\x05 \x03(\x0b\x32+.source.v1.ReadResponse.Result.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x1a\xe9\x01\n\x06Status\x12\x0b\n\x03\x65ot\x18\x01 \x01(\x08\x12\x31\n\x04\x63ode\x18\x02 \x01(\x0e\x32#.source.v1.ReadResponse.Status.Code\x12\x38\n\x05\x65rror\x18\x03 \x01(\x0e\x32$.source.v1.ReadResponse.Status.ErrorH\x00\x88\x01\x01\x12\x10\n\x03msg\x18\x04 \x01(\tH\x01\x88\x01\x01\" \n\x04\x43ode\x12\x0b\n\x07SUCCESS\x10\x00\x12\x0b\n\x07\x46\x41ILURE\x10\x01\"\x1f\n\x05\x45rror\x12\x0b\n\x07UNACKED\x10\x00\x12\t\n\x05OTHER\x10\x01\x42\x08\n\x06_errorB\x06\n\x04_msgB\x0c\n\n_handshake\"\xa7\x01\n\nAckRequest\x12.\n\x07request\x18\x01 \x01(\x0b\x32\x1d.source.v1.AckRequest.Request\x12,\n\thandshake\x18\x02 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a-\n\x07Request\x12\"\n\x07offsets\x18\x01 \x03(\x0b\x32\x11.source.v1.OffsetB\x0c\n\n_handshake\"\xab\x01\n\x0b\x41\x63kResponse\x12-\n\x06result\x18\x01 \x01(\x0b\x32\x1d.source.v1.AckResponse.Result\x12,\n\thandshake\x18\x02 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a\x31\n\x06Result\x12\'\n\x07success\x18\x01 \x01(\x0b\x32\x16.google.protobuf.EmptyB\x0c\n\n_handshake\"m\n\x0bNackRequest\x12/\n\x07request\x18\x01 \x01(\x0b\x32\x1e.source.v1.NackRequest.Request\x1a-\n\x07Request\x12\"\n\x07offsets\x18\x01 \x03(\x0b\x32\x11.source.v1.Offset\"q\n\x0cNackResponse\x12.\n\x06result\x18\x01 \x01(\x0b\x32\x1e.source.v1.NackResponse.Result\x1a\x31\n\x06Result\x12\'\n\x07success\x18\x01 \x01(\x0b\x32\x16.google.protobuf.Empty\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\"]\n\x0fPendingResponse\x12\x31\n\x06result\x18\x01 \x01(\x0b\x32!.source.v1.PendingResponse.Result\x1a\x17\n\x06Result\x12\r\n\x05\x63ount\x18\x01 \x01(\x03\"h\n\x12PartitionsResponse\x12\x34\n\x06result\x18\x01 \x01(\x0b\x32$.source.v1.PartitionsResponse.Result\x1a\x1c\n\x06Result\x12\x12\n\npartitions\x18\x01 \x03(\x05\".\n\x06Offset\x12\x0e\n\x06offset\x18\x01 \x01(\x0c\x12\x14\n\x0cpartition_id\x18\x02 \x01(\x05\x32\x83\x03\n\x06Source\x12=\n\x06ReadFn\x12\x16.source.v1.ReadRequest\x1a\x17.source.v1.ReadResponse(\x01\x30\x01\x12:\n\x05\x41\x63kFn\x12\x15.source.v1.AckRequest\x1a\x16.source.v1.AckResponse(\x01\x30\x01\x12\x39\n\x06NackFn\x12\x16.source.v1.NackRequest\x1a\x17.source.v1.NackResponse\x12?\n\tPendingFn\x12\x16.google.protobuf.Empty\x1a\x1a.source.v1.PendingResponse\x12\x45\n\x0cPartitionsFn\x12\x16.google.protobuf.Empty\x1a\x1d.source.v1.PartitionsResponse\x12;\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x18.source.v1.ReadyResponseb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n%pynumaflow/proto/sourcer/source.proto\x12\tsource.v1\x1a\x1fgoogle/protobuf/timestamp.proto\x1a\x1bgoogle/protobuf/empty.proto\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"\xb1\x01\n\x0bReadRequest\x12/\n\x07request\x18\x01 \x01(\x0b\x32\x1e.source.v1.ReadRequest.Request\x12,\n\thandshake\x18\x02 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a\x35\n\x07Request\x12\x13\n\x0bnum_records\x18\x01 \x01(\x04\x12\x15\n\rtimeout_in_ms\x18\x02 \x01(\rB\x0c\n\n_handshake\"\x81\x05\n\x0cReadResponse\x12.\n\x06result\x18\x01 \x01(\x0b\x32\x1e.source.v1.ReadResponse.Result\x12.\n\x06status\x18\x02 \x01(\x0b\x32\x1e.source.v1.ReadResponse.Status\x12,\n\thandshake\x18\x03 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a\xe8\x01\n\x06Result\x12\x0f\n\x07payload\x18\x01 \x01(\x0c\x12!\n\x06offset\x18\x02 \x01(\x0b\x32\x11.source.v1.Offset\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04keys\x18\x04 \x03(\t\x12<\n\x07headers\x18\x05 \x03(\x0b\x32+.source.v1.ReadResponse.Result.HeadersEntry\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x1a\xe9\x01\n\x06Status\x12\x0b\n\x03\x65ot\x18\x01 \x01(\x08\x12\x31\n\x04\x63ode\x18\x02 \x01(\x0e\x32#.source.v1.ReadResponse.Status.Code\x12\x38\n\x05\x65rror\x18\x03 \x01(\x0e\x32$.source.v1.ReadResponse.Status.ErrorH\x00\x88\x01\x01\x12\x10\n\x03msg\x18\x04 \x01(\tH\x01\x88\x01\x01\" \n\x04\x43ode\x12\x0b\n\x07SUCCESS\x10\x00\x12\x0b\n\x07\x46\x41ILURE\x10\x01\"\x1f\n\x05\x45rror\x12\x0b\n\x07UNACKED\x10\x00\x12\t\n\x05OTHER\x10\x01\x42\x08\n\x06_errorB\x06\n\x04_msgB\x0c\n\n_handshake\"\xa7\x01\n\nAckRequest\x12.\n\x07request\x18\x01 \x01(\x0b\x32\x1d.source.v1.AckRequest.Request\x12,\n\thandshake\x18\x02 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a-\n\x07Request\x12\"\n\x07offsets\x18\x01 \x03(\x0b\x32\x11.source.v1.OffsetB\x0c\n\n_handshake\"\xab\x01\n\x0b\x41\x63kResponse\x12-\n\x06result\x18\x01 \x01(\x0b\x32\x1d.source.v1.AckResponse.Result\x12,\n\thandshake\x18\x02 \x01(\x0b\x32\x14.source.v1.HandshakeH\x00\x88\x01\x01\x1a\x31\n\x06Result\x12\'\n\x07success\x18\x01 \x01(\x0b\x32\x16.google.protobuf.EmptyB\x0c\n\n_handshake\"m\n\x0bNackRequest\x12/\n\x07request\x18\x01 \x01(\x0b\x32\x1e.source.v1.NackRequest.Request\x1a-\n\x07Request\x12\"\n\x07offsets\x18\x01 \x03(\x0b\x32\x11.source.v1.Offset\"q\n\x0cNackResponse\x12.\n\x06result\x18\x01 \x01(\x0b\x32\x1e.source.v1.NackResponse.Result\x1a\x31\n\x06Result\x12\'\n\x07success\x18\x01 \x01(\x0b\x32\x16.google.protobuf.Empty\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\"]\n\x0fPendingResponse\x12\x31\n\x06result\x18\x01 \x01(\x0b\x32!.source.v1.PendingResponse.Result\x1a\x17\n\x06Result\x12\r\n\x05\x63ount\x18\x01 \x01(\x03\"h\n\x12PartitionsResponse\x12\x34\n\x06result\x18\x01 \x01(\x0b\x32$.source.v1.PartitionsResponse.Result\x1a\x1c\n\x06Result\x12\x12\n\npartitions\x18\x01 \x03(\x05\".\n\x06Offset\x12\x0e\n\x06offset\x18\x01 \x01(\x0c\x12\x14\n\x0cpartition_id\x18\x02 \x01(\x05\x32\x83\x03\n\x06Source\x12=\n\x06ReadFn\x12\x16.source.v1.ReadRequest\x1a\x17.source.v1.ReadResponse(\x01\x30\x01\x12:\n\x05\x41\x63kFn\x12\x15.source.v1.AckRequest\x1a\x16.source.v1.AckResponse(\x01\x30\x01\x12\x39\n\x06NackFn\x12\x16.source.v1.NackRequest\x1a\x17.source.v1.NackResponse\x12?\n\tPendingFn\x12\x16.google.protobuf.Empty\x1a\x1a.source.v1.PendingResponse\x12\x45\n\x0cPartitionsFn\x12\x16.google.protobuf.Empty\x1a\x1d.source.v1.PartitionsResponse\x12;\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a\x18.source.v1.ReadyResponseb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) -_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'source_pb2', _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'pynumaflow.proto.sourcer.source_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: DESCRIPTOR._loaded_options = None _globals['_READRESPONSE_RESULT_HEADERSENTRY']._loaded_options = None _globals['_READRESPONSE_RESULT_HEADERSENTRY']._serialized_options = b'8\001' - _globals['_HANDSHAKE']._serialized_start=89 - _globals['_HANDSHAKE']._serialized_end=113 - _globals['_READREQUEST']._serialized_start=116 - _globals['_READREQUEST']._serialized_end=293 - _globals['_READREQUEST_REQUEST']._serialized_start=226 - _globals['_READREQUEST_REQUEST']._serialized_end=279 - _globals['_READRESPONSE']._serialized_start=296 - _globals['_READRESPONSE']._serialized_end=937 - _globals['_READRESPONSE_RESULT']._serialized_start=455 - _globals['_READRESPONSE_RESULT']._serialized_end=687 - _globals['_READRESPONSE_RESULT_HEADERSENTRY']._serialized_start=641 - _globals['_READRESPONSE_RESULT_HEADERSENTRY']._serialized_end=687 - _globals['_READRESPONSE_STATUS']._serialized_start=690 - _globals['_READRESPONSE_STATUS']._serialized_end=923 - _globals['_READRESPONSE_STATUS_CODE']._serialized_start=840 - _globals['_READRESPONSE_STATUS_CODE']._serialized_end=872 - _globals['_READRESPONSE_STATUS_ERROR']._serialized_start=874 - _globals['_READRESPONSE_STATUS_ERROR']._serialized_end=905 - _globals['_ACKREQUEST']._serialized_start=940 - _globals['_ACKREQUEST']._serialized_end=1107 - _globals['_ACKREQUEST_REQUEST']._serialized_start=1048 - _globals['_ACKREQUEST_REQUEST']._serialized_end=1093 - _globals['_ACKRESPONSE']._serialized_start=1110 - _globals['_ACKRESPONSE']._serialized_end=1281 - _globals['_ACKRESPONSE_RESULT']._serialized_start=1218 - _globals['_ACKRESPONSE_RESULT']._serialized_end=1267 - _globals['_NACKREQUEST']._serialized_start=1283 - _globals['_NACKREQUEST']._serialized_end=1392 - _globals['_NACKREQUEST_REQUEST']._serialized_start=1048 - _globals['_NACKREQUEST_REQUEST']._serialized_end=1093 - _globals['_NACKRESPONSE']._serialized_start=1394 - _globals['_NACKRESPONSE']._serialized_end=1507 - _globals['_NACKRESPONSE_RESULT']._serialized_start=1218 - _globals['_NACKRESPONSE_RESULT']._serialized_end=1267 - _globals['_READYRESPONSE']._serialized_start=1509 - _globals['_READYRESPONSE']._serialized_end=1539 - _globals['_PENDINGRESPONSE']._serialized_start=1541 - _globals['_PENDINGRESPONSE']._serialized_end=1634 - _globals['_PENDINGRESPONSE_RESULT']._serialized_start=1611 - _globals['_PENDINGRESPONSE_RESULT']._serialized_end=1634 - _globals['_PARTITIONSRESPONSE']._serialized_start=1636 - _globals['_PARTITIONSRESPONSE']._serialized_end=1740 - _globals['_PARTITIONSRESPONSE_RESULT']._serialized_start=1712 - _globals['_PARTITIONSRESPONSE_RESULT']._serialized_end=1740 - _globals['_OFFSET']._serialized_start=1742 - _globals['_OFFSET']._serialized_end=1788 - _globals['_SOURCE']._serialized_start=1791 - _globals['_SOURCE']._serialized_end=2178 + _globals['_HANDSHAKE']._serialized_start=114 + _globals['_HANDSHAKE']._serialized_end=138 + _globals['_READREQUEST']._serialized_start=141 + _globals['_READREQUEST']._serialized_end=318 + _globals['_READREQUEST_REQUEST']._serialized_start=251 + _globals['_READREQUEST_REQUEST']._serialized_end=304 + _globals['_READRESPONSE']._serialized_start=321 + _globals['_READRESPONSE']._serialized_end=962 + _globals['_READRESPONSE_RESULT']._serialized_start=480 + _globals['_READRESPONSE_RESULT']._serialized_end=712 + _globals['_READRESPONSE_RESULT_HEADERSENTRY']._serialized_start=666 + _globals['_READRESPONSE_RESULT_HEADERSENTRY']._serialized_end=712 + _globals['_READRESPONSE_STATUS']._serialized_start=715 + _globals['_READRESPONSE_STATUS']._serialized_end=948 + _globals['_READRESPONSE_STATUS_CODE']._serialized_start=865 + _globals['_READRESPONSE_STATUS_CODE']._serialized_end=897 + _globals['_READRESPONSE_STATUS_ERROR']._serialized_start=899 + _globals['_READRESPONSE_STATUS_ERROR']._serialized_end=930 + _globals['_ACKREQUEST']._serialized_start=965 + _globals['_ACKREQUEST']._serialized_end=1132 + _globals['_ACKREQUEST_REQUEST']._serialized_start=1073 + _globals['_ACKREQUEST_REQUEST']._serialized_end=1118 + _globals['_ACKRESPONSE']._serialized_start=1135 + _globals['_ACKRESPONSE']._serialized_end=1306 + _globals['_ACKRESPONSE_RESULT']._serialized_start=1243 + _globals['_ACKRESPONSE_RESULT']._serialized_end=1292 + _globals['_NACKREQUEST']._serialized_start=1308 + _globals['_NACKREQUEST']._serialized_end=1417 + _globals['_NACKREQUEST_REQUEST']._serialized_start=1073 + _globals['_NACKREQUEST_REQUEST']._serialized_end=1118 + _globals['_NACKRESPONSE']._serialized_start=1419 + _globals['_NACKRESPONSE']._serialized_end=1532 + _globals['_NACKRESPONSE_RESULT']._serialized_start=1243 + _globals['_NACKRESPONSE_RESULT']._serialized_end=1292 + _globals['_READYRESPONSE']._serialized_start=1534 + _globals['_READYRESPONSE']._serialized_end=1564 + _globals['_PENDINGRESPONSE']._serialized_start=1566 + _globals['_PENDINGRESPONSE']._serialized_end=1659 + _globals['_PENDINGRESPONSE_RESULT']._serialized_start=1636 + _globals['_PENDINGRESPONSE_RESULT']._serialized_end=1659 + _globals['_PARTITIONSRESPONSE']._serialized_start=1661 + _globals['_PARTITIONSRESPONSE']._serialized_end=1765 + _globals['_PARTITIONSRESPONSE_RESULT']._serialized_start=1737 + _globals['_PARTITIONSRESPONSE_RESULT']._serialized_end=1765 + _globals['_OFFSET']._serialized_start=1767 + _globals['_OFFSET']._serialized_end=1813 + _globals['_SOURCE']._serialized_start=1816 + _globals['_SOURCE']._serialized_end=2203 # @@protoc_insertion_point(module_scope) diff --git a/pynumaflow/proto/sourcer/source_pb2_grpc.py b/pynumaflow/proto/sourcer/source_pb2_grpc.py index 1e4f6c15..6dd4103d 100644 --- a/pynumaflow/proto/sourcer/source_pb2_grpc.py +++ b/pynumaflow/proto/sourcer/source_pb2_grpc.py @@ -4,7 +4,7 @@ import warnings from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -from . import source_pb2 as source__pb2 +from pynumaflow.proto.sourcer import source_pb2 as pynumaflow_dot_proto_dot_sourcer_dot_source__pb2 GRPC_GENERATED_VERSION = '1.75.0' GRPC_VERSION = grpc.__version__ @@ -19,7 +19,7 @@ if _version_not_supported: raise RuntimeError( f'The grpc package installed is at version {GRPC_VERSION},' - + f' but the generated code in source_pb2_grpc.py depends on' + + f' but the generated code in pynumaflow/proto/sourcer/source_pb2_grpc.py depends on' + f' grpcio>={GRPC_GENERATED_VERSION}.' + f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}' + f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.' @@ -37,33 +37,33 @@ def __init__(self, channel): """ self.ReadFn = channel.stream_stream( '/source.v1.Source/ReadFn', - request_serializer=source__pb2.ReadRequest.SerializeToString, - response_deserializer=source__pb2.ReadResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadResponse.FromString, _registered_method=True) self.AckFn = channel.stream_stream( '/source.v1.Source/AckFn', - request_serializer=source__pb2.AckRequest.SerializeToString, - response_deserializer=source__pb2.AckResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.AckRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.AckResponse.FromString, _registered_method=True) self.NackFn = channel.unary_unary( '/source.v1.Source/NackFn', - request_serializer=source__pb2.NackRequest.SerializeToString, - response_deserializer=source__pb2.NackResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.NackRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.NackResponse.FromString, _registered_method=True) self.PendingFn = channel.unary_unary( '/source.v1.Source/PendingFn', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=source__pb2.PendingResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.PendingResponse.FromString, _registered_method=True) self.PartitionsFn = channel.unary_unary( '/source.v1.Source/PartitionsFn', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=source__pb2.PartitionsResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.PartitionsResponse.FromString, _registered_method=True) self.IsReady = channel.unary_unary( '/source.v1.Source/IsReady', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=source__pb2.ReadyResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadyResponse.FromString, _registered_method=True) @@ -127,33 +127,33 @@ def add_SourceServicer_to_server(servicer, server): rpc_method_handlers = { 'ReadFn': grpc.stream_stream_rpc_method_handler( servicer.ReadFn, - request_deserializer=source__pb2.ReadRequest.FromString, - response_serializer=source__pb2.ReadResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadResponse.SerializeToString, ), 'AckFn': grpc.stream_stream_rpc_method_handler( servicer.AckFn, - request_deserializer=source__pb2.AckRequest.FromString, - response_serializer=source__pb2.AckResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.AckRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.AckResponse.SerializeToString, ), 'NackFn': grpc.unary_unary_rpc_method_handler( servicer.NackFn, - request_deserializer=source__pb2.NackRequest.FromString, - response_serializer=source__pb2.NackResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.NackRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.NackResponse.SerializeToString, ), 'PendingFn': grpc.unary_unary_rpc_method_handler( servicer.PendingFn, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=source__pb2.PendingResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.PendingResponse.SerializeToString, ), 'PartitionsFn': grpc.unary_unary_rpc_method_handler( servicer.PartitionsFn, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=source__pb2.PartitionsResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.PartitionsResponse.SerializeToString, ), 'IsReady': grpc.unary_unary_rpc_method_handler( servicer.IsReady, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=source__pb2.ReadyResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadyResponse.SerializeToString, ), } generic_handler = grpc.method_handlers_generic_handler( @@ -181,8 +181,8 @@ def ReadFn(request_iterator, request_iterator, target, '/source.v1.Source/ReadFn', - source__pb2.ReadRequest.SerializeToString, - source__pb2.ReadResponse.FromString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadRequest.SerializeToString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadResponse.FromString, options, channel_credentials, insecure, @@ -208,8 +208,8 @@ def AckFn(request_iterator, request_iterator, target, '/source.v1.Source/AckFn', - source__pb2.AckRequest.SerializeToString, - source__pb2.AckResponse.FromString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.AckRequest.SerializeToString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.AckResponse.FromString, options, channel_credentials, insecure, @@ -235,8 +235,8 @@ def NackFn(request, request, target, '/source.v1.Source/NackFn', - source__pb2.NackRequest.SerializeToString, - source__pb2.NackResponse.FromString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.NackRequest.SerializeToString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.NackResponse.FromString, options, channel_credentials, insecure, @@ -263,7 +263,7 @@ def PendingFn(request, target, '/source.v1.Source/PendingFn', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - source__pb2.PendingResponse.FromString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.PendingResponse.FromString, options, channel_credentials, insecure, @@ -290,7 +290,7 @@ def PartitionsFn(request, target, '/source.v1.Source/PartitionsFn', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - source__pb2.PartitionsResponse.FromString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.PartitionsResponse.FromString, options, channel_credentials, insecure, @@ -317,7 +317,7 @@ def IsReady(request, target, '/source.v1.Source/IsReady', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - source__pb2.ReadyResponse.FromString, + pynumaflow_dot_proto_dot_sourcer_dot_source__pb2.ReadyResponse.FromString, options, channel_credentials, insecure, diff --git a/pynumaflow/proto/sourcetransformer/transform_pb2.py b/pynumaflow/proto/sourcetransformer/transform_pb2.py index 1cafc035..aebfb85a 100644 --- a/pynumaflow/proto/sourcetransformer/transform_pb2.py +++ b/pynumaflow/proto/sourcetransformer/transform_pb2.py @@ -1,7 +1,7 @@ # -*- coding: utf-8 -*- # Generated by the protocol buffer compiler. DO NOT EDIT! # NO CHECKED-IN PROTOBUF GENCODE -# source: transform.proto +# source: pynumaflow/proto/sourcetransformer/transform.proto # Protobuf Python Version: 6.31.1 """Generated protocol buffer code.""" from google.protobuf import descriptor as _descriptor @@ -15,7 +15,7 @@ 31, 1, '', - 'transform.proto' + 'pynumaflow/proto/sourcetransformer/transform.proto' ) # @@protoc_insertion_point(imports) @@ -26,29 +26,29 @@ from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x0ftransform.proto\x12\x14sourcetransformer.v1\x1a\x1fgoogle/protobuf/timestamp.proto\x1a\x1bgoogle/protobuf/empty.proto\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"\xbe\x03\n\x16SourceTransformRequest\x12\x45\n\x07request\x18\x01 \x01(\x0b\x32\x34.sourcetransformer.v1.SourceTransformRequest.Request\x12\x37\n\thandshake\x18\x02 \x01(\x0b\x32\x1f.sourcetransformer.v1.HandshakeH\x00\x88\x01\x01\x1a\x95\x02\n\x07Request\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12R\n\x07headers\x18\x05 \x03(\x0b\x32\x41.sourcetransformer.v1.SourceTransformRequest.Request.HeadersEntry\x12\n\n\x02id\x18\x06 \x01(\t\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x0c\n\n_handshake\"\x98\x02\n\x17SourceTransformResponse\x12\x45\n\x07results\x18\x01 \x03(\x0b\x32\x34.sourcetransformer.v1.SourceTransformResponse.Result\x12\n\n\x02id\x18\x02 \x01(\t\x12\x37\n\thandshake\x18\x03 \x01(\x0b\x32\x1f.sourcetransformer.v1.HandshakeH\x00\x88\x01\x01\x1a\x63\n\x06Result\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04tags\x18\x04 \x03(\tB\x0c\n\n_handshake\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\xcf\x01\n\x0fSourceTransform\x12t\n\x11SourceTransformFn\x12,.sourcetransformer.v1.SourceTransformRequest\x1a-.sourcetransformer.v1.SourceTransformResponse(\x01\x30\x01\x12\x46\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a#.sourcetransformer.v1.ReadyResponseb\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n2pynumaflow/proto/sourcetransformer/transform.proto\x12\x14sourcetransformer.v1\x1a\x1fgoogle/protobuf/timestamp.proto\x1a\x1bgoogle/protobuf/empty.proto\"\x18\n\tHandshake\x12\x0b\n\x03sot\x18\x01 \x01(\x08\"\xbe\x03\n\x16SourceTransformRequest\x12\x45\n\x07request\x18\x01 \x01(\x0b\x32\x34.sourcetransformer.v1.SourceTransformRequest.Request\x12\x37\n\thandshake\x18\x02 \x01(\x0b\x32\x1f.sourcetransformer.v1.HandshakeH\x00\x88\x01\x01\x1a\x95\x02\n\x07Request\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12-\n\twatermark\x18\x04 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12R\n\x07headers\x18\x05 \x03(\x0b\x32\x41.sourcetransformer.v1.SourceTransformRequest.Request.HeadersEntry\x12\n\n\x02id\x18\x06 \x01(\t\x1a.\n\x0cHeadersEntry\x12\x0b\n\x03key\x18\x01 \x01(\t\x12\r\n\x05value\x18\x02 \x01(\t:\x02\x38\x01\x42\x0c\n\n_handshake\"\x98\x02\n\x17SourceTransformResponse\x12\x45\n\x07results\x18\x01 \x03(\x0b\x32\x34.sourcetransformer.v1.SourceTransformResponse.Result\x12\n\n\x02id\x18\x02 \x01(\t\x12\x37\n\thandshake\x18\x03 \x01(\x0b\x32\x1f.sourcetransformer.v1.HandshakeH\x00\x88\x01\x01\x1a\x63\n\x06Result\x12\x0c\n\x04keys\x18\x01 \x03(\t\x12\r\n\x05value\x18\x02 \x01(\x0c\x12.\n\nevent_time\x18\x03 \x01(\x0b\x32\x1a.google.protobuf.Timestamp\x12\x0c\n\x04tags\x18\x04 \x03(\tB\x0c\n\n_handshake\"\x1e\n\rReadyResponse\x12\r\n\x05ready\x18\x01 \x01(\x08\x32\xcf\x01\n\x0fSourceTransform\x12t\n\x11SourceTransformFn\x12,.sourcetransformer.v1.SourceTransformRequest\x1a-.sourcetransformer.v1.SourceTransformResponse(\x01\x30\x01\x12\x46\n\x07IsReady\x12\x16.google.protobuf.Empty\x1a#.sourcetransformer.v1.ReadyResponseb\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) -_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'transform_pb2', _globals) +_builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'pynumaflow.proto.sourcetransformer.transform_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: DESCRIPTOR._loaded_options = None _globals['_SOURCETRANSFORMREQUEST_REQUEST_HEADERSENTRY']._loaded_options = None _globals['_SOURCETRANSFORMREQUEST_REQUEST_HEADERSENTRY']._serialized_options = b'8\001' - _globals['_HANDSHAKE']._serialized_start=103 - _globals['_HANDSHAKE']._serialized_end=127 - _globals['_SOURCETRANSFORMREQUEST']._serialized_start=130 - _globals['_SOURCETRANSFORMREQUEST']._serialized_end=576 - _globals['_SOURCETRANSFORMREQUEST_REQUEST']._serialized_start=285 - _globals['_SOURCETRANSFORMREQUEST_REQUEST']._serialized_end=562 - _globals['_SOURCETRANSFORMREQUEST_REQUEST_HEADERSENTRY']._serialized_start=516 - _globals['_SOURCETRANSFORMREQUEST_REQUEST_HEADERSENTRY']._serialized_end=562 - _globals['_SOURCETRANSFORMRESPONSE']._serialized_start=579 - _globals['_SOURCETRANSFORMRESPONSE']._serialized_end=859 - _globals['_SOURCETRANSFORMRESPONSE_RESULT']._serialized_start=746 - _globals['_SOURCETRANSFORMRESPONSE_RESULT']._serialized_end=845 - _globals['_READYRESPONSE']._serialized_start=861 - _globals['_READYRESPONSE']._serialized_end=891 - _globals['_SOURCETRANSFORM']._serialized_start=894 - _globals['_SOURCETRANSFORM']._serialized_end=1101 + _globals['_HANDSHAKE']._serialized_start=138 + _globals['_HANDSHAKE']._serialized_end=162 + _globals['_SOURCETRANSFORMREQUEST']._serialized_start=165 + _globals['_SOURCETRANSFORMREQUEST']._serialized_end=611 + _globals['_SOURCETRANSFORMREQUEST_REQUEST']._serialized_start=320 + _globals['_SOURCETRANSFORMREQUEST_REQUEST']._serialized_end=597 + _globals['_SOURCETRANSFORMREQUEST_REQUEST_HEADERSENTRY']._serialized_start=551 + _globals['_SOURCETRANSFORMREQUEST_REQUEST_HEADERSENTRY']._serialized_end=597 + _globals['_SOURCETRANSFORMRESPONSE']._serialized_start=614 + _globals['_SOURCETRANSFORMRESPONSE']._serialized_end=894 + _globals['_SOURCETRANSFORMRESPONSE_RESULT']._serialized_start=781 + _globals['_SOURCETRANSFORMRESPONSE_RESULT']._serialized_end=880 + _globals['_READYRESPONSE']._serialized_start=896 + _globals['_READYRESPONSE']._serialized_end=926 + _globals['_SOURCETRANSFORM']._serialized_start=929 + _globals['_SOURCETRANSFORM']._serialized_end=1136 # @@protoc_insertion_point(module_scope) diff --git a/pynumaflow/proto/sourcetransformer/transform_pb2_grpc.py b/pynumaflow/proto/sourcetransformer/transform_pb2_grpc.py index 98dbee5f..942c4450 100644 --- a/pynumaflow/proto/sourcetransformer/transform_pb2_grpc.py +++ b/pynumaflow/proto/sourcetransformer/transform_pb2_grpc.py @@ -4,7 +4,7 @@ import warnings from google.protobuf import empty_pb2 as google_dot_protobuf_dot_empty__pb2 -from . import transform_pb2 as transform__pb2 +from pynumaflow.proto.sourcetransformer import transform_pb2 as pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2 GRPC_GENERATED_VERSION = '1.75.0' GRPC_VERSION = grpc.__version__ @@ -19,7 +19,7 @@ if _version_not_supported: raise RuntimeError( f'The grpc package installed is at version {GRPC_VERSION},' - + f' but the generated code in transform_pb2_grpc.py depends on' + + f' but the generated code in pynumaflow/proto/sourcetransformer/transform_pb2_grpc.py depends on' + f' grpcio>={GRPC_GENERATED_VERSION}.' + f' Please upgrade your grpc module to grpcio>={GRPC_GENERATED_VERSION}' + f' or downgrade your generated code using grpcio-tools<={GRPC_VERSION}.' @@ -37,13 +37,13 @@ def __init__(self, channel): """ self.SourceTransformFn = channel.stream_stream( '/sourcetransformer.v1.SourceTransform/SourceTransformFn', - request_serializer=transform__pb2.SourceTransformRequest.SerializeToString, - response_deserializer=transform__pb2.SourceTransformResponse.FromString, + request_serializer=pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.SourceTransformRequest.SerializeToString, + response_deserializer=pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.SourceTransformResponse.FromString, _registered_method=True) self.IsReady = channel.unary_unary( '/sourcetransformer.v1.SourceTransform/IsReady', request_serializer=google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - response_deserializer=transform__pb2.ReadyResponse.FromString, + response_deserializer=pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.ReadyResponse.FromString, _registered_method=True) @@ -71,13 +71,13 @@ def add_SourceTransformServicer_to_server(servicer, server): rpc_method_handlers = { 'SourceTransformFn': grpc.stream_stream_rpc_method_handler( servicer.SourceTransformFn, - request_deserializer=transform__pb2.SourceTransformRequest.FromString, - response_serializer=transform__pb2.SourceTransformResponse.SerializeToString, + request_deserializer=pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.SourceTransformRequest.FromString, + response_serializer=pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.SourceTransformResponse.SerializeToString, ), 'IsReady': grpc.unary_unary_rpc_method_handler( servicer.IsReady, request_deserializer=google_dot_protobuf_dot_empty__pb2.Empty.FromString, - response_serializer=transform__pb2.ReadyResponse.SerializeToString, + response_serializer=pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.ReadyResponse.SerializeToString, ), } generic_handler = grpc.method_handlers_generic_handler( @@ -105,8 +105,8 @@ def SourceTransformFn(request_iterator, request_iterator, target, '/sourcetransformer.v1.SourceTransform/SourceTransformFn', - transform__pb2.SourceTransformRequest.SerializeToString, - transform__pb2.SourceTransformResponse.FromString, + pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.SourceTransformRequest.SerializeToString, + pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.SourceTransformResponse.FromString, options, channel_credentials, insecure, @@ -133,7 +133,7 @@ def IsReady(request, target, '/sourcetransformer.v1.SourceTransform/IsReady', google_dot_protobuf_dot_empty__pb2.Empty.SerializeToString, - transform__pb2.ReadyResponse.FromString, + pynumaflow_dot_proto_dot_sourcetransformer_dot_transform__pb2.ReadyResponse.FromString, options, channel_credentials, insecure,