Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 9 additions & 5 deletions src/OneScript.Web.Server/PropertyWrappersCollection.cs
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,16 @@ public class PropertyWrappersCollection

public T Get<T>(string propName, Func<T> 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;
}
}
}
30 changes: 23 additions & 7 deletions src/OneScript.Web.Server/WebServer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -28,7 +30,9 @@ namespace OneScript.Web.Server
public class WebServer: AutoContext<WebServer>
{
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;
Expand Down Expand Up @@ -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);
}
}

Expand Down
48 changes: 25 additions & 23 deletions src/OneScript.Web.Server/WebSockets/WebSocketWrapper.cs
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,10 @@ public class WebSocketWrapper: AutoContext<WebSocketWrapper>
{
private readonly WebSocket _webSocket;

// WebSocket допускает только одно ожидающее получение, а сообщение из нескольких частей
// должен прочитать целиком один поток
private readonly object _receiveLock = new object();

public WebSocketWrapper(WebSocket webSocket)
{
_webSocket = webSocket;
Expand Down Expand Up @@ -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);
}
Expand All @@ -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()));
}

/// <summary>
Expand All @@ -133,21 +128,28 @@ public BslStringValue ReceiveString()
/// <returns></returns>
[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();
}

/// <summary>
Expand Down
189 changes: 189 additions & 0 deletions src/Tests/OneScript.Core.Tests/WebServerConcurrencyTests.cs
Original file line number Diff line number Diff line change
@@ -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;
}
}
}