-
Notifications
You must be signed in to change notification settings - Fork 2
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
feat: Move out concrete working implementation and add working modes #307
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,33 @@ | ||
/* | ||
* Copyright (C) 2024 Hedera Hashgraph, LLC | ||
* | ||
* Licensed under the Apache License, Version 2.0 (the "License"); | ||
* you may not use this file except in compliance with the License. | ||
* You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, | ||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
* See the License for the specific language governing permissions and | ||
* limitations under the License. | ||
*/ | ||
|
||
package com.hedera.block.simulator.config.types; | ||
|
||
/** The SimulatorMode enum defines the work modes of the block stream simulator. */ | ||
public enum SimulatorMode { | ||
/** | ||
* Indicates a work mode in which the simulator is working as both consumer and publisher. | ||
*/ | ||
BOTH, | ||
/** | ||
* Indicates a work mode in which the simulator is working in consumer mode. | ||
*/ | ||
CONSUMER, | ||
/** | ||
* Indicates a work mode in which the simulator is working in publisher mode. | ||
*/ | ||
PUBLISHER | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -24,6 +24,11 @@ | |
* The PublishStreamGrpcClient interface provides the methods to stream the block and block item. | ||
*/ | ||
public interface PublishStreamGrpcClient { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I like the new |
||
/** | ||
* Initialize, opens a gRPC channel and creates the needed stubs with the passed configuration. | ||
*/ | ||
void init(); | ||
|
||
/** | ||
* Streams the block item. | ||
* | ||
|
@@ -39,4 +44,9 @@ public interface PublishStreamGrpcClient { | |
* @return true if the block is streamed successfully, false otherwise | ||
*/ | ||
boolean streamBlock(Block block); | ||
|
||
/** | ||
* Shutdowns the channel. | ||
*/ | ||
void shutdown(); | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -37,9 +37,11 @@ | |
*/ | ||
public class PublishStreamGrpcClientImpl implements PublishStreamGrpcClient { | ||
|
||
private final BlockStreamServiceGrpc.BlockStreamServiceStub stub; | ||
private final StreamObserver<PublishStreamRequest> requestStreamObserver; | ||
private BlockStreamServiceGrpc.BlockStreamServiceStub stub; | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think you can move this declaration of |
||
private StreamObserver<PublishStreamRequest> requestStreamObserver; | ||
private final BlockStreamConfig blockStreamConfig; | ||
private final GrpcConfig grpcConfig; | ||
private ManagedChannel channel; | ||
|
||
/** | ||
* Creates a new PublishStreamGrpcClientImpl instance. | ||
|
@@ -50,14 +52,22 @@ public class PublishStreamGrpcClientImpl implements PublishStreamGrpcClient { | |
@Inject | ||
public PublishStreamGrpcClientImpl( | ||
@NonNull GrpcConfig grpcConfig, @NonNull BlockStreamConfig blockStreamConfig) { | ||
ManagedChannel channel = | ||
this.grpcConfig = grpcConfig; | ||
this.blockStreamConfig = blockStreamConfig; | ||
} | ||
|
||
/** | ||
* Initialize the channel and stub for publishBlockStream with the desired configuration. | ||
*/ | ||
@Override | ||
public void init() { | ||
channel = | ||
ManagedChannelBuilder.forAddress(grpcConfig.serverAddress(), grpcConfig.port()) | ||
.usePlaintext() | ||
.build(); | ||
stub = BlockStreamServiceGrpc.newStub(channel); | ||
PublishStreamObserver publishStreamObserver = new PublishStreamObserver(); | ||
requestStreamObserver = stub.publishBlockStream(publishStreamObserver); | ||
this.blockStreamConfig = blockStreamConfig; | ||
} | ||
|
||
/** | ||
|
@@ -99,4 +109,9 @@ public boolean streamBlock(Block block) { | |
|
||
return true; | ||
} | ||
|
||
@Override | ||
public void shutdown() { | ||
channel.shutdown(); | ||
} | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,62 @@ | ||
/* | ||
* Copyright (C) 2024 Hedera Hashgraph, LLC | ||
* | ||
* Licensed under the Apache License, Version 2.0 (the "License"); | ||
* you may not use this file except in compliance with the License. | ||
* You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, | ||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
* See the License for the specific language governing permissions and | ||
* limitations under the License. | ||
*/ | ||
|
||
package com.hedera.block.simulator.mode; | ||
|
||
import static java.util.Objects.requireNonNull; | ||
|
||
import com.hedera.block.simulator.config.data.BlockStreamConfig; | ||
import com.hedera.block.simulator.generator.BlockStreamManager; | ||
import edu.umd.cs.findbugs.annotations.NonNull; | ||
|
||
/** | ||
* The {@code CombinedModeHandler} class implements the {@link SimulatorModeHandler} interface | ||
* and provides the behavior for a mode where both consuming and publishing of block data | ||
* occur simultaneously. | ||
* | ||
* <p>This mode handles dual operations in the block streaming process, utilizing the | ||
* {@link BlockStreamConfig} for configuration settings. It is designed for scenarios where | ||
* the simulator needs to handle both the consumption and publication of blocks in parallel. | ||
* | ||
* <p>For now, the actual start behavior is not implemented, as indicated by the | ||
* {@link UnsupportedOperationException}. | ||
*/ | ||
public class CombinedModeHandler implements SimulatorModeHandler { | ||
private final BlockStreamConfig blockStreamConfig; | ||
|
||
/** | ||
* Constructs a new {@code CombinedModeHandler} with the specified block stream configuration. | ||
* | ||
* @param blockStreamConfig the configuration data for managing block streams | ||
*/ | ||
public CombinedModeHandler(@NonNull final BlockStreamConfig blockStreamConfig) { | ||
requireNonNull(blockStreamConfig); | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Does |
||
this.blockStreamConfig = blockStreamConfig; | ||
} | ||
|
||
/** | ||
* Starts the simulator in combined mode, handling both consumption and publication | ||
* of block stream. However, this method is currently not implemented, and will throw | ||
* an {@link UnsupportedOperationException}. | ||
* | ||
* @param blockStreamManager the {@link BlockStreamManager} responsible for managing block streams | ||
* @throws UnsupportedOperationException as the method is not yet implemented | ||
*/ | ||
@Override | ||
public void start(@NonNull BlockStreamManager blockStreamManager) { | ||
throw new UnsupportedOperationException(); | ||
} | ||
} |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
You might leverage a switch statement here like: