This is an automated email from the ASF dual-hosted git repository. lcwik pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/beam.git
commit c85f0f51634e65fe332e8bbe18fbf0689eb44d0f Merge: fc24deb 89b08e1 Author: Lukasz Cwik <lc...@google.com> AuthorDate: Fri Jun 7 18:01:26 2019 -0700 [BEAM-7470] Update proto and all SDKs to make the logical data stream over the data plane identified solely by instruction id and transform id. .../fn-execution/src/main/proto/beam_fn_api.proto | 32 +- .../dataflow/worker/BatchDataflowWorker.java | 7 +- .../worker/BeamFnMapTaskExecutorFactory.java | 16 +- .../dataflow/worker/FnApiWindowMappingFn.java | 16 +- .../dataflow/worker/StreamingDataflowWorker.java | 7 +- .../fn/data/RemoteGrpcPortReadOperation.java | 13 +- .../fn/data/RemoteGrpcPortWriteOperation.java | 21 +- .../graph/CreateRegisterFnOperationFunction.java | 13 +- .../beam/runners/dataflow/worker/graph/Nodes.java | 15 +- .../worker/graph/RegisterNodeFunction.java | 6 +- .../worker/fn/control/TimerReceiverTest.java | 22 +- .../worker/fn/data/BeamFnDataGrpcServiceTest.java | 18 +- .../fn/data/RemoteGrpcPortReadOperationTest.java | 15 +- .../fn/data/RemoteGrpcPortWriteOperationTest.java | 14 +- .../CreateRegisterFnOperationFunctionTest.java | 11 +- .../control/DefaultJobBundleFactory.java | 15 +- .../control/ProcessBundleDescriptors.java | 69 +- .../fnexecution/control/SdkHarnessClient.java | 27 +- .../SingleEnvironmentInstanceJobBundleFactory.java | 13 +- .../runners/fnexecution/data/GrpcDataService.java | 8 +- .../fnexecution/data/RemoteInputDestination.java | 6 +- .../fnexecution/control/RemoteExecutionTest.java | 122 +-- .../fnexecution/control/SdkHarnessClientTest.java | 61 +- .../fnexecution/data/GrpcDataServiceTest.java | 16 +- .../data/RemoteInputDestinationTest.java | 14 +- sdks/go/pkg/beam/core/runtime/exec/data.go | 16 +- sdks/go/pkg/beam/core/runtime/exec/datasource.go | 5 +- sdks/go/pkg/beam/core/runtime/exec/plan.go | 2 +- sdks/go/pkg/beam/core/runtime/exec/sideinput.go | 17 +- sdks/go/pkg/beam/core/runtime/exec/translate.go | 17 +- sdks/go/pkg/beam/core/runtime/harness/datamgr.go | 24 +- .../pkg/beam/core/runtime/harness/datamgr_test.go | 8 +- sdks/go/pkg/beam/core/runtime/harness/statemgr.go | 10 +- .../beam/model/fnexecution_v1/beam_fn_api.pb.go | 855 +++++++++------------ .../model/fnexecution_v1/beam_provision_api.pb.go | 18 +- .../model/jobmanagement_v1/beam_artifact_api.pb.go | 34 +- .../jobmanagement_v1/beam_expansion_api.pb.go | 8 +- .../beam/model/jobmanagement_v1/beam_job_api.pb.go | 52 +- .../beam/model/pipeline_v1/beam_runner_api.pb.go | 356 ++++----- sdks/go/pkg/beam/model/pipeline_v1/endpoints.pb.go | 22 +- .../model/pipeline_v1/external_transforms.pb.go | 8 +- sdks/go/pkg/beam/model/pipeline_v1/metrics.pb.go | 104 +-- .../model/pipeline_v1/standard_window_fns.pb.go | 20 +- sdks/go/pkg/beam/transforms/stats/stats.shims.go | 143 ++-- .../data/BeamFnDataBufferingOutboundObserver.java | 8 +- .../sdk/fn/data/BeamFnDataGrpcMultiplexer.java | 10 +- .../sdk/fn/data/BeamFnDataInboundObserver.java | 4 +- .../apache/beam/sdk/fn/data/LogicalEndpoint.java | 10 +- .../BeamFnDataBufferingOutboundObserverTest.java | 10 +- .../sdk/fn/data/BeamFnDataGrpcMultiplexerTest.java | 12 +- .../beam/fn/harness/BeamFnDataReadRunner.java | 19 +- .../beam/fn/harness/BeamFnDataWriteRunner.java | 20 +- .../beam/fn/harness/data/BeamFnDataGrpcClient.java | 8 +- .../fn/harness/data/QueueingBeamFnDataClient.java | 8 +- .../beam/fn/harness/BeamFnDataReadRunnerTest.java | 17 +- .../beam/fn/harness/BeamFnDataWriteRunnerTest.java | 21 +- .../fn/harness/data/BeamFnDataGrpcClientTest.java | 24 +- .../data/BeamFnDataInboundObserverTest.java | 5 +- .../harness/data/QueueingBeamFnDataClientTest.java | 26 +- .../runners/portability/fn_api_runner.py | 45 +- .../apache_beam/runners/worker/bundle_processor.py | 44 +- .../apache_beam/runners/worker/data_plane.py | 43 +- .../apache_beam/runners/worker/data_plane_test.py | 38 +- 63 files changed, 1138 insertions(+), 1530 deletions(-)