diff --git a/proto.lock b/proto.lock
index 07fd846f2e..5ea263371e 100644
--- a/proto.lock
+++ b/proto.lock
@@ -1,1506 +1,5 @@
{
"definitions": [
- {
- "protopath": "ClientAPI:/:ClientMessageDtos.proto",
- "def": {
- "enums": [
- {
- "name": "OperationResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "PrepareTimeout",
- "integer": 1
- },
- {
- "name": "CommitTimeout",
- "integer": 2
- },
- {
- "name": "ForwardTimeout",
- "integer": 3
- },
- {
- "name": "WrongExpectedVersion",
- "integer": 4
- },
- {
- "name": "StreamDeleted",
- "integer": 5
- },
- {
- "name": "InvalidTransaction",
- "integer": 6
- },
- {
- "name": "AccessDenied",
- "integer": 7
- }
- ]
- },
- {
- "name": "ReadEventCompleted.ReadEventResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "NotFound",
- "integer": 1
- },
- {
- "name": "NoStream",
- "integer": 2
- },
- {
- "name": "StreamDeleted",
- "integer": 3
- },
- {
- "name": "Error",
- "integer": 4
- },
- {
- "name": "AccessDenied",
- "integer": 5
- }
- ]
- },
- {
- "name": "ReadStreamEventsCompleted.ReadStreamResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "NoStream",
- "integer": 1
- },
- {
- "name": "StreamDeleted",
- "integer": 2
- },
- {
- "name": "NotModified",
- "integer": 3
- },
- {
- "name": "Error",
- "integer": 4
- },
- {
- "name": "AccessDenied",
- "integer": 5
- }
- ]
- },
- {
- "name": "ReadAllEventsCompleted.ReadAllResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "NotModified",
- "integer": 1
- },
- {
- "name": "Error",
- "integer": 2
- },
- {
- "name": "AccessDenied",
- "integer": 3
- }
- ]
- },
- {
- "name": "Filter.FilterContext",
- "enum_fields": [
- {
- "name": "StreamId"
- },
- {
- "name": "EventType",
- "integer": 1
- }
- ]
- },
- {
- "name": "Filter.FilterType",
- "enum_fields": [
- {
- "name": "Regex"
- },
- {
- "name": "Prefix",
- "integer": 1
- }
- ]
- },
- {
- "name": "FilteredReadAllEventsCompleted.FilteredReadAllResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "NotModified",
- "integer": 1
- },
- {
- "name": "Error",
- "integer": 2
- },
- {
- "name": "AccessDenied",
- "integer": 3
- }
- ]
- },
- {
- "name": "UpdatePersistentSubscriptionCompleted.UpdatePersistentSubscriptionResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "DoesNotExist",
- "integer": 1
- },
- {
- "name": "Fail",
- "integer": 2
- },
- {
- "name": "AccessDenied",
- "integer": 3
- }
- ]
- },
- {
- "name": "CreatePersistentSubscriptionCompleted.CreatePersistentSubscriptionResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "AlreadyExists",
- "integer": 1
- },
- {
- "name": "Fail",
- "integer": 2
- },
- {
- "name": "AccessDenied",
- "integer": 3
- }
- ]
- },
- {
- "name": "DeletePersistentSubscriptionCompleted.DeletePersistentSubscriptionResult",
- "enum_fields": [
- {
- "name": "Success"
- },
- {
- "name": "DoesNotExist",
- "integer": 1
- },
- {
- "name": "Fail",
- "integer": 2
- },
- {
- "name": "AccessDenied",
- "integer": 3
- }
- ]
- },
- {
- "name": "PersistentSubscriptionNakEvents.NakAction",
- "enum_fields": [
- {
- "name": "Unknown"
- },
- {
- "name": "Park",
- "integer": 1
- },
- {
- "name": "Retry",
- "integer": 2
- },
- {
- "name": "Skip",
- "integer": 3
- },
- {
- "name": "Stop",
- "integer": 4
- }
- ]
- },
- {
- "name": "SubscriptionDropped.SubscriptionDropReason",
- "enum_fields": [
- {
- "name": "Unsubscribed"
- },
- {
- "name": "AccessDenied",
- "integer": 1
- },
- {
- "name": "NotFound",
- "integer": 2
- },
- {
- "name": "PersistentSubscriptionDeleted",
- "integer": 3
- },
- {
- "name": "SubscriberMaxCountReached",
- "integer": 4
- }
- ]
- },
- {
- "name": "NotHandled.NotHandledReason",
- "enum_fields": [
- {
- "name": "NotReady"
- },
- {
- "name": "TooBusy",
- "integer": 1
- },
- {
- "name": "NotLeader",
- "integer": 2
- },
- {
- "name": "IsReadOnly",
- "integer": 3
- }
- ]
- },
- {
- "name": "ScavengeDatabaseResponse.ScavengeResult",
- "enum_fields": [
- {
- "name": "Started"
- },
- {
- "name": "InProgress",
- "integer": 1
- },
- {
- "name": "Unauthorized",
- "integer": 2
- }
- ]
- }
- ],
- "messages": [
- {
- "name": "NewEvent",
- "fields": [
- {
- "id": 1,
- "name": "event_id",
- "type": "bytes"
- },
- {
- "id": 2,
- "name": "event_type",
- "type": "string"
- },
- {
- "id": 3,
- "name": "data_content_type",
- "type": "int32"
- },
- {
- "id": 4,
- "name": "metadata_content_type",
- "type": "int32"
- },
- {
- "id": 5,
- "name": "data",
- "type": "bytes"
- },
- {
- "id": 6,
- "name": "metadata",
- "type": "bytes"
- }
- ]
- },
- {
- "name": "EventRecord",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "event_number",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "event_id",
- "type": "bytes"
- },
- {
- "id": 4,
- "name": "event_type",
- "type": "string"
- },
- {
- "id": 5,
- "name": "data_content_type",
- "type": "int32"
- },
- {
- "id": 6,
- "name": "metadata_content_type",
- "type": "int32"
- },
- {
- "id": 7,
- "name": "data",
- "type": "bytes"
- },
- {
- "id": 8,
- "name": "metadata",
- "type": "bytes"
- },
- {
- "id": 9,
- "name": "created",
- "type": "int64"
- },
- {
- "id": 10,
- "name": "created_epoch",
- "type": "int64"
- }
- ]
- },
- {
- "name": "ResolvedIndexedEvent",
- "fields": [
- {
- "id": 1,
- "name": "event",
- "type": "EventRecord"
- },
- {
- "id": 2,
- "name": "link",
- "type": "EventRecord"
- }
- ]
- },
- {
- "name": "ResolvedEvent",
- "fields": [
- {
- "id": 1,
- "name": "event",
- "type": "EventRecord"
- },
- {
- "id": 2,
- "name": "link",
- "type": "EventRecord"
- },
- {
- "id": 3,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 4,
- "name": "prepare_position",
- "type": "int64"
- }
- ]
- },
- {
- "name": "WriteEvents",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "expected_version",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "events",
- "type": "NewEvent",
- "is_repeated": true
- },
- {
- "id": 4,
- "name": "require_leader",
- "type": "bool"
- }
- ]
- },
- {
- "name": "WriteEventsCompleted",
- "fields": [
- {
- "id": 1,
- "name": "result",
- "type": "OperationResult"
- },
- {
- "id": 2,
- "name": "message",
- "type": "string"
- },
- {
- "id": 3,
- "name": "first_event_number",
- "type": "int64"
- },
- {
- "id": 4,
- "name": "last_event_number",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "prepare_position",
- "type": "int64"
- },
- {
- "id": 6,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 7,
- "name": "current_version",
- "type": "int64"
- }
- ]
- },
- {
- "name": "DeleteStream",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "expected_version",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "require_leader",
- "type": "bool"
- },
- {
- "id": 4,
- "name": "hard_delete",
- "type": "bool"
- }
- ]
- },
- {
- "name": "DeleteStreamCompleted",
- "fields": [
- {
- "id": 1,
- "name": "result",
- "type": "OperationResult"
- },
- {
- "id": 2,
- "name": "message",
- "type": "string"
- },
- {
- "id": 3,
- "name": "prepare_position",
- "type": "int64"
- },
- {
- "id": 4,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "current_version",
- "type": "int64"
- }
- ]
- },
- {
- "name": "TransactionStart",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "expected_version",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "require_leader",
- "type": "bool"
- }
- ]
- },
- {
- "name": "TransactionStartCompleted",
- "fields": [
- {
- "id": 1,
- "name": "transaction_id",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "result",
- "type": "OperationResult"
- },
- {
- "id": 3,
- "name": "message",
- "type": "string"
- }
- ]
- },
- {
- "name": "TransactionWrite",
- "fields": [
- {
- "id": 1,
- "name": "transaction_id",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "events",
- "type": "NewEvent",
- "is_repeated": true
- },
- {
- "id": 3,
- "name": "require_leader",
- "type": "bool"
- }
- ]
- },
- {
- "name": "TransactionWriteCompleted",
- "fields": [
- {
- "id": 1,
- "name": "transaction_id",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "result",
- "type": "OperationResult"
- },
- {
- "id": 3,
- "name": "message",
- "type": "string"
- }
- ]
- },
- {
- "name": "TransactionCommit",
- "fields": [
- {
- "id": 1,
- "name": "transaction_id",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "require_leader",
- "type": "bool"
- }
- ]
- },
- {
- "name": "TransactionCommitCompleted",
- "fields": [
- {
- "id": 1,
- "name": "transaction_id",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "result",
- "type": "OperationResult"
- },
- {
- "id": 3,
- "name": "message",
- "type": "string"
- },
- {
- "id": 4,
- "name": "first_event_number",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "last_event_number",
- "type": "int64"
- },
- {
- "id": 6,
- "name": "prepare_position",
- "type": "int64"
- },
- {
- "id": 7,
- "name": "commit_position",
- "type": "int64"
- }
- ]
- },
- {
- "name": "ReadEvent",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "event_number",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "resolve_link_tos",
- "type": "bool"
- },
- {
- "id": 4,
- "name": "require_leader",
- "type": "bool"
- }
- ]
- },
- {
- "name": "ReadEventCompleted",
- "fields": [
- {
- "id": 1,
- "name": "result",
- "type": "ReadEventResult"
- },
- {
- "id": 2,
- "name": "event",
- "type": "ResolvedIndexedEvent"
- },
- {
- "id": 3,
- "name": "error",
- "type": "string"
- }
- ]
- },
- {
- "name": "ReadStreamEvents",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "from_event_number",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "max_count",
- "type": "int32"
- },
- {
- "id": 4,
- "name": "resolve_link_tos",
- "type": "bool"
- },
- {
- "id": 5,
- "name": "require_leader",
- "type": "bool"
- }
- ]
- },
- {
- "name": "ReadStreamEventsCompleted",
- "fields": [
- {
- "id": 1,
- "name": "events",
- "type": "ResolvedIndexedEvent",
- "is_repeated": true
- },
- {
- "id": 2,
- "name": "result",
- "type": "ReadStreamResult"
- },
- {
- "id": 3,
- "name": "next_event_number",
- "type": "int64"
- },
- {
- "id": 4,
- "name": "last_event_number",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "is_end_of_stream",
- "type": "bool"
- },
- {
- "id": 6,
- "name": "last_commit_position",
- "type": "int64"
- },
- {
- "id": 7,
- "name": "error",
- "type": "string"
- }
- ]
- },
- {
- "name": "ReadAllEvents",
- "fields": [
- {
- "id": 1,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "prepare_position",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "max_count",
- "type": "int32"
- },
- {
- "id": 4,
- "name": "resolve_link_tos",
- "type": "bool"
- },
- {
- "id": 5,
- "name": "require_leader",
- "type": "bool"
- }
- ]
- },
- {
- "name": "ReadAllEventsCompleted",
- "fields": [
- {
- "id": 1,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "prepare_position",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "events",
- "type": "ResolvedEvent",
- "is_repeated": true
- },
- {
- "id": 4,
- "name": "next_commit_position",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "next_prepare_position",
- "type": "int64"
- },
- {
- "id": 6,
- "name": "result",
- "type": "ReadAllResult"
- },
- {
- "id": 7,
- "name": "error",
- "type": "string"
- }
- ]
- },
- {
- "name": "Filter",
- "fields": [
- {
- "id": 1,
- "name": "context",
- "type": "FilterContext"
- },
- {
- "id": 2,
- "name": "type",
- "type": "FilterType"
- },
- {
- "id": 3,
- "name": "data",
- "type": "string",
- "is_repeated": true
- }
- ]
- },
- {
- "name": "FilteredReadAllEvents",
- "fields": [
- {
- "id": 1,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "prepare_position",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "max_count",
- "type": "int32"
- },
- {
- "id": 4,
- "name": "max_search_window",
- "type": "int32"
- },
- {
- "id": 5,
- "name": "resolve_link_tos",
- "type": "bool"
- },
- {
- "id": 6,
- "name": "require_leader",
- "type": "bool"
- },
- {
- "id": 7,
- "name": "filter",
- "type": "Filter"
- }
- ]
- },
- {
- "name": "FilteredReadAllEventsCompleted",
- "fields": [
- {
- "id": 1,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "prepare_position",
- "type": "int64"
- },
- {
- "id": 3,
- "name": "events",
- "type": "ResolvedEvent",
- "is_repeated": true
- },
- {
- "id": 4,
- "name": "next_commit_position",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "next_prepare_position",
- "type": "int64"
- },
- {
- "id": 6,
- "name": "is_end_of_stream",
- "type": "bool"
- },
- {
- "id": 7,
- "name": "result",
- "type": "FilteredReadAllResult"
- },
- {
- "id": 8,
- "name": "error",
- "type": "string"
- }
- ]
- },
- {
- "name": "CreatePersistentSubscription",
- "fields": [
- {
- "id": 1,
- "name": "subscription_group_name",
- "type": "string"
- },
- {
- "id": 2,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 3,
- "name": "resolve_link_tos",
- "type": "bool"
- },
- {
- "id": 4,
- "name": "start_from",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "message_timeout_milliseconds",
- "type": "int32"
- },
- {
- "id": 6,
- "name": "record_statistics",
- "type": "bool"
- },
- {
- "id": 7,
- "name": "live_buffer_size",
- "type": "int32"
- },
- {
- "id": 8,
- "name": "read_batch_size",
- "type": "int32"
- },
- {
- "id": 9,
- "name": "buffer_size",
- "type": "int32"
- },
- {
- "id": 10,
- "name": "max_retry_count",
- "type": "int32"
- },
- {
- "id": 11,
- "name": "prefer_round_robin",
- "type": "bool"
- },
- {
- "id": 12,
- "name": "checkpoint_after_time",
- "type": "int32"
- },
- {
- "id": 13,
- "name": "checkpoint_max_count",
- "type": "int32"
- },
- {
- "id": 14,
- "name": "checkpoint_min_count",
- "type": "int32"
- },
- {
- "id": 15,
- "name": "subscriber_max_count",
- "type": "int32"
- },
- {
- "id": 16,
- "name": "named_consumer_strategy",
- "type": "string"
- }
- ]
- },
- {
- "name": "DeletePersistentSubscription",
- "fields": [
- {
- "id": 1,
- "name": "subscription_group_name",
- "type": "string"
- },
- {
- "id": 2,
- "name": "event_stream_id",
- "type": "string"
- }
- ]
- },
- {
- "name": "UpdatePersistentSubscription",
- "fields": [
- {
- "id": 1,
- "name": "subscription_group_name",
- "type": "string"
- },
- {
- "id": 2,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 3,
- "name": "resolve_link_tos",
- "type": "bool"
- },
- {
- "id": 4,
- "name": "start_from",
- "type": "int64"
- },
- {
- "id": 5,
- "name": "message_timeout_milliseconds",
- "type": "int32"
- },
- {
- "id": 6,
- "name": "record_statistics",
- "type": "bool"
- },
- {
- "id": 7,
- "name": "live_buffer_size",
- "type": "int32"
- },
- {
- "id": 8,
- "name": "read_batch_size",
- "type": "int32"
- },
- {
- "id": 9,
- "name": "buffer_size",
- "type": "int32"
- },
- {
- "id": 10,
- "name": "max_retry_count",
- "type": "int32"
- },
- {
- "id": 11,
- "name": "prefer_round_robin",
- "type": "bool"
- },
- {
- "id": 12,
- "name": "checkpoint_after_time",
- "type": "int32"
- },
- {
- "id": 13,
- "name": "checkpoint_max_count",
- "type": "int32"
- },
- {
- "id": 14,
- "name": "checkpoint_min_count",
- "type": "int32"
- },
- {
- "id": 15,
- "name": "subscriber_max_count",
- "type": "int32"
- },
- {
- "id": 16,
- "name": "named_consumer_strategy",
- "type": "string"
- }
- ]
- },
- {
- "name": "UpdatePersistentSubscriptionCompleted",
- "fields": [
- {
- "id": 1,
- "name": "result",
- "type": "UpdatePersistentSubscriptionResult"
- },
- {
- "id": 2,
- "name": "reason",
- "type": "string"
- }
- ]
- },
- {
- "name": "CreatePersistentSubscriptionCompleted",
- "fields": [
- {
- "id": 1,
- "name": "result",
- "type": "CreatePersistentSubscriptionResult"
- },
- {
- "id": 2,
- "name": "reason",
- "type": "string"
- }
- ]
- },
- {
- "name": "DeletePersistentSubscriptionCompleted",
- "fields": [
- {
- "id": 1,
- "name": "result",
- "type": "DeletePersistentSubscriptionResult"
- },
- {
- "id": 2,
- "name": "reason",
- "type": "string"
- }
- ]
- },
- {
- "name": "ConnectToPersistentSubscription",
- "fields": [
- {
- "id": 1,
- "name": "subscription_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 3,
- "name": "allowed_in_flight_messages",
- "type": "int32"
- }
- ]
- },
- {
- "name": "PersistentSubscriptionAckEvents",
- "fields": [
- {
- "id": 1,
- "name": "subscription_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "processed_event_ids",
- "type": "bytes",
- "is_repeated": true
- }
- ]
- },
- {
- "name": "PersistentSubscriptionNakEvents",
- "fields": [
- {
- "id": 1,
- "name": "subscription_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "processed_event_ids",
- "type": "bytes",
- "is_repeated": true
- },
- {
- "id": 3,
- "name": "message",
- "type": "string"
- },
- {
- "id": 4,
- "name": "action",
- "type": "NakAction"
- }
- ]
- },
- {
- "name": "PersistentSubscriptionConfirmation",
- "fields": [
- {
- "id": 1,
- "name": "last_commit_position",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "subscription_id",
- "type": "string"
- },
- {
- "id": 3,
- "name": "last_event_number",
- "type": "int64"
- }
- ]
- },
- {
- "name": "PersistentSubscriptionStreamEventAppeared",
- "fields": [
- {
- "id": 1,
- "name": "event",
- "type": "ResolvedIndexedEvent"
- },
- {
- "id": 2,
- "name": "retryCount",
- "type": "int32"
- }
- ]
- },
- {
- "name": "SubscribeToStream",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "resolve_link_tos",
- "type": "bool"
- }
- ]
- },
- {
- "name": "FilteredSubscribeToStream",
- "fields": [
- {
- "id": 1,
- "name": "event_stream_id",
- "type": "string"
- },
- {
- "id": 2,
- "name": "resolve_link_tos",
- "type": "bool"
- },
- {
- "id": 3,
- "name": "filter",
- "type": "Filter"
- },
- {
- "id": 4,
- "name": "checkpoint_interval",
- "type": "int32"
- }
- ]
- },
- {
- "name": "CheckpointReached",
- "fields": [
- {
- "id": 1,
- "name": "commit_position",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "prepare_position",
- "type": "int64"
- }
- ]
- },
- {
- "name": "SubscriptionConfirmation",
- "fields": [
- {
- "id": 1,
- "name": "last_commit_position",
- "type": "int64"
- },
- {
- "id": 2,
- "name": "last_event_number",
- "type": "int64"
- }
- ]
- },
- {
- "name": "StreamEventAppeared",
- "fields": [
- {
- "id": 1,
- "name": "event",
- "type": "ResolvedEvent"
- }
- ]
- },
- {
- "name": "UnsubscribeFromStream"
- },
- {
- "name": "SubscriptionDropped",
- "fields": [
- {
- "id": 1,
- "name": "reason",
- "type": "SubscriptionDropReason"
- }
- ]
- },
- {
- "name": "NotHandled",
- "fields": [
- {
- "id": 1,
- "name": "reason",
- "type": "NotHandledReason"
- },
- {
- "id": 2,
- "name": "additional_info",
- "type": "bytes"
- }
- ],
- "messages": [
- {
- "name": "LeaderInfo",
- "fields": [
- {
- "id": 1,
- "name": "external_tcp_address",
- "type": "string"
- },
- {
- "id": 2,
- "name": "external_tcp_port",
- "type": "int32"
- },
- {
- "id": 3,
- "name": "http_address",
- "type": "string"
- },
- {
- "id": 4,
- "name": "http_port",
- "type": "int32"
- },
- {
- "id": 5,
- "name": "external_secure_tcp_address",
- "type": "string"
- },
- {
- "id": 6,
- "name": "external_secure_tcp_port",
- "type": "int32"
- }
- ]
- }
- ]
- },
- {
- "name": "ScavengeDatabase"
- },
- {
- "name": "ScavengeDatabaseResponse",
- "fields": [
- {
- "id": 1,
- "name": "result",
- "type": "ScavengeResult"
- },
- {
- "id": 2,
- "name": "scavengeId",
- "type": "string"
- }
- ]
- },
- {
- "name": "IdentifyClient",
- "fields": [
- {
- "id": 1,
- "name": "version",
- "type": "int32"
- },
- {
- "id": 2,
- "name": "connection_name",
- "type": "string"
- }
- ]
- },
- {
- "name": "ClientIdentified"
- }
- ],
- "package": {
- "name": "EventStore.Client.Messages"
- }
- }
- },
{
"protopath": "Grpc:/:cluster.proto",
"def": {
@@ -1909,26 +408,6 @@
"name": "http_end_point",
"type": "EndPoint"
},
- {
- "id": 6,
- "name": "internal_tcp",
- "type": "EndPoint"
- },
- {
- "id": 7,
- "name": "external_tcp",
- "type": "EndPoint"
- },
- {
- "id": 8,
- "name": "internal_tcp_uses_tls",
- "type": "bool"
- },
- {
- "id": 9,
- "name": "external_tcp_uses_tls",
- "type": "bool"
- },
{
"id": 10,
"name": "last_commit_position",
@@ -1979,11 +458,6 @@
"name": "advertise_http_port_to_client_as",
"type": "uint32"
},
- {
- "id": 20,
- "name": "advertise_tcp_port_to_client_as",
- "type": "uint32"
- },
{
"id": 21,
"name": "es_version",
@@ -1994,6 +468,20 @@
"name": "replication_end_point",
"type": "EndPoint"
}
+ ],
+ "reserved_ids": [
+ 6,
+ 7,
+ 8,
+ 9,
+ 20
+ ],
+ "reserved_names": [
+ "internal_tcp",
+ "external_tcp",
+ "internal_tcp_uses_tls",
+ "external_tcp_uses_tls",
+ "advertise_tcp_port_to_client_as"
]
}
],
@@ -2789,21 +1277,19 @@
{
"name": "LeaderInfo",
"fields": [
- {
- "id": 1,
- "name": "external_tcp",
- "type": "EndPoint"
- },
- {
- "id": 2,
- "name": "is_secure",
- "type": "bool"
- },
{
"id": 3,
"name": "http",
"type": "EndPoint"
}
+ ],
+ "reserved_ids": [
+ 1,
+ 2
+ ],
+ "reserved_names": [
+ "external_tcp",
+ "is_secure"
]
},
{
diff --git a/src/EventStore.ClusterNode/Components/Pages/Cluster.razor b/src/EventStore.ClusterNode/Components/Pages/Cluster.razor
index 47d3c0279e..e72be103d8 100644
--- a/src/EventStore.ClusterNode/Components/Pages/Cluster.razor
+++ b/src/EventStore.ClusterNode/Components/Pages/Cluster.razor
@@ -43,8 +43,7 @@
Status |
Timestamp (UTC) |
Checkpoints |
- TCP |
- HTTP |
+ HTTP / gRPC |
Actions |
@@ -52,7 +51,7 @@
@if (!ClusterMembers.Any())
{
- | @ClusterEmptyMessage |
+ @ClusterEmptyMessage |
}
else
@@ -77,10 +76,6 @@
@EpochLabel(member)
}
-
- Internal: @InternalTcpEndpoint(member)
- External: @ExternalTcpEndpoint(member)
- |
@HttpEndpoint(member) |
@@ -307,7 +302,7 @@
@@ -370,9 +365,7 @@
.Append("Snapshot taken at ")
.Append(TimestampLabel(ClusterReadAt ?? DateTime.UtcNow))
.AppendLine()
- .Append(PadRight("Internal Tcp", 31)).Append(' ')
- .Append(PadRight("External Tcp", 31)).Append(' ')
- .Append(PadRight("Http", 23)).Append(' ')
+ .Append(PadRight("HTTP / gRPC", 23)).Append(' ')
.Append(PadRight("Status", 11)).Append(' ')
.Append(PadRight("State", 18)).Append(' ')
.Append(PadRight("Timestamp (UTC)", 19)).Append(" Checkpoints");
@@ -380,8 +373,6 @@
foreach (var member in ClusterMembers)
{
builder.AppendLine()
- .Append(PadRight(InternalTcpEndpoint(member), 31)).Append(' ')
- .Append(PadRight(ExternalTcpEndpoint(member), 31)).Append(' ')
.Append(PadRight(HttpEndpoint(member), 23)).Append(' ')
.Append(PadRight(MemberStatus(member), 11)).Append(' ')
.Append(PadRight(member.State.ToString(), 18)).Append(' ')
@@ -500,16 +491,6 @@
private static string MemberStatus(ClientClusterInfo.ClientMemberInfo member) =>
member.IsAlive ? "Alive" : "Unreachable";
- private static string InternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
- Endpoint(
- member.InternalTcpIp,
- member.InternalSecureTcpPort != 0 ? member.InternalSecureTcpPort : member.InternalTcpPort);
-
- private static string ExternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
- Endpoint(
- member.ExternalTcpIp,
- member.ExternalSecureTcpPort != 0 ? member.ExternalSecureTcpPort : member.ExternalTcpPort);
-
private static string HttpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(member.HttpEndPointIp, member.HttpEndPointPort);
diff --git a/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs b/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs
index c8362e60cc..e2a467d337 100644
--- a/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs
+++ b/src/EventStore.ClusterNode/Components/Services/ClusterStatusService.cs
@@ -171,7 +171,7 @@ private ClusterReplicaRow ParseReplicaRow(
: Guid.Empty;
var totalBytesSent = row.TotalBytesSent;
var previousRow = _previousReplicas.GetValueOrDefault(connectionId);
- var replicaNode = FindMemberByInternalEndpoint(members, row.SubscriptionEndpoint);
+ var replicaNode = FindMemberByEndpoint(members, row.SubscriptionEndpoint);
var isCatchingUp = replicaNode?.State == VNodeState.CatchingUp;
var catchupStartTime = now;
var catchupStartBytesSent = totalBytesSent;
@@ -206,14 +206,14 @@ private ClusterReplicaRow ParseReplicaRow(
private ClaimsPrincipal CurrentUser =>
httpContextAccessor.HttpContext?.User ?? new ClaimsPrincipal(new ClaimsIdentity());
- private static ClientClusterInfo.ClientMemberInfo FindMemberByInternalEndpoint(
+ private static ClientClusterInfo.ClientMemberInfo FindMemberByEndpoint(
IReadOnlyList members,
string endpoint)
{
var cleaned = endpoint.Replace("Unspecified/", "", StringComparison.OrdinalIgnoreCase);
return members.FirstOrDefault(x =>
- string.Equals(ReplicationEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase) ||
- string.Equals(InternalTcpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase));
+ string.Equals(ClusterEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase) ||
+ string.Equals(HttpEndpoint(x), cleaned, StringComparison.OrdinalIgnoreCase));
}
private static Uri BuildLeaderAddress(
@@ -221,14 +221,8 @@ private static Uri BuildLeaderAddress(
ClientClusterInfo.ClientMemberInfo leader) =>
new UriBuilder(request.Scheme, leader.HttpEndPointIp, leader.HttpEndPointPort).Uri;
- private static string InternalTcpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
- Endpoint(
- member.InternalTcpIp,
- member.InternalSecureTcpPort != 0 ? member.InternalSecureTcpPort : member.InternalTcpPort);
-
- private static string ReplicationEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
+ private static string ClusterEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(member.ClusterEndPointIp, member.ClusterEndPointPort);
-
private static string HttpEndpoint(ClientClusterInfo.ClientMemberInfo member) =>
Endpoint(member.HttpEndPointIp, member.HttpEndPointPort);
diff --git a/src/EventStore.Core.Tests/ClientAPI/Helpers/EventDataComparer.cs b/src/EventStore.Core.Tests/ClientAPI/Helpers/EventDataComparer.cs
deleted file mode 100644
index 15a424be6c..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/Helpers/EventDataComparer.cs
+++ /dev/null
@@ -1,46 +0,0 @@
-using EventStore.ClientAPI;
-using EventStore.Common.Utils;
-
-namespace EventStore.Core.Tests.ClientAPI.Helpers;
-
-internal static class EventDataComparer
-{
- public static bool Equal(EventData expected, RecordedEvent actual)
- {
- if (expected.EventId != actual.EventId)
- {
- return false;
- }
-
- if (expected.Type != actual.EventType)
- {
- return false;
- }
-
- var expectedDataString = Helper.UTF8NoBom.GetString(expected.Data ?? new byte[0]);
- var expectedMetadataString = Helper.UTF8NoBom.GetString(expected.Metadata ?? new byte[0]);
-
- var actualDataString = Helper.UTF8NoBom.GetString(actual.Data ?? new byte[0]);
- var actualMetadataDataString = Helper.UTF8NoBom.GetString(actual.Metadata ?? new byte[0]);
-
- return expectedDataString == actualDataString && expectedMetadataString == actualMetadataDataString;
- }
-
- public static bool Equal(EventData[] expected, RecordedEvent[] actual)
- {
- if (expected.Length != actual.Length)
- {
- return false;
- }
-
- for (var i = 0; i < expected.Length; i++)
- {
- if (!Equal(expected[i], actual[i]))
- {
- return false;
- }
- }
-
- return true;
- }
-}
diff --git a/src/EventStore.Core.Tests/ClientAPI/Helpers/EventsStream.cs b/src/EventStore.Core.Tests/ClientAPI/Helpers/EventsStream.cs
deleted file mode 100644
index da4bb62ce2..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/Helpers/EventsStream.cs
+++ /dev/null
@@ -1,25 +0,0 @@
-using System.Threading.Tasks;
-using EventStore.ClientAPI;
-
-namespace EventStore.Core.Tests.ClientAPI.Helpers;
-
-internal class EventsStream
-{
- private const int SliceSize = 10;
-
- public static async Task Count(IEventStoreConnection store, string stream)
- {
- var result = 0;
- while (true)
- {
- var slice = await store.ReadStreamEventsForwardAsync(stream, result, SliceSize, false);
- result += slice.Events.Length;
- if (slice.IsEndOfStream)
- {
- break;
- }
- }
-
- return result;
- }
-}
diff --git a/src/EventStore.Core.Tests/ClientAPI/Helpers/TcpType.cs b/src/EventStore.Core.Tests/ClientAPI/Helpers/TcpType.cs
deleted file mode 100644
index b860eb63ac..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/Helpers/TcpType.cs
+++ /dev/null
@@ -1,7 +0,0 @@
-namespace EventStore.Core.Tests.ClientAPI.Helpers;
-
-public enum TcpType
-{
- Normal,
- Ssl
-}
diff --git a/src/EventStore.Core.Tests/ClientAPI/Helpers/TestConnection.cs b/src/EventStore.Core.Tests/ClientAPI/Helpers/TestConnection.cs
deleted file mode 100644
index 674773889c..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/Helpers/TestConnection.cs
+++ /dev/null
@@ -1,72 +0,0 @@
-using System;
-using System.Net;
-using System.Threading;
-using EventStore.ClientAPI;
-using EventStore.ClientAPI.Internal;
-using EventStore.ClientAPI.SystemData;
-using EventStore.Core.Tests.Helpers;
-using NUnit.Framework;
-
-namespace EventStore.Core.Tests.ClientAPI.Helpers;
-
-public static class TestConnection
-{
- private static int _nextConnId = -1;
-
- public static IEventStoreConnection Create(IPEndPoint endPoint, TcpType tcpType = TcpType.Ssl,
- UserCredentials userCredentials = null)
- {
- return EventStoreConnection.Create(Settings(tcpType, userCredentials),
- endPoint.ToESTcpUri(),
- $"ESC-{Interlocked.Increment(ref _nextConnId)}");
- }
-
- public static IEventStoreConnection CreateMiniNodeClient(IPEndPoint endPoint, TcpType tcpType = TcpType.Ssl,
- UserCredentials userCredentials = null)
- {
- return EventStoreConnection.Create(Settings(
- tcpType,
- userCredentials,
- limitAttemptsForOperationTo: 10,
- reconnectionDelay: TimeSpan.FromMilliseconds(100)),
- endPoint.ToESTcpUri(),
- $"ESC-{Interlocked.Increment(ref _nextConnId)}");
- }
-
- public static IEventStoreConnection To(MiniNode miniNode, TcpType tcpType,
- UserCredentials userCredentials = null)
- {
- return EventStoreConnection.Create(Settings(tcpType, userCredentials),
- miniNode.TcpEndPoint.ToESTcpUri(),
- $"ESC-{Interlocked.Increment(ref _nextConnId)}");
- }
-
- private static ConnectionSettingsBuilder Settings(
- TcpType tcpType,
- UserCredentials userCredentials,
- int limitAttemptsForOperationTo = 1,
- TimeSpan? reconnectionDelay = null)
- {
- var settings = ConnectionSettings.Create()
- .SetDefaultUserCredentials(userCredentials)
- .UseCustomLogger(ClientApiLoggerBridge.Default)
- .EnableVerboseLogging()
- .LimitReconnectionsTo(10)
- .LimitAttemptsForOperationTo(limitAttemptsForOperationTo)
- .SetTimeoutCheckPeriodTo(TimeSpan.FromMilliseconds(100))
- .SetReconnectionDelayTo(reconnectionDelay ?? TimeSpan.Zero)
- .FailOnNoServerResponse()
- //.SetOperationTimeoutTo(TimeSpan.FromDays(1))
- ;
- if (tcpType == TcpType.Ssl)
- {
- settings.DisableServerCertificateValidation();
- }
- else
- {
- settings.DisableTls();
- }
-
- return settings;
- }
-}
diff --git a/src/EventStore.Core.Tests/ClientAPI/Helpers/TestConnectionLifecycle.cs b/src/EventStore.Core.Tests/ClientAPI/Helpers/TestConnectionLifecycle.cs
deleted file mode 100644
index 9ea1365f78..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/Helpers/TestConnectionLifecycle.cs
+++ /dev/null
@@ -1,77 +0,0 @@
-using System;
-using System.Threading.Tasks;
-using EventStore.ClientAPI;
-
-namespace EventStore.Core.Tests.ClientAPI.Helpers;
-
-public static class TestConnectionLifecycle
-{
- public static async Task ReconnectUntilReady(
- Func createConnection,
- Func readinessProbe,
- TimeSpan timeout)
- {
- var deadline = DateTime.UtcNow + timeout;
-
- while (true)
- {
- IEventStoreConnection connection = null;
-
- try
- {
- connection = createConnection();
- await connection.ConnectAsync();
- await readinessProbe(connection);
- return connection;
- }
- catch (Exception ex)
- {
- if (connection != null)
- {
- TryCloseConnection(connection);
- }
-
- if (IsTransientConnectionFailure(ex) && DateTime.UtcNow < deadline)
- {
- await Task.Delay(250);
- continue;
- }
-
- throw;
- }
- }
- }
-
- public static async Task CloseConnectionAndWait(IEventStoreConnection connection, TimeSpan timeout)
- {
- var closed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
- connection.Closed += (_, _) => closed.TrySetResult();
- connection.Close();
- await closed.Task.WithTimeout(timeout);
- }
-
- public static bool IsTransientConnectionFailure(Exception ex) =>
- ex.GetType().Name is "ConnectionClosedException"
- or "RetriesLimitReachedException"
- or "NotAuthenticatedException"
- or "AccessDeniedException";
-
- public static void TryCloseConnection(IEventStoreConnection connection)
- {
- try
- {
- connection.Close();
- }
- catch
- {
- }
- }
-
- public static void DisposeIfNeeded(object candidate)
- {
- if (candidate is IDisposable disposable)
- {
- disposable.Dispose();
- }
- }
-}
diff --git a/src/EventStore.Core.Tests/ClientAPI/Helpers/TestEvent.cs b/src/EventStore.Core.Tests/ClientAPI/Helpers/TestEvent.cs
deleted file mode 100644
index 8fb92bedf3..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/Helpers/TestEvent.cs
+++ /dev/null
@@ -1,22 +0,0 @@
-using System;
-using System.Text;
-using EventStore.ClientAPI;
-using EventStore.Common.Utils;
-
-namespace EventStore.Core.Tests.ClientAPI.Helpers;
-
-public class TestEvent
-{
- public static EventData NewTestEvent(string data = null, string metadata = null, string eventName = "TestEvent")
- {
- return NewTestEvent(Guid.NewGuid(), data, metadata, eventName);
- }
-
- public static EventData NewTestEvent(Guid eventId, string data = null, string metadata = null, string eventName = "TestEvent")
- {
- var encodedData = Helper.UTF8NoBom.GetBytes(data ?? eventId.ToString());
- var encodedMetadata = Helper.UTF8NoBom.GetBytes(metadata ?? "metadata");
-
- return new EventData(eventId, eventName, false, encodedData, encodedMetadata);
- }
-}
diff --git a/src/EventStore.Core.Tests/ClientAPI/Helpers/Writer.cs b/src/EventStore.Core.Tests/ClientAPI/Helpers/Writer.cs
deleted file mode 100644
index 236e47ccb6..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/Helpers/Writer.cs
+++ /dev/null
@@ -1,98 +0,0 @@
-using System.Threading.Tasks;
-using EventStore.ClientAPI;
-using NUnit.Framework;
-
-namespace EventStore.Core.Tests.ClientAPI.Helpers;
-
-internal class StreamWriter
-{
- private readonly IEventStoreConnection _store;
- private readonly string _stream;
- private readonly long _version;
-
- public StreamWriter(IEventStoreConnection store, string stream, long version)
- {
- _store = store;
- _stream = stream;
- _version = version;
- }
-
- public async Task Append(params EventData[] events)
- {
- for (var i = 0; i < events.Length; i++)
- {
- var expVer = _version == ExpectedVersion.Any ? ExpectedVersion.Any : _version + i;
- var nextExpVer = (await _store.AppendToStreamAsync(_stream, expVer, new[] { events[i] })).NextExpectedVersion;
- if (_version != ExpectedVersion.Any)
- {
- Assert.AreEqual(expVer + 1, nextExpVer);
- }
- }
-
- return new TailWriter(_store, _stream);
- }
-}
-
-internal class TailWriter
-{
- private readonly IEventStoreConnection _store;
- private readonly string _stream;
-
- public TailWriter(IEventStoreConnection store, string stream)
- {
- _store = store;
- _stream = stream;
- }
-
- public async Task Then(EventData @event, long expectedVersion)
- {
- await _store.AppendToStreamAsync(_stream, expectedVersion, new[] { @event });
- return this;
- }
-}
-
-internal class TransactionalWriter
-{
- private readonly IEventStoreConnection _store;
- private readonly string _stream;
-
- public TransactionalWriter(IEventStoreConnection store, string stream)
- {
- _store = store;
- _stream = stream;
- }
-
- public async Task StartTransaction(long expectedVersion)
- {
- return new OngoingTransaction(await _store.StartTransactionAsync(_stream, expectedVersion));
- }
-
- public OngoingTransaction ContinueTransaction(long transactionId)
- {
- return new OngoingTransaction(_store.ContinueTransaction(transactionId));
- }
-}
-
-//TODO GFY this should be removed and merged with the public idea of a transaction.
-internal class OngoingTransaction
-{
- private readonly EventStoreTransaction _transaction;
-
- public long TransactionId => _transaction.TransactionId;
-
- public OngoingTransaction(EventStoreTransaction transaction)
- {
- _transaction = transaction;
- }
-
- public async Task Write(params EventData[] events)
- {
- await _transaction.WriteAsync(events);
- return this;
- }
-
- public Task Commit()
- {
- return _transaction.CommitAsync();
- }
-}
diff --git a/src/EventStore.Core.Tests/ClientAPI/SpecificationWithMiniNode.cs b/src/EventStore.Core.Tests/ClientAPI/SpecificationWithMiniNode.cs
deleted file mode 100644
index ddd01c25d0..0000000000
--- a/src/EventStore.Core.Tests/ClientAPI/SpecificationWithMiniNode.cs
+++ /dev/null
@@ -1,102 +0,0 @@
-using System;
-using System.Threading.Tasks;
-using EventStore.ClientAPI;
-using EventStore.Core.Tests.ClientAPI.Helpers;
-using EventStore.Core.Tests.Helpers;
-using NUnit.Framework;
-
-namespace EventStore.Core.Tests.ClientAPI;
-
-public abstract class SpecificationWithMiniNode : SpecificationWithDirectoryPerTestFixture
-{
- private readonly int _chunkSize;
- protected MiniNode _node;
- protected IEventStoreConnection _conn;
- protected virtual TimeSpan Timeout { get; } = TimeSpan.FromMinutes(1);
- protected virtual TimeSpan StartupTimeout => TimeSpan.FromMinutes(5);
-
- protected virtual Task Given() => Task.CompletedTask;
-
- protected abstract Task When();
-
- protected virtual IEventStoreConnection BuildConnection(MiniNode node)
- {
- return TestConnection.CreateMiniNodeClient(node.TcpEndPoint, TcpType.Ssl);
- }
-
- protected Task CloseConnectionAndWait(IEventStoreConnection connection) =>
- TestConnectionLifecycle.CloseConnectionAndWait(connection, Timeout);
-
- protected SpecificationWithMiniNode() : this(chunkSize: 1024 * 1024) { }
-
- protected SpecificationWithMiniNode(int chunkSize)
- {
- _chunkSize = chunkSize;
- }
-
- [OneTimeSetUp]
- public override async Task TestFixtureSetUp()
- {
-
- MiniNodeLogging.Setup();
-
- try
- {
- await base.TestFixtureSetUp();
- }
- catch (Exception ex)
- {
- throw new Exception("TestFixtureSetUp Failed", ex);
- }
-
- try
- {
- _node = new MiniNode(PathName, chunkSize: _chunkSize);
- await _node.Start(StartupTimeout);
- await _node.WaitForTcpEndPoint().WithTimeout(StartupTimeout);
- _conn = await TestConnectionLifecycle.ReconnectUntilReady(
- () => BuildConnection(_node),
- connection => connection.ReadAllEventsForwardAsync(Position.Start, 1, false, DefaultData.AdminCredentials),
- StartupTimeout);
- }
- catch (Exception ex)
- {
- MiniNodeLogging.WriteLogs();
- throw new Exception("MiniNodeSetUp Failed", ex);
- }
-
- try
- {
- await Given().WithTimeout(Timeout);
- }
- catch (Exception ex)
- {
- MiniNodeLogging.WriteLogs();
- throw new Exception("Given Failed", ex);
- }
-
- try
- {
- await When().WithTimeout(Timeout);
- }
- catch (Exception ex)
- {
- MiniNodeLogging.WriteLogs();
- throw new Exception("When Failed", ex);
- }
- }
-
- [OneTimeTearDown]
- public override async Task TestFixtureTearDown()
- {
- if (_conn != null)
- {
- await TestConnectionLifecycle.CloseConnectionAndWait(_conn, Timeout);
- }
-
- await _node.Shutdown();
- await base.TestFixtureTearDown();
-
- MiniNodeLogging.Clear();
- }
-}
diff --git a/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs b/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs
index 5852fa8077..1a5b4da733 100644
--- a/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs
+++ b/src/EventStore.Core.Tests/Cluster/MemberInfoTests.cs
@@ -9,33 +9,19 @@ namespace EventStore.Core.Tests.Cluster;
[TestFixture]
public class MemberInfoTests
{
- private static readonly DnsEndPoint InternalTcp = new("internal", 1112);
- private static readonly DnsEndPoint InternalSecureTcp = new("internal-secure", 2112);
- private static readonly DnsEndPoint ExternalTcp = new("external", 1113);
- private static readonly DnsEndPoint ExternalSecureTcp = new("external-secure", 2113);
private static readonly DnsEndPoint Http = new("http", 2113);
- private static readonly DnsEndPoint Replication = new("replication", 3113);
+ private static readonly DnsEndPoint Cluster = new("cluster", 1112);
[Test]
public void member_with_dns_endpoint_should_equal()
{
var ipAddress = "127.0.0.1";
var port = 1113;
- var memberWithDnsEndPoint = EventStore.Core.Cluster.MemberInfo.Initial(Guid.Empty, DateTime.UtcNow,
- VNodeState.Unknown, true,
- new DnsEndPoint(ipAddress, port),
- new DnsEndPoint(ipAddress, port),
- new DnsEndPoint(ipAddress, port),
- new DnsEndPoint(ipAddress, port),
- new DnsEndPoint(ipAddress, port),
- null, 0, 0,
- 0, false);
-
- var ipEndPoint = new IPEndPoint(IPAddress.Parse(ipAddress), port);
- var dnsEndPoint = new DnsEndPoint(ipAddress, port);
-
- Assert.True(memberWithDnsEndPoint.Is(ipEndPoint));
- Assert.True(memberWithDnsEndPoint.Is(dnsEndPoint));
+ var member = EventStore.Core.Cluster.MemberInfo.Initial(Guid.Empty, DateTime.UtcNow,
+ VNodeState.Unknown, true, new DnsEndPoint(ipAddress, port), null, 0, 0, false);
+
+ Assert.That(member.Is(new IPEndPoint(IPAddress.Parse(ipAddress), port)), Is.True);
+ Assert.That(member.Is(new DnsEndPoint(ipAddress, port)), Is.True);
}
[Test]
@@ -43,87 +29,54 @@ public void member_with_ip_endpoint_should_equal()
{
var ipAddress = "127.0.0.1";
var port = 1113;
- var memberWithDnsEndPoint = EventStore.Core.Cluster.MemberInfo.Initial(Guid.Empty, DateTime.UtcNow,
- VNodeState.Unknown, true,
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- null, 0, 0, 0, false);
-
- var ipEndPoint = new IPEndPoint(IPAddress.Parse(ipAddress), port);
- var dnsEndPoint = new DnsEndPoint(ipAddress, port);
-
- Assert.True(memberWithDnsEndPoint.Is(ipEndPoint));
- Assert.True(memberWithDnsEndPoint.Is(dnsEndPoint));
+ var member = EventStore.Core.Cluster.MemberInfo.Initial(Guid.Empty, DateTime.UtcNow,
+ VNodeState.Unknown, true, new IPEndPoint(IPAddress.Parse(ipAddress), port), null, 0, 0, false);
+
+ Assert.That(member.Is(new IPEndPoint(IPAddress.Parse(ipAddress), port)), Is.True);
+ Assert.That(member.Is(new DnsEndPoint(ipAddress, port)), Is.True);
}
[Test]
- public void grpc_round_trip_preserves_tcp_and_cluster_endpoints()
+ public void grpc_round_trip_preserves_client_and_cluster_endpoints()
{
- var member = EventStore.Core.Cluster.MemberInfo.Initial(Guid.NewGuid(), DateTime.UtcNow,
- VNodeState.Unknown, true,
- InternalTcp, null, null, ExternalSecureTcp, Http,
- "client", 2113, 1113, 0, false, clusterEndPoint: Replication);
-
var result = FromGrpcClusterInfo(ToGrpcClusterInfo(
- new EventStore.Core.Cluster.ClusterInfo(member))).Members[0];
+ new EventStore.Core.Cluster.ClusterInfo(CreateMember(Cluster)))).Members[0];
- Assert.That(result.InternalTcpEndPoint, Is.EqualTo(InternalTcp));
- Assert.That(result.InternalSecureTcpEndPoint, Is.Null);
- Assert.That(result.ExternalTcpEndPoint, Is.Null);
- Assert.That(result.ExternalSecureTcpEndPoint, Is.EqualTo(ExternalSecureTcp));
Assert.That(result.HttpEndPoint, Is.EqualTo(Http));
- Assert.That(result.ClusterEndPoint, Is.EqualTo(Replication));
+ Assert.That(result.ClusterEndPoint, Is.EqualTo(Cluster));
+ Assert.That(result.ReplicationEndPoint, Is.EqualTo(Cluster));
}
[Test]
- public void explicit_cluster_endpoint_is_recognized_without_replacing_tcp_endpoints()
+ public void explicit_cluster_endpoint_is_recognized()
{
- var member = CreateMember(Replication);
- var vnode = new VNodeInfo(Guid.NewGuid(), 0,
- new IPEndPoint(IPAddress.Loopback, 1112), null,
- new IPEndPoint(IPAddress.Loopback, 1113), null,
- Http, false, Replication);
- var advertise = new GossipAdvertiseInfo(
- InternalTcp, InternalSecureTcp, ExternalTcp, ExternalSecureTcp, Http,
- null, null, 0, null, 0, 0, Replication);
-
- Assert.That(member.Is(Replication), Is.True);
- Assert.That(member.InternalTcpEndPoint, Is.EqualTo(InternalTcp));
- Assert.That(member.InternalSecureTcpEndPoint, Is.EqualTo(InternalSecureTcp));
- Assert.That(member.ExternalTcpEndPoint, Is.EqualTo(ExternalTcp));
- Assert.That(member.ExternalSecureTcpEndPoint, Is.EqualTo(ExternalSecureTcp));
- Assert.That(vnode.ClusterEndPoint, Is.SameAs(Replication));
- Assert.That(advertise.ClusterEndPoint, Is.SameAs(Replication));
+ var member = CreateMember(Cluster);
+ var vnode = new VNodeInfo(Guid.NewGuid(), 0, Http, false, Cluster);
+ var advertise = new GossipAdvertiseInfo(Http, null, 0, Cluster);
+
+ Assert.That(member.Is(Cluster), Is.True);
+ Assert.That(vnode.ClusterEndPoint, Is.SameAs(Cluster));
+ Assert.That(advertise.ClusterEndPoint, Is.SameAs(Cluster));
}
[Test]
public void client_member_preserves_the_cluster_endpoint()
{
- var clientMember = new EventStore.Core.Cluster.ClientClusterInfo.ClientMemberInfo(
- CreateMember(Replication));
+ var clientMember = new EventStore.Core.Cluster.ClientClusterInfo.ClientMemberInfo(CreateMember(Cluster));
- Assert.That(clientMember.ClusterEndPointIp, Is.EqualTo(Replication.Host));
- Assert.That(clientMember.ClusterEndPointPort, Is.EqualTo(Replication.Port));
+ Assert.That(clientMember.ClusterEndPointIp, Is.EqualTo(Cluster.Host));
+ Assert.That(clientMember.ClusterEndPointPort, Is.EqualTo(Cluster.Port));
}
[Test]
public void client_cluster_info_excludes_internal_discovery_placeholders()
{
- var member = CreateMember(Replication);
+ var member = CreateMember(Cluster);
var seed = EventStore.Core.Cluster.MemberInfo.ForManager(
- Guid.Empty,
- DateTime.UtcNow,
- true,
- Replication,
- clusterEndPoint: Replication);
+ Guid.Empty, DateTime.UtcNow, true, Cluster, clusterEndPoint: Cluster);
var clientCluster = new EventStore.Core.Cluster.ClientClusterInfo(
- new EventStore.Core.Cluster.ClusterInfo(member, seed),
- Http.Host,
- Http.Port);
+ new EventStore.Core.Cluster.ClusterInfo(member, seed), Http.Host, Http.Port);
Assert.That(clientCluster.Members, Has.Length.EqualTo(1));
Assert.That(clientCluster.Members[0].InstanceId, Is.EqualTo(member.InstanceId));
@@ -133,13 +86,8 @@ public void client_cluster_info_excludes_internal_discovery_placeholders()
public void missing_cluster_endpoint_falls_back_to_http_endpoint()
{
var member = CreateMember();
- var vnode = new VNodeInfo(Guid.NewGuid(), 0,
- new IPEndPoint(IPAddress.Loopback, 1112), null,
- new IPEndPoint(IPAddress.Loopback, 1113), null,
- Http, false);
- var advertise = new GossipAdvertiseInfo(
- InternalTcp, InternalSecureTcp, ExternalTcp, ExternalSecureTcp, Http,
- null, null, 0, null, 0, 0);
+ var vnode = new VNodeInfo(Guid.NewGuid(), 0, Http, false);
+ var advertise = new GossipAdvertiseInfo(Http, null, 0);
Assert.That(member.ClusterEndPoint, Is.SameAs(Http));
Assert.That(vnode.ClusterEndPoint, Is.SameAs(Http));
@@ -150,7 +98,7 @@ public void missing_cluster_endpoint_falls_back_to_http_endpoint()
public void grpc_member_without_cluster_endpoint_falls_back_to_http_endpoint()
{
var grpcCluster = ToGrpcClusterInfo(
- new EventStore.Core.Cluster.ClusterInfo(CreateMember(Replication)));
+ new EventStore.Core.Cluster.ClusterInfo(CreateMember(Cluster)));
grpcCluster.Members[0].ReplicationEndPoint = null;
var result = FromGrpcClusterInfo(grpcCluster).Members[0];
@@ -160,9 +108,7 @@ public void grpc_member_without_cluster_endpoint_falls_back_to_http_endpoint()
private static EventStore.Core.Cluster.MemberInfo CreateMember(DnsEndPoint clusterEndPoint = null) =>
EventStore.Core.Cluster.MemberInfo.Initial(Guid.NewGuid(), DateTime.UtcNow,
- VNodeState.Unknown, true,
- InternalTcp, InternalSecureTcp, ExternalTcp, ExternalSecureTcp, Http,
- "client", 2113, 1113, 0, false, clusterEndPoint: clusterEndPoint);
+ VNodeState.Unknown, true, Http, null, 0, 0, false, clusterEndPoint: clusterEndPoint);
private static EventStore.Cluster.ClusterInfo ToGrpcClusterInfo(
EventStore.Core.Cluster.ClusterInfo clusterInfo) =>
diff --git a/src/EventStore.Core.Tests/DefaultData.cs b/src/EventStore.Core.Tests/DefaultData.cs
index 3c647b08fc..262f7312dc 100644
--- a/src/EventStore.Core.Tests/DefaultData.cs
+++ b/src/EventStore.Core.Tests/DefaultData.cs
@@ -1,5 +1,4 @@
using System.Net;
-using EventStore.ClientAPI.SystemData;
using EventStore.Core.Services;
namespace EventStore.Core.Tests;
@@ -8,7 +7,6 @@ public class DefaultData
{
public static string AdminUsername = SystemUsers.Admin;
public static string AdminPassword = SystemUsers.DefaultAdminPassword;
- public static UserCredentials AdminCredentials = new UserCredentials(AdminUsername, AdminPassword);
public static NetworkCredential AdminNetworkCredentials = new NetworkCredential(AdminUsername, AdminPassword);
public static ClusterVNodeOptions.DefaultUserOptions DefaultUserOptions = new ClusterVNodeOptions.DefaultUserOptions()
{
diff --git a/src/EventStore.Core.Tests/EventStore.Core.Tests.csproj b/src/EventStore.Core.Tests/EventStore.Core.Tests.csproj
index 3be13ef772..4aee9b0985 100644
--- a/src/EventStore.Core.Tests/EventStore.Core.Tests.csproj
+++ b/src/EventStore.Core.Tests/EventStore.Core.Tests.csproj
@@ -4,7 +4,6 @@
true
-
diff --git a/src/EventStore.Core.Tests/Helpers/ClientApiLoggerBridge.cs b/src/EventStore.Core.Tests/Helpers/ClientApiLoggerBridge.cs
deleted file mode 100644
index 862b03493f..0000000000
--- a/src/EventStore.Core.Tests/Helpers/ClientApiLoggerBridge.cs
+++ /dev/null
@@ -1,92 +0,0 @@
-using System;
-using EventStore.Common.Utils;
-using ILogger = Serilog.ILogger;
-
-namespace EventStore.Core.Tests.Helpers;
-
-public class ClientApiLoggerBridge : EventStore.ClientAPI.ILogger
-{
- public static readonly ClientApiLoggerBridge Default =
- new ClientApiLoggerBridge(Serilog.Log.ForContext(Serilog.Core.Constants.SourceContextPropertyName,
- "client-api"));
-
- private readonly Serilog.ILogger _log;
-
- public ClientApiLoggerBridge(ILogger log)
- {
- Ensure.NotNull(log, "log");
- _log = log;
- }
-
- public void Error(string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Error(format);
- }
- else
- {
- _log.Error(format, args);
- }
- }
-
- public void Error(Exception ex, string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Error(ex, format);
- }
- else
- {
- _log.Error(ex, format, args);
- }
- }
-
- public void Info(string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Information(format);
- }
- else
- {
- _log.Information(format, args);
- }
- }
-
- public void Info(Exception ex, string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Information(ex, format);
- }
- else
- {
- _log.Information(ex, format, args);
- }
- }
-
- public void Debug(string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Debug(format);
- }
- else
- {
- _log.Debug(format, args);
- }
- }
-
- public void Debug(Exception ex, string format, params object[] args)
- {
- if (args.Length == 0)
- {
- _log.Debug(ex, format);
- }
- else
- {
- _log.Debug(ex, format, args);
- }
- }
-}
diff --git a/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs b/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs
index e84c40b2c8..fe6634efea 100644
--- a/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs
+++ b/src/EventStore.Core.Tests/Helpers/MiniClusterNode.cs
@@ -20,10 +20,8 @@
using EventStore.Core.Services.Monitoring;
using EventStore.Core.Services.PersistentSubscription.ConsumerStrategy;
using EventStore.Core.Services.Storage.ReaderIndex;
-using EventStore.Core.Tests.Services.Transport.Tcp;
using EventStore.Core.TransactionLog.Chunks;
using EventStore.Plugins.Subsystems;
-using EventStore.TcpUnitTestPlugin;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Server.Kestrel.Core;
using Microsoft.AspNetCore.Server.Kestrel.Https;
@@ -43,7 +41,6 @@ public class MiniClusterNode
private static readonly ILogger Log = Serilog.Log.ForContext>();
- public IPEndPoint ExternalTcpEndPoint { get; }
public IPEndPoint HttpEndPoint { get; }
public IPEndPoint ClusterEndPoint { get; }
@@ -63,33 +60,31 @@ public class MiniClusterNode
public VNodeState NodeState = VNodeState.Unknown;
private readonly IHost _host;
- public MiniClusterNode(string pathname, int debugIndex, IPEndPoint clusterEndPoint, IPEndPoint externalTcp,
- IPEndPoint httpEndPoint, EndPoint[] gossipSeeds, ISubsystem[] subsystems = null,
+ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint httpEndPoint, IPEndPoint clusterEndPoint,
+ EndPoint[] gossipSeeds, ISubsystem[] subsystems = null,
bool enableTrustedAuth = false, int memTableSize = 1000,
bool disableFlushToDisk = false, bool readOnlyReplica = false, int nodePriority = 0,
- string intHostAdvertiseAs = null, IExpiryStrategy expiryStrategy = null,
+ string replicationHostAdvertiseAs = null,
+ IExpiryStrategy expiryStrategy = null,
ArchiveOptions archiveOptions = null, bool archiver = false,
- int clusterSize = 3, bool unsafeAllowSurplusNodes = false,
- string replicationHostAdvertiseAs = null)
+ int clusterSize = 3, bool unsafeAllowSurplusNodes = false)
{
RunningTime.Start();
RunCount += 1;
DebugIndex = debugIndex;
- ExternalTcpEndPoint = externalTcp;
HttpEndPoint = httpEndPoint;
ClusterEndPoint = clusterEndPoint;
_dbPath = Path.Combine(
pathname,
- $"mini-cluster-node-db-{externalTcp.Port}-{httpEndPoint.Port}");
+ $"mini-cluster-node-db-{httpEndPoint.Port}");
Directory.CreateDirectory(_dbPath);
FileStreamExtensions.ConfigureFlush(disableFlushToDisk);
subsystems ??= [];
- subsystems = [.. subsystems, new TcpApiTestPlugin()];
var options = new ClusterVNodeOptions
{
@@ -119,14 +114,14 @@ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint clusterEndPoi
},
Interface = new()
{
+ NodeIp = HttpEndPoint.Address,
+ NodePort = HttpEndPoint.Port,
ReplicationIp = ClusterEndPoint.Address,
- NodeIp = ExternalTcpEndPoint.Address,
ReplicationPort = ClusterEndPoint.Port,
- NodePort = HttpEndPoint.Port,
ReplicationHeartbeatTimeout = 2_000,
ReplicationHeartbeatInterval = 2_000,
- EnableTrustedAuth = enableTrustedAuth,
- ReplicationHostAdvertiseAs = replicationHostAdvertiseAs ?? intHostAdvertiseAs
+ ReplicationHostAdvertiseAs = replicationHostAdvertiseAs,
+ EnableTrustedAuth = enableTrustedAuth
},
Database = new()
{
@@ -151,14 +146,7 @@ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint clusterEndPoi
PlugableComponents = subsystems
};
- var configuration = new List> {
- new("EventStore:TcpPlugin:NodeTcpPort", externalTcp.Port.ToString()),
- new("EventStore:TcpPlugin:EnableExternalTcp", "true"),
- new("EventStore:TcpUnitTestPlugin:NodeTcpPort", externalTcp.Port.ToString()),
- new("EventStore:TcpUnitTestPlugin:NodeHeartbeatInterval", "10000"),
- new("EventStore:TcpUnitTestPlugin:NodeHeartbeatTimeout", "10000"),
- new("EventStore:TcpUnitTestPlugin:Insecure", options.Application.Insecure.ToString()),
- };
+ var configuration = new List>();
if (archiveOptions is not null)
{
@@ -178,9 +166,9 @@ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint clusterEndPoi
var inMemConf = new ConfigurationBuilder()
.AddInMemoryCollection(configuration)
.Build();
- var serverCertificate = ssl_connections.GetServerCertificate();
+ var serverCertificate = TestCertificates.GetServerCertificate();
var trustedRootCertificates =
- new X509Certificate2Collection(ssl_connections.GetRootCertificate());
+ new X509Certificate2Collection(TestCertificates.GetRootCertificate());
options = options.Secure(trustedRootCertificates, serverCertificate);
_isReadOnlyReplica = readOnlyReplica;
@@ -193,8 +181,8 @@ public MiniClusterNode(string pathname, int debugIndex, IPEndPoint clusterEndPoi
RuntimeInformation.RuntimeMode, "GC:",
GC.MaxGeneration == 0
? "NON-GENERATION (PROBABLY BOEHM)"
- : $"{GC.MaxGeneration + 1} GENERATIONS", "DBPATH:", _dbPath, "ExTCP ENDPOINT:",
- ExternalTcpEndPoint, "ExHTTP ENDPOINT:", HttpEndPoint);
+ : $"{GC.MaxGeneration + 1} GENERATIONS", "DBPATH:", _dbPath, "NODE ENDPOINT:",
+ HttpEndPoint, "HTTP ENDPOINT:", HttpEndPoint);
var logFormatFactory = LogFormatHelper.LogFormatFactory;
Node = new ClusterVNode(options, logFormatFactory, new AuthenticationProviderFactory(
diff --git a/src/EventStore.Core.Tests/Helpers/MiniNode.cs b/src/EventStore.Core.Tests/Helpers/MiniNode.cs
index 9a41dc772b..95dd4d7d5e 100644
--- a/src/EventStore.Core.Tests/Helpers/MiniNode.cs
+++ b/src/EventStore.Core.Tests/Helpers/MiniNode.cs
@@ -5,7 +5,6 @@
using System.Linq;
using System.Net;
using System.Net.Http;
-using System.Net.Sockets;
using System.Security.Cryptography.X509Certificates;
using System.Threading.Tasks;
using EventStore.Common.Utils;
@@ -20,13 +19,11 @@
using EventStore.Core.Services.Monitoring;
using EventStore.Core.Services.Storage.ReaderIndex;
using EventStore.Core.Tests.Index.Hashers;
-using EventStore.Core.Tests.Services.Transport.Tcp;
using EventStore.Core.TransactionLog.Chunks;
using EventStore.Plugins.Authentication;
using EventStore.Plugins.Authorization;
using EventStore.Plugins.Subsystems;
using EventStore.Plugins.Transforms;
-using EventStore.TcpUnitTestPlugin;
using Microsoft.AspNetCore.Hosting;
using Microsoft.AspNetCore.Server.Kestrel.Core;
using Microsoft.AspNetCore.Server.Kestrel.Https;
@@ -45,8 +42,6 @@ public class MiniNode
public const int CachedChunkSize = ChunkSize + ChunkHeader.Size + ChunkFooter.Size;
protected static readonly ILogger Log = Serilog.Log.ForContext();
- public IPEndPoint TcpEndPoint { get; protected set; }
- public IPEndPoint IntTcpEndPoint { get; protected set; }
public IPEndPoint HttpEndPoint { get; protected set; }
}
@@ -71,7 +66,7 @@ public class MiniNode : MiniNode, IAsyncDisposable
public Task AdminUserCreated => _adminUserCreated.Task;
public MiniNode(string pathname,
- int? tcpPort = null, int? httpPort = null,
+ int? httpPort = null,
ISubsystem[] subsystems = null,
int chunkSize = ChunkSize, int cachedChunkSize = CachedChunkSize, bool enableTrustedAuth = false,
int memTableSize = 1000,
@@ -97,26 +92,21 @@ public MiniNode(string pathname,
var ip = IPAddress.Loopback;
- int extTcpPort = tcpPort ?? PortsHelper.GetAvailablePort(ip);
int httpEndPointPort = httpPort ?? PortsHelper.GetAvailablePort(ip);
- int intTcpPort = PortsHelper.GetAvailablePort(ip);
if (string.IsNullOrEmpty(dbPath))
{
DbPath = Path.Combine(pathname,
- $"mini-node-db-{extTcpPort}-{httpEndPointPort}");
+ $"mini-node-db-{httpEndPointPort}");
}
else
{
DbPath = dbPath;
}
- TcpEndPoint = new IPEndPoint(ip, extTcpPort);
- IntTcpEndPoint = new IPEndPoint(ip, intTcpPort);
HttpEndPoint = new IPEndPoint(ip, httpEndPointPort);
subsystems ??= [];
- subsystems = [.. subsystems, new TcpApiTestPlugin()];
var options = new ClusterVNodeOptions
{
@@ -130,8 +120,6 @@ public MiniNode(string pathname,
},
Interface = new()
{
- ReplicationHeartbeatInterval = 10_000,
- ReplicationHeartbeatTimeout = 10_000,
EnableTrustedAuth = enableTrustedAuth
},
Cluster = new()
@@ -162,21 +150,11 @@ public MiniNode(string pathname,
LoadedOptions = ClusterVNodeOptions.GetLoadedOptions(new ConfigurationBuilder()
.AddEventStoreDefaultValues()
.Build()),
- }.Secure(new X509Certificate2Collection(ssl_connections.GetRootCertificate()),
- ssl_connections.GetServerCertificate())
- .WithReplicationEndpointOn(IntTcpEndPoint)
- .WithExternalTcpOn(TcpEndPoint)
+ }.Secure(new X509Certificate2Collection(TestCertificates.GetRootCertificate()),
+ TestCertificates.GetServerCertificate())
.WithNodeEndpointOn(HttpEndPoint);
- var inMemConf = new ConfigurationBuilder()
- .AddInMemoryCollection(new KeyValuePair[] {
- new("EventStore:TcpPlugin:NodeTcpPort", extTcpPort.ToString()),
- new("EventStore:TcpPlugin:EnableExternalTcp", "true"),
- new("EventStore:TcpUnitTestPlugin:NodeTcpPort", extTcpPort.ToString()),
- new("EventStore:TcpUnitTestPlugin:NodeHeartbeatInterval", "10000"),
- new("EventStore:TcpUnitTestPlugin:NodeHeartbeatTimeout", "10000"),
- new("EventStore:TcpUnitTestPlugin:Insecure", options.Application.Insecure.ToString()),
- }).Build();
+ var inMemConf = new ConfigurationBuilder().Build();
if (advertisedExtHostAddress != null)
{
@@ -200,7 +178,7 @@ public MiniNode(string pathname,
? "NON-GENERATION (PROBABLY BOEHM)"
: $"{GC.MaxGeneration + 1} GENERATIONS",
"DBPATH:", DbPath,
- "TCP ENDPOINT:", TcpEndPoint,
+ "NODE ENDPOINT:", HttpEndPoint,
"HTTP ENDPOINT:", HttpEndPoint);
var logFormatFactory = LogFormatHelper.LogFormatFactory
@@ -244,12 +222,12 @@ public MiniNode(string pathname,
{
options.UseHttps(new HttpsConnectionAdapterOptions
{
- ServerCertificate = ssl_connections.GetServerCertificate(),
+ ServerCertificate = TestCertificates.GetServerCertificate(),
ClientCertificateMode = ClientCertificateMode.AllowCertificate,
ClientCertificateValidation = (certificate, chain, sslPolicyErrors) =>
{
var (isValid, error) =
- ClusterVNode.ValidateClientCertificate(certificate, chain, sslPolicyErrors, () => null, () => new X509Certificate2Collection(ssl_connections.GetRootCertificate()));
+ ClusterVNode.ValidateClientCertificate(certificate, chain, sslPolicyErrors, () => null, () => new X509Certificate2Collection(TestCertificates.GetRootCertificate()));
if (!isValid && error != null)
{
Log.Error("Client certificate validation error: {e}", error);
@@ -335,25 +313,6 @@ void WaitForAdminUser(StorageMessage.EventCommitted m)
Log.Information("MiniNode successfully started!");
}
- public async Task WaitForTcpEndPoint()
- {
- while (true)
- {
- using var client = new TcpClient();
-
- try
- {
- await client.ConnectAsync(TcpEndPoint.Address, TcpEndPoint.Port)
- .WaitAsync(TimeSpan.FromMilliseconds(250));
- return;
- }
- catch (Exception ex) when (ex is SocketException or TimeoutException)
- {
- await Task.Delay(100);
- }
- }
- }
-
public async Task Shutdown(bool keepDb = false)
{
diff --git a/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs b/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs
index 888cac4935..564a800c2e 100644
--- a/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs
+++ b/src/EventStore.Core.Tests/Integration/Archive/when_archiving_and_restoring_a_cluster.cs
@@ -2,11 +2,12 @@
using System.IO;
using System.Linq;
using System.Net;
+using System.Net.Http;
using System.Threading;
using System.Threading.Tasks;
using Amazon.S3;
using Amazon.S3.Model;
-using EventStore.ClientAPI;
+using EventStore.Client.Streams;
using EventStore.Core.Data;
using EventStore.Core.Messages;
using EventStore.Core.Messaging;
@@ -18,7 +19,11 @@
using EventStore.Core.Tests.Helpers;
using EventStore.Core.TransactionLog.Chunks.TFChunk;
using EventStore.Core.TransactionLog.FileNamingStrategy;
+using Google.Protobuf;
+using Grpc.Net.Client;
using NUnit.Framework;
+using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata;
+using StreamsClient = EventStore.Client.Streams.Streams.StreamsClient;
namespace EventStore.Core.Tests.Integration.Archive;
@@ -44,6 +49,8 @@ public class when_archiving_and_restoring_a_cluster
private long _archivedCheckpoint;
private int _restoredNodeIndex;
private int _completedIterations;
+ private GrpcChannel _channel;
+ private StreamsClient _client;
protected override int NodeCount => 4;
protected override TimeSpan GivenTimeout => SoakTimeout;
@@ -117,19 +124,13 @@ protected override MiniClusterNode CreateNode(
new(
PathName,
index,
- endpoints.ClusterEndPoint,
- endpoints.ExternalTcp,
endpoints.HttpEndPoint,
+ endpoints.ClusterEndPoint,
gossipSeeds,
readOnlyReplica: index == ArchiverNodeIndex,
archiveOptions: _archiveOptions.Enabled ? _archiveOptions : null,
archiver: index == ArchiverNodeIndex);
- protected override IEventStoreConnection CreateConnection() =>
- EventStoreConnection.Create(
- ConnectionSettings.Create().DisableServerCertificateValidation(),
- GetLeader().ExternalTcpEndPoint);
-
protected override async Task Given()
{
var payload = new byte[256 * 1024];
@@ -141,10 +142,30 @@ protected override async Task Given()
var leader = await ReconnectToLeader();
for (var eventNumber = 0; eventNumber < EventsPerIteration; eventNumber++)
{
- await _conn.AppendToStreamAsync(
- Stream,
- EventStore.ClientAPI.ExpectedVersion.Any,
- new EventData(Guid.NewGuid(), "archive-event", isJson: false, payload, Array.Empty()));
+ using var call = _client.Append();
+ await call.RequestStream.WriteAsync(new AppendReq
+ {
+ Options = new()
+ {
+ Any = new(),
+ StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(Stream) }
+ }
+ });
+ await call.RequestStream.WriteAsync(new AppendReq
+ {
+ ProposedMessage = new()
+ {
+ Id = Core.Services.Transport.Grpc.Uuid.NewUuid().ToDto(),
+ Data = ByteString.CopyFrom(payload),
+ CustomMetadata = ByteString.Empty,
+ Metadata = {
+ { GrpcMetadata.Type, "archive-event" },
+ { GrpcMetadata.ContentType, GrpcMetadata.ContentTypes.ApplicationOctetStream }
+ }
+ }
+ });
+ await call.RequestStream.CompleteAsync();
+ await call.ResponseAsync;
}
AssertEx.IsOrBecomesTrue(
@@ -175,11 +196,16 @@ await _conn.AppendToStreamAsync(
private async Task> ReconnectToLeader()
{
var leader = GetLeader();
- _conn?.Close();
- _conn = EventStoreConnection.Create(
- ConnectionSettings.Create().DisableServerCertificateValidation(),
- leader.ExternalTcpEndPoint);
- await _conn.ConnectAsync();
+ _channel?.Dispose();
+ _channel = GrpcChannel.ForAddress(new Uri($"https://{leader.HttpEndPoint}"),
+ new GrpcChannelOptions
+ {
+ HttpHandler = new SocketsHttpHandler
+ {
+ SslOptions = { RemoteCertificateValidationCallback = delegate { return true; } }
+ }
+ });
+ _client = new StreamsClient(_channel);
return leader;
}
@@ -258,6 +284,7 @@ private async Task WaitForArchiveCheckpoint(long minimum)
[OneTimeTearDown]
public override async Task TestFixtureTearDown()
{
+ _channel?.Dispose();
await base.TestFixtureTearDown();
if (_s3Client is null)
{
diff --git a/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs b/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs
index a5b20c79c8..4bfde56f1e 100644
--- a/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs
+++ b/src/EventStore.Core.Tests/Integration/specification_with_cluster.cs
@@ -4,7 +4,6 @@
using System.Net;
using System.Net.Sockets;
using System.Threading.Tasks;
-using EventStore.ClientAPI;
using EventStore.Core.Data;
using EventStore.Core.Tests.Helpers;
using EventStore.Plugins.Subsystems;
@@ -17,7 +16,6 @@ public abstract class specification_with_cluster : Specif
{
protected MiniClusterNode[] _nodes;
protected Endpoints[] _nodeEndpoints;
- protected IEventStoreConnection _conn;
protected virtual TimeSpan GivenTimeout { get; } = TimeSpan.FromMinutes(2);
protected virtual int NodeCount => 3;
@@ -26,13 +24,11 @@ public abstract class specification_with_cluster : Specif
protected class Endpoints
{
public readonly IPEndPoint ClusterEndPoint;
- public readonly IPEndPoint ExternalTcp;
public readonly IPEndPoint HttpEndPoint;
public IEnumerable Ports()
{
yield return ClusterEndPoint.Port;
- yield return ExternalTcp.Port;
yield return HttpEndPoint.Port;
}
@@ -48,16 +44,11 @@ public Endpoints()
cluster.Bind(defaultLoopBack);
_sockets.Add(cluster);
- var externalTcp = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
- externalTcp.Bind(defaultLoopBack);
- _sockets.Add(externalTcp);
-
var httpEndPoint = new Socket(AddressFamily.InterNetwork, SocketType.Stream, ProtocolType.Tcp);
httpEndPoint.Bind(defaultLoopBack);
_sockets.Add(httpEndPoint);
ClusterEndPoint = CopyEndpoint((IPEndPoint)cluster.LocalEndPoint);
- ExternalTcp = CopyEndpoint((IPEndPoint)externalTcp.LocalEndPoint);
HttpEndPoint = CopyEndpoint((IPEndPoint)httpEndPoint.LocalEndPoint);
}
@@ -144,9 +135,6 @@ public override async Task TestFixtureSetUp()
onFail: MiniNodeLogging.WriteLogs,
msg: $"Waiting for followers timed out! States={string.Join(", ", _nodes.Select(n => n.NodeState))}");
- _conn = CreateConnection();
- await _conn.ConnectAsync();
-
try
{
await Given().WithTimeout(GivenTimeout);
@@ -158,9 +146,6 @@ public override async Task TestFixtureSetUp()
}
}
- protected virtual IEventStoreConnection CreateConnection() =>
- EventStoreConnection.Create(_nodes[0].ExternalTcpEndPoint);
-
protected virtual void BeforeNodesStart()
{
}
@@ -171,8 +156,7 @@ protected virtual void BeforeNodesStart()
protected virtual MiniClusterNode CreateNode(int index, Endpoints endpoints, EndPoint[] gossipSeeds,
bool wait = true) => new(
- PathName, index, endpoints.ClusterEndPoint,
- endpoints.ExternalTcp, endpoints.HttpEndPoint,
+ PathName, index, endpoints.HttpEndPoint, endpoints.ClusterEndPoint,
subsystems: Array.Empty(), gossipSeeds: gossipSeeds);
[TearDown]
@@ -187,7 +171,6 @@ public void AfterEachTest()
[OneTimeTearDown]
public override async Task TestFixtureTearDown()
{
- _conn?.Close();
if (_nodes is not null)
{
await Task.WhenAll(_nodes.Where(node => node is not null).Select(node => node.Shutdown()));
diff --git a/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs b/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs
index 2bef053672..dd4de5c12f 100644
--- a/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs
+++ b/src/EventStore.Core.Tests/Integration/when_node_becomes_leader_with_unindexed_data.cs
@@ -164,8 +164,7 @@ private async Task> ReadAllEvents(IPEndPoint endpoint
private MiniClusterNode CreateNode(int index, Endpoints endpoints, EndPoint[] gossipSeeds,
int nodePriority, string replicationHostAdvertiseAs) => new(
- PathName, index, endpoints.ClusterEndPoint,
- endpoints.ExternalTcp, endpoints.HttpEndPoint,
+ PathName, index, endpoints.HttpEndPoint, endpoints.ClusterEndPoint,
subsystems: Array.Empty(), gossipSeeds: gossipSeeds,
nodePriority: nodePriority, replicationHostAdvertiseAs: replicationHostAdvertiseAs);
diff --git a/src/EventStore.Core.Tests/Regression/ClusterStatusServiceTests.cs b/src/EventStore.Core.Tests/Regression/ClusterStatusServiceTests.cs
index b82dc0fe58..e462947b86 100644
--- a/src/EventStore.Core.Tests/Regression/ClusterStatusServiceTests.cs
+++ b/src/EventStore.Core.Tests/Regression/ClusterStatusServiceTests.cs
@@ -19,7 +19,7 @@ public void replication_statistics_match_the_member_cluster_endpoint()
};
var result = (ClientClusterInfo.ClientMemberInfo)typeof(ClusterStatusService)
- .GetMethod("FindMemberByInternalEndpoint", BindingFlags.NonPublic | BindingFlags.Static)!
+ .GetMethod("FindMemberByEndpoint", BindingFlags.NonPublic | BindingFlags.Static)!
.Invoke(null, [new List { member }, "replica.internal:1112"])!;
Assert.That(result, Is.SameAs(member));
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/ClusterSettingsFactory.cs b/src/EventStore.Core.Tests/Services/ElectionsService/ClusterSettingsFactory.cs
index f63c8c9319..b072443420 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/ClusterSettingsFactory.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/ClusterSettingsFactory.cs
@@ -1,9 +1,7 @@
using System;
using System.Linq;
using System.Net;
-using System.Runtime.InteropServices;
using EventStore.Core.Cluster.Settings;
-using EventStore.Core.Tests.Services.Transport.Tcp;
namespace EventStore.Core.Tests.Services.ElectionsService;
@@ -14,13 +12,9 @@ public class ClusterSettingsFactory
private static ClusterVNodeSettings CreateVNode(int nodeNumber, bool isReadOnlyReplica)
{
- int tcpIntPort = StartingPort + nodeNumber * 2,
- tcpExtPort = tcpIntPort + 1,
- httpPort = tcpIntPort + 11;
+ var httpPort = StartingPort + nodeNumber;
return new ClusterVNodeSettings(Guid.NewGuid(), 0,
- GetLoopbackForPort(tcpIntPort), null,
- GetLoopbackForPort(tcpExtPort), null,
GetLoopbackForPort(httpPort), 0,
isReadOnlyReplica);
}
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/ClusterVNodeSettings.cs b/src/EventStore.Core.Tests/Services/ElectionsService/ClusterVNodeSettings.cs
index e1fb14badc..b9c4a888d9 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/ClusterVNodeSettings.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/ClusterVNodeSettings.cs
@@ -16,24 +16,14 @@ public class ClusterVNodeSettings
public readonly bool ReadOnlyReplica;
public ClusterVNodeSettings(Guid instanceId, int debugIndex,
- IPEndPoint internalTcpEndPoint,
- IPEndPoint internalSecureTcpEndPoint,
- IPEndPoint externalTcpEndPoint,
- IPEndPoint externalSecureTcpEndPoint,
IPEndPoint httpEndPoint,
int nodePriority,
bool readOnlyReplica)
{
Ensure.NotEmptyGuid(instanceId, "instanceId");
- Ensure.Equal(false, internalTcpEndPoint == null && internalSecureTcpEndPoint == null, "Both internal TCP endpoints are null");
-
Ensure.NotNull(httpEndPoint, nameof(httpEndPoint));
- NodeInfo = new VNodeInfo(instanceId, debugIndex,
- internalTcpEndPoint, internalSecureTcpEndPoint,
- externalTcpEndPoint, externalSecureTcpEndPoint,
- httpEndPoint,
- readOnlyReplica);
+ NodeInfo = new VNodeInfo(instanceId, debugIndex, httpEndPoint, readOnlyReplica);
NodePriority = nodePriority;
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/ElectionServiceUnit.cs b/src/EventStore.Core.Tests/Services/ElectionsService/ElectionServiceUnit.cs
index b24a1e875e..a07c837474 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/ElectionServiceUnit.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/ElectionServiceUnit.cs
@@ -47,11 +47,7 @@ public ElectionsServiceUnit(ClusterSettings clusterSettings)
_bus = new(GetType().Name);
var memberInfo = MemberInfo.Initial(clusterSettings.Self.NodeInfo.InstanceId, InitialDate,
VNodeState.Unknown, true,
- clusterSettings.Self.NodeInfo.InternalTcp,
- clusterSettings.Self.NodeInfo.InternalSecureTcp,
- clusterSettings.Self.NodeInfo.ExternalTcp,
- clusterSettings.Self.NodeInfo.ExternalSecureTcp,
- clusterSettings.Self.NodeInfo.HttpEndPoint, null, 0, 0,
+ clusterSettings.Self.NodeInfo.HttpEndPoint, null, 0,
clusterSettings.Self.NodePriority,
clusterSettings.Self.ReadOnlyReplica);
ElectionsService = new Core.Services.ElectionsService(Publisher,
@@ -83,11 +79,7 @@ private ClusterInfo BuildClusterInfo(ClusterSettings clusterSettings)
InitialDate,
VNodeState.Unknown,
true,
- clusterSettings.Self.NodeInfo.InternalTcp,
- clusterSettings.Self.NodeInfo.InternalSecureTcp,
- clusterSettings.Self.NodeInfo.ExternalTcp,
- clusterSettings.Self.NodeInfo.ExternalSecureTcp,
- clusterSettings.Self.NodeInfo.HttpEndPoint, null, 0, 0,
+ clusterSettings.Self.NodeInfo.HttpEndPoint, null, 0,
LastCommitPosition, WriterCheckpoint, ChaserCheckpoint,
-1,
-1,
@@ -98,11 +90,7 @@ private ClusterInfo BuildClusterInfo(ClusterSettings clusterSettings)
InitialDate,
VNodeState.Unknown,
true,
- x.NodeInfo.InternalTcp,
- x.NodeInfo.InternalSecureTcp,
- x.NodeInfo.ExternalTcp,
- x.NodeInfo.ExternalSecureTcp,
- x.NodeInfo.HttpEndPoint, null, 0, 0,
+ x.NodeInfo.HttpEndPoint, null, 0,
LastCommitPosition, WriterCheckpoint, ChaserCheckpoint,
-1,
-1,
@@ -196,9 +184,7 @@ public IEnumerable ListMembers(Func predicate = nu
? MemberInfo.ForManager(x.InstanceId, x.TimeStamp, x.IsAlive,
x.HttpEndPoint)
: MemberInfo.ForVNode(x.InstanceId, x.TimeStamp, x.State, x.IsAlive,
- x.InternalTcpEndPoint, x.InternalSecureTcpEndPoint,
- x.ExternalTcpEndPoint, x.ExternalSecureTcpEndPoint,
- x.HttpEndPoint, null, 0, 0,
+ x.HttpEndPoint, null, 0,
x.LastCommitPosition, x.WriterCheckpoint, x.ChaserCheckpoint,
x.EpochPosition, x.EpochNumber, x.EpochId, x.NodePriority, x.IsReadOnlyReplica));
}
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs b/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs
index 16a861db2f..e54f644bce 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/ElectionsServiceTests.cs
@@ -29,19 +29,13 @@ public abstract class ElectionsFixture
protected static Func NodeFactory = (id) => new VNodeInfo(
Guid.Parse($"00000000-0000-0000-0000-00000000000{id}"), id,
- new IPEndPoint(IPAddress.Loopback, id),
- new IPEndPoint(IPAddress.Loopback, id),
- new IPEndPoint(IPAddress.Loopback, id),
- new IPEndPoint(IPAddress.Loopback, id),
new IPEndPoint(IPAddress.Loopback, id), false,
new IPEndPoint(IPAddress.Loopback, 10_000 + id));
protected static readonly Func MemberInfoFromVNode =
(nodeInfo, timestamp, state, isAlive, epochNumber, epochId, priority) => MemberInfo.ForVNode(
nodeInfo.InstanceId, timestamp, state, isAlive,
- nodeInfo.InternalTcp,
- nodeInfo.InternalSecureTcp, nodeInfo.ExternalTcp, nodeInfo.ExternalSecureTcp,
- nodeInfo.HttpEndPoint, null, 0, 0,
+ nodeInfo.HttpEndPoint, null, 0,
0, 0, 0, 0, epochNumber, epochId, priority,
nodeInfo.IsReadOnlyReplica, clusterEndPoint: nodeInfo.ClusterEndPoint);
@@ -925,9 +919,7 @@ public void should_send_an_acceptance_to_other_members()
new ElectionMessage.ElectionsDone(0,0,
MemberInfo.ForVNode(
_nodeThree.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeThree.InternalTcp,
- _nodeThree.InternalSecureTcp, _nodeThree.ExternalTcp, _nodeThree.ExternalSecureTcp,
- _nodeThree.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, 0, _epochId, 0,
+ _nodeThree.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, _epochId, 0,
_nodeThree.IsReadOnlyReplica)),
new GrpcMessage.SendOverGrpc(_nodeThree.ClusterEndPoint,
new ElectionMessage.Accept(_node.InstanceId, _node.HttpEndPoint,
@@ -1078,9 +1070,7 @@ public void should_complete_elections()
new ElectionMessage.ElectionsDone(0,0,
MemberInfo.ForVNode(
_nodeTwo.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeTwo.InternalTcp,
- _nodeTwo.InternalSecureTcp, _nodeTwo.ExternalTcp, _nodeTwo.ExternalSecureTcp,
- _nodeTwo.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, 0, _epochId, 0,
+ _nodeTwo.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, _epochId, 0,
_nodeTwo.IsReadOnlyReplica)),
};
_publisher.Messages.Should().BeEquivalentTo(expected);
@@ -1349,9 +1339,7 @@ public void should_attempt_not_to_elect_previously_elected_leader()
new ElectionMessage.ElectionsDone(3,1,
MemberInfo.ForVNode(
_nodeTwo.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeTwo.InternalTcp,
- _nodeTwo.InternalSecureTcp, _nodeTwo.ExternalTcp, _nodeTwo.ExternalSecureTcp,
- _nodeTwo.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, 0, _epochId, 0,
+ _nodeTwo.HttpEndPoint, null, 0, 0, 0, 0, 0, 0, _epochId, 0,
_nodeTwo.IsReadOnlyReplica)),
};
_publisher.Messages.Should().BeEquivalentTo(expected);
@@ -1392,11 +1380,7 @@ public void should_throw_argument_exception()
var endpoint = new IPEndPoint(IPAddress.Loopback, 1234);
var nodeInfo = MemberInfo.Initial(Guid.NewGuid(),
DateTime.UtcNow, VNodeState.ReadOnlyLeaderless, true,
- endpoint,
- endpoint,
- endpoint,
- endpoint,
- endpoint, null, 0, 0,
+ endpoint, null, 0,
0,
true);
@@ -1440,9 +1424,7 @@ public void previous_leader_should_be_elected()
new ElectionMessage.ElectionsDone(0,0,
MemberInfo.ForVNode(
_nodeThree.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeThree.InternalTcp,
- _nodeThree.InternalSecureTcp, _nodeThree.ExternalTcp, _nodeThree.ExternalSecureTcp,
- _nodeThree.HttpEndPoint, null, 0, 0,
+ _nodeThree.HttpEndPoint, null, 0,
0, 0, 0, 0, 0, _epochId, 0,
_nodeThree.IsReadOnlyReplica)),
};
@@ -1540,9 +1522,7 @@ public void previous_leader_should_not_be_elected()
new ElectionMessage.ElectionsDone(0,0,
MemberInfo.ForVNode(
_nodeTwo.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeTwo.InternalTcp,
- _nodeTwo.InternalSecureTcp, _nodeTwo.ExternalTcp, _nodeTwo.ExternalSecureTcp,
- _nodeTwo.HttpEndPoint, null, 0, 0,
+ _nodeTwo.HttpEndPoint, null, 0,
0, 0, 0, 0, 0, _epochId, 0,
_nodeTwo.IsReadOnlyReplica)),
};
@@ -1576,9 +1556,7 @@ public void previous_leader_should_not_be_elected()
new ElectionMessage.ElectionsDone(0,0,
MemberInfo.ForVNode(
_nodeTwo.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeTwo.InternalTcp,
- _nodeTwo.InternalSecureTcp, _nodeTwo.ExternalTcp, _nodeTwo.ExternalSecureTcp,
- _nodeTwo.HttpEndPoint, null, 0, 0,
+ _nodeTwo.HttpEndPoint, null, 0,
0, 0, 0, 0, 0, _epochId, 0,
_nodeTwo.IsReadOnlyReplica)),
};
@@ -1633,9 +1611,7 @@ public void previous_leader_should_not_be_elected()
new ElectionMessage.ElectionsDone(0,0,
MemberInfo.ForVNode(
_nodeTwo.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeTwo.InternalTcp,
- _nodeTwo.InternalSecureTcp, _nodeTwo.ExternalTcp, _nodeTwo.ExternalSecureTcp,
- _nodeTwo.HttpEndPoint, null, 0, 0,
+ _nodeTwo.HttpEndPoint, null, 0,
0, 0, 0, 0, 0, _epochId, 0,
_nodeTwo.IsReadOnlyReplica)),
};
@@ -1669,9 +1645,7 @@ public void previous_leader_should_not_be_elected()
new ElectionMessage.ElectionsDone(0,0,
MemberInfo.ForVNode(
_nodeTwo.InstanceId, _timeProvider.UtcNow, VNodeState.Unknown, true,
- _nodeTwo.InternalTcp,
- _nodeTwo.InternalSecureTcp, _nodeTwo.ExternalTcp, _nodeTwo.ExternalSecureTcp,
- _nodeTwo.HttpEndPoint, null, 0, 0,
+ _nodeTwo.HttpEndPoint, null, 0,
0, 0, 0, 0, 0, _epochId, 0,
_nodeTwo.IsReadOnlyReplica)),
};
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/LeaderNode/ElectionsServiceUnitTests.cs b/src/EventStore.Core.Tests/Services/ElectionsService/LeaderNode/ElectionsServiceUnitTests.cs
index 11c8462e3d..6e3a013cf5 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/LeaderNode/ElectionsServiceUnitTests.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/LeaderNode/ElectionsServiceUnitTests.cs
@@ -40,7 +40,7 @@ public void Setup()
seeds.Add(endPoint);
var instanceId = Guid.Parse($"101EFD13-F9CD-49BE-9C6D-E6AF9AF5540{i}");
var memberInfo = MemberInfo.ForVNode(instanceId, DateTime.UtcNow, VNodeState.Unknown, true,
- endPoint, null, endPoint, null, endPoint, null, 0, 0, -1, 0, 0, -1, -1, Guid.Empty, 0, false);
+ endPoint, null, 0, -1, 0, 0, -1, -1, Guid.Empty, 0, false);
members.Add(memberInfo);
_fakeTimeProvider = new FakeTimeProvider();
_scheduler = new FakeScheduler(new FakeTimer(), _fakeTimeProvider);
@@ -379,7 +379,7 @@ Func epochNumber
{
var id = IdForNode(i);
var ep = EndpointForNode(i);
- return MemberInfo.ForVNode(id, DateTime.Now, VNodeState.Follower, true, ep, ep, ep, ep, ep, null, 0, 0,
+ return MemberInfo.ForVNode(id, DateTime.Now, VNodeState.Follower, true, ep, null, 0,
-1, writerCheckpoint(i), chaserCheckpoint(i), 1, epochNumber(i), epochId, nodePriority(i), false);
}
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/RandomizedElectionsTestCase.cs b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/RandomizedElectionsTestCase.cs
index 68b6ce5fd9..168f65a466 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/RandomizedElectionsTestCase.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/RandomizedElectionsTestCase.cs
@@ -74,7 +74,7 @@ public virtual void Init()
var outputBus = new SynchronousScheduler($"ELECTIONS-OUTPUT-BUS-{i}");
var endPoint = new IPEndPoint(BaseEndPoint.Address, BaseEndPoint.Port + i);
var memberInfo = MemberInfo.Initial(Guid.NewGuid(), DateTime.UtcNow, VNodeState.Unknown, true,
- endPoint, endPoint, endPoint, endPoint, endPoint, null, 0, 0, 0, false);
+ endPoint, null, 0, 0, false);
_instances.Add(new ElectionsInstance(memberInfo.InstanceId, endPoint, inputBus, outputBus));
sendOverHttpHandler.RegisterEndPoint(endPoint, inputBus);
@@ -125,8 +125,7 @@ protected virtual GossipMessage.GossipUpdated GetInitialGossipFor(ElectionsInsta
{
var members = allInstances.Select(
x => MemberInfo.ForVNode(x.InstanceId, DateTime.UtcNow, VNodeState.Unknown, true,
- x.EndPoint, null, x.EndPoint, null,
- x.EndPoint, null, 0, 0, -1, 0, 0, -1, -1, Guid.Empty, 0, false));
+ x.EndPoint, null, 0, -1, 0, 0, -1, -1, Guid.Empty, 0, false));
var gossip = new GossipMessage.GossipUpdated(new ClusterInfo(members.ToArray()));
return gossip;
}
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/UpdateGossipProcessor.cs b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/UpdateGossipProcessor.cs
index 7ed5cd5fc9..b9381ca8dc 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/UpdateGossipProcessor.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/UpdateGossipProcessor.cs
@@ -63,7 +63,7 @@ public void Process(int iteration, RandTestQueueItem item)
previousMembers[leaderIndex] =
MemberInfo.ForVNode(previousLeaderInfo.InstanceId, DateTime.UtcNow, VNodeState.Leader,
- previousLeaderInfo.IsAlive, leaderEndPoint, null, leaderEndPoint, null, leaderEndPoint, null, 0, 0,
+ previousLeaderInfo.IsAlive, leaderEndPoint, null, 0,
-1, 0, 0, -1, -1, Guid.Empty, 0, false);
}
}
@@ -82,7 +82,7 @@ public void Process(int iteration, RandTestQueueItem item)
foreach (var memberInfo in updatedGossip)
{
- _sendOverGrpcProcessor.RegisterEndpointToSkip(memberInfo.ExternalTcpEndPoint, !memberInfo.IsAlive);
+ _sendOverGrpcProcessor.RegisterEndpointToSkip(memberInfo.HttpEndPoint, !memberInfo.IsAlive);
}
var updateGossipMessage = new GossipMessage.GossipUpdated(new ClusterInfo(updatedGossip));
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_full_imediately.cs b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_full_imediately.cs
index 339ca873bd..6722010f0f 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_full_imediately.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_full_imediately.cs
@@ -37,7 +37,7 @@ private MemberInfo[] CreateInitialGossip(ElectionsInstance instance, ElectionsIn
{
return new[] {
MemberInfo.ForVNode(instance.InstanceId, DateTime.UtcNow, VNodeState.Unknown, true,
- instance.EndPoint, null, instance.EndPoint, null, instance.EndPoint, null, 0, 0,
+ instance.EndPoint, null, 0,
-1, 0, 0, -1, -1, Guid.Empty, 0, false)
};
}
@@ -53,7 +53,7 @@ private MemberInfo[] CreateUpdatedGossip(int iteration,
Console.WriteLine("Update item: {0} : {1}", iteration, item.EndPoint.GetPort());
return instances.Select((x, i) =>
MemberInfo.ForVNode(x.InstanceId, DateTime.UtcNow, VNodeState.Unknown, true,
- x.EndPoint, null, x.EndPoint, null, x.EndPoint, null, 0, 0,
+ x.EndPoint, null, 0,
-1, 0, 0, -1, -1, Guid.Empty, 0, false)).ToArray();
}
diff --git a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_to_full_later.cs b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_to_full_later.cs
index 05f165cfab..cdf621e6ab 100644
--- a/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_to_full_later.cs
+++ b/src/EventStore.Core.Tests/Services/ElectionsService/Randomized/elections_service_5_nodes_with_1_known_when_started_and_set_to_full_later.cs
@@ -37,7 +37,7 @@ private MemberInfo[] CreateInitialGossip(ElectionsInstance instance, ElectionsIn
{
return new[] {
MemberInfo.ForVNode(instance.InstanceId, DateTime.UtcNow, VNodeState.Unknown, true,
- instance.EndPoint, null, instance.EndPoint, null, instance.EndPoint, null, 0, 0,
+ instance.EndPoint, null, 0,
-1, 0, 0, -1, -1, Guid.Empty, 0, false)
};
}
@@ -57,7 +57,7 @@ private MemberInfo[] CreateUpdatedGossip(int iteration,
{
return instances.Select((x, i) =>
MemberInfo.ForVNode(x.InstanceId, DateTime.UtcNow, VNodeState.Unknown, true,
- x.EndPoint, null, x.EndPoint, null, x.EndPoint, null, 0, 0,
+ x.EndPoint, null, 0,
-1, 0, 0, -1, -1, Guid.Empty, 0, false))
.ToArray();
}
diff --git a/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs b/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs
index 7f230e2f06..f8d6d3419f 100644
--- a/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs
+++ b/src/EventStore.Core.Tests/Services/GossipService/NodeGossipServiceTests.cs
@@ -44,34 +44,22 @@ public NodeGossipServiceTestFixture()
_currentNode = new VNodeInfo(
Guid.Parse("00000000-0000-0000-0000-000000000001"), 1,
new IPEndPoint(IPAddress.Loopback, 1111),
- new IPEndPoint(IPAddress.Loopback, 1111),
- new IPEndPoint(IPAddress.Loopback, 1111),
- new IPEndPoint(IPAddress.Loopback, 1111),
- new IPEndPoint(IPAddress.Loopback, 1111), false,
+ false,
new IPEndPoint(IPAddress.Loopback, 11112));
_nodeTwo = new VNodeInfo(
Guid.Parse("00000000-0000-0000-0000-000000000002"), 2,
new IPEndPoint(IPAddress.Loopback, 2222),
- new IPEndPoint(IPAddress.Loopback, 2222),
- new IPEndPoint(IPAddress.Loopback, 2222),
- new IPEndPoint(IPAddress.Loopback, 2222),
- new IPEndPoint(IPAddress.Loopback, 2222), false,
+ false,
new IPEndPoint(IPAddress.Loopback, 22212));
_nodeThree = new VNodeInfo(
Guid.Parse("00000000-0000-0000-0000-000000000003"), 3,
new IPEndPoint(IPAddress.Loopback, 3333),
- new IPEndPoint(IPAddress.Loopback, 3333),
- new IPEndPoint(IPAddress.Loopback, 3333),
- new IPEndPoint(IPAddress.Loopback, 3333),
- new IPEndPoint(IPAddress.Loopback, 3333), false,
+ false,
new IPEndPoint(IPAddress.Loopback, 33312));
_nodeFour = new VNodeInfo(
Guid.Parse("00000000-0000-0000-0000-000000000004"), 4,
new IPEndPoint(IPAddress.Loopback, 4444),
- new IPEndPoint(IPAddress.Loopback, 4444),
- new IPEndPoint(IPAddress.Loopback, 4444),
- new IPEndPoint(IPAddress.Loopback, 4444),
- new IPEndPoint(IPAddress.Loopback, 4444), false,
+ false,
new IPEndPoint(IPAddress.Loopback, 44412));
_getNodeToGossipTo = infos => infos.First(x => Equals(x.ClusterEndPoint, _nodeTwo.ClusterEndPoint));
@@ -124,9 +112,8 @@ protected static MemberInfo MemberInfoForVNode(VNodeInfo nodeInfo, DateTime utcN
VNodeState nodeState = VNodeState.Initializing, string esVersion = VersionInfo.DefaultVersion, bool isAlive = true)
{
return MemberInfo.ForVNode(nodeInfo.InstanceId, utcNow, nodeState, isAlive,
- nodeInfo.InternalTcp, nodeInfo.InternalSecureTcp, nodeInfo.ExternalTcp,
- nodeInfo.ExternalSecureTcp, nodeInfo.HttpEndPoint, null, 0, 0,
- 0, writerCheckpoint ?? 0, 0, -1, epochNumber ?? -1, Guid.Empty, nodePriority ?? 0, false, esVersion,
+ nodeInfo.HttpEndPoint, null, 0, 0,
+ writerCheckpoint ?? 0, 0, -1, epochNumber ?? -1, Guid.Empty, nodePriority ?? 0, false, esVersion,
nodeInfo.ClusterEndPoint);
}
@@ -257,10 +244,6 @@ public when_got_gossip_seed_sources_with_distinct_cluster_endpoint()
{
_currentNode = new VNodeInfo(
Guid.Parse("00000000-0000-0000-0000-000000000001"), 1,
- new IPEndPoint(IPAddress.Loopback, 1111),
- new IPEndPoint(IPAddress.Loopback, 1111),
- new IPEndPoint(IPAddress.Loopback, 1111),
- new IPEndPoint(IPAddress.Loopback, 1111),
new IPEndPoint(IPAddress.Loopback, 1111), false,
new IPEndPoint(IPAddress.Loopback, 1112));
}
@@ -797,7 +780,7 @@ protected override Message[] Given() =>
_nodeTwo.HttpEndPoint);
[Test]
- public void should_ignore_message_and_wait_for_tcp_to_decide()
+ public void should_ignore_message_and_wait_for_connection_state_to_change()
{
ExpectNoMessages();
}
@@ -1085,7 +1068,7 @@ private static MemberInfo TestNodeFor(int identifier, bool isAlive, DateTime tim
{
var ipEndpoint = new IPEndPoint(IPAddress.Loopback, identifier);
return MemberInfo.ForVNode(Guid.NewGuid(), timeStamp, VNodeState.Initializing, isAlive,
- ipEndpoint, ipEndpoint, ipEndpoint, ipEndpoint, ipEndpoint, null, 0, 0,
+ ipEndpoint, null, 0,
0, 0, 0, -1, -1, Guid.Empty, 0, false);
}
@@ -1162,7 +1145,7 @@ private static MemberInfo TestNodeFor(int identifier, bool isAlive, DateTime tim
{
var ipEndpoint = new IPEndPoint(IPAddress.Loopback, identifier);
return MemberInfo.ForVNode(Guid.NewGuid(), timeStamp, nodeState, isAlive,
- ipEndpoint, ipEndpoint, ipEndpoint, ipEndpoint, ipEndpoint, null, 0, 0,
+ ipEndpoint, null, 0,
0, 0, 0, -1, -1, Guid.Empty, 0, false);
}
@@ -1177,8 +1160,7 @@ public void should_never_replace_self_with_a_newer_cluster_endpoint_seed()
var clusterEndPoint = new IPEndPoint(IPAddress.Loopback, 1112);
var me = MemberInfo.ForVNode(
Guid.NewGuid(), now, VNodeState.Initializing, true,
- clusterEndPoint, null, httpEndPoint, null, httpEndPoint,
- null, 0, 0, -1, -1, -1, -1, -1, Guid.Empty, 0, false,
+ httpEndPoint, null, 0, 0, -1, -1, -1, -1, Guid.Empty, 0, false,
clusterEndPoint: clusterEndPoint);
var selfSeed = MemberInfo.ForManager(
Guid.Empty, now.AddSeconds(1), true, clusterEndPoint,
diff --git a/src/EventStore.Core.Tests/Services/Replication/LogReplication/LogReplicationFixture.cs b/src/EventStore.Core.Tests/Services/Replication/LogReplication/LogReplicationFixture.cs
index a7514a3a51..1dfac124f9 100644
--- a/src/EventStore.Core.Tests/Services/Replication/LogReplication/LogReplicationFixture.cs
+++ b/src/EventStore.Core.Tests/Services/Replication/LogReplication/LogReplicationFixture.cs
@@ -157,14 +157,9 @@ private async ValueTask> CreateLeader(TFChunkDb db, Cancel
timeStamp: DateTime.Now,
state: VNodeState.Leader,
isAlive: true,
- internalTcpEndPoint: FakeEndPoint,
- internalSecureTcpEndPoint: null,
- externalTcpEndPoint: null,
- externalSecureTcpEndPoint: null,
httpEndPoint: FakeEndPoint,
advertiseHostToClientAs: null,
advertiseHttpPortToClientAs: 0,
- advertiseTcpPortToClientAs: 0,
lastCommitPosition: 0,
writerCheckpoint: 0,
chaserCheckpoint: 0,
diff --git a/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs b/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs
index e59e68f072..2296f920c7 100644
--- a/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs
+++ b/src/EventStore.Core.Tests/Services/Replication/ReadOnlyReplica/connecting_to_read_only_replica.cs
@@ -1,11 +1,21 @@
+using System;
+using System.Linq;
using System.Net;
+using System.Net.Http;
+using System.Text;
using System.Threading.Tasks;
-using EventStore.ClientAPI;
-using EventStore.ClientAPI.Exceptions;
-using EventStore.Core.Tests.ClientAPI.Helpers;
+using EventStore.Client.Streams;
+using EventStore.Core.Data;
+using EventStore.Core.Services.Transport.Grpc;
using EventStore.Core.Tests.Helpers;
using EventStore.Core.Tests.Integration;
+using Google.Protobuf;
+using Grpc.Core;
+using Grpc.Net.Client;
using NUnit.Framework;
+using Empty = EventStore.Client.Empty;
+using GrpcExceptions = EventStore.Core.Services.Transport.Grpc.Constants.Exceptions;
+using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata;
namespace EventStore.Core.Tests.Replication.ReadOnlyReplica;
@@ -13,13 +23,20 @@ namespace EventStore.Core.Tests.Replication.ReadOnlyReplica;
[TestFixture(typeof(LogFormat.V2), typeof(string))]
public class connecting_to_read_only_replica : specification_with_cluster
{
+ protected override async Task Given()
+ {
+ await _nodes[2].AdminUserCreated.WithTimeout(TimeSpan.FromSeconds(30));
+ AssertEx.IsOrBecomesTrue(() => _nodes[2].NodeState == VNodeState.ReadOnlyReplica,
+ timeout: TimeSpan.FromSeconds(30),
+ onFail: MiniNodeLogging.WriteLogs);
+ }
+
protected override MiniClusterNode CreateNode(int index, Endpoints endpoints, EndPoint[] gossipSeeds,
bool wait = true)
{
var isReadOnly = index == 2;
var node = new MiniClusterNode(
- PathName, index, endpoints.ClusterEndPoint,
- endpoints.ExternalTcp, endpoints.HttpEndPoint, gossipSeeds,
+ PathName, index, endpoints.HttpEndPoint, endpoints.ClusterEndPoint, gossipSeeds,
readOnlyReplica: isReadOnly);
if (wait && !isReadOnly)
{
@@ -29,35 +46,185 @@ protected override MiniClusterNode CreateNode(int index,
return node;
}
- protected override IEventStoreConnection CreateConnection()
+ private static CallOptions GetCallOptions()
+ {
+ var credentials = CallCredentials.FromInterceptor((_, metadata) =>
+ {
+ metadata.Add("authorization",
+ $"Basic {Convert.ToBase64String(Encoding.ASCII.GetBytes("admin:changeit"))}");
+ return Task.CompletedTask;
+ });
+ return new CallOptions(credentials: credentials, deadline: DateTime.UtcNow.AddSeconds(30));
+ }
+
+ private static Streams.StreamsClient CreateClient(
+ MiniClusterNode node,
+ out GrpcChannel channel,
+ out HttpClient httpClient)
{
- var settings = ConnectionSettings.Create()
- .DisableServerCertificateValidation()
- .PerformOnAnyNode();
- return EventStoreConnection.Create(settings, _nodes[2].ExternalTcpEndPoint);
+ httpClient = new HttpClient(new SocketsHttpHandler
+ {
+ SslOptions = { RemoteCertificateValidationCallback = delegate { return true; } }
+ });
+ channel = GrpcChannel.ForAddress(new Uri($"https://{node.HttpEndPoint}"),
+ new GrpcChannelOptions { HttpClient = httpClient });
+ return new Streams.StreamsClient(channel);
}
[Test]
- public async Task append_to_stream_should_fail_with_not_supported_exception()
+ public async Task append_to_stream_is_rejected()
{
- const string stream = "append_to_stream_should_fail_with_not_supported_exception";
- await AssertEx.ThrowsAsync(
- () => _conn.AppendToStreamAsync(stream, ExpectedVersion.Any, TestEvent.NewTestEvent()));
+ var client = CreateClient(_nodes[2], out var channel, out var httpClient);
+ using (channel)
+ using (httpClient)
+ using (var call = client.Append(GetCallOptions()))
+ {
+ await call.RequestStream.WriteAsync(new AppendReq
+ {
+ Options = new()
+ {
+ Any = new Empty(),
+ StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(nameof(append_to_stream_is_rejected)) }
+ }
+ });
+ await call.RequestStream.WriteAsync(new AppendReq
+ {
+ ProposedMessage = new()
+ {
+ Id = Uuid.NewUuid().ToDto(),
+ Data = ByteString.Empty,
+ CustomMetadata = ByteString.Empty,
+ Metadata =
+ {
+ [GrpcMetadata.Type] = "test",
+ [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson
+ }
+ }
+ });
+ await call.RequestStream.CompleteAsync();
+
+ var exception = Assert.ThrowsAsync(async () => await call.ResponseAsync);
+ Assert.That(exception.StatusCode, Is.EqualTo(StatusCode.NotFound));
+ Assert.That(exception.Trailers.Select(x => (x.Key, x.Value)),
+ Does.Contain((GrpcExceptions.ExceptionKey, GrpcExceptions.NotLeader)));
+ }
}
[Test]
- public async Task delete_stream_should_fail_with_not_supported_exception()
+ public async Task batch_append_is_rejected()
{
- const string stream = "delete_stream_should_fail_with_not_supported_exception";
- await AssertEx.ThrowsAsync(() =>
- _conn.DeleteStreamAsync(stream, ExpectedVersion.Any));
+ var client = CreateClient(_nodes[2], out var channel, out var httpClient);
+ using (channel)
+ using (httpClient)
+ using (var call = client.BatchAppend(GetCallOptions()))
+ {
+ await call.RequestStream.WriteAsync(new BatchAppendReq
+ {
+ CorrelationId = Uuid.NewUuid().ToDto(),
+ Options = new()
+ {
+ Any = new Google.Protobuf.WellKnownTypes.Empty(),
+ StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(nameof(batch_append_is_rejected)) }
+ },
+ IsFinal = true,
+ ProposedMessages =
+ {
+ new BatchAppendReq.Types.ProposedMessage
+ {
+ Id = Uuid.NewUuid().ToDto(),
+ Metadata =
+ {
+ [GrpcMetadata.Type] = "test",
+ [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson
+ }
+ },
+ new BatchAppendReq.Types.ProposedMessage
+ {
+ Id = Uuid.NewUuid().ToDto(),
+ Metadata =
+ {
+ [GrpcMetadata.Type] = "test",
+ [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson
+ }
+ }
+ }
+ });
+ await call.RequestStream.CompleteAsync();
+
+ var exception = Assert.ThrowsAsync(async () => await call.ResponseStream.MoveNext());
+ Assert.That(exception.StatusCode, Is.EqualTo(StatusCode.NotFound));
+ Assert.That(exception.Trailers.Select(x => (x.Key, x.Value)),
+ Does.Contain((GrpcExceptions.ExceptionKey, GrpcExceptions.NotLeader)));
+ }
}
[Test]
- public async Task start_transaction_should_fail_with_not_supported_exception()
+ public async Task delete_stream_is_rejected()
+ {
+ const string stream = nameof(delete_stream_is_rejected);
+ var leader = GetLeader();
+ await AppendToStream(leader, stream);
+ var leaderWriterPosition = leader.Db.Config.WriterCheckpoint.Read();
+ AssertEx.IsOrBecomesTrue(
+ () => _nodes[2].Db.Config.WriterCheckpoint.Read() >= leaderWriterPosition,
+ timeout: TimeSpan.FromSeconds(30),
+ onFail: MiniNodeLogging.WriteLogs,
+ msg: "The stream was not replicated to the read-only replica.");
+
+ var client = CreateClient(_nodes[2], out var channel, out var httpClient);
+ using (channel)
+ using (httpClient)
+ using (var call = client.DeleteAsync(new DeleteReq
+ {
+ Options = new()
+ {
+ Any = new Empty(),
+ StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) }
+ }
+ }, GetCallOptions()))
+ {
+ var exception = Assert.ThrowsAsync(async () => await call.ResponseAsync);
+ Assert.That(exception.StatusCode, Is.EqualTo(StatusCode.NotFound));
+ Assert.That(exception.Trailers.Select(x => (x.Key, x.Value)),
+ Does.Contain((GrpcExceptions.ExceptionKey, GrpcExceptions.NotLeader)));
+ }
+ }
+
+ private static async Task AppendToStream(
+ MiniClusterNode node,
+ string stream)
{
- const string stream = "start_transaction_should_fail_with_not_supported_exception";
- await AssertEx.ThrowsAsync(() =>
- _conn.StartTransactionAsync(stream, ExpectedVersion.Any));
+ var client = CreateClient(node, out var channel, out var httpClient);
+ using (channel)
+ using (httpClient)
+ using (var call = client.Append(GetCallOptions()))
+ {
+ await call.RequestStream.WriteAsync(new AppendReq
+ {
+ Options = new()
+ {
+ NoStream = new Empty(),
+ StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(stream) }
+ }
+ });
+ await call.RequestStream.WriteAsync(new AppendReq
+ {
+ ProposedMessage = new()
+ {
+ Id = Uuid.NewUuid().ToDto(),
+ Data = ByteString.Empty,
+ CustomMetadata = ByteString.Empty,
+ Metadata =
+ {
+ [GrpcMetadata.Type] = "test",
+ [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson
+ }
+ }
+ });
+ await call.RequestStream.CompleteAsync();
+
+ var response = await call.ResponseAsync;
+ Assert.That(response.ResultCase, Is.EqualTo(AppendResp.ResultOneofCase.Success));
+ }
}
}
diff --git a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs
index 6534f123e2..f756f98deb 100644
--- a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs
+++ b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingServiceTests.cs
@@ -812,10 +812,6 @@ private static MemberInfo CreateLeader(
DateTime.UtcNow,
VNodeState.Leader,
true,
- new DnsEndPoint("leader-replication.internal", 1112),
- new DnsEndPoint("leader-replication.internal", 1113),
- null,
- null,
new DnsEndPoint("leader.internal", httpPort),
null,
0,
@@ -824,7 +820,6 @@ private static MemberInfo CreateLeader(
0,
0,
0,
- 0,
Guid.NewGuid(),
0,
false,
diff --git a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs
index 9851290496..5c4c720501 100644
--- a/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs
+++ b/src/EventStore.Core.Tests/Services/RequestForwarding/GrpcRequestForwardingTransportSecurityTests.cs
@@ -53,7 +53,7 @@ public async Task cleartext_credential_rejection_completes_the_proxy_correlation
supervisor.Handle(new ClientMessage.ForwardMessage(request));
- var response = publisher.Messages.OfType().Single();
+ var response = publisher.Messages.OfType().Single();
Assert.That(response.CorrelationId, Is.EqualTo(request.InternalCorrId));
}
@@ -119,10 +119,6 @@ private static MemberInfo CreateLeader() => MemberInfo.ForVNode(
DateTime.UtcNow,
VNodeState.Leader,
true,
- new DnsEndPoint("leader-replication.internal", 1112),
- new DnsEndPoint("leader-replication.internal", 1113),
- null,
- null,
new DnsEndPoint("leader.internal", 2113),
null,
0,
@@ -131,7 +127,6 @@ private static MemberInfo CreateLeader() => MemberInfo.ForVNode(
0,
0,
0,
- 0,
Guid.NewGuid(),
0,
false,
diff --git a/src/EventStore.Core.Tests/Services/RequestForwarding/RequestForwardingServiceTests.cs b/src/EventStore.Core.Tests/Services/RequestForwarding/RequestForwardingServiceTests.cs
index d0ed4eb5e5..8eb4dccb3c 100644
--- a/src/EventStore.Core.Tests/Services/RequestForwarding/RequestForwardingServiceTests.cs
+++ b/src/EventStore.Core.Tests/Services/RequestForwarding/RequestForwardingServiceTests.cs
@@ -55,13 +55,13 @@ public void not_authenticated_survives_the_client_correlation_rewrite()
clientCorrelationId,
new CallbackEnvelope(message => response = message),
TimeSpan.FromMinutes(1),
- new TcpMessage.NotAuthenticated(clientCorrelationId, "timeout"));
+ new ClientMessage.NotAuthenticated(clientCorrelationId, "timeout"));
var service = new RequestForwardingService(
new NoopPublisher(), forwardingProxy, TimeSpan.FromSeconds(1));
- service.Handle(new TcpMessage.NotAuthenticated(internalCorrelationId, "not authenticated"));
+ service.Handle(new ClientMessage.NotAuthenticated(internalCorrelationId, "not authenticated"));
- var completion = (TcpMessage.NotAuthenticated)response;
+ var completion = (ClientMessage.NotAuthenticated)response;
Assert.Multiple(() =>
{
Assert.That(completion.CorrelationId, Is.EqualTo(clientCorrelationId));
diff --git a/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader.cs b/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader.cs
index 8ef3d8d03b..6f3bf4ff30 100644
--- a/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader.cs
+++ b/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader.cs
@@ -42,10 +42,6 @@ private static MemberInfo FakeMemberInfo()
return EventStore.Core.Cluster.MemberInfo.Initial(Guid.Empty, DateTime.UtcNow,
VNodeState.Unknown, true,
new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- null, 0, 0, 0, false);
+ null, 0, 0, false);
}
}
diff --git a/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader_and_replica_moves_forward.cs b/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader_and_replica_moves_forward.cs
index 994a62f28d..4b7ba26312 100644
--- a/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader_and_replica_moves_forward.cs
+++ b/src/EventStore.Core.Tests/Services/RequestManagement/Service/when_writing_and_deposed_as_leader_and_replica_moves_forward.cs
@@ -37,10 +37,6 @@ private static MemberInfo FakeMemberInfo()
return EventStore.Core.Cluster.MemberInfo.Initial(Guid.Empty, DateTime.UtcNow,
VNodeState.Unknown, true,
new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- new IPEndPoint(IPAddress.Parse(ipAddress), port),
- null, 0, 0, 0, false);
+ null, 0, 0, false);
}
}
diff --git a/src/EventStore.Core.Tests/Services/Transport/Enumerators/SubscriptionDisposalOrderingTests.cs b/src/EventStore.Core.Tests/Services/Transport/Enumerators/SubscriptionDisposalOrderingTests.cs
new file mode 100644
index 0000000000..45df456276
--- /dev/null
+++ b/src/EventStore.Core.Tests/Services/Transport/Enumerators/SubscriptionDisposalOrderingTests.cs
@@ -0,0 +1,94 @@
+using System;
+using System.Collections.Concurrent;
+using System.Collections.Generic;
+using System.Threading;
+using System.Threading.Tasks;
+using EventStore.Core.Bus;
+using EventStore.Core.Messages;
+using EventStore.Core.Messaging;
+using EventStore.Core.Services.Storage.ReaderIndex;
+using EventStore.Core.Services.Transport.Common;
+using EventStore.Core.Services.Transport.Enumerators;
+using EventStore.Core.Services.UserManagement;
+using NUnit.Framework;
+
+namespace EventStore.Core.Tests.Services.Transport.Enumerators;
+
+[TestFixture]
+public class SubscriptionDisposalOrderingTests
+{
+ [TestCase("stream")]
+ [TestCase("all")]
+ [TestCase("filtered-all")]
+ public async Task disposing_while_live_subscription_is_being_registered_does_not_leave_it_active(string kind)
+ {
+ var publisher = new DelayedSubscriptionPublisher();
+ var enumerator = CreateEnumerator(kind, publisher);
+ try
+ {
+ await publisher.SubscribeEntered.Task.WaitAsync(TimeSpan.FromSeconds(10));
+ var disposal = enumerator.DisposeAsync().AsTask();
+ publisher.ReleaseSubscribe();
+ await disposal.WaitAsync(TimeSpan.FromSeconds(10));
+ await publisher.SubscribeCompleted.Task.WaitAsync(TimeSpan.FromSeconds(10));
+
+ Assert.That(publisher.ActiveSubscriptions, Is.Empty);
+ }
+ finally
+ {
+ publisher.ReleaseSubscribe();
+ await enumerator.DisposeAsync();
+ }
+ }
+
+ private static IAsyncEnumerator CreateEnumerator(string kind, IPublisher publisher) => kind switch
+ {
+ "stream" => new Enumerator.StreamSubscription(
+ publisher, new DefaultExpiryStrategy(), "subscription-disposal-ordering", StreamRevision.End,
+ false, SystemAccounts.System, false, CancellationToken.None),
+ "all" => new Enumerator.AllSubscription(
+ publisher, new DefaultExpiryStrategy(), Position.End,
+ false, SystemAccounts.System, false, CancellationToken.None),
+ "filtered-all" => new Enumerator.AllSubscriptionFiltered(
+ publisher, new DefaultExpiryStrategy(), Position.End, false,
+ EventFilter.EventType.Prefixes(false, "matching"), SystemAccounts.System, false,
+ null, 1, CancellationToken.None),
+ _ => throw new ArgumentOutOfRangeException(nameof(kind))
+ };
+
+ private sealed class DelayedSubscriptionPublisher : IPublisher
+ {
+ private readonly ManualResetEventSlim _continueSubscribe = new(false);
+ private readonly ConcurrentDictionary _activeSubscriptions = new();
+ public TaskCompletionSource SubscribeEntered { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ public TaskCompletionSource SubscribeCompleted { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously);
+ public Guid[] ActiveSubscriptions => [.. _activeSubscriptions.Keys];
+
+ public void ReleaseSubscribe() => _continueSubscribe.Set();
+
+ public void Publish(Message message)
+ {
+ switch (message)
+ {
+ case ClientMessage.SubscribeToStream subscription:
+ Register(subscription.CorrelationId, subscription.Envelope);
+ break;
+ case ClientMessage.FilteredSubscribeToStream subscription:
+ Register(subscription.CorrelationId, subscription.Envelope);
+ break;
+ case ClientMessage.UnsubscribeFromStream unsubscribe:
+ _activeSubscriptions.TryRemove(unsubscribe.CorrelationId, out _);
+ break;
+ }
+ }
+
+ private void Register(Guid id, IEnvelope envelope)
+ {
+ SubscribeEntered.TrySetResult(true);
+ _continueSubscribe.Wait(TimeSpan.FromSeconds(10));
+ _activeSubscriptions.TryAdd(id, 0);
+ envelope.ReplyWith(new ClientMessage.SubscriptionConfirmation(id, 0, 0));
+ SubscribeCompleted.TrySetResult(true);
+ }
+ }
+}
diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/Forwarding/ForwardingGrpcCodecTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/Forwarding/ForwardingGrpcCodecTests.cs
index 828810fab5..e4396f5693 100644
--- a/src/EventStore.Core.Tests/Services/Transport/Grpc/Forwarding/ForwardingGrpcCodecTests.cs
+++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/Forwarding/ForwardingGrpcCodecTests.cs
@@ -415,8 +415,6 @@ public void not_handled_leader_info_round_trips()
Guid.NewGuid(),
ClientMessage.NotHandled.Types.NotHandledReason.NotLeader,
new ClientMessage.NotHandled.Types.LeaderInfo(
- new DnsEndPoint("leader-tcp.internal", 1113),
- true,
new DnsEndPoint("leader-http.internal", 2113)));
var decoded = RoundTripResponse(message);
@@ -425,28 +423,6 @@ public void not_handled_leader_info_round_trips()
{
Assert.That(decoded.CorrelationId, Is.EqualTo(message.CorrelationId));
Assert.That(decoded.Reason, Is.EqualTo(message.Reason));
- Assert.That(decoded.LeaderInfo.IsSecure, Is.True);
- Assert.That(decoded.LeaderInfo.ExternalTcp, Is.EqualTo(message.LeaderInfo.ExternalTcp));
- Assert.That(decoded.LeaderInfo.Http, Is.EqualTo(message.LeaderInfo.Http));
- });
- }
-
- [Test]
- public void not_handled_leader_info_without_external_tcp_round_trips_as_null()
- {
- var message = new ClientMessage.NotHandled(
- Guid.NewGuid(),
- ClientMessage.NotHandled.Types.NotHandledReason.NotLeader,
- new ClientMessage.NotHandled.Types.LeaderInfo(
- null,
- false,
- new DnsEndPoint("leader-http.internal", 2113)));
-
- var decoded = RoundTripResponse(message);
-
- Assert.Multiple(() =>
- {
- Assert.That(decoded.LeaderInfo.ExternalTcp, Is.Null);
Assert.That(decoded.LeaderInfo.Http, Is.EqualTo(message.LeaderInfo.Http));
});
}
@@ -454,9 +430,9 @@ public void not_handled_leader_info_without_external_tcp_round_trips_as_null()
[Test]
public void not_authenticated_round_trips()
{
- var message = new TcpMessage.NotAuthenticated(Guid.NewGuid(), "not authenticated");
+ var message = new ClientMessage.NotAuthenticated(Guid.NewGuid(), "not authenticated");
- var decoded = RoundTripResponse(message);
+ var decoded = RoundTripResponse(message);
Assert.Multiple(() =>
{
diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs
index 8c3c96bb6f..03703e4916 100644
--- a/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs
+++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/Replication/GrpcReplicaServiceSupervisorTests.cs
@@ -432,10 +432,6 @@ private static MemberInfo CreateLeader() => MemberInfo.ForVNode(
DateTime.UtcNow,
VNodeState.Leader,
true,
- new DnsEndPoint("leader-replication.internal", 1112),
- null,
- null,
- null,
new DnsEndPoint("leader.internal", 2113),
null,
0,
@@ -444,7 +440,6 @@ private static MemberInfo CreateLeader() => MemberInfo.ForVNode(
0,
0,
0,
- 0,
Guid.NewGuid(),
0,
false,
diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/SubscriptionDisconnectTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/SubscriptionDisconnectTests.cs
new file mode 100644
index 0000000000..ee17e023ea
--- /dev/null
+++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/SubscriptionDisconnectTests.cs
@@ -0,0 +1,83 @@
+using System;
+using System.Threading;
+using System.Threading.Tasks;
+using EventStore.Client.Streams;
+using EventStore.Core.Bus;
+using EventStore.Core.Messages;
+using Google.Protobuf;
+using NUnit.Framework;
+
+namespace EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests;
+
+[TestFixture(typeof(LogFormat.V2), typeof(string))]
+public class SubscriptionDisconnectTests : GrpcSpecification
+{
+ protected override Task Given() => Task.CompletedTask;
+ protected override Task When() => Task.CompletedTask;
+
+ [TestCase("stream")]
+ [TestCase("all")]
+ [TestCase("filtered-all")]
+ public async Task disposing_the_grpc_read_call_unsubscribes_its_live_subscription(string kind)
+ {
+ var unsubscribed = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
+ var handler = new AdHocHandler(message =>
+ unsubscribed.TrySetResult(message.CorrelationId));
+ Node.Node.MainBus.Subscribe(handler);
+ try
+ {
+ using var call = StreamsClient.Read(CreateReadRequest(kind), GetCallOptions(AdminCredentials));
+ Assert.That(await call.ResponseStream.MoveNext(CancellationToken.None), Is.True);
+ var confirmation = call.ResponseStream.Current;
+ Assert.That(confirmation.ContentCase, Is.EqualTo(ReadResp.ContentOneofCase.Confirmation));
+ var subscriptionId = Guid.Parse(confirmation.Confirmation.SubscriptionId);
+
+ call.Dispose();
+
+ Assert.That(await unsubscribed.Task.WaitAsync(TimeSpan.FromSeconds(10)), Is.EqualTo(subscriptionId));
+ }
+ finally
+ {
+ Node.Node.MainBus.Unsubscribe(handler);
+ }
+ }
+
+ private static ReadReq CreateReadRequest(string kind)
+ {
+ var options = new ReadReq.Types.Options
+ {
+ Subscription = new(),
+ UuidOption = new() { Structured = new() },
+ ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards,
+ };
+
+ switch (kind)
+ {
+ case "stream":
+ options.Stream = new()
+ {
+ End = new(),
+ StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8("subscription-disconnect") }
+ };
+ options.NoFilter = new();
+ break;
+ case "all":
+ options.All = new() { End = new() };
+ options.NoFilter = new();
+ break;
+ case "filtered-all":
+ options.All = new() { End = new() };
+ options.Filter = new()
+ {
+ Max = 32,
+ CheckpointIntervalMultiplier = 1,
+ StreamIdentifier = new() { Prefix = { "subscription-disconnect" } }
+ };
+ break;
+ default:
+ throw new ArgumentOutOfRangeException(nameof(kind));
+ }
+
+ return new ReadReq { Options = options };
+ }
+}
diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/core_tcp_package.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/core_tcp_package.cs
deleted file mode 100644
index ffcb22e037..0000000000
--- a/src/EventStore.Core.Tests/Services/Transport/Tcp/core_tcp_package.cs
+++ /dev/null
@@ -1,199 +0,0 @@
-using System;
-using EventStore.Core.Services.Transport.Tcp;
-using NUnit.Framework;
-
-namespace EventStore.Core.Tests.Services.Transport.Tcp;
-
-[TestFixture]
-public class core_tcp_package
-{
- [Test]
- public void should_throw_argument_null_exception_when_created_as_authorized_but_login_not_provided()
- {
- Assert.Throws(() =>
- new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated, Guid.NewGuid(), null, "pa$$",
- new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void should_throw_argument_null_exception_when_created_as_authorized_but_password_not_provided()
- {
- Assert.Throws(() =>
- new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated, Guid.NewGuid(), "login", null,
- new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void should_throw_argument_null_exception_when_created_as_authorized_but_token_not_provided()
- {
- Assert.Throws(() =>
- new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated, Guid.NewGuid(), null,
- new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void should_throw_argument_exception_when_created_as_not_authorized_but_login_is_provided()
- {
- Assert.Throws(() =>
- new TcpPackage(TcpCommand.BadRequest, TcpFlags.None, Guid.NewGuid(), "login", null,
- new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void should_throw_argument_exception_when_created_as_not_authorized_but_password_is_provided()
- {
- Assert.Throws(() =>
- new TcpPackage(TcpCommand.BadRequest, TcpFlags.None, Guid.NewGuid(), null, "pa$$",
- new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void should_throw_argument_exception_when_created_as_not_authorized_but_token_is_provided()
- {
- Assert.Throws(() =>
- new TcpPackage(TcpCommand.BadRequest, TcpFlags.None, Guid.NewGuid(), "token",
- new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void not_authorized_with_data_should_serialize_and_deserialize_correctly()
- {
- var corrId = Guid.NewGuid();
- var refPkg = new TcpPackage(TcpCommand.BadRequest, TcpFlags.None, corrId, null, null, new byte[] { 1, 2, 3 });
- var bytes = refPkg.AsArraySegment();
-
- var pkg = TcpPackage.FromArraySegment(bytes);
- Assert.AreEqual(TcpCommand.BadRequest, pkg.Command);
- Assert.AreEqual(TcpFlags.None, pkg.Flags);
- Assert.AreEqual(corrId, pkg.CorrelationId);
- Assert.False(pkg.Tokens.TryGetValue("uid", out _));
- Assert.False(pkg.Tokens.TryGetValue("pwd", out _));
- Assert.False(pkg.Tokens.TryGetValue("jwt", out _));
-
- Assert.AreEqual(3, pkg.Data.Count);
- Assert.AreEqual(1, pkg.Data.Array[pkg.Data.Offset + 0]);
- Assert.AreEqual(2, pkg.Data.Array[pkg.Data.Offset + 1]);
- Assert.AreEqual(3, pkg.Data.Array[pkg.Data.Offset + 2]);
- }
-
- [Test]
- public void not_authorized_with_empty_data_should_serialize_and_deserialize_correctly()
- {
- var corrId = Guid.NewGuid();
- var refPkg = new TcpPackage(TcpCommand.BadRequest, TcpFlags.None, corrId, null, null, new byte[0]);
- var bytes = refPkg.AsArraySegment();
-
- var pkg = TcpPackage.FromArraySegment(bytes);
- Assert.AreEqual(TcpCommand.BadRequest, pkg.Command);
- Assert.AreEqual(TcpFlags.None, pkg.Flags);
- Assert.AreEqual(corrId, pkg.CorrelationId);
- Assert.False(pkg.Tokens.TryGetValue("uid", out _));
- Assert.False(pkg.Tokens.TryGetValue("pwd", out _));
- Assert.False(pkg.Tokens.TryGetValue("jwt", out _));
-
- Assert.AreEqual(0, pkg.Data.Count);
- }
-
- [Test]
- public void authorized_with_data_should_serialize_and_deserialize_correctly()
- {
- var corrId = Guid.NewGuid();
- var refPkg = new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated, corrId, "login", "pa$$",
- new byte[] { 1, 2, 3 });
- var bytes = refPkg.AsArraySegment();
-
- var pkg = TcpPackage.FromArraySegment(bytes);
- Assert.AreEqual(TcpCommand.BadRequest, pkg.Command);
- Assert.AreEqual(TcpFlags.Authenticated, pkg.Flags);
- Assert.AreEqual(corrId, pkg.CorrelationId);
- Assert.AreEqual("login", pkg.Tokens["uid"]);
- Assert.AreEqual("pa$$", pkg.Tokens["pwd"]);
- Assert.False(pkg.Tokens.TryGetValue("jwt", out _));
-
- Assert.AreEqual(3, pkg.Data.Count);
- Assert.AreEqual(1, pkg.Data.Array[pkg.Data.Offset + 0]);
- Assert.AreEqual(2, pkg.Data.Array[pkg.Data.Offset + 1]);
- Assert.AreEqual(3, pkg.Data.Array[pkg.Data.Offset + 2]);
- }
-
- [Test]
- public void authorized_with_empty_data_should_serialize_and_deserialize_correctly()
- {
- var corrId = Guid.NewGuid();
- var refPkg = new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated, corrId, "login", "pa$$",
- new byte[0]);
- var bytes = refPkg.AsArraySegment();
-
- var pkg = TcpPackage.FromArraySegment(bytes);
- Assert.AreEqual(TcpCommand.BadRequest, pkg.Command);
- Assert.AreEqual(TcpFlags.Authenticated, pkg.Flags);
- Assert.AreEqual(corrId, pkg.CorrelationId);
- Assert.AreEqual("login", pkg.Tokens["uid"]);
- Assert.AreEqual("pa$$", pkg.Tokens["pwd"]);
- Assert.False(pkg.Tokens.TryGetValue("jwt", out _));
-
- Assert.AreEqual(0, pkg.Data.Count);
- }
-
- [Test]
- public void token_authorized_with_data_should_serialize_and_deserialize_correctly()
- {
- var corrId = Guid.NewGuid();
- var refPkg = new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated, corrId, "token",
- new byte[] { 1, 2, 3 });
- var bytes = refPkg.AsArraySegment();
-
- var pkg = TcpPackage.FromArraySegment(bytes);
- Assert.AreEqual(TcpCommand.BadRequest, pkg.Command);
- Assert.AreEqual(TcpFlags.Authenticated, pkg.Flags);
- Assert.AreEqual(corrId, pkg.CorrelationId);
- Assert.AreEqual("token", pkg.Tokens["jwt"]);
- Assert.False(pkg.Tokens.TryGetValue("uid", out _));
- Assert.False(pkg.Tokens.TryGetValue("pwd", out _));
-
- Assert.AreEqual(3, pkg.Data.Count);
- Assert.AreEqual(1, pkg.Data.Array[pkg.Data.Offset + 0]);
- Assert.AreEqual(2, pkg.Data.Array[pkg.Data.Offset + 1]);
- Assert.AreEqual(3, pkg.Data.Array[pkg.Data.Offset + 2]);
- }
-
- [Test]
- public void token_authorized_with_empty_data_should_serialize_and_deserialize_correctly()
- {
- var corrId = Guid.NewGuid();
- var refPkg = new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated, corrId, "token",
- new byte[0]);
- var bytes = refPkg.AsArraySegment();
-
- var pkg = TcpPackage.FromArraySegment(bytes);
- Assert.AreEqual(TcpCommand.BadRequest, pkg.Command);
- Assert.AreEqual(TcpFlags.Authenticated, pkg.Flags);
- Assert.AreEqual(corrId, pkg.CorrelationId);
- Assert.AreEqual("token", pkg.Tokens["jwt"]);
- Assert.False(pkg.Tokens.TryGetValue("uid", out _));
- Assert.False(pkg.Tokens.TryGetValue("pwd", out _));
-
- Assert.AreEqual(0, pkg.Data.Count);
- }
-
- [Test]
- public void should_throw_argument_exception_when_login_too_long()
- {
- Assert.Throws(() => new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated,
- Guid.NewGuid(), new string('*', TcpPackage.MaxLoginLength + 1), "pa$$", new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void should_throw_argument_exception_when_password_too_long()
- {
- Assert.Throws(() => new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated,
- Guid.NewGuid(), "login", new string('*', TcpPackage.MaxPasswordLength + 1), new byte[] { 1, 2, 3 }));
- }
-
- [Test]
- public void should_throw_argument_exception_when_token_too_long()
- {
- Assert.Throws(() => new TcpPackage(TcpCommand.BadRequest, TcpFlags.Authenticated,
- Guid.NewGuid(), new string('*', TcpPackage.MaxTokenLength + 1), new byte[] { 1, 2, 3 }));
- }
-}
diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connection.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connection.cs
index f6165d0f82..40a61a7f5d 100644
--- a/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connection.cs
+++ b/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connection.cs
@@ -8,7 +8,6 @@
using System.Text;
using System.Threading;
using EventStore.Common.Utils;
-using EventStore.Core.Services.Transport.Tcp;
using EventStore.Core.Tests.Helpers;
using EventStore.Transport.Tcp;
using NUnit.Framework;
@@ -89,7 +88,7 @@ public void should_connect_to_each_other_and_send_data()
{ return (true, null); },
null,
new TcpClientConnector(),
- TcpConnectionManager.ConnectionTimeout,
+ TimeSpan.FromSeconds(1),
conn =>
{
Log.Information("Sending bytes...");
diff --git a/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connections_mutual_auth.cs b/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connections_mutual_auth.cs
index 1a80ae435e..3a0c30b772 100644
--- a/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connections_mutual_auth.cs
+++ b/src/EventStore.Core.Tests/Services/Transport/Tcp/ssl_connections_mutual_auth.cs
@@ -8,7 +8,6 @@
using System.Security.Cryptography.X509Certificates;
using System.Threading;
using EventStore.Common.Utils;
-using EventStore.Core.Services.Transport.Tcp;
using EventStore.Core.Tests.Helpers;
using EventStore.Transport.Tcp;
using NUnit.Framework;
@@ -116,7 +115,7 @@ bool shouldConnectSuccessfully
(cert, chain, err, _) => validateServerCertificate ? ClusterVNode.ValidateServerCertificate(cert, chain, err, () => null, () => rootCertificates, null) : (true, null),
() => new X509CertificateCollection { clientCertificate },
new TcpClientConnector(),
- TcpConnectionManager.ConnectionTimeout,
+ TimeSpan.FromSeconds(1),
conn =>
{
Log.Information("Sending bytes...");
diff --git a/src/EventStore.Core.Tests/Services/VNode/InaugurationManager/InaugurationManagerTests.cs b/src/EventStore.Core.Tests/Services/VNode/InaugurationManager/InaugurationManagerTests.cs
index 8ae7429f0c..fc9ce572a8 100644
--- a/src/EventStore.Core.Tests/Services/VNode/InaugurationManager/InaugurationManagerTests.cs
+++ b/src/EventStore.Core.Tests/Services/VNode/InaugurationManager/InaugurationManagerTests.cs
@@ -20,8 +20,7 @@ public abstract class InaugurationManagerTests
protected readonly MemberInfo _leader =
MemberInfo.ForVNode(
default, default, default, default,
- new DnsEndPoint("localhost", default), default, default, default,
- new DnsEndPoint("localhost", default), default, default, default,
+ new DnsEndPoint("localhost", default), default, default,
default, default, default, default, default, default, default, default);
protected readonly long _replicationTarget = 400;
protected readonly long _indexTarget = 400;
diff --git a/src/EventStore.Core.Tests/Services/VNode/ShutdownServiceTests.cs b/src/EventStore.Core.Tests/Services/VNode/ShutdownServiceTests.cs
index f2adb06c98..d1011fa934 100644
--- a/src/EventStore.Core.Tests/Services/VNode/ShutdownServiceTests.cs
+++ b/src/EventStore.Core.Tests/Services/VNode/ShutdownServiceTests.cs
@@ -1,5 +1,4 @@
using System;
-using System.Net;
using DotNext.Net.Http;
using EventStore.Core.Data;
using EventStore.Core.Messages;
@@ -16,10 +15,6 @@ public class ShutdownServiceTests
= new(
Guid.NewGuid(),
0,
- new IPEndPoint(0, 0),
- new IPEndPoint(IPAddress.Loopback, 1),
- new IPEndPoint(IPAddress.Loopback, 2),
- new IPEndPoint(IPAddress.Loopback, 3),
new HttpEndPoint(new Uri("http://www.trogondb.com")), true);
[Test]
diff --git a/src/EventStore.Core.Tests/Services/VNode/leader_info_provider.cs b/src/EventStore.Core.Tests/Services/VNode/leader_info_provider.cs
index 68de1f0112..9a34d143df 100644
--- a/src/EventStore.Core.Tests/Services/VNode/leader_info_provider.cs
+++ b/src/EventStore.Core.Tests/Services/VNode/leader_info_provider.cs
@@ -1,8 +1,6 @@
using System;
using System.Collections.Generic;
-using System.Linq;
using System.Net;
-using EventStore.Common.Utils;
using EventStore.Core.Cluster;
using EventStore.Core.Data;
using EventStore.Core.Services.VNode;
@@ -13,252 +11,45 @@ namespace EventStore.Core.Tests.Services.VNode;
[TestFixture]
public class leader_info_provider
{
- private const string DefaultHttpEndPoint = "9.9.9.9:9";
-
- public static IEnumerable |