Commit da5ba245 by Michael Brachmann

web socket server replace

parent 8e5b49f2
unit WebSocketServer;
interface
uses
System.SysUtils, System.Generics.Collections,
IdCustomTCPServer, IdTCPConnection, IdContext, IdIOHandler, IdGlobal, IdCoderMIME, IdHashSHA,
IdSSL, IdSSLOpenSSL;
type
TWebSocketServer = class(TIdCustomTCPServer)
private
IdServerIOHandlerSSLOpenSSL: TIdServerIOHandlerSSLOpenSSL;
HashSHA1: TIdHashSHA1;
protected
procedure DoConnect(AContext: TIdContext); override;
function DoExecute(AContext: TIdContext): Boolean; override;
public
procedure InitSSL(AIdServerIOHandlerSSLOpenSSL: TIdServerIOHandlerSSLOpenSSL);
property OnExecute;
constructor Create;
destructor Destroy; override;
end;
TWebSocketIOHandlerHelper = class(TIdIOHandler)
public
function ReadBytes: TArray<byte>;
function ReadString: string;
procedure WriteBytes(RawData: TArray<byte>);
procedure WriteString(const str: string);
end;
implementation
function HeadersParse(const msg: string): TDictionary<string, string>;
var
lines: TArray<string>;
line: string;
SplittedLine: TArray<string>;
begin
result := TDictionary<string, string>.Create;
lines := msg.Split([#13#10]);
for line in lines do
begin
SplittedLine := line.Split([': ']);
if Length(SplittedLine) > 1 then
result.AddOrSetValue(Trim(SplittedLine[0]), Trim(SplittedLine[1]));
end;
end;
{ TWebSocketServer }
constructor TWebSocketServer.Create;
begin
inherited Create;
HashSHA1 := TIdHashSHA1.Create;
IdServerIOHandlerSSLOpenSSL := nil;
end;
destructor TWebSocketServer.Destroy;
begin
HashSHA1.DisposeOf;
inherited;
end;
procedure TWebSocketServer.InitSSL(AIdServerIOHandlerSSLOpenSSL: TIdServerIOHandlerSSLOpenSSL);
var
CurrentActive: boolean;
begin
CurrentActive := Active;
if CurrentActive then
Active := false;
IdServerIOHandlerSSLOpenSSL := AIdServerIOHandlerSSLOpenSSL;
IOHandler := AIdServerIOHandlerSSLOpenSSL;
if CurrentActive then
Active := true;
end;
procedure TWebSocketServer.DoConnect(AContext: TIdContext);
begin
if AContext.Connection.IOHandler is TIdSSLIOHandlerSocketBase then
TIdSSLIOHandlerSocketBase(AContext.Connection.IOHandler).PassThrough := false;
// Mark connection as "not handshaked"
AContext.Connection.IOHandler.Tag := -1;
inherited;
end;
function TWebSocketServer.DoExecute(AContext: TIdContext): Boolean;
var
c: TIdIOHandler;
Bytes: TArray<byte>;
msg, SecWebSocketKey, Hash: string;
ParsedHeaders: TDictionary<string, string>;
begin
c := AContext.Connection.IOHandler;
// Handshake
if c.Tag = -1 then
begin
c.CheckForDataOnSource(10);
if not c.InputBufferIsEmpty then
begin
// Read string and parse HTTP headers
try
c.InputBuffer.ExtractToBytes(TIdBytes(Bytes));
msg := IndyTextEncoding_UTF8.GetString(TIdBytes(Bytes));
except
end;
ParsedHeaders := HeadersParse(msg);
if ParsedHeaders.ContainsKey('Upgrade') and (ParsedHeaders['Upgrade'] = 'websocket') and
ParsedHeaders.ContainsKey('Sec-WebSocket-Key') then
begin
// Handle handshake request
// https://developer.mozilla.org/en-US/docs/Web/API/WebSockets_API/Writing_WebSocket_servers
SecWebSocketKey := ParsedHeaders['Sec-WebSocket-Key'];
// Send handshake response
Hash := TIdEncoderMIME.EncodeBytes(
HashSHA1.HashString(SecWebSocketKey + '258EAFA5-E914-47DA-95CA-C5AB0DC85B11'));
try
c.Write('HTTP/1.1 101 Switching Protocols'#13#10
+ 'Upgrade: websocket'#13#10
+ 'Connection: Upgrade'#13#10
+ 'Sec-WebSocket-Accept: ' + Hash
+ #13#10#13#10, IndyTextEncoding_UTF8);
except
end;
// Mark IOHandler as handshaked
c.Tag := 1;
end;
ParsedHeaders.DisposeOf;
end;
end;
Result := inherited;
end;
{ TWebSocketIOHandlerHelper }
function TWebSocketIOHandlerHelper.ReadBytes: TArray<byte>;
var
l: byte;
b: array [0..7] of byte;
i, DecodedSize: int64;
Mask: array [0..3] of byte;
begin
// https://stackoverflow.com/questions/8125507/how-can-i-send-and-receive-websocket-messages-on-the-server-side
try
if ReadByte = $81 then
begin
l := ReadByte;
case l of
$FE:
begin
b[1] := ReadByte; b[0] := ReadByte;
b[2] := 0; b[3] := 0; b[4] := 0; b[5] := 0; b[6] := 0; b[7] := 0;
DecodedSize := Int64(b);
end;
$FF:
begin
b[7] := ReadByte; b[6] := ReadByte; b[5] := ReadByte; b[4] := ReadByte;
b[3] := ReadByte; b[2] := ReadByte; b[1] := ReadByte; b[0] := ReadByte;
DecodedSize := Int64(b);
end;
else
DecodedSize := l - 128;
end;
Mask[0] := ReadByte; Mask[1] := ReadByte; Mask[2] := ReadByte; Mask[3] := ReadByte;
if DecodedSize < 1 then
begin
result := [];
exit;
end;
SetLength(result, DecodedSize);
inherited ReadBytes(TIdBytes(result), DecodedSize, False);
for i := 0 to DecodedSize - 1 do
result[i] := result[i] xor Mask[i mod 4];
end;
except
end;
end;
procedure TWebSocketIOHandlerHelper.WriteBytes(RawData: TArray<byte>);
var
Msg: TArray<byte>;
begin
// https://stackoverflow.com/questions/8125507/how-can-i-send-and-receive-websocket-messages-on-the-server-side
Msg := [$81];
if Length(RawData) <= 125 then
Msg := Msg + [Length(RawData)]
else if (Length(RawData) >= 126) and (Length(RawData) <= 65535) then
Msg := Msg + [126, (Length(RawData) shr 8) and 255, Length(RawData) and 255]
else
Msg := Msg + [127, (int64(Length(RawData)) shr 56) and 255, (int64(Length(RawData)) shr 48) and 255,
(int64(Length(RawData)) shr 40) and 255, (int64(Length(RawData)) shr 32) and 255,
(Length(RawData) shr 24) and 255, (Length(RawData) shr 16) and 255, (Length(RawData) shr 8) and 255, Length(RawData) and 255];
Msg := Msg + RawData;
try
Write(TIdBytes(Msg), Length(Msg));
except
end;
end;
function TWebSocketIOHandlerHelper.ReadString: string;
begin
result := IndyTextEncoding_UTF8.GetString(TIdBytes(ReadBytes));
end;
procedure TWebSocketIOHandlerHelper.WriteString(const str: string);
begin
WriteBytes(TArray<byte>(IndyTextEncoding_UTF8.GetBytes(str)));
end;
end.
object WsServerModule: TWsServerModule object WsServerModule: TWsServerModule
Height = 273 OldCreateOrder = False
Width = 230 OnCreate = DataModuleCreate
object SparkleHttpSysDispatcher3: TSparkleHttpSysDispatcher OnDestroy = DataModuleDestroy
Left = 84 Height = 150
Top = 30 Width = 150
end
object XDataServer1: TXDataServer
Dispatcher = SparkleHttpSysDispatcher3
EntitySetPermissions = <>
SwaggerOptions.Enabled = True
SwaggerUIOptions.Enabled = True
SwaggerUIOptions.ShowFilter = True
SwaggerUIOptions.TryItOutEnabled = True
Left = 85
Top = 110
object XDataServer1Logging: TSparkleGenericMiddleware
OnMiddlewareCreate = XDataServer1LoggingMiddlewareCreate
end
object XDataServer1CORS: TSparkleCorsMiddleware
Methods = 'Get'
end
end
end end
...@@ -4,33 +4,27 @@ interface ...@@ -4,33 +4,27 @@ interface
uses uses
System.SysUtils, System.Classes, System.SysUtils, System.Classes,
XData.Server.Module, IdContext,
XData.Comp.Server, WebSocketServer;
Sparkle.Comp.Server,
Sparkle.Comp.HttpSysDispatcher,
Sparkle.Comp.CorsMiddleware,
Sparkle.Comp.GenericMiddleware,
Sparkle.HttpServer.Module,
Sparkle.HttpServer.Context;
type type
TWsServerModule = class(TDataModule) TWsServerModule = class(TDataModule)
SparkleHttpSysDispatcher3: TSparkleHttpSysDispatcher; procedure DataModuleCreate(Sender: TObject);
XDataServer1: TXDataServer; procedure DataModuleDestroy(Sender: TObject);
XDataServer1Logging: TSparkleGenericMiddleware;
XDataServer1CORS: TSparkleCorsMiddleware;
procedure XDataServer1LoggingMiddlewareCreate(Sender: TObject;
var Middleware: IHttpServerMiddleware);
private private
{ Private declarations } FServer: TWebSocketServer;
procedure DoConnect(AContext: TIdContext);
procedure DoDisconnect(AContext: TIdContext);
procedure DoExecute(AContext: TIdContext);
function PeerId(AContext: TIdContext): string;
public public
{ Public declarations }
procedure StartWsServer(ABaseUrl: string; AModelName: string); procedure StartWsServer(ABaseUrl: string; AModelName: string);
procedure Broadcast(const AMessage: string);
property Server: TWebSocketServer read FServer;
end; end;
const const
SERVER_PATH_SEGMENT = 'ws'; WS_PORT = 2008;
var var
WsServerModule: TWsServerModule; WsServerModule: TWsServerModule;
...@@ -38,87 +32,97 @@ var ...@@ -38,87 +32,97 @@ var
implementation implementation
uses uses
Sparkle.HttpServer.Request, Common.Logging;
Sparkle.Middleware.Cors,
Sparkle.Middleware.Compress,
XData.OpenApi.Service,
XData.Sys.Exceptions,
Common.Logging,
Common.Middleware.Logging,
Common.Config, Vcl.Forms, IniFiles,
System.Rtti,
Ws.Service,
Ws.ServiceImpl;
{%CLASSGROUP 'Vcl.Controls.TControl'}
{$R *.dfm} {$R *.dfm}
{ TWsServerModule } function TWsServerModule.PeerId(AContext: TIdContext): string;
begin
try
Result := AContext.Binding.PeerIP + ':' + IntToStr(AContext.Binding.PeerPort);
except
Result := 'unknown';
end;
end;
procedure TWsServerModule.DataModuleCreate(Sender: TObject);
begin
FServer := TWebSocketServer.Create;
FServer.OnConnect := DoConnect;
FServer.OnDisconnect := DoDisconnect;
FServer.OnExecute := DoExecute;
end;
procedure TWsServerModule.DataModuleDestroy(Sender: TObject);
begin
if Assigned(FServer) then
begin
FServer.Active := False;
FServer.Free;
FServer := nil;
end;
end;
procedure TWsServerModule.DoConnect(AContext: TIdContext);
begin
Logger.Log(1, 'WS: Client connected [' + PeerId(AContext) + ']');
end;
procedure TWsServerModule.DoDisconnect(AContext: TIdContext);
begin
Logger.Log(1, 'WS: Client disconnected [' + PeerId(AContext) + ']');
end;
procedure TWsServerModule.DoExecute(AContext: TIdContext);
var
io: TWebSocketIOHandlerHelper;
msg: string;
begin
// Tag = 1 means the WebSocket handshake has completed; skip until then
if AContext.Connection.IOHandler.Tag <> 1 then
Exit;
io := TWebSocketIOHandlerHelper(AContext.Connection.IOHandler);
io.CheckForDataOnSource(10);
msg := io.ReadString;
if msg = '' then
Exit;
Logger.Log(1, 'WS: Message from [' + PeerId(AContext) + ']: ' + msg);
// Dispatch or handle incoming message here
end;
procedure TWsServerModule.StartWsServer(ABaseUrl: string; AModelName: string); procedure TWsServerModule.StartWsServer(ABaseUrl: string; AModelName: string);
begin
FServer.DefaultPort := WS_PORT;
FServer.Active := True;
Logger.Log(1, Format('WebSocket server listening on ws://0.0.0.0:%d/', [WS_PORT]));
end;
procedure TWsServerModule.Broadcast(const AMessage: string);
var var
Url: string; List: TList;
ctx: TRttiContext; i: Integer;
t: TRttiType; ctx: TIdContext;
attr: TCustomAttribute;
m: TRttiMethod;
s: string;
begin begin
Logger.Log(1, 'WS-DIAG: StartWsServer enter, AModelName=[' + AModelName + ']'); if not (Assigned(FServer) and FServer.Active) then
Logger.Log(1, 'WS-DIAG: XDataServer1.ModelName before=[' + XDataServer1.ModelName + ']'); Exit;
ctx := TRttiContext.Create; List := FServer.Contexts.LockList;
try try
t := ctx.GetType(TypeInfo(IWebSocketService)); for i := 0 to List.Count - 1 do
if t = nil then
Logger.Log(1, 'WS-DIAG: IWebSocketService RTTI=NIL')
else
begin
Logger.Log(1, 'WS-DIAG: IWebSocketService RTTI.Name=[' + t.Name + ']');
s := '';
for attr in t.GetAttributes do
s := s + attr.ClassName + ' ';
Logger.Log(1, 'WS-DIAG: IWebSocketService attrs=[' + s + ']');
s := '';
for m in t.GetMethods do
begin begin
s := s + m.Name + '('; ctx := TIdContext(List[i]);
var ma := ''; if ctx.Connection.Connected and (ctx.Connection.IOHandler.Tag = 1) then
for attr in m.GetAttributes do try
ma := ma + attr.ClassName + ' '; TWebSocketIOHandlerHelper(ctx.Connection.IOHandler).WriteString(AMessage);
s := s + ma + ') '; except
// ignore dead connections
end; end;
Logger.Log(1, 'WS-DIAG: IWebSocketService methods=[' + s + ']');
end; end;
finally finally
ctx.Free; FServer.Contexts.UnlockList;
end;
RegisterOpenApiService;
Logger.Log(1, 'WS-DIAG: RegisterOpenApiService done');
Url := ABaseUrl;
if not Url.EndsWith('/') then
Url := Url + '/';
Url := Url + SERVER_PATH_SEGMENT;
XDataServer1.BaseUrl := Url;
Logger.Log(1, 'WS-DIAG: about to set ModelName to [' + AModelName + ']');
try
XDataServer1.ModelName := AModelName;
Logger.Log(1, 'WS-DIAG: ModelName set, current value=[' + XDataServer1.ModelName + ']');
except
on E: Exception do
Logger.Log(1, 'WS-DIAG: ModelName assignment FAILED: ' + E.ClassName + ': ' + E.Message);
end; end;
SparkleHttpSysDispatcher3.Start;
Logger.Log(1, Format('Ws server module listening at "%s"', [Url]));
end;
procedure TWsServerModule.XDataServer1LoggingMiddlewareCreate(Sender: TObject;
var Middleware: IHttpServerMiddleware);
begin
Middleware := TLoggingMiddleware.Create(Logger, 1);
end; end;
end. end.
...@@ -2,20 +2,8 @@ unit Ws.Service; ...@@ -2,20 +2,8 @@ unit Ws.Service;
interface interface
uses
System.JSON,
XData.Service.Common;
const const
WS_MODEL = 'WsApi'; WS_MODEL = 'WsApi'; // retained for call-site compatibility in Main.pas
type
[ServiceContract, Model(WS_MODEL)]
IWebSocketService = interface(IInvokable)
['{673FE678-D9EF-468D-89CB-CEF26E8758BC}']
[HttpGet] function Ping: TJSONObject;
[HttpGet] function WebSockerConnectionHandler: TJSONObject;
end;
implementation implementation
......
unit Ws.ServiceImpl; unit Ws.ServiceImpl;
interface // Stub retained for project compatibility. Implementation moved to Ws.Server.Module.
uses
XData.Service.Common,
XData.Server.Module,
System.JSON,
Ws.Service;
type interface
[ServiceImplementation]
TWebSocketService = class(TInterfacedObject, IWebSocketService)
public
function Ping: TJSONObject;
function WebSockerConnectionHandler: TJSONObject;
end;
implementation implementation
function TWebSocketService.Ping: TJSONObject;
begin
Result := nil;
end;
function TWebSocketService.WebSockerConnectionHandler: TJSONObject;
begin
Result := nil;
end;
initialization
RegisterServiceType(TWebSocketService);
end. end.
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment