diff --git a/src/OneScript.Web.Server/PropertyWrappersCollection.cs b/src/OneScript.Web.Server/PropertyWrappersCollection.cs index e6dfa95e4..741b881e7 100644 --- a/src/OneScript.Web.Server/PropertyWrappersCollection.cs +++ b/src/OneScript.Web.Server/PropertyWrappersCollection.cs @@ -19,12 +19,16 @@ public class PropertyWrappersCollection public T Get(string propName, Func factory) { - if (_objects.TryGetValue(propName, out var value)) - return (T)value; + // Контекст запроса могут передать в фоновые задания + lock (_objects) + { + if (_objects.TryGetValue(propName, out var value)) + return (T)value; - value = factory(); - _objects.Add(propName, value); + value = factory(); + _objects.Add(propName, value); - return (T)value; + return (T)value; + } } } \ No newline at end of file diff --git a/src/OneScript.Web.Server/WebServer.cs b/src/OneScript.Web.Server/WebServer.cs index d06a876b9..c2dd59a1b 100644 --- a/src/OneScript.Web.Server/WebServer.cs +++ b/src/OneScript.Web.Server/WebServer.cs @@ -9,6 +9,8 @@ This Source Code Form is subject to the terms of the using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Http.Features; using OneScript.Contexts; +using OneScript.Exceptions; +using OneScript.Localization; using OneScript.Types; using ScriptEngine.Machine; using ScriptEngine.Machine.Contexts; @@ -28,7 +30,9 @@ namespace OneScript.Web.Server public class WebServer: AutoContext { private readonly ExecutionContext _executionContext; - private WebApplication _app; + // Остановить вызывают из другого потока, чем Запустить + private volatile WebApplication _app; + private int _isRunning; private readonly List<(IRuntimeContextInstance Target, string MethodName)> _middlewares = new List<(IRuntimeContextInstance Target, string MethodName)>(); private string _contentRoot = null; @@ -75,19 +79,31 @@ public static WebServer Constructor(TypeActivationContext typeActivationContext, [ContextMethod("Запустить", "Run")] public void Run() { - ConfigureApp(); + // Второй запуск того же сервера (например, из фонового задания) подменил бы приложение: + // Остановить и освобождение первого запуска пришлись бы на чужое + if (System.Threading.Interlocked.Exchange(ref _isRunning, 1) != 0) + throw new RuntimeException(new BilingualString("Веб-сервер уже запущен", "Web server is already running")); try { - _app.Start(); - if (Port == 0) - Port = new Uri(_app.Urls.First()).Port; + ConfigureApp(); - _app.WaitForShutdown(); + try + { + _app.Start(); + if (Port == 0) + Port = new Uri(_app.Urls.First()).Port; + + _app.WaitForShutdown(); + } + finally + { + _app.DisposeAsync().AsTask().Wait(); + } } finally { - _app.DisposeAsync().AsTask().Wait(); + System.Threading.Volatile.Write(ref _isRunning, 0); } } diff --git a/src/OneScript.Web.Server/WebSockets/WebSocketWrapper.cs b/src/OneScript.Web.Server/WebSockets/WebSocketWrapper.cs index e8dd2de48..c8d3639d6 100644 --- a/src/OneScript.Web.Server/WebSockets/WebSocketWrapper.cs +++ b/src/OneScript.Web.Server/WebSockets/WebSocketWrapper.cs @@ -24,6 +24,10 @@ public class WebSocketWrapper: AutoContext { private readonly WebSocket _webSocket; + // WebSocket допускает только одно ожидающее получение, а сообщение из нескольких частей + // должен прочитать целиком один поток + private readonly object _receiveLock = new object(); + public WebSocketWrapper(WebSocket webSocket) { _webSocket = webSocket; @@ -99,7 +103,11 @@ public void CloseOutput(WebSocketCloseStatusWrapper status, string statusDescrip [ContextMethod("Получить", "Receive")] public WebSocketReceiveResultWrapper Receive(BinaryDataBuffer buffer) { - var result = _webSocket.ReceiveAsync(buffer.Bytes, default).Result; + WebSocketReceiveResult result; + lock (_receiveLock) + { + result = _webSocket.ReceiveAsync(buffer.Bytes, default).Result; + } return new WebSocketReceiveResultWrapper(result); } @@ -111,20 +119,7 @@ public WebSocketReceiveResultWrapper Receive(BinaryDataBuffer buffer) [ContextMethod("ПолучитьСтроку", "ReceiveString")] public BslStringValue ReceiveString() { - var buffer = new byte[1024]; - using var stream = new MemoryStream(); - - WebSocketReceiveResult result; - do - { - result = _webSocket.ReceiveAsync(buffer, default).Result; - stream.Write(buffer); - } - while (!result.EndOfMessage); - - var data = stream.GetBuffer(); - - return BslStringValue.Create(Encoding.UTF8.GetString(data)); + return BslStringValue.Create(Encoding.UTF8.GetString(ReceiveMessage())); } /// @@ -133,21 +128,28 @@ public BslStringValue ReceiveString() /// [ContextMethod("ПолучитьДвоичныеДанные", "ReceiveBinary")] public byte[] ReceiveBinary() + { + return ReceiveMessage(); + } + + private byte[] ReceiveMessage() { var buffer = new byte[1024]; using var stream = new MemoryStream(); - WebSocketReceiveResult result; - do + lock (_receiveLock) { - result = _webSocket.ReceiveAsync(buffer, default).Result; - stream.Write(buffer); + WebSocketReceiveResult result; + do + { + result = _webSocket.ReceiveAsync(buffer, default).Result; + // Только полученные байты: буфер заполнен не весь + stream.Write(buffer, 0, result.Count); + } + while (!result.EndOfMessage); } - while (!result.EndOfMessage); - - var data = stream.GetBuffer(); - return data; + return stream.ToArray(); } /// diff --git a/src/Tests/OneScript.Core.Tests/WebServerConcurrencyTests.cs b/src/Tests/OneScript.Core.Tests/WebServerConcurrencyTests.cs new file mode 100644 index 000000000..007d00992 --- /dev/null +++ b/src/Tests/OneScript.Core.Tests/WebServerConcurrencyTests.cs @@ -0,0 +1,189 @@ +/*---------------------------------------------------------- +This Source Code Form is subject to the terms of the +Mozilla Public License, v.2.0. If a copy of the MPL +was not distributed with this file, You can obtain one +at http://mozilla.org/MPL/2.0/. +----------------------------------------------------------*/ + +using System; +using System.Net; +using System.Net.Sockets; +using System.Net.WebSockets; +using System.Text; +using System.Threading; +using FluentAssertions; +using OneScript.StandardLibrary; +using OneScript.Web.Server; +using ScriptEngine.HostedScript; +using ScriptEngine.HostedScript.Extensions; +using ScriptEngine.Hosting; +using ScriptEngine.Machine; +using Xunit; + +namespace OneScript.Core.Tests +{ + public class WebServerConcurrencyTests + { + private const string WebSocketScript = + "Перем мСервер Экспорт;\n" + + "\n" + + "Функция Прочитать(Сокет) Экспорт\n" + + " Возврат Сокет.ПолучитьСтроку();\n" + + "КонецФункции\n" + + "\n" + + "Функция ОбработчикЗапроса(Контекст, СледующийОбработчик) Экспорт\n" + + " Сокет = Контекст.ВебСокеты.ПодключитьВебСокет();\n" + + " Параметры = Новый Массив;\n" + + " Параметры.Добавить(Сокет);\n" + + " Задания = Новый Массив;\n" + + " // Два задания ждут сообщения из одного сокета одновременно\n" + + " Задания.Добавить(ФоновыеЗадания.Выполнить(ЭтотОбъект, \"Прочитать\", Параметры, Истина));\n" + + " Задания.Добавить(ФоновыеЗадания.Выполнить(ЭтотОбъект, \"Прочитать\", Параметры, Истина));\n" + + " ФоновыеЗадания.ОжидатьВсе(Задания);\n" + + " Ответы = Новый Массив;\n" + + " Для Каждого Задание Из Задания Цикл\n" + + " Если Задание.ИнформацияОбОшибке = Неопределено Тогда\n" + + " Ответы.Добавить(Задание.Результат);\n" + + " Иначе\n" + + " Ответы.Добавить(\"ошибка: \" + Задание.ИнформацияОбОшибке.Описание);\n" + + " КонецЕсли;\n" + + " КонецЦикла;\n" + + " Сокет.ОтправитьСтроку(СтрСоединить(Ответы, \"|\"));\n" + + "КонецФункции\n" + + "\n" + + "Процедура ЗапуститьСервер() Экспорт\n" + + " мСервер.Запустить();\n" + + "КонецПроцедуры\n" + + "\n" + + "мСервер = Новый ВебСервер(Порт);\n" + + "мСервер.ИспользоватьВебСокеты();\n" + + "мСервер.ДобавитьОбработчикЗапросов(ЭтотОбъект, \"ОбработчикЗапроса\");\n" + + "Задание = ФоновыеЗадания.Выполнить(ЭтотОбъект, \"ЗапуститьСервер\");\n"; + + private const string SecondRunScript = + "Перем мСервер;\n" + + "Перем ОшибкаПовторногоЗапуска Экспорт;\n" + + "Перем Остановлен Экспорт;\n" + + "\n" + + "Функция ОбработчикЗапроса(Контекст, СледующийОбработчик) Экспорт\n" + + " Контекст.Ответ.КодСостояния = 200;\n" + + "КонецФункции\n" + + "\n" + + "Процедура ЗапуститьСервер() Экспорт\n" + + " мСервер.Запустить();\n" + + "КонецПроцедуры\n" + + "\n" + + "мСервер = Новый ВебСервер(0);\n" + + "мСервер.ДобавитьОбработчикЗапросов(ЭтотОбъект, \"ОбработчикЗапроса\");\n" + + "Задания = Новый Массив;\n" + + "Задания.Добавить(ФоновыеЗадания.Выполнить(ЭтотОбъект, \"ЗапуститьСервер\"));\n" + + "Для Номер = 1 По 200 Цикл\n" + + " Если мСервер.Порт <> 0 Тогда\n" + + " Прервать;\n" + + " КонецЕсли;\n" + + " Приостановить(50);\n" + + "КонецЦикла;\n" + + "Попытка\n" + + " мСервер.Запустить();\n" + + " ОшибкаПовторногоЗапуска = \"\";\n" + + "Исключение\n" + + " ОшибкаПовторногоЗапуска = ИнформацияОбОшибке().Описание;\n" + + "КонецПопытки;\n" + + "мСервер.Остановить();\n" + + "Остановлен = ФоновыеЗадания.ОжидатьВсе(Задания, 10000);\n"; + + [Fact] + public void WebSocketMessagesAreReceivedWholeFromManyTasks() + { + var engine = CreateEngine(); + var port = FreePort(); + var context = new ExternalContextData + { + { "Порт", ValueFactory.Create(port) } + }; + var instance = engine.Engine.AttachedScriptsFactory.LoadFromString( + engine.GetCompilerService(), WebSocketScript, engine.Engine.NewProcess(), context); + var server = (WebServer)instance.GetPropValue(instance.GetPropertyNumber("мСервер")); + try + { + using var client = Connect(new Uri($"ws://127.0.0.1:{port}/")); + Send(client, "раз"); + Send(client, "два"); + + var answers = Receive(client).Split('|'); + + answers.Should().BeEquivalentTo("раз", "два"); + } + finally + { + server.Stop(); + } + } + + [Fact] + public void SecondRunOfSameServerIsRejected() + { + var engine = CreateEngine(); + var instance = engine.Engine.AttachedScriptsFactory.LoadFromString( + engine.GetCompilerService(), SecondRunScript, engine.Engine.NewProcess()); + + instance.GetPropValue(instance.GetPropertyNumber("ОшибкаПовторногоЗапуска")).ToString() + .Should().Contain("уже запущен"); + instance.GetPropValue(instance.GetPropertyNumber("Остановлен")).AsBoolean() + .Should().BeTrue("Остановить должен остановить первый запуск"); + } + + private static ClientWebSocket Connect(Uri uri) + { + for (var attempt = 0; ; attempt++) + { + var client = new ClientWebSocket(); + try + { + client.ConnectAsync(uri, CancellationToken.None).Wait(); + return client; + } + catch (AggregateException) when (attempt < 100) + { + // Сервер еще запускается в фоновом задании + client.Dispose(); + Thread.Sleep(100); + } + } + } + + private static void Send(ClientWebSocket client, string text) + { + client.SendAsync(Encoding.UTF8.GetBytes(text), WebSocketMessageType.Text, true, CancellationToken.None).Wait(); + } + + private static string Receive(ClientWebSocket client) + { + var buffer = new byte[4096]; + using var cancellation = new CancellationTokenSource(TimeSpan.FromSeconds(10)); + var result = client.ReceiveAsync(buffer, cancellation.Token).Result; + return Encoding.UTF8.GetString(buffer, 0, result.Count); + } + + private static HostedScriptEngine CreateEngine() + { + var builder = DefaultEngineBuilder.Create() + .SetDefaultOptions() + .UseImports() + .UseDefaultHosting() + .SetupEnvironment(e => e.AddStandardLibrary().AddWebServer()); + var engine = new HostedScriptEngine(builder.Build()); + engine.Initialize(); + return engine; + } + + private static int FreePort() + { + var listener = new TcpListener(IPAddress.Loopback, 0); + listener.Start(); + var port = ((IPEndPoint)listener.LocalEndpoint).Port; + listener.Stop(); + return port; + } + } +}