This repository was archived by the owner on Aug 12, 2025. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathProgram.cs
More file actions
170 lines (150 loc) · 5.99 KB
/
Copy pathProgram.cs
File metadata and controls
170 lines (150 loc) · 5.99 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
using System.Net;
using InteractiveConsole;
using Serilog;
using TaskFlux.Transport.Tcp.Client;
using TaskFlux.Utils.Network;
// ReSharper disable AccessToDisposedClosure
/*
* Пример интерактивного клиента.
* Взаимодействие осуществляется через консоль.
*
* В данном примере используется кластер из одного узла.
* Запустить узел можно из docker-compose.yaml, расположенного в директории проекта.
*/
using var cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, eventArgs) =>
{
cts.Cancel();
eventArgs.Cancel = true;
};
if (!TryGetEndpoints(out var endpoints))
{
return 1;
}
var clientFactory = new TaskFluxClientFactory(endpoints);
Log.Logger.Debug($"Создаю клиента");
await using var client = await clientFactory.ConnectAsync(cts.Token);
Log.Logger.Debug("Клиент создан");
while (!cts.IsCancellationRequested)
{
// 1. Читаю команду
string commandString;
try
{
var input = await ReadInputAsync(cts.Token);
if (input is null)
{
break;
}
commandString = input;
}
catch (OperationCanceledException)
{
break;
}
if (commandString.StartsWith("help", StringComparison.InvariantCultureIgnoreCase))
{
PrintHelp();
continue;
}
if (commandString.Equals("exit", StringComparison.InvariantCultureIgnoreCase))
{
break;
}
try
{
var command = StringCommandParser.ParseCommand(commandString);
try
{
await command.Execute(client, cts.Token);
}
catch (OperationCanceledException)
{
break;
}
}
catch (Exception e)
{
Console.WriteLine($"Ошибка выполнения команды: {e.Message}");
}
}
return 0;
static void PrintHelp()
{
Console.WriteLine($"Использование: COMMAND [ARGS...]");
Console.WriteLine($"Основные команды:");
foreach (var commandDescription in GetCommandDescriptions())
{
Console.WriteLine(commandDescription);
}
}
static IEnumerable<string> GetCommandDescriptions()
{
yield return """
- enqueue [QUEUE_NAME] KEY VALUES... - Вставить элемент в очередь
QUEUE_NAME - название очереди. Пропустить, если использовать стандартную
KEY - ключ для вставляемого значения.
VALUES... - разделенные пробелом слова, которые будут добавлены в нагрузку
""";
yield return """
- dequeue [QUEUE_NAME] - получить элемент из очереди
QUEUE_NAME - название очереди. Пропустить, если использовать стандартную
""";
yield return """
- create QUEUE_NAME [WITHMAXSIZE max_size] [WITHMAXPAYLOAD max_payload] [WITHPRIORITYRANGE min max] [TYPE code] - создать новую очередь с указанным названием
QUEUE_NAME - название очереди
WITHMAXSIZE max_size - выставить ограничение на максимальный размер очереди в max_size
WITHMAXPAYLOAD max_payload - выставить ограничение на максимальный размер сообщения в max_payload байтов
WITHPRIORITYRANGE min max - ограничить допустимый диапазон выставляемых ключей с min до max включительно
TYPE code - использовать указанную реализацию структуры для хранения. code - код структуры
""";
yield return """
- delete QUEUE_NAME - удалить очередь с указанным названием
QUEUE_NAME - название очереди
""";
yield return """
- count [QUEUE_NAME] - получить размер очереди
QUEUE_NAME - название очереди. Пропустить, если использовать очередь по умолчанию
""";
yield return """
- list - получить список всех очередей и их данных
""";
yield return """
- help - вывести это сообщение
""";
}
async Task<string?> ReadInputAsync(CancellationToken token)
{
// Консольные ReadAsync всегда блокирующие, нужны отдельные таски
Console.Write("$> ");
var readTask = Task.Run(Console.ReadLine);
var waitTask = Task.Delay(Timeout.Infinite, token);
await Task.WhenAny(readTask, waitTask);
return token.IsCancellationRequested
? null
: readTask.Result;
}
bool TryGetEndpoints(out EndPoint[] endpoints)
{
if (args.Length == 0)
{
Console.WriteLine(
$"Необходимо передать адреса узлов кластера: {AppDomain.CurrentDomain.FriendlyName} [АДРЕС ...]");
endpoints = default!;
return false;
}
endpoints = new EndPoint[args.Length];
for (var i = 0; i < endpoints.Length; i++)
{
try
{
endpoints[i] = EndPointHelpers.ParseEndPoint(args[i]);
}
catch (ArgumentException ae)
{
Console.WriteLine($"Ошибка парсинга адреса {args[i]}");
return false;
}
}
return true;
}